diff --git a/backend/pipeline_worker.py b/backend/pipeline_worker.py index 1ff5d1d..c6573d2 100644 --- a/backend/pipeline_worker.py +++ b/backend/pipeline_worker.py @@ -81,11 +81,13 @@ async def _run() -> None: # are visible to any API instance tailing the event log. pipeline_svc.set_broker(event_bridge.GCSEventBroker(storage, user_id)) - # Fresh runs wipe the prior event log so the SSE consumer doesn't - # mix old events into the new run. Resume keeps the prior log so - # users see the full history. - if not resume: - pipeline_svc.broker.clear_history(project_id) + # Always wipe the prior event log. Reprocess uses resume=True, and + # the old log still contains ``pipeline_complete``; the SSE tail + # would stop there and the UI would show a finished run with no live + # log while the worker is still reviewing. Pause-resume also hits a + # terminal ``pipeline_paused``. A fresh seq from 0 is the only safe + # option — completed_review_refs still skip paid ICs. + pipeline_svc.broker.clear_history(project_id) if mode == "run": await pipeline_svc.run_pipeline( diff --git a/tests/test_event_bridge.py b/tests/test_event_bridge.py new file mode 100644 index 0000000..5e8de62 --- /dev/null +++ b/tests/test_event_bridge.py @@ -0,0 +1,24 @@ +"""GCS event log: a new run must not replay a prior pipeline_complete.""" + +from backend.services.event_bridge import GCSEventBroker + + +def test_clear_history_wipes_terminal_events(tmp_path): + from backend.services.storage import LocalStorageBackend + + storage = LocalStorageBackend(tmp_path) + broker = GCSEventBroker(storage, "local") + broker.publish("p1", "step_update", {"stage": "validation"}) + broker.publish("p1", "pipeline_complete", {"summary": {}}) + + prefix = "users/local/projects/p1/events/" + assert len(storage.list_prefix(prefix)) == 2 + + broker.clear_history("p1") + assert storage.list_prefix(prefix) == [] + + broker.publish("p1", "step_update", {"stage": "bom_parse"}) + keys = storage.list_prefix(prefix) + assert len(keys) == 1 + assert keys[0].endswith("0000000000.json") + assert storage.read_json(keys[0])["event"] == "step_update"