Refactored stream observations and config refresh handling.
CI / Formatting (push) Successful in 6s
CI / Linting (push) Successful in 6s
CI / Tests (push) Successful in 27s
CI / Type Checking (push) Successful in 23s
CI / Spelling (push) Successful in 19s

This commit is contained in:
2026-05-21 19:44:29 -04:00
parent c837b916c0
commit 38e91add32
14 changed files with 1503 additions and 512 deletions
+390
View File
@@ -0,0 +1,390 @@
# OwncastSentry Package
This package contains the Maubot plugin runtime code for OwncastSentry. It
tracks room subscriptions to Owncast instances, polls subscribed streams,
detects stream state changes, and sends Matrix notifications.
This README gives maintainers a high-level map of how the package fits
together and where the important behavior lives. User-facing setup and command
usage belong outside this package. Function-level details belong in docstrings
or tests.
## Files
- `__init__.py`: Maubot plugin entry point, service wiring, command
decorators, scheduler loop, shutdown, and metrics endpoint.
- `config.py`: Maubot configuration proxy for the Prometheus metrics endpoint
setting.
- `types.py`: immutable value objects, Owncast API response validation, derived
stream status, update results, and subscription domain errors.
- `repository.py`: database schema migrations plus stream and subscription
repositories.
- `owncast_client.py`: bounded aiohttp client for Owncast `/api/status` and
`/api/config` requests.
- `subscription_manager.py`: subscription use cases, domain normalization, and
first-subscription Owncast validation.
- `stream_monitor.py`: stream update orchestration, state transition detection,
progressive failure backoff, cleanup thresholds, notification decisions, and
stream metrics.
- `notification_service.py`: Matrix notification formatting, broadcast,
per-domain notification cooldowns, cleanup notices, and delivery metrics.
- `commands.py`: Maubot command handlers for subscribe, unsubscribe,
subscriptions, and live listings.
- `metrics.py`: isolated Prometheus registry, counters, gauges, and response
timing helpers.
The package has a service-layer shape: `OwncastSentry.start()` builds the
services, commands call `CommandHandler`, command handlers call
`SubscriptionManager`, scheduled updates call `StreamMonitor`, and
`StreamMonitor` coordinates repositories, the Owncast client, notifications,
and metrics.
```mermaid
flowchart LR
Maubot["Maubot loads OwncastSentry"] --> Start["OwncastSentry.start()"]
Start --> Repos["StreamRepository and SubscriptionRepository"]
Start --> Owncast["OwncastClient"]
Start --> Metrics["MetricsService"]
Start --> Notify["NotificationService"]
Start --> Monitor["StreamMonitor"]
Start --> Manager["SubscriptionManager"]
Start --> Commands["CommandHandler"]
Commands --> Manager
Manager --> Repos
Manager --> Owncast
Monitor --> Repos
Monitor --> Owncast
Monitor --> Notify
Monitor --> Metrics
Notify --> Repos
Notify --> Matrix["Matrix client"]
```
## Core Concepts
A few terms are used throughout the package:
- A stream is a normalized Owncast domain tracked by the plugin. The database
stores one stream row per domain, even if multiple Matrix rooms subscribe to
it.
- A subscription is a Matrix room's request to receive notifications for one
stream domain.
- `StreamStateObservation` is one successful Owncast status API sample. The
monitor decides whether the observation starts a new persisted status period.
- `StreamConfigObservation` is metadata fetched from Owncast's config API,
currently the instance name and tags used in notifications and periodic
display-name refreshes.
- `StreamState` is an immutable snapshot of the latest persisted state for a
stream: display metadata, the latest Owncast online flag, when that status was
first observed by the bot, failure count, and the derived online, offline, or
unknown status.
Durable state lives in the database. The `streams` table stores the latest
known display and state fields for each Owncast domain, and the `subscriptions`
table stores which Matrix rooms follow each domain.
Subscriptions are room-scoped, while stream state is shared per domain.
Notification delivery fans out from one stored stream state to the rooms
currently subscribed to that domain.
Stream identity is the normalized bare domain. Repositories, metrics,
notification state, and Owncast API calls all use that domain rather than the
original user-supplied URL.
Failures stay isolated at plugin boundaries. A bad command, failed stream
check, or room delivery is logged and counted without stopping unrelated
commands, stream updates, or notifications.
`StreamStateObservation`, `StreamConfigObservation`, `StreamState`,
`UpdateResult`, and `RoomSubscription` are immutable snapshots. State changes
are represented by new values that are persisted through repositories instead of
mutating existing objects.
## Service Wiring
`OwncastSentry.start()` is the composition root for one plugin instance. It
loads config, creates the metrics service, Owncast API client, repositories,
notification service, stream monitor, subscription manager, and command handler,
then passes dependencies into each service explicitly.
This keeps lower-level modules independent of Maubot globals. Commands do not
construct repositories, `StreamMonitor` does not know about the Maubot
scheduler, and notification delivery does not fetch its own subscription state
from the plugin instance.
The plugin boundary owns scheduling. The minute loop loads subscribed domains
from `SubscriptionRepository` and calls `StreamMonitor.update_all_streams()`.
The same pattern applies to new background work: schedule it from the plugin
boundary, keep the worker service callable directly, and close owned external
resources from `OwncastSentry.stop()`.
## Subscription Management
`SubscriptionManager` coordinates the room-to-domain relationship. Command
handlers pass it user-supplied stream targets, and it turns those targets into
normalized domains before touching repository state.
Subscribing a room creates a subscription for the normalized domain and creates
the shared stream row when the domain is first seen. First-time domains are
validated through `OwncastClient`; domains that already have subscribers reuse
the existing stream record instead of revalidating.
Unsubscribing removes one room's subscription to a domain. It does not delete
the shared stream row or make remote Owncast requests. Long-term dead stream
cleanup is owned by `StreamMonitor`.
Listing methods stay room-scoped. They return the subscriptions relevant to one
Matrix room, with the shared stream state attached for display.
## Owncast API Client
`OwncastClient` owns the aiohttp session used for Owncast API requests. It
fetches stream status and instance config, validates response shape, and turns
accepted responses into the value objects used by the rest of the package.
The session uses a plugin-specific User-Agent, a dummy cookie jar, no DNS cache,
a global connection limit, a per-host connection limit of one, and connect/read
socket timeouts. aiohttp's own DNS cache is disabled so DNS caching and policy
stay with the resolver configured for the deployment.
The per-host connection limit and keepalive settings are intentional. Each
Owncast instance gets at most one reusable connection, so minute-by-minute
status checks can reuse the existing TLS session instead of opening a fresh
connection for every poll.
Responses are accepted only when they are HTTP 200 JSON objects under the
configured size limit. Invalid JSON, non-object JSON, oversized bodies,
unexpected status codes, malformed API fields, connection errors, and timeouts
are logged and ignored.
The response-time metric is recorded only when a request completes and the
parsed response shape is valid. Failed requests remove any stale timing label
for that domain.
## Stream Monitoring
`StreamMonitor` owns the core per-domain state machine. `update_stream()`
compares the stored stream state with the latest Owncast status response,
decides whether to send a notification, updates persistent state when needed,
and records stream metrics. Each poll first handles missing rows, polling
backoff, and failed Owncast requests. Only a valid status response enters the
stream transition policy. `update_all_streams()` wraps this flow for many
domains, isolates per-domain exceptions, and records subscription-count metrics.
```mermaid
flowchart TD
Start["update_stream(domain)"] --> Old["load stored StreamState"]
Old --> Missing{"row exists?"}
Missing -->|no| Noop["return success"]
Missing -->|yes| Backoff{"backoff allows poll?"}
Backoff -->|no| Skip["skip HTTP call<br/>increment failure counter"]
Skip --> CleanupSkip["run warning/deletion cleanup checks"]
CleanupSkip --> MetricsSkip["if row remains, update<br/>status and failure metrics"]
MetricsSkip --> SuccessSkip["return success"]
Backoff -->|yes| Fetch["fetch and validate /api/status"]
Fetch --> Fetched{"got valid state?"}
Fetched -->|no| Failure["increment failure counter"]
Failure --> CleanupFail["run warning/deletion cleanup checks"]
CleanupFail --> MetricsFail["if row remains, update<br/>status and failure metrics"]
MetricsFail --> Failed["return failure"]
Fetched -->|yes| Reset["reset failure counter and metric"]
Reset --> Timer["ensure offline timer<br/>cache entry exists"]
Timer --> Transition["classify stream transition"]
Transition --> First{"first observation?"}
First -->|yes| FirstLog["log suppressed notification"]
First -->|no| LiveCheck{"went live?"}
FirstLog --> ConfigDecision
LiveCheck -->|yes| Brief{"observed offline less than<br/>temporary cooldown?"}
LiveCheck -->|no| TitleChanged{"title changed<br/>while online?"}
Brief -->|no| LiveNotify["go-live notification needed"]
Brief -->|yes| BriefTitle{"title changed?"}
BriefTitle -->|yes| BriefNotify["title-change notification needed"]
BriefTitle -->|no| NoNotify["no notification"]
TitleChanged -->|yes| OfflineMarker{"offline marker newer than<br/>last notification?"}
TitleChanged -->|no| WentOffline{"went offline?"}
WentOffline -->|yes| Offline["record offline timer"]
WentOffline -->|no| NoNotify
OfflineMarker -->|yes| ResumeNotify["go-live notification needed"]
OfflineMarker -->|no| TitleNotify["title-change notification needed"]
NoNotify --> ConfigDecision
Offline --> ConfigDecision
LiveNotify --> ConfigDecision
BriefNotify --> ConfigDecision
ResumeNotify --> ConfigDecision
TitleNotify --> ConfigDecision
ConfigDecision{"notification, first observation,<br/>or hourly refresh?"}
ConfigDecision -->|yes| Config["fetch /api/config"]
ConfigDecision -->|no| NotifyGate{"notification needed?"}
Config --> NotifyGate
NotifyGate -->|yes| Send["attempt live/title notification"]
NotifyGate -->|no| Save["save stream row<br/>repository ignores unchanged values"]
Send --> Save
Save --> FinalMetrics["set current status metric"]
FinalMetrics --> Done["return success"]
```
Failure counters drive both status and polling behavior. A stream is displayed
as unknown when its failure counter is above `UNKNOWN_STATUS_THRESHOLD`.
Counters 0 through 4 poll every minute; later tiers poll only on selected
counter values and eventually only every fifteenth minute. Skipped cycles still
increment the counter so cleanup can progress.
Transition decisions use Owncast's `online` field. The stream status timestamp is
the bot's local UTC observation time for the first successful poll, a missing
timestamp repair, or an online/offline transition. Owncast's remote connect,
disconnect, and server timestamps are not required for status parsing.
First observation is special. If a new stream row has no status timestamp, the
first successful update fetches config and stores state but suppresses
notification attempts, even if the stream is already live.
Go-live notification attempts happen when a stream moves from offline to online
and the monitor's in-memory offline timer shows it was not merely a brief
outage. A stream that returns within the temporary offline cooldown only attempts
a notification when the title changed, and then it uses a title change
notification.
Mid-session title changes attempt title change notifications. If the monitor has
an offline marker that is newer than the last notification, it attempts a normal
live notification instead so rooms do not see a title change for what is
effectively a new session.
Offline transitions do not notify rooms. They update the database and record a
monotonic offline timestamp used by later live/title notification policy.
The live/title notification send step is an attempt, not a guaranteed Matrix message.
`NotificationService` still applies the shared live/title cooldown and records
the cooldown only when at least one room receives the message.
`/api/config` is fetched when the monitor is about to attempt a live or
title-change notification, on a stream's first successful observation, and
during the hourly refresh window for successfully polled streams. The config
response supplies the display name and tags for notifications. If config is
fetched successfully during an update, the persisted stream name is refreshed
from it; otherwise the stored name is left unchanged.
## Cleanup Policy
A domain that remains unreachable is eventually cleaned up. The failure counter
represents minute ticks, including skipped backoff ticks.
At the 83-day threshold, `StreamMonitor` sends a cleanup warning to all rooms
subscribed to the domain. At the 90-day threshold or beyond, it sends a deletion
notice, deletes all subscriptions for the domain, deletes the stream row, clears
the monitor and notification caches for that domain, and removes per-domain
metrics.
```mermaid
flowchart TD
Counter["failure counter after failed or skipped check"] --> Warn{"counter == 83 days?"}
Warn -->|yes| Warning["send cleanup warning"]
Warn -->|no| Delete
Warning --> Delete{"counter >= 90 days?"}
Delete -->|no| Done["done"]
Delete -->|yes| Notice["send cleanup deletion notice"]
Notice --> Subs["delete subscriptions"]
Subs --> Stream["delete stream row"]
Stream --> Caches["clear local caches"]
Caches --> Metrics["remove per-domain metrics"]
Metrics --> DoneCleanup["done"]
```
## Notification Delivery
`NotificationService` sends plain text Matrix messages to every room subscribed
to a domain. Broadcasts run concurrently and use `asyncio.gather(...,
return_exceptions=True)` so one failed room does not block delivery to other
rooms.
Live and title-change notifications share a per-domain cooldown. The cooldown
is recorded only when at least one room receives the message. If there are no
subscribed rooms, or every delivery fails, no cooldown is recorded.
Notification text is whitespace-normalized before sending. Stream names fall
back to the domain when config lookup fails or the instance name is empty. Tags
are appended as hash tags after filtering out empty tags and tags that start
with a dot.
Cleanup warning and deletion notices bypass the live/title cooldown, but still
record delivery metrics.
## Command Handling
`OwncastSentry` exposes Maubot command decorators in `__init__.py`, but command
behavior lives in `CommandHandler`. The plugin entry points are thin wrappers
that catch unexpected exceptions and record command error metrics.
`CommandHandler.subscribe()` and `CommandHandler.unsubscribe()` translate
domain errors into user-facing Matrix replies. `subscriptions()` and `live()`
read resolved room subscriptions and render Markdown output with escaped
untrusted text.
Command listings are room-scoped and ordered by domain through repository
queries. `subscriptions()` includes online, offline, and unknown streams.
`live()` includes only rows marked online with a failure count at or below the
unknown threshold.
## Persistence
The database schema is managed by `repository.py` through Maubot's upgrade
table.
- Revision 1 creates `streams` and `subscriptions`.
- Revision 2 fixes `subscriptions.stream_domain` from `INTEGER` to `TEXT`.
- Revision 3 adds `streams.failure_counter`.
- Revision 4 replaces separate connect/disconnect timestamps with
`streams.online` and `streams.status_since`.
Repositories own SQL access. `StreamRepository.update()` writes display
and state fields; failure counters use dedicated methods.
`SubscriptionRepository` raises domain-specific errors for duplicate adds and
missing removes.
Resolved subscription listings join `subscriptions` to `streams`, which means
orphaned subscription rows without a stream row are skipped in room display
queries.
## Metrics
`MetricsService` uses a private `CollectorRegistry` so plugin metrics do not mix
with process-global Prometheus collectors.
The registered metrics cover:
- notification delivery attempts by type and result;
- current stream status per domain (`1` online, `0` offline, `-1` unknown);
- subscription count per domain;
- consecutive check failures per domain;
- last successful Owncast API response duration per domain;
- build/version information;
- open aiohttp connection count;
- internal scheduler and command errors.
Known counter label combinations are initialized to zero. This keeps the
Prometheus output stable before any events have occurred.
## Test Map
Start with these tests when changing package behavior:
- `tests/test_types.py`: immutable value objects, API response shape checks,
truncation, stream status derivation, and subscription errors.
- `tests/test_repository.py`: schema-backed stream and subscription repository
behavior.
- `tests/test_subscription_manager.py`: domain normalization, subscribe and
unsubscribe workflows, validation skipping, and room listing delegation.
- `tests/test_owncast_client.py`: Owncast API request handling, response body
limits, validation failures, User-Agent, response timing, and connection
counts.
- `tests/test_stream_monitor.py`: polling backoff, first-update suppression,
live/offline/title transitions, cleanup thresholds, exception isolation, and
monitor metrics.
- `tests/test_notification_service.py`: message formatting, sanitization,
cooldown behavior, broadcast failure accounting, cleanup notices, and
notification metrics.
- `tests/test_commands.py`: end-to-end Maubot command behavior, Markdown
escaping, duration formatting, room-scoped subscription listings, and live
listings.
- `tests/test_metrics.py`: Prometheus counters, gauges, response timer,
per-domain label removal, build info, open connections, and isolated registry
output.
+1 -1
View File
@@ -78,7 +78,7 @@ def _format_duration(timestamp_str: str, now: datetime) -> str:
if seconds < _SECONDS_PER_DAY: if seconds < _SECONDS_PER_DAY:
hours = seconds // _SECONDS_PER_HOUR hours = seconds // _SECONDS_PER_HOUR
return f"{hours} hour{'s' if hours != 1 else ''}" return f"{hours} hour{'s' if hours != 1 else ''}"
except (TypeError, ValueError): except TypeError, ValueError:
return "unknown duration" return "unknown duration"
else: else:
days = seconds // _SECONDS_PER_DAY days = seconds // _SECONDS_PER_DAY
+3 -3
View File
@@ -23,6 +23,9 @@ from prometheus_client import CollectorRegistry, Counter, Gauge, Info
from .types import StreamStatus from .types import StreamStatus
if TYPE_CHECKING:
from collections.abc import Callable, Generator
class NotificationType(StrEnum): class NotificationType(StrEnum):
"""Notification type labels for the delivery counter.""" """Notification type labels for the delivery counter."""
@@ -40,9 +43,6 @@ class ErrorSource(StrEnum):
COMMAND = "command" COMMAND = "command"
if TYPE_CHECKING:
from collections.abc import Callable, Generator
# Mapping from StreamStatus enum to numeric gauge values # Mapping from StreamStatus enum to numeric gauge values
_STATUS_VALUES: dict[StreamStatus, float] = { _STATUS_VALUES: dict[StreamStatus, float] = {
StreamStatus.ONLINE: 1.0, StreamStatus.ONLINE: 1.0,
+13 -11
View File
@@ -15,13 +15,16 @@
"""HTTP client for querying Owncast instance APIs.""" """HTTP client for querying Owncast instance APIs."""
import json import json
from datetime import UTC, datetime
from http import HTTPStatus from http import HTTPStatus
from typing import TYPE_CHECKING, Any from typing import TYPE_CHECKING, Any
import aiohttp import aiohttp
from .types import InvalidApiResponseError, StreamConfig, StreamState from .types import (
InvalidApiResponseError,
StreamConfigObservation,
StreamStateObservation,
)
if TYPE_CHECKING: if TYPE_CHECKING:
import logging import logging
@@ -110,14 +113,14 @@ class OwncastClient:
connector=connector, connector=connector,
) )
async def get_stream_state(self, domain: str) -> StreamState | None: async def get_stream_state(self, domain: str) -> StreamStateObservation | None:
"""Get the current stream state for a given domain. """Get the current stream state for a given domain.
HTTPS on port 443 is assumed, no other protocols or ports HTTPS on port 443 is assumed, no other protocols or ports
are supported. are supported.
:param domain: The domain (not URL) where the stream is hosted. :param domain: The domain (not URL) where the stream is hosted.
:return: A StreamState if available, None on error. :return: A StreamStateObservation if available, None on error.
""" """
self.log.debug("[%s] Fetching current stream state...", domain) self.log.debug("[%s] Fetching current stream state...", domain)
with self.metrics.response_timer(domain) as timer: with self.metrics.response_timer(domain) as timer:
@@ -126,10 +129,9 @@ class OwncastClient:
if new_state is None: if new_state is None:
return None return None
observed_at = datetime.now(UTC)
try: try:
stream_state = StreamState.from_api_response( stream_observation = StreamStateObservation.from_api_response(
new_state, domain, observed_at new_state, domain
) )
except InvalidApiResponseError as e: except InvalidApiResponseError as e:
self.log.warning( self.log.warning(
@@ -142,16 +144,16 @@ class OwncastClient:
return None return None
timer.success() timer.success()
return stream_state return stream_observation
async def get_stream_config(self, domain: str) -> StreamConfig | None: async def get_stream_config(self, domain: str) -> StreamConfigObservation | None:
"""Get the current stream config for a given domain. """Get the current stream config for a given domain.
HTTPS on port 443 is assumed, no other protocols or ports HTTPS on port 443 is assumed, no other protocols or ports
are supported. are supported.
:param domain: The domain (not URL) where the stream is hosted. :param domain: The domain (not URL) where the stream is hosted.
:return: A StreamConfig, or None if fetch failed. :return: A StreamConfigObservation, or None if fetch failed.
""" """
self.log.debug("[%s] Fetching current stream config...", domain) self.log.debug("[%s] Fetching current stream config...", domain)
with self.metrics.response_timer(domain) as timer: with self.metrics.response_timer(domain) as timer:
@@ -160,7 +162,7 @@ class OwncastClient:
return None return None
try: try:
stream_config = StreamConfig.from_api_response(config) stream_config = StreamConfigObservation.from_api_response(config)
except InvalidApiResponseError as e: except InvalidApiResponseError as e:
self.log.warning( self.log.warning(
"[%s] Rejecting response to request on %s as response " "[%s] Rejecting response to request on %s as response "
+64 -20
View File
@@ -12,7 +12,10 @@
# See the License for the specific language governing permissions and # See the License for the specific language governing permissions and
# limitations under the License. # limitations under the License.
"""Repository and schema upgrade definitions for OwncastSentry.""" """Repository and schema upgrade definitions for OwncastSentry.
Only SQLite is supported as a database backend for now.
"""
from datetime import UTC, datetime from datetime import UTC, datetime
from typing import TYPE_CHECKING, Any from typing import TYPE_CHECKING, Any
@@ -35,9 +38,13 @@ if TYPE_CHECKING:
upgrade_table = UpgradeTable() upgrade_table = UpgradeTable()
def _has_legacy_timestamp(value: Any) -> bool: class _Unset:
"""Return whether a legacy timestamp value carries usable content.""" """Sentinel type for omitted partial stream update fields."""
return value is not None and str(value).strip() != ""
__slots__ = ()
_UNSET = _Unset()
def _normalize_legacy_status_since(value: Any) -> str | None: def _normalize_legacy_status_since(value: Any) -> str | None:
@@ -165,10 +172,9 @@ async def upgrade_v4(conn: Connection) -> None:
FROM streams""" FROM streams"""
) )
for row in rows: for row in rows:
online = _has_legacy_timestamp(row["last_connect_time"]) last_connect_time = row["last_connect_time"]
legacy_timestamp = ( online = last_connect_time is not None and str(last_connect_time).strip() != ""
row["last_connect_time"] if online else row["last_disconnect_time"] legacy_timestamp = last_connect_time if online else row["last_disconnect_time"]
)
await conn.execute( await conn.execute(
"""INSERT INTO streams_new ( """INSERT INTO streams_new (
domain, name, title, online, status_since, failure_counter domain, name, title, online, status_since, failure_counter
@@ -234,25 +240,63 @@ class StreamRepository:
result = await self.get_by_domain(domain) result = await self.get_by_domain(domain)
return result is not None return result is not None
async def update(self, state: StreamState) -> None: async def update(
"""Update a stream's state in the database. self,
domain: str,
*,
name: str | None | _Unset = _UNSET,
title: str | None | _Unset = _UNSET,
online: bool | _Unset = _UNSET,
status_since: str | None | _Unset = _UNSET,
) -> None:
"""Update only the supplied stream fields in the database.
This updates display/state fields only. Failure counters are Passing None writes NULL for nullable fields. Omitting a field leaves
updated through dedicated methods. that column unchanged.
:param state: The StreamState to save. :param domain: The stream domain.
:param name: Optional stream display name update.
:param title: Optional stream title update.
:param online: Optional stream online state update.
:param status_since: Optional current status timestamp update.
""" """
# _UNSET marks omitted fields so None can still be written as SQL NULL.
update_name = not isinstance(name, _Unset)
update_title = not isinstance(title, _Unset)
update_online = not isinstance(online, _Unset)
update_status_since = not isinstance(status_since, _Unset)
if not (update_name or update_title or update_online or update_status_since):
return
name_value = None if isinstance(name, _Unset) else name
title_value = None if isinstance(title, _Unset) else title
online_value = None if isinstance(online, _Unset) else online
status_since_value = None if isinstance(status_since, _Unset) else status_since
# SQLite's IS NOT gives null-safe comparisons for the supported backend.
query = """UPDATE streams query = """UPDATE streams
SET name=$1, title=$2, online=$3, status_since=$4 SET name = CASE WHEN $1 THEN $2 ELSE name END,
WHERE domain=$5""" title = CASE WHEN $3 THEN $4 ELSE title END,
online = CASE WHEN $5 THEN $6 ELSE online END,
status_since = CASE WHEN $7 THEN $8 ELSE status_since END
WHERE domain=$9
AND (($1 AND name IS NOT $2)
OR ($3 AND title IS NOT $4)
OR ($5 AND online IS NOT $6)
OR ($7 AND status_since IS NOT $8))"""
async with self.db.acquire() as conn: async with self.db.acquire() as conn:
await conn.execute( await conn.execute(
query, query,
state.name, update_name,
state.title, name_value,
state.online, update_title,
state.status_since, title_value,
state.domain, update_online,
online_value,
update_status_since,
status_since_value,
domain,
) )
async def delete(self, domain: str) -> None: async def delete(self, domain: str) -> None:
+194 -130
View File
@@ -16,9 +16,18 @@
import asyncio import asyncio
import time import time
from dataclasses import dataclass
from datetime import UTC, datetime
from enum import Enum, auto
from typing import TYPE_CHECKING from typing import TYPE_CHECKING
from .types import StreamState, UpdateResult from .types import (
StreamConfigObservation,
StreamState,
StreamStateObservation,
StreamStatus,
UpdateResult,
)
if TYPE_CHECKING: if TYPE_CHECKING:
import logging import logging
@@ -56,6 +65,102 @@ def _should_query_stream(failure_counter: int) -> bool:
return failure_counter % _QUERY_EVERY_15_MINUTES_INTERVAL == 0 return failure_counter % _QUERY_EVERY_15_MINUTES_INTERVAL == 0
class _StreamTransitionKind(Enum):
"""High-level status transition from stored state to fresh observation."""
FIRST_OBSERVATION = auto()
WENT_LIVE = auto()
WENT_OFFLINE = auto()
TITLE_CHANGED = auto()
STATUS_UNCHANGED = auto()
@dataclass(frozen=True, slots=True)
class _StreamTransition:
"""Classified stream transition and timestamp metadata."""
kind: _StreamTransitionKind
status_since: str | None
class _NotificationKind(Enum):
"""Live/title notification action for a transition."""
NONE = auto()
LIVE = auto()
TITLE_CHANGE = auto()
@property
def should_notify(self) -> bool:
"""Return whether a live/title-change notification should be attempted."""
return self is not _NotificationKind.NONE
def _classify_transition(
old_state: StreamState, observation: StreamStateObservation
) -> _StreamTransition:
"""Classify the transition from the stored state to the latest observation."""
if old_state.status_since is None:
return _StreamTransition(
kind=_StreamTransitionKind.FIRST_OBSERVATION,
status_since=observation.observed_at,
)
if old_state.online != observation.online:
return _StreamTransition(
kind=(
_StreamTransitionKind.WENT_LIVE
if observation.online
else _StreamTransitionKind.WENT_OFFLINE
),
status_since=observation.observed_at,
)
if old_state.online and old_state.title != observation.title:
return _StreamTransition(
kind=_StreamTransitionKind.TITLE_CHANGED,
status_since=old_state.status_since,
)
return _StreamTransition(
kind=_StreamTransitionKind.STATUS_UNCHANGED,
status_since=old_state.status_since,
)
def _classify_notification(
transition: _StreamTransition,
old_state: StreamState,
observation: StreamStateObservation,
*,
offline_duration_seconds: int | None = None,
went_offline_after_last_notification: bool = False,
) -> _NotificationKind:
"""Classify the live/title notification policy for a transition.
The caller supplies timing context because it comes from the monitor's
monotonic offline cache and notification-service cache, not the transition.
"""
match transition.kind:
case _StreamTransitionKind.WENT_LIVE:
if offline_duration_seconds is not None and (
offline_duration_seconds < _TEMPORARY_OFFLINE_NOTIFICATION_COOLDOWN
):
if old_state.title != observation.title:
return _NotificationKind.TITLE_CHANGE
return _NotificationKind.NONE
return _NotificationKind.LIVE
case _StreamTransitionKind.TITLE_CHANGED:
if went_offline_after_last_notification:
return _NotificationKind.LIVE
return _NotificationKind.TITLE_CHANGE
case _:
return _NotificationKind.NONE
class StreamMonitor: class StreamMonitor:
"""Monitors Owncast streams and detects state changes.""" """Monitors Owncast streams and detects state changes."""
@@ -178,22 +283,11 @@ class StreamMonitor:
# Backoff is expected behavior, not a failure # Backoff is expected behavior, not a failure
return True return True
# Flag: no status timestamp has been recorded yet, so suppress
# notifications for a stream's first observed state.
first_update = False
# Flag: whether to update the stream's state in the database.
# Used to avoid writes when state hasn't changed at all.
update_database = False
# The stream's latest configuration, if fetched during update.
stream_config = None
# Fetch the latest stream state from the server # Fetch the latest stream state from the server
new_state = await self.owncast_client.get_stream_state(domain) observation = await self.owncast_client.get_stream_state(domain)
# If the fetch failed, increment failure counter and skip the update # If the fetch failed, increment failure counter and skip the update
if new_state is None: if observation is None:
await self.stream_repo.increment_failure_counter(domain) await self.stream_repo.increment_failure_counter(domain)
self.log.warning( self.log.warning(
"[%s] Connection failure (counter=%s)", "[%s] Connection failure (counter=%s)",
@@ -216,138 +310,108 @@ class StreamMonitor:
# Initialize timer cache entries to prevent KeyError on first access # Initialize timer cache entries to prevent KeyError on first access
self.offline_timer_cache.setdefault(domain, 0) self.offline_timer_cache.setdefault(domain, 0)
if old_state.status_since is None: transition = _classify_transition(old_state, observation)
# No stream history has been recorded yet. Don't send notifications. # Notification policy may need runtime timing context in addition to
update_database = True # the stored state and latest observation.
first_update = True offline_duration_seconds: int | None = None
went_offline_after_last_notification = False
if first_update: match transition.kind:
self.log.info( case _StreamTransitionKind.FIRST_OBSERVATION:
"[%s] Not sending notifications. This is the first state " self.log.info(
"update for this stream.", "[%s] Not sending notifications. This is the first state "
domain, "update for this stream.",
) domain,
)
# Did the stream become publicly online? case _StreamTransitionKind.WENT_LIVE:
elif new_state.online and not old_state.online: self.log.info("[%s] Stream is now live!", domain)
# Yes! This stream is now live. offline_duration_seconds = round(
update_database = True time.monotonic() - self.offline_timer_cache[domain]
stream_config = await self.owncast_client.get_stream_config(domain) )
# Brief reconnects with the same title are treated as transient
self.log.info("[%s] Stream is now live!", domain) # outages, not new live events.
if (
# Calculate seconds since the stream last went offline offline_duration_seconds < _TEMPORARY_OFFLINE_NOTIFICATION_COOLDOWN
seconds_since_last_offline = round( and old_state.title == observation.title
time.monotonic() - self.offline_timer_cache[domain] ):
)
# Use fallback values if config fetch failed
stream_name = stream_config.name if stream_config else domain
stream_tags = stream_config.tags if stream_config else ()
# Has this stream been offline for a short time?
if seconds_since_last_offline < _TEMPORARY_OFFLINE_NOTIFICATION_COOLDOWN:
# Did the stream title change?
if old_state.title != new_state.title:
# Stream was briefly down; send title change notification.
await self.notification_service.notify_stream_live(
domain,
stream_name,
new_state.title or "",
stream_tags,
title_change=True,
)
else:
# Briefly offline, no title change. Skip.
self.log.info( self.log.info(
"[%s] Not sending notifications. Stream was only " "[%s] Not sending notifications. Stream was only offline for "
"offline for %s of %s seconds and did not change " "%s of %s seconds and did not change its title.",
"its title.",
domain, domain,
seconds_since_last_offline, offline_duration_seconds,
_TEMPORARY_OFFLINE_NOTIFICATION_COOLDOWN, _TEMPORARY_OFFLINE_NOTIFICATION_COOLDOWN,
) )
else: case _StreamTransitionKind.TITLE_CHANGED:
# Offline for a while. Send a normal notification.
await self.notification_service.notify_stream_live(
domain,
stream_name,
new_state.title or "",
stream_tags,
title_change=False,
)
elif new_state.online and old_state.online:
# Did the stream title change mid-session?
if old_state.title != new_state.title:
self.log.info("[%s] Stream title was changed!", domain) self.log.info("[%s] Stream title was changed!", domain)
update_database = True # If we saw an offline event after the last notification, send
stream_config = await self.owncast_client.get_stream_config(domain) # a go-live notice instead of only a title-change notice.
went_offline_after_last_notification = self.offline_timer_cache[
# Use fallback values if config fetch failed
stream_name = stream_config.name if stream_config else domain
stream_tags = stream_config.tags if stream_config else ()
# If the offline marker is newer than the last recorded
# notification, send a regular go-live instead of a title
# change to avoid confusion.
if self.offline_timer_cache[
domain domain
] > self.notification_service.get_last_notification_time(domain): ] > self.notification_service.get_last_notification_time(domain)
await self.notification_service.notify_stream_live( case _StreamTransitionKind.WENT_OFFLINE:
domain, self.offline_timer_cache[domain] = time.monotonic()
stream_name, self.log.info("[%s] Stream is now offline.", domain)
new_state.title or "", case _StreamTransitionKind.STATUS_UNCHANGED:
stream_tags, pass
title_change=False,
)
else:
# No. Send a normal title change notification.
await self.notification_service.notify_stream_live(
domain,
stream_name,
new_state.title or "",
stream_tags,
title_change=True,
)
# Did the stream go offline? notification_kind = _classify_notification(
elif not new_state.online and old_state.online: transition,
# Yep. This stream is now offline. Log it. old_state,
update_database = True observation,
self.offline_timer_cache[domain] = time.monotonic() offline_duration_seconds=offline_duration_seconds,
self.log.info("[%s] Stream is now offline.", domain) went_offline_after_last_notification=went_offline_after_last_notification,
)
# Update the database with current stream state, if needed. stream_config: StreamConfigObservation | None = None
if update_database: observed_at = datetime.fromisoformat(observation.observed_at)
# Ensure we have the stream config before updating hourly_config_refresh_due = (
if stream_config is None: observed_at.tzinfo is not None and observed_at.astimezone(UTC).minute == 0
stream_config = await self.owncast_client.get_stream_config(domain) )
should_fetch_config = (
notification_kind.should_notify
or transition.kind is _StreamTransitionKind.FIRST_OBSERVATION
or hourly_config_refresh_due
)
if should_fetch_config:
stream_config = await self.owncast_client.get_stream_config(domain)
# Use fallback value if config fetch failed if notification_kind.should_notify:
stream_name = stream_config.name if stream_config else "" stream_name = stream_config.name if stream_config else domain
stream_tags = stream_config.tags if stream_config else ()
is_title_change = notification_kind is _NotificationKind.TITLE_CHANGE
self.log.debug("[%s] Updating stream state in database...", domain) await self.notification_service.notify_stream_live(
domain,
if first_update or old_state.online != new_state.online: stream_name,
status_since = new_state.status_since observation.title or "",
else: stream_tags,
status_since = old_state.status_since title_change=is_title_change,
# Create updated state object (title already truncated in new_state)
updated_state = StreamState(
domain=domain,
name=stream_name,
title=new_state.title,
online=new_state.online,
status_since=status_since,
) )
await self.stream_repo.update(updated_state) self.log.debug("[%s] Saving stream state if changed...", domain)
if stream_config is not None:
await self.stream_repo.update(
domain,
name=stream_config.name,
title=observation.title,
online=observation.online,
status_since=transition.status_since,
)
else:
await self.stream_repo.update(
domain,
title=observation.title,
online=observation.online,
status_since=transition.status_since,
)
# All done. # All done.
self.log.debug("[%s] State update completed.", domain) self.log.debug("[%s] State update completed.", domain)
self.metrics.set_stream_status(domain, new_state.status) stream_status = (
StreamStatus.ONLINE if observation.online else StreamStatus.OFFLINE
)
self.metrics.set_stream_status(domain, stream_status)
return True return True
async def _check_cleanup_thresholds(self, domain: str, counter: int) -> None: async def _check_cleanup_thresholds(self, domain: str, counter: int) -> None:
+60 -31
View File
@@ -14,7 +14,7 @@
"""Data containers and domain errors for OwncastSentry.""" """Data containers and domain errors for OwncastSentry."""
from dataclasses import dataclass from dataclasses import dataclass, field
from datetime import UTC, datetime from datetime import UTC, datetime
from enum import Enum from enum import Enum
from typing import Any from typing import Any
@@ -80,6 +80,11 @@ def format_status_since(timestamp: datetime) -> str:
return timestamp.astimezone(UTC).isoformat(timespec="seconds") return timestamp.astimezone(UTC).isoformat(timespec="seconds")
def _current_observed_at() -> str:
"""Return the current UTC time in the package's timestamp format."""
return format_status_since(datetime.now(UTC))
class StreamStatus(Enum): class StreamStatus(Enum):
"""Represents the status of a stream.""" """Represents the status of a stream."""
@@ -88,9 +93,51 @@ class StreamStatus(Enum):
UNKNOWN = "unknown" UNKNOWN = "unknown"
@dataclass(frozen=True, slots=True)
class StreamStateObservation:
"""Status API sample fetched from an Owncast instance.
This records what `/api/status` reported during one request. `observed_at`
belongs to that fetch; the monitor decides whether that time starts a new
persisted online/offline period in `StreamState`.
"""
domain: str
title: str | None = None
online: bool = False
observed_at: str = field(default_factory=_current_observed_at)
@classmethod
def from_api_response(
cls, response: dict[str, Any], domain: str
) -> StreamStateObservation:
"""Create a StreamStateObservation from a status API response.
:param response: API response as a dictionary (camelCase keys).
:param domain: The stream domain.
:return: StreamStateObservation instance.
:raises InvalidApiResponseError: If the response shape is invalid.
"""
stream_title = _require_str(response, "streamTitle")
online = _require_field(response, "online")
if not isinstance(online, bool):
raise InvalidApiResponseError("online must be a boolean")
return cls(
domain=domain,
title=_truncate(stream_title, _MAX_STREAM_TITLE_LENGTH),
online=online,
)
@dataclass(frozen=True, slots=True) @dataclass(frozen=True, slots=True)
class StreamState: class StreamState:
"""Represents the state of an Owncast stream.""" """Stream record remembered by the repository.
States include persisted display data, the current online/offline period,
and failure counters used to derive user-facing status. `status_since`
records when the persisted online/offline period began.
"""
domain: str domain: str
name: str | None = None name: str | None = None
@@ -112,30 +159,6 @@ class StreamState:
return StreamStatus.ONLINE return StreamStatus.ONLINE
return StreamStatus.OFFLINE return StreamStatus.OFFLINE
@classmethod
def from_api_response(
cls, response: dict[str, Any], domain: str, observed_at: datetime
) -> StreamState:
"""Create a StreamState from an API response.
:param response: API response as a dictionary (camelCase keys).
:param domain: The stream domain.
:param observed_at: Local time when this status was observed.
:return: StreamState instance.
:raises InvalidApiResponseError: If the response shape is invalid.
"""
stream_title = _require_str(response, "streamTitle")
online = _require_field(response, "online")
if not isinstance(online, bool):
raise InvalidApiResponseError("online must be a boolean")
return cls(
domain=domain,
title=_truncate(stream_title, _MAX_STREAM_TITLE_LENGTH),
online=online,
status_since=format_status_since(observed_at),
)
@classmethod @classmethod
def from_db_row(cls, row: dict[str, Any]) -> StreamState: def from_db_row(cls, row: dict[str, Any]) -> StreamState:
"""Create a StreamState from a database row. """Create a StreamState from a database row.
@@ -154,18 +177,24 @@ class StreamState:
@dataclass(frozen=True, slots=True) @dataclass(frozen=True, slots=True)
class StreamConfig: class StreamConfigObservation:
"""Represents the configuration of an Owncast stream.""" """Config API sample fetched from an Owncast instance.
This records display metadata from one `/api/config` request. It is used for
notification text and may refresh persisted display fields independently of
the stream's online/offline state.
"""
name: str = "" name: str = ""
tags: tuple[str, ...] = () tags: tuple[str, ...] = ()
observed_at: str = field(default_factory=_current_observed_at)
@classmethod @classmethod
def from_api_response(cls, response: dict[str, Any]) -> StreamConfig: def from_api_response(cls, response: dict[str, Any]) -> StreamConfigObservation:
"""Create a StreamConfig from an API response. """Create a StreamConfigObservation from an API response.
:param response: API response as a dictionary. :param response: API response as a dictionary.
:return: StreamConfig instance. :return: StreamConfigObservation instance.
:raises InvalidApiResponseError: If the response shape is invalid. :raises InvalidApiResponseError: If the response shape is invalid.
""" """
# Apply Owncast's instance name limit. # Apply Owncast's instance name limit.
+6 -6
View File
@@ -38,7 +38,7 @@ if TYPE_CHECKING:
from mautrix.util.config import BaseProxyConfig from mautrix.util.config import BaseProxyConfig
from owncastsentry.metrics import MetricsService from owncastsentry.metrics import MetricsService
from owncastsentry.types import StreamConfig, StreamState from owncastsentry.types import StreamConfigObservation, StreamStateObservation
def generate_metrics_output(metrics: MetricsService) -> str: def generate_metrics_output(metrics: MetricsService) -> str:
@@ -212,19 +212,19 @@ class _StubMatrixClient:
class _StubOwncastClient: class _StubOwncastClient:
"""Configurable stub for the Owncast HTTP client.""" """Configurable stub for the Owncast HTTP client."""
stream_state: StreamState | None = None stream_state: StreamStateObservation | None = None
stream_config: StreamConfig | None = None stream_config: StreamConfigObservation | None = None
state_call_count: int = field(default=0, init=False) state_call_count: int = field(default=0, init=False)
config_call_count: int = field(default=0, init=False) config_call_count: int = field(default=0, init=False)
queried_domains: list[str] = field(default_factory=list, init=False) queried_domains: list[str] = field(default_factory=list, init=False)
async def get_stream_state(self, domain: str) -> StreamState | None: async def get_stream_state(self, domain: str) -> StreamStateObservation | None:
"""Return the configured stream state.""" """Return the configured stream observation."""
self.state_call_count += 1 self.state_call_count += 1
self.queried_domains.append(domain) self.queried_domains.append(domain)
return self.stream_state return self.stream_state
async def get_stream_config(self, _domain: str) -> StreamConfig | None: async def get_stream_config(self, _domain: str) -> StreamConfigObservation | None:
"""Return the configured stream config.""" """Return the configured stream config."""
self.config_call_count += 1 self.config_call_count += 1
return self.stream_config return self.stream_config
+40 -58
View File
@@ -31,7 +31,7 @@ from owncastsentry.commands import (
_sanitize_for_markdown, _sanitize_for_markdown,
) )
from owncastsentry.owncast_client import _OWNCAST_STATUS_PATH from owncastsentry.owncast_client import _OWNCAST_STATUS_PATH
from owncastsentry.types import UNKNOWN_STATUS_THRESHOLD, StreamState from owncastsentry.types import UNKNOWN_STATUS_THRESHOLD
from tests.conftest import VALID_STATUS_RESPONSE from tests.conftest import VALID_STATUS_RESPONSE
if TYPE_CHECKING: if TYPE_CHECKING:
@@ -258,13 +258,11 @@ class TestSubscriptionsCommand:
# Update stream state to be online with a name # Update stream state to be online with a name
await maubot_plugin.stream_repo.update( await maubot_plugin.stream_repo.update(
StreamState( "stream.logal.dev",
domain="stream.logal.dev", name="Test Stream",
name="Test Stream", title="Playing Games",
title="Playing Games", online=True,
online=True, status_since="2026-01-01T12:00:00+00:00",
status_since="2026-01-01T12:00:00+00:00",
)
) )
await maubot_test_bot.send("!subscriptions") await maubot_test_bot.send("!subscriptions")
@@ -294,13 +292,11 @@ class TestSubscriptionsCommand:
await maubot_test_bot.send("!subscribe stream.logal.dev") await maubot_test_bot.send("!subscribe stream.logal.dev")
await maubot_plugin.stream_repo.update( await maubot_plugin.stream_repo.update(
StreamState( "stream.logal.dev",
domain="stream.logal.dev", name="*Bold* [link](https://evil.example)\nName",
name="*Bold* [link](https://evil.example)\nName", title="`code` > quote #tag",
title="`code` > quote #tag", online=True,
online=True, status_since="2026-01-01T12:00:00+00:00",
status_since="2026-01-01T12:00:00+00:00",
)
) )
await maubot_test_bot.send("!subscriptions") await maubot_test_bot.send("!subscriptions")
@@ -332,11 +328,9 @@ class TestSubscriptionsCommand:
# Update stream state to be offline # Update stream state to be offline
await maubot_plugin.stream_repo.update( await maubot_plugin.stream_repo.update(
StreamState( "stream.logal.dev",
domain="stream.logal.dev", name="Test Stream",
name="Test Stream", status_since="2026-01-01T10:00:00+00:00",
status_since="2026-01-01T10:00:00+00:00",
)
) )
await maubot_test_bot.send("!subscriptions") await maubot_test_bot.send("!subscriptions")
@@ -426,20 +420,16 @@ class TestSubscriptionsCommand:
# Set alpha online, beta offline # Set alpha online, beta offline
await maubot_plugin.stream_repo.update( await maubot_plugin.stream_repo.update(
StreamState( "alpha.com",
domain="alpha.com", name="Alpha Stream",
name="Alpha Stream", title="Streaming Live",
title="Streaming Live", online=True,
online=True, status_since="2026-03-13T10:00:00+00:00",
status_since="2026-03-13T10:00:00+00:00",
)
) )
await maubot_plugin.stream_repo.update( await maubot_plugin.stream_repo.update(
StreamState( "beta.com",
domain="beta.com", name="Beta Stream",
name="Beta Stream", status_since="2026-03-12T18:00:00+00:00",
status_since="2026-03-12T18:00:00+00:00",
)
) )
await maubot_test_bot.send("!subscriptions") await maubot_test_bot.send("!subscriptions")
@@ -491,11 +481,9 @@ class TestLiveCommand:
# Update stream state to offline # Update stream state to offline
await maubot_plugin.stream_repo.update( await maubot_plugin.stream_repo.update(
StreamState( "stream.logal.dev",
domain="stream.logal.dev", name="Test Stream",
name="Test Stream", status_since="2026-01-01T10:00:00+00:00",
status_since="2026-01-01T10:00:00+00:00",
)
) )
await maubot_test_bot.send("!live") await maubot_test_bot.send("!live")
@@ -523,13 +511,11 @@ class TestLiveCommand:
# Update stream state to online # Update stream state to online
await maubot_plugin.stream_repo.update( await maubot_plugin.stream_repo.update(
StreamState( "stream.logal.dev",
domain="stream.logal.dev", name="Test Stream",
name="Test Stream", title="Playing Games",
title="Playing Games", online=True,
online=True, status_since="2026-01-01T12:00:00+00:00",
status_since="2026-01-01T12:00:00+00:00",
)
) )
await maubot_test_bot.send("!live") await maubot_test_bot.send("!live")
@@ -563,22 +549,18 @@ class TestLiveCommand:
# Set both streams online with different status timestamps # Set both streams online with different status timestamps
await maubot_plugin.stream_repo.update( await maubot_plugin.stream_repo.update(
StreamState( "alpha.com",
domain="alpha.com", name="Alpha Stream",
name="Alpha Stream", title="Morning Show",
title="Morning Show", online=True,
online=True, status_since="2026-03-13T10:00:00+00:00",
status_since="2026-03-13T10:00:00+00:00",
)
) )
await maubot_plugin.stream_repo.update( await maubot_plugin.stream_repo.update(
StreamState( "beta.com",
domain="beta.com", name="Beta Stream",
name="Beta Stream", title="Evening Vibes",
title="Evening Vibes", online=True,
online=True, status_since="2026-03-13T06:00:00+00:00",
status_since="2026-03-13T06:00:00+00:00",
)
) )
await maubot_test_bot.send("!live") await maubot_test_bot.send("!live")
+4 -4
View File
@@ -148,13 +148,13 @@ class TestReadLimitedResponseBody:
class TestGetStreamState: class TestGetStreamState:
"""Stream state retrieval from the status API.""" """Stream observation retrieval from the status API."""
@time_machine.travel(datetime(2026, 3, 13, 12, 0, 0, tzinfo=UTC)) @time_machine.travel(datetime(2026, 3, 13, 12, 0, 0, tzinfo=UTC))
async def test_returns_state_on_success( async def test_returns_state_on_success(
self, owncast_client: OwncastClient self, owncast_client: OwncastClient
) -> None: ) -> None:
"""Return a StreamState with correct fields on a valid 200 response.""" """Return a StreamStateObservation with correct fields on a valid response."""
with aioresponses() as mocked: with aioresponses() as mocked:
mocked.get( mocked.get(
"https://stream.logal.dev/api/status", "https://stream.logal.dev/api/status",
@@ -169,7 +169,7 @@ class TestGetStreamState:
== "I think I can do this... Let's start a nuclear reaction - Playing Nucleares!" # noqa: E501 == "I think I can do this... Let's start a nuclear reaction - Playing Nucleares!" # noqa: E501
) )
assert result.online is False assert result.online is False
assert result.status_since == "2026-03-13T12:00:00+00:00" assert result.observed_at == "2026-03-13T12:00:00+00:00"
async def test_returns_none_on_missing_field( async def test_returns_none_on_missing_field(
self, owncast_client: OwncastClient self, owncast_client: OwncastClient
@@ -272,7 +272,7 @@ class TestGetStreamConfig:
async def test_returns_config_on_success( async def test_returns_config_on_success(
self, owncast_client: OwncastClient self, owncast_client: OwncastClient
) -> None: ) -> None:
"""Return a StreamConfig with correct fields on a valid 200 response.""" """Return a StreamConfigObservation with fields from a valid response."""
with aioresponses() as mocked: with aioresponses() as mocked:
mocked.get( mocked.get(
"https://stream.logal.dev/api/config", "https://stream.logal.dev/api/config",
+77 -20
View File
@@ -23,7 +23,6 @@ from owncastsentry.types import (
UNKNOWN_STATUS_THRESHOLD, UNKNOWN_STATUS_THRESHOLD,
AlreadySubscribedError, AlreadySubscribedError,
NotSubscribedError, NotSubscribedError,
StreamState,
) )
if TYPE_CHECKING: if TYPE_CHECKING:
@@ -129,6 +128,70 @@ class TestStreamDelete:
assert await stream_repo.get_by_domain("example.com") is None assert await stream_repo.get_by_domain("example.com") is None
class TestStreamUpdate:
"""Partial stream state updates."""
async def test_updates_only_supplied_fields(
self, stream_repo: StreamRepository
) -> None:
"""Update supplied fields while leaving omitted fields unchanged."""
await stream_repo.create("example.com")
await stream_repo.update(
"example.com",
name="Original Name",
title="Original Title",
online=True,
status_since="2026-01-01T12:00:00+00:00",
)
await stream_repo.update("example.com", title="Updated Title")
state = await stream_repo.get_by_domain("example.com")
assert state is not None
assert state.name == "Original Name"
assert state.title == "Updated Title"
assert state.online is True
assert state.status_since == "2026-01-01T12:00:00+00:00"
async def test_can_set_nullable_field_to_null(
self, stream_repo: StreamRepository
) -> None:
"""Treat None as a supplied NULL value, not an omitted field."""
await stream_repo.create("example.com")
await stream_repo.update(
"example.com",
name="Original Name",
title="Original Title",
status_since="2026-01-01T12:00:00+00:00",
)
await stream_repo.update("example.com", title=None)
state = await stream_repo.get_by_domain("example.com")
assert state is not None
assert state.name == "Original Name"
assert state.title is None
assert state.status_since == "2026-01-01T12:00:00+00:00"
async def test_accepts_no_supplied_fields(
self, stream_repo: StreamRepository
) -> None:
"""Leave the row unchanged when no partial fields are supplied."""
await stream_repo.create("example.com")
await stream_repo.update(
"example.com",
title="Original Title",
status_since="2026-01-01T12:00:00+00:00",
)
await stream_repo.update("example.com")
state = await stream_repo.get_by_domain("example.com")
assert state is not None
assert state.title == "Original Title"
assert state.status_since == "2026-01-01T12:00:00+00:00"
class TestGetSubscribedStreamsForRoom: class TestGetSubscribedStreamsForRoom:
"""Subscribed stream lookup by room.""" """Subscribed stream lookup by room."""
@@ -186,9 +249,9 @@ class TestGetRoomSubscriptions:
) -> None: ) -> None:
"""Return sorted resolved subscriptions and skip missing stream rows.""" """Return sorted resolved subscriptions and skip missing stream rows."""
await stream_repo.create("beta.example") await stream_repo.create("beta.example")
await stream_repo.update(StreamState(domain="beta.example", name="Beta")) await stream_repo.update("beta.example", name="Beta")
await stream_repo.create("alpha.example") await stream_repo.create("alpha.example")
await stream_repo.update(StreamState(domain="alpha.example", name="Alpha")) await stream_repo.update("alpha.example", name="Alpha")
await subscription_repo.add("beta.example", "!room:example.com") await subscription_repo.add("beta.example", "!room:example.com")
await subscription_repo.add("missing.example", "!room:example.com") await subscription_repo.add("missing.example", "!room:example.com")
await subscription_repo.add("alpha.example", "!room:example.com") await subscription_repo.add("alpha.example", "!room:example.com")
@@ -227,29 +290,23 @@ class TestGetLiveRoomSubscriptions:
"""Return only streams whose derived status is ONLINE.""" """Return only streams whose derived status is ONLINE."""
await stream_repo.create("offline.example") await stream_repo.create("offline.example")
await stream_repo.update( await stream_repo.update(
StreamState( "offline.example",
domain="offline.example", name="Offline",
name="Offline", status_since="2026-01-01T00:00:00+00:00",
status_since="2026-01-01T00:00:00+00:00",
)
) )
await stream_repo.create("online.example") await stream_repo.create("online.example")
await stream_repo.update( await stream_repo.update(
StreamState( "online.example",
domain="online.example", name="Online",
name="Online", online=True,
online=True, status_since="2026-01-01T00:00:00+00:00",
status_since="2026-01-01T00:00:00+00:00",
)
) )
await stream_repo.create("unknown.example") await stream_repo.create("unknown.example")
await stream_repo.update( await stream_repo.update(
StreamState( "unknown.example",
domain="unknown.example", name="Unknown",
name="Unknown", online=True,
online=True, status_since="2026-01-01T00:00:00+00:00",
status_since="2026-01-01T00:00:00+00:00",
)
) )
for _ in range(UNKNOWN_STATUS_THRESHOLD + 1): for _ in range(UNKNOWN_STATUS_THRESHOLD + 1):
await stream_repo.increment_failure_counter("unknown.example") await stream_repo.increment_failure_counter("unknown.example")
+553 -141
View File
@@ -16,9 +16,11 @@
import logging import logging
import time import time
from datetime import datetime
from typing import TYPE_CHECKING from typing import TYPE_CHECKING
import pytest import pytest
import time_machine
from owncastsentry.metrics import MetricsService from owncastsentry.metrics import MetricsService
from owncastsentry.notification_service import ( from owncastsentry.notification_service import (
@@ -38,9 +40,18 @@ from owncastsentry.stream_monitor import (
_QUERY_EVERY_MINUTE_FAILURE_LIMIT, _QUERY_EVERY_MINUTE_FAILURE_LIMIT,
_TEMPORARY_OFFLINE_NOTIFICATION_COOLDOWN, _TEMPORARY_OFFLINE_NOTIFICATION_COOLDOWN,
StreamMonitor, StreamMonitor,
_classify_notification,
_classify_transition,
_NotificationKind,
_should_query_stream, _should_query_stream,
_StreamTransitionKind,
)
from owncastsentry.types import (
StreamConfigObservation,
StreamState,
StreamStateObservation,
StreamStatus,
) )
from owncastsentry.types import StreamConfig, StreamState, StreamStatus
from tests.conftest import ( from tests.conftest import (
_StubMatrixClient, _StubMatrixClient,
_StubOwncastClient, _StubOwncastClient,
@@ -92,17 +103,28 @@ async def _seed_stream(
) -> None: ) -> None:
"""Insert a stream and subscription into the database.""" """Insert a stream and subscription into the database."""
await stream_repo.create(domain) await stream_repo.create(domain)
state = StreamState( await stream_repo.update(
domain=domain, domain,
name=name, name=name,
title=title, title=title,
online=online, online=online,
status_since=status_since, status_since=status_since,
) )
await stream_repo.update(state)
await subscription_repo.add(domain, room_id) await subscription_repo.add(domain, room_id)
def _stream_observation(
domain: str,
*,
observed_at: str,
title: str | None = None,
online: bool = False,
) -> StreamStateObservation:
"""Create a stream observation at a deterministic time."""
with time_machine.travel(datetime.fromisoformat(observed_at), tick=False):
return StreamStateObservation(domain=domain, title=title, online=online)
def _make_monitor_with_metrics( def _make_monitor_with_metrics(
*, *,
owncast_client: _StubOwncastClient, owncast_client: _StubOwncastClient,
@@ -213,6 +235,254 @@ class TestShouldQueryStream:
assert _should_query_stream(counter) == expected assert _should_query_stream(counter) == expected
class TestClassifyTransition:
"""Pure stream state transition classification."""
def test_first_observation(self) -> None:
"""Classify a missing stored timestamp as the first observation."""
old_state = StreamState(domain="example.com")
observation = _stream_observation(
domain="example.com",
title="Live",
online=True,
observed_at="2026-01-01T12:00:00+00:00",
)
transition = _classify_transition(old_state, observation)
assert transition.kind is _StreamTransitionKind.FIRST_OBSERVATION
assert transition.status_since == observation.observed_at
def test_went_live(self) -> None:
"""Classify an offline to online status change."""
old_state = StreamState(
domain="example.com",
title="Title",
status_since="2026-01-01T10:00:00+00:00",
)
observation = _stream_observation(
domain="example.com",
title="Title",
online=True,
observed_at="2026-01-01T12:00:00+00:00",
)
transition = _classify_transition(old_state, observation)
assert transition.kind is _StreamTransitionKind.WENT_LIVE
assert transition.status_since == observation.observed_at
def test_went_offline(self) -> None:
"""Classify an online to offline status change."""
old_state = StreamState(
domain="example.com",
title="Title",
online=True,
status_since="2026-01-01T10:00:00+00:00",
)
observation = _stream_observation(
domain="example.com",
title="Title",
observed_at="2026-01-01T12:00:00+00:00",
)
transition = _classify_transition(old_state, observation)
assert transition.kind is _StreamTransitionKind.WENT_OFFLINE
assert transition.status_since == observation.observed_at
def test_online_title_changed(self) -> None:
"""Classify an online title change without resetting status duration."""
old_state = StreamState(
domain="example.com",
title="Old",
online=True,
status_since="2026-01-01T10:00:00+00:00",
)
observation = _stream_observation(
domain="example.com",
title="New",
online=True,
observed_at="2026-01-01T12:00:00+00:00",
)
transition = _classify_transition(old_state, observation)
assert transition.kind is _StreamTransitionKind.TITLE_CHANGED
assert transition.status_since == old_state.status_since
def test_unchanged(self) -> None:
"""Classify identical status and title as unchanged."""
old_state = StreamState(
domain="example.com",
title="Same",
status_since="2026-01-01T10:00:00+00:00",
)
observation = _stream_observation(
domain="example.com",
title="Same",
observed_at="2026-01-01T12:00:00+00:00",
)
transition = _classify_transition(old_state, observation)
assert transition.kind is _StreamTransitionKind.STATUS_UNCHANGED
assert transition.status_since == old_state.status_since
def test_offline_title_change_keeps_status_timestamp(self) -> None:
"""Keep offline status duration when only the title changes."""
old_state = StreamState(
domain="example.com",
title="Old",
status_since="2026-01-01T10:00:00+00:00",
)
observation = _stream_observation(
domain="example.com",
title="New",
observed_at="2026-01-01T12:00:00+00:00",
)
transition = _classify_transition(old_state, observation)
assert transition.kind is _StreamTransitionKind.STATUS_UNCHANGED
assert transition.status_since == old_state.status_since
class TestClassifyNotification:
"""Live/title notification policy classification."""
def test_first_observation_suppresses_notification(self) -> None:
"""Do not notify for the first observed stream state."""
old_state = StreamState(domain="example.com")
observation = _stream_observation(
domain="example.com",
online=True,
observed_at="2026-01-01T12:00:00+00:00",
)
transition = _classify_transition(old_state, observation)
notification_kind = _classify_notification(transition, old_state, observation)
assert notification_kind is _NotificationKind.NONE
def test_normal_went_live_notifies_live(self) -> None:
"""Notify live after a non-brief offline gap."""
old_state = StreamState(
domain="example.com",
title="Title",
status_since="2026-01-01T10:00:00+00:00",
)
observation = _stream_observation(
domain="example.com",
title="Title",
online=True,
observed_at="2026-01-01T12:00:00+00:00",
)
transition = _classify_transition(old_state, observation)
notification_kind = _classify_notification(
transition,
old_state,
observation,
offline_duration_seconds=_TEMPORARY_OFFLINE_NOTIFICATION_COOLDOWN + 1,
)
assert notification_kind is _NotificationKind.LIVE
def test_brief_outage_same_title_suppresses_notification(self) -> None:
"""Do not notify when a stream briefly returns with the same title."""
old_state = StreamState(
domain="example.com",
title="Same",
status_since="2026-01-01T10:00:00+00:00",
)
observation = _stream_observation(
domain="example.com",
title="Same",
online=True,
observed_at="2026-01-01T12:00:00+00:00",
)
transition = _classify_transition(old_state, observation)
notification_kind = _classify_notification(
transition,
old_state,
observation,
offline_duration_seconds=60,
)
assert notification_kind is _NotificationKind.NONE
def test_brief_outage_title_change_notifies_title_change(self) -> None:
"""Notify title change when a brief return changes title."""
old_state = StreamState(
domain="example.com",
title="Old",
status_since="2026-01-01T10:00:00+00:00",
)
observation = _stream_observation(
domain="example.com",
title="New",
online=True,
observed_at="2026-01-01T12:00:00+00:00",
)
transition = _classify_transition(old_state, observation)
notification_kind = _classify_notification(
transition,
old_state,
observation,
offline_duration_seconds=60,
)
assert notification_kind is _NotificationKind.TITLE_CHANGE
def test_mid_session_title_change_notifies_title_change(self) -> None:
"""Notify title change for a normal online title change."""
old_state = StreamState(
domain="example.com",
title="Old",
online=True,
status_since="2026-01-01T10:00:00+00:00",
)
observation = _stream_observation(
domain="example.com",
title="New",
online=True,
observed_at="2026-01-01T12:00:00+00:00",
)
transition = _classify_transition(old_state, observation)
notification_kind = _classify_notification(transition, old_state, observation)
assert notification_kind is _NotificationKind.TITLE_CHANGE
def test_went_offline_after_last_notification_notifies_live(self) -> None:
"""Notify live when stream went offline after the last notification."""
old_state = StreamState(
domain="example.com",
title="Old",
online=True,
status_since="2026-01-01T10:00:00+00:00",
)
observation = _stream_observation(
domain="example.com",
title="New",
online=True,
observed_at="2026-01-01T12:00:00+00:00",
)
transition = _classify_transition(old_state, observation)
notification_kind = _classify_notification(
transition,
old_state,
observation,
went_offline_after_last_notification=True,
)
assert notification_kind is _NotificationKind.LIVE
class TestUpdateAllStreams: class TestUpdateAllStreams:
"""Parallel stream update orchestration.""" """Parallel stream update orchestration."""
@@ -223,12 +493,12 @@ class TestUpdateAllStreams:
) -> None: ) -> None:
"""Return an UpdateResult with correct success and failure counts.""" """Return an UpdateResult with correct success and failure counts."""
owncast = _StubOwncastClient( owncast = _StubOwncastClient(
stream_state=StreamState( stream_state=_stream_observation(
domain="ok.com", domain="ok.com",
online=True, online=True,
status_since="2026-01-01T00:00:00+00:00", observed_at="2026-01-01T00:00:00+00:00",
), ),
stream_config=StreamConfig(name="OK"), stream_config=StreamConfigObservation(name="OK"),
) )
client = _StubMatrixClient() client = _StubMatrixClient()
monitor, _ = _make_monitor( monitor, _ = _make_monitor(
@@ -240,9 +510,9 @@ class TestUpdateAllStreams:
await _seed_stream(stream_repo, subscription_repo, domain="ok.com") await _seed_stream(stream_repo, subscription_repo, domain="ok.com")
result = await monitor.update_all_streams(["ok.com"]) result = await monitor.update_all_streams(["ok.com", "missing.com"])
assert result.total_streams == 1 assert result.total_streams == 2
assert result.successful_checks == 1 assert result.successful_checks == 2
assert result.failed_checks == 0 assert result.failed_checks == 0
assert owncast.queried_domains == ["ok.com"] assert owncast.queried_domains == ["ok.com"]
@@ -274,6 +544,7 @@ class TestUpdateStreamBackoff:
result = await monitor.update_stream("fail.com") result = await monitor.update_stream("fail.com")
assert result is True assert result is True
assert owncast.state_call_count == 0 assert owncast.state_call_count == 0
assert owncast.config_call_count == 0
# Counter should have been incremented to 6 # Counter should have been incremented to 6
state = await stream_repo.get_by_domain("fail.com") state = await stream_repo.get_by_domain("fail.com")
@@ -291,12 +562,12 @@ class TestUpdateStreamFirstUpdate:
) -> None: ) -> None:
"""Skip notifications when a stream is already live on first update.""" """Skip notifications when a stream is already live on first update."""
owncast = _StubOwncastClient( owncast = _StubOwncastClient(
stream_state=StreamState( stream_state=_stream_observation(
domain="new.com", domain="new.com",
online=True, online=True,
status_since="2026-01-01T00:00:00+00:00", observed_at="2026-01-01T12:01:00+00:00",
), ),
stream_config=StreamConfig(name="New Stream"), stream_config=StreamConfigObservation(name="New Stream"),
) )
client = _StubMatrixClient() client = _StubMatrixClient()
monitor, _ = _make_monitor( monitor, _ = _make_monitor(
@@ -313,6 +584,45 @@ class TestUpdateStreamFirstUpdate:
result = await monitor.update_stream("new.com") result = await monitor.update_stream("new.com")
assert result is True assert result is True
assert len(client.sent_messages) == 0 assert len(client.sent_messages) == 0
assert owncast.config_call_count == 1
async def test_fetches_config_and_saves_name_on_first_update(
self,
stream_repo: StreamRepository,
subscription_repo: SubscriptionRepository,
) -> None:
"""Fetch config on first update and persist the stream display name."""
owncast = _StubOwncastClient(
stream_state=_stream_observation(
domain="new.com",
title="First Title",
online=True,
observed_at="2026-01-01T12:01:00+00:00",
),
stream_config=StreamConfigObservation(name="New Stream"),
)
client = _StubMatrixClient()
monitor, _ = _make_monitor(
owncast_client=owncast,
stream_repo=stream_repo,
subscription_repo=subscription_repo,
client=client,
)
await stream_repo.create("new.com")
await subscription_repo.add("new.com", "!room:matrix.org")
result = await monitor.update_stream("new.com")
assert result is True
assert len(client.sent_messages) == 0
assert owncast.config_call_count == 1
state = await stream_repo.get_by_domain("new.com")
assert state is not None
assert state.name == "New Stream"
assert state.title == "First Title"
assert state.online is True
assert state.status_since == "2026-01-01T12:01:00+00:00"
async def test_no_notification_on_first_offline( async def test_no_notification_on_first_offline(
self, self,
@@ -321,11 +631,11 @@ class TestUpdateStreamFirstUpdate:
) -> None: ) -> None:
"""Skip notifications when a stream is offline on first update.""" """Skip notifications when a stream is offline on first update."""
owncast = _StubOwncastClient( owncast = _StubOwncastClient(
stream_state=StreamState( stream_state=_stream_observation(
domain="new.com", domain="new.com",
status_since="2025-12-31T00:00:00+00:00", observed_at="2025-12-31T12:01:00+00:00",
), ),
stream_config=StreamConfig(name="New Stream"), stream_config=StreamConfigObservation(name="New Stream"),
) )
client = _StubMatrixClient() client = _StubMatrixClient()
monitor, _ = _make_monitor( monitor, _ = _make_monitor(
@@ -341,6 +651,7 @@ class TestUpdateStreamFirstUpdate:
result = await monitor.update_stream("new.com") result = await monitor.update_stream("new.com")
assert result is True assert result is True
assert len(client.sent_messages) == 0 assert len(client.sent_messages) == 0
assert owncast.config_call_count == 1
class TestUpdateStreamGoesLive: class TestUpdateStreamGoesLive:
@@ -353,13 +664,13 @@ class TestUpdateStreamGoesLive:
) -> None: ) -> None:
"""Send go-live with stream name and tags.""" """Send go-live with stream name and tags."""
owncast = _StubOwncastClient( owncast = _StubOwncastClient(
stream_state=StreamState( stream_state=_stream_observation(
domain="live.com", domain="live.com",
title="Now Streaming", title="Now Streaming",
online=True, online=True,
status_since="2026-01-01T12:00:00+00:00", observed_at="2026-01-01T12:00:00+00:00",
), ),
stream_config=StreamConfig(name="Live Stream", tags=("gaming",)), stream_config=StreamConfigObservation(name="Live Stream", tags=("gaming",)),
) )
client = _StubMatrixClient() client = _StubMatrixClient()
monitor, _ = _make_monitor( monitor, _ = _make_monitor(
@@ -393,6 +704,7 @@ class TestUpdateStreamGoesLive:
"\n" "\n"
"#gaming" "#gaming"
) )
assert owncast.config_call_count == 1
async def test_falls_back_when_config_fetch_fails( async def test_falls_back_when_config_fetch_fails(
self, self,
@@ -401,11 +713,11 @@ class TestUpdateStreamGoesLive:
) -> None: ) -> None:
"""Use domain as name in go-live when config fetch fails.""" """Use domain as name in go-live when config fetch fails."""
owncast = _StubOwncastClient( owncast = _StubOwncastClient(
stream_state=StreamState( stream_state=_stream_observation(
domain="live.com", domain="live.com",
title="Now Streaming", title="Now Streaming",
online=True, online=True,
status_since="2026-01-01T12:00:00+00:00", observed_at="2026-01-01T12:00:00+00:00",
), ),
stream_config=None, stream_config=None,
) )
@@ -437,6 +749,7 @@ class TestUpdateStreamGoesLive:
"\n" "\n"
"To tune in, visit: https://live.com/" "To tune in, visit: https://live.com/"
) )
assert owncast.config_call_count == 1
class TestUpdateStreamBriefOffline: class TestUpdateStreamBriefOffline:
@@ -449,13 +762,13 @@ class TestUpdateStreamBriefOffline:
) -> None: ) -> None:
"""Skip notification for a brief outage with no title change.""" """Skip notification for a brief outage with no title change."""
owncast = _StubOwncastClient( owncast = _StubOwncastClient(
stream_state=StreamState( stream_state=_stream_observation(
domain="brief.com", domain="brief.com",
title="Same Title", title="Same Title",
online=True, online=True,
status_since="2026-01-01T12:00:00+00:00", observed_at="2026-01-01T12:01:00+00:00",
), ),
stream_config=StreamConfig(name="Brief Stream"), stream_config=StreamConfigObservation(name="Brief Stream"),
) )
client = _StubMatrixClient() client = _StubMatrixClient()
monitor, _ = _make_monitor( monitor, _ = _make_monitor(
@@ -480,6 +793,11 @@ class TestUpdateStreamBriefOffline:
result = await monitor.update_stream("brief.com") result = await monitor.update_stream("brief.com")
assert result is True assert result is True
assert len(client.sent_messages) == 0 assert len(client.sent_messages) == 0
assert owncast.config_call_count == 0
state = await stream_repo.get_by_domain("brief.com")
assert state is not None
assert state.online is True
assert state.status_since == "2026-01-01T12:01:00+00:00"
async def test_title_change_notification( async def test_title_change_notification(
self, self,
@@ -488,13 +806,13 @@ class TestUpdateStreamBriefOffline:
) -> None: ) -> None:
"""Send a title-change notification for a brief outage with a new title.""" """Send a title-change notification for a brief outage with a new title."""
owncast = _StubOwncastClient( owncast = _StubOwncastClient(
stream_state=StreamState( stream_state=_stream_observation(
domain="brief.com", domain="brief.com",
title="New Title", title="New Title",
online=True, online=True,
status_since="2026-01-01T12:00:00+00:00", observed_at="2026-01-01T12:00:00+00:00",
), ),
stream_config=StreamConfig(name="Brief Stream"), stream_config=StreamConfigObservation(name="Brief Stream"),
) )
client = _StubMatrixClient() client = _StubMatrixClient()
monitor, _ = _make_monitor( monitor, _ = _make_monitor(
@@ -525,6 +843,10 @@ class TestUpdateStreamBriefOffline:
"\n" "\n"
"To tune in, visit: https://brief.com/" "To tune in, visit: https://brief.com/"
) )
assert owncast.config_call_count == 1
state = await stream_repo.get_by_domain("brief.com")
assert state is not None
assert state.name == "Brief Stream"
class TestUpdateStreamTitleChange: class TestUpdateStreamTitleChange:
@@ -537,11 +859,11 @@ class TestUpdateStreamTitleChange:
) -> None: ) -> None:
"""Use domain as name in title-change when config fails.""" """Use domain as name in title-change when config fails."""
owncast = _StubOwncastClient( owncast = _StubOwncastClient(
stream_state=StreamState( stream_state=_stream_observation(
domain="title.com", domain="title.com",
title="Updated Title", title="Updated Title",
online=True, online=True,
status_since="2026-01-01T12:00:00+00:00", observed_at="2026-01-01T12:00:00+00:00",
), ),
stream_config=None, stream_config=None,
) )
@@ -580,6 +902,7 @@ class TestUpdateStreamTitleChange:
"\n" "\n"
"To tune in, visit: https://title.com/" "To tune in, visit: https://title.com/"
) )
assert owncast.config_call_count == 1
async def test_sends_title_change( async def test_sends_title_change(
self, self,
@@ -588,13 +911,13 @@ class TestUpdateStreamTitleChange:
) -> None: ) -> None:
"""Send a title-change notification when the title changes mid-stream.""" """Send a title-change notification when the title changes mid-stream."""
owncast = _StubOwncastClient( owncast = _StubOwncastClient(
stream_state=StreamState( stream_state=_stream_observation(
domain="title.com", domain="title.com",
title="Updated Title", title="Updated Title",
online=True, online=True,
status_since="2026-01-01T12:00:00+00:00", observed_at="2026-01-01T12:00:00+00:00",
), ),
stream_config=StreamConfig(name="Title Stream"), stream_config=StreamConfigObservation(name="Title Stream"),
) )
client = _StubMatrixClient() client = _StubMatrixClient()
monitor, notification_service = _make_monitor( monitor, notification_service = _make_monitor(
@@ -634,6 +957,7 @@ class TestUpdateStreamTitleChange:
"\n" "\n"
"To tune in, visit: https://title.com/" "To tune in, visit: https://title.com/"
) )
assert owncast.config_call_count == 1
async def test_sends_go_live_when_last_notification_before_offline( async def test_sends_go_live_when_last_notification_before_offline(
self, self,
@@ -642,13 +966,13 @@ class TestUpdateStreamTitleChange:
) -> None: ) -> None:
"""Send go-live instead of title-change after an offline gap.""" """Send go-live instead of title-change after an offline gap."""
owncast = _StubOwncastClient( owncast = _StubOwncastClient(
stream_state=StreamState( stream_state=_stream_observation(
domain="title.com", domain="title.com",
title="Updated Title", title="Updated Title",
online=True, online=True,
status_since="2026-01-01T12:00:00+00:00", observed_at="2026-01-01T12:00:00+00:00",
), ),
stream_config=StreamConfig(name="Title Stream"), stream_config=StreamConfigObservation(name="Title Stream"),
) )
client = _StubMatrixClient() client = _StubMatrixClient()
monitor, notification_service = _make_monitor( monitor, notification_service = _make_monitor(
@@ -686,6 +1010,7 @@ class TestUpdateStreamTitleChange:
"\n" "\n"
"To tune in, visit: https://title.com/" "To tune in, visit: https://title.com/"
) )
assert owncast.config_call_count == 1
class TestUpdateStreamGoesOffline: class TestUpdateStreamGoesOffline:
@@ -698,12 +1023,12 @@ class TestUpdateStreamGoesOffline:
) -> None: ) -> None:
"""Record the offline time in the cache when a stream goes offline.""" """Record the offline time in the cache when a stream goes offline."""
owncast = _StubOwncastClient( owncast = _StubOwncastClient(
stream_state=StreamState( stream_state=_stream_observation(
domain="offline.com", domain="offline.com",
title="Title", title="Title",
status_since="2026-01-01T12:00:00+00:00", observed_at="2026-01-01T12:01:00+00:00",
), ),
stream_config=StreamConfig(name="Offline Stream"), stream_config=StreamConfigObservation(name="Offline Stream"),
) )
client = _StubMatrixClient() client = _StubMatrixClient()
monitor, _ = _make_monitor( monitor, _ = _make_monitor(
@@ -730,6 +1055,180 @@ class TestUpdateStreamGoesOffline:
assert result is True assert result is True
assert before <= monitor.offline_timer_cache["offline.com"] <= after assert before <= monitor.offline_timer_cache["offline.com"] <= after
assert len(client.sent_messages) == 0 assert len(client.sent_messages) == 0
assert owncast.config_call_count == 0
state = await stream_repo.get_by_domain("offline.com")
assert state is not None
assert state.online is False
assert state.name == "Test Stream"
class TestUpdateStreamNoStateChange:
"""Offline stream persistence and config refresh behavior."""
async def test_does_not_fetch_config_for_unchanged_non_hourly_state(
self,
stream_repo: StreamRepository,
subscription_repo: SubscriptionRepository,
) -> None:
"""Do not fetch config when state is unchanged outside hourly refresh."""
owncast = _StubOwncastClient(
stream_state=_stream_observation(
domain="stable.com",
title="Same Title",
observed_at="2026-01-01T12:01:00+00:00",
),
)
client = _StubMatrixClient()
monitor, _ = _make_monitor(
owncast_client=owncast,
stream_repo=stream_repo,
subscription_repo=subscription_repo,
client=client,
)
# Seed as offline with same status timestamp and title.
await _seed_stream(
stream_repo,
subscription_repo,
domain="stable.com",
title="Same Title",
status_since="2026-01-01T12:00:00+00:00",
)
result = await monitor.update_stream("fail.com")
assert result is True
assert len(client.sent_messages) == 0
# Config should not have been fetched since this was not an hourly refresh.
assert owncast.config_call_count == 0
async def test_refreshes_config_at_top_of_hour(
self,
stream_repo: StreamRepository,
subscription_repo: SubscriptionRepository,
) -> None:
"""Refresh and persist stream config during the hourly refresh window."""
owncast = _StubOwncastClient(
stream_state=_stream_observation(
domain="stable.com",
title="Same Title",
observed_at="2026-01-01T12:00:30+00:00",
),
stream_config=StreamConfigObservation(name="Updated Stream"),
)
client = _StubMatrixClient()
monitor, _ = _make_monitor(
owncast_client=owncast,
stream_repo=stream_repo,
subscription_repo=subscription_repo,
client=client,
)
await _seed_stream(
stream_repo,
subscription_repo,
domain="stable.com",
name="Old Stream",
title="Same Title",
status_since="2026-01-01T10:00:00+00:00",
)
result = await monitor.update_stream("stable.com")
assert result is True
assert len(client.sent_messages) == 0
assert owncast.config_call_count == 1
state = await stream_repo.get_by_domain("stable.com")
assert state is not None
assert state.name == "Updated Stream"
assert state.title == "Same Title"
assert state.status_since == "2026-01-01T10:00:00+00:00"
async def test_saves_offline_title_change(
self,
stream_repo: StreamRepository,
subscription_repo: SubscriptionRepository,
) -> None:
"""Save title-only changes while the stream remains offline."""
owncast = _StubOwncastClient(
stream_state=_stream_observation(
domain="stable.com",
title="New Offline Title",
observed_at="2026-01-01T12:30:00+00:00",
),
)
client = _StubMatrixClient()
monitor, _ = _make_monitor(
owncast_client=owncast,
stream_repo=stream_repo,
subscription_repo=subscription_repo,
client=client,
)
await _seed_stream(
stream_repo,
subscription_repo,
domain="stable.com",
title="Old Offline Title",
status_since="2026-01-01T12:00:00+00:00",
)
result = await monitor.update_stream("stable.com")
assert result is True
assert len(client.sent_messages) == 0
assert owncast.config_call_count == 0
state = await stream_repo.get_by_domain("stable.com")
assert state is not None
assert state.title == "New Offline Title"
assert state.status_since == "2026-01-01T12:00:00+00:00"
class TestUpdateStreamFailureCounterReset:
"""Failure counter reset on successful fetch."""
async def test_resets_counter_on_success(
self,
stream_repo: StreamRepository,
subscription_repo: SubscriptionRepository,
) -> None:
"""Reset the failure counter to zero after a successful fetch."""
owncast = _StubOwncastClient(
stream_state=_stream_observation(
domain="recover.com",
title="Title",
observed_at="2026-01-01T12:00:00+00:00",
),
)
client = _StubMatrixClient()
monitor, _ = _make_monitor(
owncast_client=owncast,
stream_repo=stream_repo,
subscription_repo=subscription_repo,
client=client,
)
await _seed_stream(
stream_repo,
subscription_repo,
domain="recover.com",
title="Title",
status_since="2026-01-01T12:00:00+00:00",
)
# Simulate prior failures (counter=4 still passes backoff)
for _ in range(4):
await stream_repo.increment_failure_counter("recover.com")
state = await stream_repo.get_by_domain("recover.com")
assert state is not None
assert state.failure_counter == 4
result = await monitor.update_stream("recover.com")
assert result is True
state = await stream_repo.get_by_domain("recover.com")
assert state is not None
assert state.failure_counter == 0
class TestUpdateStreamConnectionFailure: class TestUpdateStreamConnectionFailure:
@@ -759,6 +1258,7 @@ class TestUpdateStreamConnectionFailure:
result = await monitor.update_stream("fail.com") result = await monitor.update_stream("fail.com")
assert result is False assert result is False
assert owncast.config_call_count == 0
state = await stream_repo.get_by_domain("fail.com") state = await stream_repo.get_by_domain("fail.com")
assert state is not None assert state is not None
@@ -863,94 +1363,6 @@ class TestCheckCleanupThresholds:
assert await stream_repo.get_by_domain("ok.com") is not None assert await stream_repo.get_by_domain("ok.com") is not None
class TestUpdateStreamNoStateChange:
"""Stream that stays offline across updates with no state change."""
async def test_skips_database_write(
self,
stream_repo: StreamRepository,
subscription_repo: SubscriptionRepository,
) -> None:
"""Skip the database write when stream state has not changed."""
owncast = _StubOwncastClient(
stream_state=StreamState(
domain="stable.com",
title="Same Title",
status_since="2026-01-01T12:00:00+00:00",
),
)
client = _StubMatrixClient()
monitor, _ = _make_monitor(
owncast_client=owncast,
stream_repo=stream_repo,
subscription_repo=subscription_repo,
client=client,
)
# Seed as offline with same status timestamp and title.
await _seed_stream(
stream_repo,
subscription_repo,
domain="stable.com",
title="Same Title",
status_since="2026-01-01T12:00:00+00:00",
)
result = await monitor.update_stream("stable.com")
assert result is True
assert len(client.sent_messages) == 0
# Config should not have been fetched since no DB update was needed
assert owncast.config_call_count == 0
class TestUpdateStreamFailureCounterReset:
"""Failure counter reset on successful fetch."""
async def test_resets_counter_on_success(
self,
stream_repo: StreamRepository,
subscription_repo: SubscriptionRepository,
) -> None:
"""Reset the failure counter to zero after a successful fetch."""
owncast = _StubOwncastClient(
stream_state=StreamState(
domain="recover.com",
title="Title",
status_since="2026-01-01T12:00:00+00:00",
),
)
client = _StubMatrixClient()
monitor, _ = _make_monitor(
owncast_client=owncast,
stream_repo=stream_repo,
subscription_repo=subscription_repo,
client=client,
)
await _seed_stream(
stream_repo,
subscription_repo,
domain="recover.com",
title="Title",
status_since="2026-01-01T12:00:00+00:00",
)
# Simulate prior failures (counter=4 still passes backoff)
for _ in range(4):
await stream_repo.increment_failure_counter("recover.com")
state = await stream_repo.get_by_domain("recover.com")
assert state is not None
assert state.failure_counter == 4
result = await monitor.update_stream("recover.com")
assert result is True
state = await stream_repo.get_by_domain("recover.com")
assert state is not None
assert state.failure_counter == 0
class TestUpdateAllStreamsMixed: class TestUpdateAllStreamsMixed:
"""Mixed success/failure results from parallel updates.""" """Mixed success/failure results from parallel updates."""
@@ -1001,16 +1413,16 @@ class TestUpdateAllStreamsMixed:
) -> None: ) -> None:
"""Count an unhandled exception as a failure without aborting other updates.""" """Count an unhandled exception as a failure without aborting other updates."""
owncast = _StubOwncastClient( owncast = _StubOwncastClient(
stream_state=StreamState( stream_state=_stream_observation(
domain="ok.com", domain="ok.com",
title="Title", title="Title",
status_since="2026-01-01T00:00:00+00:00", observed_at="2026-01-01T00:00:00+00:00",
), ),
) )
# Patch get_stream_state to raise for one specific domain # Patch get_stream_state to raise for one specific domain
original_get_state = owncast.get_stream_state original_get_state = owncast.get_stream_state
async def raising_get_state(domain: str) -> StreamState | None: async def raising_get_state(domain: str) -> StreamStateObservation | None:
if domain == "raise.com": if domain == "raise.com":
msg = "simulated DB error" msg = "simulated DB error"
raise RuntimeError(msg) raise RuntimeError(msg)
@@ -1057,13 +1469,13 @@ class TestStreamMonitorMetrics:
) -> None: ) -> None:
"""Record online stream status gauge.""" """Record online stream status gauge."""
owncast = _StubOwncastClient( owncast = _StubOwncastClient(
stream_state=StreamState( stream_state=_stream_observation(
domain="live.com", domain="live.com",
title="Title", title="Title",
online=True, online=True,
status_since="2026-01-01T12:00:00+00:00", observed_at="2026-01-01T12:00:00+00:00",
), ),
stream_config=StreamConfig(name="Live Stream"), stream_config=StreamConfigObservation(name="Live Stream"),
) )
client = _StubMatrixClient() client = _StubMatrixClient()
monitor, _, metrics = _make_monitor_with_metrics( monitor, _, metrics = _make_monitor_with_metrics(
@@ -1090,12 +1502,12 @@ class TestStreamMonitorMetrics:
) -> None: ) -> None:
"""Record offline stream status gauge.""" """Record offline stream status gauge."""
owncast = _StubOwncastClient( owncast = _StubOwncastClient(
stream_state=StreamState( stream_state=_stream_observation(
domain="off.com", domain="off.com",
title="Title", title="Title",
status_since="2026-01-01T12:00:00+00:00", observed_at="2026-01-01T12:00:00+00:00",
), ),
stream_config=StreamConfig(name="Off Stream"), stream_config=StreamConfigObservation(name="Off Stream"),
) )
client = _StubMatrixClient() client = _StubMatrixClient()
monitor, _, metrics = _make_monitor_with_metrics( monitor, _, metrics = _make_monitor_with_metrics(
@@ -1146,12 +1558,12 @@ class TestStreamMonitorMetrics:
) -> None: ) -> None:
"""Reset failure counter gauge to zero on successful check.""" """Reset failure counter gauge to zero on successful check."""
owncast = _StubOwncastClient( owncast = _StubOwncastClient(
stream_state=StreamState( stream_state=_stream_observation(
domain="recover.com", domain="recover.com",
title="Title", title="Title",
status_since="2026-01-01T12:00:00+00:00", observed_at="2026-01-01T12:00:00+00:00",
), ),
stream_config=StreamConfig(name="Recover"), stream_config=StreamConfigObservation(name="Recover"),
) )
client = _StubMatrixClient() client = _StubMatrixClient()
monitor, _, metrics = _make_monitor_with_metrics( monitor, _, metrics = _make_monitor_with_metrics(
@@ -1228,12 +1640,12 @@ class TestStreamMonitorMetrics:
) -> None: ) -> None:
"""Record per-domain subscription counts after update_all_streams.""" """Record per-domain subscription counts after update_all_streams."""
owncast = _StubOwncastClient( owncast = _StubOwncastClient(
stream_state=StreamState( stream_state=_stream_observation(
domain="pop.com", domain="pop.com",
online=True, online=True,
status_since="2026-01-01T00:00:00+00:00", observed_at="2026-01-01T00:00:00+00:00",
), ),
stream_config=StreamConfig(name="Popular"), stream_config=StreamConfigObservation(name="Popular"),
) )
client = _StubMatrixClient() client = _StubMatrixClient()
monitor, _, metrics = _make_monitor_with_metrics( monitor, _, metrics = _make_monitor_with_metrics(
+48 -55
View File
@@ -25,13 +25,47 @@ from owncastsentry.types import (
AlreadySubscribedError, AlreadySubscribedError,
InvalidOwncastInstanceError, InvalidOwncastInstanceError,
NotSubscribedError, NotSubscribedError,
StreamState,
) )
if TYPE_CHECKING: if TYPE_CHECKING:
from owncastsentry.repository import StreamRepository, SubscriptionRepository from owncastsentry.repository import StreamRepository, SubscriptionRepository
class _StubOwncastClient:
"""Owncast client stub for validation-only manager tests."""
def __init__(self, *, valid: bool = True) -> None:
"""Initialize the stub with a fixed validation result."""
self.valid = valid
self.validated_domains: list[str] = []
async def validate_instance(self, domain: str) -> bool:
"""Record the domain and return the configured validation result."""
self.validated_domains.append(domain)
return self.valid
@pytest.fixture
def owncast_client() -> _StubOwncastClient:
"""Return a validation-only Owncast client stub."""
return _StubOwncastClient()
@pytest.fixture
def manager(
owncast_client: _StubOwncastClient,
stream_repo: StreamRepository,
subscription_repo: SubscriptionRepository,
) -> SubscriptionManager:
"""SubscriptionManager built directly for unit tests."""
return SubscriptionManager(
owncast_client=owncast_client, # type: ignore[arg-type]
stream_repo=stream_repo,
subscription_repo=subscription_repo,
logger=logging.getLogger("test"),
)
class TestDomainify: class TestDomainify:
"""Domain extraction and sanitization from user input.""" """Domain extraction and sanitization from user input."""
@@ -74,41 +108,6 @@ class TestDomainify:
assert _domainify(input_url) == expected assert _domainify(input_url) == expected
class _StubOwncastClient:
"""Owncast client stub for validation-only manager tests."""
def __init__(self, *, valid: bool = True) -> None:
"""Initialize the stub with a fixed validation result."""
self.valid = valid
self.validated_domains: list[str] = []
async def validate_instance(self, domain: str) -> bool:
"""Record the domain and return the configured validation result."""
self.validated_domains.append(domain)
return self.valid
@pytest.fixture
def owncast_client() -> _StubOwncastClient:
"""Return a validation-only Owncast client stub."""
return _StubOwncastClient()
@pytest.fixture
def manager(
owncast_client: _StubOwncastClient,
stream_repo: StreamRepository,
subscription_repo: SubscriptionRepository,
) -> SubscriptionManager:
"""SubscriptionManager built directly for unit tests."""
return SubscriptionManager(
owncast_client=owncast_client, # type: ignore[arg-type]
stream_repo=stream_repo,
subscription_repo=subscription_repo,
logger=logging.getLogger("test"),
)
class TestManagerSubscribe: class TestManagerSubscribe:
"""SubscriptionManager subscribe workflow.""" """SubscriptionManager subscribe workflow."""
@@ -226,9 +225,9 @@ class TestManagerListings:
) -> None: ) -> None:
"""Return sorted room subscriptions and skip missing stream rows.""" """Return sorted room subscriptions and skip missing stream rows."""
await stream_repo.create("beta.example") await stream_repo.create("beta.example")
await stream_repo.update(StreamState(domain="beta.example", name="Beta")) await stream_repo.update("beta.example", name="Beta")
await stream_repo.create("alpha.example") await stream_repo.create("alpha.example")
await stream_repo.update(StreamState(domain="alpha.example", name="Alpha")) await stream_repo.update("alpha.example", name="Alpha")
await subscription_repo.add("beta.example", "!room:example.com") await subscription_repo.add("beta.example", "!room:example.com")
await subscription_repo.add("missing.example", "!room:example.com") await subscription_repo.add("missing.example", "!room:example.com")
await subscription_repo.add("alpha.example", "!room:example.com") await subscription_repo.add("alpha.example", "!room:example.com")
@@ -253,29 +252,23 @@ class TestManagerListings:
"""Live listing includes only subscriptions with ONLINE stream status.""" """Live listing includes only subscriptions with ONLINE stream status."""
await stream_repo.create("offline.example") await stream_repo.create("offline.example")
await stream_repo.update( await stream_repo.update(
StreamState( "offline.example",
domain="offline.example", name="Offline",
name="Offline", status_since="2026-01-01T00:00:00+00:00",
status_since="2026-01-01T00:00:00+00:00",
)
) )
await stream_repo.create("online.example") await stream_repo.create("online.example")
await stream_repo.update( await stream_repo.update(
StreamState( "online.example",
domain="online.example", name="Online",
name="Online", online=True,
online=True, status_since="2026-01-01T00:00:00+00:00",
status_since="2026-01-01T00:00:00+00:00",
)
) )
await stream_repo.create("unknown.example") await stream_repo.create("unknown.example")
await stream_repo.update( await stream_repo.update(
StreamState( "unknown.example",
domain="unknown.example", name="Unknown",
name="Unknown", online=True,
online=True, status_since="2026-01-01T00:00:00+00:00",
status_since="2026-01-01T00:00:00+00:00",
)
) )
for _ in range(UNKNOWN_STATUS_THRESHOLD + 1): for _ in range(UNKNOWN_STATUS_THRESHOLD + 1):
await stream_repo.increment_failure_counter("unknown.example") await stream_repo.increment_failure_counter("unknown.example")
+50 -32
View File
@@ -18,6 +18,7 @@ from dataclasses import FrozenInstanceError
from datetime import UTC, datetime, timedelta, timezone from datetime import UTC, datetime, timedelta, timezone
import pytest import pytest
import time_machine
from owncastsentry.types import ( from owncastsentry.types import (
_MAX_INSTANCE_TITLE_LENGTH, _MAX_INSTANCE_TITLE_LENGTH,
@@ -29,8 +30,9 @@ from owncastsentry.types import (
InvalidOwncastInstanceError, InvalidOwncastInstanceError,
NotSubscribedError, NotSubscribedError,
RoomSubscription, RoomSubscription,
StreamConfig, StreamConfigObservation,
StreamState, StreamState,
StreamStateObservation,
StreamStatus, StreamStatus,
SubscriptionError, SubscriptionError,
UpdateResult, UpdateResult,
@@ -140,29 +142,29 @@ class TestFormatStatusSince:
assert format_status_since(timestamp) == "2026-01-01T12:34:56+00:00" assert format_status_since(timestamp) == "2026-01-01T12:34:56+00:00"
class TestStreamStateFromApiResponse: class TestStreamStateObservationFromApiResponse:
"""StreamState construction from an API response dictionary.""" """StreamStateObservation construction from an API response dictionary."""
def test_typical_response(self) -> None: def test_typical_response(self) -> None:
"""Populate API-derived fields from a complete stream state response.""" """Populate API-derived fields from a complete status response."""
observed_at = datetime(2026, 1, 1, 0, 0, 1, tzinfo=UTC) observed_at = datetime(2026, 1, 1, 0, 0, 1, tzinfo=UTC)
response = { response = {
"streamTitle": "My Stream", "streamTitle": "My Stream",
"online": True, "online": True,
} }
state = StreamState.from_api_response(response, "example.com", observed_at) with time_machine.travel(observed_at, tick=False):
assert state.domain == "example.com" observation = StreamStateObservation.from_api_response(
assert state.title == "My Stream" response, "example.com"
assert state.online is True )
assert state.status_since == "2026-01-01T00:00:01+00:00" assert observation.domain == "example.com"
assert state.name is None assert observation.title == "My Stream"
assert state.failure_counter == 0 assert observation.online is True
assert observation.observed_at == "2026-01-01T00:00:01+00:00"
def test_missing_required_field_raises(self) -> None: def test_missing_required_field_raises(self) -> None:
"""Reject API responses without required stream state fields.""" """Reject API responses without required stream state fields."""
observed_at = datetime(2026, 1, 1, 0, 0, 1, tzinfo=UTC)
with pytest.raises(InvalidApiResponseError): with pytest.raises(InvalidApiResponseError):
StreamState.from_api_response({}, "bare.example.com", observed_at) StreamStateObservation.from_api_response({}, "bare.example.com")
def test_offline_response(self) -> None: def test_offline_response(self) -> None:
"""Populate offline state from a status response.""" """Populate offline state from a status response."""
@@ -171,9 +173,12 @@ class TestStreamStateFromApiResponse:
"streamTitle": "Offline Stream", "streamTitle": "Offline Stream",
"online": False, "online": False,
} }
state = StreamState.from_api_response(response, "example.com", observed_at) with time_machine.travel(observed_at, tick=False):
assert state.online is False observation = StreamStateObservation.from_api_response(
assert state.status_since == "2026-01-01T00:00:01+00:00" response, "example.com"
)
assert observation.online is False
assert observation.observed_at == "2026-01-01T00:00:01+00:00"
def test_title_truncation(self) -> None: def test_title_truncation(self) -> None:
"""Truncate the stream title to _MAX_STREAM_TITLE_LENGTH.""" """Truncate the stream title to _MAX_STREAM_TITLE_LENGTH."""
@@ -183,9 +188,12 @@ class TestStreamStateFromApiResponse:
"streamTitle": long_title, "streamTitle": long_title,
"online": True, "online": True,
} }
state = StreamState.from_api_response(response, "example.com", observed_at) with time_machine.travel(observed_at, tick=False):
assert len(state.title) == _MAX_STREAM_TITLE_LENGTH observation = StreamStateObservation.from_api_response(
assert state.title == "A" * _MAX_STREAM_TITLE_LENGTH response, "example.com"
)
assert len(observation.title) == _MAX_STREAM_TITLE_LENGTH
assert observation.title == "A" * _MAX_STREAM_TITLE_LENGTH
@pytest.mark.parametrize( @pytest.mark.parametrize(
("field", "value"), ("field", "value"),
@@ -196,7 +204,6 @@ class TestStreamStateFromApiResponse:
) )
def test_invalid_field_type_raises(self, field: str, value: object) -> None: def test_invalid_field_type_raises(self, field: str, value: object) -> None:
"""Reject stream state responses with malformed field types.""" """Reject stream state responses with malformed field types."""
observed_at = datetime(2026, 1, 1, 0, 0, 1, tzinfo=UTC)
response: dict[str, object] = { response: dict[str, object] = {
"streamTitle": "My Stream", "streamTitle": "My Stream",
"online": True, "online": True,
@@ -204,7 +211,7 @@ class TestStreamStateFromApiResponse:
response[field] = value response[field] = value
with pytest.raises(InvalidApiResponseError): with pytest.raises(InvalidApiResponseError):
StreamState.from_api_response(response, "example.com", observed_at) StreamStateObservation.from_api_response(response, "example.com")
class TestStreamStateFromDbRow: class TestStreamStateFromDbRow:
@@ -247,19 +254,22 @@ class TestStreamStateFromDbRow:
assert state.failure_counter == 0 assert state.failure_counter == 0
class TestStreamConfigFromApiResponse: class TestStreamConfigObservationFromApiResponse:
"""StreamConfig construction from an API response dictionary.""" """StreamConfigObservation construction from an API response dictionary."""
def test_typical_response(self) -> None: def test_typical_response(self) -> None:
"""Populate name and tags from a complete API response.""" """Populate name and tags from a complete API response."""
response = {"name": "My Instance", "tags": ["gaming", "music"]} response = {"name": "My Instance", "tags": ["gaming", "music"]}
config = StreamConfig.from_api_response(response) observed_at = datetime(2026, 1, 1, 0, 0, 1, tzinfo=UTC)
with time_machine.travel(observed_at, tick=False):
config = StreamConfigObservation.from_api_response(response)
assert config.name == "My Instance" assert config.name == "My Instance"
assert config.tags == ("gaming", "music") assert config.tags == ("gaming", "music")
assert config.observed_at == "2026-01-01T00:00:01+00:00"
def test_missing_keys_defaults(self) -> None: def test_missing_keys_defaults(self) -> None:
"""Use defaults when name and tags keys are missing.""" """Use defaults when name and tags keys are missing."""
config = StreamConfig.from_api_response({}) config = StreamConfigObservation.from_api_response({})
assert config.name == "" assert config.name == ""
assert config.tags == () assert config.tags == ()
@@ -267,7 +277,7 @@ class TestStreamConfigFromApiResponse:
"""Truncate the instance name to _MAX_INSTANCE_TITLE_LENGTH.""" """Truncate the instance name to _MAX_INSTANCE_TITLE_LENGTH."""
long_name = "B" * (_MAX_INSTANCE_TITLE_LENGTH + 50) long_name = "B" * (_MAX_INSTANCE_TITLE_LENGTH + 50)
response = {"name": long_name, "tags": []} response = {"name": long_name, "tags": []}
config = StreamConfig.from_api_response(response) config = StreamConfigObservation.from_api_response(response)
assert len(config.name) == _MAX_INSTANCE_TITLE_LENGTH assert len(config.name) == _MAX_INSTANCE_TITLE_LENGTH
assert config.name == "B" * _MAX_INSTANCE_TITLE_LENGTH assert config.name == "B" * _MAX_INSTANCE_TITLE_LENGTH
@@ -275,7 +285,7 @@ class TestStreamConfigFromApiResponse:
"""Truncate each tag to _MAX_TAG_LENGTH.""" """Truncate each tag to _MAX_TAG_LENGTH."""
long_tag = "C" * (_MAX_TAG_LENGTH + 10) long_tag = "C" * (_MAX_TAG_LENGTH + 10)
response = {"name": "", "tags": [long_tag, "short"]} response = {"name": "", "tags": [long_tag, "short"]}
config = StreamConfig.from_api_response(response) config = StreamConfigObservation.from_api_response(response)
assert len(config.tags[0]) == _MAX_TAG_LENGTH assert len(config.tags[0]) == _MAX_TAG_LENGTH
assert config.tags[0] == "C" * _MAX_TAG_LENGTH assert config.tags[0] == "C" * _MAX_TAG_LENGTH
assert config.tags[1] == "short" assert config.tags[1] == "short"
@@ -294,7 +304,7 @@ class TestStreamConfigFromApiResponse:
response[field] = value response[field] = value
with pytest.raises(InvalidApiResponseError): with pytest.raises(InvalidApiResponseError):
StreamConfig.from_api_response(response) StreamConfigObservation.from_api_response(response)
class TestValueTypeImmutability: class TestValueTypeImmutability:
@@ -306,15 +316,23 @@ class TestValueTypeImmutability:
with pytest.raises(FrozenInstanceError): with pytest.raises(FrozenInstanceError):
state.title = "Changed" # type: ignore[misc] state.title = "Changed" # type: ignore[misc]
def test_stream_observation_is_immutable(self) -> None:
"""StreamStateObservation cannot be mutated in place."""
observation = StreamStateObservation(
domain="stream.example",
)
with pytest.raises(FrozenInstanceError):
observation.title = "Changed" # type: ignore[misc]
def test_stream_config_is_immutable(self) -> None: def test_stream_config_is_immutable(self) -> None:
"""StreamConfig cannot be mutated in place.""" """StreamConfigObservation cannot be mutated in place."""
config = StreamConfig(name="Stream") config = StreamConfigObservation(name="Stream")
with pytest.raises(FrozenInstanceError): with pytest.raises(FrozenInstanceError):
config.name = "Changed" # type: ignore[misc] config.name = "Changed" # type: ignore[misc]
def test_stream_config_tags_are_immutable(self) -> None: def test_stream_config_tags_are_immutable(self) -> None:
"""StreamConfig tags are stored in an immutable tuple.""" """StreamConfigObservation tags are stored in an immutable tuple."""
config = StreamConfig(name="Stream", tags=("gaming",)) config = StreamConfigObservation(name="Stream", tags=("gaming",))
assert config.tags == ("gaming",) assert config.tags == ("gaming",)
def test_update_result_is_immutable(self) -> None: def test_update_result_is_immutable(self) -> None: