Clarified comments and docstrings.
This commit is contained in:
@@ -51,10 +51,10 @@ class OwncastSentry(Plugin):
|
||||
return Config
|
||||
|
||||
async def start(self) -> None:
|
||||
"""Initialize all services and register recurring tasks.
|
||||
"""Initialize services and schedule stream state updates.
|
||||
|
||||
Registers a recurring task every minute to update the state of
|
||||
all subscribed streams.
|
||||
Starts a scheduler coroutine that updates subscribed streams at the top
|
||||
of each wall-clock minute.
|
||||
"""
|
||||
# Load configuration
|
||||
config: Config = self.config # type: ignore[assignment]
|
||||
@@ -170,7 +170,7 @@ class OwncastSentry(Plugin):
|
||||
await self._update_all_stream_states()
|
||||
|
||||
async def _update_all_stream_states(self) -> None:
|
||||
"""Update all stream states."""
|
||||
"""Update stream states for domains with active subscriptions."""
|
||||
try:
|
||||
# Get list of all stream domains with active subscriptions
|
||||
subscribed_domains = (
|
||||
|
||||
@@ -102,7 +102,7 @@ class CommandHandler:
|
||||
"""Subscribe a room to a stream's notifications.
|
||||
|
||||
:param evt: MessageEvent of the message calling the command.
|
||||
:param url: User supplied URL to a stream to subscribe to.
|
||||
:param url: User-supplied Owncast URL, domain, or Fediverse-style address.
|
||||
"""
|
||||
try:
|
||||
stream_domain = await self.subscription_manager.subscribe(evt.room_id, url)
|
||||
@@ -129,7 +129,7 @@ class CommandHandler:
|
||||
"""Unsubscribe a room from a stream's notifications.
|
||||
|
||||
:param evt: MessageEvent of the message calling the command.
|
||||
:param url: User supplied URL to a stream to unsubscribe from.
|
||||
:param url: User-supplied Owncast URL, domain, or Fediverse-style address.
|
||||
"""
|
||||
try:
|
||||
stream_domain = await self.subscription_manager.unsubscribe(
|
||||
@@ -249,7 +249,7 @@ class CommandHandler:
|
||||
# Start building this stream's entry with stream name as main bullet
|
||||
parts.append(f"- **{safe_stream_name}** \n")
|
||||
|
||||
# Add title (should be present for live streams)
|
||||
# Add title when available
|
||||
if stream_state.title:
|
||||
safe_title = _sanitize_for_markdown(stream_state.title)
|
||||
parts.append(f" - Title: {safe_title} \n")
|
||||
|
||||
@@ -110,7 +110,7 @@ class MetricsService:
|
||||
)
|
||||
self.api_response_seconds = Gauge(
|
||||
"owncastsentry_api_response_seconds",
|
||||
"Last successful HTTP response time in seconds per stream domain",
|
||||
"Last successful Owncast API query duration in seconds per stream domain",
|
||||
["domain"],
|
||||
registry=self.registry,
|
||||
)
|
||||
@@ -121,7 +121,7 @@ class MetricsService:
|
||||
)
|
||||
self.open_connections = Gauge(
|
||||
"owncastsentry_http_connections_open",
|
||||
"Total number of open HTTP connections (idle and active)",
|
||||
"Current number of open HTTP connections (idle and active)",
|
||||
registry=self.registry,
|
||||
)
|
||||
self.errors_total = Counter(
|
||||
@@ -131,7 +131,7 @@ class MetricsService:
|
||||
registry=self.registry,
|
||||
)
|
||||
|
||||
# Initialize all known label combinations so they start at 0
|
||||
# Initialize known counter label combinations so they start at 0
|
||||
for notification_type in NotificationType:
|
||||
for result in ("success", "failure"):
|
||||
self.notification_delivery_total.labels(
|
||||
@@ -186,10 +186,11 @@ class MetricsService:
|
||||
|
||||
@contextmanager
|
||||
def response_timer(self, domain: str) -> Generator[_ResponseTimer]:
|
||||
"""Return a context manager that times an HTTP request.
|
||||
"""Return a context manager that times an Owncast API query.
|
||||
|
||||
Call ``timer.success()`` inside the block to record the duration.
|
||||
If ``success()`` is never called, nothing is recorded.
|
||||
If ``success()`` is never called, any existing timing label for the
|
||||
domain is removed.
|
||||
|
||||
:param domain: The stream domain being queried.
|
||||
"""
|
||||
|
||||
@@ -183,10 +183,11 @@ class NotificationService:
|
||||
)
|
||||
|
||||
def get_last_notification_time(self, domain: str) -> float:
|
||||
"""Get the timestamp of the last notification sent for a domain.
|
||||
"""Get the cached live/title-change notification timestamp for a domain.
|
||||
|
||||
:param domain: The stream domain.
|
||||
:return: Unix timestamp of last notification, or 0 if never notified.
|
||||
:return: Monotonic timestamp of the last successful live/title-change
|
||||
notification, or 0 if none has been recorded.
|
||||
"""
|
||||
return self.notification_timers_cache.get(domain, 0)
|
||||
|
||||
@@ -205,7 +206,7 @@ class NotificationService:
|
||||
:param room_id: The Matrix room ID to send to.
|
||||
:param body_text: The message body text.
|
||||
:param domain: The stream domain (for logging).
|
||||
:raises Exception: If sending fails.
|
||||
Propagates send errors so broadcast delivery accounting can count failures.
|
||||
"""
|
||||
try:
|
||||
content = TextMessageEventContent(msgtype=MessageType.TEXT, body=body_text)
|
||||
|
||||
@@ -100,8 +100,7 @@ class StreamMonitor:
|
||||
|
||||
total_streams = len(subscribed_domains)
|
||||
|
||||
# Build a list of async tasks for each stream domain
|
||||
# Run all stream updates in parallel and collect results
|
||||
# Run all stream updates concurrently and collect results.
|
||||
results = await asyncio.gather(
|
||||
*(self.update_stream(domain) for domain in subscribed_domains),
|
||||
return_exceptions=True,
|
||||
@@ -147,9 +146,9 @@ class StreamMonitor:
|
||||
async def update_stream(self, domain: str) -> bool:
|
||||
"""Update the state of a stream and send notifications as needed.
|
||||
|
||||
Sends notifications to subscribed Matrix rooms if a stream goes
|
||||
live. Implements progressive backoff for connection failures and
|
||||
auto-cleanup for dead instances.
|
||||
Sends notifications to subscribed Matrix rooms when a stream goes
|
||||
live or changes title. Implements progressive backoff for connection
|
||||
failures and auto-cleanup for dead instances.
|
||||
|
||||
:param domain: The domain of the stream to update.
|
||||
:return: True if check succeeded or was skipped, False on failure.
|
||||
@@ -179,8 +178,8 @@ class StreamMonitor:
|
||||
# Backoff is expected behavior, not a failure
|
||||
return True
|
||||
|
||||
# Flag: first state update of a brand-new stream to avoid
|
||||
# sending notifications if it's already live.
|
||||
# Flag: no connect/disconnect time has been recorded yet, so suppress
|
||||
# notifications for a stream whose initial history state is live.
|
||||
first_update = False
|
||||
|
||||
# Flag: whether to update the stream's state in the database.
|
||||
@@ -222,7 +221,7 @@ class StreamMonitor:
|
||||
old_state.last_connect_time is None
|
||||
and old_state.last_disconnect_time is None
|
||||
):
|
||||
# Yes, this is the first update. Don't send any notifications.
|
||||
# No stream history has been recorded yet. Don't send notifications.
|
||||
update_database = True
|
||||
first_update = True
|
||||
|
||||
@@ -242,7 +241,7 @@ class StreamMonitor:
|
||||
time.monotonic() - self.offline_timer_cache[domain]
|
||||
)
|
||||
|
||||
# Have we queried this stream before?
|
||||
# Has a prior connect/disconnect time been recorded?
|
||||
if not first_update:
|
||||
# Use fallback values if config fetch failed
|
||||
stream_name = stream_config.name if stream_config else domain
|
||||
@@ -284,7 +283,7 @@ class StreamMonitor:
|
||||
title_change=False,
|
||||
)
|
||||
else:
|
||||
# No, this is the first time we're querying
|
||||
# No stream history has been recorded yet.
|
||||
self.log.info(
|
||||
"[%s] Not sending notifications. This is the first state "
|
||||
"update for this stream.",
|
||||
@@ -305,9 +304,9 @@ class StreamMonitor:
|
||||
stream_name = stream_config.name if stream_config else domain
|
||||
stream_tags = stream_config.tags if stream_config else ()
|
||||
|
||||
# Was the last notification sent before the stream
|
||||
# last went offline? If so, send a regular go-live
|
||||
# instead of a title change to avoid confusion.
|
||||
# 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
|
||||
] > self.notification_service.get_last_notification_time(domain):
|
||||
|
||||
+9
-10
@@ -20,12 +20,11 @@ from typing import Any
|
||||
|
||||
UNKNOWN_STATUS_THRESHOLD = 15
|
||||
|
||||
# Maximum field lengths based on Owncast's configuration
|
||||
# Source: https://github.com/owncast/owncast/blob/master/
|
||||
# web/utils/config-constants.tsx
|
||||
_MAX_INSTANCE_TITLE_LENGTH = 255 # Server Name (line 81)
|
||||
_MAX_STREAM_TITLE_LENGTH = 100 # Stream Title (line 91)
|
||||
_MAX_TAG_LENGTH = 24 # Per tag (line 208)
|
||||
# Maximum field lengths based on Owncast's admin config constants.
|
||||
# Source: https://github.com/owncast/owncast/blob/master/web/utils/config-constants.tsx
|
||||
_MAX_INSTANCE_TITLE_LENGTH = 255 # TEXTFIELD_PROPS_SERVER_NAME maxLength
|
||||
_MAX_STREAM_TITLE_LENGTH = 100 # TEXTFIELD_PROPS_STREAM_TITLE maxLength
|
||||
_MAX_TAG_LENGTH = 24 # tag maxLength
|
||||
|
||||
|
||||
class InvalidApiResponseError(ValueError):
|
||||
@@ -102,10 +101,10 @@ class StreamState:
|
||||
|
||||
@property
|
||||
def status(self) -> StreamStatus:
|
||||
"""Derive stream status from failure count and connect times.
|
||||
"""Derive stream status from failure count and last connect time.
|
||||
|
||||
Returns UNKNOWN if failures exceed the threshold, ONLINE if a
|
||||
connect time is present, or OFFLINE otherwise.
|
||||
last connect time is present, or OFFLINE otherwise.
|
||||
"""
|
||||
if self.failure_counter > UNKNOWN_STATUS_THRESHOLD:
|
||||
return StreamStatus.UNKNOWN
|
||||
@@ -168,12 +167,12 @@ class StreamConfig:
|
||||
:return: StreamConfig instance.
|
||||
:raises InvalidApiResponseError: If the response shape is invalid.
|
||||
"""
|
||||
# Truncate instance name to max length
|
||||
# Apply Owncast's instance name limit.
|
||||
name = _truncate(
|
||||
_optional_config_str(response, "name"), _MAX_INSTANCE_TITLE_LENGTH
|
||||
)
|
||||
|
||||
# Truncate each tag to max length
|
||||
# Apply Owncast's per-tag limit.
|
||||
raw_tags = _optional_tag_list(response)
|
||||
tags = tuple([_truncate(tag, _MAX_TAG_LENGTH) for tag in raw_tags])
|
||||
|
||||
|
||||
+11
-11
@@ -60,15 +60,15 @@ class TestEscapeMarkdown:
|
||||
],
|
||||
)
|
||||
def test_escapes_special_chars(self, input_text: str, expected: str) -> None:
|
||||
"""Escape the given Markdown special character."""
|
||||
"""Escape Markdown special characters in input text."""
|
||||
assert _escape_markdown(input_text) == expected
|
||||
|
||||
|
||||
class TestSanitizeForMarkdown:
|
||||
"""Markdown sanitization combining newline removal and escaping."""
|
||||
"""Markdown sanitization combining whitespace normalization and escaping."""
|
||||
|
||||
def test_removes_newlines_and_escapes(self) -> None:
|
||||
"""Remove newlines and escape Markdown special characters."""
|
||||
"""Normalize newlines and escape Markdown special characters."""
|
||||
result = _sanitize_for_markdown("*bold*\nnew line")
|
||||
assert result == "\\*bold\\* new line"
|
||||
|
||||
@@ -171,7 +171,7 @@ class TestSubscribeCommand:
|
||||
async def test_subscribe_existing_stream_new_room(
|
||||
self, maubot_test_bot: TestBot
|
||||
) -> None:
|
||||
"""Skip instance validation when subscribing from a new room."""
|
||||
"""Skip validation when the domain already has subscriptions."""
|
||||
status_url = f"https://stream.logal.dev{_OWNCAST_STATUS_PATH}"
|
||||
with aioresponses() as mocked:
|
||||
mocked.get(
|
||||
@@ -180,9 +180,9 @@ class TestSubscribeCommand:
|
||||
)
|
||||
await maubot_test_bot.send("!subscribe stream.logal.dev")
|
||||
|
||||
# Subscribe from a different room; skips validation and should not
|
||||
# query the remote instance. An empty aioresponses context will raise
|
||||
# ConnectionError if any HTTP request is attempted.
|
||||
# Subscribe from a different room. The existing subscribed domain skips
|
||||
# validation, so an empty aioresponses context will raise ConnectionError
|
||||
# if any HTTP request is attempted.
|
||||
with aioresponses():
|
||||
await maubot_test_bot.send(
|
||||
"!subscribe stream.logal.dev", room_id="!otherroom:example.com"
|
||||
@@ -408,8 +408,8 @@ class TestSubscriptionsCommand:
|
||||
async def test_shows_multiple_subscriptions(
|
||||
self, maubot_test_bot: TestBot, maubot_plugin: OwncastSentry
|
||||
) -> None:
|
||||
"""List subscriptions alphabetically with mixed statuses."""
|
||||
# Subscribe in reverse alphabetical order to verify sorted output
|
||||
"""List subscriptions ordered by domain with mixed statuses."""
|
||||
# Subscribe in reverse domain order to verify domain-sorted output
|
||||
with aioresponses() as mocked:
|
||||
mocked.get(
|
||||
f"https://beta.com{_OWNCAST_STATUS_PATH}",
|
||||
@@ -543,8 +543,8 @@ class TestLiveCommand:
|
||||
async def test_shows_multiple_live_streams(
|
||||
self, maubot_test_bot: TestBot, maubot_plugin: OwncastSentry
|
||||
) -> None:
|
||||
"""List live streams alphabetically with different durations."""
|
||||
# Subscribe in reverse alphabetical order to verify sorted output
|
||||
"""List live streams ordered by domain with different durations."""
|
||||
# Subscribe in reverse domain order to verify domain-sorted output
|
||||
with aioresponses() as mocked:
|
||||
mocked.get(
|
||||
f"https://beta.com{_OWNCAST_STATUS_PATH}",
|
||||
|
||||
@@ -142,7 +142,7 @@ class TestSetSubscriptionCount:
|
||||
|
||||
|
||||
class TestSetCheckFailures:
|
||||
"""Check failure counter gauge per domain."""
|
||||
"""Consecutive check failure gauge per domain."""
|
||||
|
||||
def test_sets_count(self) -> None:
|
||||
"""Set the failure count for a domain."""
|
||||
@@ -196,7 +196,7 @@ class TestResponseTimer:
|
||||
assert len(matches) == 1
|
||||
|
||||
def test_does_not_record_on_exception(self) -> None:
|
||||
"""Do not record when the block raises an exception."""
|
||||
"""Do not record when the block raises before success() is called."""
|
||||
service = MetricsService()
|
||||
with (
|
||||
pytest.raises(ValueError, match="boom"),
|
||||
|
||||
@@ -272,7 +272,7 @@ class TestFormatMessage:
|
||||
|
||||
|
||||
class TestNotifyStreamLive:
|
||||
"""End-to-end notification sending."""
|
||||
"""Live notification sending behavior."""
|
||||
|
||||
async def test_sends_to_all_subscribed_rooms(
|
||||
self,
|
||||
@@ -393,7 +393,7 @@ class TestNotifyStreamLive:
|
||||
stream_repo: StreamRepository,
|
||||
subscription_repo: SubscriptionRepository,
|
||||
) -> None:
|
||||
"""Send to remaining rooms after a delivery failure."""
|
||||
"""Deliver to other subscribed rooms when one room delivery fails."""
|
||||
client = _StubMatrixClient()
|
||||
client.should_fail_for_rooms.add("!bad:matrix.org")
|
||||
service = _make_service(client=client, subscription_repo=subscription_repo)
|
||||
|
||||
@@ -405,7 +405,7 @@ class TestValidateInstance:
|
||||
|
||||
|
||||
class TestResponseTimeMetrics:
|
||||
"""Response time histogram recording."""
|
||||
"""Response time gauge recording."""
|
||||
|
||||
async def test_records_on_success(self) -> None:
|
||||
"""Record response time on a successful request."""
|
||||
|
||||
@@ -173,7 +173,7 @@ class TestManagerSubscribe:
|
||||
owncast_client: _StubOwncastClient,
|
||||
subscription_repo: SubscriptionRepository,
|
||||
) -> None:
|
||||
"""Known domains skip remote validation for additional room subscriptions."""
|
||||
"""Domains with existing subscriptions skip remote validation."""
|
||||
await manager.subscribe("!room1:example.com", "stream.example")
|
||||
owncast_client.valid = False
|
||||
|
||||
|
||||
+2
-2
@@ -50,7 +50,7 @@ class TestTruncate:
|
||||
],
|
||||
)
|
||||
def test_truncates(self, text: str, max_length: int, expected: str) -> None:
|
||||
"""Truncate text that exceeds the maximum length."""
|
||||
"""Return text unchanged within the limit and truncate over-limit text."""
|
||||
assert _truncate(text, max_length) == expected
|
||||
|
||||
|
||||
@@ -117,7 +117,7 @@ class TestStreamStateFromApiResponse:
|
||||
"""StreamState construction from an API response dictionary."""
|
||||
|
||||
def test_typical_response(self) -> None:
|
||||
"""Populate all fields from a complete API response."""
|
||||
"""Populate API-derived fields from a complete stream state response."""
|
||||
response = {
|
||||
"streamTitle": "My Stream",
|
||||
"lastConnectTime": "2026-01-01T00:00:00Z",
|
||||
|
||||
Reference in New Issue
Block a user