Delayed chat archive finalization after stream stops.
CI / Formatting (push) Successful in 4s
CI / Linting (push) Successful in 4s
CI / Tests (Python 3.12) (push) Successful in 7s
CI / Tests (Python 3.13) (push) Successful in 7s
CI / Tests (Python 3.14) (push) Successful in 6s
CI / Type Checking (push) Successful in 7s
CI / Spelling (push) Successful in 4s

This commit is contained in:
2026-05-27 11:19:12 -04:00
parent 5d113eecaa
commit c93ffc551d
3 changed files with 198 additions and 19 deletions
+11 -6
View File
@@ -13,9 +13,13 @@ 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 parts, name changes, stream title updates, and message visibility updates in the
order Owlbot dispatches them. Events outside an active stream are ignored. order Owlbot dispatches them. Events outside an active stream are ignored.
Archive files are stream-scoped. When a stream ends, the active archive is Archive files are stream-scoped. When a stream ends, the module keeps that
finalized, and the next live stream starts a new archive file. Module unload stream's archive open for a five-minute grace period so it can capture chat
also finalizes the active archive. that Owncast still accepts after the stream stops. If no new stream starts
during that window, the archive is finalized when the grace period expires. If a
new stream starts first, the stopped stream's archive is finalized immediately
and a new archive is opened for the new stream. Module unload also finalizes the
active archive.
## Configuration ## Configuration
@@ -51,11 +55,12 @@ Each archive is a single JSON document with these top-level fields:
- `archive_stopped_at`: ISO 8601 timestamp for when the archive was finalized, - `archive_stopped_at`: ISO 8601 timestamp for when the archive was finalized,
or `null` while the archive is active. or `null` while the archive is active.
- `archive_stop_reason`: reason the archive was finalized, or `null` while the - `archive_stop_reason`: reason the archive was finalized, or `null` while the
archive is active. Current values are `stream_stopped_event`, archive is active or waiting through the post-stream grace period. Current
`module_teardown`, and `replaced_by_stream_started_event`. values are `stream_stopped_event`, `module_teardown`, and
`replaced_by_stream_started_event`.
- `stream_title`: stream title when the archive started, or `null` if - `stream_title`: stream title when the archive started, or `null` if
unavailable. unavailable.
- `stream_stopped_at`: stream stop timestamp recorded when finalized by a - `stream_stopped_at`: stream stop timestamp recorded when Owlbot receives a
stream stop event, or `null` otherwise. stream stop event, or `null` otherwise.
- `event_count`: number of objects in `events`. - `event_count`: number of objects in `events`.
- `events`: ordered list of captured activity events. - `events`: ordered list of captured activity events.
+57 -7
View File
@@ -17,6 +17,7 @@
from __future__ import annotations from __future__ import annotations
import asyncio import asyncio
import contextlib
from dataclasses import asdict, dataclass, field from dataclasses import asdict, dataclass, field
from datetime import UTC, datetime from datetime import UTC, datetime
from pathlib import Path from pathlib import Path
@@ -43,6 +44,7 @@ from owlbot.api import (
) )
_STATE_KEY = "chat_archive_state" _STATE_KEY = "chat_archive_state"
_STOP_GRACE_SECONDS = 300.0
type ArchiveStartReason = Literal["module_startup", "stream_started_event"] type ArchiveStartReason = Literal["module_startup", "stream_started_event"]
type ArchiveStopReason = Literal[ type ArchiveStopReason = Literal[
@@ -61,6 +63,7 @@ class ArchiveState:
lock: asyncio.Lock = field(default_factory=asyncio.Lock) lock: asyncio.Lock = field(default_factory=asyncio.Lock)
current_path: Path | None = None current_path: Path | None = None
current_data: dict[str, Any] | None = None current_data: dict[str, Any] | None = None
stop_task: asyncio.Task[None] | None = None
def _get_state(ctx: ModuleContext) -> ArchiveState: def _get_state(ctx: ModuleContext) -> ArchiveState:
@@ -134,6 +137,7 @@ async def handle_stream_started(ctx: EventContext[StreamStartedEvent]) -> None:
started_at = ctx.event.timestamp or datetime.now(UTC) started_at = ctx.event.timestamp or datetime.now(UTC)
async with state.lock: async with state.lock:
await _cancel_pending_stop(state)
path = await _start_archive( path = await _start_archive(
ctx.module, ctx.module,
state, state,
@@ -187,19 +191,31 @@ async def archive_stream_title_update(
@on_event(EventType.STREAM_STOPPED) @on_event(EventType.STREAM_STOPPED)
async def handle_stream_stopped(ctx: EventContext[StreamStoppedEvent]) -> None: async def handle_stream_stopped(ctx: EventContext[StreamStoppedEvent]) -> None:
"""Finalize the active archive when Owncast reports that the stream stopped.""" """Schedule active archive finalization after Owncast's chat grace period."""
state = _get_state(ctx.module) state = _get_state(ctx.module)
async with state.lock: async with state.lock:
path = await _finalize_archive( await _cancel_pending_stop(state)
ctx.module, data = state.current_data
state, path = state.current_path
reason="stream_stopped_event", if data is not None and path is not None:
stream_stopped_at=ctx.event.timestamp or datetime.now(UTC), updated_at = datetime.now(UTC)
data["stream_stopped_at"] = (
ctx.event.timestamp or datetime.now(UTC)
).isoformat()
data["archive_updated_at"] = updated_at.isoformat()
await _write_json(path, data, pretty_json=state.pretty_json)
state.stop_task = asyncio.create_task(
_delayed_finalize_after_stream_stop(ctx.module, state),
name="Chat Archive Module - Stream stop grace period",
) )
if path is not None: if path is not None:
ctx.logger.info("Chat archive finalized after stream stopped: %s", path) ctx.logger.info(
"Stream stopped. Chat archive will finalize in %.0f seconds: %s",
_STOP_GRACE_SECONDS,
path,
)
@on_teardown @on_teardown
@@ -208,6 +224,7 @@ async def teardown(ctx: ModuleContext) -> None:
state = _get_state(ctx) state = _get_state(ctx)
async with state.lock: async with state.lock:
await _cancel_pending_stop(state)
await _finalize_archive(ctx, state, reason="module_teardown") await _finalize_archive(ctx, state, reason="module_teardown")
@@ -255,6 +272,39 @@ async def _start_archive(
return path return path
async def _cancel_pending_stop(state: ArchiveState) -> None:
"""Cancel a delayed stream-stop finalization task, if one exists."""
task = state.stop_task
if task is None:
return
state.stop_task = None
task.cancel()
with contextlib.suppress(asyncio.CancelledError):
await task
async def _delayed_finalize_after_stream_stop(
ctx: ModuleContext,
state: ArchiveState,
) -> None:
"""Wait for post-stream chat to close, then finalize the active archive."""
await asyncio.sleep(_STOP_GRACE_SECONDS)
async with state.lock:
if state.stop_task is not asyncio.current_task():
return
try:
await _finalize_archive(ctx, state, reason="stream_stopped_event")
except Exception:
ctx.logger.exception(
"Failed to finalize chat archive after stream stop grace period."
)
finally:
state.stop_task = None
async def _archive_event( async def _archive_event(
ctx: ModuleContext, ctx: ModuleContext,
event_type: EventType, event_type: EventType,
+129 -5
View File
@@ -287,12 +287,12 @@ class TestChatArchive:
"chat-archive-20260526T120000Z.json" "chat-archive-20260526T120000Z.json"
] ]
async def test_stream_stopped_finalizes_archive_and_ignores_later_events( async def test_stream_stopped_keeps_archive_active_during_grace(
self, self,
archive_dir: Path, archive_dir: Path,
event_dispatcher: EventDispatcher, event_dispatcher: EventDispatcher,
) -> None: ) -> None:
"""STREAM_STOPPED finalizes the active archive and closes capture.""" """STREAM_STOPPED records stop metadata but keeps capturing grace chat."""
await event_dispatcher.dispatch( await event_dispatcher.dispatch(
EventType.STREAM_STARTED, EventType.STREAM_STARTED,
make_stream_started_event( make_stream_started_event(
@@ -312,17 +312,59 @@ class TestChatArchive:
path = _archive_files(archive_dir)[0] path = _archive_files(archive_dir)[0]
data = _read_archive(path) data = _read_archive(path)
assert data["archive_stop_reason"] == "stream_stopped_event" assert data["archive_stop_reason"] is None
assert data["archive_stopped_at"] is None
assert data["stream_stopped_at"] == "2026-05-26T13:30:00+00:00" assert data["stream_stopped_at"] == "2026-05-26T13:30:00+00:00"
assert data["event_count"] == 1 assert data["event_count"] == 1
await event_dispatcher.dispatch(
EventType.CHAT,
make_chat_event(raw_body="grace chat", message_id="msg-grace"),
)
data = _read_archive(path)
assert data["archive_stop_reason"] is None
assert data["event_count"] == 2
assert _events(data)[1]["event_data"]["message_id"] == "msg-grace"
async def test_stream_stopped_finalizes_after_grace_and_ignores_later_events(
self,
archive_dir: Path,
event_dispatcher: EventDispatcher,
module_context: ModuleContext,
monkeypatch: pytest.MonkeyPatch,
) -> None:
"""The delayed stop task closes capture after the grace period."""
monkeypatch.setattr(chat_archive, "_STOP_GRACE_SECONDS", 0.0)
await event_dispatcher.dispatch(
EventType.STREAM_STARTED,
make_stream_started_event(
timestamp=datetime(2026, 5, 26, 12, tzinfo=UTC),
),
)
stopped_at = datetime(2026, 5, 26, 13, 30, tzinfo=UTC)
await event_dispatcher.dispatch(
EventType.STREAM_STOPPED,
make_stream_stopped_event(timestamp=stopped_at),
)
state = chat_archive._get_state(module_context)
assert state.stop_task is not None
await asyncio.wait_for(state.stop_task, timeout=1)
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_iso_datetime(data["archive_stopped_at"]) _assert_iso_datetime(data["archive_stopped_at"])
await event_dispatcher.dispatch( await event_dispatcher.dispatch(
EventType.CHAT, EventType.CHAT,
make_chat_event(raw_body="after stop", message_id="msg-after"), make_chat_event(raw_body="after grace", message_id="msg-after-grace"),
) )
assert _read_archive(path)["event_count"] == 1 assert _read_archive(path)["event_count"] == 0
async def test_new_stream_replaces_active_archive_and_uses_filename_suffix( async def test_new_stream_replaces_active_archive_and_uses_filename_suffix(
self, self,
@@ -364,6 +406,54 @@ class TestChatArchive:
assert new_archive["event_count"] == 1 assert new_archive["event_count"] == 1
assert _events(new_archive)[0]["event_data"]["message_id"] == "msg-new" assert _events(new_archive)[0]["event_data"]["message_id"] == "msg-new"
async def test_new_stream_during_grace_replaces_stopped_archive(
self,
archive_dir: Path,
event_dispatcher: EventDispatcher,
) -> None:
"""A new stream during the grace period cancels the delayed stop."""
first_started_at = datetime(2026, 5, 26, 12, tzinfo=UTC)
second_started_at = datetime(2026, 5, 26, 14, tzinfo=UTC)
await event_dispatcher.dispatch(
EventType.STREAM_STARTED,
make_stream_started_event(stream_title="First", timestamp=first_started_at),
)
await event_dispatcher.dispatch(
EventType.STREAM_STOPPED,
make_stream_stopped_event(
stream_title="First",
timestamp=datetime(2026, 5, 26, 13, 30, tzinfo=UTC),
),
)
await event_dispatcher.dispatch(
EventType.CHAT,
make_chat_event(raw_body="grace chat", message_id="msg-grace"),
)
await event_dispatcher.dispatch(
EventType.STREAM_STARTED,
make_stream_started_event(
stream_title="Second",
timestamp=second_started_at,
),
)
await event_dispatcher.dispatch(
EventType.CHAT,
make_chat_event(raw_body="new stream message", message_id="msg-new"),
)
old_archive = _read_archive(archive_dir / "chat-archive-20260526T120000Z.json")
assert old_archive["archive_stop_reason"] == "replaced_by_stream_started_event"
assert old_archive["stream_stopped_at"] == "2026-05-26T13:30:00+00:00"
assert old_archive["event_count"] == 1
assert _events(old_archive)[0]["event_data"]["message_id"] == "msg-grace"
new_archive = _read_archive(archive_dir / "chat-archive-20260526T140000Z.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: class TestSetupAndTeardown:
"""Cases that need direct control around setup and teardown timing.""" """Cases that need direct control around setup and teardown timing."""
@@ -487,6 +577,40 @@ class TestSetupAndTeardown:
assert data["stream_stopped_at"] is None assert data["stream_stopped_at"] is None
_assert_iso_datetime(data["archive_stopped_at"]) _assert_iso_datetime(data["archive_stopped_at"])
@pytest.mark.usefixtures("registered_handlers")
async def test_teardown_cancels_pending_stream_stop_grace(
self,
archive_dir: Path,
module_context: ModuleContext,
event_dispatcher: EventDispatcher,
) -> None:
"""Module teardown cancels delayed stop and finalizes immediately."""
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),
),
)
await event_dispatcher.dispatch(
EventType.STREAM_STOPPED,
make_stream_stopped_event(
timestamp=datetime(2026, 5, 26, 13, 30, tzinfo=UTC),
),
)
state = chat_archive._get_state(module_context)
assert state.stop_task is not None
path = _archive_files(archive_dir)[0]
await chat_archive.teardown(module_context)
assert chat_archive._get_state(module_context).stop_task is None
data = _read_archive(path)
assert data["archive_stop_reason"] == "module_teardown"
assert data["stream_stopped_at"] == "2026-05-26T13:30:00+00:00"
_assert_iso_datetime(data["archive_stopped_at"])
class TestConfiguredOutput: class TestConfiguredOutput:
"""Configuration-dependent archive output behavior.""" """Configuration-dependent archive output behavior."""