Initial commit.
Audit / Dependencies (push) Successful in 8s
CI / Formatting (push) Successful in 4s
CI / Linting (push) Successful in 5s
CI / Tests (Python 3.12) (push) Successful in 9s
CI / Tests (Python 3.13) (push) Successful in 12s
CI / Tests (Python 3.14) (push) Successful in 10s
CI / Type Checking (push) Successful in 10s
CI / Spelling (push) Successful in 7s

This commit is contained in:
2026-05-27 10:32:49 -04:00
commit 5d113eecaa
13 changed files with 3087 additions and 0 deletions
+38
View File
@@ -0,0 +1,38 @@
name: Audit
on:
schedule:
- cron: "0 0 * * 1"
push:
branches: [master]
paths: [uv.lock]
pull_request:
paths: [uv.lock]
permissions:
contents: read
env:
UV_PYTHON_DOWNLOADS: never
jobs:
audit:
name: Dependencies
runs-on: logaldeveloper-archlinux-ci
steps:
- name: Checkout repository
uses: actions/checkout@de0fac2e4500dabe0009e67214ff5f5447ce83dd # v6.0.2
with:
fetch-depth: 0
- name: Cache uv packages
uses: actions/cache@27d5ce7f107fe9357f9df03efb73ab90386fccae # v5.0.5
with:
path: ~/.cache/uv
key: uv-${{ hashFiles('uv.lock') }}
- name: Install dependencies
run: uv sync --locked
- name: Audit dependencies with pip-audit
run: uv run pip-audit --skip-editable
+134
View File
@@ -0,0 +1,134 @@
name: CI
on:
push:
branches: [master]
pull_request:
permissions:
contents: read
env:
UV_PYTHON_DOWNLOADS: never
jobs:
formatting:
name: Formatting
runs-on: logaldeveloper-archlinux-ci
steps:
- name: Checkout repository
uses: actions/checkout@de0fac2e4500dabe0009e67214ff5f5447ce83dd # v6.0.2
with:
fetch-depth: 0
- name: Cache uv packages
uses: actions/cache@27d5ce7f107fe9357f9df03efb73ab90386fccae # v5.0.5
with:
path: ~/.cache/uv
key: uv-${{ hashFiles('uv.lock') }}
- name: Install dependencies
run: uv sync --locked
- name: Check formatting with Ruff
run: uv run ruff format --check --diff .
linting:
name: Linting
runs-on: logaldeveloper-archlinux-ci
steps:
- name: Checkout repository
uses: actions/checkout@de0fac2e4500dabe0009e67214ff5f5447ce83dd # v6.0.2
with:
fetch-depth: 0
- name: Cache uv packages
uses: actions/cache@27d5ce7f107fe9357f9df03efb73ab90386fccae # v5.0.5
with:
path: ~/.cache/uv
key: uv-${{ hashFiles('uv.lock') }}
- name: Install dependencies
run: uv sync --locked
- name: Check linting with Ruff
run: uv run ruff check .
tests:
name: Tests (Python ${{ matrix.python-version }})
runs-on: logaldeveloper-archlinux-ci
strategy:
matrix:
python-version: ["3.12", "3.13", "3.14"]
env:
UV_PYTHON: ${{ matrix.python-version }}
UV_PYTHON_PREFERENCE: system
UV_PYTHON_DOWNLOADS: automatic
steps:
- name: Checkout repository
uses: actions/checkout@de0fac2e4500dabe0009e67214ff5f5447ce83dd # v6.0.2
with:
fetch-depth: 0
- name: Cache uv managed Python installs
uses: actions/cache@27d5ce7f107fe9357f9df03efb73ab90386fccae # v5.0.5
with:
path: ~/.local/share/uv/python
key: uv-python-${{ matrix.python-version }}-${{ runner.os }}-${{ runner.arch }}
- name: Cache uv packages
uses: actions/cache@27d5ce7f107fe9357f9df03efb73ab90386fccae # v5.0.5
with:
path: ~/.cache/uv
key: uv-${{ matrix.python-version }}-${{ hashFiles('uv.lock') }}
- name: Install dependencies
run: uv sync --locked
- name: Run unit tests with Pytest
run: uv run pytest -v --cov --cov-report=
- name: Report code coverage
run: uv run coverage report
type-checking:
name: Type Checking
runs-on: logaldeveloper-archlinux-ci
steps:
- name: Checkout repository
uses: actions/checkout@de0fac2e4500dabe0009e67214ff5f5447ce83dd # v6.0.2
with:
fetch-depth: 0
- name: Cache uv packages
uses: actions/cache@27d5ce7f107fe9357f9df03efb73ab90386fccae # v5.0.5
with:
path: ~/.cache/uv
key: uv-${{ hashFiles('uv.lock') }}
- name: Install dependencies
run: uv sync --locked
- name: Check types with Mypy
run: uv run mypy .
spelling:
name: Spelling
runs-on: logaldeveloper-archlinux-ci
steps:
- name: Checkout repository
uses: actions/checkout@de0fac2e4500dabe0009e67214ff5f5447ce83dd # v6.0.2
with:
fetch-depth: 0
- name: Cache uv packages
uses: actions/cache@27d5ce7f107fe9357f9df03efb73ab90386fccae # v5.0.5
with:
path: ~/.cache/uv
key: uv-${{ hashFiles('uv.lock') }}
- name: Install dependencies
run: uv sync --locked
- name: Check spelling with codespell
run: uv run codespell
+36
View File
@@ -0,0 +1,36 @@
# Byte-compiled / optimized / DLL files
__pycache__/
*.py[cod]
*$py.class
# Distribution / packaging
dist/
build/
*.egg-info/
*.egg
# Virtual environments
venv/
.venv/
env/
# IDE
.idea/
.vscode/
*.swp
*.swo
# Runtime data
data/
logs/
# Test / coverage artifacts
.coverage
coverage.xml
htmlcov/
.pytest_cache/
.mypy_cache/
.ruff_cache/
# Local configuration
config.yaml
+1
View File
@@ -0,0 +1 @@
3.14
+202
View File
@@ -0,0 +1,202 @@
Apache License
Version 2.0, January 2004
http://www.apache.org/licenses/
TERMS AND CONDITIONS FOR USE, REPRODUCTION, AND DISTRIBUTION
1. Definitions.
"License" shall mean the terms and conditions for use, reproduction,
and distribution as defined by Sections 1 through 9 of this document.
"Licensor" shall mean the copyright owner or entity authorized by
the copyright owner that is granting the License.
"Legal Entity" shall mean the union of the acting entity and all
other entities that control, are controlled by, or are under common
control with that entity. For the purposes of this definition,
"control" means (i) the power, direct or indirect, to cause the
direction or management of such entity, whether by contract or
otherwise, or (ii) ownership of fifty percent (50%) or more of the
outstanding shares, or (iii) beneficial ownership of such entity.
"You" (or "Your") shall mean an individual or Legal Entity
exercising permissions granted by this License.
"Source" form shall mean the preferred form for making modifications,
including but not limited to software source code, documentation
source, and configuration files.
"Object" form shall mean any form resulting from mechanical
transformation or translation of a Source form, including but
not limited to compiled object code, generated documentation,
and conversions to other media types.
"Work" shall mean the work of authorship, whether in Source or
Object form, made available under the License, as indicated by a
copyright notice that is included in or attached to the work
(an example is provided in the Appendix below).
"Derivative Works" shall mean any work, whether in Source or Object
form, that is based on (or derived from) the Work and for which the
editorial revisions, annotations, elaborations, or other modifications
represent, as a whole, an original work of authorship. For the purposes
of this License, Derivative Works shall not include works that remain
separable from, or merely link (or bind by name) to the interfaces of,
the Work and Derivative Works thereof.
"Contribution" shall mean any work of authorship, including
the original version of the Work and any modifications or additions
to that Work or Derivative Works thereof, that is intentionally
submitted to Licensor for inclusion in the Work by the copyright owner
or by an individual or Legal Entity authorized to submit on behalf of
the copyright owner. For the purposes of this definition, "submitted"
means any form of electronic, verbal, or written communication sent
to the Licensor or its representatives, including but not limited to
communication on electronic mailing lists, source code control systems,
and issue tracking systems that are managed by, or on behalf of, the
Licensor for the purpose of discussing and improving the Work, but
excluding communication that is conspicuously marked or otherwise
designated in writing by the copyright owner as "Not a Contribution."
"Contributor" shall mean Licensor and any individual or Legal Entity
on behalf of whom a Contribution has been received by Licensor and
subsequently incorporated within the Work.
2. Grant of Copyright License. Subject to the terms and conditions of
this License, each Contributor hereby grants to You a perpetual,
worldwide, non-exclusive, no-charge, royalty-free, irrevocable
copyright license to reproduce, prepare Derivative Works of,
publicly display, publicly perform, sublicense, and distribute the
Work and such Derivative Works in Source or Object form.
3. Grant of Patent License. Subject to the terms and conditions of
this License, each Contributor hereby grants to You a perpetual,
worldwide, non-exclusive, no-charge, royalty-free, irrevocable
(except as stated in this section) patent license to make, have made,
use, offer to sell, sell, import, and otherwise transfer the Work,
where such license applies only to those patent claims licensable
by such Contributor that are necessarily infringed by their
Contribution(s) alone or by combination of their Contribution(s)
with the Work to which such Contribution(s) was submitted. If You
institute patent litigation against any entity (including a
cross-claim or counterclaim in a lawsuit) alleging that the Work
or a Contribution incorporated within the Work constitutes direct
or contributory patent infringement, then any patent licenses
granted to You under this License for that Work shall terminate
as of the date such litigation is filed.
4. Redistribution. You may reproduce and distribute copies of the
Work or Derivative Works thereof in any medium, with or without
modifications, and in Source or Object form, provided that You
meet the following conditions:
(a) You must give any other recipients of the Work or
Derivative Works a copy of this License; and
(b) You must cause any modified files to carry prominent notices
stating that You changed the files; and
(c) You must retain, in the Source form of any Derivative Works
that You distribute, all copyright, patent, trademark, and
attribution notices from the Source form of the Work,
excluding those notices that do not pertain to any part of
the Derivative Works; and
(d) If the Work includes a "NOTICE" text file as part of its
distribution, then any Derivative Works that You distribute must
include a readable copy of the attribution notices contained
within such NOTICE file, excluding those notices that do not
pertain to any part of the Derivative Works, in at least one
of the following places: within a NOTICE text file distributed
as part of the Derivative Works; within the Source form or
documentation, if provided along with the Derivative Works; or,
within a display generated by the Derivative Works, if and
wherever such third-party notices normally appear. The contents
of the NOTICE file are for informational purposes only and
do not modify the License. You may add Your own attribution
notices within Derivative Works that You distribute, alongside
or as an addendum to the NOTICE text from the Work, provided
that such additional attribution notices cannot be construed
as modifying the License.
You may add Your own copyright statement to Your modifications and
may provide additional or different license terms and conditions
for use, reproduction, or distribution of Your modifications, or
for any such Derivative Works as a whole, provided Your use,
reproduction, and distribution of the Work otherwise complies with
the conditions stated in this License.
5. Submission of Contributions. Unless You explicitly state otherwise,
any Contribution intentionally submitted for inclusion in the Work
by You to the Licensor shall be under the terms and conditions of
this License, without any additional terms or conditions.
Notwithstanding the above, nothing herein shall supersede or modify
the terms of any separate license agreement you may have executed
with Licensor regarding such Contributions.
6. Trademarks. This License does not grant permission to use the trade
names, trademarks, service marks, or product names of the Licensor,
except as required for reasonable and customary use in describing the
origin of the Work and reproducing the content of the NOTICE file.
7. Disclaimer of Warranty. Unless required by applicable law or
agreed to in writing, Licensor provides the Work (and each
Contributor provides its Contributions) on an "AS IS" BASIS,
WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or
implied, including, without limitation, any warranties or conditions
of TITLE, NON-INFRINGEMENT, MERCHANTABILITY, or FITNESS FOR A
PARTICULAR PURPOSE. You are solely responsible for determining the
appropriateness of using or redistributing the Work and assume any
risks associated with Your exercise of permissions under this License.
8. Limitation of Liability. In no event and under no legal theory,
whether in tort (including negligence), contract, or otherwise,
unless required by applicable law (such as deliberate and grossly
negligent acts) or agreed to in writing, shall any Contributor be
liable to You for damages, including any direct, indirect, special,
incidental, or consequential damages of any character arising as a
result of this License or out of the use or inability to use the
Work (including but not limited to damages for loss of goodwill,
work stoppage, computer failure or malfunction, or any and all
other commercial damages or losses), even if such Contributor
has been advised of the possibility of such damages.
9. Accepting Warranty or Additional Liability. While redistributing
the Work or Derivative Works thereof, You may choose to offer,
and charge a fee for, acceptance of support, warranty, indemnity,
or other liability obligations and/or rights consistent with this
License. However, in accepting such obligations, You may act only
on Your own behalf and on Your sole responsibility, not on behalf
of any other Contributor, and only if You agree to indemnify,
defend, and hold each Contributor harmless for any liability
incurred by, or claims asserted against, such Contributor by reason
of your accepting any such warranty or additional liability.
END OF TERMS AND CONDITIONS
APPENDIX: How to apply the Apache License to your work.
To apply the Apache License to your work, attach the following
boilerplate notice, with the fields enclosed by brackets "[]"
replaced with your own identifying information. (Don't include
the brackets!) The text should be enclosed in the appropriate
comment syntax for the file format. We also recommend that a
file or class name and description of purpose be included on the
same "printed page" as the copyright notice for easier
identification within third-party archives.
Copyright [yyyy] [name of copyright owner]
Licensed under the Apache License, Version 2.0 (the "License");
you may not use this file except in compliance with the License.
You may obtain a copy of the License at
http://www.apache.org/licenses/LICENSE-2.0
Unless required by applicable law or agreed to in writing, software
distributed under the License is distributed on an "AS IS" BASIS,
WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
See the License for the specific language governing permissions and
limitations under the License.
+42
View File
@@ -0,0 +1,42 @@
# Extra Owlbot Modules
This repository contains standalone Owlbot modules that are useful but not
built into upstream Owlbot.
It also acts as a reference project for developing modules outside the main
Owlbot repository. Each module is a top-level Python package directory,
dependencies are declared in `pyproject.toml`, the environment is managed with
`uv`, and CI runs the same style, type, test, and spelling checks used by
Owlbot.
## Modules
- [`chat_archive`](chat_archive/): archives Owncast chat activity and stream
context updates to JSON files.
Each module package is kept at the repository root so it can be copied or
symlinked independently into the directory Owlbot scans for user modules.
## Loading Modules
Place the module folders you want to use inside your Owlbot user modules
directory:
```text
modules/
chat_archive/
__init__.py
README.md
```
If you do not already have a user modules directory, create one and configure
Owlbot to scan it with `owlbot.modules_dir`:
```yaml
owlbot:
modules_dir: "/path/to/modules"
```
Module-specific configuration belongs under `modules.<module_name>` in
`config.yaml`. See the README inside each module package for its config and
runtime behavior.
+104
View File
@@ -0,0 +1,104 @@
# chat_archive
This package contains an Owlbot module that writes Owncast chat activity
and stream context updates to JSON files for later processing.
## Behavior
An archive is the JSON file this module writes for one live stream. The module
starts an archive when Owlbot sees a live stream, including streams that are
already live when the module loads.
While an archive is active, the module records chat messages, user joins, user
parts, name changes, stream title updates, and message visibility updates in the
order Owlbot dispatches them. Events outside an active stream are ignored.
Archive files are stream-scoped. When a stream ends, the active archive is
finalized, and the next live stream starts a new archive file. Module unload
also finalizes the active archive.
## Configuration
```yaml
modules:
chat_archive:
enabled: true
archive_dir: "data/chat_archives"
pretty_json: true
```
- `archive_dir`: directory where archive files are written. The module creates
it if needed.
- `pretty_json`: when true, write indented JSON. When false, write compact JSON.
## Archive Format
Archive filenames are based on the detected stream start time when available:
```text
chat-archive-20260526T120000Z.json
```
If a file already exists for that timestamp, the module adds a suffix such as
`chat-archive-20260526T120000Z-2.json`.
Each archive is a single JSON document with these top-level fields:
- `archive_start_reason`: one of `module_startup` or `stream_started_event`.
- `archive_started_at`: ISO 8601 timestamp for when the archive file was
started.
- `archive_updated_at`: ISO 8601 timestamp for the most recent archive write.
- `archive_stopped_at`: ISO 8601 timestamp for when the archive was finalized,
or `null` while the archive is active.
- `archive_stop_reason`: reason the archive was finalized, or `null` while the
archive is active. Current values are `stream_stopped_event`,
`module_teardown`, and `replaced_by_stream_started_event`.
- `stream_title`: stream title when the archive started, or `null` if
unavailable.
- `stream_stopped_at`: stream stop timestamp recorded when finalized by a
stream stop event, or `null` otherwise.
- `event_count`: number of objects in `events`.
- `events`: ordered list of captured activity events.
Each object in `events` contains:
- `event_type`: one of `CHAT`, `USER_JOINED`, `USER_PARTED`, `NAME_CHANGE`,
`STREAM_TITLE_UPDATED`, or `VISIBILITY-UPDATE`.
- `archived_at`: ISO 8601 timestamp for when the event was written.
- `event_data`: parsed Owlbot event data converted to JSON-compatible values.
The exact fields depend on `event_type`.
Abridged example archive:
```json
{
"archive_start_reason": "stream_started_event",
"archive_started_at": "2026-05-26T12:00:01.000000+00:00",
"archive_updated_at": "2026-05-26T12:01:00.000000+00:00",
"archive_stopped_at": null,
"archive_stop_reason": null,
"stream_title": "Game Night",
"stream_stopped_at": null,
"event_count": 1,
"events": [
{
"event_type": "CHAT",
"archived_at": "2026-05-26T12:01:00.000000+00:00",
"event_data": {
"message_id": "msg-1",
"raw_body": "hello chat",
"body": "hello chat",
"is_visible": true
}
}
]
}
```
Owlbot modules receive parsed event dataclasses rather than the untouched
Owncast webhook payload. `event_data` stores that parsed event data after
converting values such as datetimes and sets into JSON-compatible values.
The archive is rewritten after each captured event through a temporary sibling
file and atomic replace. For an existing archive file, readers should see either
the previous complete JSON document or the next complete JSON document.
+332
View File
@@ -0,0 +1,332 @@
# Copyright 2026 Logan Fick
#
# Licensed under the Apache License, Version 2.0 (the "License");
# you may not use this file except in compliance with the License.
# You may obtain a copy of the License at
#
# http://www.apache.org/licenses/LICENSE-2.0
#
# Unless required by applicable law or agreed to in writing, software
# distributed under the License is distributed on an "AS IS" BASIS,
# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
# See the License for the specific language governing permissions and
# limitations under the License.
"""Archive Owncast chat activity to JSON files."""
from __future__ import annotations
import asyncio
from dataclasses import asdict, dataclass, field
from datetime import UTC, datetime
from pathlib import Path
from typing import Any, Literal, cast
import orjson
from owlbot.api import (
ChatEvent,
Event,
EventContext,
EventType,
ModuleContext,
NameChangedEvent,
Priority,
StreamStartedEvent,
StreamStoppedEvent,
StreamTitleUpdatedEvent,
UserJoinedEvent,
UserPartedEvent,
VisibilityUpdateEvent,
on_event,
on_setup,
on_teardown,
)
_STATE_KEY = "chat_archive_state"
type ArchiveStartReason = Literal["module_startup", "stream_started_event"]
type ArchiveStopReason = Literal[
"module_teardown",
"replaced_by_stream_started_event",
"stream_stopped_event",
]
@dataclass(slots=True)
class ArchiveState:
"""Runtime state for the chat archive module."""
archive_dir: Path
pretty_json: bool
lock: asyncio.Lock = field(default_factory=asyncio.Lock)
current_path: Path | None = None
current_data: dict[str, Any] | None = None
def _get_state(ctx: ModuleContext) -> ArchiveState:
"""Return the typed ArchiveState stored in Owlbot's ctx.state dict."""
state = ctx.state.get(_STATE_KEY)
if not isinstance(state, ArchiveState):
msg = "chat_archive state is not initialized"
raise RuntimeError(msg) # noqa: TRY004 # state error, not a type error
return state
@on_setup
async def setup(ctx: ModuleContext) -> None:
"""Initialize config/runtime state and start an archive if the stream is live."""
ctx.config.register_defaults(
{
"archive_dir": "data/chat_archives",
"pretty_json": True,
}
)
archive_dir_value = ctx.config.get("archive_dir", "data/chat_archives")
pretty_json = bool(ctx.config.get("pretty_json", True))
state = ArchiveState(
archive_dir=Path(str(archive_dir_value)),
pretty_json=pretty_json,
)
state.archive_dir.mkdir(parents=True, exist_ok=True)
ctx.state[_STATE_KEY] = state
try:
status = await ctx.owncast_client.get_status()
except Exception: # noqa: BLE001 # best-effort startup check; stream events still work
ctx.logger.warning(
"Could not check stream status. Chat archiving will start on stream start.",
exc_info=True,
)
return
if status.get("online", False):
last_connect_time = status.get("lastConnectTime")
started_at = datetime.now(UTC)
if isinstance(last_connect_time, str) and last_connect_time:
try:
started_at = datetime.fromisoformat(last_connect_time)
except ValueError:
ctx.logger.warning(
"Owncast status returned lastConnectTime in unexpected "
"format: %r. Falling back to the current time.",
last_connect_time,
)
async with state.lock:
path = await _start_archive(
ctx,
state,
archive_start_reason="module_startup",
stream_started_at=started_at,
stream_title=status.get("streamTitle"),
)
ctx.logger.info("Stream is already live. Chat archive started: %s", path)
else:
ctx.logger.info("Stream is offline. Chat archiving will start on stream start.")
@on_event(EventType.STREAM_STARTED)
async def handle_stream_started(ctx: EventContext[StreamStartedEvent]) -> None:
"""Create a new archive when Owncast reports that the stream went live."""
state = _get_state(ctx.module)
started_at = ctx.event.timestamp or datetime.now(UTC)
async with state.lock:
path = await _start_archive(
ctx.module,
state,
archive_start_reason="stream_started_event",
stream_started_at=started_at,
stream_title=ctx.event.stream_title,
)
ctx.logger.info("Chat archive started: %s", path)
@on_event(EventType.CHAT, priority=Priority.HIGHEST)
async def archive_chat_message(ctx: EventContext[ChatEvent]) -> None:
"""Append an incoming chat message to the active archive file."""
await _archive_event(ctx.module, EventType.CHAT, ctx.event)
@on_event(EventType.VISIBILITY_UPDATE, priority=Priority.HIGHEST)
async def archive_visibility_update(
ctx: EventContext[VisibilityUpdateEvent],
) -> None:
"""Append a message visibility update to the active archive file."""
await _archive_event(ctx.module, EventType.VISIBILITY_UPDATE, ctx.event)
@on_event(EventType.USER_JOINED, priority=Priority.HIGHEST)
async def archive_user_joined(ctx: EventContext[UserJoinedEvent]) -> None:
"""Append a user join event to the active archive file."""
await _archive_event(ctx.module, EventType.USER_JOINED, ctx.event)
@on_event(EventType.USER_PARTED, priority=Priority.HIGHEST)
async def archive_user_parted(ctx: EventContext[UserPartedEvent]) -> None:
"""Append a user part event to the active archive file."""
await _archive_event(ctx.module, EventType.USER_PARTED, ctx.event)
@on_event(EventType.NAME_CHANGE, priority=Priority.HIGHEST)
async def archive_name_change(ctx: EventContext[NameChangedEvent]) -> None:
"""Append a name change event to the active archive file."""
await _archive_event(ctx.module, EventType.NAME_CHANGE, ctx.event)
@on_event(EventType.STREAM_TITLE_UPDATED, priority=Priority.HIGHEST)
async def archive_stream_title_update(
ctx: EventContext[StreamTitleUpdatedEvent],
) -> None:
"""Append a stream title update event to the active archive file."""
await _archive_event(ctx.module, EventType.STREAM_TITLE_UPDATED, ctx.event)
@on_event(EventType.STREAM_STOPPED)
async def handle_stream_stopped(ctx: EventContext[StreamStoppedEvent]) -> None:
"""Finalize the active archive when Owncast reports that the stream stopped."""
state = _get_state(ctx.module)
async with state.lock:
path = await _finalize_archive(
ctx.module,
state,
reason="stream_stopped_event",
stream_stopped_at=ctx.event.timestamp or datetime.now(UTC),
)
if path is not None:
ctx.logger.info("Chat archive finalized after stream stopped: %s", path)
@on_teardown
async def teardown(ctx: ModuleContext) -> None:
"""Finalize any active archive when Owlbot unloads the module."""
state = _get_state(ctx)
async with state.lock:
await _finalize_archive(ctx, state, reason="module_teardown")
async def _start_archive(
ctx: ModuleContext,
state: ArchiveState,
*,
archive_start_reason: ArchiveStartReason,
stream_started_at: datetime,
stream_title: str | None,
) -> Path:
"""Open a new archive file, finalizing any active file first."""
if state.current_data is not None:
await _finalize_archive(ctx, state, reason="replaced_by_stream_started_event")
archive_started_at = datetime.now(UTC)
if stream_started_at.tzinfo is None:
stream_started_at = stream_started_at.replace(tzinfo=UTC)
timestamp = stream_started_at.astimezone(UTC).strftime("%Y%m%dT%H%M%SZ")
stem = f"chat-archive-{timestamp}"
path = state.archive_dir / f"{stem}.json"
index = 2
while path.exists():
path = state.archive_dir / f"{stem}-{index}.json"
index += 1
data: dict[str, Any] = {
"archive_start_reason": archive_start_reason,
"archive_started_at": archive_started_at.isoformat(),
"archive_updated_at": archive_started_at.isoformat(),
"archive_stopped_at": None,
"archive_stop_reason": None,
"stream_title": stream_title,
"stream_stopped_at": None,
"event_count": 0,
"events": [],
}
await _write_json(path, data, pretty_json=state.pretty_json)
state.current_path = path
state.current_data = data
return path
async def _archive_event(
ctx: ModuleContext,
event_type: EventType,
event: Event,
) -> None:
"""Record one event and rewrite the active archive, if a stream is live."""
state = _get_state(ctx)
async with state.lock:
data = state.current_data
path = state.current_path
if data is None or path is None:
return
archived_at = datetime.now(UTC)
events = cast("list[dict[str, Any]]", data["events"])
events.append(
{
"event_type": event_type.value,
"archived_at": archived_at.isoformat(),
"event_data": asdict(event),
},
)
data["event_count"] = len(events)
data["archive_updated_at"] = archived_at.isoformat()
await _write_json(path, data, pretty_json=state.pretty_json)
async def _finalize_archive(
ctx: ModuleContext,
state: ArchiveState,
*,
reason: ArchiveStopReason,
stream_stopped_at: datetime | None = None,
) -> Path | None:
"""Mark the active archive complete and write its final JSON state."""
data = state.current_data
path = state.current_path
if data is None or path is None:
return None
stopped_at = datetime.now(UTC)
data["archive_stopped_at"] = stopped_at.isoformat()
data["archive_stop_reason"] = reason
data["archive_updated_at"] = stopped_at.isoformat()
if stream_stopped_at is not None:
data["stream_stopped_at"] = stream_stopped_at.isoformat()
await _write_json(path, data, pretty_json=state.pretty_json)
ctx.logger.info("Wrote chat archive: %s", path)
state.current_path = None
state.current_data = None
return path
async def _write_json(path: Path, data: dict[str, Any], *, pretty_json: bool) -> None:
"""Atomically write a valid JSON document to disk without blocking the loop."""
def json_default(value: object) -> object:
if isinstance(value, (set, frozenset)):
return sorted(value, key=str)
raise TypeError(f"Type is not JSON serializable: {type(value).__name__}")
def write_sync() -> None:
path.parent.mkdir(parents=True, exist_ok=True)
option = orjson.OPT_INDENT_2 if pretty_json else 0
payload = orjson.dumps(data, option=option, default=json_default) + b"\n"
tmp_path = path.with_name(f".{path.name}.tmp")
tmp_path.write_bytes(payload)
tmp_path.replace(path)
await asyncio.to_thread(write_sync)
+83
View File
@@ -0,0 +1,83 @@
[project]
name = "owlbot-extra-modules"
version = "0.0.0"
description = "Standalone Owlbot modules that are useful but not built into upstream Owlbot."
requires-python = ">=3.12"
license = "Apache-2.0"
authors = [
{ name = "Logan Fick" },
]
dependencies = [
"orjson>=3.11.8",
"owlbot==1.2.0",
]
[project.urls]
Repository = "https://git.logal.dev/LogalDeveloper/Owlbot-Extra-Modules"
[dependency-groups]
dev = [
"codespell>=2.4.2",
"mypy>=1.20.2",
"owlbot[testing]==1.2.0",
"pip-audit>=2.10.0",
"pytest-cov>=7.1.0",
"ruff>=0.15.12",
]
[tool.uv]
package = false
extra-index-url = ["https://git.logal.dev/api/packages/LogalDeveloper/pypi/simple/"]
[tool.mypy]
python_version = "3.12"
strict = true
warn_unreachable = true
explicit_package_bases = true
[tool.pytest.ini_options]
asyncio_mode = "auto"
asyncio_default_fixture_loop_scope = "function"
pythonpath = ["."]
[tool.ruff]
target-version = "py312"
[tool.ruff.lint]
select = ["ALL"]
ignore = [
"ANN401", # Any is valid at system boundaries; mypy strict handles real issues
"C901", # McCabe complexity: noisy and not actionable
"COM812", # handled by the formatter
"D203", # incompatible with D211 (no blank line before class docstring)
"D213", # incompatible with D212 (summary on first line)
"EM", # exception message style: inline literals are fine
"PLR0911", # too many return statements: flat early-returns are clear
"PLR0912", # too many branches: inherent in parsers, validators, CLI
"PLR0913", # too many arguments: API surfaces and constructors need them
"PLR0915", # too many statements: inherent in parsers, validators, CLI
"TRY003", # inline exception messages are fine (complements EM ignore)
"TRY301", # raise inside try: guard clauses don't need helper functions
]
[tool.ruff.lint.per-file-ignores]
"tests/**" = [
"S101", # assert is standard for pytest
"SLF001", # tests legitimately inspect module state and test fixture internals
"ARG001", # unused args are normal for fixtures
"PLR2004", # magic values are clear in test assertions
]
[tool.codespell]
skip = "uv.lock"
ignore-words-list = "ECT,IST,PONT,TAHT,THA,WAST,WIT"
[tool.coverage.run]
source = ["."]
omit = [
"tests/**",
]
[tool.coverage.report]
show_missing = true
skip_empty = true
+1
View File
@@ -0,0 +1 @@
"""Extra Owlbot modules test suite."""
+3
View File
@@ -0,0 +1,3 @@
"""Pytest configuration for extra Owlbot modules."""
pytest_plugins = ["owlbot.testing.plugin"]
+529
View File
@@ -0,0 +1,529 @@
# Copyright 2026 Logan Fick
#
# Licensed under the Apache License, Version 2.0 (the "License");
# you may not use this file except in compliance with the License.
# You may obtain a copy of the License at
#
# http://www.apache.org/licenses/LICENSE-2.0
#
# Unless required by applicable law or agreed to in writing, software
# distributed under the License is distributed on an "AS IS" BASIS,
# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
# See the License for the specific language governing permissions and
# limitations under the License.
"""Tests for the chat_archive module."""
from __future__ import annotations
import asyncio
import logging
from datetime import UTC, datetime
from typing import TYPE_CHECKING, Any, cast
from unittest.mock import AsyncMock
import orjson
import pytest
from owlbot.api import EventType
from owlbot.testing import (
make_chat_event,
make_name_changed_event,
make_stream_started_event,
make_stream_stopped_event,
make_stream_title_updated_event,
make_user,
make_user_joined_event,
make_user_parted_event,
make_visibility_update_event,
)
import chat_archive
if TYPE_CHECKING:
from pathlib import Path
from types import ModuleType
from owlbot.api import ModuleContext
from owlbot.registries.events import EventDispatcher
from owlbot.testing import RecordingOwncastClient
def _config_data(archive_dir: Path, *, pretty_json: bool = True) -> dict[str, Any]:
"""Build plugin config data with chat_archive pointed at a temp directory."""
return {
"owncast": {"url": "http://localhost:8080"},
"owlbot": {"public_base_url": "http://localhost:8081"},
"modules": {
"chat_archive": {
"archive_dir": str(archive_dir),
"pretty_json": pretty_json,
},
},
}
@pytest.fixture
def module_pkg() -> ModuleType:
"""Auto-wire chat_archive handler registration and lifecycle."""
return chat_archive
@pytest.fixture
def archive_dir(tmp_path: Path) -> Path:
"""Return the per-test archive output directory."""
return tmp_path / "archives"
@pytest.fixture
def config_data(archive_dir: Path) -> dict[str, Any]:
"""Seed module config so archive writes stay inside tmp_path."""
return _config_data(archive_dir)
def _archive_files(archive_dir: Path) -> list[Path]:
"""Return archive JSON files in stable filename order."""
return sorted(archive_dir.glob("chat-archive-*.json"))
def _read_archive(path: Path) -> dict[str, Any]:
"""Read one archive JSON document."""
data = orjson.loads(path.read_bytes())
assert isinstance(data, dict)
return data
def _events(data: dict[str, Any]) -> list[dict[str, Any]]:
"""Return the typed events list from an archive document."""
events = data["events"]
assert isinstance(events, list)
return cast("list[dict[str, Any]]", events)
def _assert_iso_datetime(value: object) -> None:
"""Assert that a value is a timezone-aware ISO datetime string."""
assert isinstance(value, str)
parsed = datetime.fromisoformat(value)
assert parsed.tzinfo is not None
@pytest.mark.usefixtures("module_lifecycle")
class TestChatArchive:
"""End-to-end behavior through Owlbot's real event dispatch pipeline."""
async def test_setup_creates_archive_dir_when_stream_is_offline(
self,
archive_dir: Path,
) -> None:
"""Offline startup initializes the output directory without an archive."""
assert await asyncio.to_thread(archive_dir.is_dir)
assert _archive_files(archive_dir) == []
async def test_stream_started_creates_active_archive(
self,
archive_dir: Path,
event_dispatcher: EventDispatcher,
) -> None:
"""STREAM_STARTED creates an active stream-scoped archive file."""
started_at = datetime(2026, 5, 26, 12, tzinfo=UTC)
await event_dispatcher.dispatch(
EventType.STREAM_STARTED,
make_stream_started_event(
stream_title="Game Night",
timestamp=started_at,
),
)
files = _archive_files(archive_dir)
assert [path.name for path in files] == ["chat-archive-20260526T120000Z.json"]
data = _read_archive(files[0])
assert data["archive_start_reason"] == "stream_started_event"
assert data["archive_stop_reason"] is None
assert data["archive_stopped_at"] is None
assert data["stream_title"] == "Game Night"
assert data["stream_stopped_at"] is None
assert data["event_count"] == 0
assert data["events"] == []
_assert_iso_datetime(data["archive_started_at"])
_assert_iso_datetime(data["archive_updated_at"])
async def test_archives_events_while_stream_is_active(
self,
archive_dir: Path,
event_dispatcher: EventDispatcher,
) -> None:
"""Supported events append JSON-compatible snapshots in dispatch order."""
await event_dispatcher.dispatch(
EventType.STREAM_STARTED,
make_stream_started_event(
stream_title="Event Test",
timestamp=datetime(2026, 5, 26, 12, tzinfo=UTC),
),
)
user = make_user(
id="alice-id",
display_name="Alice",
is_authenticated=True,
is_moderator=True,
scopes=frozenset({"CUSTOM"}),
)
await event_dispatcher.dispatch(
EventType.CHAT,
make_chat_event(
user=user,
raw_body="hello chat",
body="<p>hello chat</p>",
client_id=7,
message_id="msg-1",
timestamp=datetime(2026, 5, 26, 12, 1, tzinfo=UTC),
),
)
await event_dispatcher.dispatch(
EventType.USER_JOINED,
make_user_joined_event(
user=make_user(id="bob-id", display_name="Bob"),
client_id=8,
event_id="join-1",
timestamp=datetime(2026, 5, 26, 12, 2, tzinfo=UTC),
),
)
await event_dispatcher.dispatch(
EventType.VISIBILITY_UPDATE,
make_visibility_update_event(
event_id="vis-1",
message_ids=["msg-1"],
is_visible=False,
timestamp=datetime(2026, 5, 26, 12, 3, tzinfo=UTC),
),
)
data = _read_archive(_archive_files(archive_dir)[0])
events = _events(data)
assert data["event_count"] == 3
assert [event["event_type"] for event in events] == [
EventType.CHAT.value,
EventType.USER_JOINED.value,
EventType.VISIBILITY_UPDATE.value,
]
chat_data = events[0]["event_data"]
assert chat_data["message_id"] == "msg-1"
assert chat_data["raw_body"] == "hello chat"
assert chat_data["body"] == "<p>hello chat</p>"
assert chat_data["timestamp"] == "2026-05-26T12:01:00+00:00"
assert chat_data["user"]["id"] == "alice-id"
assert chat_data["user"]["display_name"] == "Alice"
assert chat_data["user"]["scopes"] == ["CUSTOM", "MODERATOR"]
joined_data = events[1]["event_data"]
assert joined_data["event_id"] == "join-1"
assert joined_data["user"]["display_name"] == "Bob"
visibility_data = events[2]["event_data"]
assert visibility_data["message_ids"] == ["msg-1"]
assert visibility_data["is_visible"] is False
_assert_iso_datetime(events[2]["archived_at"])
assert data["archive_updated_at"] == events[2]["archived_at"]
async def test_archives_remaining_supported_event_wrappers(
self,
archive_dir: Path,
event_dispatcher: EventDispatcher,
) -> None:
"""User part, name change, and title update handlers all archive events."""
await event_dispatcher.dispatch(
EventType.STREAM_STARTED,
make_stream_started_event(
timestamp=datetime(2026, 5, 26, 12, tzinfo=UTC),
),
)
await event_dispatcher.dispatch(
EventType.USER_PARTED,
make_user_parted_event(
user=make_user(id="alice-id", display_name="Alice"),
event_id="part-1",
),
)
await event_dispatcher.dispatch(
EventType.NAME_CHANGE,
make_name_changed_event(
user=make_user(id="alice-id", display_name="AliceNew"),
new_name="AliceNew",
event_id="name-1",
),
)
await event_dispatcher.dispatch(
EventType.STREAM_TITLE_UPDATED,
make_stream_title_updated_event(stream_title="New Title"),
)
events = _events(_read_archive(_archive_files(archive_dir)[0]))
assert [event["event_type"] for event in events] == [
EventType.USER_PARTED.value,
EventType.NAME_CHANGE.value,
EventType.STREAM_TITLE_UPDATED.value,
]
assert events[0]["event_data"]["event_id"] == "part-1"
assert events[1]["event_data"]["new_name"] == "AliceNew"
assert events[2]["event_data"]["stream_title"] == "New Title"
async def test_naive_stream_start_timestamp_is_treated_as_utc(
self,
archive_dir: Path,
event_dispatcher: EventDispatcher,
) -> None:
"""Naive stream start timestamps are normalized as UTC for filenames."""
await event_dispatcher.dispatch(
EventType.STREAM_STARTED,
make_stream_started_event(
timestamp=datetime(2026, 5, 26, 12), # noqa: DTZ001
),
)
assert [path.name for path in _archive_files(archive_dir)] == [
"chat-archive-20260526T120000Z.json"
]
async def test_stream_stopped_finalizes_archive_and_ignores_later_events(
self,
archive_dir: Path,
event_dispatcher: EventDispatcher,
) -> None:
"""STREAM_STOPPED finalizes the active archive and closes capture."""
await event_dispatcher.dispatch(
EventType.STREAM_STARTED,
make_stream_started_event(
timestamp=datetime(2026, 5, 26, 12, tzinfo=UTC),
),
)
await event_dispatcher.dispatch(
EventType.CHAT,
make_chat_event(raw_body="before stop", message_id="msg-before"),
)
stopped_at = datetime(2026, 5, 26, 13, 30, tzinfo=UTC)
await event_dispatcher.dispatch(
EventType.STREAM_STOPPED,
make_stream_stopped_event(timestamp=stopped_at),
)
path = _archive_files(archive_dir)[0]
data = _read_archive(path)
assert data["archive_stop_reason"] == "stream_stopped_event"
assert data["stream_stopped_at"] == "2026-05-26T13:30:00+00:00"
assert data["event_count"] == 1
_assert_iso_datetime(data["archive_stopped_at"])
await event_dispatcher.dispatch(
EventType.CHAT,
make_chat_event(raw_body="after stop", message_id="msg-after"),
)
assert _read_archive(path)["event_count"] == 1
async def test_new_stream_replaces_active_archive_and_uses_filename_suffix(
self,
archive_dir: Path,
event_dispatcher: EventDispatcher,
) -> None:
"""A second stream start finalizes the first archive and opens a suffix file."""
started_at = datetime(2026, 5, 26, 12, tzinfo=UTC)
await event_dispatcher.dispatch(
EventType.STREAM_STARTED,
make_stream_started_event(stream_title="First", timestamp=started_at),
)
await event_dispatcher.dispatch(
EventType.STREAM_STARTED,
make_stream_started_event(stream_title="Second", timestamp=started_at),
)
await event_dispatcher.dispatch(
EventType.CHAT,
make_chat_event(raw_body="new stream message", message_id="msg-new"),
)
files = _archive_files(archive_dir)
assert {path.name for path in files} == {
"chat-archive-20260526T120000Z.json",
"chat-archive-20260526T120000Z-2.json",
}
old_archive = _read_archive(archive_dir / "chat-archive-20260526T120000Z.json")
assert old_archive["archive_stop_reason"] == "replaced_by_stream_started_event"
assert old_archive["event_count"] == 0
_assert_iso_datetime(old_archive["archive_stopped_at"])
new_archive = _read_archive(
archive_dir / "chat-archive-20260526T120000Z-2.json"
)
assert new_archive["archive_start_reason"] == "stream_started_event"
assert new_archive["archive_stop_reason"] is None
assert new_archive["stream_title"] == "Second"
assert new_archive["event_count"] == 1
assert _events(new_archive)[0]["event_data"]["message_id"] == "msg-new"
class TestSetupAndTeardown:
"""Cases that need direct control around setup and teardown timing."""
def test_get_state_raises_before_setup(self, module_context: ModuleContext) -> None:
"""Accessing module state before setup reports a clear runtime error."""
with pytest.raises(RuntimeError, match="chat_archive state is not initialized"):
chat_archive._get_state(module_context)
async def test_setup_continues_when_status_check_fails(
self,
archive_dir: Path,
module_context: ModuleContext,
owncast_client: RecordingOwncastClient,
monkeypatch: pytest.MonkeyPatch,
caplog: pytest.LogCaptureFixture,
) -> None:
"""A startup status failure leaves the module ready for future events."""
monkeypatch.setattr(
owncast_client,
"get_status",
AsyncMock(side_effect=RuntimeError("status failed")),
)
with caplog.at_level(logging.WARNING):
await chat_archive.setup(module_context)
try:
assert await asyncio.to_thread(archive_dir.is_dir)
assert _archive_files(archive_dir) == []
assert "Could not check stream status" in caplog.text
finally:
await chat_archive.teardown(module_context)
async def test_setup_starts_archive_when_stream_is_already_online(
self,
archive_dir: Path,
module_context: ModuleContext,
owncast_client: RecordingOwncastClient,
monkeypatch: pytest.MonkeyPatch,
) -> None:
"""Direct setup can observe a stubbed live Owncast status response."""
monkeypatch.setattr(
owncast_client,
"get_status",
AsyncMock(
return_value={
"online": True,
"lastConnectTime": "2026-05-26T12:00:00+00:00",
"streamTitle": "Already Live",
}
),
)
await chat_archive.setup(module_context)
try:
files = _archive_files(archive_dir)
assert [path.name for path in files] == [
"chat-archive-20260526T120000Z.json"
]
data = _read_archive(files[0])
assert data["archive_start_reason"] == "module_startup"
assert data["stream_title"] == "Already Live"
assert data["archive_stop_reason"] is None
finally:
await chat_archive.teardown(module_context)
async def test_setup_logs_bad_last_connect_time_and_starts_archive(
self,
archive_dir: Path,
module_context: ModuleContext,
owncast_client: RecordingOwncastClient,
monkeypatch: pytest.MonkeyPatch,
caplog: pytest.LogCaptureFixture,
) -> None:
"""Malformed lastConnectTime falls back to the current time."""
monkeypatch.setattr(
owncast_client,
"get_status",
AsyncMock(
return_value={
"online": True,
"lastConnectTime": "not-a-date",
"streamTitle": "Bad Timestamp",
}
),
)
with caplog.at_level(logging.WARNING):
await chat_archive.setup(module_context)
try:
files = _archive_files(archive_dir)
assert len(files) == 1
data = _read_archive(files[0])
assert data["archive_start_reason"] == "module_startup"
assert data["stream_title"] == "Bad Timestamp"
assert "lastConnectTime in unexpected format" in caplog.text
finally:
await chat_archive.teardown(module_context)
@pytest.mark.usefixtures("registered_handlers")
async def test_teardown_finalizes_active_archive(
self,
archive_dir: Path,
module_context: ModuleContext,
event_dispatcher: EventDispatcher,
) -> None:
"""Module teardown writes the final archive state for an active stream."""
await chat_archive.setup(module_context)
await event_dispatcher.dispatch(
EventType.STREAM_STARTED,
make_stream_started_event(
timestamp=datetime(2026, 5, 26, 12, tzinfo=UTC),
),
)
path = _archive_files(archive_dir)[0]
await chat_archive.teardown(module_context)
data = _read_archive(path)
assert data["archive_stop_reason"] == "module_teardown"
assert data["stream_stopped_at"] is None
_assert_iso_datetime(data["archive_stopped_at"])
class TestConfiguredOutput:
"""Configuration-dependent archive output behavior."""
@pytest.fixture
def config_data(self, archive_dir: Path) -> dict[str, Any]:
"""Disable pretty JSON for tests in this class."""
return _config_data(archive_dir, pretty_json=False)
@pytest.mark.usefixtures("module_lifecycle")
async def test_pretty_json_false_writes_compact_json(
self,
archive_dir: Path,
event_dispatcher: EventDispatcher,
) -> None:
"""pretty_json=False writes one compact JSON line plus a trailing newline."""
await event_dispatcher.dispatch(
EventType.STREAM_STARTED,
make_stream_started_event(
timestamp=datetime(2026, 5, 26, 12, tzinfo=UTC),
),
)
content = _archive_files(archive_dir)[0].read_bytes()
assert content.startswith(b'{"archive_start_reason"')
assert b'\n "' not in content
assert content.endswith(b"\n")
assert _read_archive(_archive_files(archive_dir)[0])["event_count"] == 0
async def test_write_json_rejects_unserializable_objects(
self,
tmp_path: Path,
) -> None:
"""The JSON writer surfaces unsupported value types."""
with pytest.raises(TypeError, match="Type is not JSON serializable: object"):
await chat_archive._write_json(
tmp_path / "bad.json",
{"bad": object()},
pretty_json=True,
)
Generated
+1582
View File
File diff suppressed because it is too large Load Diff