Compare commits
2
Commits
9be932c0c7
...
47c7e97ac0
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
47c7e97ac0
|
||
|
|
c690702113
|
@@ -80,7 +80,7 @@ class NotificationService:
|
|||||||
# Has enough time passed since the last notification was sent?
|
# Has enough time passed since the last notification was sent?
|
||||||
if not self._can_notify(domain):
|
if not self._can_notify(domain):
|
||||||
seconds_since_last = round(
|
seconds_since_last = round(
|
||||||
time.time() - self.notification_timers_cache[domain]
|
time.monotonic() - self.notification_timers_cache[domain]
|
||||||
)
|
)
|
||||||
self.log.info(
|
self.log.info(
|
||||||
f"[{domain}] Not sending notifications. Only "
|
f"[{domain}] Not sending notifications. Only "
|
||||||
|
|||||||
@@ -82,12 +82,24 @@ class StreamMonitor:
|
|||||||
# Build a list of async tasks for each stream domain
|
# Build a list of async tasks for each stream domain
|
||||||
# Run all stream updates in parallel and collect results
|
# Run all stream updates in parallel and collect results
|
||||||
results = await asyncio.gather(
|
results = await asyncio.gather(
|
||||||
*(self.update_stream(domain) for domain in subscribed_domains)
|
*(self.update_stream(domain) for domain in subscribed_domains),
|
||||||
|
return_exceptions=True,
|
||||||
)
|
)
|
||||||
|
|
||||||
# Count successes and failures
|
# Count successes and failures
|
||||||
successful_checks = results.count(True)
|
successful_checks = 0
|
||||||
failed_checks = results.count(False)
|
failed_checks = 0
|
||||||
|
for domain, result in zip(subscribed_domains, results, strict=True):
|
||||||
|
if isinstance(result, BaseException):
|
||||||
|
self.log.exception(
|
||||||
|
f"[{domain}] Unhandled exception during stream update.",
|
||||||
|
exc_info=result,
|
||||||
|
)
|
||||||
|
failed_checks += 1
|
||||||
|
elif result:
|
||||||
|
successful_checks += 1
|
||||||
|
else:
|
||||||
|
failed_checks += 1
|
||||||
|
|
||||||
self.log.debug(
|
self.log.debug(
|
||||||
f"Update complete. {successful_checks}/{total_streams} succeeded, "
|
f"Update complete. {successful_checks}/{total_streams} succeeded, "
|
||||||
|
|||||||
@@ -864,6 +864,58 @@ class TestUpdateAllStreamsMixed:
|
|||||||
assert result.successful_checks == 1
|
assert result.successful_checks == 1
|
||||||
assert result.failed_checks == 1
|
assert result.failed_checks == 1
|
||||||
|
|
||||||
|
async def test_exception_isolated_to_failing_domain(
|
||||||
|
self,
|
||||||
|
stream_repo: StreamRepository,
|
||||||
|
subscription_repo: SubscriptionRepository,
|
||||||
|
) -> None:
|
||||||
|
"""Count an unhandled exception as a failure without aborting other updates."""
|
||||||
|
owncast = _StubOwncastClient(
|
||||||
|
stream_state=StreamState(
|
||||||
|
domain="ok.com",
|
||||||
|
title="Title",
|
||||||
|
last_disconnect_time="2026-01-01T00:00:00Z",
|
||||||
|
),
|
||||||
|
)
|
||||||
|
# Patch get_stream_state to raise for one specific domain
|
||||||
|
original_get_state = owncast.get_stream_state
|
||||||
|
|
||||||
|
async def raising_get_state(domain: str) -> StreamState | None:
|
||||||
|
if domain == "raise.com":
|
||||||
|
msg = "simulated DB error"
|
||||||
|
raise RuntimeError(msg)
|
||||||
|
return await original_get_state(domain)
|
||||||
|
|
||||||
|
owncast.get_stream_state = raising_get_state # type: ignore[assignment]
|
||||||
|
|
||||||
|
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="raise.com",
|
||||||
|
last_disconnect_time="2026-01-01T00:00:00Z",
|
||||||
|
)
|
||||||
|
await _seed_stream(
|
||||||
|
stream_repo,
|
||||||
|
subscription_repo,
|
||||||
|
domain="ok.com",
|
||||||
|
room_id="!room2:matrix.org",
|
||||||
|
title="Title",
|
||||||
|
last_disconnect_time="2026-01-01T00:00:00Z",
|
||||||
|
)
|
||||||
|
|
||||||
|
result = await monitor.update_all_streams(["raise.com", "ok.com"])
|
||||||
|
assert result.total_streams == 2
|
||||||
|
assert result.failed_checks == 1
|
||||||
|
assert result.successful_checks == 1
|
||||||
|
|
||||||
|
|
||||||
class TestStreamMonitorMetrics:
|
class TestStreamMonitorMetrics:
|
||||||
"""Stream monitor metrics recording."""
|
"""Stream monitor metrics recording."""
|
||||||
|
|||||||
Reference in New Issue
Block a user