mirror of
https://github.com/rennf93/roboco.git
synced 2026-08-03 07:23:24 +02:00
* feat(metrics): capture per-session turns + tool_calls (phase 1)
Persist LLM iterations (turns) and tool invocations per agent spawn session,
the raw signal the granular per-member performance metrics build on (real
effort/iterations vs wall-clock).
- sum_transcript_usage returns a 5-tuple adding turns = unique assistant
message-id count; _usage_from_transcript + _resolve_active_tokens updated to
the 5-tuple (active-tokens keeps its 4-tuple contract by slicing).
- SDK: _SessionState.turns, set by /usage/sync; /usage/status (TokenUsageStatus)
now carries turns + tool_calls (= total_calls).
- orchestrator: new _resolve_final_turns_tools (SDK primary, transcript fallback
for turns only; Grok -> 0/0) wired into _finalize_spawn_session, which writes
turns + tool_calls to agent_spawn_sessions.
- migration 055 adds turns + tool_calls (BigInteger DEFAULT 0 -> historical/Grok
rows read 0, surfaced as n/a). Verified real alembic upgrade/downgrade.
Part of metrics-granularity (v0.15.0); recon-adjusted plan on disk.
* feat(metrics): pure compute_stage_effort helper (phase 2, part 1)
Foundation-layer overlap math (no DB): split each task status window into
active (merged wall-clock overlap of spawn stints — concurrent stints counted
once, so active <= window) vs wait (queue/review idle). Distinct from summed
effort. The per-task metrics service will feed it audit-log windows + spawn
stints. 9 unit tests (disjoint/nested/partial/merged/clamped/zero/multi-window).
* feat(metrics): per-task live metrics + GET /metrics/task/{id} (phase 2)
TaskMetrics dataclass + MetricsService.get_task_metrics: summed spawn effort
(vs wall-clock), turns/tool_calls/tokens/cost, per-stage active-vs-wait
(compute_stage_effort over audit windows x spawn stints), and who-caused-rework
(revision_count + named qa/pr fail events). Open stints and the open final
stage window close at completed_at for a terminal task (else now), so stages
don't grow past completion. Exposed at GET /dashboard/metrics/task/{task_id}
(404 if absent). Real-PG tests (compose/none/in-flight) + route tests (200/404).
* feat(metrics): CEO-as-member scorecard + ceo_reject audit regression (phase 3)
The human CEO is a measured member, read purely from audit_log (agent_role='ceo'
serializes from the CEO StrEnum): approval dwell (awaiting_ceo_approval -> a CEO
decision, incl. the coordination-root reject that lands in pending), unblock
dwell (blocked -> a CEO revive), and god-mode action count (every CEO-attributed
transition). CeoScorecard + MetricsService.get_ceo_scorecard (p50/p90 via
PERCENTILE_CONT, expanding IN for the decision sets) + GET
/dashboard/metrics/member/ceo (declared before any future member/{id} route).
The ceo_reject coordination-root audit gap the plan meant to close was already
closed by the gap-sweep (routes through admin_set_status -> agent_role='ceo'
audit); locked with a regression assertion in the existing coordination-reject
test. Real-PG tests: approval/unblock/godmode, non-ceo exclusion, empty->zeros.
* feat(metrics): audit instrumentation for escalations/blocked-others/idle (phase 4a)
The three extra per-member metrics that had no data source get durable,
in-session audit events (additive; never gate the underlying action):
- apply_escalation -> task.escalated (details.escalator_slug) on both the
normal block path and the pool-divert path -> escalations count.
- _unblock_dependents -> task.unblocked_dependents (details.count) on the
completed BLOCKER task, captured before the dependency edges are pruned ->
blocked-others count (sweeper attributes to the blocker's owner).
- mark_agent_idle -> agent.idle (details.agent_slug) -> idle/utilization (the
sweeper pairs an idle mark to the member's next spawn for idle duration).
(QA pass-rate needs no new event — reuses task.awaiting_documentation[qa] +
task.qa_fail.) Real-PG tests for each; 111 transition tests still green.
* feat(metrics): member_performance_daily rollup table + migration 056 (phase 4b)
The per-member scorecard rollup: one row per (date, member_kind, agent_slug),
CEO as a first-class member_kind='ceo' row (agent_slug='' NOT NULL so the
NULL-distinct UNIQUE keeps it unique). Full column set + the four CEO-approved
extras (qa_reviews_total/passed, escalations, blocked_others, idle_seconds) plus
blocked_seconds. Overwrite-upsert on (date, member_kind, agent_slug) for an
idempotent sweep. Migration 056 verified real up/down (24 cols, 4 indexes).
* feat(metrics): _sweep_member_performance rollup sweeper (phase 4c)
The daily per-member rollup sweep (mirrors _sweep_daily_rollup): a trailing
7-day, idempotent overwrite-upsert wired into _run_sweep. One focused query per
metric merges into a (date, agent_slug) accumulator — spawn effort/turns/tokens/
cost, completed/first-pass/revisions-received, revisions-caused (qa/pr fails),
QA pass-rate (passed + total), escalations (by escalator_slug), blocked-others
(unblocked_dependents by blocker owner), idle_seconds (idle mark -> next spawn),
blocked_seconds (blocked dwell) — plus one CEO row/day (approval/unblock dwell +
god-mode). Real-PG test asserts every facet + idempotency (a 2nd sweep
overwrites, never doubles); spawn-day != completion-day split is by-design.
* feat(metrics): member/org rollup scorecards + endpoints + live overlay (phase 5)
MemberScorecard + OrgScorecard with derived rates (FPY, effort-throughput,
turns/tool-calls per task, QA pass-rate, utilization) — all division-guarded to
None. get_member_scorecard reads member_performance_daily by slug and overlays
the member's live in-flight (non-terminal) tasks' effort via get_task_metrics
(disjoint by status: completion counts stay rollup-only, overlay only enriches
effort/turns/cost; includes_live_inflight flags it). get_org_scorecard
aggregates the cell (?team=) or whole org. Routes: GET /metrics/member/{agent_id}
(404 if absent, after the ceo literal route) + GET /metrics/org?team=. Real-PG
tests (derived rates, overlay no double-count, guards, org) + route tests.
* feat(metrics): granular CEO completion notification (phase 6)
There was no CEO completion notification at all (EventType.TASK_COMPLETED was
defined but never emitted). Add notify_ceo_of_completion in
NotificationDeliveryService — a granular body (real effort vs wall-clock +
stints/turns/tool-calls/revisions[QA/PR]/cost from get_task_metrics; degrades to
wall-clock-only, turns 'n/a', when there are no spawn sessions). Reuses the
existing ALERT type (no enum migration; the notificationtype PG enum is fixed at
001). ceo_approve now emits TASK_COMPLETED + fires the notification (best-effort
via _notify_completion — never blocks completion); complete() emits
TASK_COMPLETED too (closes the dead-code gap; the WS bridge can forward it).
Pure formatter tests + real-PG notification test.
* [metrics-granularity] Phase 7: panel Scorecards tab + dashboard overview
Add the CEO-facing metrics surfaces for the granularity feature:
- New "Scorecards" tab on the Metrics page: org rollup headline, the
CEO-as-member card (approval/unblock dwell + god-mode count), and a
per-member table (completed, first-pass yield, active effort, turns/task,
QA pass-rate, escalations, blocked-others, utilization). Each member row
self-fetches its rollup scorecard; live in-flight rows carry a "live" badge.
- New dashboard overview card (ScorecardOverviewPanel): org-wide 30-day
headline (completed, FPY, throughput/hr, active effort, cost) deep-linking
into the Scorecards tab.
- Plumbing: TaskMetrics/MemberScorecard/OrgScorecard/CeoScorecard types,
observability API client methods + empty fallbacks, and the four
useCeoScorecard/useMemberScorecard/useOrgScorecard/useTaskMetrics hooks.
Panel gate green: tsc, eslint, prettier, vitest (175 tests, +6 new).
* [metrics-granularity] test: make completion-notification robust to shared-DB CEO
test_notify_ceo_of_completion_creates_alert errored in the full suite (passed
in isolation): the session-scoped test DB is shared across the run, and the
sibling real-DB board-gate test commits a role=CEO agent (slug="ceo") without
cleanup — so my env fixture's hardcoded slug="ceo" insert hit a unique-constraint
violation, and a second role=CEO row would also make _get_ceo_agent()'s
scalar_one_or_none() raise. Reuse an existing CEO when present (the singleton the
production system actually has), else create one with a unique slug. Order-
independent. Also reflow test_metrics_instrumentation.py to ruff format.
* chore(release): 0.15.0
Metrics granularity: per-member/per-task/org + CEO-as-member scorecards,
turn/tool-call capture (migration 055), member_performance_daily rollup
(migration 056) with QA pass-rate / escalations / blocked-others / utilization,
per-task active-vs-wait metrics, granular completion notification, panel
Scorecards tab + dashboard Performance card, and the ceo_reject audit fix.
Version bump across the canonical set + CHANGELOG.
* [metrics-granularity] fix pre-tag audit findings (overlay double-count + panel error states)
Adversarial review before the v0.15.0 tag surfaced two real logical gaps:
- MAJOR (backend): the live in-flight overlay re-summed ALL sessions of every
non-terminal task via get_task_metrics, but _msweep_spawn already rolls up
every CLOSED session regardless of task status — so a closed session on a
still-open task was counted twice (rollup + overlay), permanently inflating a
member's effort/turns/tokens/cost on the common reap/respawn path. The overlay
now sums only OPEN sessions (ended_at IS NULL), which the closed-only rollup
can never contain — disjoint by construction. A just-closed session lands in
the rollup on the next ~60s sweep (no gap of note). Aggregated in SQL to mirror
_msweep_spawn. Regression test reproduces the double-count (turns 10→5).
- MAJOR (panel): the four new scorecard surfaces used `isLoading || !data` with
no isError branch, so a failed query span forever on a skeleton. They now
surface a load error. Tests added.
Also: OrgSummary active-effort formatting no longer round-trips hours→seconds→
hours; dashboard grid uses xl:grid-cols-4 (was 2xl) so 4 panels show at 1280px;
corrected the inaccurate "NULL distinct" CEO-row uniqueness comment (agent_slug
is NOT NULL; the '' tuple is simply distinct from agent rows).
make quality GREEN (cov 95.31%); panel GREEN (vitest 178).
* [metrics-granularity] fix: decode bytes stream message-id before XCLAIM
StreamEventBus._recover_stream passed the pending message id to XCLAIM via
str() on the raw bytes the client returns (redis client has no
decode_responses), producing "b'1782066556728-0'". Redis rejects that with
"Unrecognized XCLAIM option", so pending-message recovery threw on every
reclaim tick and unacked messages from crashed/slow consumers were never
reclaimed (leaking in the PEL on every stream, spamming the error log). Decode
via the existing _to_str helper — the fix the sibling claim path already uses.
Pre-existing in v0.14.0 (unrelated to metrics granularity); folded into this
release per CEO. TDD regression test + CHANGELOG entry. make quality GREEN.
---------
Co-authored-by: Renn F <rennf93@users.noreply.github.com>
406 lines
13 KiB
Python
406 lines
13 KiB
Python
"""Coverage for roboco.events.bus thin wrapper functions."""
|
|
|
|
from __future__ import annotations
|
|
|
|
import asyncio
|
|
import contextlib
|
|
import json
|
|
from typing import TYPE_CHECKING, cast
|
|
from unittest.mock import AsyncMock, MagicMock, patch
|
|
|
|
import pytest
|
|
from roboco.events.bus import EventBus, get_event_bus, init_event_bus
|
|
from roboco.events.stream_bus import StreamEventBus
|
|
from roboco.models.events import Event, EventType
|
|
|
|
if TYPE_CHECKING:
|
|
from redis.asyncio import Redis
|
|
|
|
|
|
def test_get_event_bus_delegates() -> None:
|
|
"""get_event_bus() returns the underlying stream event bus singleton."""
|
|
fake_bus = MagicMock()
|
|
with patch(
|
|
"roboco.events.bus.get_stream_event_bus", return_value=fake_bus
|
|
) as mock_get:
|
|
result = get_event_bus()
|
|
mock_get.assert_called_once()
|
|
assert result is fake_bus
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_init_event_bus_delegates() -> None:
|
|
"""init_event_bus() forwards args to init_stream_event_bus (line 53)."""
|
|
fake_bus = MagicMock()
|
|
with patch(
|
|
"roboco.events.bus.init_stream_event_bus",
|
|
new_callable=AsyncMock,
|
|
return_value=fake_bus,
|
|
) as mock_init:
|
|
result = await init_event_bus(consumer_name="custom", recover_pending=False)
|
|
mock_init.assert_awaited_once_with(consumer_name="custom", recover_pending=False)
|
|
assert result is fake_bus
|
|
|
|
|
|
def test_event_bus_alias_is_stream_event_bus() -> None:
|
|
"""EventBus is an alias for StreamEventBus."""
|
|
|
|
assert EventBus is StreamEventBus
|
|
|
|
|
|
# --- #19: replayed events must not re-run already-succeeded handlers ---
|
|
|
|
|
|
class _FakeRedis:
|
|
"""In-memory stand-in for the redis client's SET-NX + DELETE surface.
|
|
|
|
``set(..., nx=True)`` returns True the first time a key is set, None if it
|
|
already exists (matches redis-py). ``delete`` removes a key. ``get`` is
|
|
unused but kept for completeness.
|
|
"""
|
|
|
|
def __init__(self) -> None:
|
|
self.keys: dict[str, str] = {}
|
|
self.set_calls: list[tuple[str, bool]] = []
|
|
|
|
async def set(
|
|
self, key: str, value: str, *, nx: bool = False, ex: int | None = None
|
|
) -> bool | None:
|
|
del ex
|
|
self.set_calls.append((key, nx))
|
|
if nx:
|
|
if key in self.keys:
|
|
return None
|
|
self.keys[key] = value
|
|
return True
|
|
self.keys[key] = value
|
|
return True
|
|
|
|
async def delete(self, key: str) -> int:
|
|
return 1 if self.keys.pop(key, None) is not None else 0
|
|
|
|
async def get(self, key: str) -> str | None:
|
|
return self.keys.get(key)
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_dispatch_skips_handler_that_already_succeeded_on_replay() -> None:
|
|
"""#19: ``recover_pending`` re-delivers a message whose ACK was blocked by a
|
|
sibling handler's failure. A handler that already succeeded for that event
|
|
must NOT re-run (no duplicate side effects) — the bus marks each
|
|
(event.id, handler) processed via a SET-NX guard."""
|
|
|
|
bus = StreamEventBus()
|
|
bus._redis = cast("Redis", _FakeRedis())
|
|
|
|
calls: list[str] = []
|
|
|
|
async def _handler(_event: Event) -> None:
|
|
calls.append("ran")
|
|
|
|
bus.subscribe(EventType.NOTIFICATION_SENT, _handler)
|
|
|
|
event = Event(type=EventType.NOTIFICATION_SENT, data={"task_id": "t1"})
|
|
|
|
await bus._dispatch_event(event) # first delivery: handler runs
|
|
assert calls == ["ran"]
|
|
|
|
await bus._dispatch_event(event) # replay: handler already succeeded → skip
|
|
assert calls == ["ran"] # not re-run
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_dispatch_reruns_handler_that_failed_on_first_attempt() -> None:
|
|
"""#19: a handler that failed on the first delivery must NOT be marked
|
|
processed — a replay re-runs it (the SET-NX key is cleared on failure)."""
|
|
|
|
bus = StreamEventBus()
|
|
bus._redis = cast("Redis", _FakeRedis())
|
|
|
|
attempts: list[str] = []
|
|
|
|
async def _flaky(_event: Event) -> None:
|
|
attempts.append("ran")
|
|
if len(attempts) == 1:
|
|
raise RuntimeError("transient blow-up")
|
|
|
|
bus.subscribe(EventType.NOTIFICATION_SENT, _flaky)
|
|
|
|
event = Event(type=EventType.NOTIFICATION_SENT, data={"task_id": "t2"})
|
|
first = await bus._dispatch_event(event) # fails → not marked processed
|
|
assert first is False
|
|
assert attempts == ["ran"]
|
|
|
|
await bus._dispatch_event(event) # replay: re-run (failed before)
|
|
assert attempts == ["ran", "ran"]
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_dispatch_runs_handler_when_redis_guard_unavailable() -> None:
|
|
"""#19: the idempotency guard is best-effort — with no redis the bus must
|
|
still run the handler (fail-open: never skip a handler because the dedup
|
|
infra is down)."""
|
|
|
|
bus = StreamEventBus()
|
|
# No redis connected — guard is skipped, handler runs normally.
|
|
assert bus._redis is None
|
|
|
|
calls: list[str] = []
|
|
|
|
async def _handler(_event: Event) -> None:
|
|
calls.append("ran")
|
|
|
|
bus.subscribe(EventType.NOTIFICATION_SENT, _handler)
|
|
event = Event(type=EventType.NOTIFICATION_SENT, data={"task_id": "t3"})
|
|
|
|
ok = await bus._dispatch_event(event)
|
|
assert ok is True
|
|
assert calls == ["ran"]
|
|
|
|
|
|
# --- poison-pill: an undecodable message must be ACKed, not retried forever ---
|
|
|
|
|
|
class _FakeRedisStream:
|
|
"""In-memory redis surface for the xack/xadd/group path.
|
|
|
|
Tracks ACKs and dead-letter xadds so the poison-pill test can assert a
|
|
malformed message is acknowledged (not left pending) and parked on the
|
|
dead-letter stream. ``xreadgroup`` yields nothing so a listen loop never
|
|
spins.
|
|
"""
|
|
|
|
def __init__(self) -> None:
|
|
self.xack_calls: list[tuple[str, str, tuple[str, ...]]] = []
|
|
self.xadd_calls: list[tuple[str, dict]] = []
|
|
|
|
async def xack(self, stream: str, group: str, *ids: str) -> int:
|
|
self.xack_calls.append((stream, group, ids))
|
|
return len(ids)
|
|
|
|
async def xadd(
|
|
self,
|
|
stream: str,
|
|
fields: dict,
|
|
maxlen: int | None = None,
|
|
approximate: bool = True,
|
|
) -> bytes:
|
|
del maxlen, approximate
|
|
self.xadd_calls.append((stream, dict(fields)))
|
|
return b"1-0"
|
|
|
|
async def xreadgroup(self, *args: object, **kwargs: object) -> list:
|
|
del args, kwargs
|
|
return []
|
|
|
|
async def xgroup_create(self, *args: object, **kwargs: object) -> bool:
|
|
del args, kwargs
|
|
return True
|
|
|
|
async def xpending(self, *args: object, **kwargs: object) -> dict:
|
|
del args, kwargs
|
|
return {"pending": 0}
|
|
|
|
async def xpending_range(self, *args: object, **kwargs: object) -> list:
|
|
del args, kwargs
|
|
return []
|
|
|
|
async def xclaim(self, *args: object, **kwargs: object) -> list:
|
|
del args, kwargs
|
|
return []
|
|
|
|
async def set(
|
|
self, key: str, value: str, *, nx: bool = False, ex: int | None = None
|
|
) -> bool:
|
|
del key, value, nx, ex
|
|
return True
|
|
|
|
async def delete(self, key: str) -> int:
|
|
del key
|
|
return 1
|
|
|
|
async def get(self, key: str) -> None:
|
|
del key
|
|
|
|
async def close(self) -> None:
|
|
"""No-op close for the fake client."""
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_undecodable_message_is_acked_and_dead_lettered() -> None:
|
|
"""A message whose payload fails Event.from_json (unknown EventType value,
|
|
bad UUID, malformed JSON) is a poison pill: no handler could ever process
|
|
it, so retrying is pointless. The bus must ACK it (and dead-letter it) so
|
|
the stream doesn't wedge on an unkillable pending message re-failing on
|
|
every reclaim."""
|
|
|
|
bus = StreamEventBus()
|
|
fake = _FakeRedisStream()
|
|
bus._redis = cast("Redis", fake)
|
|
|
|
invoked: list[str] = []
|
|
|
|
async def _handler(_event: Event) -> None:
|
|
invoked.append("ran")
|
|
|
|
bus.subscribe(EventType.NOTIFICATION_SENT, _handler)
|
|
|
|
# type="task.bogus" is not a real EventType → EventType(...) raises ValueError
|
|
# inside Event.from_json.
|
|
malformed = json.dumps(
|
|
{
|
|
"id": "not-a-uuid",
|
|
"type": "task.bogus_unknown",
|
|
"data": {},
|
|
"timestamp": "2026-06-30T00:00:00+00:00",
|
|
}
|
|
)
|
|
await bus._handle_message("roboco:stream:task", "1234-0", {b"data": malformed})
|
|
|
|
# ACKed exactly once — not left pending for reclaim to re-fail forever.
|
|
assert len(fake.xack_calls) == 1
|
|
assert fake.xack_calls[0][2] == ("1234-0",)
|
|
# Dead-lettered for inspection before the ACK.
|
|
assert len(fake.xadd_calls) == 1
|
|
assert fake.xadd_calls[0][0] == StreamEventBus.DEAD_LETTER_STREAM
|
|
# No handler could run — the event never decoded.
|
|
assert invoked == []
|
|
|
|
|
|
class _FakeRecoverRedis:
|
|
"""Fake whose xpending_range returns the message id as BYTES (the real
|
|
client has no decode_responses), and which captures the ids XCLAIM gets."""
|
|
|
|
def __init__(self, message_id: bytes) -> None:
|
|
self._message_id = message_id
|
|
self.claimed_ids: list[object] = []
|
|
|
|
async def xpending(self, *args: object, **kwargs: object) -> dict:
|
|
del args, kwargs
|
|
return {"pending": 1}
|
|
|
|
async def xpending_range(self, *args: object, **kwargs: object) -> list:
|
|
del args, kwargs
|
|
return [{"message_id": self._message_id, "time_since_delivered": 10_000}]
|
|
|
|
async def xclaim(self, *args: object, **kwargs: object) -> list:
|
|
del args
|
|
self.claimed_ids = cast("list[object]", kwargs.get("message_ids") or [])
|
|
return [] # nothing claimed back → no handling
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_recover_stream_decodes_bytes_message_id_for_xclaim() -> None:
|
|
"""xpending_range returns the message id as bytes; _recover_stream must
|
|
decode it before XCLAIM. A raw ``str(bytes)`` yields ``"b'1782..-0'"``,
|
|
which Redis rejects with "Unrecognized XCLAIM option", so pending-message
|
|
recovery silently fails every reclaim tick."""
|
|
bus = StreamEventBus()
|
|
fake = _FakeRecoverRedis(b"1782066556728-0")
|
|
bus._redis = cast("Redis", fake)
|
|
|
|
await bus._recover_stream("roboco:stream:usage", idle_time_ms=0)
|
|
|
|
assert fake.claimed_ids == ["1782066556728-0"] # decoded, not "b'...'"
|
|
|
|
|
|
# --- periodic reclaim: a runtime handler failure is retried without a restart ---
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_reclaim_loop_periodically_calls_recover_pending() -> None:
|
|
"""XREADGROUP '>' delivers only NEW messages, so a handler that fails at
|
|
runtime leaves its message pending and unretried until the orchestrator
|
|
restarts. A periodic reclaim loop must call recover_pending so the
|
|
idempotency-guarded replay actually fires."""
|
|
|
|
bus = StreamEventBus()
|
|
bus._running = True
|
|
bus._reclaim_interval = 60
|
|
|
|
calls: list[int] = []
|
|
|
|
async def _fake_recover(idle_time_ms: int = 60000) -> int:
|
|
calls.append(idle_time_ms)
|
|
bus._running = False # break the loop after the first reclaim
|
|
return 0
|
|
|
|
bus.recover_pending = _fake_recover # type: ignore[method-assign]
|
|
|
|
async def _no_sleep(_seconds: float) -> None:
|
|
return
|
|
|
|
with patch("roboco.events.stream_bus.asyncio.sleep", new=_no_sleep):
|
|
await bus._reclaim_loop()
|
|
|
|
# Reclaim ran once with the interval-aligned idle window, then the loop exited.
|
|
assert calls == [60000]
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_start_listening_spawns_reclaim_task_alongside_listen() -> None:
|
|
"""start_listening must spawn the reclaim task, not just the listen task —
|
|
otherwise pending messages are never re-delivered at runtime."""
|
|
|
|
bus = StreamEventBus()
|
|
bus._redis = cast("Redis", _FakeRedisStream())
|
|
|
|
async def _noop(self: StreamEventBus) -> None:
|
|
del self
|
|
|
|
async def _handler(_event: Event) -> None: ...
|
|
|
|
bus.subscribe(EventType.NOTIFICATION_SENT, _handler)
|
|
|
|
with (
|
|
patch.object(StreamEventBus, "_listen_loop", _noop),
|
|
patch.object(StreamEventBus, "_reclaim_loop", _noop),
|
|
):
|
|
await bus.start_listening()
|
|
|
|
try:
|
|
assert bus._listen_task is not None
|
|
assert bus._reclaim_task is not None
|
|
finally:
|
|
await bus.disconnect()
|
|
|
|
|
|
# --- cancellation mid-handler must clear the idempotency marker ---
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_cancelled_handler_clears_idempotency_marker() -> None:
|
|
"""A handler cancelled mid-flight (shutdown / sibling gather cancellation)
|
|
is BaseException-cancelled, not Exception-raised, so the old ``except
|
|
Exception`` left the SET-NX marker set: the message stayed pending but the
|
|
guard then suppressed the very redelivery that would complete the work.
|
|
The cleanup must catch BaseException so the marker is cleared and reclaim
|
|
re-runs the handler."""
|
|
|
|
bus = StreamEventBus()
|
|
fake = _FakeRedis()
|
|
bus._redis = cast("Redis", fake)
|
|
|
|
started = asyncio.Event()
|
|
proceed = asyncio.Event()
|
|
|
|
async def _blocking(_event: Event) -> None:
|
|
started.set()
|
|
await proceed.wait() # block until the dispatch task is cancelled
|
|
|
|
bus.subscribe(EventType.NOTIFICATION_SENT, _blocking)
|
|
event = Event(type=EventType.NOTIFICATION_SENT, data={"task_id": "tc"})
|
|
|
|
task = asyncio.create_task(bus._dispatch_event(event))
|
|
await started.wait() # handler is now blocked → marker is set
|
|
|
|
key = f"bus:processed:{event.id}:_blocking"
|
|
assert key in fake.keys # marker acquired before the handler blocked
|
|
|
|
task.cancel()
|
|
with contextlib.suppress(asyncio.CancelledError):
|
|
await task
|
|
|
|
# Marker cleared despite cancellation → a replay re-runs the handler.
|
|
assert key not in fake.keys
|