Compare commits
3
Commits
v1.1.1
..
9be932c0c7
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
9be932c0c7
|
||
|
|
295a8c90a9
|
||
|
|
653c27c97c
|
+4
-4
@@ -1,4 +1,4 @@
|
||||
# Health check endpoint URL.
|
||||
# If configured, a GET request will be sent to this URL after each successful update cycle.
|
||||
# Leave empty to disable health check reporting.
|
||||
health_check_endpoint: ""
|
||||
# Prometheus metrics endpoint.
|
||||
# When enabled, a /metrics endpoint is exposed via the maubot webapp.
|
||||
# Disabled by default.
|
||||
metrics_enabled: false
|
||||
|
||||
@@ -8,5 +8,8 @@ main_class: OwncastSentry
|
||||
database: true
|
||||
database_type: asyncpg
|
||||
config: true
|
||||
webapp: true
|
||||
dependencies:
|
||||
- prometheus_client>=0.24.1
|
||||
extra_files:
|
||||
- base-config.yaml
|
||||
|
||||
+61
-21
@@ -16,13 +16,15 @@
|
||||
|
||||
from typing import TYPE_CHECKING
|
||||
|
||||
from aiohttp import web
|
||||
from maubot import MessageEvent, Plugin # type: ignore[attr-defined]
|
||||
from maubot.handlers import command
|
||||
from prometheus_client.exposition import choose_encoder
|
||||
|
||||
from .commands import CommandHandler
|
||||
from .config import Config
|
||||
from .database import StreamRepository, SubscriptionRepository
|
||||
from .health_checker import HealthChecker
|
||||
from .metrics import ErrorSource, MetricsService
|
||||
from .migrations import get_upgrade_table
|
||||
from .notification_service import NotificationService
|
||||
from .owncast_client import OwncastClient
|
||||
@@ -57,8 +59,21 @@ class OwncastSentry(Plugin):
|
||||
config.load_and_update()
|
||||
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
|
||||
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
|
||||
self.stream_repo = StreamRepository(db)
|
||||
@@ -66,7 +81,10 @@ class OwncastSentry(Plugin):
|
||||
|
||||
# Initialize notification service
|
||||
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
|
||||
@@ -76,13 +94,7 @@ class OwncastSentry(Plugin):
|
||||
self.subscription_repo,
|
||||
self.notification_service,
|
||||
self.log,
|
||||
)
|
||||
|
||||
# Initialize health checker
|
||||
self.health_checker = HealthChecker(
|
||||
db,
|
||||
self.owncast_client,
|
||||
self.log,
|
||||
metrics=self.metrics_service,
|
||||
)
|
||||
|
||||
# Initialize command handler
|
||||
@@ -97,41 +109,69 @@ class OwncastSentry(Plugin):
|
||||
self.sched.run_periodically(60, self._update_all_stream_states)
|
||||
|
||||
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
|
||||
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,
|
||||
subscribed_domains = (
|
||||
await self.subscription_repo.get_all_subscribed_domains()
|
||||
)
|
||||
|
||||
# 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.argument("url")
|
||||
async def subscribe(self, evt: MessageEvent, url: str) -> None:
|
||||
"""Delegate subscribe command to CommandHandler."""
|
||||
try:
|
||||
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.argument("url")
|
||||
async def unsubscribe(self, evt: MessageEvent, url: str) -> None:
|
||||
"""Delegate unsubscribe command to CommandHandler."""
|
||||
try:
|
||||
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.")
|
||||
async def subscriptions(self, evt: MessageEvent) -> None:
|
||||
"""Delegate subscriptions command to CommandHandler."""
|
||||
try:
|
||||
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.")
|
||||
async def live(self, evt: MessageEvent) -> None:
|
||||
"""Delegate live command to CommandHandler."""
|
||||
try:
|
||||
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:
|
||||
"""Clean up resources by closing the HTTP session."""
|
||||
|
||||
+32
-40
@@ -83,9 +83,7 @@ class CommandHandler:
|
||||
except sqlite3.IntegrityError:
|
||||
# Room is already subscribed.
|
||||
await evt.reply(
|
||||
"This room is already subscribed to notifications for "
|
||||
+ stream_domain
|
||||
+ "."
|
||||
f"This room is already subscribed to notifications for {stream_domain}."
|
||||
)
|
||||
return
|
||||
|
||||
@@ -101,9 +99,8 @@ class CommandHandler:
|
||||
# All went well! Tell the user.
|
||||
self.log.info(f"[{stream_domain}] Subscription added for room {evt.room_id}.")
|
||||
await evt.reply(
|
||||
"Subscription added! This room will receive notifications when "
|
||||
+ stream_domain
|
||||
+ " goes live."
|
||||
f"Subscription added! This room will receive "
|
||||
f"notifications when {stream_domain} goes live."
|
||||
)
|
||||
|
||||
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}."
|
||||
)
|
||||
await evt.reply(
|
||||
"Subscription removed! This room will no "
|
||||
"longer receive notifications for " + stream_domain + "."
|
||||
f"Subscription removed! This room will no "
|
||||
f"longer receive notifications for {stream_domain}."
|
||||
)
|
||||
else:
|
||||
# No, nothing changed. Tell the user.
|
||||
await evt.reply(
|
||||
"This room is already not subscribed to notifications for "
|
||||
+ stream_domain
|
||||
+ "."
|
||||
"This room is already not subscribed to "
|
||||
f"notifications for {stream_domain}."
|
||||
)
|
||||
|
||||
def _format_duration(self, timestamp_str: str) -> str:
|
||||
@@ -143,7 +139,7 @@ class CommandHandler:
|
||||
:return: Formatted duration string (e.g., "1 hour", "2 days").
|
||||
"""
|
||||
try:
|
||||
timestamp = datetime.fromisoformat(timestamp_str.replace("Z", "+00:00"))
|
||||
timestamp = datetime.fromisoformat(timestamp_str)
|
||||
now = datetime.now(UTC)
|
||||
delta = now - timestamp
|
||||
|
||||
@@ -158,7 +154,7 @@ class CommandHandler:
|
||||
return f"{hours} hour{'s' if hours != 1 else ''}"
|
||||
days = seconds // 86400
|
||||
return f"{days} day{'s' if days != 1 else ''}"
|
||||
except Exception:
|
||||
except ValueError:
|
||||
return "unknown duration"
|
||||
|
||||
async def subscriptions(self, evt: MessageEvent) -> None:
|
||||
@@ -183,7 +179,7 @@ class CommandHandler:
|
||||
|
||||
# Build the response message body as Markdown
|
||||
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:
|
||||
# Get the stream state from the database
|
||||
@@ -192,45 +188,41 @@ class CommandHandler:
|
||||
continue
|
||||
|
||||
# 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)
|
||||
|
||||
# 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)
|
||||
if stream_state.status == StreamStatus.ONLINE and 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)
|
||||
if stream_state.status == StreamStatus.ONLINE:
|
||||
# Stream is online - use last_connect_time
|
||||
if stream_state.last_connect_time:
|
||||
match stream_state.status:
|
||||
case StreamStatus.ONLINE if stream_state.last_connect_time:
|
||||
duration = self._format_duration(stream_state.last_connect_time)
|
||||
body_text += f" - Status: Online for {duration} \n"
|
||||
elif stream_state.status == StreamStatus.UNKNOWN:
|
||||
# Stream status is unknown - instance unreachable
|
||||
body_text += " - Status: Unknown (instance unreachable) \n"
|
||||
else:
|
||||
# Stream is offline - use last_disconnect_time
|
||||
if stream_state.last_disconnect_time:
|
||||
parts.append(f" - Status: Online for {duration} \n")
|
||||
case StreamStatus.UNKNOWN:
|
||||
parts.append(" - Status: Unknown (instance unreachable) \n")
|
||||
case StreamStatus.OFFLINE if stream_state.last_disconnect_time:
|
||||
duration = self._format_duration(stream_state.last_disconnect_time)
|
||||
body_text += f" - Status: Offline for {duration} \n"
|
||||
else:
|
||||
body_text += " - Status: Offline \n"
|
||||
parts.append(f" - Status: Offline for {duration} \n")
|
||||
case StreamStatus.OFFLINE:
|
||||
parts.append(" - Status: Offline \n")
|
||||
|
||||
# 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
|
||||
body_text += (
|
||||
parts.append(
|
||||
"\nTo unsubscribe from any of these Owncast "
|
||||
"instances, use `!unsubscribe <domain>`"
|
||||
)
|
||||
|
||||
# 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:
|
||||
"""List currently live streams in the current room.
|
||||
@@ -271,28 +263,28 @@ class CommandHandler:
|
||||
|
||||
# Build the response message body as Markdown
|
||||
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:
|
||||
# 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)
|
||||
|
||||
# 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)
|
||||
if 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
|
||||
if 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
|
||||
body_text += f" - Link: https://{domain}\n\n"
|
||||
parts.append(f" - Link: https://{domain}\n\n")
|
||||
|
||||
# Send the response as Markdown
|
||||
await evt.reply(body_text.rstrip(), markdown=True)
|
||||
await evt.reply("".join(parts).rstrip(), markdown=True)
|
||||
|
||||
@@ -25,9 +25,9 @@ class Config(BaseProxyConfig):
|
||||
|
||||
:param helper: ConfigUpdateHelper for copying values.
|
||||
"""
|
||||
helper.copy("health_check_endpoint")
|
||||
helper.copy("metrics_enabled")
|
||||
|
||||
@property
|
||||
def health_check_endpoint(self) -> str:
|
||||
"""Return the configured health check endpoint URL."""
|
||||
return self["health_check_endpoint"] # type: ignore[no-any-return]
|
||||
def metrics_enabled(self) -> bool:
|
||||
"""Return whether the Prometheus metrics endpoint is enabled."""
|
||||
return self["metrics_enabled"] # type: ignore[no-any-return]
|
||||
|
||||
@@ -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}")
|
||||
@@ -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)
|
||||
@@ -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
|
||||
class StreamConfig:
|
||||
"""Represents the configuration of an Owncast stream."""
|
||||
|
||||
@@ -20,9 +20,10 @@ from typing import TYPE_CHECKING, Any
|
||||
|
||||
from mautrix.types import MessageType, TextMessageEventContent
|
||||
|
||||
from .metrics import NotificationType
|
||||
from .utils import (
|
||||
CLEANUP_DELETE_THRESHOLD,
|
||||
CLEANUP_WARNING_THRESHOLD,
|
||||
CLEANUP_DELETE_DAYS,
|
||||
CLEANUP_WARNING_DAYS,
|
||||
SECONDS_BETWEEN_NOTIFICATIONS,
|
||||
sanitize_for_plain_text,
|
||||
)
|
||||
@@ -31,6 +32,7 @@ if TYPE_CHECKING:
|
||||
import logging
|
||||
|
||||
from .database import SubscriptionRepository
|
||||
from .metrics import MetricsService
|
||||
|
||||
|
||||
class NotificationService:
|
||||
@@ -41,16 +43,19 @@ class NotificationService:
|
||||
client: Any,
|
||||
subscription_repo: SubscriptionRepository,
|
||||
logger: logging.Logger,
|
||||
metrics: MetricsService,
|
||||
) -> None:
|
||||
"""Initialize the notification service.
|
||||
|
||||
:param client: The Matrix client for sending messages.
|
||||
:param subscription_repo: Repository for managing subscriptions.
|
||||
:param logger: Logger instance for debugging.
|
||||
:param metrics: Metrics service for recording counters.
|
||||
"""
|
||||
self.client = client
|
||||
self.subscription_repo = subscription_repo
|
||||
self.log = logger
|
||||
self.metrics = metrics
|
||||
|
||||
# Cache for tracking when notifications were last sent
|
||||
self.notification_timers_cache: dict[str, float] = {}
|
||||
@@ -88,36 +93,24 @@ class NotificationService:
|
||||
# Record that we're sending a notification now
|
||||
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
|
||||
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
|
||||
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
|
||||
for result in results:
|
||||
if isinstance(result, Exception):
|
||||
failed_notifications += 1
|
||||
else:
|
||||
successful_notifications += 1
|
||||
successful, failed = await self._broadcast_to_rooms(domain, body_text)
|
||||
|
||||
# Log completion
|
||||
notification_type = "title change" if title_change else "going live"
|
||||
self.log.info(
|
||||
f"[{domain}] Completed sending {notification_type} "
|
||||
f"notifications! {successful_notifications} succeeded, "
|
||||
f"{failed_notifications} failed."
|
||||
f"notifications! {successful} succeeded, "
|
||||
f"{failed} failed."
|
||||
)
|
||||
|
||||
self.metrics.record_delivery(
|
||||
NotificationType.TITLE_CHANGE if title_change else NotificationType.LIVE,
|
||||
successful=successful,
|
||||
failed=failed,
|
||||
)
|
||||
|
||||
async def _send_notification(
|
||||
@@ -153,36 +146,36 @@ class NotificationService:
|
||||
:return: Formatted message body.
|
||||
"""
|
||||
# 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)
|
||||
|
||||
# Choose message based on notification type
|
||||
if title_change:
|
||||
body_text = "📝 " + safe_stream_name + " has changed its stream title!"
|
||||
parts = [f"📝 {safe_stream_name} has changed its stream title!"]
|
||||
else:
|
||||
body_text = "🎥 " + safe_stream_name + " is now live!"
|
||||
parts = [f"🎥 {safe_stream_name} is now live!"]
|
||||
|
||||
# Add title if present
|
||||
if title != "":
|
||||
if title:
|
||||
safe_title = sanitize_for_plain_text(title)
|
||||
body_text += "\nStream Title: " + safe_title
|
||||
parts.append(f"\nStream Title: {safe_title}")
|
||||
|
||||
# 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
|
||||
if tags:
|
||||
safe_tags = []
|
||||
for tag in tags:
|
||||
safe_tag = sanitize_for_plain_text(tag)
|
||||
if safe_tag and not safe_tag.startswith("."):
|
||||
safe_tags.append(safe_tag)
|
||||
safe_tags = [
|
||||
safe_tag
|
||||
for tag in tags
|
||||
if (safe_tag := sanitize_for_plain_text(tag))
|
||||
and not safe_tag.startswith(".")
|
||||
]
|
||||
|
||||
if safe_tags:
|
||||
body_text += "\n\n"
|
||||
body_text += " ".join("#" + tag for tag in safe_tags)
|
||||
parts.append(f"\n\n{' '.join(f'#{tag}' for tag in safe_tags)}")
|
||||
|
||||
return body_text
|
||||
return "".join(parts)
|
||||
|
||||
def get_last_notification_time(self, domain: str) -> float:
|
||||
"""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:
|
||||
return True
|
||||
|
||||
seconds_since_last = round(time.time() - self.notification_timers_cache[domain])
|
||||
return bool(seconds_since_last >= SECONDS_BETWEEN_NOTIFICATIONS)
|
||||
seconds_since_last = round(
|
||||
time.monotonic() - self.notification_timers_cache[domain]
|
||||
)
|
||||
return seconds_since_last >= SECONDS_BETWEEN_NOTIFICATIONS
|
||||
|
||||
def _record_notification(self, domain: str) -> None:
|
||||
"""Record that a notification was sent at the current time.
|
||||
|
||||
: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:
|
||||
"""Send cleanup warning notification to all subscribed rooms.
|
||||
|
||||
:param domain: The stream domain.
|
||||
"""
|
||||
# Get all subscribed rooms
|
||||
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
|
||||
remaining_days = CLEANUP_DELETE_DAYS - CLEANUP_WARNING_DAYS
|
||||
body_text = (
|
||||
"⚠️ Warning: Subscription Cleanup Scheduled\n\n"
|
||||
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"({delete_days} days total), this subscription "
|
||||
f"({CLEANUP_DELETE_DAYS} days total), this subscription "
|
||||
f"will be automatically removed."
|
||||
)
|
||||
|
||||
# Send to all rooms in parallel
|
||||
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))
|
||||
|
||||
successful, failed = await self._broadcast_to_rooms(domain, body_text)
|
||||
self.log.info(
|
||||
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:
|
||||
"""Send cleanup deletion notification to all subscribed rooms.
|
||||
|
||||
: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 = (
|
||||
"🗑️ Subscription Automatically Removed\n\n"
|
||||
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"room.\n\n"
|
||||
f"If the instance comes online again and you want to "
|
||||
f"resubscribe, run `!subscribe {domain}`."
|
||||
)
|
||||
|
||||
# Send to all rooms in parallel
|
||||
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))
|
||||
|
||||
successful, failed = await self._broadcast_to_rooms(domain, body_text)
|
||||
self.log.info(
|
||||
f"[{domain}] Sent cleanup deletion notice to "
|
||||
f"{successful} rooms ({failed} failed)."
|
||||
)
|
||||
|
||||
self.metrics.record_delivery(
|
||||
NotificationType.CLEANUP_DELETION, successful=successful, failed=failed
|
||||
)
|
||||
|
||||
@@ -14,28 +14,41 @@
|
||||
|
||||
"""HTTP client for querying Owncast instance APIs."""
|
||||
|
||||
import json
|
||||
from typing import TYPE_CHECKING
|
||||
from typing import TYPE_CHECKING, Any
|
||||
|
||||
import aiohttp
|
||||
|
||||
from .models import StreamConfig, StreamState
|
||||
from .utils import (
|
||||
OWNCAST_CONFIG_PATH,
|
||||
OWNCAST_STATUS_PATH,
|
||||
REQUIRED_STATUS_FIELDS,
|
||||
user_agent,
|
||||
)
|
||||
|
||||
if TYPE_CHECKING:
|
||||
import logging
|
||||
|
||||
from .models import StreamConfig, StreamState
|
||||
from .utils import OWNCAST_CONFIG_PATH, OWNCAST_STATUS_PATH, user_agent
|
||||
from .metrics import MetricsService
|
||||
|
||||
|
||||
class OwncastClient:
|
||||
"""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.
|
||||
|
||||
:param logger: Logger instance for debugging
|
||||
:param version: Plugin version string for the User-Agent header
|
||||
:param metrics: Metrics service for recording response times.
|
||||
"""
|
||||
self.log = logger
|
||||
self.metrics = metrics
|
||||
|
||||
# Set up HTTP session configuration
|
||||
headers = {"User-Agent": user_agent(version)}
|
||||
@@ -55,6 +68,37 @@ class OwncastClient:
|
||||
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:
|
||||
"""Get the current stream state for a given domain.
|
||||
|
||||
@@ -65,54 +109,23 @@ class OwncastClient:
|
||||
:return: A StreamState if available, None on error.
|
||||
"""
|
||||
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
|
||||
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}"
|
||||
)
|
||||
if new_state is None:
|
||||
return None
|
||||
|
||||
# Validate the response contains all basic info needed
|
||||
required_fields = [
|
||||
"lastConnectTime",
|
||||
"lastDisconnectTime",
|
||||
"streamTitle",
|
||||
"online",
|
||||
]
|
||||
for field in required_fields:
|
||||
if field not in new_state:
|
||||
missing = REQUIRED_STATUS_FIELDS - new_state.keys()
|
||||
if missing:
|
||||
self.log.warning(
|
||||
f"[{domain}] Rejecting response to request "
|
||||
f"on {OWNCAST_STATUS_PATH} as it does not "
|
||||
f"have {field} field."
|
||||
f"[{domain}] Rejecting response to request on "
|
||||
f"{OWNCAST_STATUS_PATH} as it is missing "
|
||||
f"fields: {', '.join(sorted(missing))}"
|
||||
)
|
||||
return None
|
||||
|
||||
timer.success()
|
||||
return StreamState.from_api_response(new_state, domain)
|
||||
|
||||
async def get_stream_config(self, domain: str) -> StreamConfig | None:
|
||||
@@ -125,39 +138,12 @@ class OwncastClient:
|
||||
:return: A StreamConfig, or None if fetch failed.
|
||||
"""
|
||||
self.log.debug(f"[{domain}] Fetching current stream config...")
|
||||
config_url = "https://" + domain + OWNCAST_CONFIG_PATH
|
||||
|
||||
# Make a request to the endpoint
|
||||
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."
|
||||
)
|
||||
with self.metrics.response_timer(domain) as timer:
|
||||
config = await self._fetch_json(domain, OWNCAST_CONFIG_PATH)
|
||||
if config is None:
|
||||
return None
|
||||
|
||||
# Try to interpret the response as JSON
|
||||
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)
|
||||
timer.success()
|
||||
return StreamConfig.from_api_response(config)
|
||||
|
||||
async def validate_instance(self, domain: str) -> bool:
|
||||
@@ -169,6 +155,16 @@ class OwncastClient:
|
||||
state = await self.get_stream_state(domain)
|
||||
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:
|
||||
"""Close the HTTP session."""
|
||||
await self.session.close()
|
||||
|
||||
@@ -18,8 +18,7 @@ import asyncio
|
||||
import time
|
||||
from typing import TYPE_CHECKING
|
||||
|
||||
from .health_checker import UpdateResult
|
||||
from .models import StreamState
|
||||
from .models import StreamState, StreamStatus, UpdateResult
|
||||
from .utils import (
|
||||
CLEANUP_DELETE_THRESHOLD,
|
||||
CLEANUP_WARNING_THRESHOLD,
|
||||
@@ -31,6 +30,7 @@ if TYPE_CHECKING:
|
||||
import logging
|
||||
|
||||
from .database import StreamRepository, SubscriptionRepository
|
||||
from .metrics import MetricsService
|
||||
from .notification_service import NotificationService
|
||||
from .owncast_client import OwncastClient
|
||||
|
||||
@@ -45,6 +45,7 @@ class StreamMonitor:
|
||||
subscription_repo: SubscriptionRepository,
|
||||
notification_service: NotificationService,
|
||||
logger: logging.Logger,
|
||||
metrics: MetricsService,
|
||||
) -> None:
|
||||
"""Initialize the stream monitor.
|
||||
|
||||
@@ -53,12 +54,14 @@ class StreamMonitor:
|
||||
:param subscription_repo: Repository for subscription data.
|
||||
:param notification_service: Service for sending notifications.
|
||||
:param logger: Logger instance for debugging.
|
||||
:param metrics: Metrics service for recording Prometheus metrics.
|
||||
"""
|
||||
self.owncast_client = owncast_client
|
||||
self.stream_repo = stream_repo
|
||||
self.subscription_repo = subscription_repo
|
||||
self.notification_service = notification_service
|
||||
self.log = logger
|
||||
self.metrics = metrics
|
||||
|
||||
# Cache for tracking when streams last went offline
|
||||
self.offline_timer_cache: dict[str, float] = {}
|
||||
@@ -77,23 +80,24 @@ class StreamMonitor:
|
||||
total_streams = len(subscribed_domains)
|
||||
|
||||
# Build a list of async tasks for each stream domain
|
||||
tasks = [
|
||||
asyncio.create_task(self.update_stream(domain))
|
||||
for domain in subscribed_domains
|
||||
]
|
||||
|
||||
# Run the tasks in parallel and collect results
|
||||
results = await asyncio.gather(*tasks)
|
||||
# Run all stream updates in parallel and collect results
|
||||
results = await asyncio.gather(
|
||||
*(self.update_stream(domain) for domain in subscribed_domains)
|
||||
)
|
||||
|
||||
# Count successes and failures
|
||||
successful_checks = sum(1 for result in results if result is True)
|
||||
failed_checks = sum(1 for result in results if result is False)
|
||||
successful_checks = results.count(True)
|
||||
failed_checks = results.count(False)
|
||||
|
||||
self.log.debug(
|
||||
f"Update complete. {successful_checks}/{total_streams} succeeded, "
|
||||
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(
|
||||
total_streams=total_streams,
|
||||
successful_checks=successful_checks,
|
||||
@@ -127,6 +131,10 @@ class StreamMonitor:
|
||||
)
|
||||
# Check cleanup thresholds even when skipping query
|
||||
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
|
||||
return True
|
||||
|
||||
@@ -152,15 +160,19 @@ class StreamMonitor:
|
||||
)
|
||||
# Check cleanup thresholds after connection failure
|
||||
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
|
||||
return False
|
||||
|
||||
# Fetch succeeded! Reset failure counter
|
||||
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
|
||||
if domain not in self.offline_timer_cache:
|
||||
self.offline_timer_cache[domain] = 0
|
||||
self.offline_timer_cache.setdefault(domain, 0)
|
||||
|
||||
# Does the last known stream state lack connect/disconnect?
|
||||
if (
|
||||
@@ -184,7 +196,7 @@ class StreamMonitor:
|
||||
|
||||
# Calculate seconds since the stream last went offline
|
||||
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?
|
||||
@@ -278,7 +290,7 @@ class StreamMonitor:
|
||||
):
|
||||
# Yep. This stream is now offline. Log it.
|
||||
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.")
|
||||
|
||||
# Update the database with current stream state, if needed.
|
||||
@@ -305,6 +317,10 @@ class StreamMonitor:
|
||||
|
||||
# All done.
|
||||
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
|
||||
|
||||
async def _check_cleanup_thresholds(self, domain: str, counter: int) -> None:
|
||||
@@ -340,3 +356,4 @@ class StreamMonitor:
|
||||
f"Deleted {deleted_count} subscriptions "
|
||||
f"and stream record."
|
||||
)
|
||||
self.metrics.remove_stream(domain)
|
||||
|
||||
+23
-44
@@ -23,6 +23,16 @@ OWNCAST_STATUS_PATH = "/api/status"
|
||||
# Path to GetWebConfig API call on Owncast instances
|
||||
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:
|
||||
"""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.
|
||||
TEMPORARY_OFFLINE_NOTIFICATION_COOLDOWN = 7 * 60 # 7 min in seconds
|
||||
|
||||
# Counter thresholds for auto-cleanup (60-second polling intervals)
|
||||
CLEANUP_WARNING_THRESHOLD = 83 * 24 * 60 # 119,520 cycles = 83 days
|
||||
CLEANUP_DELETE_THRESHOLD = 90 * 24 * 60 # 129,600 cycles = 90 days
|
||||
# Auto-cleanup timing (days of continuous unreachability)
|
||||
CLEANUP_WARNING_DAYS = 83
|
||||
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"
|
||||
UNKNOWN_STATUS_THRESHOLD = 15
|
||||
@@ -109,7 +123,7 @@ def domainify(url: str) -> str:
|
||||
|
||||
# Prepend // if no scheme so urlparse treats input as netloc
|
||||
if not url.startswith(("http://", "https://", "//")):
|
||||
url = "//" + url
|
||||
url = f"//{url}"
|
||||
|
||||
parsed = urlparse(url)
|
||||
domain = (parsed.netloc or parsed.path).lower()
|
||||
@@ -133,6 +147,9 @@ def truncate(text: str, max_length: int) -> str:
|
||||
return text[:max_length]
|
||||
|
||||
|
||||
_MARKDOWN_ESCAPE_TABLE = str.maketrans({c: f"\\{c}" for c in r"\*_[]()~`#+-=|{}.!<>&"})
|
||||
|
||||
|
||||
def escape_markdown(text: str) -> str:
|
||||
"""Escape Markdown special characters to prevent injection attacks.
|
||||
|
||||
@@ -146,38 +163,7 @@ def escape_markdown(text: str) -> str:
|
||||
if not text:
|
||||
return text
|
||||
|
||||
# Escape Markdown special characters by prefixing with backslash
|
||||
# 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
|
||||
return text.translate(_MARKDOWN_ESCAPE_TABLE)
|
||||
|
||||
|
||||
def sanitize_for_plain_text(text: str) -> str:
|
||||
@@ -216,11 +202,4 @@ def sanitize_for_markdown(text: str) -> str:
|
||||
if not text:
|
||||
return text
|
||||
|
||||
# Remove newlines and carriage returns to prevent multi-line injection
|
||||
sanitized = text.replace("\n", " ").replace("\r", " ")
|
||||
|
||||
# Collapse multiple spaces into single space
|
||||
sanitized = " ".join(sanitized.split())
|
||||
|
||||
# Escape Markdown special characters
|
||||
return escape_markdown(sanitized)
|
||||
return escape_markdown(sanitize_for_plain_text(text))
|
||||
|
||||
+3
-3
@@ -7,13 +7,13 @@ authors = [
|
||||
]
|
||||
license = "Apache-2.0"
|
||||
requires-python = ">=3.14"
|
||||
dependencies = []
|
||||
dependencies = ["prometheus_client>=0.24.1"]
|
||||
|
||||
[project.urls]
|
||||
Repository = "https://git.logal.dev/LogalDeveloper/OwncastSentry"
|
||||
|
||||
[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"
|
||||
|
||||
[dependency-groups]
|
||||
@@ -27,7 +27,7 @@ dev = [
|
||||
"pytest>=9.0.2",
|
||||
"pytest-asyncio>=1.3.0",
|
||||
"pytest-cov>=7.0.0",
|
||||
"ruff>=0.15.5",
|
||||
"ruff>=0.15.6",
|
||||
"time-machine>=3.2.0",
|
||||
]
|
||||
|
||||
|
||||
+7
-2
@@ -14,13 +14,12 @@
|
||||
|
||||
"""Shared test fixtures and stubs for OwncastSentry tests."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from dataclasses import dataclass, field
|
||||
from typing import TYPE_CHECKING, Any
|
||||
|
||||
import pytest
|
||||
from mautrix.util.async_db import Database
|
||||
from prometheus_client import generate_latest
|
||||
|
||||
from owncastsentry import OwncastSentry
|
||||
from owncastsentry.config import Config
|
||||
@@ -31,9 +30,15 @@ if TYPE_CHECKING:
|
||||
from collections.abc import AsyncIterator
|
||||
from pathlib import Path
|
||||
|
||||
from owncastsentry.metrics import MetricsService
|
||||
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
|
||||
async def database(tmp_path: Path) -> AsyncIterator[Database]:
|
||||
"""Yield a real SQLite-backed mautrix Database with migrations applied."""
|
||||
|
||||
@@ -14,8 +14,6 @@
|
||||
|
||||
"""Tests for bot command handlers."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
import logging
|
||||
from datetime import UTC, datetime, timedelta
|
||||
|
||||
@@ -14,8 +14,6 @@
|
||||
|
||||
"""Tests for database repository classes."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from typing import TYPE_CHECKING
|
||||
|
||||
if TYPE_CHECKING:
|
||||
|
||||
@@ -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")
|
||||
@@ -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
|
||||
@@ -14,8 +14,6 @@
|
||||
|
||||
"""Tests for data models."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import pytest
|
||||
|
||||
from owncastsentry.models import StreamConfig, StreamState, StreamStatus
|
||||
|
||||
@@ -14,17 +14,16 @@
|
||||
|
||||
"""Tests for the notification service."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import logging
|
||||
import time
|
||||
from typing import TYPE_CHECKING
|
||||
|
||||
import pytest
|
||||
|
||||
from owncastsentry.metrics import MetricsService
|
||||
from owncastsentry.notification_service import NotificationService
|
||||
from owncastsentry.utils import SECONDS_BETWEEN_NOTIFICATIONS
|
||||
from tests.conftest import _StubMatrixClient
|
||||
from tests.conftest import _StubMatrixClient, generate_metrics_output
|
||||
|
||||
if TYPE_CHECKING:
|
||||
from owncastsentry.database import StreamRepository, SubscriptionRepository
|
||||
@@ -34,12 +33,14 @@ def _make_service(
|
||||
*,
|
||||
client: _StubMatrixClient,
|
||||
subscription_repo: SubscriptionRepository,
|
||||
metrics: MetricsService | None = None,
|
||||
) -> NotificationService:
|
||||
"""Build a NotificationService with a stub client and real repo."""
|
||||
return NotificationService(
|
||||
client=client,
|
||||
subscription_repo=subscription_repo,
|
||||
logger=logging.getLogger("test"),
|
||||
metrics=metrics or MetricsService(),
|
||||
)
|
||||
|
||||
|
||||
@@ -62,7 +63,7 @@ class TestCanNotify:
|
||||
service = _make_service(
|
||||
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
|
||||
|
||||
def test_after_cooldown_allowed(
|
||||
@@ -74,7 +75,7 @@ class TestCanNotify:
|
||||
)
|
||||
# Subtract an extra second to ensure the cooldown has fully elapsed
|
||||
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
|
||||
|
||||
@@ -253,7 +254,7 @@ class TestNotifyStreamLive:
|
||||
"""Skip sending when the domain is within the rate-limit cooldown."""
|
||||
client = _StubMatrixClient()
|
||||
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 subscription_repo.add("example.com", "!room1:matrix.org")
|
||||
@@ -346,3 +347,111 @@ class TestSendCleanupDeletion:
|
||||
"If the instance comes online again and you want to "
|
||||
"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
|
||||
|
||||
@@ -14,8 +14,6 @@
|
||||
|
||||
"""Tests for the Owncast HTTP client."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
import logging
|
||||
from typing import TYPE_CHECKING
|
||||
@@ -23,8 +21,13 @@ from typing import TYPE_CHECKING
|
||||
import pytest
|
||||
from aioresponses import aioresponses
|
||||
|
||||
from owncastsentry.metrics import MetricsService
|
||||
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:
|
||||
from collections.abc import AsyncIterator
|
||||
@@ -33,7 +36,11 @@ if TYPE_CHECKING:
|
||||
@pytest.fixture
|
||||
async def owncast_client() -> AsyncIterator[OwncastClient]:
|
||||
"""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
|
||||
await client.close()
|
||||
|
||||
@@ -205,3 +212,69 @@ class TestValidateInstance:
|
||||
result = await owncast_client.validate_instance("invalid.com")
|
||||
|
||||
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
@@ -14,13 +14,12 @@
|
||||
|
||||
"""Tests for the stream monitor."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import logging
|
||||
import time
|
||||
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.stream_monitor import StreamMonitor
|
||||
from owncastsentry.utils import (
|
||||
@@ -28,7 +27,11 @@ from owncastsentry.utils import (
|
||||
CLEANUP_WARNING_THRESHOLD,
|
||||
SECONDS_BETWEEN_NOTIFICATIONS,
|
||||
)
|
||||
from tests.conftest import _StubMatrixClient, _StubOwncastClient
|
||||
from tests.conftest import (
|
||||
_StubMatrixClient,
|
||||
_StubOwncastClient,
|
||||
generate_metrics_output,
|
||||
)
|
||||
|
||||
if TYPE_CHECKING:
|
||||
from owncastsentry.database import StreamRepository, SubscriptionRepository
|
||||
@@ -40,13 +43,16 @@ def _make_monitor(
|
||||
stream_repo: StreamRepository,
|
||||
subscription_repo: SubscriptionRepository,
|
||||
client: _StubMatrixClient,
|
||||
metrics: MetricsService | None = None,
|
||||
) -> tuple[StreamMonitor, NotificationService]:
|
||||
"""Build a StreamMonitor with stubs and a real NotificationService."""
|
||||
logger = logging.getLogger("test")
|
||||
metrics = metrics or MetricsService()
|
||||
notification_service = NotificationService(
|
||||
client=client,
|
||||
subscription_repo=subscription_repo,
|
||||
logger=logger,
|
||||
metrics=metrics,
|
||||
)
|
||||
monitor = StreamMonitor(
|
||||
owncast_client=owncast_client,
|
||||
@@ -54,6 +60,7 @@ def _make_monitor(
|
||||
subscription_repo=subscription_repo,
|
||||
notification_service=notification_service,
|
||||
logger=logger,
|
||||
metrics=metrics,
|
||||
)
|
||||
return monitor, notification_service
|
||||
|
||||
@@ -82,6 +89,25 @@ async def _seed_stream(
|
||||
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:
|
||||
"""Parallel stream update orchestration."""
|
||||
|
||||
@@ -337,7 +363,7 @@ class TestUpdateStreamBriefOffline:
|
||||
)
|
||||
|
||||
# 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")
|
||||
assert result is True
|
||||
@@ -375,7 +401,7 @@ class TestUpdateStreamBriefOffline:
|
||||
)
|
||||
|
||||
# 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")
|
||||
assert result is True
|
||||
@@ -424,7 +450,7 @@ class TestUpdateStreamTitleChange:
|
||||
monitor.offline_timer_cache["title.com"] = 0
|
||||
# Subtract an extra second to ensure the cooldown has fully elapsed
|
||||
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")
|
||||
@@ -473,7 +499,7 @@ class TestUpdateStreamTitleChange:
|
||||
monitor.offline_timer_cache["title.com"] = 0
|
||||
# Subtract an extra second to ensure the cooldown has fully elapsed
|
||||
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")
|
||||
@@ -518,7 +544,7 @@ class TestUpdateStreamTitleChange:
|
||||
|
||||
# Offline timer is MORE recent than last notification,
|
||||
# and both are old enough to pass rate limiting
|
||||
now = time.time()
|
||||
now = time.monotonic()
|
||||
monitor.offline_timer_cache["title.com"] = (
|
||||
now - SECONDS_BETWEEN_NOTIFICATIONS - 100
|
||||
)
|
||||
@@ -571,9 +597,9 @@ class TestUpdateStreamGoesOffline:
|
||||
)
|
||||
monitor.offline_timer_cache["offline.com"] = 0
|
||||
|
||||
before = time.time()
|
||||
before = time.monotonic()
|
||||
result = await monitor.update_stream("offline.com")
|
||||
after = time.time()
|
||||
after = time.monotonic()
|
||||
|
||||
assert result is True
|
||||
assert before <= monitor.offline_timer_cache["offline.com"] <= after
|
||||
@@ -837,3 +863,178 @@ class TestUpdateAllStreamsMixed:
|
||||
assert result.total_streams == 2
|
||||
assert result.successful_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
|
||||
|
||||
@@ -14,8 +14,6 @@
|
||||
|
||||
"""Tests for utility functions and constants."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import pytest
|
||||
|
||||
from owncastsentry.utils import (
|
||||
|
||||
@@ -924,6 +924,9 @@ wheels = [
|
||||
[[package]]
|
||||
name = "owncastsentry"
|
||||
source = { editable = "." }
|
||||
dependencies = [
|
||||
{ name = "prometheus-client" },
|
||||
]
|
||||
|
||||
[package.dev-dependencies]
|
||||
dev = [
|
||||
@@ -941,6 +944,7 @@ dev = [
|
||||
]
|
||||
|
||||
[package.metadata]
|
||||
requires-dist = [{ name = "prometheus-client", specifier = ">=0.24.1" }]
|
||||
|
||||
[package.metadata.requires-dev]
|
||||
dev = [
|
||||
@@ -953,7 +957,7 @@ dev = [
|
||||
{ name = "pytest", specifier = ">=9.0.2" },
|
||||
{ name = "pytest-asyncio", specifier = ">=1.3.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" },
|
||||
]
|
||||
|
||||
@@ -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" },
|
||||
]
|
||||
|
||||
[[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]]
|
||||
name = "prompt-toolkit"
|
||||
version = "3.0.52"
|
||||
@@ -1341,27 +1354,27 @@ wheels = [
|
||||
|
||||
[[package]]
|
||||
name = "ruff"
|
||||
version = "0.15.5"
|
||||
version = "0.15.6"
|
||||
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 = [
|
||||
{ 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/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/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/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/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/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/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/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/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/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/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/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/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/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/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/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/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/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/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/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/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/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/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/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/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/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/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/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/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/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/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/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/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/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]]
|
||||
|
||||
Reference in New Issue
Block a user