Fixed single stream update exception aborting all parallel updates.
CD / Build (push) Successful in 8s
CI / Formatting (push) Successful in 6s
CI / Linting (push) Successful in 6s
CI / Tests (push) Successful in 11s
CI / Type Checking (push) Successful in 10s
CI / Spelling (push) Successful in 5s
Audit / Dependencies (push) Failing after 9s
CD / Build (push) Successful in 8s
CI / Formatting (push) Successful in 6s
CI / Linting (push) Successful in 6s
CI / Tests (push) Successful in 11s
CI / Type Checking (push) Successful in 10s
CI / Spelling (push) Successful in 5s
Audit / Dependencies (push) Failing after 9s
This commit is contained in:
@@ -82,12 +82,24 @@ class StreamMonitor:
|
||||
# Build a list of async tasks for each stream domain
|
||||
# Run all stream updates in parallel and collect results
|
||||
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
|
||||
successful_checks = results.count(True)
|
||||
failed_checks = results.count(False)
|
||||
successful_checks = 0
|
||||
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(
|
||||
f"Update complete. {successful_checks}/{total_streams} succeeded, "
|
||||
|
||||
@@ -864,6 +864,58 @@ class TestUpdateAllStreamsMixed:
|
||||
assert result.successful_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:
|
||||
"""Stream monitor metrics recording."""
|
||||
|
||||
Reference in New Issue
Block a user