Added comprehensive test suite with CI workflow.
Audit / Dependencies (push) Successful in 8s
CD / Build (push) Successful in 8s
CI / Formatting (push) Successful in 5s
CI / Linting (push) Successful in 5s
CI / Tests (push) Successful in 11s
CI / Type Checking (push) Successful in 9s
CI / Spelling (push) Successful in 6s
Audit / Dependencies (push) Successful in 8s
CD / Build (push) Successful in 8s
CI / Formatting (push) Successful in 5s
CI / Linting (push) Successful in 5s
CI / Tests (push) Successful in 11s
CI / Type Checking (push) Successful in 9s
CI / Spelling (push) Successful in 6s
This commit is contained in:
@@ -0,0 +1,15 @@
|
||||
# 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.
|
||||
|
||||
"""OwncastSentry test suite."""
|
||||
@@ -0,0 +1,218 @@
|
||||
# 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.
|
||||
|
||||
"""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 owncastsentry import OwncastSentry
|
||||
from owncastsentry.config import Config
|
||||
from owncastsentry.database import StreamRepository, SubscriptionRepository
|
||||
from owncastsentry.migrations import get_upgrade_table
|
||||
|
||||
if TYPE_CHECKING:
|
||||
from collections.abc import AsyncIterator
|
||||
from pathlib import Path
|
||||
|
||||
from owncastsentry.models import StreamConfig, StreamState
|
||||
|
||||
|
||||
@pytest.fixture
|
||||
async def database(tmp_path: Path) -> AsyncIterator[Database]:
|
||||
"""Yield a real SQLite-backed mautrix Database with migrations applied."""
|
||||
db_path = tmp_path / "test.db"
|
||||
db = Database.create(
|
||||
f"sqlite:///{db_path}",
|
||||
upgrade_table=get_upgrade_table(),
|
||||
)
|
||||
await db.start()
|
||||
yield db
|
||||
await db.stop()
|
||||
|
||||
|
||||
@pytest.fixture
|
||||
def stream_repo(database: Database) -> StreamRepository:
|
||||
"""Return a StreamRepository backed by the test database."""
|
||||
return StreamRepository(database)
|
||||
|
||||
|
||||
@pytest.fixture
|
||||
def subscription_repo(database: Database) -> SubscriptionRepository:
|
||||
"""Return a SubscriptionRepository backed by the test database."""
|
||||
return SubscriptionRepository(database)
|
||||
|
||||
|
||||
# These fixtures use names expected by the maubot.testing framework.
|
||||
# They look unused but are discovered automatically by maubot's pytest plugin.
|
||||
|
||||
|
||||
@pytest.fixture
|
||||
def maubot_plugin_class():
|
||||
"""Use OwncastSentry as the plugin class for maubot integration tests."""
|
||||
return OwncastSentry
|
||||
|
||||
|
||||
@pytest.fixture
|
||||
def maubot_plugin_config_class():
|
||||
"""Use our Config class for maubot integration tests."""
|
||||
return Config
|
||||
|
||||
|
||||
@pytest.fixture
|
||||
def maubot_upgrade_table():
|
||||
"""Provide the database upgrade table for maubot integration tests."""
|
||||
return get_upgrade_table()
|
||||
|
||||
|
||||
VALID_CONFIG_RESPONSE: dict[str, object] = {
|
||||
"appearanceVariables": {},
|
||||
"name": "LogalDeveloper's Live Stream",
|
||||
"customStyles": "",
|
||||
"streamTitle": "I think I can do this... Let's start a nuclear reaction - Playing Nucleares!", # noqa: E501
|
||||
"offlineMessage": "<p>HTTP/1.1 204 No Content</p>\n<p>You've reached the right place, but I'm not live right now.</p>", # noqa: E501
|
||||
"logo": "/logo",
|
||||
"version": "Owncast v0.2.4-linux-64bit (8e89391309dd3aa4d0db4361a2ba1e144f42f8c8)",
|
||||
"extraPageContent": (
|
||||
"<h1>About</h1>\n<p>Hi there! I'm Logan, a cybersecurity professional "
|
||||
"with a background in Linux systems and network administration. This is my "
|
||||
"little corner of the internet where I exclusively run my live streams. While "
|
||||
"gaming is my most common stream topic, I occasionally stream other things "
|
||||
"which interest me, such as cybersecurity challenges. If you are looking for "
|
||||
"my main website, where my blog and project directory is hosted, please visit: "
|
||||
'<a href="https://logal.dev/">https://logal.dev/</a></p>\n<p>Any views and '
|
||||
"opinions expressed in my live streams and videos are my own and do not "
|
||||
"necessarily reflect those of my employer or any affiliated organizations."
|
||||
"</p>\n<h1>Following</h1>\n<p>I don\u2019t stick to a strict streaming "
|
||||
"schedule, but there are a few ways to stay updated so you know when I go "
|
||||
"live:</p>\n<ul>\n<li>\u2b50 <strong>Matrix</strong>: Send the message "
|
||||
"<code>!subscribe stream.logal.dev</code> to "
|
||||
'<a href="https://matrix.to/#/@owncastsentry:logal.dev">'
|
||||
"@owncastsentry:logal.dev</a>.</li>\n<li><strong>Fediverse</strong>: Follow "
|
||||
"<code>@notify@stream.logal.dev</code>.</li>\n</ul>\n<h1>Chat Commands</h1>"
|
||||
"\n<p>To make things more interactive, there are several commands you can "
|
||||
"send in chat to play sound effects live on stream:</p>\n<ul>\n"
|
||||
"<li>!boom</li>\n<li>!bluetooth</li>\n<li>!bruh</li>\n<li>!creeper</li>\n"
|
||||
"<li>!directed</li>\n<li>!fart</li>\n<li>!icetea</li>\n<li>!oof</li>\n"
|
||||
"<li>!perfect</li>\n<li>!spare</li>\n<li>!thatsit</li>\n<li>!usb</li>\n"
|
||||
"<li>!whocares</li>\n<li>!yoda</li>\n</ul>"
|
||||
),
|
||||
"summary": "Video games, cybersecurity, and more!",
|
||||
"tags": [
|
||||
"video games",
|
||||
"chatting",
|
||||
"casual",
|
||||
"english",
|
||||
"streaming",
|
||||
"owncast",
|
||||
"variety",
|
||||
],
|
||||
"socialHandles": None,
|
||||
"externalActions": [
|
||||
{
|
||||
"url": "https://ko-fi.com/logaldeveloper",
|
||||
"html": "",
|
||||
"title": "Tip",
|
||||
"description": "",
|
||||
"icon": "https://stream.logal.dev/img/platformlogos/ko-fi.svg",
|
||||
"color": "",
|
||||
"openExternally": True,
|
||||
},
|
||||
{
|
||||
"url": "https://tubefree.org/c/logaldeveloper_stream_archive",
|
||||
"html": "",
|
||||
"title": "Previous Live Stream Recordings",
|
||||
"description": "",
|
||||
"icon": "https://stream.logal.dev/img/platformlogos/fediverse.svg",
|
||||
"color": "",
|
||||
"openExternally": True,
|
||||
},
|
||||
],
|
||||
"notifications": {
|
||||
"browser": {
|
||||
"publicKey": "BI9BhIY6c7nfP6ZSIu7T53Lta5sGYqDqwSpiCabZ0XQxPQOmnUHRSjInaa3HX9XmYE-bV8SFmbFk4stZy2jnJ3M", # noqa: E501
|
||||
"enabled": True,
|
||||
},
|
||||
},
|
||||
"federation": {
|
||||
"account": "notify@stream.logal.dev",
|
||||
"followerCount": 74,
|
||||
"enabled": True,
|
||||
},
|
||||
"maxSocketPayloadSize": 2048,
|
||||
"hideViewerCount": False,
|
||||
"chatDisabled": False,
|
||||
"chatSpamProtectionDisabled": False,
|
||||
"nsfw": False,
|
||||
"authentication": {"indieAuthEnabled": True},
|
||||
}
|
||||
|
||||
VALID_STATUS_RESPONSE: dict[str, object] = {
|
||||
"serverTime": "2026-03-13T15:16:14.354962696-04:00",
|
||||
"lastConnectTime": None,
|
||||
"lastDisconnectTime": "2026-03-04T21:05:32-05:00",
|
||||
"versionNumber": "0.2.4",
|
||||
"streamTitle": "I think I can do this... Let's start a nuclear reaction - Playing Nucleares!", # noqa: E501
|
||||
"online": False,
|
||||
}
|
||||
|
||||
|
||||
@dataclass
|
||||
class _SentMessage:
|
||||
"""A message recorded by _StubMatrixClient."""
|
||||
|
||||
room_id: str
|
||||
content: Any
|
||||
|
||||
|
||||
class _StubMatrixClient:
|
||||
"""Recording stub for the Matrix client used in notification tests."""
|
||||
|
||||
def __init__(self) -> None:
|
||||
self.sent_messages: list[_SentMessage] = []
|
||||
self.should_fail_for_rooms: set[str] = set()
|
||||
|
||||
async def send_message(self, room_id: str, content: Any) -> None:
|
||||
"""Record a sent message, or raise if room is in the fail set."""
|
||||
if room_id in self.should_fail_for_rooms:
|
||||
msg = f"Stubbed failure for room {room_id}"
|
||||
raise RuntimeError(msg)
|
||||
self.sent_messages.append(_SentMessage(room_id=room_id, content=content))
|
||||
|
||||
|
||||
@dataclass
|
||||
class _StubOwncastClient:
|
||||
"""Configurable stub for the Owncast HTTP client."""
|
||||
|
||||
stream_state: StreamState | None = None
|
||||
stream_config: StreamConfig | None = None
|
||||
state_call_count: int = field(default=0, init=False)
|
||||
config_call_count: int = field(default=0, init=False)
|
||||
queried_domains: list[str] = field(default_factory=list, init=False)
|
||||
|
||||
async def get_stream_state(self, domain: str) -> StreamState | None:
|
||||
"""Return the configured stream state."""
|
||||
self.state_call_count += 1
|
||||
self.queried_domains.append(domain)
|
||||
return self.stream_state
|
||||
|
||||
async def get_stream_config(self, domain: str) -> StreamConfig | None:
|
||||
"""Return the configured stream config."""
|
||||
self.config_call_count += 1
|
||||
return self.stream_config
|
||||
@@ -0,0 +1,455 @@
|
||||
# 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 bot command handlers."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
import logging
|
||||
from datetime import UTC, datetime, timedelta
|
||||
from unittest.mock import MagicMock
|
||||
|
||||
import pytest
|
||||
import time_machine
|
||||
from aioresponses import aioresponses
|
||||
|
||||
from owncastsentry.commands import CommandHandler
|
||||
from owncastsentry.models import StreamState
|
||||
from owncastsentry.utils import OWNCAST_STATUS_PATH
|
||||
from tests.conftest import VALID_STATUS_RESPONSE
|
||||
|
||||
|
||||
def _make_command_handler() -> CommandHandler:
|
||||
"""Build a CommandHandler with dummy dependencies for pure logic tests."""
|
||||
return CommandHandler(
|
||||
owncast_client=MagicMock(),
|
||||
stream_repo=MagicMock(),
|
||||
subscription_repo=MagicMock(),
|
||||
logger=logging.getLogger("test"),
|
||||
)
|
||||
|
||||
|
||||
class TestFormatDuration:
|
||||
"""Elapsed time calculation from ISO timestamps."""
|
||||
|
||||
_NOW = datetime(2026, 3, 13, 12, 0, 0, tzinfo=UTC)
|
||||
|
||||
@pytest.mark.parametrize(
|
||||
("seconds_ago", "expected"),
|
||||
[
|
||||
pytest.param(1, "1 second", id="singular-second"),
|
||||
pytest.param(30, "30 seconds", id="plural-seconds"),
|
||||
pytest.param(60, "1 minute", id="singular-minute"),
|
||||
pytest.param(120, "2 minutes", id="plural-minutes"),
|
||||
pytest.param(3600, "1 hour", id="singular-hour"),
|
||||
pytest.param(7200, "2 hours", id="plural-hours"),
|
||||
pytest.param(86400, "1 day", id="singular-day"),
|
||||
pytest.param(172800, "2 days", id="plural-days"),
|
||||
],
|
||||
)
|
||||
@time_machine.travel(_NOW)
|
||||
def test_formats_duration(self, seconds_ago: int, expected: str) -> None:
|
||||
"""Format a timestamp into a human-readable duration."""
|
||||
handler = _make_command_handler()
|
||||
timestamp = (self._NOW - timedelta(seconds=seconds_ago)).isoformat()
|
||||
result = handler._format_duration(timestamp)
|
||||
assert result == expected
|
||||
|
||||
def test_invalid_timestamp(self) -> None:
|
||||
"""Return 'unknown duration' for unparsable timestamps."""
|
||||
handler = _make_command_handler()
|
||||
assert handler._format_duration("not-a-timestamp") == "unknown duration"
|
||||
|
||||
|
||||
class TestSubscribeCommand:
|
||||
"""Subscribe command end-to-end via maubot."""
|
||||
|
||||
async def test_subscribe_valid_stream(self, maubot_test_bot, maubot_plugin) -> None:
|
||||
"""Subscribe to a valid Owncast stream."""
|
||||
status_url = f"https://stream.logal.dev{OWNCAST_STATUS_PATH}"
|
||||
with aioresponses() as mocked:
|
||||
mocked.get(
|
||||
status_url,
|
||||
body=json.dumps(VALID_STATUS_RESPONSE).encode(),
|
||||
)
|
||||
await maubot_test_bot.send("!subscribe stream.logal.dev")
|
||||
|
||||
assert len(maubot_test_bot.responded) == 1
|
||||
assert maubot_test_bot.responded[0].content.body == (
|
||||
"Subscription added! This room will receive notifications when "
|
||||
"stream.logal.dev goes live."
|
||||
)
|
||||
|
||||
async def test_subscribe_invalid_stream(
|
||||
self, maubot_test_bot, maubot_plugin
|
||||
) -> None:
|
||||
"""Reject subscription to an invalid Owncast instance."""
|
||||
status_url = f"https://invalid.com{OWNCAST_STATUS_PATH}"
|
||||
with aioresponses() as mocked:
|
||||
mocked.get(status_url, status=404)
|
||||
await maubot_test_bot.send("!subscribe invalid.com")
|
||||
|
||||
assert len(maubot_test_bot.responded) == 1
|
||||
assert maubot_test_bot.responded[0].content.body == (
|
||||
"The URL you supplied does not appear to "
|
||||
"be a valid Owncast instance. You may have "
|
||||
"specified an invalid domain, or the "
|
||||
"instance is offline."
|
||||
)
|
||||
|
||||
async def test_subscribe_already_subscribed(
|
||||
self, maubot_test_bot, maubot_plugin
|
||||
) -> None:
|
||||
"""Reject duplicate subscription in the same room."""
|
||||
status_url = f"https://stream.logal.dev{OWNCAST_STATUS_PATH}"
|
||||
with aioresponses() as mocked:
|
||||
mocked.get(
|
||||
status_url,
|
||||
body=json.dumps(VALID_STATUS_RESPONSE).encode(),
|
||||
)
|
||||
await maubot_test_bot.send("!subscribe stream.logal.dev")
|
||||
|
||||
# Second subscribe; stream already exists so validation is skipped
|
||||
await maubot_test_bot.send("!subscribe stream.logal.dev")
|
||||
|
||||
assert len(maubot_test_bot.responded) == 2
|
||||
assert maubot_test_bot.responded[1].content.body == (
|
||||
"This room is already subscribed to notifications for stream.logal.dev."
|
||||
)
|
||||
|
||||
async def test_subscribe_existing_stream_new_room(
|
||||
self, maubot_test_bot, maubot_plugin
|
||||
) -> None:
|
||||
"""Skip instance validation when subscribing from a new room."""
|
||||
status_url = f"https://stream.logal.dev{OWNCAST_STATUS_PATH}"
|
||||
with aioresponses() as mocked:
|
||||
mocked.get(
|
||||
status_url,
|
||||
body=json.dumps(VALID_STATUS_RESPONSE).encode(),
|
||||
)
|
||||
await maubot_test_bot.send("!subscribe stream.logal.dev")
|
||||
|
||||
# Subscribe from a different room; skips validation and should not
|
||||
# query the remote instance. An empty aioresponses context will raise
|
||||
# ConnectionError if any HTTP request is attempted.
|
||||
with aioresponses():
|
||||
await maubot_test_bot.send(
|
||||
"!subscribe stream.logal.dev", room_id="!otherroom:example.com"
|
||||
)
|
||||
|
||||
assert len(maubot_test_bot.responded) == 2
|
||||
assert maubot_test_bot.responded[1].content.body == (
|
||||
"Subscription added! This room will receive notifications when "
|
||||
"stream.logal.dev goes live."
|
||||
)
|
||||
|
||||
|
||||
class TestUnsubscribeCommand:
|
||||
"""Unsubscribe command end-to-end via maubot."""
|
||||
|
||||
async def test_unsubscribe_existing(self, maubot_test_bot, maubot_plugin) -> None:
|
||||
"""Unsubscribe from a subscribed stream."""
|
||||
status_url = f"https://stream.logal.dev{OWNCAST_STATUS_PATH}"
|
||||
with aioresponses() as mocked:
|
||||
mocked.get(
|
||||
status_url,
|
||||
body=json.dumps(VALID_STATUS_RESPONSE).encode(),
|
||||
)
|
||||
await maubot_test_bot.send("!subscribe stream.logal.dev")
|
||||
|
||||
await maubot_test_bot.send("!unsubscribe stream.logal.dev")
|
||||
|
||||
assert len(maubot_test_bot.responded) == 2
|
||||
assert maubot_test_bot.responded[1].content.body == (
|
||||
"Subscription removed! This room will no "
|
||||
"longer receive notifications for stream.logal.dev."
|
||||
)
|
||||
|
||||
async def test_unsubscribe_not_subscribed(
|
||||
self, maubot_test_bot, maubot_plugin
|
||||
) -> None:
|
||||
"""Reject unsubscribe from a non-subscribed stream."""
|
||||
await maubot_test_bot.send("!unsubscribe unknown.com")
|
||||
|
||||
assert len(maubot_test_bot.responded) == 1
|
||||
assert maubot_test_bot.responded[0].content.body == (
|
||||
"This room is already not subscribed to notifications for unknown.com."
|
||||
)
|
||||
|
||||
|
||||
class TestSubscriptionsCommand:
|
||||
"""Subscriptions listing command end-to-end via maubot."""
|
||||
|
||||
async def test_no_subscriptions(self, maubot_test_bot, maubot_plugin) -> None:
|
||||
"""Show help text when no subscriptions exist."""
|
||||
await maubot_test_bot.send("!subscriptions")
|
||||
|
||||
assert len(maubot_test_bot.responded) == 1
|
||||
assert maubot_test_bot.responded[0].content.body == (
|
||||
"This room is not subscribed to any Owncast "
|
||||
"instances.\n\nTo subscribe to an Owncast "
|
||||
"instance, use `!subscribe <domain>`"
|
||||
)
|
||||
|
||||
@time_machine.travel(datetime(2026, 3, 13, 12, 0, 0, tzinfo=UTC))
|
||||
async def test_shows_online_stream(self, maubot_test_bot, maubot_plugin) -> None:
|
||||
"""Show stream details including title and duration."""
|
||||
# Subscribe first
|
||||
status_url = f"https://stream.logal.dev{OWNCAST_STATUS_PATH}"
|
||||
with aioresponses() as mocked:
|
||||
mocked.get(
|
||||
status_url,
|
||||
body=json.dumps(VALID_STATUS_RESPONSE).encode(),
|
||||
)
|
||||
await maubot_test_bot.send("!subscribe stream.logal.dev")
|
||||
|
||||
# Update stream state to be online with a name
|
||||
await maubot_plugin.stream_repo.update(
|
||||
StreamState(
|
||||
domain="stream.logal.dev",
|
||||
name="Test Stream",
|
||||
title="Playing Games",
|
||||
last_connect_time="2026-01-01T12:00:00Z",
|
||||
)
|
||||
)
|
||||
|
||||
await maubot_test_bot.send("!subscriptions")
|
||||
|
||||
assert len(maubot_test_bot.responded) == 2
|
||||
assert maubot_test_bot.responded[1].content.body == (
|
||||
"**Subscriptions for this room (1):**\n\n"
|
||||
"● **Test Stream** \n"
|
||||
" ○ Title: Playing Games\n"
|
||||
" ○ Status: Online for 71 days\n"
|
||||
" ○ Link: https://stream.logal.dev\n"
|
||||
"To unsubscribe from any of these Owncast "
|
||||
"instances, use `!unsubscribe <domain>`"
|
||||
)
|
||||
|
||||
@time_machine.travel(datetime(2026, 3, 13, 12, 0, 0, tzinfo=UTC))
|
||||
async def test_shows_offline_stream(self, maubot_test_bot, maubot_plugin) -> None:
|
||||
"""Show offline status for non-live streams."""
|
||||
# Subscribe first
|
||||
status_url = f"https://stream.logal.dev{OWNCAST_STATUS_PATH}"
|
||||
with aioresponses() as mocked:
|
||||
mocked.get(
|
||||
status_url,
|
||||
body=json.dumps(VALID_STATUS_RESPONSE).encode(),
|
||||
)
|
||||
await maubot_test_bot.send("!subscribe stream.logal.dev")
|
||||
|
||||
# Update stream state to be offline
|
||||
await maubot_plugin.stream_repo.update(
|
||||
StreamState(
|
||||
domain="stream.logal.dev",
|
||||
name="Test Stream",
|
||||
last_disconnect_time="2026-01-01T10:00:00Z",
|
||||
)
|
||||
)
|
||||
|
||||
await maubot_test_bot.send("!subscriptions")
|
||||
|
||||
assert len(maubot_test_bot.responded) == 2
|
||||
assert maubot_test_bot.responded[1].content.body == (
|
||||
"**Subscriptions for this room (1):**\n\n"
|
||||
"● **Test Stream** \n"
|
||||
" ○ Status: Offline for 71 days\n"
|
||||
" ○ Link: https://stream.logal.dev\n"
|
||||
"To unsubscribe from any of these Owncast "
|
||||
"instances, use `!unsubscribe <domain>`"
|
||||
)
|
||||
|
||||
@time_machine.travel(datetime(2026, 3, 13, 12, 0, 0, tzinfo=UTC))
|
||||
async def test_shows_multiple_subscriptions(
|
||||
self, maubot_test_bot, maubot_plugin
|
||||
) -> None:
|
||||
"""List subscriptions alphabetically with mixed statuses."""
|
||||
# Subscribe in reverse alphabetical order to verify sorted output
|
||||
with aioresponses() as mocked:
|
||||
mocked.get(
|
||||
f"https://beta.com{OWNCAST_STATUS_PATH}",
|
||||
body=json.dumps(VALID_STATUS_RESPONSE).encode(),
|
||||
)
|
||||
mocked.get(
|
||||
f"https://alpha.com{OWNCAST_STATUS_PATH}",
|
||||
body=json.dumps(VALID_STATUS_RESPONSE).encode(),
|
||||
)
|
||||
await maubot_test_bot.send("!subscribe beta.com")
|
||||
await maubot_test_bot.send("!subscribe alpha.com")
|
||||
|
||||
# Set alpha online, beta offline
|
||||
await maubot_plugin.stream_repo.update(
|
||||
StreamState(
|
||||
domain="alpha.com",
|
||||
name="Alpha Stream",
|
||||
title="Streaming Live",
|
||||
last_connect_time="2026-03-13T10:00:00Z",
|
||||
)
|
||||
)
|
||||
await maubot_plugin.stream_repo.update(
|
||||
StreamState(
|
||||
domain="beta.com",
|
||||
name="Beta Stream",
|
||||
last_disconnect_time="2026-03-12T18:00:00Z",
|
||||
)
|
||||
)
|
||||
|
||||
await maubot_test_bot.send("!subscriptions")
|
||||
|
||||
assert len(maubot_test_bot.responded) == 3
|
||||
assert maubot_test_bot.responded[2].content.body == (
|
||||
"**Subscriptions for this room (2):**\n\n"
|
||||
"● **Alpha Stream**\n"
|
||||
" \n"
|
||||
" ○ Title: Streaming Live\n"
|
||||
" ○ Status: Online for 2 hours\n"
|
||||
" ○ Link: https://alpha.com\n"
|
||||
"● **Beta Stream**\n"
|
||||
" \n"
|
||||
" ○ Status: Offline for 18 hours\n"
|
||||
" ○ Link: https://beta.com\n"
|
||||
"To unsubscribe from any of these Owncast "
|
||||
"instances, use `!unsubscribe <domain>`"
|
||||
)
|
||||
|
||||
|
||||
class TestLiveCommand:
|
||||
"""Live streams listing command end-to-end via maubot."""
|
||||
|
||||
async def test_no_subscriptions(self, maubot_test_bot, maubot_plugin) -> None:
|
||||
"""Show help text when no subscriptions exist."""
|
||||
await maubot_test_bot.send("!live")
|
||||
|
||||
assert len(maubot_test_bot.responded) == 1
|
||||
assert maubot_test_bot.responded[0].content.body == (
|
||||
"This room is not subscribed to any Owncast "
|
||||
"instances.\n\nTo subscribe to an Owncast "
|
||||
"instance, use `!subscribe <domain>`"
|
||||
)
|
||||
|
||||
async def test_no_live_streams(self, maubot_test_bot, maubot_plugin) -> None:
|
||||
"""Show 'no live' message when all streams are offline."""
|
||||
# Subscribe first
|
||||
status_url = f"https://stream.logal.dev{OWNCAST_STATUS_PATH}"
|
||||
with aioresponses() as mocked:
|
||||
mocked.get(
|
||||
status_url,
|
||||
body=json.dumps(VALID_STATUS_RESPONSE).encode(),
|
||||
)
|
||||
await maubot_test_bot.send("!subscribe stream.logal.dev")
|
||||
|
||||
# Update stream state to offline
|
||||
await maubot_plugin.stream_repo.update(
|
||||
StreamState(
|
||||
domain="stream.logal.dev",
|
||||
name="Test Stream",
|
||||
last_disconnect_time="2026-01-01T10:00:00Z",
|
||||
)
|
||||
)
|
||||
|
||||
await maubot_test_bot.send("!live")
|
||||
|
||||
assert len(maubot_test_bot.responded) == 2
|
||||
assert maubot_test_bot.responded[1].content.body == (
|
||||
"No subscribed Owncast instances are currently "
|
||||
"live.\n\nUse `!subscriptions` to list all "
|
||||
"subscriptions."
|
||||
)
|
||||
|
||||
@time_machine.travel(datetime(2026, 3, 13, 12, 0, 0, tzinfo=UTC))
|
||||
async def test_shows_live_stream(self, maubot_test_bot, maubot_plugin) -> None:
|
||||
"""Show live stream with title and duration."""
|
||||
# Subscribe first
|
||||
status_url = f"https://stream.logal.dev{OWNCAST_STATUS_PATH}"
|
||||
with aioresponses() as mocked:
|
||||
mocked.get(
|
||||
status_url,
|
||||
body=json.dumps(VALID_STATUS_RESPONSE).encode(),
|
||||
)
|
||||
await maubot_test_bot.send("!subscribe stream.logal.dev")
|
||||
|
||||
# Update stream state to online
|
||||
await maubot_plugin.stream_repo.update(
|
||||
StreamState(
|
||||
domain="stream.logal.dev",
|
||||
name="Test Stream",
|
||||
title="Playing Games",
|
||||
last_connect_time="2026-01-01T12:00:00Z",
|
||||
)
|
||||
)
|
||||
|
||||
await maubot_test_bot.send("!live")
|
||||
|
||||
assert len(maubot_test_bot.responded) == 2
|
||||
assert maubot_test_bot.responded[1].content.body == (
|
||||
"**Live Owncast instances (1):**\n\n"
|
||||
"● **Test Stream** \n"
|
||||
" ○ Title: Playing Games\n"
|
||||
" ○ Online for 71 days\n"
|
||||
" ○ Link: https://stream.logal.dev"
|
||||
)
|
||||
|
||||
@time_machine.travel(datetime(2026, 3, 13, 12, 0, 0, tzinfo=UTC))
|
||||
async def test_shows_multiple_live_streams(
|
||||
self, maubot_test_bot, maubot_plugin
|
||||
) -> None:
|
||||
"""List live streams alphabetically with different durations."""
|
||||
# Subscribe in reverse alphabetical order to verify sorted output
|
||||
with aioresponses() as mocked:
|
||||
mocked.get(
|
||||
f"https://beta.com{OWNCAST_STATUS_PATH}",
|
||||
body=json.dumps(VALID_STATUS_RESPONSE).encode(),
|
||||
)
|
||||
mocked.get(
|
||||
f"https://alpha.com{OWNCAST_STATUS_PATH}",
|
||||
body=json.dumps(VALID_STATUS_RESPONSE).encode(),
|
||||
)
|
||||
await maubot_test_bot.send("!subscribe beta.com")
|
||||
await maubot_test_bot.send("!subscribe alpha.com")
|
||||
|
||||
# Set both streams online with different connect times
|
||||
await maubot_plugin.stream_repo.update(
|
||||
StreamState(
|
||||
domain="alpha.com",
|
||||
name="Alpha Stream",
|
||||
title="Morning Show",
|
||||
last_connect_time="2026-03-13T10:00:00Z",
|
||||
)
|
||||
)
|
||||
await maubot_plugin.stream_repo.update(
|
||||
StreamState(
|
||||
domain="beta.com",
|
||||
name="Beta Stream",
|
||||
title="Evening Vibes",
|
||||
last_connect_time="2026-03-13T06:00:00Z",
|
||||
)
|
||||
)
|
||||
|
||||
await maubot_test_bot.send("!live")
|
||||
|
||||
assert len(maubot_test_bot.responded) == 3
|
||||
assert maubot_test_bot.responded[2].content.body == (
|
||||
"**Live Owncast instances (2):**\n\n"
|
||||
"● **Alpha Stream**\n"
|
||||
" \n"
|
||||
" ○ Title: Morning Show\n"
|
||||
" ○ Online for 2 hours\n"
|
||||
" ○ Link: https://alpha.com\n"
|
||||
"● **Beta Stream**\n"
|
||||
" \n"
|
||||
" ○ Title: Evening Vibes\n"
|
||||
" ○ Online for 6 hours\n"
|
||||
" ○ Link: https://beta.com"
|
||||
)
|
||||
@@ -0,0 +1,150 @@
|
||||
# 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 database repository classes."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from typing import TYPE_CHECKING
|
||||
|
||||
if TYPE_CHECKING:
|
||||
from owncastsentry.database import StreamRepository, SubscriptionRepository
|
||||
|
||||
|
||||
class TestStreamExists:
|
||||
"""Stream existence checks."""
|
||||
|
||||
async def test_returns_true_for_existing_stream(
|
||||
self, stream_repo: StreamRepository
|
||||
) -> None:
|
||||
"""Return True when the stream exists in the database."""
|
||||
await stream_repo.create("example.com")
|
||||
assert await stream_repo.exists("example.com") is True
|
||||
|
||||
async def test_returns_false_for_missing_stream(
|
||||
self, stream_repo: StreamRepository
|
||||
) -> None:
|
||||
"""Return False when the stream does not exist in the database."""
|
||||
assert await stream_repo.exists("missing.com") is False
|
||||
|
||||
|
||||
class TestStreamDelete:
|
||||
"""Stream record deletion."""
|
||||
|
||||
async def test_removes_stream_record(self, stream_repo: StreamRepository) -> None:
|
||||
"""Remove the stream record so get_by_domain returns None."""
|
||||
await stream_repo.create("example.com")
|
||||
await stream_repo.delete("example.com")
|
||||
assert await stream_repo.get_by_domain("example.com") is None
|
||||
|
||||
|
||||
class TestGetSubscribedStreamsForRoom:
|
||||
"""Subscribed stream lookup by room."""
|
||||
|
||||
async def test_returns_all_domains_for_room(
|
||||
self,
|
||||
stream_repo: StreamRepository,
|
||||
subscription_repo: SubscriptionRepository,
|
||||
) -> None:
|
||||
"""Return all domains a room is subscribed to."""
|
||||
await stream_repo.create("alpha.com")
|
||||
await stream_repo.create("beta.com")
|
||||
await subscription_repo.add("alpha.com", "!room1:example.com")
|
||||
await subscription_repo.add("beta.com", "!room1:example.com")
|
||||
|
||||
result = await subscription_repo.get_subscribed_streams_for_room(
|
||||
"!room1:example.com"
|
||||
)
|
||||
assert sorted(result) == ["alpha.com", "beta.com"]
|
||||
|
||||
async def test_returns_empty_list_for_unsubscribed_room(
|
||||
self, subscription_repo: SubscriptionRepository
|
||||
) -> None:
|
||||
"""Return an empty list when the room has no subscriptions."""
|
||||
result = await subscription_repo.get_subscribed_streams_for_room(
|
||||
"!nobody:example.com"
|
||||
)
|
||||
assert result == []
|
||||
|
||||
|
||||
class TestGetAllSubscribedDomains:
|
||||
"""Unique subscribed domain retrieval."""
|
||||
|
||||
async def test_returns_each_domain_once(
|
||||
self,
|
||||
stream_repo: StreamRepository,
|
||||
subscription_repo: SubscriptionRepository,
|
||||
) -> None:
|
||||
"""Return each domain once even with multiple subscriptions."""
|
||||
await stream_repo.create("alpha.com")
|
||||
await subscription_repo.add("alpha.com", "!room1:example.com")
|
||||
await subscription_repo.add("alpha.com", "!room2:example.com")
|
||||
|
||||
result = await subscription_repo.get_all_subscribed_domains()
|
||||
assert result == ["alpha.com"]
|
||||
|
||||
async def test_returns_empty_list_with_no_subscriptions(
|
||||
self, subscription_repo: SubscriptionRepository
|
||||
) -> None:
|
||||
"""Return an empty list when there are no subscriptions."""
|
||||
result = await subscription_repo.get_all_subscribed_domains()
|
||||
assert result == []
|
||||
|
||||
|
||||
class TestCountByDomain:
|
||||
"""Subscription count by domain."""
|
||||
|
||||
async def test_returns_correct_count(
|
||||
self,
|
||||
stream_repo: StreamRepository,
|
||||
subscription_repo: SubscriptionRepository,
|
||||
) -> None:
|
||||
"""Return the correct subscription count for a domain."""
|
||||
await stream_repo.create("alpha.com")
|
||||
await subscription_repo.add("alpha.com", "!room1:example.com")
|
||||
await subscription_repo.add("alpha.com", "!room2:example.com")
|
||||
|
||||
assert await subscription_repo.count_by_domain("alpha.com") == 2
|
||||
|
||||
async def test_returns_zero_for_unknown_domain(
|
||||
self, subscription_repo: SubscriptionRepository
|
||||
) -> None:
|
||||
"""Return 0 for a domain with no subscriptions."""
|
||||
assert await subscription_repo.count_by_domain("unknown.com") == 0
|
||||
|
||||
|
||||
class TestDeleteAllForDomain:
|
||||
"""Bulk subscription deletion by domain."""
|
||||
|
||||
async def test_deletes_all_subscriptions_and_returns_count(
|
||||
self,
|
||||
stream_repo: StreamRepository,
|
||||
subscription_repo: SubscriptionRepository,
|
||||
) -> None:
|
||||
"""Delete all subscriptions for the domain and return the count."""
|
||||
await stream_repo.create("alpha.com")
|
||||
await subscription_repo.add("alpha.com", "!room1:example.com")
|
||||
await subscription_repo.add("alpha.com", "!room2:example.com")
|
||||
|
||||
deleted = await subscription_repo.delete_all_for_domain("alpha.com")
|
||||
assert deleted == 2
|
||||
|
||||
rooms = await subscription_repo.get_subscribed_rooms("alpha.com")
|
||||
assert rooms == []
|
||||
|
||||
async def test_returns_zero_for_unknown_domain(
|
||||
self, subscription_repo: SubscriptionRepository
|
||||
) -> None:
|
||||
"""Return 0 when deleting subscriptions for an unknown domain."""
|
||||
assert await subscription_repo.delete_all_for_domain("unknown.com") == 0
|
||||
@@ -0,0 +1,171 @@
|
||||
# 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,195 @@
|
||||
# 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 data models."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import pytest
|
||||
|
||||
from owncastsentry.models import StreamConfig, StreamState, StreamStatus
|
||||
from owncastsentry.utils import (
|
||||
MAX_INSTANCE_TITLE_LENGTH,
|
||||
MAX_STREAM_TITLE_LENGTH,
|
||||
MAX_TAG_LENGTH,
|
||||
UNKNOWN_STATUS_THRESHOLD,
|
||||
)
|
||||
|
||||
|
||||
class TestStreamStateStatus:
|
||||
"""Stream status derivation from state fields."""
|
||||
|
||||
@pytest.mark.parametrize(
|
||||
("failure_counter", "last_connect_time", "expected"),
|
||||
[
|
||||
pytest.param(
|
||||
UNKNOWN_STATUS_THRESHOLD + 1,
|
||||
None,
|
||||
StreamStatus.UNKNOWN,
|
||||
id="above-threshold-offline-returns-unknown",
|
||||
),
|
||||
pytest.param(
|
||||
UNKNOWN_STATUS_THRESHOLD + 1,
|
||||
"2026-01-01T00:00:00Z",
|
||||
StreamStatus.UNKNOWN,
|
||||
id="above-threshold-online-returns-unknown",
|
||||
),
|
||||
pytest.param(
|
||||
0,
|
||||
"2026-01-01T00:00:00Z",
|
||||
StreamStatus.ONLINE,
|
||||
id="zero-failures-with-connect-time-returns-online",
|
||||
),
|
||||
pytest.param(
|
||||
0,
|
||||
None,
|
||||
StreamStatus.OFFLINE,
|
||||
id="zero-failures-no-connect-time-returns-offline",
|
||||
),
|
||||
pytest.param(
|
||||
UNKNOWN_STATUS_THRESHOLD,
|
||||
"2026-01-01T00:00:00Z",
|
||||
StreamStatus.ONLINE,
|
||||
id="at-threshold-with-connect-time-returns-online",
|
||||
),
|
||||
pytest.param(
|
||||
UNKNOWN_STATUS_THRESHOLD,
|
||||
None,
|
||||
StreamStatus.OFFLINE,
|
||||
id="at-threshold-no-connect-time-returns-offline",
|
||||
),
|
||||
],
|
||||
)
|
||||
def test_status(
|
||||
self,
|
||||
failure_counter: int,
|
||||
last_connect_time: str | None,
|
||||
expected: StreamStatus,
|
||||
) -> None:
|
||||
"""Return the correct status based on failure counter and connect time."""
|
||||
state = StreamState(
|
||||
domain="example.com",
|
||||
failure_counter=failure_counter,
|
||||
last_connect_time=last_connect_time,
|
||||
)
|
||||
assert state.status is expected
|
||||
|
||||
|
||||
class TestStreamStateFromApiResponse:
|
||||
"""StreamState construction from an API response dictionary."""
|
||||
|
||||
def test_typical_response(self) -> None:
|
||||
"""Populate all fields from a complete API response."""
|
||||
response = {
|
||||
"streamTitle": "My Stream",
|
||||
"lastConnectTime": "2026-01-01T00:00:00Z",
|
||||
"lastDisconnectTime": "2025-12-31T23:00:00Z",
|
||||
}
|
||||
state = StreamState.from_api_response(response, "example.com")
|
||||
assert state.domain == "example.com"
|
||||
assert state.title == "My Stream"
|
||||
assert state.last_connect_time == "2026-01-01T00:00:00Z"
|
||||
assert state.last_disconnect_time == "2025-12-31T23:00:00Z"
|
||||
assert state.name is None
|
||||
assert state.failure_counter == 0
|
||||
|
||||
def test_empty_response_defaults(self) -> None:
|
||||
"""Use defaults when optional fields are missing."""
|
||||
state = StreamState.from_api_response({}, "bare.example.com")
|
||||
assert state.domain == "bare.example.com"
|
||||
assert state.title == ""
|
||||
assert state.last_connect_time is None
|
||||
assert state.last_disconnect_time is None
|
||||
|
||||
def test_title_truncation(self) -> None:
|
||||
"""Truncate the stream title to MAX_STREAM_TITLE_LENGTH."""
|
||||
long_title = "A" * (MAX_STREAM_TITLE_LENGTH + 50)
|
||||
response = {"streamTitle": long_title}
|
||||
state = StreamState.from_api_response(response, "example.com")
|
||||
assert len(state.title) == MAX_STREAM_TITLE_LENGTH
|
||||
assert state.title == "A" * MAX_STREAM_TITLE_LENGTH
|
||||
|
||||
|
||||
class TestStreamStateFromDbRow:
|
||||
"""StreamState construction from a database row dictionary."""
|
||||
|
||||
def test_typical_row(self) -> None:
|
||||
"""Populate all fields from a complete database row."""
|
||||
row = {
|
||||
"domain": "example.com",
|
||||
"name": "Test Instance",
|
||||
"title": "Live Now",
|
||||
"last_connect_time": "2026-01-01T00:00:00Z",
|
||||
"last_disconnect_time": "2025-12-31T23:00:00Z",
|
||||
"failure_counter": 3,
|
||||
}
|
||||
state = StreamState.from_db_row(row)
|
||||
assert state.domain == "example.com"
|
||||
assert state.name == "Test Instance"
|
||||
assert state.title == "Live Now"
|
||||
assert state.last_connect_time == "2026-01-01T00:00:00Z"
|
||||
assert state.last_disconnect_time == "2025-12-31T23:00:00Z"
|
||||
assert state.failure_counter == 3
|
||||
|
||||
def test_row_with_none_optional_fields(self) -> None:
|
||||
"""Accept None for optional fields in a database row."""
|
||||
row = {
|
||||
"domain": "example.com",
|
||||
"name": None,
|
||||
"title": None,
|
||||
"last_connect_time": None,
|
||||
"last_disconnect_time": None,
|
||||
"failure_counter": 0,
|
||||
}
|
||||
state = StreamState.from_db_row(row)
|
||||
assert state.domain == "example.com"
|
||||
assert state.name is None
|
||||
assert state.title is None
|
||||
assert state.last_connect_time is None
|
||||
assert state.last_disconnect_time is None
|
||||
assert state.failure_counter == 0
|
||||
|
||||
|
||||
class TestStreamConfigFromApiResponse:
|
||||
"""StreamConfig construction from an API response dictionary."""
|
||||
|
||||
def test_typical_response(self) -> None:
|
||||
"""Populate name and tags from a complete API response."""
|
||||
response = {"name": "My Instance", "tags": ["gaming", "music"]}
|
||||
config = StreamConfig.from_api_response(response)
|
||||
assert config.name == "My Instance"
|
||||
assert config.tags == ["gaming", "music"]
|
||||
|
||||
def test_missing_keys_defaults(self) -> None:
|
||||
"""Use defaults when name and tags keys are missing."""
|
||||
config = StreamConfig.from_api_response({})
|
||||
assert config.name == ""
|
||||
assert config.tags == []
|
||||
|
||||
def test_name_truncation(self) -> None:
|
||||
"""Truncate the instance name to MAX_INSTANCE_TITLE_LENGTH."""
|
||||
long_name = "B" * (MAX_INSTANCE_TITLE_LENGTH + 50)
|
||||
response = {"name": long_name, "tags": []}
|
||||
config = StreamConfig.from_api_response(response)
|
||||
assert len(config.name) == MAX_INSTANCE_TITLE_LENGTH
|
||||
assert config.name == "B" * MAX_INSTANCE_TITLE_LENGTH
|
||||
|
||||
def test_tag_truncation(self) -> None:
|
||||
"""Truncate each tag to MAX_TAG_LENGTH."""
|
||||
long_tag = "C" * (MAX_TAG_LENGTH + 10)
|
||||
response = {"name": "", "tags": [long_tag, "short"]}
|
||||
config = StreamConfig.from_api_response(response)
|
||||
assert len(config.tags[0]) == MAX_TAG_LENGTH
|
||||
assert config.tags[0] == "C" * MAX_TAG_LENGTH
|
||||
assert config.tags[1] == "short"
|
||||
@@ -0,0 +1,348 @@
|
||||
# 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 notification service."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import logging
|
||||
import time
|
||||
from typing import TYPE_CHECKING
|
||||
|
||||
import pytest
|
||||
|
||||
from owncastsentry.notification_service import NotificationService
|
||||
from owncastsentry.utils import SECONDS_BETWEEN_NOTIFICATIONS
|
||||
from tests.conftest import _StubMatrixClient
|
||||
|
||||
if TYPE_CHECKING:
|
||||
from owncastsentry.database import StreamRepository, SubscriptionRepository
|
||||
|
||||
|
||||
def _make_service(
|
||||
*,
|
||||
client: _StubMatrixClient,
|
||||
subscription_repo: SubscriptionRepository,
|
||||
) -> NotificationService:
|
||||
"""Build a NotificationService with a stub client and real repo."""
|
||||
return NotificationService(
|
||||
client=client,
|
||||
subscription_repo=subscription_repo,
|
||||
logger=logging.getLogger("test"),
|
||||
)
|
||||
|
||||
|
||||
class TestCanNotify:
|
||||
"""Rate-limiting logic for notification cooldowns."""
|
||||
|
||||
def test_first_notification_allowed(
|
||||
self, subscription_repo: SubscriptionRepository
|
||||
) -> None:
|
||||
"""Allow the first notification for an unseen domain."""
|
||||
service = _make_service(
|
||||
client=_StubMatrixClient(), subscription_repo=subscription_repo
|
||||
)
|
||||
assert service._can_notify("example.com") is True
|
||||
|
||||
def test_within_cooldown_blocked(
|
||||
self, subscription_repo: SubscriptionRepository
|
||||
) -> None:
|
||||
"""Block notifications within the cooldown window."""
|
||||
service = _make_service(
|
||||
client=_StubMatrixClient(), subscription_repo=subscription_repo
|
||||
)
|
||||
service.notification_timers_cache["example.com"] = time.time()
|
||||
assert service._can_notify("example.com") is False
|
||||
|
||||
def test_after_cooldown_allowed(
|
||||
self, subscription_repo: SubscriptionRepository
|
||||
) -> None:
|
||||
"""Allow notifications after the cooldown window expires."""
|
||||
service = _make_service(
|
||||
client=_StubMatrixClient(), subscription_repo=subscription_repo
|
||||
)
|
||||
# Subtract an extra second to ensure the cooldown has fully elapsed
|
||||
service.notification_timers_cache["example.com"] = (
|
||||
time.time() - SECONDS_BETWEEN_NOTIFICATIONS - 1
|
||||
)
|
||||
assert service._can_notify("example.com") is True
|
||||
|
||||
|
||||
class TestGetLastNotificationTime:
|
||||
"""Last notification timestamp retrieval."""
|
||||
|
||||
def test_returns_cached_value(
|
||||
self, subscription_repo: SubscriptionRepository
|
||||
) -> None:
|
||||
"""Return the cached timestamp for a known domain."""
|
||||
service = _make_service(
|
||||
client=_StubMatrixClient(), subscription_repo=subscription_repo
|
||||
)
|
||||
service.notification_timers_cache["example.com"] = 12345.0
|
||||
assert service.get_last_notification_time("example.com") == 12345.0
|
||||
|
||||
def test_returns_zero_for_unknown(
|
||||
self, subscription_repo: SubscriptionRepository
|
||||
) -> None:
|
||||
"""Return 0 for a domain that has never been notified."""
|
||||
service = _make_service(
|
||||
client=_StubMatrixClient(), subscription_repo=subscription_repo
|
||||
)
|
||||
assert service.get_last_notification_time("unknown.com") == 0
|
||||
|
||||
|
||||
class TestFormatMessage:
|
||||
"""Notification message formatting."""
|
||||
|
||||
@pytest.mark.parametrize(
|
||||
("name", "title", "domain", "tags", "title_change", "expected"),
|
||||
[
|
||||
pytest.param(
|
||||
"My Stream",
|
||||
"Playing Games",
|
||||
"example.com",
|
||||
[],
|
||||
False,
|
||||
"🎥 My Stream is now live!\n"
|
||||
"Stream Title: Playing Games\n"
|
||||
"\n"
|
||||
"To tune in, visit: https://example.com/",
|
||||
id="go-live-with-title",
|
||||
),
|
||||
pytest.param(
|
||||
"My Stream",
|
||||
"New Title",
|
||||
"example.com",
|
||||
[],
|
||||
True,
|
||||
"📝 My Stream has changed its stream title!\n"
|
||||
"Stream Title: New Title\n"
|
||||
"\n"
|
||||
"To tune in, visit: https://example.com/",
|
||||
id="title-change",
|
||||
),
|
||||
pytest.param(
|
||||
"My Stream",
|
||||
"",
|
||||
"example.com",
|
||||
[],
|
||||
False,
|
||||
"🎥 My Stream is now live!\n\nTo tune in, visit: https://example.com/",
|
||||
id="go-live-no-title",
|
||||
),
|
||||
pytest.param(
|
||||
"",
|
||||
"Title",
|
||||
"example.com",
|
||||
[],
|
||||
False,
|
||||
"🎥 example.com is now live!\n"
|
||||
"Stream Title: Title\n"
|
||||
"\n"
|
||||
"To tune in, visit: https://example.com/",
|
||||
id="name-fallback-to-domain",
|
||||
),
|
||||
pytest.param(
|
||||
"Stream",
|
||||
"Title",
|
||||
"example.com",
|
||||
["gaming", "fun"],
|
||||
False,
|
||||
"🎥 Stream is now live!\n"
|
||||
"Stream Title: Title\n"
|
||||
"\n"
|
||||
"To tune in, visit: https://example.com/\n"
|
||||
"\n"
|
||||
"#gaming #fun",
|
||||
id="with-tags",
|
||||
),
|
||||
pytest.param(
|
||||
"Stream",
|
||||
"Title",
|
||||
"example.com",
|
||||
[".hidden", "visible"],
|
||||
False,
|
||||
"🎥 Stream is now live!\n"
|
||||
"Stream Title: Title\n"
|
||||
"\n"
|
||||
"To tune in, visit: https://example.com/\n"
|
||||
"\n"
|
||||
"#visible",
|
||||
id="dot-prefix-tag-filtered",
|
||||
),
|
||||
pytest.param(
|
||||
"Stream",
|
||||
"Title",
|
||||
"example.com",
|
||||
[".secret"],
|
||||
False,
|
||||
"🎥 Stream is now live!\n"
|
||||
"Stream Title: Title\n"
|
||||
"\n"
|
||||
"To tune in, visit: https://example.com/",
|
||||
id="all-tags-dot-prefixed",
|
||||
),
|
||||
],
|
||||
)
|
||||
def test_format_message(
|
||||
self,
|
||||
name: str,
|
||||
title: str,
|
||||
domain: str,
|
||||
tags: list[str],
|
||||
title_change: bool,
|
||||
expected: str,
|
||||
subscription_repo: SubscriptionRepository,
|
||||
) -> None:
|
||||
"""Format the notification message with expected content."""
|
||||
service = _make_service(
|
||||
client=_StubMatrixClient(), subscription_repo=subscription_repo
|
||||
)
|
||||
result = service._format_message(name, title, domain, tags, title_change)
|
||||
assert result == expected
|
||||
|
||||
|
||||
class TestNotifyStreamLive:
|
||||
"""End-to-end notification sending."""
|
||||
|
||||
async def test_sends_to_all_subscribed_rooms(
|
||||
self,
|
||||
stream_repo: StreamRepository,
|
||||
subscription_repo: SubscriptionRepository,
|
||||
) -> None:
|
||||
"""Send notifications to every room subscribed to the domain."""
|
||||
client = _StubMatrixClient()
|
||||
service = _make_service(client=client, subscription_repo=subscription_repo)
|
||||
|
||||
await stream_repo.create("example.com")
|
||||
await subscription_repo.add("example.com", "!room1:matrix.org")
|
||||
await subscription_repo.add("example.com", "!room2:matrix.org")
|
||||
|
||||
await service.notify_stream_live("example.com", "Stream", "Title", ["tag"])
|
||||
|
||||
assert len(client.sent_messages) == 2
|
||||
room_ids = {msg.room_id for msg in client.sent_messages}
|
||||
assert room_ids == {"!room1:matrix.org", "!room2:matrix.org"}
|
||||
expected_body = (
|
||||
"🎥 Stream is now live!\n"
|
||||
"Stream Title: Title\n"
|
||||
"\n"
|
||||
"To tune in, visit: https://example.com/\n"
|
||||
"\n"
|
||||
"#tag"
|
||||
)
|
||||
for msg in client.sent_messages:
|
||||
assert msg.content.body == expected_body
|
||||
|
||||
async def test_skips_when_rate_limited(
|
||||
self,
|
||||
stream_repo: StreamRepository,
|
||||
subscription_repo: SubscriptionRepository,
|
||||
) -> None:
|
||||
"""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()
|
||||
|
||||
await stream_repo.create("example.com")
|
||||
await subscription_repo.add("example.com", "!room1:matrix.org")
|
||||
|
||||
await service.notify_stream_live("example.com", "Stream", "Title", [])
|
||||
|
||||
assert len(client.sent_messages) == 0
|
||||
|
||||
async def test_counts_failures(
|
||||
self,
|
||||
stream_repo: StreamRepository,
|
||||
subscription_repo: SubscriptionRepository,
|
||||
) -> None:
|
||||
"""Send to remaining rooms after a delivery failure."""
|
||||
client = _StubMatrixClient()
|
||||
client.should_fail_for_rooms.add("!bad:matrix.org")
|
||||
service = _make_service(client=client, subscription_repo=subscription_repo)
|
||||
|
||||
await stream_repo.create("example.com")
|
||||
await subscription_repo.add("example.com", "!bad:matrix.org")
|
||||
await subscription_repo.add("example.com", "!good:matrix.org")
|
||||
|
||||
await service.notify_stream_live("example.com", "Stream", "Title", [])
|
||||
|
||||
assert len(client.sent_messages) == 1
|
||||
assert client.sent_messages[0].room_id == "!good:matrix.org"
|
||||
assert client.sent_messages[0].content.body == (
|
||||
"🎥 Stream is now live!\n"
|
||||
"Stream Title: Title\n"
|
||||
"\n"
|
||||
"To tune in, visit: https://example.com/"
|
||||
)
|
||||
|
||||
|
||||
class TestSendCleanupWarning:
|
||||
"""Cleanup warning notification sending."""
|
||||
|
||||
async def test_sends_warning_to_all_rooms(
|
||||
self,
|
||||
stream_repo: StreamRepository,
|
||||
subscription_repo: SubscriptionRepository,
|
||||
) -> None:
|
||||
"""Send cleanup warning to all subscribed rooms."""
|
||||
client = _StubMatrixClient()
|
||||
service = _make_service(client=client, subscription_repo=subscription_repo)
|
||||
|
||||
await stream_repo.create("example.com")
|
||||
await subscription_repo.add("example.com", "!room1:matrix.org")
|
||||
|
||||
await service.send_cleanup_warning("example.com")
|
||||
|
||||
assert len(client.sent_messages) == 1
|
||||
assert client.sent_messages[0].content.body == (
|
||||
"⚠️ Warning: Subscription Cleanup Scheduled\n"
|
||||
"\n"
|
||||
"The Owncast instance at example.com has been "
|
||||
"unreachable for 83 days. If it remains "
|
||||
"unreachable for 7 more days "
|
||||
"(90 days total), this subscription "
|
||||
"will be automatically removed."
|
||||
)
|
||||
|
||||
|
||||
class TestSendCleanupDeletion:
|
||||
"""Cleanup deletion notification sending."""
|
||||
|
||||
async def test_sends_deletion_to_all_rooms(
|
||||
self,
|
||||
stream_repo: StreamRepository,
|
||||
subscription_repo: SubscriptionRepository,
|
||||
) -> None:
|
||||
"""Send cleanup deletion notice to all subscribed rooms."""
|
||||
client = _StubMatrixClient()
|
||||
service = _make_service(client=client, subscription_repo=subscription_repo)
|
||||
|
||||
await stream_repo.create("example.com")
|
||||
await subscription_repo.add("example.com", "!room1:matrix.org")
|
||||
|
||||
await service.send_cleanup_deletion("example.com")
|
||||
|
||||
assert len(client.sent_messages) == 1
|
||||
assert client.sent_messages[0].content.body == (
|
||||
"🗑️ Subscription Automatically Removed\n"
|
||||
"\n"
|
||||
"The Owncast instance at example.com has been "
|
||||
"unreachable for 90 days and has been "
|
||||
"automatically removed from subscriptions in this "
|
||||
"room.\n"
|
||||
"\n"
|
||||
"If the instance comes online again and you want to "
|
||||
"resubscribe, run `!subscribe example.com`."
|
||||
)
|
||||
@@ -0,0 +1,207 @@
|
||||
# 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 Owncast HTTP client."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
import logging
|
||||
from typing import TYPE_CHECKING
|
||||
|
||||
import pytest
|
||||
from aioresponses import aioresponses
|
||||
|
||||
from owncastsentry.owncast_client import OwncastClient
|
||||
from tests.conftest import VALID_CONFIG_RESPONSE, VALID_STATUS_RESPONSE
|
||||
|
||||
if TYPE_CHECKING:
|
||||
from collections.abc import AsyncIterator
|
||||
|
||||
|
||||
@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")
|
||||
yield client
|
||||
await client.close()
|
||||
|
||||
|
||||
class TestGetStreamState:
|
||||
"""Stream state retrieval from the status API."""
|
||||
|
||||
async def test_returns_state_on_success(
|
||||
self, owncast_client: OwncastClient
|
||||
) -> None:
|
||||
"""Return a StreamState with correct fields on a valid 200 response."""
|
||||
with aioresponses() as mocked:
|
||||
mocked.get(
|
||||
"https://stream.logal.dev/api/status",
|
||||
body=json.dumps(VALID_STATUS_RESPONSE).encode(),
|
||||
)
|
||||
result = await owncast_client.get_stream_state("stream.logal.dev")
|
||||
|
||||
assert result is not None
|
||||
assert result.domain == "stream.logal.dev"
|
||||
assert (
|
||||
result.title
|
||||
== "I think I can do this... Let's start a nuclear reaction - Playing Nucleares!" # noqa: E501
|
||||
)
|
||||
assert result.last_connect_time is None
|
||||
assert result.last_disconnect_time == "2026-03-04T21:05:32-05:00"
|
||||
|
||||
async def test_returns_none_on_missing_field(
|
||||
self, owncast_client: OwncastClient
|
||||
) -> None:
|
||||
"""Return None when the response is missing required fields."""
|
||||
incomplete = {"streamTitle": "Test Stream", "online": True}
|
||||
with aioresponses() as mocked:
|
||||
mocked.get(
|
||||
"https://stream.logal.dev/api/status",
|
||||
body=json.dumps(incomplete).encode(),
|
||||
)
|
||||
result = await owncast_client.get_stream_state("stream.logal.dev")
|
||||
|
||||
assert result is None
|
||||
|
||||
async def test_returns_none_on_invalid_json(
|
||||
self, owncast_client: OwncastClient
|
||||
) -> None:
|
||||
"""Return None when the response body is not valid JSON."""
|
||||
with aioresponses() as mocked:
|
||||
mocked.get(
|
||||
"https://stream.logal.dev/api/status",
|
||||
body=b"not json",
|
||||
)
|
||||
result = await owncast_client.get_stream_state("stream.logal.dev")
|
||||
|
||||
assert result is None
|
||||
|
||||
async def test_returns_none_on_non_200(self, owncast_client: OwncastClient) -> None:
|
||||
"""Return None when the response status is not 200."""
|
||||
with aioresponses() as mocked:
|
||||
mocked.get(
|
||||
"https://stream.logal.dev/api/status",
|
||||
status=404,
|
||||
)
|
||||
result = await owncast_client.get_stream_state("stream.logal.dev")
|
||||
|
||||
assert result is None
|
||||
|
||||
async def test_returns_none_on_connection_error(
|
||||
self, owncast_client: OwncastClient
|
||||
) -> None:
|
||||
"""Return None when a connection error occurs."""
|
||||
with aioresponses() as mocked:
|
||||
mocked.get(
|
||||
"https://stream.logal.dev/api/status",
|
||||
exception=ConnectionError(),
|
||||
)
|
||||
result = await owncast_client.get_stream_state("stream.logal.dev")
|
||||
|
||||
assert result is None
|
||||
|
||||
|
||||
class TestGetStreamConfig:
|
||||
"""Stream configuration retrieval from the config API."""
|
||||
|
||||
async def test_returns_config_on_success(
|
||||
self, owncast_client: OwncastClient
|
||||
) -> None:
|
||||
"""Return a StreamConfig with correct fields on a valid 200 response."""
|
||||
with aioresponses() as mocked:
|
||||
mocked.get(
|
||||
"https://stream.logal.dev/api/config",
|
||||
body=json.dumps(VALID_CONFIG_RESPONSE).encode(),
|
||||
)
|
||||
result = await owncast_client.get_stream_config("stream.logal.dev")
|
||||
|
||||
assert result is not None
|
||||
assert result.name == "LogalDeveloper's Live Stream"
|
||||
assert result.tags == [
|
||||
"video games",
|
||||
"chatting",
|
||||
"casual",
|
||||
"english",
|
||||
"streaming",
|
||||
"owncast",
|
||||
"variety",
|
||||
]
|
||||
|
||||
async def test_returns_none_on_invalid_json(
|
||||
self, owncast_client: OwncastClient
|
||||
) -> None:
|
||||
"""Return None when the response body is not valid JSON."""
|
||||
with aioresponses() as mocked:
|
||||
mocked.get(
|
||||
"https://stream.logal.dev/api/config",
|
||||
body=b"not json",
|
||||
)
|
||||
result = await owncast_client.get_stream_config("stream.logal.dev")
|
||||
|
||||
assert result is None
|
||||
|
||||
async def test_returns_none_on_non_200(self, owncast_client: OwncastClient) -> None:
|
||||
"""Return None when the response status is not 200."""
|
||||
with aioresponses() as mocked:
|
||||
mocked.get(
|
||||
"https://stream.logal.dev/api/config",
|
||||
status=500,
|
||||
)
|
||||
result = await owncast_client.get_stream_config("stream.logal.dev")
|
||||
|
||||
assert result is None
|
||||
|
||||
async def test_returns_none_on_connection_error(
|
||||
self, owncast_client: OwncastClient
|
||||
) -> None:
|
||||
"""Return None when a connection error occurs."""
|
||||
with aioresponses() as mocked:
|
||||
mocked.get(
|
||||
"https://stream.logal.dev/api/config",
|
||||
exception=ConnectionError(),
|
||||
)
|
||||
result = await owncast_client.get_stream_config("stream.logal.dev")
|
||||
|
||||
assert result is None
|
||||
|
||||
|
||||
class TestValidateInstance:
|
||||
"""Owncast instance validation via the status API."""
|
||||
|
||||
async def test_returns_true_for_valid_instance(
|
||||
self, owncast_client: OwncastClient
|
||||
) -> None:
|
||||
"""Return True when the status endpoint returns a valid response."""
|
||||
with aioresponses() as mocked:
|
||||
mocked.get(
|
||||
"https://valid.com/api/status",
|
||||
body=json.dumps(VALID_STATUS_RESPONSE).encode(),
|
||||
)
|
||||
result = await owncast_client.validate_instance("valid.com")
|
||||
|
||||
assert result is True
|
||||
|
||||
async def test_returns_false_for_invalid_instance(
|
||||
self, owncast_client: OwncastClient
|
||||
) -> None:
|
||||
"""Return False when the status endpoint returns a non-200 response."""
|
||||
with aioresponses() as mocked:
|
||||
mocked.get(
|
||||
"https://invalid.com/api/status",
|
||||
status=404,
|
||||
)
|
||||
result = await owncast_client.validate_instance("invalid.com")
|
||||
|
||||
assert result is False
|
||||
@@ -0,0 +1,839 @@
|
||||
# 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 stream monitor."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import logging
|
||||
import time
|
||||
from typing import TYPE_CHECKING
|
||||
|
||||
from owncastsentry.models import StreamConfig, StreamState
|
||||
from owncastsentry.notification_service import NotificationService
|
||||
from owncastsentry.stream_monitor import StreamMonitor
|
||||
from owncastsentry.utils import (
|
||||
CLEANUP_DELETE_THRESHOLD,
|
||||
CLEANUP_WARNING_THRESHOLD,
|
||||
SECONDS_BETWEEN_NOTIFICATIONS,
|
||||
)
|
||||
from tests.conftest import _StubMatrixClient, _StubOwncastClient
|
||||
|
||||
if TYPE_CHECKING:
|
||||
from owncastsentry.database import StreamRepository, SubscriptionRepository
|
||||
|
||||
|
||||
def _make_monitor(
|
||||
*,
|
||||
owncast_client: _StubOwncastClient,
|
||||
stream_repo: StreamRepository,
|
||||
subscription_repo: SubscriptionRepository,
|
||||
client: _StubMatrixClient,
|
||||
) -> tuple[StreamMonitor, NotificationService]:
|
||||
"""Build a StreamMonitor with stubs and a real NotificationService."""
|
||||
logger = logging.getLogger("test")
|
||||
notification_service = NotificationService(
|
||||
client=client,
|
||||
subscription_repo=subscription_repo,
|
||||
logger=logger,
|
||||
)
|
||||
monitor = StreamMonitor(
|
||||
owncast_client=owncast_client,
|
||||
stream_repo=stream_repo,
|
||||
subscription_repo=subscription_repo,
|
||||
notification_service=notification_service,
|
||||
logger=logger,
|
||||
)
|
||||
return monitor, notification_service
|
||||
|
||||
|
||||
async def _seed_stream(
|
||||
stream_repo: StreamRepository,
|
||||
subscription_repo: SubscriptionRepository,
|
||||
*,
|
||||
domain: str = "example.com",
|
||||
room_id: str = "!room:matrix.org",
|
||||
name: str | None = "Test Stream",
|
||||
title: str | None = "Test Title",
|
||||
last_connect_time: str | None = None,
|
||||
last_disconnect_time: str | None = None,
|
||||
) -> None:
|
||||
"""Insert a stream and subscription into the database."""
|
||||
await stream_repo.create(domain)
|
||||
state = StreamState(
|
||||
domain=domain,
|
||||
name=name,
|
||||
title=title,
|
||||
last_connect_time=last_connect_time,
|
||||
last_disconnect_time=last_disconnect_time,
|
||||
)
|
||||
await stream_repo.update(state)
|
||||
await subscription_repo.add(domain, room_id)
|
||||
|
||||
|
||||
class TestUpdateAllStreams:
|
||||
"""Parallel stream update orchestration."""
|
||||
|
||||
async def test_returns_correct_counts(
|
||||
self,
|
||||
stream_repo: StreamRepository,
|
||||
subscription_repo: SubscriptionRepository,
|
||||
) -> None:
|
||||
"""Return an UpdateResult with correct success and failure counts."""
|
||||
owncast = _StubOwncastClient(
|
||||
stream_state=StreamState(
|
||||
domain="ok.com", last_connect_time="2026-01-01T00:00:00Z"
|
||||
),
|
||||
stream_config=StreamConfig(name="OK"),
|
||||
)
|
||||
client = _StubMatrixClient()
|
||||
monitor, _ = _make_monitor(
|
||||
owncast_client=owncast,
|
||||
stream_repo=stream_repo,
|
||||
subscription_repo=subscription_repo,
|
||||
client=client,
|
||||
)
|
||||
|
||||
await _seed_stream(stream_repo, subscription_repo, domain="ok.com")
|
||||
|
||||
result = await monitor.update_all_streams(["ok.com"])
|
||||
assert result.total_streams == 1
|
||||
assert result.successful_checks == 1
|
||||
assert result.failed_checks == 0
|
||||
assert owncast.queried_domains == ["ok.com"]
|
||||
|
||||
|
||||
class TestUpdateStreamBackoff:
|
||||
"""Backoff logic that skips queries for failing streams."""
|
||||
|
||||
async def test_skips_query_on_backoff(
|
||||
self,
|
||||
stream_repo: StreamRepository,
|
||||
subscription_repo: SubscriptionRepository,
|
||||
) -> None:
|
||||
"""Skip the HTTP query and increment the counter during backoff."""
|
||||
owncast = _StubOwncastClient()
|
||||
client = _StubMatrixClient()
|
||||
monitor, _ = _make_monitor(
|
||||
owncast_client=owncast,
|
||||
stream_repo=stream_repo,
|
||||
subscription_repo=subscription_repo,
|
||||
client=client,
|
||||
)
|
||||
|
||||
await _seed_stream(stream_repo, subscription_repo, domain="fail.com")
|
||||
|
||||
# Set failure counter to 5 (should skip odd cycles)
|
||||
for _ in range(5):
|
||||
await stream_repo.increment_failure_counter("fail.com")
|
||||
|
||||
result = await monitor.update_stream("fail.com")
|
||||
assert result is True
|
||||
assert owncast.state_call_count == 0
|
||||
|
||||
# Counter should have been incremented to 6
|
||||
state = await stream_repo.get_by_domain("fail.com")
|
||||
assert state is not None
|
||||
assert state.failure_counter == 6
|
||||
|
||||
|
||||
class TestUpdateStreamFirstUpdate:
|
||||
"""First state update of a newly subscribed stream."""
|
||||
|
||||
async def test_no_notification_on_first_live(
|
||||
self,
|
||||
stream_repo: StreamRepository,
|
||||
subscription_repo: SubscriptionRepository,
|
||||
) -> None:
|
||||
"""Skip notifications when a stream is already live on first update."""
|
||||
owncast = _StubOwncastClient(
|
||||
stream_state=StreamState(
|
||||
domain="new.com",
|
||||
last_connect_time="2026-01-01T00:00:00Z",
|
||||
last_disconnect_time="2025-12-31T00:00:00Z",
|
||||
),
|
||||
stream_config=StreamConfig(name="New Stream"),
|
||||
)
|
||||
client = _StubMatrixClient()
|
||||
monitor, _ = _make_monitor(
|
||||
owncast_client=owncast,
|
||||
stream_repo=stream_repo,
|
||||
subscription_repo=subscription_repo,
|
||||
client=client,
|
||||
)
|
||||
|
||||
# Seed with no connect/disconnect times (brand new)
|
||||
await stream_repo.create("new.com")
|
||||
await subscription_repo.add("new.com", "!room:matrix.org")
|
||||
|
||||
result = await monitor.update_stream("new.com")
|
||||
assert result is True
|
||||
assert len(client.sent_messages) == 0
|
||||
|
||||
async def test_no_notification_on_first_offline(
|
||||
self,
|
||||
stream_repo: StreamRepository,
|
||||
subscription_repo: SubscriptionRepository,
|
||||
) -> None:
|
||||
"""Skip notifications when a stream is offline on first update."""
|
||||
owncast = _StubOwncastClient(
|
||||
stream_state=StreamState(
|
||||
domain="new.com",
|
||||
last_disconnect_time="2025-12-31T00:00:00Z",
|
||||
),
|
||||
stream_config=StreamConfig(name="New Stream"),
|
||||
)
|
||||
client = _StubMatrixClient()
|
||||
monitor, _ = _make_monitor(
|
||||
owncast_client=owncast,
|
||||
stream_repo=stream_repo,
|
||||
subscription_repo=subscription_repo,
|
||||
client=client,
|
||||
)
|
||||
|
||||
await stream_repo.create("new.com")
|
||||
await subscription_repo.add("new.com", "!room:matrix.org")
|
||||
|
||||
result = await monitor.update_stream("new.com")
|
||||
assert result is True
|
||||
assert len(client.sent_messages) == 0
|
||||
|
||||
|
||||
class TestUpdateStreamGoesLive:
|
||||
"""Stream transitioning from offline to online."""
|
||||
|
||||
async def test_sends_go_live_notification(
|
||||
self,
|
||||
stream_repo: StreamRepository,
|
||||
subscription_repo: SubscriptionRepository,
|
||||
) -> None:
|
||||
"""Send go-live with stream name and tags."""
|
||||
owncast = _StubOwncastClient(
|
||||
stream_state=StreamState(
|
||||
domain="live.com",
|
||||
title="Now Streaming",
|
||||
last_connect_time="2026-01-01T12:00:00Z",
|
||||
last_disconnect_time="2026-01-01T10:00:00Z",
|
||||
),
|
||||
stream_config=StreamConfig(name="Live Stream", tags=["gaming"]),
|
||||
)
|
||||
client = _StubMatrixClient()
|
||||
monitor, _ = _make_monitor(
|
||||
owncast_client=owncast,
|
||||
stream_repo=stream_repo,
|
||||
subscription_repo=subscription_repo,
|
||||
client=client,
|
||||
)
|
||||
|
||||
# Seed as offline (has disconnect but no connect)
|
||||
await _seed_stream(
|
||||
stream_repo,
|
||||
subscription_repo,
|
||||
domain="live.com",
|
||||
last_disconnect_time="2026-01-01T10:00:00Z",
|
||||
)
|
||||
|
||||
# Set offline timer to long ago so it's not a brief outage
|
||||
monitor.offline_timer_cache["live.com"] = 0
|
||||
|
||||
result = await monitor.update_stream("live.com")
|
||||
assert result is True
|
||||
assert len(client.sent_messages) == 1
|
||||
assert client.sent_messages[0].content.body == (
|
||||
"🎥 Live Stream is now live!\n"
|
||||
"Stream Title: Now Streaming\n"
|
||||
"\n"
|
||||
"To tune in, visit: https://live.com/\n"
|
||||
"\n"
|
||||
"#gaming"
|
||||
)
|
||||
|
||||
async def test_falls_back_when_config_fetch_fails(
|
||||
self,
|
||||
stream_repo: StreamRepository,
|
||||
subscription_repo: SubscriptionRepository,
|
||||
) -> None:
|
||||
"""Use domain as name in go-live when config fetch fails."""
|
||||
owncast = _StubOwncastClient(
|
||||
stream_state=StreamState(
|
||||
domain="live.com",
|
||||
title="Now Streaming",
|
||||
last_connect_time="2026-01-01T12:00:00Z",
|
||||
last_disconnect_time="2026-01-01T10:00:00Z",
|
||||
),
|
||||
stream_config=None,
|
||||
)
|
||||
client = _StubMatrixClient()
|
||||
monitor, _ = _make_monitor(
|
||||
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
|
||||
|
||||
result = await monitor.update_stream("live.com")
|
||||
assert result is True
|
||||
assert len(client.sent_messages) == 1
|
||||
assert client.sent_messages[0].content.body == (
|
||||
"🎥 live.com is now live!\n"
|
||||
"Stream Title: Now Streaming\n"
|
||||
"\n"
|
||||
"To tune in, visit: https://live.com/"
|
||||
)
|
||||
|
||||
|
||||
class TestUpdateStreamBriefOffline:
|
||||
"""Stream that was briefly offline and comes back."""
|
||||
|
||||
async def test_no_notification_same_title(
|
||||
self,
|
||||
stream_repo: StreamRepository,
|
||||
subscription_repo: SubscriptionRepository,
|
||||
) -> None:
|
||||
"""Skip notification for a brief outage with no title change."""
|
||||
owncast = _StubOwncastClient(
|
||||
stream_state=StreamState(
|
||||
domain="brief.com",
|
||||
title="Same Title",
|
||||
last_connect_time="2026-01-01T12:00:00Z",
|
||||
),
|
||||
stream_config=StreamConfig(name="Brief Stream"),
|
||||
)
|
||||
client = _StubMatrixClient()
|
||||
monitor, _ = _make_monitor(
|
||||
owncast_client=owncast,
|
||||
stream_repo=stream_repo,
|
||||
subscription_repo=subscription_repo,
|
||||
client=client,
|
||||
)
|
||||
|
||||
# Seed as offline with same title
|
||||
await _seed_stream(
|
||||
stream_repo,
|
||||
subscription_repo,
|
||||
domain="brief.com",
|
||||
title="Same Title",
|
||||
last_disconnect_time="2026-01-01T11:55:00Z",
|
||||
)
|
||||
|
||||
# Recently offline (within cooldown)
|
||||
monitor.offline_timer_cache["brief.com"] = time.time() - 60
|
||||
|
||||
result = await monitor.update_stream("brief.com")
|
||||
assert result is True
|
||||
assert len(client.sent_messages) == 0
|
||||
|
||||
async def test_title_change_notification(
|
||||
self,
|
||||
stream_repo: StreamRepository,
|
||||
subscription_repo: SubscriptionRepository,
|
||||
) -> None:
|
||||
"""Send a title-change notification for a brief outage with a new title."""
|
||||
owncast = _StubOwncastClient(
|
||||
stream_state=StreamState(
|
||||
domain="brief.com",
|
||||
title="New Title",
|
||||
last_connect_time="2026-01-01T12:00:00Z",
|
||||
),
|
||||
stream_config=StreamConfig(name="Brief Stream"),
|
||||
)
|
||||
client = _StubMatrixClient()
|
||||
monitor, _ = _make_monitor(
|
||||
owncast_client=owncast,
|
||||
stream_repo=stream_repo,
|
||||
subscription_repo=subscription_repo,
|
||||
client=client,
|
||||
)
|
||||
|
||||
# Seed as offline with old title
|
||||
await _seed_stream(
|
||||
stream_repo,
|
||||
subscription_repo,
|
||||
domain="brief.com",
|
||||
title="Old Title",
|
||||
last_disconnect_time="2026-01-01T11:55:00Z",
|
||||
)
|
||||
|
||||
# Recently offline (within cooldown)
|
||||
monitor.offline_timer_cache["brief.com"] = time.time() - 60
|
||||
|
||||
result = await monitor.update_stream("brief.com")
|
||||
assert result is True
|
||||
assert len(client.sent_messages) == 1
|
||||
assert client.sent_messages[0].content.body == (
|
||||
"📝 Brief Stream has changed its stream title!\n"
|
||||
"Stream Title: New Title\n"
|
||||
"\n"
|
||||
"To tune in, visit: https://brief.com/"
|
||||
)
|
||||
|
||||
|
||||
class TestUpdateStreamTitleChange:
|
||||
"""Mid-session title change while stream stays online."""
|
||||
|
||||
async def test_falls_back_when_config_fetch_fails(
|
||||
self,
|
||||
stream_repo: StreamRepository,
|
||||
subscription_repo: SubscriptionRepository,
|
||||
) -> None:
|
||||
"""Use domain as name in title-change when config fails."""
|
||||
owncast = _StubOwncastClient(
|
||||
stream_state=StreamState(
|
||||
domain="title.com",
|
||||
title="Updated Title",
|
||||
last_connect_time="2026-01-01T12:00:00Z",
|
||||
),
|
||||
stream_config=None,
|
||||
)
|
||||
client = _StubMatrixClient()
|
||||
monitor, notification_service = _make_monitor(
|
||||
owncast_client=owncast,
|
||||
stream_repo=stream_repo,
|
||||
subscription_repo=subscription_repo,
|
||||
client=client,
|
||||
)
|
||||
|
||||
await _seed_stream(
|
||||
stream_repo,
|
||||
subscription_repo,
|
||||
domain="title.com",
|
||||
title="Original Title",
|
||||
last_connect_time="2026-01-01T12:00:00Z",
|
||||
)
|
||||
|
||||
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
|
||||
)
|
||||
|
||||
result = await monitor.update_stream("title.com")
|
||||
assert result is True
|
||||
assert len(client.sent_messages) == 1
|
||||
assert client.sent_messages[0].content.body == (
|
||||
"📝 title.com has changed its stream title!\n"
|
||||
"Stream Title: Updated Title\n"
|
||||
"\n"
|
||||
"To tune in, visit: https://title.com/"
|
||||
)
|
||||
|
||||
async def test_sends_title_change(
|
||||
self,
|
||||
stream_repo: StreamRepository,
|
||||
subscription_repo: SubscriptionRepository,
|
||||
) -> None:
|
||||
"""Send a title-change notification when the title changes mid-stream."""
|
||||
owncast = _StubOwncastClient(
|
||||
stream_state=StreamState(
|
||||
domain="title.com",
|
||||
title="Updated Title",
|
||||
last_connect_time="2026-01-01T12:00:00Z",
|
||||
),
|
||||
stream_config=StreamConfig(name="Title Stream"),
|
||||
)
|
||||
client = _StubMatrixClient()
|
||||
monitor, notification_service = _make_monitor(
|
||||
owncast_client=owncast,
|
||||
stream_repo=stream_repo,
|
||||
subscription_repo=subscription_repo,
|
||||
client=client,
|
||||
)
|
||||
|
||||
# Seed as online with different title
|
||||
await _seed_stream(
|
||||
stream_repo,
|
||||
subscription_repo,
|
||||
domain="title.com",
|
||||
title="Original Title",
|
||||
last_connect_time="2026-01-01T12:00:00Z",
|
||||
)
|
||||
|
||||
# Last notification was long enough ago to pass rate limiting,
|
||||
# but more recent than the offline timer (so title-change fires)
|
||||
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
|
||||
)
|
||||
|
||||
result = await monitor.update_stream("title.com")
|
||||
assert result is True
|
||||
assert len(client.sent_messages) == 1
|
||||
assert client.sent_messages[0].content.body == (
|
||||
"📝 Title Stream has changed its stream title!\n"
|
||||
"Stream Title: Updated Title\n"
|
||||
"\n"
|
||||
"To tune in, visit: https://title.com/"
|
||||
)
|
||||
|
||||
async def test_sends_go_live_when_last_notification_before_offline(
|
||||
self,
|
||||
stream_repo: StreamRepository,
|
||||
subscription_repo: SubscriptionRepository,
|
||||
) -> None:
|
||||
"""Send go-live instead of title-change after an offline gap."""
|
||||
owncast = _StubOwncastClient(
|
||||
stream_state=StreamState(
|
||||
domain="title.com",
|
||||
title="Updated Title",
|
||||
last_connect_time="2026-01-01T12:00:00Z",
|
||||
),
|
||||
stream_config=StreamConfig(name="Title Stream"),
|
||||
)
|
||||
client = _StubMatrixClient()
|
||||
monitor, notification_service = _make_monitor(
|
||||
owncast_client=owncast,
|
||||
stream_repo=stream_repo,
|
||||
subscription_repo=subscription_repo,
|
||||
client=client,
|
||||
)
|
||||
|
||||
await _seed_stream(
|
||||
stream_repo,
|
||||
subscription_repo,
|
||||
domain="title.com",
|
||||
title="Original Title",
|
||||
last_connect_time="2026-01-01T12:00:00Z",
|
||||
)
|
||||
|
||||
# Offline timer is MORE recent than last notification,
|
||||
# and both are old enough to pass rate limiting
|
||||
now = time.time()
|
||||
monitor.offline_timer_cache["title.com"] = (
|
||||
now - SECONDS_BETWEEN_NOTIFICATIONS - 100
|
||||
)
|
||||
notification_service.notification_timers_cache["title.com"] = (
|
||||
now - SECONDS_BETWEEN_NOTIFICATIONS - 200
|
||||
)
|
||||
|
||||
result = await monitor.update_stream("title.com")
|
||||
assert result is True
|
||||
assert len(client.sent_messages) == 1
|
||||
assert client.sent_messages[0].content.body == (
|
||||
"🎥 Title Stream is now live!\n"
|
||||
"Stream Title: Updated Title\n"
|
||||
"\n"
|
||||
"To tune in, visit: https://title.com/"
|
||||
)
|
||||
|
||||
|
||||
class TestUpdateStreamGoesOffline:
|
||||
"""Stream transitioning from online to offline."""
|
||||
|
||||
async def test_records_offline_time(
|
||||
self,
|
||||
stream_repo: StreamRepository,
|
||||
subscription_repo: SubscriptionRepository,
|
||||
) -> None:
|
||||
"""Record the offline time in the cache when a stream goes offline."""
|
||||
owncast = _StubOwncastClient(
|
||||
stream_state=StreamState(
|
||||
domain="offline.com",
|
||||
title="Title",
|
||||
last_disconnect_time="2026-01-01T12:00:00Z",
|
||||
),
|
||||
stream_config=StreamConfig(name="Offline Stream"),
|
||||
)
|
||||
client = _StubMatrixClient()
|
||||
monitor, _ = _make_monitor(
|
||||
owncast_client=owncast,
|
||||
stream_repo=stream_repo,
|
||||
subscription_repo=subscription_repo,
|
||||
client=client,
|
||||
)
|
||||
|
||||
# Seed as online
|
||||
await _seed_stream(
|
||||
stream_repo,
|
||||
subscription_repo,
|
||||
domain="offline.com",
|
||||
last_connect_time="2026-01-01T10:00:00Z",
|
||||
)
|
||||
monitor.offline_timer_cache["offline.com"] = 0
|
||||
|
||||
before = time.time()
|
||||
result = await monitor.update_stream("offline.com")
|
||||
after = time.time()
|
||||
|
||||
assert result is True
|
||||
assert before <= monitor.offline_timer_cache["offline.com"] <= after
|
||||
assert len(client.sent_messages) == 0
|
||||
|
||||
|
||||
class TestUpdateStreamConnectionFailure:
|
||||
"""Connection failure handling."""
|
||||
|
||||
async def test_returns_false_and_increments_counter(
|
||||
self,
|
||||
stream_repo: StreamRepository,
|
||||
subscription_repo: SubscriptionRepository,
|
||||
) -> None:
|
||||
"""Return False and increment the failure counter on connection failure."""
|
||||
owncast = _StubOwncastClient(stream_state=None)
|
||||
client = _StubMatrixClient()
|
||||
monitor, _ = _make_monitor(
|
||||
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",
|
||||
)
|
||||
|
||||
result = await monitor.update_stream("fail.com")
|
||||
assert result is False
|
||||
|
||||
state = await stream_repo.get_by_domain("fail.com")
|
||||
assert state is not None
|
||||
assert state.failure_counter == 1
|
||||
|
||||
|
||||
class TestCheckCleanupThresholds:
|
||||
"""Auto-cleanup warning and deletion thresholds."""
|
||||
|
||||
async def test_sends_warning_at_threshold(
|
||||
self,
|
||||
stream_repo: StreamRepository,
|
||||
subscription_repo: SubscriptionRepository,
|
||||
) -> None:
|
||||
"""Send a cleanup warning when the failure counter hits the 83-day threshold."""
|
||||
owncast = _StubOwncastClient()
|
||||
client = _StubMatrixClient()
|
||||
monitor, _ = _make_monitor(
|
||||
owncast_client=owncast,
|
||||
stream_repo=stream_repo,
|
||||
subscription_repo=subscription_repo,
|
||||
client=client,
|
||||
)
|
||||
|
||||
await _seed_stream(stream_repo, subscription_repo, domain="warn.com")
|
||||
|
||||
await monitor._check_cleanup_thresholds("warn.com", CLEANUP_WARNING_THRESHOLD)
|
||||
|
||||
assert len(client.sent_messages) == 1
|
||||
assert client.sent_messages[0].content.body == (
|
||||
"⚠️ Warning: Subscription Cleanup Scheduled\n"
|
||||
"\n"
|
||||
"The Owncast instance at warn.com has been "
|
||||
"unreachable for 83 days. If it remains "
|
||||
"unreachable for 7 more days "
|
||||
"(90 days total), this subscription "
|
||||
"will be automatically removed."
|
||||
)
|
||||
|
||||
async def test_deletes_at_threshold(
|
||||
self,
|
||||
stream_repo: StreamRepository,
|
||||
subscription_repo: SubscriptionRepository,
|
||||
) -> None:
|
||||
"""Delete all subscriptions and the stream record at the 90-day threshold."""
|
||||
owncast = _StubOwncastClient()
|
||||
client = _StubMatrixClient()
|
||||
monitor, _ = _make_monitor(
|
||||
owncast_client=owncast,
|
||||
stream_repo=stream_repo,
|
||||
subscription_repo=subscription_repo,
|
||||
client=client,
|
||||
)
|
||||
|
||||
await _seed_stream(stream_repo, subscription_repo, domain="delete.com")
|
||||
|
||||
await monitor._check_cleanup_thresholds("delete.com", CLEANUP_DELETE_THRESHOLD)
|
||||
|
||||
# Deletion notification sent
|
||||
assert len(client.sent_messages) == 1
|
||||
assert client.sent_messages[0].content.body == (
|
||||
"🗑️ Subscription Automatically Removed\n"
|
||||
"\n"
|
||||
"The Owncast instance at delete.com has been "
|
||||
"unreachable for 90 days and has been "
|
||||
"automatically removed from subscriptions in this "
|
||||
"room.\n"
|
||||
"\n"
|
||||
"If the instance comes online again and you want to "
|
||||
"resubscribe, run `!subscribe delete.com`."
|
||||
)
|
||||
|
||||
# Stream and subscriptions removed from DB
|
||||
assert await stream_repo.get_by_domain("delete.com") is None
|
||||
rooms = await subscription_repo.get_subscribed_rooms("delete.com")
|
||||
assert rooms == []
|
||||
|
||||
async def test_no_action_below_thresholds(
|
||||
self,
|
||||
stream_repo: StreamRepository,
|
||||
subscription_repo: SubscriptionRepository,
|
||||
) -> None:
|
||||
"""Take no action when the counter is below both thresholds."""
|
||||
owncast = _StubOwncastClient()
|
||||
client = _StubMatrixClient()
|
||||
monitor, _ = _make_monitor(
|
||||
owncast_client=owncast,
|
||||
stream_repo=stream_repo,
|
||||
subscription_repo=subscription_repo,
|
||||
client=client,
|
||||
)
|
||||
|
||||
await _seed_stream(stream_repo, subscription_repo, domain="ok.com")
|
||||
|
||||
await monitor._check_cleanup_thresholds("ok.com", 100)
|
||||
|
||||
assert len(client.sent_messages) == 0
|
||||
assert await stream_repo.get_by_domain("ok.com") is not None
|
||||
|
||||
|
||||
class TestUpdateStreamNoStateChange:
|
||||
"""Stream that stays offline across updates with no state change."""
|
||||
|
||||
async def test_skips_database_write(
|
||||
self,
|
||||
stream_repo: StreamRepository,
|
||||
subscription_repo: SubscriptionRepository,
|
||||
) -> None:
|
||||
"""Skip the database write when stream state has not changed."""
|
||||
owncast = _StubOwncastClient(
|
||||
stream_state=StreamState(
|
||||
domain="stable.com",
|
||||
title="Same Title",
|
||||
last_disconnect_time="2026-01-01T12:00:00Z",
|
||||
),
|
||||
)
|
||||
client = _StubMatrixClient()
|
||||
monitor, _ = _make_monitor(
|
||||
owncast_client=owncast,
|
||||
stream_repo=stream_repo,
|
||||
subscription_repo=subscription_repo,
|
||||
client=client,
|
||||
)
|
||||
|
||||
# Seed as offline with same disconnect time and title
|
||||
await _seed_stream(
|
||||
stream_repo,
|
||||
subscription_repo,
|
||||
domain="stable.com",
|
||||
title="Same Title",
|
||||
last_disconnect_time="2026-01-01T12:00:00Z",
|
||||
)
|
||||
|
||||
result = await monitor.update_stream("stable.com")
|
||||
assert result is True
|
||||
assert len(client.sent_messages) == 0
|
||||
# Config should not have been fetched since no DB update was needed
|
||||
assert owncast.config_call_count == 0
|
||||
|
||||
|
||||
class TestUpdateStreamFailureCounterReset:
|
||||
"""Failure counter reset on successful fetch."""
|
||||
|
||||
async def test_resets_counter_on_success(
|
||||
self,
|
||||
stream_repo: StreamRepository,
|
||||
subscription_repo: SubscriptionRepository,
|
||||
) -> None:
|
||||
"""Reset the failure counter to zero after a successful fetch."""
|
||||
owncast = _StubOwncastClient(
|
||||
stream_state=StreamState(
|
||||
domain="recover.com",
|
||||
title="Title",
|
||||
last_disconnect_time="2026-01-01T12:00:00Z",
|
||||
),
|
||||
)
|
||||
client = _StubMatrixClient()
|
||||
monitor, _ = _make_monitor(
|
||||
owncast_client=owncast,
|
||||
stream_repo=stream_repo,
|
||||
subscription_repo=subscription_repo,
|
||||
client=client,
|
||||
)
|
||||
|
||||
await _seed_stream(
|
||||
stream_repo,
|
||||
subscription_repo,
|
||||
domain="recover.com",
|
||||
title="Title",
|
||||
last_disconnect_time="2026-01-01T12:00:00Z",
|
||||
)
|
||||
|
||||
# Simulate prior failures (counter=4 still passes backoff)
|
||||
for _ in range(4):
|
||||
await stream_repo.increment_failure_counter("recover.com")
|
||||
|
||||
state = await stream_repo.get_by_domain("recover.com")
|
||||
assert state is not None
|
||||
assert state.failure_counter == 4
|
||||
|
||||
result = await monitor.update_stream("recover.com")
|
||||
assert result is True
|
||||
|
||||
state = await stream_repo.get_by_domain("recover.com")
|
||||
assert state is not None
|
||||
assert state.failure_counter == 0
|
||||
|
||||
|
||||
class TestUpdateAllStreamsMixed:
|
||||
"""Mixed success/failure results from parallel updates."""
|
||||
|
||||
async def test_counts_mixed_results(
|
||||
self,
|
||||
stream_repo: StreamRepository,
|
||||
subscription_repo: SubscriptionRepository,
|
||||
) -> None:
|
||||
"""Return correct counts with one failure and one backoff skip."""
|
||||
# Stub returns None, so all actual fetches fail
|
||||
owncast = _StubOwncastClient(stream_state=None)
|
||||
client = _StubMatrixClient()
|
||||
monitor, _ = _make_monitor(
|
||||
owncast_client=owncast,
|
||||
stream_repo=stream_repo,
|
||||
subscription_repo=subscription_repo,
|
||||
client=client,
|
||||
)
|
||||
|
||||
# "fail.com" will be fetched and fail (counter=0, no backoff)
|
||||
await _seed_stream(
|
||||
stream_repo,
|
||||
subscription_repo,
|
||||
domain="fail.com",
|
||||
last_disconnect_time="2026-01-01T00:00:00Z",
|
||||
)
|
||||
|
||||
# "skip.com" will be skipped via backoff (returns True)
|
||||
await _seed_stream(
|
||||
stream_repo,
|
||||
subscription_repo,
|
||||
domain="skip.com",
|
||||
room_id="!room2:matrix.org",
|
||||
last_disconnect_time="2026-01-01T00:00:00Z",
|
||||
)
|
||||
for _ in range(5):
|
||||
await stream_repo.increment_failure_counter("skip.com")
|
||||
|
||||
result = await monitor.update_all_streams(["fail.com", "skip.com"])
|
||||
assert result.total_streams == 2
|
||||
assert result.successful_checks == 1
|
||||
assert result.failed_checks == 1
|
||||
@@ -0,0 +1,199 @@
|
||||
# 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 utility functions and constants."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import pytest
|
||||
|
||||
from owncastsentry.utils import (
|
||||
domainify,
|
||||
escape_markdown,
|
||||
sanitize_for_markdown,
|
||||
sanitize_for_plain_text,
|
||||
should_query_stream,
|
||||
truncate,
|
||||
user_agent,
|
||||
)
|
||||
|
||||
|
||||
class TestUserAgent:
|
||||
"""User-Agent header construction."""
|
||||
|
||||
@pytest.mark.parametrize(
|
||||
("version", "expected"),
|
||||
[
|
||||
pytest.param(
|
||||
"1.2.3",
|
||||
"OwncastSentry/1.2.3 (bot; +https://git.logal.dev/LogalDeveloper/OwncastSentry)",
|
||||
id="semver",
|
||||
),
|
||||
pytest.param(
|
||||
"0.0.0",
|
||||
"OwncastSentry/0.0.0 (bot; +https://git.logal.dev/LogalDeveloper/OwncastSentry)",
|
||||
id="zeroed",
|
||||
),
|
||||
pytest.param(
|
||||
"1.1.1.dev10+gf0146d061.d20260313",
|
||||
"OwncastSentry/1.1.1.dev10+gf0146d061.d20260313 (bot; +https://git.logal.dev/LogalDeveloper/OwncastSentry)",
|
||||
id="dev-version",
|
||||
),
|
||||
],
|
||||
)
|
||||
def test_user_agent(self, version: str, expected: str) -> None:
|
||||
"""Build a correctly formatted User-Agent header."""
|
||||
assert user_agent(version) == expected
|
||||
|
||||
|
||||
class TestShouldQueryStream:
|
||||
"""Progressive backoff logic for stream polling."""
|
||||
|
||||
@pytest.mark.parametrize(
|
||||
("counter", "expected"),
|
||||
[
|
||||
pytest.param(0, True, id="counter-0-always-query"),
|
||||
pytest.param(1, True, id="counter-1-always-query"),
|
||||
pytest.param(4, True, id="counter-4-always-query"),
|
||||
pytest.param(5, False, id="counter-5-skip-odd"),
|
||||
pytest.param(6, True, id="counter-6-query-even"),
|
||||
pytest.param(9, False, id="counter-9-skip-odd"),
|
||||
pytest.param(10, False, id="counter-10-skip-not-mod-3"),
|
||||
pytest.param(12, True, id="counter-12-query-mod-3"),
|
||||
pytest.param(14, False, id="counter-14-skip-not-mod-3"),
|
||||
pytest.param(15, True, id="counter-15-query-mod-5"),
|
||||
pytest.param(16, False, id="counter-16-skip-not-mod-5"),
|
||||
pytest.param(20, True, id="counter-20-query-mod-5"),
|
||||
pytest.param(29, False, id="counter-29-skip-not-mod-5"),
|
||||
pytest.param(30, True, id="counter-30-query-mod-15"),
|
||||
pytest.param(31, False, id="counter-31-skip-not-mod-15"),
|
||||
pytest.param(45, True, id="counter-45-query-mod-15"),
|
||||
pytest.param(100, False, id="counter-100-skip-not-mod-15"),
|
||||
pytest.param(105, True, id="counter-105-query-mod-15"),
|
||||
],
|
||||
)
|
||||
def test_backoff_tiers(self, counter: int, expected: bool) -> None:
|
||||
"""Return the expected query decision for each backoff tier."""
|
||||
assert should_query_stream(counter) == expected
|
||||
|
||||
|
||||
class TestDomainify:
|
||||
"""Domain extraction and sanitization from user input."""
|
||||
|
||||
@pytest.mark.parametrize(
|
||||
("input_url", "expected"),
|
||||
[
|
||||
pytest.param("example.com", "example.com", id="bare-domain"),
|
||||
pytest.param("https://example.com", "example.com", id="https-url"),
|
||||
pytest.param("http://example.com", "example.com", id="http-url"),
|
||||
pytest.param("https://example.com:8080", "example.com", id="url-with-port"),
|
||||
pytest.param(
|
||||
"https://example.com/path/to/page",
|
||||
"example.com",
|
||||
id="url-with-path",
|
||||
),
|
||||
pytest.param(
|
||||
"user@stream.logal.dev",
|
||||
"stream.logal.dev",
|
||||
id="email-style",
|
||||
),
|
||||
pytest.param("EXAMPLE.COM", "example.com", id="uppercase"),
|
||||
pytest.param("exam!ple.com", "example.com", id="special-chars-stripped"),
|
||||
pytest.param(".example.com.", "example.com", id="leading-trailing-dots"),
|
||||
pytest.param("-example.com-", "example.com", id="leading-trailing-hyphens"),
|
||||
pytest.param(
|
||||
"sub.domain.example.com",
|
||||
"sub.domain.example.com",
|
||||
id="subdomain",
|
||||
),
|
||||
],
|
||||
)
|
||||
def test_extracts_domain(self, input_url: str, expected: str) -> None:
|
||||
"""Extract and sanitize the domain from various input formats."""
|
||||
assert domainify(input_url) == expected
|
||||
|
||||
|
||||
class TestTruncate:
|
||||
"""Text truncation to a maximum length."""
|
||||
|
||||
@pytest.mark.parametrize(
|
||||
("text", "max_length", "expected"),
|
||||
[
|
||||
pytest.param("hello", 10, "hello", id="under-limit"),
|
||||
pytest.param("hello", 5, "hello", id="exact-limit"),
|
||||
pytest.param("hello world", 5, "hello", id="over-limit"),
|
||||
pytest.param("", 5, "", id="empty-string"),
|
||||
],
|
||||
)
|
||||
def test_truncates(self, text: str, max_length: int, expected: str) -> None:
|
||||
"""Truncate text that exceeds the maximum length."""
|
||||
assert truncate(text, max_length) == expected
|
||||
|
||||
|
||||
class TestEscapeMarkdown:
|
||||
"""Markdown special character escaping."""
|
||||
|
||||
@pytest.mark.parametrize(
|
||||
("input_text", "expected"),
|
||||
[
|
||||
pytest.param("hello", "hello", id="plain-text-unchanged"),
|
||||
pytest.param("*bold*", "\\*bold\\*", id="asterisks"),
|
||||
pytest.param("_italic_", "\\_italic\\_", id="underscores"),
|
||||
pytest.param("[link](url)", "\\[link\\]\\(url\\)", id="link-syntax"),
|
||||
pytest.param("`code`", "\\`code\\`", id="backticks"),
|
||||
pytest.param("# heading", "\\# heading", id="heading"),
|
||||
pytest.param("> quote", "\\> quote", id="blockquote"),
|
||||
pytest.param("<html>", "\\<html\\>", id="angle-brackets"),
|
||||
pytest.param("a & b", "a \\& b", id="ampersand"),
|
||||
pytest.param("a\\b", "a\\\\b", id="backslash"),
|
||||
pytest.param("", "", id="empty-string"),
|
||||
],
|
||||
)
|
||||
def test_escapes_special_chars(self, input_text: str, expected: str) -> None:
|
||||
"""Escape the given Markdown special character."""
|
||||
assert escape_markdown(input_text) == expected
|
||||
|
||||
|
||||
class TestSanitizeForPlainText:
|
||||
"""Plain text sanitization for notifications."""
|
||||
|
||||
@pytest.mark.parametrize(
|
||||
("input_text", "expected"),
|
||||
[
|
||||
pytest.param("hello world", "hello world", id="plain-text"),
|
||||
pytest.param("line1\nline2", "line1 line2", id="newline-removed"),
|
||||
pytest.param("line1\rline2", "line1 line2", id="carriage-return"),
|
||||
pytest.param("line1\r\nline2", "line1 line2", id="crlf-removed"),
|
||||
pytest.param(
|
||||
"too many spaces", "too many spaces", id="spaces-collapsed"
|
||||
),
|
||||
pytest.param("", "", id="empty-string"),
|
||||
],
|
||||
)
|
||||
def test_sanitizes(self, input_text: str, expected: str) -> None:
|
||||
"""Sanitize the text for safe plain-text rendering."""
|
||||
assert sanitize_for_plain_text(input_text) == expected
|
||||
|
||||
|
||||
class TestSanitizeForMarkdown:
|
||||
"""Markdown sanitization combining newline removal and escaping."""
|
||||
|
||||
def test_removes_newlines_and_escapes(self) -> None:
|
||||
"""Remove newlines and escape Markdown special characters."""
|
||||
result = sanitize_for_markdown("*bold*\nnew line")
|
||||
assert result == "\\*bold\\* new line"
|
||||
|
||||
def test_empty_string(self) -> None:
|
||||
"""Return empty string unchanged."""
|
||||
assert sanitize_for_markdown("") == ""
|
||||
Reference in New Issue
Block a user