Replaced health checker with Prometheus metrics service.
Audit / Dependencies (push) Failing after 8s
CD / Build (push) Successful in 8s
CI / Formatting (push) Successful in 6s
CI / Linting (push) Successful in 6s
CI / Tests (push) Successful in 11s
CI / Type Checking (push) Successful in 14s
CI / Spelling (push) Successful in 13s
Audit / Dependencies (push) Failing after 8s
CD / Build (push) Successful in 8s
CI / Formatting (push) Successful in 6s
CI / Linting (push) Successful in 6s
CI / Tests (push) Successful in 11s
CI / Type Checking (push) Successful in 14s
CI / Spelling (push) Successful in 13s
This commit is contained in:
+66
-26
@@ -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."""
|
||||
# Get list of all stream domains with active subscriptions
|
||||
subscribed_domains = await self.subscription_repo.get_all_subscribed_domains()
|
||||
"""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,
|
||||
)
|
||||
# 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."""
|
||||
await self.command_handler.subscribe(evt, url)
|
||||
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."""
|
||||
await self.command_handler.unsubscribe(evt, url)
|
||||
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."""
|
||||
await self.command_handler.subscriptions(evt)
|
||||
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."""
|
||||
await self.command_handler.live(evt)
|
||||
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."""
|
||||
|
||||
@@ -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,141 +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 .models import UpdateResult
|
||||
from .owncast_client import OwncastClient
|
||||
|
||||
|
||||
@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: # broad catch - DB backends raise varied errors
|
||||
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.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)
|
||||
@@ -99,20 +99,6 @@ class UpdateResult:
|
||||
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 StreamConfig:
|
||||
|
||||
@@ -20,6 +20,7 @@ from typing import TYPE_CHECKING, Any
|
||||
|
||||
from mautrix.types import MessageType, TextMessageEventContent
|
||||
|
||||
from .metrics import NotificationType
|
||||
from .utils import (
|
||||
CLEANUP_DELETE_DAYS,
|
||||
CLEANUP_WARNING_DAYS,
|
||||
@@ -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] = {}
|
||||
@@ -102,6 +107,12 @@ class NotificationService:
|
||||
f"{failed} failed."
|
||||
)
|
||||
|
||||
self.metrics.record_delivery(
|
||||
NotificationType.TITLE_CHANGE if title_change else NotificationType.LIVE,
|
||||
successful=successful,
|
||||
failed=failed,
|
||||
)
|
||||
|
||||
async def _send_notification(
|
||||
self, room_id: str, body_text: str, domain: str
|
||||
) -> None:
|
||||
@@ -231,6 +242,10 @@ class NotificationService:
|
||||
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.
|
||||
|
||||
@@ -251,3 +266,7 @@ class NotificationService:
|
||||
f"[{domain}] Sent cleanup deletion notice to "
|
||||
f"{successful} rooms ({failed} failed)."
|
||||
)
|
||||
|
||||
self.metrics.record_delivery(
|
||||
NotificationType.CLEANUP_DELETION, successful=successful, failed=failed
|
||||
)
|
||||
|
||||
@@ -29,17 +29,26 @@ from .utils import (
|
||||
if TYPE_CHECKING:
|
||||
import logging
|
||||
|
||||
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)}
|
||||
@@ -100,21 +109,24 @@ class OwncastClient:
|
||||
:return: A StreamState if available, None on error.
|
||||
"""
|
||||
self.log.debug(f"[{domain}] Fetching current stream state...")
|
||||
new_state = await self._fetch_json(domain, OWNCAST_STATUS_PATH)
|
||||
if new_state is None:
|
||||
return None
|
||||
with self.metrics.response_timer(domain) as timer:
|
||||
new_state = await self._fetch_json(domain, OWNCAST_STATUS_PATH)
|
||||
|
||||
# Validate the response contains all basic info needed
|
||||
missing = REQUIRED_STATUS_FIELDS - new_state.keys()
|
||||
if missing:
|
||||
self.log.warning(
|
||||
f"[{domain}] Rejecting response to request on "
|
||||
f"{OWNCAST_STATUS_PATH} as it is missing "
|
||||
f"fields: {', '.join(sorted(missing))}"
|
||||
)
|
||||
return None
|
||||
if new_state is None:
|
||||
return None
|
||||
|
||||
return StreamState.from_api_response(new_state, domain)
|
||||
# Validate the response contains all basic info needed
|
||||
missing = REQUIRED_STATUS_FIELDS - new_state.keys()
|
||||
if missing:
|
||||
self.log.warning(
|
||||
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:
|
||||
"""Get the current stream config for a given domain.
|
||||
@@ -126,11 +138,13 @@ class OwncastClient:
|
||||
:return: A StreamConfig, or None if fetch failed.
|
||||
"""
|
||||
self.log.debug(f"[{domain}] Fetching current stream config...")
|
||||
config = await self._fetch_json(domain, OWNCAST_CONFIG_PATH)
|
||||
if config is None:
|
||||
return None
|
||||
with self.metrics.response_timer(domain) as timer:
|
||||
config = await self._fetch_json(domain, OWNCAST_CONFIG_PATH)
|
||||
if config is None:
|
||||
return None
|
||||
|
||||
return StreamConfig.from_api_response(config)
|
||||
timer.success()
|
||||
return StreamConfig.from_api_response(config)
|
||||
|
||||
async def validate_instance(self, domain: str) -> bool:
|
||||
"""Validate that a domain is a valid Owncast instance.
|
||||
@@ -141,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,7 +18,7 @@ import asyncio
|
||||
import time
|
||||
from typing import TYPE_CHECKING
|
||||
|
||||
from .models import StreamState, UpdateResult
|
||||
from .models import StreamState, StreamStatus, UpdateResult
|
||||
from .utils import (
|
||||
CLEANUP_DELETE_THRESHOLD,
|
||||
CLEANUP_WARNING_THRESHOLD,
|
||||
@@ -30,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
|
||||
|
||||
@@ -44,6 +45,7 @@ class StreamMonitor:
|
||||
subscription_repo: SubscriptionRepository,
|
||||
notification_service: NotificationService,
|
||||
logger: logging.Logger,
|
||||
metrics: MetricsService,
|
||||
) -> None:
|
||||
"""Initialize the stream monitor.
|
||||
|
||||
@@ -52,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] = {}
|
||||
@@ -90,6 +94,10 @@ class StreamMonitor:
|
||||
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,
|
||||
@@ -123,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
|
||||
|
||||
@@ -148,11 +160,16 @@ 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
|
||||
self.offline_timer_cache.setdefault(domain, 0)
|
||||
@@ -300,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:
|
||||
@@ -335,3 +356,4 @@ class StreamMonitor:
|
||||
f"Deleted {deleted_count} subscriptions "
|
||||
f"and stream record."
|
||||
)
|
||||
self.metrics.remove_stream(domain)
|
||||
|
||||
Reference in New Issue
Block a user