3 Commits
Author SHA1 Message Date
LogalDeveloper 9be932c0c7 Replaced health checker with Prometheus metrics service.
Audit / Dependencies (push) Failing after 8s
CD / Build (push) Successful in 8s
CI / Formatting (push) Successful in 6s
CI / Linting (push) Successful in 6s
CI / Tests (push) Successful in 11s
CI / Type Checking (push) Successful in 14s
CI / Spelling (push) Successful in 13s
2026-03-25 08:54:38 -04:00
LogalDeveloper 295a8c90a9 Improved code quality with more idiomatic Python patterns.
CD / Build (push) Successful in 7s
CI / Formatting (push) Successful in 5s
CI / Linting (push) Successful in 5s
CI / Tests (push) Successful in 11s
CI / Type Checking (push) Successful in 20s
CI / Spelling (push) Successful in 32s
2026-03-24 11:02:13 -04:00
LogalDeveloper 653c27c97c Updated dependencies.
CD / Build (push) Successful in 8s
CI / Formatting (push) Successful in 5s
CI / Linting (push) Successful in 5s
CI / Tests (push) Successful in 22s
CI / Type Checking (push) Successful in 13s
CI / Spelling (push) Successful in 9s
Audit / Dependencies (push) Successful in 9s
2026-03-19 09:25:30 -04:00
24 changed files with 1322 additions and 691 deletions
+4 -4
View File
@@ -1,4 +1,4 @@
# Health check endpoint URL. # Prometheus metrics endpoint.
# If configured, a GET request will be sent to this URL after each successful update cycle. # When enabled, a /metrics endpoint is exposed via the maubot webapp.
# Leave empty to disable health check reporting. # Disabled by default.
health_check_endpoint: "" metrics_enabled: false
+3
View File
@@ -8,5 +8,8 @@ main_class: OwncastSentry
database: true database: true
database_type: asyncpg database_type: asyncpg
config: true config: true
webapp: true
dependencies:
- prometheus_client>=0.24.1
extra_files: extra_files:
- base-config.yaml - base-config.yaml
+61 -21
View File
@@ -16,13 +16,15 @@
from typing import TYPE_CHECKING from typing import TYPE_CHECKING
from aiohttp import web
from maubot import MessageEvent, Plugin # type: ignore[attr-defined] from maubot import MessageEvent, Plugin # type: ignore[attr-defined]
from maubot.handlers import command from maubot.handlers import command
from prometheus_client.exposition import choose_encoder
from .commands import CommandHandler from .commands import CommandHandler
from .config import Config from .config import Config
from .database import StreamRepository, SubscriptionRepository from .database import StreamRepository, SubscriptionRepository
from .health_checker import HealthChecker from .metrics import ErrorSource, MetricsService
from .migrations import get_upgrade_table from .migrations import get_upgrade_table
from .notification_service import NotificationService from .notification_service import NotificationService
from .owncast_client import OwncastClient from .owncast_client import OwncastClient
@@ -57,8 +59,21 @@ class OwncastSentry(Plugin):
config.load_and_update() config.load_and_update()
db: Database = self.database # type: ignore[assignment] db: Database = self.database # type: ignore[assignment]
# Initialize metrics service and register web endpoint (if enabled)
self.metrics_service = MetricsService()
self.metrics_service.set_build_info(str(self.loader.meta.version))
if config.metrics_enabled and self.webapp is not None:
self.webapp.add_route(
method="GET", path="/metrics", handler=self._metrics_endpoint
)
# Initialize the Owncast API client # Initialize the Owncast API client
self.owncast_client = OwncastClient(self.log, str(self.loader.meta.version)) self.owncast_client = OwncastClient(
self.log, str(self.loader.meta.version), metrics=self.metrics_service
)
self.metrics_service.register_open_connections_gauge(
lambda: self.owncast_client.open_connection_count
)
# Initialize repositories # Initialize repositories
self.stream_repo = StreamRepository(db) self.stream_repo = StreamRepository(db)
@@ -66,7 +81,10 @@ class OwncastSentry(Plugin):
# Initialize notification service # Initialize notification service
self.notification_service = NotificationService( self.notification_service = NotificationService(
self.client, self.subscription_repo, self.log self.client,
self.subscription_repo,
self.log,
metrics=self.metrics_service,
) )
# Initialize stream monitor # Initialize stream monitor
@@ -76,13 +94,7 @@ class OwncastSentry(Plugin):
self.subscription_repo, self.subscription_repo,
self.notification_service, self.notification_service,
self.log, self.log,
) metrics=self.metrics_service,
# Initialize health checker
self.health_checker = HealthChecker(
db,
self.owncast_client,
self.log,
) )
# Initialize command handler # Initialize command handler
@@ -97,41 +109,69 @@ class OwncastSentry(Plugin):
self.sched.run_periodically(60, self._update_all_stream_states) self.sched.run_periodically(60, self._update_all_stream_states)
async def _update_all_stream_states(self) -> None: async def _update_all_stream_states(self) -> None:
"""Update all stream states and perform health check.""" """Update all stream states."""
try:
# Get list of all stream domains with active subscriptions # Get list of all stream domains with active subscriptions
subscribed_domains = await self.subscription_repo.get_all_subscribed_domains() subscribed_domains = (
await self.subscription_repo.get_all_subscribed_domains()
# Delegate to stream monitor and get results
update_result = await self.stream_monitor.update_all_streams(subscribed_domains)
# Perform health check
config: Config = self.config # type: ignore[assignment]
await self.health_checker.perform_health_check(
update_result,
config.health_check_endpoint,
) )
# Delegate to stream monitor
await self.stream_monitor.update_all_streams(subscribed_domains)
except Exception:
self.metrics_service.record_error(ErrorSource.SCHEDULER_LOOP)
self.log.exception("Unhandled exception in scheduler loop.")
@command.new(help="Subscribes to a new Owncast stream.") @command.new(help="Subscribes to a new Owncast stream.")
@command.argument("url") @command.argument("url")
async def subscribe(self, evt: MessageEvent, url: str) -> None: async def subscribe(self, evt: MessageEvent, url: str) -> None:
"""Delegate subscribe command to CommandHandler.""" """Delegate subscribe command to CommandHandler."""
try:
await self.command_handler.subscribe(evt, url) await self.command_handler.subscribe(evt, url)
except Exception:
self.metrics_service.record_error(ErrorSource.COMMAND)
self.log.exception("Unhandled exception in subscribe command.")
await evt.reply("An unexpected error occurred. Please try again later.")
@command.new(help="Unsubscribes from an Owncast stream.") @command.new(help="Unsubscribes from an Owncast stream.")
@command.argument("url") @command.argument("url")
async def unsubscribe(self, evt: MessageEvent, url: str) -> None: async def unsubscribe(self, evt: MessageEvent, url: str) -> None:
"""Delegate unsubscribe command to CommandHandler.""" """Delegate unsubscribe command to CommandHandler."""
try:
await self.command_handler.unsubscribe(evt, url) await self.command_handler.unsubscribe(evt, url)
except Exception:
self.metrics_service.record_error(ErrorSource.COMMAND)
self.log.exception("Unhandled exception in unsubscribe command.")
await evt.reply("An unexpected error occurred. Please try again later.")
@command.new(help="Lists all stream subscriptions in this room.") @command.new(help="Lists all stream subscriptions in this room.")
async def subscriptions(self, evt: MessageEvent) -> None: async def subscriptions(self, evt: MessageEvent) -> None:
"""Delegate subscriptions command to CommandHandler.""" """Delegate subscriptions command to CommandHandler."""
try:
await self.command_handler.subscriptions(evt) await self.command_handler.subscriptions(evt)
except Exception:
self.metrics_service.record_error(ErrorSource.COMMAND)
self.log.exception("Unhandled exception in subscriptions command.")
await evt.reply("An unexpected error occurred. Please try again later.")
@command.new(help="Lists currently live streams in this room.") @command.new(help="Lists currently live streams in this room.")
async def live(self, evt: MessageEvent) -> None: async def live(self, evt: MessageEvent) -> None:
"""Delegate live command to CommandHandler.""" """Delegate live command to CommandHandler."""
try:
await self.command_handler.live(evt) await self.command_handler.live(evt)
except Exception:
self.metrics_service.record_error(ErrorSource.COMMAND)
self.log.exception("Unhandled exception in live command.")
await evt.reply("An unexpected error occurred. Please try again later.")
async def _metrics_endpoint(self, request: web.Request) -> web.Response:
"""Serve Prometheus metrics."""
accept = request.headers.get("Accept", "")
encoder, content_type = choose_encoder(accept)
output = encoder(self.metrics_service.registry)
response = web.Response(body=output)
response.headers["Content-Type"] = content_type
return response
async def stop(self) -> None: async def stop(self) -> None:
"""Clean up resources by closing the HTTP session.""" """Clean up resources by closing the HTTP session."""
+32 -40
View File
@@ -83,9 +83,7 @@ class CommandHandler:
except sqlite3.IntegrityError: except sqlite3.IntegrityError:
# Room is already subscribed. # Room is already subscribed.
await evt.reply( await evt.reply(
"This room is already subscribed to notifications for " f"This room is already subscribed to notifications for {stream_domain}."
+ stream_domain
+ "."
) )
return return
@@ -101,9 +99,8 @@ class CommandHandler:
# All went well! Tell the user. # All went well! Tell the user.
self.log.info(f"[{stream_domain}] Subscription added for room {evt.room_id}.") self.log.info(f"[{stream_domain}] Subscription added for room {evt.room_id}.")
await evt.reply( await evt.reply(
"Subscription added! This room will receive notifications when " f"Subscription added! This room will receive "
+ stream_domain f"notifications when {stream_domain} goes live."
+ " goes live."
) )
async def unsubscribe(self, evt: MessageEvent, url: str) -> None: async def unsubscribe(self, evt: MessageEvent, url: str) -> None:
@@ -125,15 +122,14 @@ class CommandHandler:
f"[{stream_domain}] Subscription removed for room {evt.room_id}." f"[{stream_domain}] Subscription removed for room {evt.room_id}."
) )
await evt.reply( await evt.reply(
"Subscription removed! This room will no " f"Subscription removed! This room will no "
"longer receive notifications for " + stream_domain + "." f"longer receive notifications for {stream_domain}."
) )
else: else:
# No, nothing changed. Tell the user. # No, nothing changed. Tell the user.
await evt.reply( await evt.reply(
"This room is already not subscribed to notifications for " "This room is already not subscribed to "
+ stream_domain f"notifications for {stream_domain}."
+ "."
) )
def _format_duration(self, timestamp_str: str) -> str: def _format_duration(self, timestamp_str: str) -> str:
@@ -143,7 +139,7 @@ class CommandHandler:
:return: Formatted duration string (e.g., "1 hour", "2 days"). :return: Formatted duration string (e.g., "1 hour", "2 days").
""" """
try: try:
timestamp = datetime.fromisoformat(timestamp_str.replace("Z", "+00:00")) timestamp = datetime.fromisoformat(timestamp_str)
now = datetime.now(UTC) now = datetime.now(UTC)
delta = now - timestamp delta = now - timestamp
@@ -158,7 +154,7 @@ class CommandHandler:
return f"{hours} hour{'s' if hours != 1 else ''}" return f"{hours} hour{'s' if hours != 1 else ''}"
days = seconds // 86400 days = seconds // 86400
return f"{days} day{'s' if days != 1 else ''}" return f"{days} day{'s' if days != 1 else ''}"
except Exception: except ValueError:
return "unknown duration" return "unknown duration"
async def subscriptions(self, evt: MessageEvent) -> None: async def subscriptions(self, evt: MessageEvent) -> None:
@@ -183,7 +179,7 @@ class CommandHandler:
# Build the response message body as Markdown # Build the response message body as Markdown
count = len(subscribed_domains) count = len(subscribed_domains)
body_text = f"**Subscriptions for this room ({count}):**\n\n" parts = [f"**Subscriptions for this room ({count}):**\n\n"]
for domain in subscribed_domains: for domain in subscribed_domains:
# Get the stream state from the database # Get the stream state from the database
@@ -192,45 +188,41 @@ class CommandHandler:
continue continue
# Determine stream name (use domain as fallback) # Determine stream name (use domain as fallback)
stream_name = stream_state.name if stream_state.name else domain stream_name = stream_state.name or domain
safe_stream_name = sanitize_for_markdown(stream_name) safe_stream_name = sanitize_for_markdown(stream_name)
# Start building this stream's entry with stream name as main bullet # Start building this stream's entry with stream name as main bullet
body_text += f"- **{safe_stream_name}** \n" parts.append(f"- **{safe_stream_name}** \n")
# Add title if stream is online (as a sub-bullet) # Add title if stream is online (as a sub-bullet)
if stream_state.status == StreamStatus.ONLINE and stream_state.title: if stream_state.status == StreamStatus.ONLINE and stream_state.title:
safe_title = sanitize_for_markdown(stream_state.title) safe_title = sanitize_for_markdown(stream_state.title)
body_text += f" - Title: {safe_title} \n" parts.append(f" - Title: {safe_title} \n")
# Determine status and duration (as a sub-bullet) # Determine status and duration (as a sub-bullet)
if stream_state.status == StreamStatus.ONLINE: match stream_state.status:
# Stream is online - use last_connect_time case StreamStatus.ONLINE if stream_state.last_connect_time:
if stream_state.last_connect_time:
duration = self._format_duration(stream_state.last_connect_time) duration = self._format_duration(stream_state.last_connect_time)
body_text += f" - Status: Online for {duration} \n" parts.append(f" - Status: Online for {duration} \n")
elif stream_state.status == StreamStatus.UNKNOWN: case StreamStatus.UNKNOWN:
# Stream status is unknown - instance unreachable parts.append(" - Status: Unknown (instance unreachable) \n")
body_text += " - Status: Unknown (instance unreachable) \n" case StreamStatus.OFFLINE if stream_state.last_disconnect_time:
else:
# Stream is offline - use last_disconnect_time
if stream_state.last_disconnect_time:
duration = self._format_duration(stream_state.last_disconnect_time) duration = self._format_duration(stream_state.last_disconnect_time)
body_text += f" - Status: Offline for {duration} \n" parts.append(f" - Status: Offline for {duration} \n")
else: case StreamStatus.OFFLINE:
body_text += " - Status: Offline \n" parts.append(" - Status: Offline \n")
# Add stream link (as a sub-bullet) # Add stream link (as a sub-bullet)
body_text += f" - Link: https://{domain}\n\n" parts.append(f" - Link: https://{domain}\n\n")
# Add help text for unsubscribing # Add help text for unsubscribing
body_text += ( parts.append(
"\nTo unsubscribe from any of these Owncast " "\nTo unsubscribe from any of these Owncast "
"instances, use `!unsubscribe <domain>`" "instances, use `!unsubscribe <domain>`"
) )
# Send the response as Markdown # Send the response as Markdown
await evt.reply(body_text, markdown=True) await evt.reply("".join(parts), markdown=True)
async def live(self, evt: MessageEvent) -> None: async def live(self, evt: MessageEvent) -> None:
"""List currently live streams in the current room. """List currently live streams in the current room.
@@ -271,28 +263,28 @@ class CommandHandler:
# Build the response message body as Markdown # Build the response message body as Markdown
count = len(live_streams) count = len(live_streams)
body_text = f"**Live Owncast instances ({count}):**\n\n" parts = [f"**Live Owncast instances ({count}):**\n\n"]
for domain, stream_state in live_streams: for domain, stream_state in live_streams:
# Determine stream name (use domain as fallback) # Determine stream name (use domain as fallback)
stream_name = stream_state.name if stream_state.name else domain stream_name = stream_state.name or domain
safe_stream_name = sanitize_for_markdown(stream_name) safe_stream_name = sanitize_for_markdown(stream_name)
# Start building this stream's entry with stream name as main bullet # Start building this stream's entry with stream name as main bullet
body_text += f"- **{safe_stream_name}** \n" parts.append(f"- **{safe_stream_name}** \n")
# Add title (should be present for live streams) # Add title (should be present for live streams)
if stream_state.title: if stream_state.title:
safe_title = sanitize_for_markdown(stream_state.title) safe_title = sanitize_for_markdown(stream_state.title)
body_text += f" - Title: {safe_title} \n" parts.append(f" - Title: {safe_title} \n")
# Add status with duration # Add status with duration
if stream_state.last_connect_time: if stream_state.last_connect_time:
duration = self._format_duration(stream_state.last_connect_time) duration = self._format_duration(stream_state.last_connect_time)
body_text += f" - Online for {duration} \n" parts.append(f" - Online for {duration} \n")
# Add stream link # Add stream link
body_text += f" - Link: https://{domain}\n\n" parts.append(f" - Link: https://{domain}\n\n")
# Send the response as Markdown # Send the response as Markdown
await evt.reply(body_text.rstrip(), markdown=True) await evt.reply("".join(parts).rstrip(), markdown=True)
+4 -4
View File
@@ -25,9 +25,9 @@ class Config(BaseProxyConfig):
:param helper: ConfigUpdateHelper for copying values. :param helper: ConfigUpdateHelper for copying values.
""" """
helper.copy("health_check_endpoint") helper.copy("metrics_enabled")
@property @property
def health_check_endpoint(self) -> str: def metrics_enabled(self) -> bool:
"""Return the configured health check endpoint URL.""" """Return whether the Prometheus metrics endpoint is enabled."""
return self["health_check_endpoint"] # type: ignore[no-any-return] return self["metrics_enabled"] # type: ignore[no-any-return]
-163
View File
@@ -1,163 +0,0 @@
# Copyright 2026 Logan Fick
#
# Licensed under the Apache License, Version 2.0 (the "License");
# you may not use this file except in compliance with the License.
# You may obtain a copy of the License at
#
# http://www.apache.org/licenses/LICENSE-2.0
#
# Unless required by applicable law or agreed to in writing, software
# distributed under the License is distributed on an "AS IS" BASIS,
# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
# See the License for the specific language governing permissions and
# limitations under the License.
"""Health checking service for OwncastSentry."""
from dataclasses import dataclass
from typing import TYPE_CHECKING
if TYPE_CHECKING:
import logging
from mautrix.util.async_db import Database
from .owncast_client import OwncastClient
@dataclass
class UpdateResult:
"""Result of a stream update cycle."""
total_streams: int
successful_checks: int
failed_checks: int
@property
def http_healthy(self) -> bool:
"""Determine HTTP health based on update results.
HTTP is considered healthy if:
- No streams are subscribed (nothing to check), OR
- At least one stream check succeeded
:return: True if HTTP is considered healthy.
"""
if self.total_streams == 0:
return True
return self.successful_checks > 0
@dataclass
class HealthStatus:
"""Represents the health status of the plugin."""
database_healthy: bool
http_healthy: bool
@property
def is_healthy(self) -> bool:
"""Check if all health components are healthy.
:return: True if all checks pass.
"""
return self.database_healthy and self.http_healthy
class HealthChecker:
"""Service for performing health checks on the plugin."""
def __init__(
self,
database: Database,
owncast_client: OwncastClient,
logger: logging.Logger,
):
"""Initialize the health checker.
:param database: The maubot database instance.
:param owncast_client: Client for making HTTP requests.
:param logger: Logger instance for debugging.
"""
self.db = database
self.owncast_client = owncast_client
self.log = logger
async def check_database(self) -> bool:
"""Check if the database is functioning by executing a simple query.
:return: True if database is healthy, False otherwise.
"""
try:
async with self.db.acquire() as conn: # type: ignore[var-annotated]
await conn.fetchval("SELECT 1")
return True
except Exception as e:
self.log.warning(f"Database health check failed: {e}")
return False
async def perform_health_check(
self,
update_result: UpdateResult,
endpoint: str,
) -> None:
"""Perform health check and report to configured endpoint if all healthy.
:param update_result: Result of the stream update cycle.
:param endpoint: Health check endpoint URL (empty string to skip reporting).
"""
# Check database health
database_healthy = await self.check_database()
# Evaluate HTTP health from update results
http_healthy = update_result.http_healthy
# Create health status
status = HealthStatus(
database_healthy=database_healthy,
http_healthy=http_healthy,
)
self.log.debug(
f"Health check: database={database_healthy}, http={http_healthy}, "
f"streams={update_result.total_streams}, "
f"succeeded={update_result.successful_checks}, "
f"failed={update_result.failed_checks}"
)
# Skip endpoint notification if not configured
if not endpoint or not endpoint.strip():
self.log.debug("Health check endpoint not configured, skipping report.")
return
# Only send to endpoint if ALL checks pass
if not status.is_healthy:
self.log.warning(
f"Health check failed, not reporting to endpoint. "
f"database={database_healthy}, http={http_healthy}"
)
return
# Send GET request to health endpoint
await self._send_health_report(endpoint)
async def _send_health_report(self, endpoint: str) -> None:
"""Send a GET request to the health check endpoint.
:param endpoint: The endpoint URL.
"""
try:
async with self.owncast_client.session.get(
endpoint, allow_redirects=True
) as response:
if 200 <= response.status < 300:
self.log.debug(
f"Health check reported successfully (status={response.status})"
)
else:
self.log.warning(
"Health check endpoint returned "
f"non-success status: {response.status}"
)
except Exception as e:
self.log.warning(f"Failed to report health check to endpoint: {e}")
+234
View File
@@ -0,0 +1,234 @@
# Copyright 2026 Logan Fick
#
# Licensed under the Apache License, Version 2.0 (the "License");
# you may not use this file except in compliance with the License.
# You may obtain a copy of the License at
#
# http://www.apache.org/licenses/LICENSE-2.0
#
# Unless required by applicable law or agreed to in writing, software
# distributed under the License is distributed on an "AS IS" BASIS,
# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
# See the License for the specific language governing permissions and
# limitations under the License.
"""Prometheus metrics service for OwncastSentry."""
import time
from contextlib import contextmanager, suppress
from enum import StrEnum
from typing import TYPE_CHECKING
from prometheus_client import CollectorRegistry, Counter, Gauge, Info
from .models import StreamStatus
class NotificationType(StrEnum):
"""Notification type labels for the delivery counter."""
LIVE = "live"
TITLE_CHANGE = "title_change"
CLEANUP_WARNING = "cleanup_warning"
CLEANUP_DELETION = "cleanup_deletion"
class ErrorSource(StrEnum):
"""Error source labels for the error counter."""
SCHEDULER_LOOP = "scheduler_loop"
COMMAND = "command"
if TYPE_CHECKING:
from collections.abc import Callable, Generator
# Mapping from StreamStatus enum to numeric gauge values
_STATUS_VALUES: dict[StreamStatus, float] = {
StreamStatus.ONLINE: 1.0,
StreamStatus.OFFLINE: 0.0,
StreamStatus.UNKNOWN: -1.0,
}
class _ResponseTimer:
"""Timer that only records an observation when explicitly marked successful."""
__slots__ = ("_domain", "_gauge", "_should_observe", "_start")
def __init__(self, gauge: Gauge, domain: str) -> None:
self._gauge = gauge
self._domain = domain
self._start = time.monotonic()
self._should_observe = False
def success(self) -> None:
"""Mark the request as successful so the duration is recorded."""
self._should_observe = True
def _finalize(self) -> None:
if self._should_observe:
self._gauge.labels(domain=self._domain).set(
max(time.monotonic() - self._start, 0)
)
else:
with suppress(KeyError):
self._gauge.remove(self._domain)
class MetricsService:
"""Manages Prometheus metrics with an isolated registry."""
def __init__(self) -> None:
"""Initialize metrics with a custom collector registry."""
self.registry = CollectorRegistry()
self.notification_delivery_total = Counter(
"owncastsentry_notification_delivery_total",
"Total notification delivery attempts to individual rooms",
["type", "result"],
registry=self.registry,
)
self.stream_status = Gauge(
"owncastsentry_stream_status",
"Current stream status (1=online, 0=offline, -1=unknown)",
["domain"],
registry=self.registry,
)
self.stream_subscriptions = Gauge(
"owncastsentry_stream_subscriptions",
"Number of room subscriptions per stream domain",
["domain"],
registry=self.registry,
)
self.check_failures = Gauge(
"owncastsentry_check_failures",
"Consecutive check failure count per stream domain",
["domain"],
registry=self.registry,
)
self.api_response_seconds = Gauge(
"owncastsentry_api_response_seconds",
"Last successful HTTP response time in seconds per stream domain",
["domain"],
registry=self.registry,
)
self.build_info = Info(
"owncastsentry",
"OwncastSentry build information",
registry=self.registry,
)
self.open_connections = Gauge(
"owncastsentry_http_connections_open",
"Total number of open HTTP connections (idle and active)",
registry=self.registry,
)
self.errors_total = Counter(
"owncastsentry_errors_total",
"Internal errors by source",
["source"],
registry=self.registry,
)
# Initialize all known label combinations so they start at 0
for notification_type in NotificationType:
for result in ("success", "failure"):
self.notification_delivery_total.labels(
type=notification_type, result=result
)
for source in ErrorSource:
self.errors_total.labels(source=source)
def record_delivery(
self,
notification_type: NotificationType,
*,
successful: int = 0,
failed: int = 0,
) -> None:
"""Record notification delivery results.
:param notification_type: The type of notification delivered.
:param successful: Number of successful room deliveries.
:param failed: Number of failed room deliveries.
"""
self.notification_delivery_total.labels(
type=notification_type, result="success"
).inc(successful)
self.notification_delivery_total.labels(
type=notification_type, result="failure"
).inc(failed)
def set_stream_status(self, domain: str, status: StreamStatus) -> None:
"""Set the status gauge for a stream domain.
:param domain: The stream domain.
:param status: The current stream status.
"""
self.stream_status.labels(domain=domain).set(_STATUS_VALUES[status])
def set_check_failures(self, domain: str, count: int) -> None:
"""Set the consecutive failure count for a stream domain.
:param domain: The stream domain.
:param count: The current failure counter value.
"""
self.check_failures.labels(domain=domain).set(count)
def set_subscription_count(self, domain: str, count: int) -> None:
"""Set the subscription count for a stream domain.
:param domain: The stream domain.
:param count: The number of room subscriptions.
"""
self.stream_subscriptions.labels(domain=domain).set(count)
@contextmanager
def response_timer(self, domain: str) -> Generator[_ResponseTimer]:
"""Return a context manager that times an HTTP request.
Call ``timer.success()`` inside the block to record the duration.
If ``success()`` is never called, nothing is recorded.
:param domain: The stream domain being queried.
"""
timer = _ResponseTimer(self.api_response_seconds, domain)
try:
yield timer
finally:
timer._finalize()
def set_build_info(self, version: str) -> None:
"""Set the build version info metric.
:param version: The plugin version string.
"""
self.build_info.info({"version": version})
def register_open_connections_gauge(self, callback: Callable[[], float]) -> None:
"""Register a gauge that reads open connection count on scrape.
:param callback: Function returning the current open count.
"""
self.open_connections.set_function(callback)
def record_error(self, source: ErrorSource) -> None:
"""Increment the internal error counter.
:param source: The source of the error.
"""
self.errors_total.labels(source=source).inc()
def remove_stream(self, domain: str) -> None:
"""Remove a stream's gauge labels after cleanup deletion.
:param domain: The stream domain to remove.
"""
with suppress(KeyError):
self.stream_status.remove(domain)
with suppress(KeyError):
self.check_failures.remove(domain)
with suppress(KeyError):
self.stream_subscriptions.remove(domain)
with suppress(KeyError):
self.api_response_seconds.remove(domain)
+9
View File
@@ -91,6 +91,15 @@ class StreamState:
) )
@dataclass
class UpdateResult:
"""Result of a stream update cycle."""
total_streams: int
successful_checks: int
failed_checks: int
@dataclass @dataclass
class StreamConfig: class StreamConfig:
"""Represents the configuration of an Owncast stream.""" """Represents the configuration of an Owncast stream."""
+65 -77
View File
@@ -20,9 +20,10 @@ from typing import TYPE_CHECKING, Any
from mautrix.types import MessageType, TextMessageEventContent from mautrix.types import MessageType, TextMessageEventContent
from .metrics import NotificationType
from .utils import ( from .utils import (
CLEANUP_DELETE_THRESHOLD, CLEANUP_DELETE_DAYS,
CLEANUP_WARNING_THRESHOLD, CLEANUP_WARNING_DAYS,
SECONDS_BETWEEN_NOTIFICATIONS, SECONDS_BETWEEN_NOTIFICATIONS,
sanitize_for_plain_text, sanitize_for_plain_text,
) )
@@ -31,6 +32,7 @@ if TYPE_CHECKING:
import logging import logging
from .database import SubscriptionRepository from .database import SubscriptionRepository
from .metrics import MetricsService
class NotificationService: class NotificationService:
@@ -41,16 +43,19 @@ class NotificationService:
client: Any, client: Any,
subscription_repo: SubscriptionRepository, subscription_repo: SubscriptionRepository,
logger: logging.Logger, logger: logging.Logger,
metrics: MetricsService,
) -> None: ) -> None:
"""Initialize the notification service. """Initialize the notification service.
:param client: The Matrix client for sending messages. :param client: The Matrix client for sending messages.
:param subscription_repo: Repository for managing subscriptions. :param subscription_repo: Repository for managing subscriptions.
:param logger: Logger instance for debugging. :param logger: Logger instance for debugging.
:param metrics: Metrics service for recording counters.
""" """
self.client = client self.client = client
self.subscription_repo = subscription_repo self.subscription_repo = subscription_repo
self.log = logger self.log = logger
self.metrics = metrics
# Cache for tracking when notifications were last sent # Cache for tracking when notifications were last sent
self.notification_timers_cache: dict[str, float] = {} self.notification_timers_cache: dict[str, float] = {}
@@ -88,36 +93,24 @@ class NotificationService:
# Record that we're sending a notification now # Record that we're sending a notification now
self._record_notification(domain) self._record_notification(domain)
# Get a list of room IDs with active subscriptions to the stream domain
room_ids = await self.subscription_repo.get_subscribed_rooms(domain)
# Build the notification message # Build the notification message
body_text = self._format_message(name, title, domain, tags, title_change) body_text = self._format_message(name, title, domain, tags, title_change)
# Set up counters for statistics
successful_notifications = 0
failed_notifications = 0
# Send notifications to all subscribed rooms in parallel # Send notifications to all subscribed rooms in parallel
tasks = [ successful, failed = await self._broadcast_to_rooms(domain, body_text)
self._send_notification(room_id, body_text, domain) for room_id in room_ids
]
results = await asyncio.gather(*tasks, return_exceptions=True)
# Count successes and failures
for result in results:
if isinstance(result, Exception):
failed_notifications += 1
else:
successful_notifications += 1
# Log completion # Log completion
notification_type = "title change" if title_change else "going live" notification_type = "title change" if title_change else "going live"
self.log.info( self.log.info(
f"[{domain}] Completed sending {notification_type} " f"[{domain}] Completed sending {notification_type} "
f"notifications! {successful_notifications} succeeded, " f"notifications! {successful} succeeded, "
f"{failed_notifications} failed." f"{failed} failed."
)
self.metrics.record_delivery(
NotificationType.TITLE_CHANGE if title_change else NotificationType.LIVE,
successful=successful,
failed=failed,
) )
async def _send_notification( async def _send_notification(
@@ -153,36 +146,36 @@ class NotificationService:
:return: Formatted message body. :return: Formatted message body.
""" """
# Use name if available, fallback to domain # Use name if available, fallback to domain
stream_name = name if name else domain stream_name = name or domain
safe_stream_name = sanitize_for_plain_text(stream_name) safe_stream_name = sanitize_for_plain_text(stream_name)
# Choose message based on notification type # Choose message based on notification type
if title_change: if title_change:
body_text = "📝 " + safe_stream_name + " has changed its stream title!" parts = [f"📝 {safe_stream_name} has changed its stream title!"]
else: else:
body_text = "🎥 " + safe_stream_name + " is now live!" parts = [f"🎥 {safe_stream_name} is now live!"]
# Add title if present # Add title if present
if title != "": if title:
safe_title = sanitize_for_plain_text(title) safe_title = sanitize_for_plain_text(title)
body_text += "\nStream Title: " + safe_title parts.append(f"\nStream Title: {safe_title}")
# Add stream URL # Add stream URL
body_text += "\n\nTo tune in, visit: https://" + domain + "/" parts.append(f"\n\nTo tune in, visit: https://{domain}/")
# Add tags if present # Add tags if present
if tags: if tags:
safe_tags = [] safe_tags = [
for tag in tags: safe_tag
safe_tag = sanitize_for_plain_text(tag) for tag in tags
if safe_tag and not safe_tag.startswith("."): if (safe_tag := sanitize_for_plain_text(tag))
safe_tags.append(safe_tag) and not safe_tag.startswith(".")
]
if safe_tags: if safe_tags:
body_text += "\n\n" parts.append(f"\n\n{' '.join(f'#{tag}' for tag in safe_tags)}")
body_text += " ".join("#" + tag for tag in safe_tags)
return body_text return "".join(parts)
def get_last_notification_time(self, domain: str) -> float: def get_last_notification_time(self, domain: str) -> float:
"""Get the timestamp of the last notification sent for a domain. """Get the timestamp of the last notification sent for a domain.
@@ -201,84 +194,79 @@ class NotificationService:
if domain not in self.notification_timers_cache: if domain not in self.notification_timers_cache:
return True return True
seconds_since_last = round(time.time() - self.notification_timers_cache[domain]) seconds_since_last = round(
return bool(seconds_since_last >= SECONDS_BETWEEN_NOTIFICATIONS) time.monotonic() - self.notification_timers_cache[domain]
)
return seconds_since_last >= SECONDS_BETWEEN_NOTIFICATIONS
def _record_notification(self, domain: str) -> None: def _record_notification(self, domain: str) -> None:
"""Record that a notification was sent at the current time. """Record that a notification was sent at the current time.
:param domain: The stream domain. :param domain: The stream domain.
""" """
self.notification_timers_cache[domain] = time.time() self.notification_timers_cache[domain] = time.monotonic()
async def _broadcast_to_rooms(self, domain: str, body_text: str) -> tuple[int, int]:
"""Send a message to all rooms subscribed to a domain.
:param domain: The stream domain.
:param body_text: The message body text.
:return: Tuple of (successful, failed) counts.
"""
room_ids = await self.subscription_repo.get_subscribed_rooms(domain)
tasks = [
self._send_notification(room_id, body_text, domain) for room_id in room_ids
]
results = await asyncio.gather(*tasks, return_exceptions=True)
failed = sum(1 for r in results if isinstance(r, Exception))
successful = len(results) - failed
return successful, failed
async def send_cleanup_warning(self, domain: str) -> None: async def send_cleanup_warning(self, domain: str) -> None:
"""Send cleanup warning notification to all subscribed rooms. """Send cleanup warning notification to all subscribed rooms.
:param domain: The stream domain. :param domain: The stream domain.
""" """
# Get all subscribed rooms remaining_days = CLEANUP_DELETE_DAYS - CLEANUP_WARNING_DAYS
room_ids = await self.subscription_repo.get_subscribed_rooms(domain)
# Build the warning message
warning_days = CLEANUP_WARNING_THRESHOLD // (24 * 60)
delete_days = CLEANUP_DELETE_THRESHOLD // (24 * 60)
remaining_days = delete_days - warning_days
body_text = ( body_text = (
"⚠️ Warning: Subscription Cleanup Scheduled\n\n" "⚠️ Warning: Subscription Cleanup Scheduled\n\n"
f"The Owncast instance at {domain} has been " f"The Owncast instance at {domain} has been "
f"unreachable for {warning_days} days. If it remains " f"unreachable for {CLEANUP_WARNING_DAYS} days. If it remains "
f"unreachable for {remaining_days} more days " f"unreachable for {remaining_days} more days "
f"({delete_days} days total), this subscription " f"({CLEANUP_DELETE_DAYS} days total), this subscription "
f"will be automatically removed." f"will be automatically removed."
) )
# Send to all rooms in parallel successful, failed = await self._broadcast_to_rooms(domain, body_text)
tasks = [
self._send_notification(room_id, body_text, domain) for room_id in room_ids
]
results = await asyncio.gather(*tasks, return_exceptions=True)
# Count successes and failures
successful = sum(1 for r in results if not isinstance(r, Exception))
failed = sum(1 for r in results if isinstance(r, Exception))
self.log.info( self.log.info(
f"[{domain}] Sent cleanup warning to {successful} rooms ({failed} failed)." f"[{domain}] Sent cleanup warning to {successful} rooms ({failed} failed)."
) )
self.metrics.record_delivery(
NotificationType.CLEANUP_WARNING, successful=successful, failed=failed
)
async def send_cleanup_deletion(self, domain: str) -> None: async def send_cleanup_deletion(self, domain: str) -> None:
"""Send cleanup deletion notification to all subscribed rooms. """Send cleanup deletion notification to all subscribed rooms.
:param domain: The stream domain. :param domain: The stream domain.
""" """
# Get all subscribed rooms
room_ids = await self.subscription_repo.get_subscribed_rooms(domain)
# Build the deletion message
delete_days = CLEANUP_DELETE_THRESHOLD // (24 * 60)
body_text = ( body_text = (
"🗑️ Subscription Automatically Removed\n\n" "🗑️ Subscription Automatically Removed\n\n"
f"The Owncast instance at {domain} has been " f"The Owncast instance at {domain} has been "
f"unreachable for {delete_days} days and has been " f"unreachable for {CLEANUP_DELETE_DAYS} days and has been "
f"automatically removed from subscriptions in this " f"automatically removed from subscriptions in this "
f"room.\n\n" f"room.\n\n"
f"If the instance comes online again and you want to " f"If the instance comes online again and you want to "
f"resubscribe, run `!subscribe {domain}`." f"resubscribe, run `!subscribe {domain}`."
) )
# Send to all rooms in parallel successful, failed = await self._broadcast_to_rooms(domain, body_text)
tasks = [
self._send_notification(room_id, body_text, domain) for room_id in room_ids
]
results = await asyncio.gather(*tasks, return_exceptions=True)
# Count successes and failures
successful = sum(1 for r in results if not isinstance(r, Exception))
failed = sum(1 for r in results if isinstance(r, Exception))
self.log.info( self.log.info(
f"[{domain}] Sent cleanup deletion notice to " f"[{domain}] Sent cleanup deletion notice to "
f"{successful} rooms ({failed} failed)." f"{successful} rooms ({failed} failed)."
) )
self.metrics.record_delivery(
NotificationType.CLEANUP_DELETION, successful=successful, failed=failed
)
+72 -76
View File
@@ -14,28 +14,41 @@
"""HTTP client for querying Owncast instance APIs.""" """HTTP client for querying Owncast instance APIs."""
import json from typing import TYPE_CHECKING, Any
from typing import TYPE_CHECKING
import aiohttp import aiohttp
from .models import StreamConfig, StreamState
from .utils import (
OWNCAST_CONFIG_PATH,
OWNCAST_STATUS_PATH,
REQUIRED_STATUS_FIELDS,
user_agent,
)
if TYPE_CHECKING: if TYPE_CHECKING:
import logging import logging
from .models import StreamConfig, StreamState from .metrics import MetricsService
from .utils import OWNCAST_CONFIG_PATH, OWNCAST_STATUS_PATH, user_agent
class OwncastClient: class OwncastClient:
"""HTTP client for communicating with Owncast instances.""" """HTTP client for communicating with Owncast instances."""
def __init__(self, logger: logging.Logger, version: str) -> None: def __init__(
self,
logger: logging.Logger,
version: str,
metrics: MetricsService,
) -> None:
"""Initialize the Owncast client with an HTTP session. """Initialize the Owncast client with an HTTP session.
:param logger: Logger instance for debugging :param logger: Logger instance for debugging
:param version: Plugin version string for the User-Agent header :param version: Plugin version string for the User-Agent header
:param metrics: Metrics service for recording response times.
""" """
self.log = logger self.log = logger
self.metrics = metrics
# Set up HTTP session configuration # Set up HTTP session configuration
headers = {"User-Agent": user_agent(version)} headers = {"User-Agent": user_agent(version)}
@@ -55,6 +68,37 @@ class OwncastClient:
connector=connector, connector=connector,
) )
async def _fetch_json(self, domain: str, path: str) -> dict[str, Any] | None:
"""Fetch and parse JSON from an Owncast API endpoint.
:param domain: The domain to query.
:param path: The API path to request.
:return: Parsed JSON response, or None on error.
"""
url = f"https://{domain}{path}"
try:
async with self.session.get(url, allow_redirects=False) as response:
if response.status != 200:
self.log.warning(
f"[{domain}] Response to request on "
f"{path} was not 200, "
f"got {response.status} instead."
)
return None
try:
result: dict[str, Any] = await response.json()
return result
except (ValueError, aiohttp.ContentTypeError) as e:
self.log.warning(
f"[{domain}] Rejecting response to request on "
f"{path} as could not be "
f"interpreted as JSON: {e}"
)
return None
except (aiohttp.ClientError, TimeoutError, OSError) as e:
self.log.warning(f"[{domain}] Error making GET request to {path}: {e}")
return None
async def get_stream_state(self, domain: str) -> StreamState | None: async def get_stream_state(self, domain: str) -> StreamState | None:
"""Get the current stream state for a given domain. """Get the current stream state for a given domain.
@@ -65,54 +109,23 @@ class OwncastClient:
:return: A StreamState if available, None on error. :return: A StreamState if available, None on error.
""" """
self.log.debug(f"[{domain}] Fetching current stream state...") self.log.debug(f"[{domain}] Fetching current stream state...")
status_url = "https://" + domain + OWNCAST_STATUS_PATH with self.metrics.response_timer(domain) as timer:
new_state = await self._fetch_json(domain, OWNCAST_STATUS_PATH)
# Make a request to the endpoint if new_state is None:
try:
async with self.session.request(
"GET", status_url, allow_redirects=False
) as response:
# Check the response code is success
if response.status != 200:
self.log.warning(
f"[{domain}] Response to request on "
f"{OWNCAST_STATUS_PATH} was not 200, "
f"got {response.status} instead."
)
return None
# Try to interpret the response as JSON
try:
new_state = json.loads(await response.read())
except Exception as e:
self.log.warning(
f"[{domain}] Rejecting response to request on "
f"{OWNCAST_STATUS_PATH} as could not be "
f"interpreted as JSON: {e}"
)
return None
except Exception as e:
self.log.warning(
f"[{domain}] Error making GET request to {OWNCAST_STATUS_PATH}: {e}"
)
return None return None
# Validate the response contains all basic info needed # Validate the response contains all basic info needed
required_fields = [ missing = REQUIRED_STATUS_FIELDS - new_state.keys()
"lastConnectTime", if missing:
"lastDisconnectTime",
"streamTitle",
"online",
]
for field in required_fields:
if field not in new_state:
self.log.warning( self.log.warning(
f"[{domain}] Rejecting response to request " f"[{domain}] Rejecting response to request on "
f"on {OWNCAST_STATUS_PATH} as it does not " f"{OWNCAST_STATUS_PATH} as it is missing "
f"have {field} field." f"fields: {', '.join(sorted(missing))}"
) )
return None return None
timer.success()
return StreamState.from_api_response(new_state, domain) return StreamState.from_api_response(new_state, domain)
async def get_stream_config(self, domain: str) -> StreamConfig | None: async def get_stream_config(self, domain: str) -> StreamConfig | None:
@@ -125,39 +138,12 @@ class OwncastClient:
:return: A StreamConfig, or None if fetch failed. :return: A StreamConfig, or None if fetch failed.
""" """
self.log.debug(f"[{domain}] Fetching current stream config...") self.log.debug(f"[{domain}] Fetching current stream config...")
config_url = "https://" + domain + OWNCAST_CONFIG_PATH with self.metrics.response_timer(domain) as timer:
config = await self._fetch_json(domain, OWNCAST_CONFIG_PATH)
# Make a request to the endpoint if config is None:
try:
async with self.session.request(
"GET", config_url, allow_redirects=False
) as response:
# Check the response code is success
if response.status != 200:
self.log.warning(
f"[{domain}] Response to request on "
f"{OWNCAST_CONFIG_PATH} was not 200, "
f"got {response.status} instead."
)
return None return None
# Try to interpret the response as JSON timer.success()
try:
config = json.loads(await response.read())
except Exception as e:
self.log.warning(
f"[{domain}] Rejecting response to request on "
f"{OWNCAST_CONFIG_PATH} as could not be "
f"interpreted as JSON: {e}"
)
return None
except Exception as e:
self.log.warning(
f"[{domain}] Error making GET request to {OWNCAST_CONFIG_PATH}: {e}"
)
return None
# Create StreamConfig from response (fields are truncated to max lengths)
return StreamConfig.from_api_response(config) return StreamConfig.from_api_response(config)
async def validate_instance(self, domain: str) -> bool: async def validate_instance(self, domain: str) -> bool:
@@ -169,6 +155,16 @@ class OwncastClient:
state = await self.get_stream_state(domain) state = await self.get_stream_state(domain)
return state is not None return state is not None
@property
def open_connection_count(self) -> int:
"""Return the total number of open HTTP connections."""
connector = self.session.connector
if connector is None:
return 0
idle = sum(len(conns) for conns in connector._conns.values())
active = len(connector._acquired)
return idle + active
async def close(self) -> None: async def close(self) -> None:
"""Close the HTTP session.""" """Close the HTTP session."""
await self.session.close() await self.session.close()
+32 -15
View File
@@ -18,8 +18,7 @@ import asyncio
import time import time
from typing import TYPE_CHECKING from typing import TYPE_CHECKING
from .health_checker import UpdateResult from .models import StreamState, StreamStatus, UpdateResult
from .models import StreamState
from .utils import ( from .utils import (
CLEANUP_DELETE_THRESHOLD, CLEANUP_DELETE_THRESHOLD,
CLEANUP_WARNING_THRESHOLD, CLEANUP_WARNING_THRESHOLD,
@@ -31,6 +30,7 @@ if TYPE_CHECKING:
import logging import logging
from .database import StreamRepository, SubscriptionRepository from .database import StreamRepository, SubscriptionRepository
from .metrics import MetricsService
from .notification_service import NotificationService from .notification_service import NotificationService
from .owncast_client import OwncastClient from .owncast_client import OwncastClient
@@ -45,6 +45,7 @@ class StreamMonitor:
subscription_repo: SubscriptionRepository, subscription_repo: SubscriptionRepository,
notification_service: NotificationService, notification_service: NotificationService,
logger: logging.Logger, logger: logging.Logger,
metrics: MetricsService,
) -> None: ) -> None:
"""Initialize the stream monitor. """Initialize the stream monitor.
@@ -53,12 +54,14 @@ class StreamMonitor:
:param subscription_repo: Repository for subscription data. :param subscription_repo: Repository for subscription data.
:param notification_service: Service for sending notifications. :param notification_service: Service for sending notifications.
:param logger: Logger instance for debugging. :param logger: Logger instance for debugging.
:param metrics: Metrics service for recording Prometheus metrics.
""" """
self.owncast_client = owncast_client self.owncast_client = owncast_client
self.stream_repo = stream_repo self.stream_repo = stream_repo
self.subscription_repo = subscription_repo self.subscription_repo = subscription_repo
self.notification_service = notification_service self.notification_service = notification_service
self.log = logger self.log = logger
self.metrics = metrics
# Cache for tracking when streams last went offline # Cache for tracking when streams last went offline
self.offline_timer_cache: dict[str, float] = {} self.offline_timer_cache: dict[str, float] = {}
@@ -77,23 +80,24 @@ class StreamMonitor:
total_streams = len(subscribed_domains) total_streams = len(subscribed_domains)
# Build a list of async tasks for each stream domain # Build a list of async tasks for each stream domain
tasks = [ # Run all stream updates in parallel and collect results
asyncio.create_task(self.update_stream(domain)) results = await asyncio.gather(
for domain in subscribed_domains *(self.update_stream(domain) for domain in subscribed_domains)
] )
# Run the tasks in parallel and collect results
results = await asyncio.gather(*tasks)
# Count successes and failures # Count successes and failures
successful_checks = sum(1 for result in results if result is True) successful_checks = results.count(True)
failed_checks = sum(1 for result in results if result is False) failed_checks = results.count(False)
self.log.debug( self.log.debug(
f"Update complete. {successful_checks}/{total_streams} succeeded, " f"Update complete. {successful_checks}/{total_streams} succeeded, "
f"{failed_checks} failed." f"{failed_checks} failed."
) )
for domain in subscribed_domains:
count = await self.subscription_repo.count_by_domain(domain)
self.metrics.set_subscription_count(domain, count)
return UpdateResult( return UpdateResult(
total_streams=total_streams, total_streams=total_streams,
successful_checks=successful_checks, successful_checks=successful_checks,
@@ -127,6 +131,10 @@ class StreamMonitor:
) )
# Check cleanup thresholds even when skipping query # Check cleanup thresholds even when skipping query
await self._check_cleanup_thresholds(domain, failure_counter + 1) await self._check_cleanup_thresholds(domain, failure_counter + 1)
updated_state = await self.stream_repo.get_by_domain(domain)
if updated_state is not None:
self.metrics.set_stream_status(domain, updated_state.status)
self.metrics.set_check_failures(domain, failure_counter + 1)
# Backoff is expected behavior, not a failure # Backoff is expected behavior, not a failure
return True return True
@@ -152,15 +160,19 @@ class StreamMonitor:
) )
# Check cleanup thresholds after connection failure # Check cleanup thresholds after connection failure
await self._check_cleanup_thresholds(domain, failure_counter + 1) await self._check_cleanup_thresholds(domain, failure_counter + 1)
updated_state = await self.stream_repo.get_by_domain(domain)
if updated_state is not None:
self.metrics.set_stream_status(domain, updated_state.status)
self.metrics.set_check_failures(domain, failure_counter + 1)
# Actual connection failure # Actual connection failure
return False return False
# Fetch succeeded! Reset failure counter # Fetch succeeded! Reset failure counter
await self.stream_repo.reset_failure_counter(domain) await self.stream_repo.reset_failure_counter(domain)
self.metrics.set_check_failures(domain, 0)
# Initialize timer cache entries to prevent KeyError on first access # Initialize timer cache entries to prevent KeyError on first access
if domain not in self.offline_timer_cache: self.offline_timer_cache.setdefault(domain, 0)
self.offline_timer_cache[domain] = 0
# Does the last known stream state lack connect/disconnect? # Does the last known stream state lack connect/disconnect?
if ( if (
@@ -184,7 +196,7 @@ class StreamMonitor:
# Calculate seconds since the stream last went offline # Calculate seconds since the stream last went offline
seconds_since_last_offline = round( seconds_since_last_offline = round(
time.time() - self.offline_timer_cache[domain] time.monotonic() - self.offline_timer_cache[domain]
) )
# Have we queried this stream before? # Have we queried this stream before?
@@ -278,7 +290,7 @@ class StreamMonitor:
): ):
# Yep. This stream is now offline. Log it. # Yep. This stream is now offline. Log it.
update_database = True update_database = True
self.offline_timer_cache[domain] = time.time() self.offline_timer_cache[domain] = time.monotonic()
self.log.info(f"[{domain}] Stream is now offline.") self.log.info(f"[{domain}] Stream is now offline.")
# Update the database with current stream state, if needed. # Update the database with current stream state, if needed.
@@ -305,6 +317,10 @@ class StreamMonitor:
# All done. # All done.
self.log.debug(f"[{domain}] State update completed.") self.log.debug(f"[{domain}] State update completed.")
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)
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:
@@ -340,3 +356,4 @@ class StreamMonitor:
f"Deleted {deleted_count} subscriptions " f"Deleted {deleted_count} subscriptions "
f"and stream record." f"and stream record."
) )
self.metrics.remove_stream(domain)
+23 -44
View File
@@ -23,6 +23,16 @@ OWNCAST_STATUS_PATH = "/api/status"
# Path to GetWebConfig API call on Owncast instances # Path to GetWebConfig API call on Owncast instances
OWNCAST_CONFIG_PATH = "/api/config" OWNCAST_CONFIG_PATH = "/api/config"
# Fields that must be present in an Owncast status API response
REQUIRED_STATUS_FIELDS = frozenset(
{
"lastConnectTime",
"lastDisconnectTime",
"streamTitle",
"online",
}
)
def user_agent(version: str) -> str: def user_agent(version: str) -> str:
"""Build the User-Agent header string for HTTP requests. """Build the User-Agent header string for HTTP requests.
@@ -50,9 +60,13 @@ SECONDS_BETWEEN_NOTIFICATIONS = 20 * 60 # 20 minutes in seconds
# online after, it's treated as regular going live. # online after, it's treated as regular going live.
TEMPORARY_OFFLINE_NOTIFICATION_COOLDOWN = 7 * 60 # 7 min in seconds TEMPORARY_OFFLINE_NOTIFICATION_COOLDOWN = 7 * 60 # 7 min in seconds
# Counter thresholds for auto-cleanup (60-second polling intervals) # Auto-cleanup timing (days of continuous unreachability)
CLEANUP_WARNING_THRESHOLD = 83 * 24 * 60 # 119,520 cycles = 83 days CLEANUP_WARNING_DAYS = 83
CLEANUP_DELETE_THRESHOLD = 90 * 24 * 60 # 129,600 cycles = 90 days CLEANUP_DELETE_DAYS = 90
# Counter thresholds derived from days (60-second polling intervals)
CLEANUP_WARNING_THRESHOLD = CLEANUP_WARNING_DAYS * 24 * 60
CLEANUP_DELETE_THRESHOLD = CLEANUP_DELETE_DAYS * 24 * 60
# Failure counter threshold for treating stream status as "unknown" # Failure counter threshold for treating stream status as "unknown"
UNKNOWN_STATUS_THRESHOLD = 15 UNKNOWN_STATUS_THRESHOLD = 15
@@ -109,7 +123,7 @@ def domainify(url: str) -> str:
# Prepend // if no scheme so urlparse treats input as netloc # Prepend // if no scheme so urlparse treats input as netloc
if not url.startswith(("http://", "https://", "//")): if not url.startswith(("http://", "https://", "//")):
url = "//" + url url = f"//{url}"
parsed = urlparse(url) parsed = urlparse(url)
domain = (parsed.netloc or parsed.path).lower() domain = (parsed.netloc or parsed.path).lower()
@@ -133,6 +147,9 @@ def truncate(text: str, max_length: int) -> str:
return text[:max_length] return text[:max_length]
_MARKDOWN_ESCAPE_TABLE = str.maketrans({c: f"\\{c}" for c in r"\*_[]()~`#+-=|{}.!<>&"})
def escape_markdown(text: str) -> str: def escape_markdown(text: str) -> str:
"""Escape Markdown special characters to prevent injection attacks. """Escape Markdown special characters to prevent injection attacks.
@@ -146,38 +163,7 @@ def escape_markdown(text: str) -> str:
if not text: if not text:
return text return text
# Escape Markdown special characters by prefixing with backslash return text.translate(_MARKDOWN_ESCAPE_TABLE)
# Covers: formatting (*_~`), links ([]()), headings (#), lists (-+),
# blockquotes (>), code blocks (```), and other special characters
special_chars = {
"\\": "\\\\", # Backslash must be first to avoid double-escaping
"*": "\\*",
"_": "\\_",
"[": "\\[",
"]": "\\]",
"(": "\\(",
")": "\\)",
"~": "\\~",
"`": "\\`",
"#": "\\#",
"+": "\\+",
"-": "\\-",
"=": "\\=",
"|": "\\|",
"{": "\\{",
"}": "\\}",
".": "\\.",
"!": "\\!",
"<": "\\<",
">": "\\>",
"&": "\\&",
}
escaped_text = text
for char, replacement in special_chars.items():
escaped_text = escaped_text.replace(char, replacement)
return escaped_text
def sanitize_for_plain_text(text: str) -> str: def sanitize_for_plain_text(text: str) -> str:
@@ -216,11 +202,4 @@ def sanitize_for_markdown(text: str) -> str:
if not text: if not text:
return text return text
# Remove newlines and carriage returns to prevent multi-line injection return escape_markdown(sanitize_for_plain_text(text))
sanitized = text.replace("\n", " ").replace("\r", " ")
# Collapse multiple spaces into single space
sanitized = " ".join(sanitized.split())
# Escape Markdown special characters
return escape_markdown(sanitized)
+3 -3
View File
@@ -7,13 +7,13 @@ authors = [
] ]
license = "Apache-2.0" license = "Apache-2.0"
requires-python = ">=3.14" requires-python = ">=3.14"
dependencies = [] dependencies = ["prometheus_client>=0.24.1"]
[project.urls] [project.urls]
Repository = "https://git.logal.dev/LogalDeveloper/OwncastSentry" Repository = "https://git.logal.dev/LogalDeveloper/OwncastSentry"
[build-system] [build-system]
requires = ["hatchling>=1.28.0", "hatch-vcs>=0.5.0"] requires = ["hatchling>=1.29.0", "hatch-vcs>=0.5.0"]
build-backend = "hatchling.build" build-backend = "hatchling.build"
[dependency-groups] [dependency-groups]
@@ -27,7 +27,7 @@ dev = [
"pytest>=9.0.2", "pytest>=9.0.2",
"pytest-asyncio>=1.3.0", "pytest-asyncio>=1.3.0",
"pytest-cov>=7.0.0", "pytest-cov>=7.0.0",
"ruff>=0.15.5", "ruff>=0.15.6",
"time-machine>=3.2.0", "time-machine>=3.2.0",
] ]
+7 -2
View File
@@ -14,13 +14,12 @@
"""Shared test fixtures and stubs for OwncastSentry tests.""" """Shared test fixtures and stubs for OwncastSentry tests."""
from __future__ import annotations
from dataclasses import dataclass, field from dataclasses import dataclass, field
from typing import TYPE_CHECKING, Any from typing import TYPE_CHECKING, Any
import pytest import pytest
from mautrix.util.async_db import Database from mautrix.util.async_db import Database
from prometheus_client import generate_latest
from owncastsentry import OwncastSentry from owncastsentry import OwncastSentry
from owncastsentry.config import Config from owncastsentry.config import Config
@@ -31,9 +30,15 @@ if TYPE_CHECKING:
from collections.abc import AsyncIterator from collections.abc import AsyncIterator
from pathlib import Path from pathlib import Path
from owncastsentry.metrics import MetricsService
from owncastsentry.models import StreamConfig, StreamState from owncastsentry.models import StreamConfig, StreamState
def generate_metrics_output(metrics: MetricsService) -> str:
"""Generate Prometheus text format output from a MetricsService registry."""
return generate_latest(metrics.registry).decode("utf-8")
@pytest.fixture @pytest.fixture
async def database(tmp_path: Path) -> AsyncIterator[Database]: async def database(tmp_path: Path) -> AsyncIterator[Database]:
"""Yield a real SQLite-backed mautrix Database with migrations applied.""" """Yield a real SQLite-backed mautrix Database with migrations applied."""
-2
View File
@@ -14,8 +14,6 @@
"""Tests for bot command handlers.""" """Tests for bot command handlers."""
from __future__ import annotations
import json import json
import logging import logging
from datetime import UTC, datetime, timedelta from datetime import UTC, datetime, timedelta
-2
View File
@@ -14,8 +14,6 @@
"""Tests for database repository classes.""" """Tests for database repository classes."""
from __future__ import annotations
from typing import TYPE_CHECKING from typing import TYPE_CHECKING
if TYPE_CHECKING: if TYPE_CHECKING:
-171
View File
@@ -1,171 +0,0 @@
# Copyright 2026 Logan Fick
#
# Licensed under the Apache License, Version 2.0 (the "License");
# you may not use this file except in compliance with the License.
# You may obtain a copy of the License at
#
# http://www.apache.org/licenses/LICENSE-2.0
#
# Unless required by applicable law or agreed to in writing, software
# distributed under the License is distributed on an "AS IS" BASIS,
# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
# See the License for the specific language governing permissions and
# limitations under the License.
"""Tests for the health checking service."""
from __future__ import annotations
import logging
from typing import TYPE_CHECKING
import pytest
from aioresponses import aioresponses
from owncastsentry.health_checker import HealthChecker, HealthStatus, UpdateResult
from owncastsentry.owncast_client import OwncastClient
if TYPE_CHECKING:
from collections.abc import AsyncIterator
from mautrix.util.async_db import Database
@pytest.fixture
async def health_checker(database: Database) -> AsyncIterator[HealthChecker]:
"""Yield a HealthChecker backed by a real OwncastClient."""
client = OwncastClient(logger=logging.getLogger("test"), version="0.0.0")
yield HealthChecker(database, client, logging.getLogger("test"))
await client.close()
class TestUpdateResultHttpHealthy:
"""HTTP health derivation from update results."""
@pytest.mark.parametrize(
("total", "successful", "failed", "expected"),
[
pytest.param(0, 0, 0, True, id="no-streams-is-healthy"),
pytest.param(3, 2, 1, True, id="some-successes-is-healthy"),
pytest.param(3, 0, 3, False, id="all-failures-is-unhealthy"),
],
)
def test_http_healthy(
self, total: int, successful: int, failed: int, expected: bool
) -> None:
"""Derive HTTP health from stream check results."""
result = UpdateResult(
total_streams=total,
successful_checks=successful,
failed_checks=failed,
)
assert result.http_healthy is expected
class TestHealthStatusIsHealthy:
"""Overall health status derivation."""
@pytest.mark.parametrize(
("db_healthy", "http_healthy", "expected"),
[
pytest.param(True, True, True, id="all-healthy"),
pytest.param(False, True, False, id="db-unhealthy"),
pytest.param(True, False, False, id="http-unhealthy"),
pytest.param(False, False, False, id="both-unhealthy"),
],
)
def test_is_healthy(
self, db_healthy: bool, http_healthy: bool, expected: bool
) -> None:
"""Derive overall health from component health."""
status = HealthStatus(database_healthy=db_healthy, http_healthy=http_healthy)
assert status.is_healthy is expected
class TestCheckDatabase:
"""Database health check."""
async def test_returns_true_for_healthy_db(
self, health_checker: HealthChecker
) -> None:
"""Return True when the database responds to queries."""
assert await health_checker.check_database() is True
async def test_returns_false_for_stopped_db(
self, health_checker: HealthChecker, database: Database
) -> None:
"""Return False when the database connection is closed."""
await database.stop()
assert await health_checker.check_database() is False
class TestPerformHealthCheck:
"""Health check orchestration and endpoint reporting."""
async def test_skips_report_when_no_endpoint(
self, health_checker: HealthChecker
) -> None:
"""Skip reporting when endpoint is empty."""
result = UpdateResult(total_streams=0, successful_checks=0, failed_checks=0)
with aioresponses():
await health_checker.perform_health_check(result, "")
async def test_skips_report_when_unhealthy(
self, health_checker: HealthChecker
) -> None:
"""Skip health report when all stream checks failed."""
result = UpdateResult(total_streams=3, successful_checks=0, failed_checks=3)
with aioresponses():
await health_checker.perform_health_check(
result, "https://health.example.com/ping"
)
async def test_sends_report_when_healthy(
self, health_checker: HealthChecker
) -> None:
"""Send GET to endpoint when all checks pass."""
result = UpdateResult(total_streams=1, successful_checks=1, failed_checks=0)
with aioresponses() as mocked:
mocked.get("https://health.example.com/ping", status=200)
await health_checker.perform_health_check(
result, "https://health.example.com/ping"
)
async def test_skips_report_for_whitespace_endpoint(
self, health_checker: HealthChecker
) -> None:
"""Skip reporting when endpoint is whitespace."""
result = UpdateResult(total_streams=0, successful_checks=0, failed_checks=0)
with aioresponses():
await health_checker.perform_health_check(result, " ")
class TestSendHealthReport:
"""Health report HTTP delivery."""
async def test_handles_success_response(
self, health_checker: HealthChecker
) -> None:
"""Complete without error on a 2xx response."""
with aioresponses() as mocked:
mocked.get("https://health.example.com/ping", status=200)
await health_checker._send_health_report("https://health.example.com/ping")
async def test_handles_non_success_response(
self, health_checker: HealthChecker
) -> None:
"""Complete without error on a non-2xx response."""
with aioresponses() as mocked:
mocked.get("https://health.example.com/ping", status=500)
await health_checker._send_health_report("https://health.example.com/ping")
async def test_handles_connection_error(
self, health_checker: HealthChecker
) -> None:
"""Complete without error on a connection failure."""
with aioresponses() as mocked:
mocked.get(
"https://health.example.com/ping",
exception=ConnectionError(),
)
await health_checker._send_health_report("https://health.example.com/ping")
+314
View File
@@ -0,0 +1,314 @@
# Copyright 2026 Logan Fick
#
# Licensed under the Apache License, Version 2.0 (the "License");
# you may not use this file except in compliance with the License.
# You may obtain a copy of the License at
#
# http://www.apache.org/licenses/LICENSE-2.0
#
# Unless required by applicable law or agreed to in writing, software
# distributed under the License is distributed on an "AS IS" BASIS,
# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
# See the License for the specific language governing permissions and
# limitations under the License.
"""Tests for the Prometheus metrics service."""
import pytest
from owncastsentry.metrics import ErrorSource, MetricsService, NotificationType
from owncastsentry.models import StreamStatus
from tests.conftest import generate_metrics_output
class TestRecordDelivery:
"""Notification delivery counter with type and result labels."""
@pytest.mark.parametrize(
("notification_type", "successful", "failed", "expected_fragments"),
[
pytest.param(
NotificationType.LIVE,
3,
0,
['result="success",type="live"} 3.0'],
id="live-success",
),
pytest.param(
NotificationType.LIVE,
0,
2,
['result="failure",type="live"} 2.0'],
id="live-failure",
),
pytest.param(
NotificationType.TITLE_CHANGE,
1,
0,
['result="success",type="title_change"} 1.0'],
id="title-change-success",
),
pytest.param(
NotificationType.CLEANUP_WARNING,
2,
0,
['result="success",type="cleanup_warning"} 2.0'],
id="cleanup-warning-success",
),
pytest.param(
NotificationType.CLEANUP_DELETION,
1,
1,
[
'result="success",type="cleanup_deletion"} 1.0',
'result="failure",type="cleanup_deletion"} 1.0',
],
id="cleanup-deletion-mixed",
),
],
)
def test_records_delivery(
self,
notification_type: NotificationType,
successful: int,
failed: int,
expected_fragments: list[str],
) -> None:
"""Record delivery results with correct type and result labels."""
service = MetricsService()
service.record_delivery(notification_type, successful=successful, failed=failed)
output = generate_metrics_output(service)
for fragment in expected_fragments:
assert fragment in output
@pytest.mark.parametrize(
("notification_type", "result"),
[
pytest.param(t, r, id=f"{t}-{r}")
for t in NotificationType
for r in ("success", "failure")
],
)
def test_all_combinations_initialized(
self, notification_type: NotificationType, result: str
) -> None:
"""All type/result label combinations exist at zero on init."""
service = MetricsService()
output = generate_metrics_output(service)
expected = (
f"owncastsentry_notification_delivery_total"
f'{{result="{result}",type="{notification_type}"}} 0.0'
)
assert expected in output
class TestSetStreamStatus:
"""Per-stream status gauge."""
@pytest.mark.parametrize(
("status", "expected_value"),
[
pytest.param(StreamStatus.ONLINE, 1.0, id="online"),
pytest.param(StreamStatus.OFFLINE, 0.0, id="offline"),
pytest.param(StreamStatus.UNKNOWN, -1.0, id="unknown"),
],
)
def test_sets_status(self, status: StreamStatus, expected_value: float) -> None:
"""Set gauge to the correct value for each stream status."""
service = MetricsService()
service.set_stream_status("test.com", status)
output = generate_metrics_output(service)
expected = f'owncastsentry_stream_status{{domain="test.com"}} {expected_value}'
assert expected in output
class TestSetSubscriptionCount:
"""Per-stream subscription count gauge."""
def test_sets_count(self) -> None:
"""Set the subscription count for a domain."""
service = MetricsService()
service.set_subscription_count("test.com", 5)
output = generate_metrics_output(service)
assert 'owncastsentry_stream_subscriptions{domain="test.com"} 5.0' in output
def test_updates_count(self) -> None:
"""Update the subscription count for a domain."""
service = MetricsService()
service.set_subscription_count("test.com", 5)
service.set_subscription_count("test.com", 3)
output = generate_metrics_output(service)
assert 'owncastsentry_stream_subscriptions{domain="test.com"} 3.0' in output
class TestSetCheckFailures:
"""Check failure counter gauge per domain."""
def test_sets_count(self) -> None:
"""Set the failure count for a domain."""
service = MetricsService()
service.set_check_failures("fail.com", 3)
output = generate_metrics_output(service)
assert 'owncastsentry_check_failures{domain="fail.com"} 3.0' in output
def test_resets_to_zero(self) -> None:
"""Reset the failure count to zero."""
service = MetricsService()
service.set_check_failures("fail.com", 5)
service.set_check_failures("fail.com", 0)
output = generate_metrics_output(service)
assert 'owncastsentry_check_failures{domain="fail.com"} 0.0' in output
class TestResponseTimer:
"""Response time gauge via context manager."""
def test_records_on_success(self) -> None:
"""Record a response time when success() is called."""
service = MetricsService()
with service.response_timer("example.com") as timer:
timer.success()
output = generate_metrics_output(service)
assert 'owncastsentry_api_response_seconds{domain="example.com"}' in output
def test_does_not_record_without_success(self) -> None:
"""Do not record when success() is never called."""
service = MetricsService()
with service.response_timer("example.com"):
pass
output = generate_metrics_output(service)
assert 'owncastsentry_api_response_seconds{domain="example.com"}' not in output
def test_overwrites_previous_value(self) -> None:
"""Overwrite previous value with the latest response time."""
service = MetricsService()
with service.response_timer("example.com") as timer:
timer.success()
with service.response_timer("example.com") as timer:
timer.success()
output = generate_metrics_output(service)
# Gauge should have exactly one line for this domain, not accumulated
matches = [
line
for line in output.splitlines()
if line.startswith("owncastsentry_api_response_seconds{")
]
assert len(matches) == 1
def test_does_not_record_on_exception(self) -> None:
"""Do not record when the block raises an exception."""
service = MetricsService()
with (
pytest.raises(ValueError, match="boom"),
service.response_timer("example.com"),
):
raise ValueError("boom")
output = generate_metrics_output(service)
assert 'owncastsentry_api_response_seconds{domain="example.com"}' not in output
class TestRemoveStream:
"""Stale stream label cleanup."""
def test_removes_stream_label(self) -> None:
"""Remove a stream's gauge labels after cleanup deletion."""
service = MetricsService()
service.set_stream_status("gone.com", StreamStatus.OFFLINE)
service.set_subscription_count("gone.com", 2)
assert 'domain="gone.com"' in generate_metrics_output(service)
service.remove_stream("gone.com")
assert 'domain="gone.com"' not in generate_metrics_output(service)
def test_remove_nonexistent_is_noop(self) -> None:
"""Removing a nonexistent stream does not raise."""
service = MetricsService()
service.remove_stream("never.com")
class TestRegisterOpenConnectionsGauge:
"""Callback-based open connection gauge."""
def test_reads_value_from_callback(self) -> None:
"""Read the open connection count from the callback at scrape time."""
service = MetricsService()
counter = [3]
service.register_open_connections_gauge(lambda: counter[0])
output = generate_metrics_output(service)
assert "owncastsentry_http_connections_open 3.0" in output
def test_reflects_updated_value(self) -> None:
"""Reflect changes in the callback value on subsequent scrapes."""
service = MetricsService()
counter = [1]
service.register_open_connections_gauge(lambda: counter[0])
counter[0] = 5
output = generate_metrics_output(service)
assert "owncastsentry_http_connections_open 5.0" in output
class TestSetBuildInfo:
"""Build version info metric."""
def test_sets_version(self) -> None:
"""Set the build version info."""
service = MetricsService()
service.set_build_info("1.2.3")
output = generate_metrics_output(service)
assert 'owncastsentry_info{version="1.2.3"} 1.0' in output
class TestRecordError:
"""Internal error counter."""
def test_increments_counter(self) -> None:
"""Increment the error counter for a source."""
service = MetricsService()
service.record_error(ErrorSource.SCHEDULER_LOOP)
output = generate_metrics_output(service)
assert 'owncastsentry_errors_total{source="scheduler_loop"} 1.0' in output
def test_increments_multiple_sources(self) -> None:
"""Increment error counters for different sources independently."""
service = MetricsService()
service.record_error(ErrorSource.SCHEDULER_LOOP)
service.record_error(ErrorSource.COMMAND)
service.record_error(ErrorSource.COMMAND)
output = generate_metrics_output(service)
assert 'owncastsentry_errors_total{source="scheduler_loop"} 1.0' in output
assert 'owncastsentry_errors_total{source="command"} 2.0' in output
@pytest.mark.parametrize(
"source",
[pytest.param(s, id=s) for s in ErrorSource],
)
def test_all_sources_initialized(self, source: ErrorSource) -> None:
"""All known source labels exist at zero on init."""
service = MetricsService()
output = generate_metrics_output(service)
expected = f'owncastsentry_errors_total{{source="{source}"}} 0.0'
assert expected in output
class TestRegistryOutput:
"""Prometheus registry output."""
def test_returns_string(self) -> None:
"""Return a string (not bytes)."""
service = MetricsService()
output = generate_metrics_output(service)
assert isinstance(output, str)
def test_contains_help_lines(self) -> None:
"""Include HELP lines for registered metrics."""
service = MetricsService()
output = generate_metrics_output(service)
assert "# HELP owncastsentry_notification_delivery_total" in output
assert "# HELP owncastsentry_errors_total" in output
assert "# HELP owncastsentry_info" in output
def test_uses_isolated_registry(self) -> None:
"""Use a custom registry, not the global default."""
service = MetricsService()
output = generate_metrics_output(service)
assert "python_gc" not in output
assert "process_" not in output
-2
View File
@@ -14,8 +14,6 @@
"""Tests for data models.""" """Tests for data models."""
from __future__ import annotations
import pytest import pytest
from owncastsentry.models import StreamConfig, StreamState, StreamStatus from owncastsentry.models import StreamConfig, StreamState, StreamStatus
+115 -6
View File
@@ -14,17 +14,16 @@
"""Tests for the notification service.""" """Tests for the notification service."""
from __future__ import annotations
import logging import logging
import time import time
from typing import TYPE_CHECKING from typing import TYPE_CHECKING
import pytest import pytest
from owncastsentry.metrics import MetricsService
from owncastsentry.notification_service import NotificationService from owncastsentry.notification_service import NotificationService
from owncastsentry.utils import SECONDS_BETWEEN_NOTIFICATIONS from owncastsentry.utils import SECONDS_BETWEEN_NOTIFICATIONS
from tests.conftest import _StubMatrixClient from tests.conftest import _StubMatrixClient, generate_metrics_output
if TYPE_CHECKING: if TYPE_CHECKING:
from owncastsentry.database import StreamRepository, SubscriptionRepository from owncastsentry.database import StreamRepository, SubscriptionRepository
@@ -34,12 +33,14 @@ def _make_service(
*, *,
client: _StubMatrixClient, client: _StubMatrixClient,
subscription_repo: SubscriptionRepository, subscription_repo: SubscriptionRepository,
metrics: MetricsService | None = None,
) -> NotificationService: ) -> NotificationService:
"""Build a NotificationService with a stub client and real repo.""" """Build a NotificationService with a stub client and real repo."""
return NotificationService( return NotificationService(
client=client, client=client,
subscription_repo=subscription_repo, subscription_repo=subscription_repo,
logger=logging.getLogger("test"), logger=logging.getLogger("test"),
metrics=metrics or MetricsService(),
) )
@@ -62,7 +63,7 @@ class TestCanNotify:
service = _make_service( service = _make_service(
client=_StubMatrixClient(), subscription_repo=subscription_repo client=_StubMatrixClient(), subscription_repo=subscription_repo
) )
service.notification_timers_cache["example.com"] = time.time() service.notification_timers_cache["example.com"] = time.monotonic()
assert service._can_notify("example.com") is False assert service._can_notify("example.com") is False
def test_after_cooldown_allowed( def test_after_cooldown_allowed(
@@ -74,7 +75,7 @@ class TestCanNotify:
) )
# Subtract an extra second to ensure the cooldown has fully elapsed # Subtract an extra second to ensure the cooldown has fully elapsed
service.notification_timers_cache["example.com"] = ( service.notification_timers_cache["example.com"] = (
time.time() - SECONDS_BETWEEN_NOTIFICATIONS - 1 time.monotonic() - SECONDS_BETWEEN_NOTIFICATIONS - 1
) )
assert service._can_notify("example.com") is True assert service._can_notify("example.com") is True
@@ -253,7 +254,7 @@ class TestNotifyStreamLive:
"""Skip sending when the domain is within the rate-limit cooldown.""" """Skip sending when the domain is within the rate-limit cooldown."""
client = _StubMatrixClient() client = _StubMatrixClient()
service = _make_service(client=client, subscription_repo=subscription_repo) service = _make_service(client=client, subscription_repo=subscription_repo)
service.notification_timers_cache["example.com"] = time.time() service.notification_timers_cache["example.com"] = time.monotonic()
await stream_repo.create("example.com") await stream_repo.create("example.com")
await subscription_repo.add("example.com", "!room1:matrix.org") await subscription_repo.add("example.com", "!room1:matrix.org")
@@ -346,3 +347,111 @@ class TestSendCleanupDeletion:
"If the instance comes online again and you want to " "If the instance comes online again and you want to "
"resubscribe, run `!subscribe example.com`." "resubscribe, run `!subscribe example.com`."
) )
class TestNotificationMetrics:
"""Notification metrics recording."""
async def test_records_live_notification(
self,
stream_repo: StreamRepository,
subscription_repo: SubscriptionRepository,
) -> None:
"""Record a live notification metric."""
client = _StubMatrixClient()
metrics = MetricsService()
service = _make_service(
client=client, subscription_repo=subscription_repo, metrics=metrics
)
await stream_repo.create("example.com")
await subscription_repo.add("example.com", "!room:matrix.org")
await service.notify_stream_live("example.com", "Stream", "Title", [])
output = generate_metrics_output(metrics)
expected = (
"owncastsentry_notification_delivery_total"
'{result="success",type="live"} 1.0'
)
assert expected in output
async def test_records_title_change_notification(
self,
stream_repo: StreamRepository,
subscription_repo: SubscriptionRepository,
) -> None:
"""Record a title_change notification metric."""
client = _StubMatrixClient()
metrics = MetricsService()
service = _make_service(
client=client, subscription_repo=subscription_repo, metrics=metrics
)
await stream_repo.create("example.com")
await subscription_repo.add("example.com", "!room:matrix.org")
await service.notify_stream_live(
"example.com", "Stream", "Title", [], title_change=True
)
output = generate_metrics_output(metrics)
expected = (
"owncastsentry_notification_delivery_total"
'{result="success",type="title_change"} 1.0'
)
assert expected in output
async def test_records_cleanup_warning_notification(
self,
stream_repo: StreamRepository,
subscription_repo: SubscriptionRepository,
) -> None:
"""Record a cleanup_warning notification metric."""
client = _StubMatrixClient()
metrics = MetricsService()
service = _make_service(
client=client, subscription_repo=subscription_repo, metrics=metrics
)
await stream_repo.create("example.com")
await subscription_repo.add("example.com", "!room:matrix.org")
await service.send_cleanup_warning("example.com")
output = generate_metrics_output(metrics)
expected = (
"owncastsentry_notification_delivery_total"
'{result="success",type="cleanup_warning"} 1.0'
)
assert expected in output
async def test_records_cleanup_deletion_notification(
self,
stream_repo: StreamRepository,
subscription_repo: SubscriptionRepository,
) -> None:
"""Record a cleanup_deletion notification metric."""
client = _StubMatrixClient()
metrics = MetricsService()
service = _make_service(
client=client, subscription_repo=subscription_repo, metrics=metrics
)
await stream_repo.create("example.com")
await subscription_repo.add("example.com", "!room:matrix.org")
await service.send_cleanup_deletion("example.com")
output = generate_metrics_output(metrics)
expected = (
"owncastsentry_notification_delivery_total"
'{result="success",type="cleanup_deletion"} 1.0'
)
assert expected in output
async def test_no_metric_when_rate_limited(
self,
stream_repo: StreamRepository,
subscription_repo: SubscriptionRepository,
) -> None:
"""Do not record metric when notification is rate-limited."""
client = _StubMatrixClient()
metrics = MetricsService()
service = _make_service(
client=client, subscription_repo=subscription_repo, metrics=metrics
)
service.notification_timers_cache["example.com"] = time.monotonic()
await stream_repo.create("example.com")
await subscription_repo.add("example.com", "!room:matrix.org")
await service.notify_stream_live("example.com", "Stream", "Title", [])
output = generate_metrics_output(metrics)
assert 'result="success",type="live"} 0.0' in output
+77 -4
View File
@@ -14,8 +14,6 @@
"""Tests for the Owncast HTTP client.""" """Tests for the Owncast HTTP client."""
from __future__ import annotations
import json import json
import logging import logging
from typing import TYPE_CHECKING from typing import TYPE_CHECKING
@@ -23,8 +21,13 @@ from typing import TYPE_CHECKING
import pytest import pytest
from aioresponses import aioresponses from aioresponses import aioresponses
from owncastsentry.metrics import MetricsService
from owncastsentry.owncast_client import OwncastClient from owncastsentry.owncast_client import OwncastClient
from tests.conftest import VALID_CONFIG_RESPONSE, VALID_STATUS_RESPONSE from tests.conftest import (
VALID_CONFIG_RESPONSE,
VALID_STATUS_RESPONSE,
generate_metrics_output,
)
if TYPE_CHECKING: if TYPE_CHECKING:
from collections.abc import AsyncIterator from collections.abc import AsyncIterator
@@ -33,7 +36,11 @@ if TYPE_CHECKING:
@pytest.fixture @pytest.fixture
async def owncast_client() -> AsyncIterator[OwncastClient]: async def owncast_client() -> AsyncIterator[OwncastClient]:
"""Create an OwncastClient and close it after the test.""" """Create an OwncastClient and close it after the test."""
client = OwncastClient(logger=logging.getLogger("test"), version="0.0.0") client = OwncastClient(
logger=logging.getLogger("test"),
version="0.0.0",
metrics=MetricsService(),
)
yield client yield client
await client.close() await client.close()
@@ -205,3 +212,69 @@ class TestValidateInstance:
result = await owncast_client.validate_instance("invalid.com") result = await owncast_client.validate_instance("invalid.com")
assert result is False assert result is False
class TestResponseTimeMetrics:
"""Response time histogram recording."""
async def test_records_on_success(self) -> None:
"""Record response time on a successful request."""
metrics = MetricsService()
client = OwncastClient(
logger=logging.getLogger("test"),
version="0.0.0",
metrics=metrics,
)
with aioresponses() as mocked:
mocked.get(
"https://example.com/api/status",
body=json.dumps(VALID_STATUS_RESPONSE).encode(),
)
await client.get_stream_state("example.com")
output = generate_metrics_output(metrics)
assert 'owncastsentry_api_response_seconds{domain="example.com"}' in output
await client.close()
async def test_no_observation_on_failure(self) -> None:
"""Do not record response time when request fails."""
metrics = MetricsService()
client = OwncastClient(
logger=logging.getLogger("test"),
version="0.0.0",
metrics=metrics,
)
with aioresponses() as mocked:
mocked.get(
"https://example.com/api/status",
status=500,
)
await client.get_stream_state("example.com")
output = generate_metrics_output(metrics)
assert 'owncastsentry_api_response_seconds{domain="example.com"}' not in output
await client.close()
async def test_no_observation_on_connection_error(self) -> None:
"""Do not record response time on connection error."""
metrics = MetricsService()
client = OwncastClient(
logger=logging.getLogger("test"),
version="0.0.0",
metrics=metrics,
)
with aioresponses() as mocked:
mocked.get(
"https://example.com/api/status",
exception=ConnectionError(),
)
await client.get_stream_state("example.com")
output = generate_metrics_output(metrics)
assert 'owncastsentry_api_response_seconds{domain="example.com"}' not in output
await client.close()
class TestOpenConnectionCount:
"""Open connection count."""
async def test_zero_with_no_requests(self, owncast_client: OwncastClient) -> None:
"""Return zero when no requests have been made."""
assert owncast_client.open_connection_count == 0
+212 -11
View File
@@ -14,13 +14,12 @@
"""Tests for the stream monitor.""" """Tests for the stream monitor."""
from __future__ import annotations
import logging import logging
import time import time
from typing import TYPE_CHECKING from typing import TYPE_CHECKING
from owncastsentry.models import StreamConfig, StreamState from owncastsentry.metrics import MetricsService
from owncastsentry.models import StreamConfig, StreamState, StreamStatus
from owncastsentry.notification_service import NotificationService from owncastsentry.notification_service import NotificationService
from owncastsentry.stream_monitor import StreamMonitor from owncastsentry.stream_monitor import StreamMonitor
from owncastsentry.utils import ( from owncastsentry.utils import (
@@ -28,7 +27,11 @@ from owncastsentry.utils import (
CLEANUP_WARNING_THRESHOLD, CLEANUP_WARNING_THRESHOLD,
SECONDS_BETWEEN_NOTIFICATIONS, SECONDS_BETWEEN_NOTIFICATIONS,
) )
from tests.conftest import _StubMatrixClient, _StubOwncastClient from tests.conftest import (
_StubMatrixClient,
_StubOwncastClient,
generate_metrics_output,
)
if TYPE_CHECKING: if TYPE_CHECKING:
from owncastsentry.database import StreamRepository, SubscriptionRepository from owncastsentry.database import StreamRepository, SubscriptionRepository
@@ -40,13 +43,16 @@ def _make_monitor(
stream_repo: StreamRepository, stream_repo: StreamRepository,
subscription_repo: SubscriptionRepository, subscription_repo: SubscriptionRepository,
client: _StubMatrixClient, client: _StubMatrixClient,
metrics: MetricsService | None = None,
) -> tuple[StreamMonitor, NotificationService]: ) -> tuple[StreamMonitor, NotificationService]:
"""Build a StreamMonitor with stubs and a real NotificationService.""" """Build a StreamMonitor with stubs and a real NotificationService."""
logger = logging.getLogger("test") logger = logging.getLogger("test")
metrics = metrics or MetricsService()
notification_service = NotificationService( notification_service = NotificationService(
client=client, client=client,
subscription_repo=subscription_repo, subscription_repo=subscription_repo,
logger=logger, logger=logger,
metrics=metrics,
) )
monitor = StreamMonitor( monitor = StreamMonitor(
owncast_client=owncast_client, owncast_client=owncast_client,
@@ -54,6 +60,7 @@ def _make_monitor(
subscription_repo=subscription_repo, subscription_repo=subscription_repo,
notification_service=notification_service, notification_service=notification_service,
logger=logger, logger=logger,
metrics=metrics,
) )
return monitor, notification_service return monitor, notification_service
@@ -82,6 +89,25 @@ async def _seed_stream(
await subscription_repo.add(domain, room_id) await subscription_repo.add(domain, room_id)
def _make_monitor_with_metrics(
*,
owncast_client: _StubOwncastClient,
stream_repo: StreamRepository,
subscription_repo: SubscriptionRepository,
client: _StubMatrixClient,
) -> tuple[StreamMonitor, NotificationService, MetricsService]:
"""Build a StreamMonitor with stubs, a real NotificationService, and metrics."""
metrics = MetricsService()
monitor, notification_service = _make_monitor(
owncast_client=owncast_client,
stream_repo=stream_repo,
subscription_repo=subscription_repo,
client=client,
metrics=metrics,
)
return monitor, notification_service, metrics
class TestUpdateAllStreams: class TestUpdateAllStreams:
"""Parallel stream update orchestration.""" """Parallel stream update orchestration."""
@@ -337,7 +363,7 @@ class TestUpdateStreamBriefOffline:
) )
# Recently offline (within cooldown) # Recently offline (within cooldown)
monitor.offline_timer_cache["brief.com"] = time.time() - 60 monitor.offline_timer_cache["brief.com"] = time.monotonic() - 60
result = await monitor.update_stream("brief.com") result = await monitor.update_stream("brief.com")
assert result is True assert result is True
@@ -375,7 +401,7 @@ class TestUpdateStreamBriefOffline:
) )
# Recently offline (within cooldown) # Recently offline (within cooldown)
monitor.offline_timer_cache["brief.com"] = time.time() - 60 monitor.offline_timer_cache["brief.com"] = time.monotonic() - 60
result = await monitor.update_stream("brief.com") result = await monitor.update_stream("brief.com")
assert result is True assert result is True
@@ -424,7 +450,7 @@ class TestUpdateStreamTitleChange:
monitor.offline_timer_cache["title.com"] = 0 monitor.offline_timer_cache["title.com"] = 0
# Subtract an extra second to ensure the cooldown has fully elapsed # Subtract an extra second to ensure the cooldown has fully elapsed
notification_service.notification_timers_cache["title.com"] = ( notification_service.notification_timers_cache["title.com"] = (
time.time() - SECONDS_BETWEEN_NOTIFICATIONS - 1 time.monotonic() - SECONDS_BETWEEN_NOTIFICATIONS - 1
) )
result = await monitor.update_stream("title.com") result = await monitor.update_stream("title.com")
@@ -473,7 +499,7 @@ class TestUpdateStreamTitleChange:
monitor.offline_timer_cache["title.com"] = 0 monitor.offline_timer_cache["title.com"] = 0
# Subtract an extra second to ensure the cooldown has fully elapsed # Subtract an extra second to ensure the cooldown has fully elapsed
notification_service.notification_timers_cache["title.com"] = ( notification_service.notification_timers_cache["title.com"] = (
time.time() - SECONDS_BETWEEN_NOTIFICATIONS - 1 time.monotonic() - SECONDS_BETWEEN_NOTIFICATIONS - 1
) )
result = await monitor.update_stream("title.com") result = await monitor.update_stream("title.com")
@@ -518,7 +544,7 @@ class TestUpdateStreamTitleChange:
# Offline timer is MORE recent than last notification, # Offline timer is MORE recent than last notification,
# and both are old enough to pass rate limiting # and both are old enough to pass rate limiting
now = time.time() 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
) )
@@ -571,9 +597,9 @@ class TestUpdateStreamGoesOffline:
) )
monitor.offline_timer_cache["offline.com"] = 0 monitor.offline_timer_cache["offline.com"] = 0
before = time.time() before = time.monotonic()
result = await monitor.update_stream("offline.com") result = await monitor.update_stream("offline.com")
after = time.time() after = time.monotonic()
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
@@ -837,3 +863,178 @@ class TestUpdateAllStreamsMixed:
assert result.total_streams == 2 assert result.total_streams == 2
assert result.successful_checks == 1 assert result.successful_checks == 1
assert result.failed_checks == 1 assert result.failed_checks == 1
class TestStreamMonitorMetrics:
"""Stream monitor metrics recording."""
async def test_records_stream_status_online(
self,
stream_repo: StreamRepository,
subscription_repo: SubscriptionRepository,
) -> None:
"""Record online stream status gauge."""
owncast = _StubOwncastClient(
stream_state=StreamState(
domain="live.com",
title="Title",
last_connect_time="2026-01-01T12:00:00Z",
last_disconnect_time="2026-01-01T10:00:00Z",
),
stream_config=StreamConfig(name="Live Stream"),
)
client = _StubMatrixClient()
monitor, _, metrics = _make_monitor_with_metrics(
owncast_client=owncast,
stream_repo=stream_repo,
subscription_repo=subscription_repo,
client=client,
)
await _seed_stream(
stream_repo,
subscription_repo,
domain="live.com",
last_disconnect_time="2026-01-01T10:00:00Z",
)
monitor.offline_timer_cache["live.com"] = 0
await monitor.update_stream("live.com")
output = generate_metrics_output(metrics)
assert 'owncastsentry_stream_status{domain="live.com"} 1.0' in output
async def test_records_stream_status_offline(
self,
stream_repo: StreamRepository,
subscription_repo: SubscriptionRepository,
) -> None:
"""Record offline stream status gauge."""
owncast = _StubOwncastClient(
stream_state=StreamState(
domain="off.com",
title="Title",
last_disconnect_time="2026-01-01T12:00:00Z",
),
stream_config=StreamConfig(name="Off Stream"),
)
client = _StubMatrixClient()
monitor, _, metrics = _make_monitor_with_metrics(
owncast_client=owncast,
stream_repo=stream_repo,
subscription_repo=subscription_repo,
client=client,
)
await _seed_stream(
stream_repo,
subscription_repo,
domain="off.com",
title="Title",
last_disconnect_time="2026-01-01T12:00:00Z",
)
await monitor.update_stream("off.com")
output = generate_metrics_output(metrics)
assert 'owncastsentry_stream_status{domain="off.com"} 0.0' in output
async def test_records_check_failures_on_connection_failure(
self,
stream_repo: StreamRepository,
subscription_repo: SubscriptionRepository,
) -> None:
"""Record failure counter gauge when a stream check fails."""
owncast = _StubOwncastClient(stream_state=None)
client = _StubMatrixClient()
monitor, _, metrics = _make_monitor_with_metrics(
owncast_client=owncast,
stream_repo=stream_repo,
subscription_repo=subscription_repo,
client=client,
)
await _seed_stream(
stream_repo,
subscription_repo,
domain="fail.com",
last_disconnect_time="2026-01-01T00:00:00Z",
)
await monitor.update_stream("fail.com")
output = generate_metrics_output(metrics)
assert 'owncastsentry_check_failures{domain="fail.com"} 1.0' in output
async def test_resets_check_failures_on_success(
self,
stream_repo: StreamRepository,
subscription_repo: SubscriptionRepository,
) -> None:
"""Reset failure counter gauge to zero on successful check."""
owncast = _StubOwncastClient(
stream_state=StreamState(
domain="recover.com",
title="Title",
last_disconnect_time="2026-01-01T12:00:00Z",
),
stream_config=StreamConfig(name="Recover"),
)
client = _StubMatrixClient()
monitor, _, metrics = _make_monitor_with_metrics(
owncast_client=owncast,
stream_repo=stream_repo,
subscription_repo=subscription_repo,
client=client,
)
await _seed_stream(
stream_repo,
subscription_repo,
domain="recover.com",
last_disconnect_time="2026-01-01T12:00:00Z",
)
# Simulate prior failures
for _ in range(3):
await stream_repo.increment_failure_counter("recover.com")
await monitor.update_stream("recover.com")
output = generate_metrics_output(metrics)
assert 'owncastsentry_check_failures{domain="recover.com"} 0.0' in output
async def test_removes_stream_on_cleanup_deletion(
self,
stream_repo: StreamRepository,
subscription_repo: SubscriptionRepository,
) -> None:
"""Remove stream gauge label on cleanup deletion."""
owncast = _StubOwncastClient()
client = _StubMatrixClient()
monitor, _, metrics = _make_monitor_with_metrics(
owncast_client=owncast,
stream_repo=stream_repo,
subscription_repo=subscription_repo,
client=client,
)
await _seed_stream(stream_repo, subscription_repo, domain="delete.com")
metrics.set_stream_status("delete.com", StreamStatus.OFFLINE)
assert 'domain="delete.com"' in generate_metrics_output(metrics)
await monitor._check_cleanup_thresholds("delete.com", CLEANUP_DELETE_THRESHOLD)
assert 'domain="delete.com"' not in generate_metrics_output(metrics)
async def test_records_subscription_counts(
self,
stream_repo: StreamRepository,
subscription_repo: SubscriptionRepository,
) -> None:
"""Record per-domain subscription counts after update_all_streams."""
owncast = _StubOwncastClient(
stream_state=StreamState(
domain="pop.com",
last_connect_time="2026-01-01T00:00:00Z",
),
stream_config=StreamConfig(name="Popular"),
)
client = _StubMatrixClient()
monitor, _, metrics = _make_monitor_with_metrics(
owncast_client=owncast,
stream_repo=stream_repo,
subscription_repo=subscription_repo,
client=client,
)
await stream_repo.create("pop.com")
await subscription_repo.add("pop.com", "!room1:matrix.org")
await subscription_repo.add("pop.com", "!room2:matrix.org")
await monitor.update_all_streams(["pop.com"])
output = generate_metrics_output(metrics)
assert 'owncastsentry_stream_subscriptions{domain="pop.com"} 2.0' in output
-2
View File
@@ -14,8 +14,6 @@
"""Tests for utility functions and constants.""" """Tests for utility functions and constants."""
from __future__ import annotations
import pytest import pytest
from owncastsentry.utils import ( from owncastsentry.utils import (
Generated
+33 -20
View File
@@ -924,6 +924,9 @@ wheels = [
[[package]] [[package]]
name = "owncastsentry" name = "owncastsentry"
source = { editable = "." } source = { editable = "." }
dependencies = [
{ name = "prometheus-client" },
]
[package.dev-dependencies] [package.dev-dependencies]
dev = [ dev = [
@@ -941,6 +944,7 @@ dev = [
] ]
[package.metadata] [package.metadata]
requires-dist = [{ name = "prometheus-client", specifier = ">=0.24.1" }]
[package.metadata.requires-dev] [package.metadata.requires-dev]
dev = [ dev = [
@@ -953,7 +957,7 @@ dev = [
{ name = "pytest", specifier = ">=9.0.2" }, { name = "pytest", specifier = ">=9.0.2" },
{ name = "pytest-asyncio", specifier = ">=1.3.0" }, { name = "pytest-asyncio", specifier = ">=1.3.0" },
{ name = "pytest-cov", specifier = ">=7.0.0" }, { name = "pytest-cov", specifier = ">=7.0.0" },
{ name = "ruff", specifier = ">=0.15.5" }, { name = "ruff", specifier = ">=0.15.6" },
{ name = "time-machine", specifier = ">=3.2.0" }, { name = "time-machine", specifier = ">=3.2.0" },
] ]
@@ -1069,6 +1073,15 @@ wheels = [
{ url = "https://files.pythonhosted.org/packages/54/20/4d324d65cc6d9205fabedc306948156824eb9f0ee1633355a8f7ec5c66bf/pluggy-1.6.0-py3-none-any.whl", hash = "sha256:e920276dd6813095e9377c0bc5566d94c932c33b27a3e3945d8389c374dd4746", size = 20538, upload-time = "2025-05-15T12:30:06.134Z" }, { url = "https://files.pythonhosted.org/packages/54/20/4d324d65cc6d9205fabedc306948156824eb9f0ee1633355a8f7ec5c66bf/pluggy-1.6.0-py3-none-any.whl", hash = "sha256:e920276dd6813095e9377c0bc5566d94c932c33b27a3e3945d8389c374dd4746", size = 20538, upload-time = "2025-05-15T12:30:06.134Z" },
] ]
[[package]]
name = "prometheus-client"
version = "0.24.1"
source = { registry = "https://pypi.org/simple" }
sdist = { url = "https://files.pythonhosted.org/packages/f0/58/a794d23feb6b00fc0c72787d7e87d872a6730dd9ed7c7b3e954637d8f280/prometheus_client-0.24.1.tar.gz", hash = "sha256:7e0ced7fbbd40f7b84962d5d2ab6f17ef88a72504dcf7c0b40737b43b2a461f9", size = 85616, upload-time = "2026-01-14T15:26:26.965Z" }
wheels = [
{ url = "https://files.pythonhosted.org/packages/74/c3/24a2f845e3917201628ecaba4f18bab4d18a337834c1df2a159ee9d22a42/prometheus_client-0.24.1-py3-none-any.whl", hash = "sha256:150db128af71a5c2482b36e588fc8a6b95e498750da4b17065947c16070f4055", size = 64057, upload-time = "2026-01-14T15:26:24.42Z" },
]
[[package]] [[package]]
name = "prompt-toolkit" name = "prompt-toolkit"
version = "3.0.52" version = "3.0.52"
@@ -1341,27 +1354,27 @@ wheels = [
[[package]] [[package]]
name = "ruff" name = "ruff"
version = "0.15.5" version = "0.15.6"
source = { registry = "https://pypi.org/simple" } source = { registry = "https://pypi.org/simple" }
sdist = { url = "https://files.pythonhosted.org/packages/77/9b/840e0039e65fcf12758adf684d2289024d6140cde9268cc59887dc55189c/ruff-0.15.5.tar.gz", hash = "sha256:7c3601d3b6d76dce18c5c824fc8d06f4eef33d6df0c21ec7799510cde0f159a2", size = 4574214, upload-time = "2026-03-05T20:06:34.946Z" } sdist = { url = "https://files.pythonhosted.org/packages/51/df/f8629c19c5318601d3121e230f74cbee7a3732339c52b21daa2b82ef9c7d/ruff-0.15.6.tar.gz", hash = "sha256:8394c7bb153a4e3811a4ecdacd4a8e6a4fa8097028119160dffecdcdf9b56ae4", size = 4597916, upload-time = "2026-03-12T23:05:47.51Z" }
wheels = [ wheels = [
{ url = "https://files.pythonhosted.org/packages/47/20/5369c3ce21588c708bcbe517a8fbe1a8dfdb5dfd5137e14790b1da71612c/ruff-0.15.5-py3-none-linux_armv6l.whl", hash = "sha256:4ae44c42281f42e3b06b988e442d344a5b9b72450ff3c892e30d11b29a96a57c", size = 10478185, upload-time = "2026-03-05T20:06:29.093Z" }, { url = "https://files.pythonhosted.org/packages/9e/2f/4e03a7e5ce99b517e98d3b4951f411de2b0fa8348d39cf446671adcce9a2/ruff-0.15.6-py3-none-linux_armv6l.whl", hash = "sha256:7c98c3b16407b2cf3d0f2b80c80187384bc92c6774d85fefa913ecd941256fff", size = 10508953, upload-time = "2026-03-12T23:05:17.246Z" },
{ url = "https://files.pythonhosted.org/packages/44/ed/e81dd668547da281e5dce710cf0bc60193f8d3d43833e8241d006720e42b/ruff-0.15.5-py3-none-macosx_10_12_x86_64.whl", hash = "sha256:6edd3792d408ebcf61adabc01822da687579a1a023f297618ac27a5b51ef0080", size = 10859201, upload-time = "2026-03-05T20:06:32.632Z" }, { url = "https://files.pythonhosted.org/packages/70/60/55bcdc3e9f80bcf39edf0cd272da6fa511a3d94d5a0dd9e0adf76ceebdb4/ruff-0.15.6-py3-none-macosx_10_12_x86_64.whl", hash = "sha256:ee7dcfaad8b282a284df4aa6ddc2741b3f4a18b0555d626805555a820ea181c3", size = 10942257, upload-time = "2026-03-12T23:05:23.076Z" },
{ url = "https://files.pythonhosted.org/packages/c4/8f/533075f00aaf19b07c5cd6aa6e5d89424b06b3b3f4583bfa9c640a079059/ruff-0.15.5-py3-none-macosx_11_0_arm64.whl", hash = "sha256:89f463f7c8205a9f8dea9d658d59eff49db05f88f89cc3047fb1a02d9f344010", size = 10184752, upload-time = "2026-03-05T20:06:40.312Z" }, { url = "https://files.pythonhosted.org/packages/e7/f9/005c29bd1726c0f492bfa215e95154cf480574140cb5f867c797c18c790b/ruff-0.15.6-py3-none-macosx_11_0_arm64.whl", hash = "sha256:3bd9967851a25f038fc8b9ae88a7fbd1b609f30349231dffaa37b6804923c4bb", size = 10322683, upload-time = "2026-03-12T23:05:33.738Z" },
{ url = "https://files.pythonhosted.org/packages/66/0e/ba49e2c3fa0395b3152bad634c7432f7edfc509c133b8f4529053ff024fb/ruff-0.15.5-py3-none-manylinux_2_17_aarch64.manylinux2014_aarch64.whl", hash = "sha256:ba786a8295c6574c1116704cf0b9e6563de3432ac888d8f83685654fe528fd65", size = 10534857, upload-time = "2026-03-05T20:06:19.581Z" }, { url = "https://files.pythonhosted.org/packages/5f/74/2f861f5fd7cbb2146bddb5501450300ce41562da36d21868c69b7a828169/ruff-0.15.6-py3-none-manylinux_2_17_aarch64.manylinux2014_aarch64.whl", hash = "sha256:13f4594b04e42cd24a41da653886b04d2ff87adbf57497ed4f728b0e8a4866f8", size = 10660986, upload-time = "2026-03-12T23:05:53.245Z" },
{ url = "https://files.pythonhosted.org/packages/59/71/39234440f27a226475a0659561adb0d784b4d247dfe7f43ffc12dd02e288/ruff-0.15.5-py3-none-manylinux_2_17_armv7l.manylinux2014_armv7l.whl", hash = "sha256:fd4b801e57955fe9f02b31d20375ab3a5c4415f2e5105b79fb94cf2642c91440", size = 10309120, upload-time = "2026-03-05T20:06:00.435Z" }, { url = "https://files.pythonhosted.org/packages/c1/a1/309f2364a424eccb763cdafc49df843c282609f47fe53aa83f38272389e0/ruff-0.15.6-py3-none-manylinux_2_17_armv7l.manylinux2014_armv7l.whl", hash = "sha256:e2ed8aea2f3fe57886d3f00ea5b8aae5bf68d5e195f487f037a955ff9fbaac9e", size = 10332177, upload-time = "2026-03-12T23:05:56.145Z" },
{ url = "https://files.pythonhosted.org/packages/f5/87/4140aa86a93df032156982b726f4952aaec4a883bb98cb6ef73c347da253/ruff-0.15.5-py3-none-manylinux_2_17_i686.manylinux2014_i686.whl", hash = "sha256:391f7c73388f3d8c11b794dbbc2959a5b5afe66642c142a6effa90b45f6f5204", size = 11047428, upload-time = "2026-03-05T20:05:51.867Z" }, { url = "https://files.pythonhosted.org/packages/30/41/7ebf1d32658b4bab20f8ac80972fb19cd4e2c6b78552be263a680edc55ac/ruff-0.15.6-py3-none-manylinux_2_17_i686.manylinux2014_i686.whl", hash = "sha256:70789d3e7830b848b548aae96766431c0dc01a6c78c13381f423bf7076c66d15", size = 11170783, upload-time = "2026-03-12T23:06:01.742Z" },
{ url = "https://files.pythonhosted.org/packages/5a/f7/4953e7e3287676f78fbe85e3a0ca414c5ca81237b7575bdadc00229ac240/ruff-0.15.5-py3-none-manylinux_2_17_ppc64le.manylinux2014_ppc64le.whl", hash = "sha256:8dc18f30302e379fe1e998548b0f5e9f4dff907f52f73ad6da419ea9c19d66c8", size = 11914251, upload-time = "2026-03-05T20:06:22.887Z" }, { url = "https://files.pythonhosted.org/packages/76/be/6d488f6adca047df82cd62c304638bcb00821c36bd4881cfca221561fdfc/ruff-0.15.6-py3-none-manylinux_2_17_ppc64le.manylinux2014_ppc64le.whl", hash = "sha256:542aaf1de3154cea088ced5a819ce872611256ffe2498e750bbae5247a8114e9", size = 12044201, upload-time = "2026-03-12T23:05:28.697Z" },
{ url = "https://files.pythonhosted.org/packages/77/46/0f7c865c10cf896ccf5a939c3e84e1cfaeed608ff5249584799a74d33835/ruff-0.15.5-py3-none-manylinux_2_17_s390x.manylinux2014_s390x.whl", hash = "sha256:1cc6e7f90087e2d27f98dc34ed1b3ab7c8f0d273cc5431415454e22c0bd2a681", size = 11333801, upload-time = "2026-03-05T20:05:57.168Z" }, { url = "https://files.pythonhosted.org/packages/71/68/e6f125df4af7e6d0b498f8d373274794bc5156b324e8ab4bf5c1b4fc0ec7/ruff-0.15.6-py3-none-manylinux_2_17_s390x.manylinux2014_s390x.whl", hash = "sha256:1c22e6f02c16cfac3888aa636e9eba857254d15bbacc9906c9689fdecb1953ab", size = 11421561, upload-time = "2026-03-12T23:05:31.236Z" },
{ url = "https://files.pythonhosted.org/packages/d3/01/a10fe54b653061585e655f5286c2662ebddb68831ed3eaebfb0eb08c0a16/ruff-0.15.5-py3-none-manylinux_2_17_x86_64.manylinux2014_x86_64.whl", hash = "sha256:c1cb7169f53c1ddb06e71a9aebd7e98fc0fea936b39afb36d8e86d36ecc2636a", size = 11206821, upload-time = "2026-03-05T20:06:03.441Z" }, { url = "https://files.pythonhosted.org/packages/f1/9f/f85ef5fd01a52e0b472b26dc1b4bd228b8f6f0435975442ffa4741278703/ruff-0.15.6-py3-none-manylinux_2_17_x86_64.manylinux2014_x86_64.whl", hash = "sha256:98893c4c0aadc8e448cfa315bd0cc343a5323d740fe5f28ef8a3f9e21b381f7e", size = 11310928, upload-time = "2026-03-12T23:05:45.288Z" },
{ url = "https://files.pythonhosted.org/packages/7a/0d/2132ceaf20c5e8699aa83da2706ecb5c5dcdf78b453f77edca7fb70f8a93/ruff-0.15.5-py3-none-manylinux_2_31_riscv64.whl", hash = "sha256:9b037924500a31ee17389b5c8c4d88874cc6ea8e42f12e9c61a3d754ff72f1ca", size = 11133326, upload-time = "2026-03-05T20:06:25.655Z" }, { url = "https://files.pythonhosted.org/packages/8c/26/b75f8c421f5654304b89471ed384ae8c7f42b4dff58fa6ce1626d7f2b59a/ruff-0.15.6-py3-none-manylinux_2_31_riscv64.whl", hash = "sha256:70d263770d234912374493e8cc1e7385c5d49376e41dfa51c5c3453169dc581c", size = 11235186, upload-time = "2026-03-12T23:05:50.677Z" },
{ url = "https://files.pythonhosted.org/packages/72/cb/2e5259a7eb2a0f87c08c0fe5bf5825a1e4b90883a52685524596bfc93072/ruff-0.15.5-py3-none-musllinux_1_2_aarch64.whl", hash = "sha256:65bb414e5b4eadd95a8c1e4804f6772bbe8995889f203a01f77ddf2d790929dd", size = 10510820, upload-time = "2026-03-05T20:06:37.79Z" }, { url = "https://files.pythonhosted.org/packages/fc/d4/d5a6d065962ff7a68a86c9b4f5500f7d101a0792078de636526c0edd40da/ruff-0.15.6-py3-none-musllinux_1_2_aarch64.whl", hash = "sha256:55a1ad63c5a6e54b1f21b7514dfadc0c7fb40093fa22e95143cf3f64ebdcd512", size = 10635231, upload-time = "2026-03-12T23:05:37.044Z" },
{ url = "https://files.pythonhosted.org/packages/ff/20/b67ce78f9e6c59ffbdb5b4503d0090e749b5f2d31b599b554698a80d861c/ruff-0.15.5-py3-none-musllinux_1_2_armv7l.whl", hash = "sha256:d20aa469ae3b57033519c559e9bc9cd9e782842e39be05b50e852c7c981fa01d", size = 10302395, upload-time = "2026-03-05T20:05:54.504Z" }, { url = "https://files.pythonhosted.org/packages/d6/56/7c3acf3d50910375349016cf33de24be021532042afbed87942858992491/ruff-0.15.6-py3-none-musllinux_1_2_armv7l.whl", hash = "sha256:8dc473ba093c5ec238bb1e7429ee676dca24643c471e11fbaa8a857925b061c0", size = 10340357, upload-time = "2026-03-12T23:06:04.748Z" },
{ url = "https://files.pythonhosted.org/packages/5f/e5/719f1acccd31b720d477751558ed74e9c88134adcc377e5e886af89d3072/ruff-0.15.5-py3-none-musllinux_1_2_i686.whl", hash = "sha256:15388dd28c9161cdb8eda68993533acc870aa4e646a0a277aa166de9ad5a8752", size = 10754069, upload-time = "2026-03-05T20:06:06.422Z" }, { url = "https://files.pythonhosted.org/packages/06/54/6faa39e9c1033ff6a3b6e76b5df536931cd30caf64988e112bbf91ef5ce5/ruff-0.15.6-py3-none-musllinux_1_2_i686.whl", hash = "sha256:85b042377c2a5561131767974617006f99f7e13c63c111b998f29fc1e58a4cfb", size = 10860583, upload-time = "2026-03-12T23:05:58.978Z" },
{ url = "https://files.pythonhosted.org/packages/c3/9c/d1db14469e32d98f3ca27079dbd30b7b44dbb5317d06ab36718dee3baf03/ruff-0.15.5-py3-none-musllinux_1_2_x86_64.whl", hash = "sha256:b30da330cbd03bed0c21420b6b953158f60c74c54c5f4c1dabbdf3a57bf355d2", size = 11304315, upload-time = "2026-03-05T20:06:10.867Z" }, { url = "https://files.pythonhosted.org/packages/cb/1e/509a201b843b4dfb0b32acdedf68d951d3377988cae43949ba4c4133a96a/ruff-0.15.6-py3-none-musllinux_1_2_x86_64.whl", hash = "sha256:cef49e30bc5a86a6a92098a7fbf6e467a234d90b63305d6f3ec01225a9d092e0", size = 11410976, upload-time = "2026-03-12T23:05:39.955Z" },
{ url = "https://files.pythonhosted.org/packages/28/3a/950367aee7c69027f4f422059227b290ed780366b6aecee5de5039d50fa8/ruff-0.15.5-py3-none-win32.whl", hash = "sha256:732e5ee1f98ba5b3679029989a06ca39a950cced52143a0ea82a2102cb592b74", size = 10551676, upload-time = "2026-03-05T20:06:13.705Z" }, { url = "https://files.pythonhosted.org/packages/6c/25/3fc9114abf979a41673ce877c08016f8e660ad6cf508c3957f537d2e9fa9/ruff-0.15.6-py3-none-win32.whl", hash = "sha256:bbf67d39832404812a2d23020dda68fee7f18ce15654e96fb1d3ad21a5fe436c", size = 10616872, upload-time = "2026-03-12T23:05:42.451Z" },
{ url = "https://files.pythonhosted.org/packages/b8/00/bf077a505b4e649bdd3c47ff8ec967735ce2544c8e4a43aba42ee9bf935d/ruff-0.15.5-py3-none-win_amd64.whl", hash = "sha256:821d41c5fa9e19117616c35eaa3f4b75046ec76c65e7ae20a333e9a8696bc7fe", size = 11678972, upload-time = "2026-03-05T20:06:45.379Z" }, { url = "https://files.pythonhosted.org/packages/89/7a/09ece68445ceac348df06e08bf75db72d0e8427765b96c9c0ffabc1be1d9/ruff-0.15.6-py3-none-win_amd64.whl", hash = "sha256:aee25bc84c2f1007ecb5037dff75cef00414fdf17c23f07dc13e577883dca406", size = 11787271, upload-time = "2026-03-12T23:05:20.168Z" },
{ url = "https://files.pythonhosted.org/packages/fe/4e/cd76eca6db6115604b7626668e891c9dd03330384082e33662fb0f113614/ruff-0.15.5-py3-none-win_arm64.whl", hash = "sha256:b498d1c60d2fe5c10c45ec3f698901065772730b411f164ae270bb6bfcc4740b", size = 10965572, upload-time = "2026-03-05T20:06:16.984Z" }, { url = "https://files.pythonhosted.org/packages/7f/d0/578c47dd68152ddddddf31cd7fc67dc30b7cdf639a86275fda821b0d9d98/ruff-0.15.6-py3-none-win_arm64.whl", hash = "sha256:c34de3dd0b0ba203be50ae70f5910b17188556630e2178fd7d79fc030eb0d837", size = 11060497, upload-time = "2026-03-12T23:05:25.968Z" },
] ]
[[package]] [[package]]