From 351453c08cd77cab000b89ab026ec5382897826b Mon Sep 17 00:00:00 2001 From: Korivi Date: Wed, 23 Sep 2026 13:04:51 +0900 Subject: [PATCH] fix: don't flag an in-flight parallel duplicate as crashed Port of #443, which landed on dev. V1.4.3 carries the identical guard code and therefore the identical bug: execute_parallel dispatches a whole batch before any call completes, so two identical irreversible calls issue two begin()s with no complete() between them. The second saw the first's INTENT row, read it as an interrupted prior attempt, downgraded it to FAILED and told the model the action "may have already happened" while the first was still running. Age cannot separate the two cases -- an immediate crash-then-restart leaves an equally young INTENT row, and that one does need the warning. Tracking in-flight keys on the guard instance distinguishes them, because a restart builds a fresh guard with an empty set. Verified on V1.4.3, not assumed: the new test fails against this branch's parent and passes with the fix. Full suite 1259 passed, 0 failed. The author's noted caveat -- a key stuck in the in-flight set if complete() never runs -- is narrower than stated here: manager.py wraps execution in try/except and calls complete() after it, including on the cancellation path, so only a BaseException outside Exception/CancelledError or process death could leak, and process death clears the set anyway. --- app/triggers/activity_log.py | 24 ++++++++++++++++++++++++ tests/test_activity_log.py | 25 +++++++++++++++++++++++++ 2 files changed, 49 insertions(+) diff --git a/app/triggers/activity_log.py b/app/triggers/activity_log.py index a7d878a2d..9c0ee3c95 100644 --- a/app/triggers/activity_log.py +++ b/app/triggers/activity_log.py @@ -240,6 +240,10 @@ class ActivityLogGuard: def __init__(self, log: ActivityLog): self._log = log + # idem_keys with an INTENT this process recorded and hasn't yet + # complete()'d, so begin() can tell a live batch mate apart from a + # crash-then-restart, which a fresh empty set never carries. + self._in_flight: set[str] = set() def begin( self, @@ -254,6 +258,7 @@ def begin( # Fresh attempt, or deliberate retake after a failed/uncertain # run — record intent BEFORE the side effect. self._log.record_intent(idem_key, action_name, session_id) + self._in_flight.add(idem_key) return GuardDecision(proceed=True, idem_key=idem_key) if row["status"] == STATUS_DONE: @@ -265,6 +270,7 @@ def begin( done_at = 0.0 if time.time() - done_at > DONE_DEDUP_WINDOW_SECONDS: self._log.record_intent(idem_key, action_name, session_id) + self._in_flight.add(idem_key) return GuardDecision(proceed=True, idem_key=idem_key) # This exact side effect just completed — return its stored # output instead of doing it again. @@ -281,6 +287,23 @@ def begin( logger.info(f"[ActivityLog] Skipped duplicate {action_name} ({idem_key})") return GuardDecision(proceed=False, idem_key=idem_key, stored_output=stored) + if idem_key in self._in_flight: + # A batch mate is still executing this exact call, not crashed; + # refuse without touching the row so it can record the outcome. + logger.info( + f"[ActivityLog] {action_name} ({idem_key}) already in " + f"flight in this batch, refusing the duplicate" + ) + return GuardDecision( + proceed=False, + idem_key=idem_key, + note=( + f"{action_name} is already running with these exact " + f"inputs. Wait for it to finish instead of calling it " + f"again." + ), + ) + # Stale INTENT: a previous attempt was interrupted between starting # the side effect and recording its outcome. Surface once; the next # identical attempt (if the LLM/user decides to retry) is allowed. @@ -306,6 +329,7 @@ def complete( status: str, outputs: Optional[Dict[str, Any]], ) -> None: + self._in_flight.discard(idem_key) ledger_status = STATUS_DONE if status == "success" else STATUS_FAILED provider_ref = None diff --git a/tests/test_activity_log.py b/tests/test_activity_log.py index aec4a12f8..5cd392f53 100644 --- a/tests/test_activity_log.py +++ b/tests/test_activity_log.py @@ -78,6 +78,31 @@ def test_crash_window_surfaces_uncertainty_once_then_allows_retry(self, tmp_path assert d3.proceed assert log2.get(d1.idem_key)["status"] == "INTENT" + def test_parallel_duplicate_in_same_process_is_not_flagged_as_crashed( + self, tmp_path + ): + log, guard = make_guard(tmp_path) + d1 = guard.begin("send_gmail", INPUTS, "task1") + assert d1.proceed + # A batch mate with identical inputs, dispatched by execute_parallel + # before d1's complete() is reached: not a crash, still running. + d2 = guard.begin("send_gmail", INPUTS, "task1") + assert not d2.proceed + assert d2.stored_output is None + assert "MAY" not in d2.note + assert "already running" in d2.note + assert log.get(d1.idem_key)["status"] == "INTENT" + + guard.complete( + d1.idem_key, "success", {"status": "success", "message_id": "msg-1"} + ) + + # Once complete()'d, a later identical call follows the ordinary + # DONE-dedup path, not the in-flight one. + d3 = guard.begin("send_gmail", INPUTS, "task1") + assert not d3.proceed + assert d3.stored_output["_idempotent_replay"] is True + def test_failed_run_can_be_retried(self, tmp_path): log, guard = make_guard(tmp_path) d1 = guard.begin("send_gmail", INPUTS, "task1")