diff --git a/CHANGELOG.md b/CHANGELOG.md index 11cac7c..bb84fe7 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -34,3 +34,7 @@ - Initial `0.1.0-alpha.2` event contract and conformance fixtures. - Initial Python reference SDK with schema/privacy validation and OTel span-event projection. + +### Fixed + +- the Python `EvidenceAccumulator` deadlocked forever if `append` or`seal` was called reentrantly on the same thread (for example, from inside a `durable_append` callback). It now raises `EvidenceError` instead, matching the TypeScript SDK's existing reentrancy protection. `snapshot` remains safely callable reentrantly. diff --git a/src/agentrust_telemetry/evidence.py b/src/agentrust_telemetry/evidence.py index aa9b1ba..01350da 100644 --- a/src/agentrust_telemetry/evidence.py +++ b/src/agentrust_telemetry/evidence.py @@ -5,7 +5,7 @@ import hashlib from copy import deepcopy from dataclasses import dataclass -from threading import Lock +from threading import Lock, get_ident from typing import Any, Callable, Literal import rfc8785 @@ -47,6 +47,15 @@ class EvidenceAccumulator: A durable callback must be idempotent by ``event_id`` and return ``True`` only after the entry is durably committed. Callback failure leaves local state unchanged so callers can retry the same event. + + ``append`` and ``seal`` are not reentrant: a durable callback (or anything + else running on the same thread while one of those calls is in progress) + must not call back into ``append`` or ``seal``. Doing so raises + :class:`EvidenceError` instead of deadlocking or silently double-writing a + sequence position. Genuine concurrent calls from other threads are still + safely serialized. ``snapshot`` has no such restriction: it only reads + already-committed entries, so it may be called reentrantly (it simply will + not observe an append that has not finished committing yet). """ def __init__( @@ -70,6 +79,15 @@ def __init__( self._sealed = False self._completeness: Completeness = "unknown" self._lock = Lock() + # Identifies the thread currently inside the append/seal critical + # section, if any. `Lock` is not reentrant: without this check, a + # durable callback (or anything else) that calls back into `append` + # or `seal` on the same thread would deadlock forever rather than + # failing loudly. Compared under the GIL, so plain attribute + # read/write here is race-free the same way `threading.RLock`'s + # pure-Python fallback tracks ownership. + self._owner: int | None = None + @property def mode(self) -> Literal["memory", "callback"]: @@ -81,55 +99,79 @@ def append(self, event: dict[str, Any]) -> EvidenceEntry: event_copy = deepcopy(event) if event_copy["run_id"] != self._run_id: raise EvidenceError("event run_id does not match accumulator run_id") + if self._owner == get_ident(): + raise EvidenceError("reentrant evidence mutation is not supported") + with self._lock: - if self._sealed: - raise EvidenceError("evidence run is sealed") - event_id = event_copy["event_id"] - if event_id in self._event_ids: - raise EvidenceError(f"duplicate event_id: {event_id}") - if len(self._entries) >= self._max_events: - raise EvidenceError(f"evidence run exceeds max_events={self._max_events}") - - sequence = len(self._entries) - previous = self._entries[-1].digest if self._entries else None + self._owner = get_ident() try: - digest = _entry_digest(sequence, previous, event_copy) - except (TypeError, ValueError, rfc8785.CanonicalizationError) as exc: - raise EvidenceError( - f"event_id={event_id} cannot be canonicalized under " - f"{CANONICALIZATION_PROFILE}" - ) from exc - entry = EvidenceEntry(sequence, event_id, previous, digest, event_copy) - - if self._durable_append is not None: + if self._sealed: + raise EvidenceError("evidence run is sealed") + event_id = event_copy["event_id"] + if event_id in self._event_ids: + raise EvidenceError(f"duplicate event_id: {event_id}") + if len(self._entries) >= self._max_events: + raise EvidenceError(f"evidence run exceeds max_events={self._max_events}") + + sequence = len(self._entries) + previous = self._entries[-1].digest if self._entries else None try: - acknowledged = self._durable_append(_copy_entry(entry)) - except Exception as exc: - raise EvidencePersistenceError( - f"durable evidence callback failed for event_id={event_id}" + digest = _entry_digest(sequence, previous, event_copy) + except (TypeError, ValueError, rfc8785.CanonicalizationError) as exc: + raise EvidenceError( + f"event_id={event_id} cannot be canonicalized under " + f"{CANONICALIZATION_PROFILE}" + ) from exc - if acknowledged is not True: - raise EvidencePersistenceError( - f"durable evidence callback did not acknowledge event_id={event_id}" - ) + entry = EvidenceEntry(sequence, event_id, previous, digest, event_copy) + + if self._durable_append is not None: + try: + acknowledged = self._durable_append(_copy_entry(entry)) + except Exception as exc: + raise EvidencePersistenceError( + f"durable evidence callback failed for event_id={event_id}" + ) from exc + if acknowledged is not True: + raise EvidencePersistenceError( + f"durable evidence callback did not acknowledge event_id={event_id}" + ) + + self._entries.append(_copy_entry(entry)) + self._event_ids.add(event_id) + return _copy_entry(entry) + finally: + self._owner = None - self._entries.append(_copy_entry(entry)) - self._event_ids.add(event_id) - return _copy_entry(entry) def seal(self, *, completeness: Completeness) -> EvidenceSnapshot: """Close the run with a caller-asserted completeness assessment.""" if completeness not in ("complete", "incomplete", "unknown"): raise EvidenceError(f"unsupported completeness: {completeness!r}") + if self._owner == get_ident(): + raise EvidenceError("reentrant evidence mutation is not supported") + with self._lock: - if self._sealed: - raise EvidenceError("evidence run is already sealed") - self._sealed = True - self._completeness = completeness - return self._snapshot() + self._owner = get_ident() + try: + if self._sealed: + raise EvidenceError("evidence run is already sealed") + self._sealed = True + self._completeness = completeness + return self._snapshot() + finally: + self._owner = None + def snapshot(self) -> EvidenceSnapshot: + # A snapshot only reads already-committed state, so it is safe to let + # it run reentrantly on the thread already holding `_lock` (it just + # will not see a mutation still in flight). Re-entering `_lock` + # itself would deadlock, since it is not a reentrant lock. + if self._owner == get_ident(): + return self._snapshot() + with self._lock: return self._snapshot() diff --git a/tests/test_evidence.py b/tests/test_evidence.py index b223ba2..4b8e8e7 100644 --- a/tests/test_evidence.py +++ b/tests/test_evidence.py @@ -173,6 +173,66 @@ def test_snapshot_is_defensive_and_unsealed_never_claims_complete(self): self.assertEqual(fresh.completeness, "unknown") self.assertEqual(fresh.entries[0].event["run_id"], "run-governed-sdlc-001") + + def test_reentrant_durable_append_call_is_refused_not_deadlocked(self): + accumulator = None + + def durable_append(entry): + # A durable backend that (accidentally or otherwise) calls back + # into the same accumulator while the outer append() is still in + # progress must be refused, not hang the caller forever. + accumulator.append(fixture("usage.json")) + return True + + accumulator = EvidenceAccumulator( + "run-governed-sdlc-001", self.validator, durable_append=durable_append + ) + with self.assertRaises(EvidencePersistenceError): + accumulator.append(fixture("policy-decision.json")) + self.assertEqual(accumulator.snapshot().entries, ()) + + # The accumulator must remain fully usable afterwards: the failed + # reentrant attempt must not leave it permanently stuck. + clean = EvidenceAccumulator("run-governed-sdlc-001", self.validator) + accepted = clean.append(fixture("policy-decision.json")) + self.assertEqual(accepted.sequence, 0) + + def test_reentrant_seal_call_is_refused_not_deadlocked(self): + accumulator = None + + def durable_append(entry): + accumulator.seal(completeness="complete") + return True + + accumulator = EvidenceAccumulator( + "run-governed-sdlc-001", self.validator, durable_append=durable_append + ) + with self.assertRaises(EvidencePersistenceError): + accumulator.append(fixture("policy-decision.json")) + snapshot = accumulator.snapshot() + self.assertEqual(snapshot.entries, ()) + self.assertFalse(snapshot.sealed) + + def test_reentrant_snapshot_call_is_safe_and_reflects_committed_state(self): + accumulator = None + observed = [] + + def durable_append(entry): + # Reading a snapshot from inside the callback is safe: it must + # not deadlock, and it must only see already-committed entries, + # not the one still being appended. + observed.append(accumulator.snapshot()) + return True + + accumulator = EvidenceAccumulator( + "run-governed-sdlc-001", self.validator, durable_append=durable_append + ) + accepted = accumulator.append(fixture("policy-decision.json")) + self.assertEqual(len(observed), 1) + self.assertEqual(observed[0].entries, ()) + self.assertEqual(accumulator.snapshot().entries, (accepted,)) + + def test_callback_and_returned_entry_cannot_mutate_retained_evidence(self): callback_entries = []