Switched stream status tracking to explicit online state.
This commit is contained in:
@@ -78,7 +78,7 @@ def _format_duration(timestamp_str: str, now: datetime) -> str:
|
||||
if seconds < _SECONDS_PER_DAY:
|
||||
hours = seconds // _SECONDS_PER_HOUR
|
||||
return f"{hours} hour{'s' if hours != 1 else ''}"
|
||||
except TypeError, ValueError:
|
||||
except (TypeError, ValueError):
|
||||
return "unknown duration"
|
||||
else:
|
||||
days = seconds // _SECONDS_PER_DAY
|
||||
@@ -185,13 +185,15 @@ class CommandHandler:
|
||||
|
||||
# Determine status and duration (as a sub-bullet)
|
||||
match stream_state.status:
|
||||
case StreamStatus.ONLINE if stream_state.last_connect_time:
|
||||
duration = _format_duration(stream_state.last_connect_time, now)
|
||||
case StreamStatus.ONLINE if stream_state.status_since:
|
||||
duration = _format_duration(stream_state.status_since, now)
|
||||
parts.append(f" - Status: Online for {duration} \n")
|
||||
case StreamStatus.ONLINE:
|
||||
parts.append(" - Status: Online \n")
|
||||
case StreamStatus.UNKNOWN:
|
||||
parts.append(" - Status: Unknown (instance unreachable) \n")
|
||||
case StreamStatus.OFFLINE if stream_state.last_disconnect_time:
|
||||
duration = _format_duration(stream_state.last_disconnect_time, now)
|
||||
case StreamStatus.OFFLINE if stream_state.status_since:
|
||||
duration = _format_duration(stream_state.status_since, now)
|
||||
parts.append(f" - Status: Offline for {duration} \n")
|
||||
case StreamStatus.OFFLINE:
|
||||
parts.append(" - Status: Offline \n")
|
||||
@@ -255,9 +257,11 @@ class CommandHandler:
|
||||
parts.append(f" - Title: {safe_title} \n")
|
||||
|
||||
# Add status with duration
|
||||
if stream_state.last_connect_time:
|
||||
duration = _format_duration(stream_state.last_connect_time, now)
|
||||
if stream_state.status_since:
|
||||
duration = _format_duration(stream_state.status_since, now)
|
||||
parts.append(f" - Online for {duration} \n")
|
||||
else:
|
||||
parts.append(" - Online \n")
|
||||
|
||||
# Add stream link
|
||||
parts.append(f" - Link: https://{domain}\n\n")
|
||||
|
||||
@@ -15,6 +15,7 @@
|
||||
"""HTTP client for querying Owncast instance APIs."""
|
||||
|
||||
import json
|
||||
from datetime import UTC, datetime
|
||||
from http import HTTPStatus
|
||||
from typing import TYPE_CHECKING, Any
|
||||
|
||||
@@ -125,8 +126,11 @@ class OwncastClient:
|
||||
if new_state is None:
|
||||
return None
|
||||
|
||||
observed_at = datetime.now(UTC)
|
||||
try:
|
||||
stream_state = StreamState.from_api_response(new_state, domain)
|
||||
stream_state = StreamState.from_api_response(
|
||||
new_state, domain, observed_at
|
||||
)
|
||||
except InvalidApiResponseError as e:
|
||||
self.log.warning(
|
||||
"[%s] Rejecting response to request on %s as response "
|
||||
|
||||
@@ -14,6 +14,7 @@
|
||||
|
||||
"""Repository and schema upgrade definitions for OwncastSentry."""
|
||||
|
||||
from datetime import UTC, datetime
|
||||
from typing import TYPE_CHECKING, Any
|
||||
|
||||
from mautrix.util.async_db import Connection, UpgradeTable
|
||||
@@ -24,6 +25,7 @@ from .types import (
|
||||
NotSubscribedError,
|
||||
RoomSubscription,
|
||||
StreamState,
|
||||
format_status_since,
|
||||
)
|
||||
|
||||
if TYPE_CHECKING:
|
||||
@@ -33,6 +35,33 @@ if TYPE_CHECKING:
|
||||
upgrade_table = UpgradeTable()
|
||||
|
||||
|
||||
def _has_legacy_timestamp(value: Any) -> bool:
|
||||
"""Return whether a legacy timestamp value carries usable content."""
|
||||
return value is not None and str(value).strip() != ""
|
||||
|
||||
|
||||
def _normalize_legacy_status_since(value: Any) -> str | None:
|
||||
"""Normalize a legacy timestamp value to the canonical UTC format."""
|
||||
if value is None:
|
||||
return None
|
||||
if isinstance(value, datetime):
|
||||
parsed = value
|
||||
elif isinstance(value, str):
|
||||
value = value.strip()
|
||||
if not value:
|
||||
return None
|
||||
try:
|
||||
parsed = datetime.fromisoformat(value)
|
||||
except ValueError:
|
||||
return None
|
||||
else:
|
||||
return None
|
||||
|
||||
if parsed.tzinfo is None:
|
||||
return None
|
||||
return format_status_since(parsed.astimezone(UTC))
|
||||
|
||||
|
||||
@upgrade_table.register( # type: ignore[arg-type, call-arg, untyped-decorator]
|
||||
description="Initial revision"
|
||||
)
|
||||
@@ -107,6 +136,56 @@ async def upgrade_v3(conn: Connection) -> None:
|
||||
)
|
||||
|
||||
|
||||
@upgrade_table.register( # type: ignore[arg-type, call-arg, untyped-decorator]
|
||||
description="Store online status and current status timestamp"
|
||||
)
|
||||
async def upgrade_v4(conn: Connection) -> None:
|
||||
"""Upgrade database schema to version 4 format.
|
||||
|
||||
Replaces separate connect and disconnect timestamp columns with an
|
||||
authoritative online flag and a timestamp for the current status.
|
||||
|
||||
:param conn: A connection to run the v4 database migration on.
|
||||
"""
|
||||
await conn.execute(
|
||||
"""CREATE TABLE "streams_new" (
|
||||
"domain" TEXT NOT NULL UNIQUE,
|
||||
"name" TEXT,
|
||||
"title" TEXT,
|
||||
"online" BOOLEAN NOT NULL DEFAULT false,
|
||||
"status_since" TEXT,
|
||||
"failure_counter" INTEGER DEFAULT 0,
|
||||
PRIMARY KEY("domain")
|
||||
)"""
|
||||
)
|
||||
|
||||
rows = await conn.fetch(
|
||||
"""SELECT domain, name, title, last_connect_time, last_disconnect_time,
|
||||
failure_counter
|
||||
FROM streams"""
|
||||
)
|
||||
for row in rows:
|
||||
online = _has_legacy_timestamp(row["last_connect_time"])
|
||||
legacy_timestamp = (
|
||||
row["last_connect_time"] if online else row["last_disconnect_time"]
|
||||
)
|
||||
await conn.execute(
|
||||
"""INSERT INTO streams_new (
|
||||
domain, name, title, online, status_since, failure_counter
|
||||
)
|
||||
VALUES ($1, $2, $3, $4, $5, $6)""",
|
||||
row["domain"],
|
||||
row["name"],
|
||||
row["title"],
|
||||
online,
|
||||
_normalize_legacy_status_since(legacy_timestamp),
|
||||
row["failure_counter"],
|
||||
)
|
||||
|
||||
await conn.execute("DROP TABLE streams")
|
||||
await conn.execute("ALTER TABLE streams_new RENAME TO streams")
|
||||
|
||||
|
||||
def get_upgrade_table() -> UpgradeTable:
|
||||
"""Return the repository upgrade table with registered migrations."""
|
||||
return upgrade_table
|
||||
@@ -164,15 +243,15 @@ class StreamRepository:
|
||||
:param state: The StreamState to save.
|
||||
"""
|
||||
query = """UPDATE streams
|
||||
SET name=$1, title=$2, last_connect_time=$3, last_disconnect_time=$4
|
||||
SET name=$1, title=$2, online=$3, status_since=$4
|
||||
WHERE domain=$5"""
|
||||
async with self.db.acquire() as conn:
|
||||
await conn.execute(
|
||||
query,
|
||||
state.name,
|
||||
state.title,
|
||||
state.last_connect_time,
|
||||
state.last_disconnect_time,
|
||||
state.online,
|
||||
state.status_since,
|
||||
state.domain,
|
||||
)
|
||||
|
||||
@@ -313,7 +392,7 @@ class SubscriptionRepository:
|
||||
FROM subscriptions
|
||||
JOIN streams ON streams.domain = subscriptions.stream_domain
|
||||
WHERE subscriptions.room_id=$1
|
||||
AND streams.last_connect_time IS NOT NULL
|
||||
AND streams.online=true
|
||||
AND streams.failure_counter <= $2
|
||||
ORDER BY streams.domain"""
|
||||
async with self.db.acquire() as conn:
|
||||
|
||||
@@ -18,7 +18,7 @@ import asyncio
|
||||
import time
|
||||
from typing import TYPE_CHECKING
|
||||
|
||||
from .types import StreamState, StreamStatus, UpdateResult
|
||||
from .types import StreamState, UpdateResult
|
||||
|
||||
if TYPE_CHECKING:
|
||||
import logging
|
||||
@@ -178,8 +178,8 @@ class StreamMonitor:
|
||||
# Backoff is expected behavior, not a failure
|
||||
return True
|
||||
|
||||
# Flag: no connect/disconnect time has been recorded yet, so suppress
|
||||
# notifications for a stream whose initial history state is live.
|
||||
# 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.
|
||||
@@ -216,20 +216,20 @@ class StreamMonitor:
|
||||
# Initialize timer cache entries to prevent KeyError on first access
|
||||
self.offline_timer_cache.setdefault(domain, 0)
|
||||
|
||||
# Does the last known stream state lack connect/disconnect?
|
||||
if (
|
||||
old_state.last_connect_time is None
|
||||
and old_state.last_disconnect_time is None
|
||||
):
|
||||
if old_state.status_since is None:
|
||||
# No stream history has been recorded yet. Don't send notifications.
|
||||
update_database = True
|
||||
first_update = True
|
||||
|
||||
# Does the new state have a connect time but the old one not?
|
||||
if (
|
||||
new_state.last_connect_time is not None
|
||||
and old_state.last_connect_time is None
|
||||
):
|
||||
if first_update:
|
||||
self.log.info(
|
||||
"[%s] Not sending notifications. This is the first state "
|
||||
"update for this stream.",
|
||||
domain,
|
||||
)
|
||||
|
||||
# Did the stream become publicly online?
|
||||
elif new_state.online and not old_state.online:
|
||||
# Yes! This stream is now live.
|
||||
update_database = True
|
||||
stream_config = await self.owncast_client.get_stream_config(domain)
|
||||
@@ -241,59 +241,43 @@ class StreamMonitor:
|
||||
time.monotonic() - self.offline_timer_cache[domain]
|
||||
)
|
||||
|
||||
# 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
|
||||
stream_tags = stream_config.tags if stream_config else ()
|
||||
# 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(
|
||||
"[%s] Not sending notifications. Stream was only "
|
||||
"offline for %s of %s seconds and did not change "
|
||||
"its title.",
|
||||
domain,
|
||||
seconds_since_last_offline,
|
||||
_TEMPORARY_OFFLINE_NOTIFICATION_COOLDOWN,
|
||||
)
|
||||
else:
|
||||
# Offline for a while. Send a normal notification.
|
||||
# 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=False,
|
||||
title_change=True,
|
||||
)
|
||||
else:
|
||||
# Briefly offline, no title change. Skip.
|
||||
self.log.info(
|
||||
"[%s] Not sending notifications. Stream was only "
|
||||
"offline for %s of %s seconds and did not change "
|
||||
"its title.",
|
||||
domain,
|
||||
seconds_since_last_offline,
|
||||
_TEMPORARY_OFFLINE_NOTIFICATION_COOLDOWN,
|
||||
)
|
||||
else:
|
||||
# No stream history has been recorded yet.
|
||||
self.log.info(
|
||||
"[%s] Not sending notifications. This is the first state "
|
||||
"update for this stream.",
|
||||
# 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,
|
||||
)
|
||||
|
||||
if (
|
||||
new_state.last_connect_time is not None
|
||||
and old_state.last_connect_time is not None
|
||||
):
|
||||
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)
|
||||
@@ -327,11 +311,8 @@ class StreamMonitor:
|
||||
title_change=True,
|
||||
)
|
||||
|
||||
# Did the stream go offline (old had connect, new doesn't)?
|
||||
elif (
|
||||
new_state.last_connect_time is None
|
||||
and old_state.last_connect_time is not None
|
||||
):
|
||||
# Did the stream go offline?
|
||||
elif not new_state.online and old_state.online:
|
||||
# Yep. This stream is now offline. Log it.
|
||||
update_database = True
|
||||
self.offline_timer_cache[domain] = time.monotonic()
|
||||
@@ -348,23 +329,25 @@ class StreamMonitor:
|
||||
|
||||
self.log.debug("[%s] Updating stream state in database...", domain)
|
||||
|
||||
if first_update or old_state.online != new_state.online:
|
||||
status_since = new_state.status_since
|
||||
else:
|
||||
status_since = old_state.status_since
|
||||
|
||||
# Create updated state object (title already truncated in new_state)
|
||||
updated_state = StreamState(
|
||||
domain=domain,
|
||||
name=stream_name,
|
||||
title=new_state.title,
|
||||
last_connect_time=new_state.last_connect_time,
|
||||
last_disconnect_time=new_state.last_disconnect_time,
|
||||
online=new_state.online,
|
||||
status_since=status_since,
|
||||
)
|
||||
|
||||
await self.stream_repo.update(updated_state)
|
||||
|
||||
# All done.
|
||||
self.log.debug("[%s] State update completed.", domain)
|
||||
if new_state.last_connect_time is not None:
|
||||
self.metrics.set_stream_status(domain, StreamStatus.ONLINE)
|
||||
else:
|
||||
self.metrics.set_stream_status(domain, StreamStatus.OFFLINE)
|
||||
self.metrics.set_stream_status(domain, new_state.status)
|
||||
return True
|
||||
|
||||
async def _check_cleanup_thresholds(self, domain: str, counter: int) -> None:
|
||||
|
||||
+22
-21
@@ -15,6 +15,7 @@
|
||||
"""Data containers and domain errors for OwncastSentry."""
|
||||
|
||||
from dataclasses import dataclass
|
||||
from datetime import UTC, datetime
|
||||
from enum import Enum
|
||||
from typing import Any
|
||||
|
||||
@@ -47,14 +48,6 @@ def _require_str(response: dict[str, Any], field: str) -> str:
|
||||
return value
|
||||
|
||||
|
||||
def _require_nullable_str(response: dict[str, Any], field: str) -> str | None:
|
||||
"""Return a required nullable string API response field."""
|
||||
value = _require_field(response, field)
|
||||
if value is not None and not isinstance(value, str):
|
||||
raise InvalidApiResponseError(f"{field} must be a string or null")
|
||||
return value
|
||||
|
||||
|
||||
def _optional_config_str(response: dict[str, Any], field: str) -> str:
|
||||
"""Return an optional config string, defaulting to empty when absent."""
|
||||
value = response.get(field, "")
|
||||
@@ -80,6 +73,13 @@ def _truncate(text: str, max_length: int) -> str:
|
||||
return text[:max_length]
|
||||
|
||||
|
||||
def format_status_since(timestamp: datetime) -> str:
|
||||
"""Format a status timestamp using the package's canonical UTC format."""
|
||||
if timestamp.tzinfo is None:
|
||||
timestamp = timestamp.replace(tzinfo=UTC)
|
||||
return timestamp.astimezone(UTC).isoformat(timespec="seconds")
|
||||
|
||||
|
||||
class StreamStatus(Enum):
|
||||
"""Represents the status of a stream."""
|
||||
|
||||
@@ -95,35 +95,36 @@ class StreamState:
|
||||
domain: str
|
||||
name: str | None = None
|
||||
title: str | None = None
|
||||
last_connect_time: str | None = None
|
||||
last_disconnect_time: str | None = None
|
||||
online: bool = False
|
||||
status_since: str | None = None
|
||||
failure_counter: int = 0
|
||||
|
||||
@property
|
||||
def status(self) -> StreamStatus:
|
||||
"""Derive stream status from failure count and last connect time.
|
||||
"""Derive stream status from failure count and online state.
|
||||
|
||||
Returns UNKNOWN if failures exceed the threshold, ONLINE if a
|
||||
last connect time is present, or OFFLINE otherwise.
|
||||
Returns UNKNOWN if failures exceed the threshold, ONLINE if the
|
||||
stream is online, or OFFLINE otherwise.
|
||||
"""
|
||||
if self.failure_counter > UNKNOWN_STATUS_THRESHOLD:
|
||||
return StreamStatus.UNKNOWN
|
||||
if self.last_connect_time is not None:
|
||||
if self.online:
|
||||
return StreamStatus.ONLINE
|
||||
return StreamStatus.OFFLINE
|
||||
|
||||
@classmethod
|
||||
def from_api_response(cls, response: dict[str, Any], domain: str) -> StreamState:
|
||||
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")
|
||||
last_connect_time = _require_nullable_str(response, "lastConnectTime")
|
||||
last_disconnect_time = _require_nullable_str(response, "lastDisconnectTime")
|
||||
online = _require_field(response, "online")
|
||||
if not isinstance(online, bool):
|
||||
raise InvalidApiResponseError("online must be a boolean")
|
||||
@@ -131,8 +132,8 @@ class StreamState:
|
||||
return cls(
|
||||
domain=domain,
|
||||
title=_truncate(stream_title, _MAX_STREAM_TITLE_LENGTH),
|
||||
last_connect_time=last_connect_time,
|
||||
last_disconnect_time=last_disconnect_time,
|
||||
online=online,
|
||||
status_since=format_status_since(observed_at),
|
||||
)
|
||||
|
||||
@classmethod
|
||||
@@ -146,8 +147,8 @@ class StreamState:
|
||||
domain=row["domain"],
|
||||
name=row["name"],
|
||||
title=row["title"],
|
||||
last_connect_time=row["last_connect_time"],
|
||||
last_disconnect_time=row["last_disconnect_time"],
|
||||
online=bool(row["online"]),
|
||||
status_since=row["status_since"],
|
||||
failure_counter=row["failure_counter"],
|
||||
)
|
||||
|
||||
|
||||
Reference in New Issue
Block a user