Compare commits
2
Commits
38e91add32
...
d02202295c
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
d02202295c
|
||
|
|
58ad476a94
|
@@ -37,7 +37,7 @@ jobs:
|
|||||||
grep "^version:" maubot.yaml
|
grep "^version:" maubot.yaml
|
||||||
|
|
||||||
- name: Build plugin
|
- name: Build plugin
|
||||||
run: 7z a -tzip -mx=9 -xr!__pycache__ owncastsentry-v${{ steps.metadata.outputs.version }}.mbp owncastsentry/ maubot.yaml base-config.yaml LICENSE.txt
|
run: 7z a -tzip -mx=9 -xr!__pycache__ -xr!README.md owncastsentry-v${{ steps.metadata.outputs.version }}.mbp owncastsentry/ maubot.yaml base-config.yaml LICENSE.txt
|
||||||
|
|
||||||
- name: Print artifact checksums
|
- name: Print artifact checksums
|
||||||
run: |
|
run: |
|
||||||
|
|||||||
+155
-228
@@ -13,64 +13,46 @@ or tests.
|
|||||||
|
|
||||||
- `__init__.py`: Maubot plugin entry point, service wiring, command
|
- `__init__.py`: Maubot plugin entry point, service wiring, command
|
||||||
decorators, scheduler loop, shutdown, and metrics endpoint.
|
decorators, scheduler loop, shutdown, and metrics endpoint.
|
||||||
|
- `commands.py`: Maubot command handlers for subscribe, unsubscribe,
|
||||||
|
subscriptions, and live listings.
|
||||||
- `config.py`: Maubot configuration proxy for the Prometheus metrics endpoint
|
- `config.py`: Maubot configuration proxy for the Prometheus metrics endpoint
|
||||||
setting.
|
setting.
|
||||||
- `types.py`: immutable value objects, Owncast API response validation, derived
|
- `metrics.py`: Prometheus counters, gauges, registry setup, and response timing
|
||||||
stream status, update results, and subscription domain errors.
|
helpers.
|
||||||
- `repository.py`: database schema migrations plus stream and subscription
|
- `notification_service.py`: Matrix notification formatting, broadcast,
|
||||||
repositories.
|
per-domain notification cooldowns, cleanup notices, and delivery metrics.
|
||||||
- `owncast_client.py`: bounded aiohttp client for Owncast `/api/status` and
|
- `owncast_client.py`: bounded aiohttp client for Owncast `/api/status` and
|
||||||
`/api/config` requests.
|
`/api/config` requests.
|
||||||
- `subscription_manager.py`: subscription use cases, domain normalization, and
|
- `repository.py`: database schema migrations plus stream and subscription
|
||||||
first-subscription Owncast validation.
|
repositories.
|
||||||
- `stream_monitor.py`: stream update orchestration, state transition detection,
|
- `stream_monitor.py`: stream update orchestration, state transition detection,
|
||||||
progressive failure backoff, cleanup thresholds, notification decisions, and
|
progressive failure backoff, cleanup thresholds, notification decisions, and
|
||||||
stream metrics.
|
stream metrics.
|
||||||
- `notification_service.py`: Matrix notification formatting, broadcast,
|
- `subscription_manager.py`: subscription use cases, domain normalization, and
|
||||||
per-domain notification cooldowns, cleanup notices, and delivery metrics.
|
first-subscription Owncast validation.
|
||||||
- `commands.py`: Maubot command handlers for subscribe, unsubscribe,
|
- `types.py`: immutable value objects, Owncast API response validation, derived
|
||||||
subscriptions, and live listings.
|
stream status, update results, and subscription domain errors.
|
||||||
- `metrics.py`: isolated Prometheus registry, counters, gauges, and response
|
|
||||||
timing helpers.
|
|
||||||
|
|
||||||
The package has a service-layer shape: `OwncastSentry.start()` builds the
|
The package has a service-layer shape: `OwncastSentry.start()` builds the
|
||||||
services, commands call `CommandHandler`, command handlers call
|
services, commands call `CommandHandler`, command handlers call
|
||||||
`SubscriptionManager`, scheduled updates call `StreamMonitor`, and
|
`SubscriptionManager`, and scheduled updates call `StreamMonitor`.
|
||||||
`StreamMonitor` coordinates repositories, the Owncast client, notifications,
|
`StreamMonitor` runs the main polling workflow: it checks Owncast streams,
|
||||||
and metrics.
|
applies backoff and cleanup rules, classifies state transitions, triggers any
|
||||||
|
live/title notifications, persists changed stream state, and records stream,
|
||||||
```mermaid
|
failure, and subscription metrics.
|
||||||
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
|
## Core Concepts
|
||||||
|
|
||||||
A few terms are used throughout the package:
|
A few terms are used throughout the package:
|
||||||
|
|
||||||
- A stream is a normalized Owncast domain tracked by the plugin. The database
|
- A stream is a tracked [Owncast](https://owncast.online/) instance, identified
|
||||||
stores one stream row per domain, even if multiple Matrix rooms subscribe to
|
by its normalized bare domain. The stored entry for that domain is the stream
|
||||||
it.
|
record. There is one stream record per domain, even if multiple Matrix rooms
|
||||||
|
subscribe to it.
|
||||||
- A subscription is a Matrix room's request to receive notifications for one
|
- A subscription is a Matrix room's request to receive notifications for one
|
||||||
stream domain.
|
stream domain.
|
||||||
- `StreamStateObservation` is one successful Owncast status API sample. The
|
- `StreamStateObservation` is one successful Owncast status API sample used by
|
||||||
monitor decides whether the observation starts a new persisted status period.
|
the monitor when updating stored stream state.
|
||||||
- `StreamConfigObservation` is metadata fetched from Owncast's config API,
|
- `StreamConfigObservation` is metadata fetched from Owncast's config API,
|
||||||
currently the instance name and tags used in notifications and periodic
|
currently the instance name and tags used in notifications and periodic
|
||||||
display-name refreshes.
|
display-name refreshes.
|
||||||
@@ -83,9 +65,11 @@ Durable state lives in the database. The `streams` table stores the latest
|
|||||||
known display and state fields for each Owncast domain, and the `subscriptions`
|
known display and state fields for each Owncast domain, and the `subscriptions`
|
||||||
table stores which Matrix rooms follow each domain.
|
table stores which Matrix rooms follow each domain.
|
||||||
|
|
||||||
Subscriptions are room-scoped, while stream state is shared per domain.
|
Subscriptions belong to individual Matrix rooms, while stream state is shared
|
||||||
Notification delivery fans out from one stored stream state to the rooms
|
per domain. Notification delivery fans out by domain to the rooms currently
|
||||||
currently subscribed to that domain.
|
subscribed to that domain. Live/title notification text is built from the fresh
|
||||||
|
status observation and, when available, fresh config metadata for the update
|
||||||
|
being processed.
|
||||||
|
|
||||||
Stream identity is the normalized bare domain. Repositories, metrics,
|
Stream identity is the normalized bare domain. Repositories, metrics,
|
||||||
notification state, and Owncast API calls all use that domain rather than the
|
notification state, and Owncast API calls all use that domain rather than the
|
||||||
@@ -100,41 +84,43 @@ commands, stream updates, or notifications.
|
|||||||
are represented by new values that are persisted through repositories instead of
|
are represented by new values that are persisted through repositories instead of
|
||||||
mutating existing objects.
|
mutating existing objects.
|
||||||
|
|
||||||
## Service Wiring
|
## Plugin Boundary
|
||||||
|
|
||||||
|
Maubot interacts with the package through `OwncastSentry`. That plugin class is
|
||||||
|
the boundary between Maubot's runtime and the internal services that implement
|
||||||
|
OwncastSentry behavior.
|
||||||
|
|
||||||
`OwncastSentry.start()` is the composition root for one plugin instance. It
|
`OwncastSentry.start()` is the composition root for one plugin instance. It
|
||||||
loads config, creates the metrics service, Owncast API client, repositories,
|
loads config, creates the Owncast client, repositories, metrics, notification
|
||||||
notification service, stream monitor, subscription manager, and command handler,
|
service, stream monitor, subscription manager, and command handler, then passes
|
||||||
then passes dependencies into each service explicitly.
|
dependencies into each service explicitly.
|
||||||
|
|
||||||
This keeps lower-level modules independent of Maubot globals. Commands do not
|
The plugin boundary owns Maubot integration: command registration, scheduled
|
||||||
construct repositories, `StreamMonitor` does not know about the Maubot
|
polling, and lifecycle cleanup. Scheduled ticks load subscribed domains and
|
||||||
scheduler, and notification delivery does not fetch its own subscription state
|
call `StreamMonitor.update_all_streams()`. `StreamMonitor` owns the per-stream
|
||||||
from the plugin instance.
|
polling workflow, while notification delivery and delivery cooldowns stay in
|
||||||
|
`NotificationService`.
|
||||||
|
|
||||||
The plugin boundary owns scheduling. The minute loop loads subscribed domains
|
Internal services do not depend on Maubot globals, which keeps command handling,
|
||||||
from `SubscriptionRepository` and calls `StreamMonitor.update_all_streams()`.
|
scheduled polling, and tests wired through explicit dependencies.
|
||||||
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
|
## Subscription Management
|
||||||
|
|
||||||
`SubscriptionManager` coordinates the room-to-domain relationship. Command
|
`SubscriptionManager` coordinates the room-to-domain relationship. Command
|
||||||
handlers pass it user-supplied stream targets, and it turns those targets into
|
handlers pass it user-supplied stream targets, and it turns those targets into
|
||||||
normalized domains before touching repository state.
|
normalized domains before creating or removing stored data.
|
||||||
|
|
||||||
Subscribing a room creates a subscription for the normalized domain and creates
|
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
|
the shared stream record when needed. Domains with no current subscribers are
|
||||||
validated through `OwncastClient`; domains that already have subscribers reuse
|
validated through `OwncastClient`; domains that already have subscribers reuse
|
||||||
the existing stream record instead of revalidating.
|
the existing stream record instead of revalidating.
|
||||||
|
|
||||||
Unsubscribing removes one room's subscription to a domain. It does not delete
|
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
|
the shared stream record or make remote Owncast requests.
|
||||||
cleanup is owned by `StreamMonitor`.
|
|
||||||
|
|
||||||
Listing methods stay room-scoped. They return the subscriptions relevant to one
|
Listing methods return the subscriptions for one Matrix room with shared stream
|
||||||
Matrix room, with the shared stream state attached for display.
|
state attached, so commands can display each instance's name, title, link,
|
||||||
|
status, and how long it has been online or offline.
|
||||||
|
|
||||||
## Owncast API Client
|
## Owncast API Client
|
||||||
|
|
||||||
@@ -144,13 +130,12 @@ 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,
|
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
|
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
|
socket timeouts. aiohttp's own DNS cache is disabled; resolver selection is left
|
||||||
stay with the resolver configured for the deployment.
|
to aiohttp and the runtime environment.
|
||||||
|
|
||||||
The per-host connection limit and keepalive settings are intentional. Each
|
The per-host connection limit and keepalive settings are intentional. Each
|
||||||
Owncast instance gets at most one reusable connection, so minute-by-minute
|
Owncast instance is intended to keep at most one long-lived connection warm, so
|
||||||
status checks can reuse the existing TLS session instead of opening a fresh
|
minute-by-minute status checks avoid unnecessary TLS renegotiation.
|
||||||
connection for every poll.
|
|
||||||
|
|
||||||
Responses are accepted only when they are HTTP 200 JSON objects under the
|
Responses are accepted only when they are HTTP 200 JSON objects under the
|
||||||
configured size limit. Invalid JSON, non-object JSON, oversized bodies,
|
configured size limit. Invalid JSON, non-object JSON, oversized bodies,
|
||||||
@@ -163,147 +148,118 @@ for that domain.
|
|||||||
|
|
||||||
## Stream Monitoring
|
## Stream Monitoring
|
||||||
|
|
||||||
`StreamMonitor` owns the core per-domain state machine. `update_stream()`
|
`StreamMonitor` runs the package's main polling workflow for tracked Owncast
|
||||||
compares the stored stream state with the latest Owncast status response,
|
instances. `update_stream()` owns the per-domain state machine: it compares the
|
||||||
decides whether to send a notification, updates persistent state when needed,
|
stored stream state with the latest Owncast status response, decides whether to
|
||||||
and records stream metrics. Each poll first handles missing rows, polling
|
send a notification, updates persistent state when needed, and records stream
|
||||||
backoff, and failed Owncast requests. Only a valid status response enters the
|
metrics. Each poll first handles stored stream lookup, polling backoff, and
|
||||||
stream transition policy. `update_all_streams()` wraps this flow for many
|
failed Owncast requests. Only a valid status response enters the stream transition
|
||||||
domains, isolates per-domain exceptions, and records subscription-count metrics.
|
policy. `update_all_streams()` wraps this flow for many domains, isolates
|
||||||
|
per-domain exceptions, and records subscription-count metrics.
|
||||||
|
|
||||||
```mermaid
|
```mermaid
|
||||||
flowchart TD
|
flowchart TD
|
||||||
Start["update_stream(domain)"] --> Old["load stored StreamState"]
|
Start["update_stream(domain)"] --> Stored["load stored stream state"]
|
||||||
Old --> Missing{"row exists?"}
|
Stored --> Query{"backoff allows query?"}
|
||||||
Missing -->|no| Noop["return success"]
|
Query -->|no| Defer["skip query for backoff"]
|
||||||
Missing -->|yes| Backoff{"backoff allows poll?"}
|
Defer --> Cleanup["apply cleanup policy"]
|
||||||
Backoff -->|no| Skip["skip HTTP call<br/>increment failure counter"]
|
Cleanup --> Done["done"]
|
||||||
Skip --> CleanupSkip["run warning/deletion cleanup checks"]
|
Query -->|yes| Status["fetch stream status from Owncast"]
|
||||||
CleanupSkip --> MetricsSkip["if row remains, update<br/>status and failure metrics"]
|
Status --> Observation{"valid observation?"}
|
||||||
MetricsSkip --> SuccessSkip["return success"]
|
Observation -->|no| Failure["record failed check"]
|
||||||
Backoff -->|yes| Fetch["fetch and validate /api/status"]
|
Failure --> Cleanup
|
||||||
Fetch --> Fetched{"got valid state?"}
|
Observation -->|yes| Compare["compare observation with stored state"]
|
||||||
Fetched -->|no| Failure["increment failure counter"]
|
Compare --> First{"first observation?"}
|
||||||
Failure --> CleanupFail["run warning/deletion cleanup checks"]
|
First -->|yes| FirstNoNotify["no notification"]
|
||||||
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?"}
|
First -->|no| LiveCheck{"went live?"}
|
||||||
FirstLog --> ConfigDecision
|
LiveCheck -->|yes| Brief{"offline less than 7 minutes?"}
|
||||||
LiveCheck -->|yes| Brief{"observed offline less than<br/>temporary cooldown?"}
|
Brief -->|no| LiveNotify["select go-live notification"]
|
||||||
LiveCheck -->|no| TitleChanged{"title changed<br/>while online?"}
|
|
||||||
Brief -->|no| LiveNotify["go-live notification needed"]
|
|
||||||
Brief -->|yes| BriefTitle{"title changed?"}
|
Brief -->|yes| BriefTitle{"title changed?"}
|
||||||
BriefTitle -->|yes| BriefNotify["title-change notification needed"]
|
BriefTitle -->|no| BriefNoNotify["no notification"]
|
||||||
BriefTitle -->|no| NoNotify["no notification"]
|
BriefTitle -->|yes| TitleNotify["select title-change notification"]
|
||||||
TitleChanged -->|yes| OfflineMarker{"offline marker newer than<br/>last notification?"}
|
LiveCheck -->|no| TitleChanged{"title changed while online?"}
|
||||||
|
TitleChanged -->|yes| TitleNotify
|
||||||
TitleChanged -->|no| WentOffline{"went offline?"}
|
TitleChanged -->|no| WentOffline{"went offline?"}
|
||||||
WentOffline -->|yes| Offline["record offline timer"]
|
WentOffline -->|yes| Offline["record when stream went offline"]
|
||||||
WentOffline -->|no| NoNotify
|
WentOffline -->|no| NoNotify["no notification"]
|
||||||
OfflineMarker -->|yes| ResumeNotify["go-live notification needed"]
|
FirstNoNotify --> Metadata
|
||||||
OfflineMarker -->|no| TitleNotify["title-change notification needed"]
|
BriefNoNotify --> Metadata
|
||||||
NoNotify --> ConfigDecision
|
NoNotify --> Metadata
|
||||||
Offline --> ConfigDecision
|
Offline --> Metadata
|
||||||
LiveNotify --> ConfigDecision
|
LiveNotify --> Metadata
|
||||||
BriefNotify --> ConfigDecision
|
TitleNotify --> Metadata
|
||||||
ResumeNotify --> ConfigDecision
|
Metadata["fetch instance config from Owncast if display metadata is needed"]
|
||||||
TitleNotify --> ConfigDecision
|
Metadata --> Notify["attempt sending notification if selected"]
|
||||||
ConfigDecision{"notification, first observation,<br/>or hourly refresh?"}
|
Notify --> Save["save stream update"]
|
||||||
ConfigDecision -->|yes| Config["fetch /api/config"]
|
Save --> Done
|
||||||
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
|
Each stream has a failure counter that tracks consecutive failed checks and
|
||||||
as unknown when its failure counter is above `UNKNOWN_STATUS_THRESHOLD`.
|
queries skipped for backoff. The counter drives both status and polling behavior:
|
||||||
Counters 0 through 4 poll every minute; later tiers poll only on selected
|
a stream is displayed as unknown after more than 15 minutes of failed or
|
||||||
counter values and eventually only every fifteenth minute. Skipped cycles still
|
skipped checks, and later failures progressively reduce query frequency until
|
||||||
increment the counter so cleanup can progress.
|
the monitor checks only every 15 minutes. Skipped queries still increment
|
||||||
|
the counter so cleanup can progress.
|
||||||
|
|
||||||
Transition decisions use Owncast's `online` field. The stream status timestamp is
|
After each successful fetch from Owncast's
|
||||||
the bot's local UTC observation time for the first successful poll, a missing
|
[`/api/status`](https://owncast.online/api/latest/#tag/Internal/operation/GetStatus)
|
||||||
timestamp repair, or an online/offline transition. Owncast's remote connect,
|
endpoint, the monitor compares the stored stream state with the latest
|
||||||
disconnect, and server timestamps are not required for status parsing.
|
observation to identify transitions, such as going online, going offline, or
|
||||||
|
changing title. Online/offline transitions are based on Owncast's `online` field.
|
||||||
|
The stream status timestamp is the bot's UTC observation time for the first
|
||||||
|
successful poll or an online/offline transition.
|
||||||
|
|
||||||
First observation is special. If a new stream row has no status timestamp, the
|
The first successful observation stores state and fetches config but suppresses
|
||||||
first successful update fetches config and stores state but suppresses
|
notifications, even if the stream is already live. After that, only
|
||||||
notification attempts, even if the stream is already live.
|
offline-to-online transitions can select go-live notifications; mid-session
|
||||||
|
title changes select title-change notifications.
|
||||||
|
|
||||||
Go-live notification attempts happen when a stream moves from offline to online
|
Brief outages are treated as transient to avoid spamming Matrix rooms when an
|
||||||
and the monitor's in-memory offline timer shows it was not merely a brief
|
instance is spotty. They do not select go-live notifications. If a stream returns
|
||||||
outage. A stream that returns within the temporary offline cooldown only attempts
|
in under seven minutes with the same title, the monitor stores the new online
|
||||||
a notification when the title changed, and then it uses a title change
|
state without notifying rooms. If the title changed during that brief outage, the
|
||||||
notification.
|
monitor selects a title-change notification instead.
|
||||||
|
|
||||||
Mid-session title changes attempt title change notifications. If the monitor has
|
Instance config is fetched from Owncast's
|
||||||
an offline marker that is newer than the last notification, it attempts a normal
|
[`/api/config`](https://owncast.online/api/latest/#tag/Internal/operation/GetWebConfig)
|
||||||
live notification instead so rooms do not see a title change for what is
|
endpoint when a live/title notification is selected, on a stream's first
|
||||||
effectively a new session.
|
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.
|
||||||
|
|
||||||
Offline transitions do not notify rooms. They update the database and record a
|
Selecting a notification does not guarantee a Matrix message. The notification
|
||||||
monotonic offline timestamp used by later live/title notification policy.
|
service can still suppress delivery when its per-domain live/title cooldown is
|
||||||
|
active, which prevents repeated messages from noisy stream state changes. That
|
||||||
The live/title notification send step is an attempt, not a guaranteed Matrix message.
|
delivery behavior is covered below.
|
||||||
`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
|
## Cleanup Policy
|
||||||
|
|
||||||
A domain that remains unreachable is eventually cleaned up. The failure counter
|
A domain that remains unreachable is eventually cleaned up. The failure counter
|
||||||
represents minute ticks, including skipped backoff ticks.
|
represents minute ticks, including skipped backoff ticks.
|
||||||
|
|
||||||
At the 83-day threshold, `StreamMonitor` sends a cleanup warning to all rooms
|
After 83 days of continued failed or skipped checks, `StreamMonitor` sends a
|
||||||
subscribed to the domain. At the 90-day threshold or beyond, it sends a deletion
|
cleanup warning to all rooms subscribed to the domain. After 90 days or more of
|
||||||
notice, deletes all subscriptions for the domain, deletes the stream row, clears
|
continued failed or skipped checks, it sends a deletion notice, deletes all
|
||||||
the monitor and notification caches for that domain, and removes per-domain
|
subscriptions for the domain, deletes the stream record, clears the monitor and
|
||||||
metrics.
|
notification caches for that domain, and asks `MetricsService` to remove
|
||||||
|
existing per-domain metric labels. The enclosing update cycle may later recreate
|
||||||
```mermaid
|
the subscription-count gauge for that domain at zero.
|
||||||
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
|
## Notification Delivery
|
||||||
|
|
||||||
`NotificationService` sends plain text Matrix messages to every room subscribed
|
`NotificationService` sends notifications to every room subscribed to a domain.
|
||||||
to a domain. Broadcasts run concurrently and use `asyncio.gather(...,
|
Delivery failures are isolated per room, so one failed Matrix send does not
|
||||||
return_exceptions=True)` so one failed room does not block delivery to other
|
block notifications to other subscribed rooms.
|
||||||
rooms.
|
|
||||||
|
|
||||||
Live and title-change notifications share a per-domain cooldown. The cooldown
|
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
|
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.
|
subscribed rooms, or every delivery fails, no cooldown is recorded.
|
||||||
|
|
||||||
Notification text is whitespace-normalized before sending. Stream names fall
|
Live/title notification text normalizes remote stream metadata before composing
|
||||||
back to the domain when config lookup fails or the instance name is empty. Tags
|
the message, so instance-provided names, titles, and tags cannot add unexpected
|
||||||
are appended as hash tags after filtering out empty tags and tags that start
|
line breaks or otherwise distort the notification. Messages fall back to the
|
||||||
with a dot.
|
domain when a display name is unavailable.
|
||||||
|
|
||||||
Cleanup warning and deletion notices bypass the live/title cooldown, but still
|
Cleanup warning and deletion notices bypass the live/title cooldown, but still
|
||||||
record delivery metrics.
|
record delivery metrics.
|
||||||
@@ -316,33 +272,30 @@ that catch unexpected exceptions and record command error metrics.
|
|||||||
|
|
||||||
`CommandHandler.subscribe()` and `CommandHandler.unsubscribe()` translate
|
`CommandHandler.subscribe()` and `CommandHandler.unsubscribe()` translate
|
||||||
domain errors into user-facing Matrix replies. `subscriptions()` and `live()`
|
domain errors into user-facing Matrix replies. `subscriptions()` and `live()`
|
||||||
read resolved room subscriptions and render Markdown output with escaped
|
read room subscription listings with stream state attached and render Markdown
|
||||||
untrusted text.
|
output with escaped untrusted text.
|
||||||
|
|
||||||
Command listings are room-scoped and ordered by domain through repository
|
Command listings are limited to the current room and ordered by domain through
|
||||||
queries. `subscriptions()` includes online, offline, and unknown streams.
|
repository queries. `subscriptions()` includes online, offline, and unknown
|
||||||
`live()` includes only rows marked online with a failure count at or below the
|
streams. `live()` includes only streams marked online with a failure count at or
|
||||||
unknown threshold.
|
below the unknown threshold.
|
||||||
|
|
||||||
## Persistence
|
## Persistence
|
||||||
|
|
||||||
The database schema is managed by `repository.py` through Maubot's upgrade
|
Maubot owns the database connection and runs the schema upgrades registered by
|
||||||
table.
|
`repository.py`. The repository classes wrap that database handle and expose the
|
||||||
|
package's stream and subscription persistence operations.
|
||||||
|
|
||||||
- Revision 1 creates `streams` and `subscriptions`.
|
The persistence model stores one stream record per normalized domain and one
|
||||||
- Revision 2 fixes `subscriptions.stream_domain` from `INTEGER` to `TEXT`.
|
subscription row per room/domain pair. `SubscriptionManager` normalizes user
|
||||||
- Revision 3 adds `streams.failure_counter`.
|
input before repository calls. `StreamRepository` writes display and state
|
||||||
- Revision 4 replaces separate connect/disconnect timestamps with
|
fields, while failure counters use dedicated methods.
|
||||||
`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
|
`SubscriptionRepository` raises domain-specific errors for duplicate adds and
|
||||||
missing removes.
|
missing removes.
|
||||||
|
|
||||||
Resolved subscription listings join `subscriptions` to `streams`, which means
|
Room subscription listings join `subscriptions` to `streams`, which means
|
||||||
orphaned subscription rows without a stream row are skipped in room display
|
orphaned subscription entries without a matching stream record are skipped in
|
||||||
queries.
|
room display queries.
|
||||||
|
|
||||||
## Metrics
|
## Metrics
|
||||||
|
|
||||||
@@ -355,36 +308,10 @@ The registered metrics cover:
|
|||||||
- current stream status per domain (`1` online, `0` offline, `-1` unknown);
|
- current stream status per domain (`1` online, `0` offline, `-1` unknown);
|
||||||
- subscription count per domain;
|
- subscription count per domain;
|
||||||
- consecutive check failures per domain;
|
- consecutive check failures per domain;
|
||||||
- last successful Owncast API response duration per domain;
|
- successful Owncast API response duration per domain;
|
||||||
- build/version information;
|
- build/version information;
|
||||||
- open aiohttp connection count;
|
- open aiohttp connection count;
|
||||||
- internal scheduler and command errors.
|
- internal scheduler and command errors.
|
||||||
|
|
||||||
Known counter label combinations are initialized to zero. This keeps the
|
Known counter label combinations are initialized to zero. This keeps the
|
||||||
Prometheus output stable before any events have occurred.
|
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.
|
|
||||||
|
|||||||
@@ -134,12 +134,11 @@ def _classify_notification(
|
|||||||
observation: StreamStateObservation,
|
observation: StreamStateObservation,
|
||||||
*,
|
*,
|
||||||
offline_duration_seconds: int | None = None,
|
offline_duration_seconds: int | None = None,
|
||||||
went_offline_after_last_notification: bool = False,
|
|
||||||
) -> _NotificationKind:
|
) -> _NotificationKind:
|
||||||
"""Classify the live/title notification policy for a transition.
|
"""Classify the live/title notification policy for a transition.
|
||||||
|
|
||||||
The caller supplies timing context because it comes from the monitor's
|
The caller supplies offline duration because it comes from the monitor's
|
||||||
monotonic offline cache and notification-service cache, not the transition.
|
monotonic offline cache, not the transition.
|
||||||
"""
|
"""
|
||||||
match transition.kind:
|
match transition.kind:
|
||||||
case _StreamTransitionKind.WENT_LIVE:
|
case _StreamTransitionKind.WENT_LIVE:
|
||||||
@@ -153,8 +152,6 @@ def _classify_notification(
|
|||||||
return _NotificationKind.LIVE
|
return _NotificationKind.LIVE
|
||||||
|
|
||||||
case _StreamTransitionKind.TITLE_CHANGED:
|
case _StreamTransitionKind.TITLE_CHANGED:
|
||||||
if went_offline_after_last_notification:
|
|
||||||
return _NotificationKind.LIVE
|
|
||||||
return _NotificationKind.TITLE_CHANGE
|
return _NotificationKind.TITLE_CHANGE
|
||||||
|
|
||||||
case _:
|
case _:
|
||||||
@@ -314,7 +311,6 @@ class StreamMonitor:
|
|||||||
# Notification policy may need runtime timing context in addition to
|
# Notification policy may need runtime timing context in addition to
|
||||||
# the stored state and latest observation.
|
# the stored state and latest observation.
|
||||||
offline_duration_seconds: int | None = None
|
offline_duration_seconds: int | None = None
|
||||||
went_offline_after_last_notification = False
|
|
||||||
|
|
||||||
match transition.kind:
|
match transition.kind:
|
||||||
case _StreamTransitionKind.FIRST_OBSERVATION:
|
case _StreamTransitionKind.FIRST_OBSERVATION:
|
||||||
@@ -343,11 +339,6 @@ class StreamMonitor:
|
|||||||
)
|
)
|
||||||
case _StreamTransitionKind.TITLE_CHANGED:
|
case _StreamTransitionKind.TITLE_CHANGED:
|
||||||
self.log.info("[%s] Stream title was changed!", domain)
|
self.log.info("[%s] Stream title was changed!", domain)
|
||||||
# If we saw an offline event after the last notification, send
|
|
||||||
# a go-live notice instead of only a title-change notice.
|
|
||||||
went_offline_after_last_notification = self.offline_timer_cache[
|
|
||||||
domain
|
|
||||||
] > self.notification_service.get_last_notification_time(domain)
|
|
||||||
case _StreamTransitionKind.WENT_OFFLINE:
|
case _StreamTransitionKind.WENT_OFFLINE:
|
||||||
self.offline_timer_cache[domain] = time.monotonic()
|
self.offline_timer_cache[domain] = time.monotonic()
|
||||||
self.log.info("[%s] Stream is now offline.", domain)
|
self.log.info("[%s] Stream is now offline.", domain)
|
||||||
@@ -359,7 +350,6 @@ class StreamMonitor:
|
|||||||
old_state,
|
old_state,
|
||||||
observation,
|
observation,
|
||||||
offline_duration_seconds=offline_duration_seconds,
|
offline_duration_seconds=offline_duration_seconds,
|
||||||
went_offline_after_last_notification=went_offline_after_last_notification,
|
|
||||||
)
|
)
|
||||||
|
|
||||||
stream_config: StreamConfigObservation | None = None
|
stream_config: StreamConfigObservation | None = None
|
||||||
|
|||||||
@@ -457,32 +457,6 @@ class TestClassifyNotification:
|
|||||||
|
|
||||||
assert notification_kind is _NotificationKind.TITLE_CHANGE
|
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."""
|
||||||
|
|
||||||
@@ -959,12 +933,12 @@ class TestUpdateStreamTitleChange:
|
|||||||
)
|
)
|
||||||
assert owncast.config_call_count == 1
|
assert owncast.config_call_count == 1
|
||||||
|
|
||||||
async def test_sends_go_live_when_last_notification_before_offline(
|
async def test_sends_title_change_after_offline_gap(
|
||||||
self,
|
self,
|
||||||
stream_repo: StreamRepository,
|
stream_repo: StreamRepository,
|
||||||
subscription_repo: SubscriptionRepository,
|
subscription_repo: SubscriptionRepository,
|
||||||
) -> None:
|
) -> None:
|
||||||
"""Send go-live instead of title-change after an offline gap."""
|
"""Send title-change when an already-online stream changes title."""
|
||||||
owncast = _StubOwncastClient(
|
owncast = _StubOwncastClient(
|
||||||
stream_state=_stream_observation(
|
stream_state=_stream_observation(
|
||||||
domain="title.com",
|
domain="title.com",
|
||||||
@@ -991,8 +965,8 @@ class TestUpdateStreamTitleChange:
|
|||||||
status_since="2026-01-01T12:00:00+00:00",
|
status_since="2026-01-01T12:00:00+00:00",
|
||||||
)
|
)
|
||||||
|
|
||||||
# Offline timer is MORE recent than last notification,
|
# Offline timer is more recent than the last notification, but title
|
||||||
# and both are old enough to pass rate limiting
|
# changes remain title changes once the stream is already online.
|
||||||
now = time.monotonic()
|
now = time.monotonic()
|
||||||
monitor.offline_timer_cache["title.com"] = (
|
monitor.offline_timer_cache["title.com"] = (
|
||||||
now - _SECONDS_BETWEEN_NOTIFICATIONS - 100
|
now - _SECONDS_BETWEEN_NOTIFICATIONS - 100
|
||||||
@@ -1005,7 +979,7 @@ class TestUpdateStreamTitleChange:
|
|||||||
assert result is True
|
assert result is True
|
||||||
assert len(client.sent_messages) == 1
|
assert len(client.sent_messages) == 1
|
||||||
assert client.sent_messages[0].content.body == (
|
assert client.sent_messages[0].content.body == (
|
||||||
"🎥 Title Stream is now live!\n"
|
"📝 Title Stream has changed its stream title!\n"
|
||||||
"Stream Title: Updated Title\n"
|
"Stream Title: Updated Title\n"
|
||||||
"\n"
|
"\n"
|
||||||
"To tune in, visit: https://title.com/"
|
"To tune in, visit: https://title.com/"
|
||||||
|
|||||||
Reference in New Issue
Block a user