Repository navigation
[Feat]: Cluster mode: detect and recover tasks abandoned by a crashed replica (heartbeat/lease) #1324
Description
Activity
- addedcomponent: serverIssues related to frameworks for agent execution, HTTP/event handling, database persistence logic.Issues related to frameworks for agent execution, HTTP/event handling, database persistence logic.
on Oct 9, 2026 Hi @KashyapNasit , Thanks for putting this together. Handling orphaned tasks during abrupt replica crashes is definitely a blind spot in cluster mode right now. Your proposal has been received by the team and is currently under review.
Hey @ishymko , PTAL
The default
failpolicy looks useful for ending orphaned client streams. Two additional fault-injection cases may help specify the optionalre-executecontract:Injected boundary Proposed assertion A downstream service accepts a non-idempotent action; the replica dies before saving its result; the sweeper records FAILED.Default recovery writes the terminal event without automatically issuing the action again. Preserve the unresolved external outcome in application evidence/status details. Permit re-execution only through a documented replay contract, such as provider-enforced idempotency for the same logical action or authoritative reconciliation establishing that it was not executed. An old executor is paused just before a downstream call; its lease expires; recovery advances the task; the old executor resumes. Assert that its stale task-store write is rejected, and separately count downstream dispatch attempts and accepted effects. If effect fencing is part of the recovery guarantee, test the resource-enforced generation or execution-boundary check that rejects the old owner. The second case follows from the current ordering in
EventConsumer._process_event: it catches the save conflict, reloads terminal state, then cancels the producer. My inference is that this fences task persistence at the next event/save boundary; downstream effects need their own enforcement contract if an executor can call a service before that boundary.For retries, keep the logical action's idempotency identity stable across attempts, and revalidate its execution authority. A new attempt counter alone does not establish that another effect is permitted.
These are proposed test oracles and a source-ordering observation, prepared with AI assistance. I have not run a multi-replica/lease recovery reproduction or verified an implementation of this proposal.
Follow-up on the stale-owner test vector above: I executed a three-case characterization on unchanged SDK
494a8ece0ad9815afd8ebd59e9a818e286eb0d80usingDefaultRequestHandlerV2/ActiveTask, the realVersionedDatabaseTaskStore, andDatabaseTaskEventStream.Two handler objects and two SQLAlchemy engines share a SQLite file in one Python process. The executor pauses after
WORKINGis durably stored. A direct versionedFAILEDsave, including its event, substitutes for the proposed sweeper's decision; then the old executor resumes.Case Dispatch attempts Accepted synthetic effects Old writer's SQL CAS conflicts Stored task / remote stream Normal completion control 1 1 0 COMPLETED / COMPLETED FAILEDwrite, task-store fence only1 1 1 FAILED / FAILED FAILEDwrite plus a synthetic resource generation fence1 0 1 FAILED / FAILED Before resuming the executor, both failure cases assert that the remote
SubscribeToTaskstream has ended asFAILED, the executor has not received cancellation, and both effect counters are still zero. After resume, the next local status save is rejected and the consumer cancels the producer; the durableFAILEDstate survives. The unfenced resource has already accepted its one synthetic effect by that point.This confirms the proposal's described next-save behavior and makes its scope measurable. A recovery contract that also promises effect fencing needs a separate assertion over accepted downstream effects. In the last control, the local resource itself enforces an independently advanced generation; that is a test double, not an SDK feature or a verified provider integration. These cases do not implement or test a heartbeat, lease expiration, actual replica death, multi-process deployment, or the separate lost-ack/re-execute vector.
All 3 characterization cases pass. Environment: macOS, Python 3.10.20, SQLAlchemy 2.0.48, aiosqlite 0.22.1, protobuf 7.36.2, pytest 9.0.3, repository
uv.lock. Required./scripts/lint.shand explicit fixturety checkexit 0. Fulluv run pytest --timeout=45reports 2,272 passed, 183 skipped, 3 xfailed, 1 xpassed. An initial run missed the legacy 0.3 server's startup deadline; the full rerun passed after warming that isolated dependency environment, with no SDK source changes. PostgreSQL/MySQL services were not configured for this local run.Complete fixture and reproduction command
Save as
tests/server/cluster/test_stale_executor_effect_boundary.pyat the pinned revision. It reuses the repository's existing context/request/handler helpers from that directory'sconftest.py."""Characterize task CAS separately from a synthetic downstream effect fence.""" import asyncio from collections.abc import Callable from pathlib import Path import pytest from a2a.helpers.proto_helpers import new_task_from_user_message from a2a.server.agent_execution.agent_executor import AgentExecutor from a2a.server.agent_execution.context import RequestContext from a2a.server.cluster import ConcurrentTaskModificationError, TaskVersion from a2a.server.cluster.database_event_stream import DatabaseTaskEventStream from a2a.server.cluster.database_task_store import VersionedDatabaseTaskStore from a2a.server.context import ServerCallContext from a2a.server.events.event_queue import Event from a2a.server.events.event_queue_v2 import EventQueue from a2a.server.tasks.task_updater import TaskUpdater from a2a.types.a2a_pb2 import ( SubscribeToTaskRequest, Task, TaskState, TaskStatus, TaskStatusUpdateEvent, ) from sqlalchemy.ext.asyncio import create_async_engine from .conftest import ( build_send_request, make_context, make_replica, wait_for_state, ) class CountingSQLStore(VersionedDatabaseTaskStore): """Count real SQL compare-and-swap rejections without replacing its logic.""" conflicts = 0 async def save( self, task: Task, *, event: Event | None = None, prev: Task | None = None, prev_version: TaskVersion, context: ServerCallContext, ) -> TaskVersion: try: return await super().save( task, event=event, prev=prev, prev_version=prev_version, context=context, ) except ConcurrentTaskModificationError: self.conflicts += 1 raise class SyntheticResource: """Local effect counter with an optional independently enforced generation.""" def __init__(self) -> None: self.generation = 0 self.attempts = 0 self.accepted = 0 def dispatch(self, expected_generation: int | None) -> None: self.attempts += 1 if ( expected_generation is None or expected_generation == self.generation ): self.accepted += 1 class PausedEffectAgent(AgentExecutor): """Pause before dispatch, then emit a status event that exercises task CAS.""" def __init__( self, resource: SyntheticResource, expected_generation: int | None ) -> None: self.resource = resource self.expected_generation = expected_generation self.ready = asyncio.Event() self.resume = asyncio.Event() self.dispatched = asyncio.Event() self.cancelled = asyncio.Event() self.hold = asyncio.Event() async def execute( self, context: RequestContext, event_queue: EventQueue ) -> None: if context.current_task is None: assert context.message is not None await event_queue.enqueue_event( new_task_from_user_message(context.message) ) updater = TaskUpdater( event_queue, str(context.task_id or ''), str(context.context_id or ''), ) await updater.start_work() self.ready.set() try: await self.resume.wait() self.resource.dispatch(self.expected_generation) self.dispatched.set() await updater.complete() # Stay cancellable so the stale-write consumer's cancellation is observed. await self.hold.wait() except asyncio.CancelledError: self.cancelled.set() raise async def cancel( self, context: RequestContext, event_queue: EventQueue ) -> None: del context, event_queue @pytest.mark.asyncio @pytest.mark.timeout(15) @pytest.mark.parametrize( 'advance_terminal,enforce_effect_generation,expected_accepted', [(False, False, 1), (True, False, 1), (True, True, 0)], ids=[ 'normal-control', 'terminal-task-only', 'terminal-with-resource-fence', ], ) async def test_stale_executor_effect_boundary( tmp_path: Path, record_property: Callable[[str, object], None], advance_terminal: bool, enforce_effect_generation: bool, expected_accepted: int, ) -> None: dsn = f'sqlite+aiosqlite:///{tmp_path / "shared.sqlite"}' engine_a = create_async_engine(dsn) engine_b = create_async_engine(dsn) store_a = CountingSQLStore(engine=engine_a) store_b = VersionedDatabaseTaskStore(engine=engine_b) await store_a.initialize() await store_b.initialize() stream_a = DatabaseTaskEventStream(engine_a, poll_interval_s=0.01) stream_b = DatabaseTaskEventStream(engine_b, poll_interval_s=0.01) resource = SyntheticResource() agent = PausedEffectAgent( resource, 0 if enforce_effect_generation else None ) handler_a = make_replica(store_a, stream_a, agent) handler_b = make_replica(store_b, stream_b, agent) context = make_context() task_seen = asyncio.Event() remote_snapshot_seen = asyncio.Event() local_events: list[Event] = [] remote_events: list[Event] = [] task_id = '' remote_task: asyncio.Task[None] | None = None async def collect_local() -> None: nonlocal task_id async for event in handler_a.on_message_send_stream( build_send_request('go', message_id='mission-stale-boundary'), context, ): local_events.append(event) if isinstance(event, Task): task_id = event.id task_seen.set() async def collect_remote() -> None: async for event in handler_b.on_subscribe_to_task( SubscribeToTaskRequest(id=task_id), context ): remote_events.append(event) if isinstance(event, Task): remote_snapshot_seen.set() local_task = asyncio.create_task(collect_local()) try: await asyncio.wait_for(agent.ready.wait(), timeout=5) await asyncio.wait_for(task_seen.wait(), timeout=5) await wait_for_state( store_a, task_id, TaskState.TASK_STATE_WORKING, context ) remote_task = asyncio.create_task(collect_remote()) await asyncio.wait_for(remote_snapshot_seen.wait(), timeout=5) snapshot = await store_b.get(task_id, context) assert snapshot is not None assert snapshot.task.status.state == TaskState.TASK_STATE_WORKING if advance_terminal: failed = Task() failed.CopyFrom(snapshot.task) failed.status.CopyFrom( TaskStatus(state=TaskState.TASK_STATE_FAILED) ) failure_event = TaskStatusUpdateEvent( task_id=task_id, context_id=failed.context_id, status=failed.status, ) await store_b.save( failed, event=failure_event, prev=snapshot.task, prev_version=snapshot.version, context=context, ) resource.generation += 1 # The remote stream ends while the old executor is still paused. await asyncio.wait_for(remote_task, timeout=5) remote_terminal = remote_events[-1] assert isinstance(remote_terminal, Task | TaskStatusUpdateEvent) assert remote_terminal.status.state == TaskState.TASK_STATE_FAILED assert not agent.cancelled.is_set() assert resource.attempts == resource.accepted == 0 agent.resume.set() await asyncio.wait_for(agent.dispatched.wait(), timeout=5) await asyncio.wait_for(local_task, timeout=5) await asyncio.wait_for(remote_task, timeout=5) final = await store_b.get(task_id, context) assert final is not None expected_state = ( TaskState.TASK_STATE_FAILED if advance_terminal else TaskState.TASK_STATE_COMPLETED ) assert final.task.status.state == expected_state assert resource.attempts == 1 assert resource.accepted == expected_accepted assert store_a.conflicts == int(advance_terminal) if advance_terminal: await asyncio.wait_for(agent.cancelled.wait(), timeout=5) record_property('dispatch_attempts', resource.attempts) record_property('accepted_synthetic_effects', resource.accepted) record_property('sql_cas_conflicts', store_a.conflicts) record_property('stored_state', TaskState.Name(final.task.status.state)) remote_final = remote_events[-1] assert isinstance(remote_final, Task | TaskStatusUpdateEvent) record_property( 'remote_final_state', TaskState.Name(remote_final.status.state) ) finally: agent.hold.set() await handler_a.aclose() await handler_b.aclose() for pending in (local_task, remote_task): if pending is not None and not pending.done(): pending.cancel() await asyncio.gather(pending, return_exceptions=True) await engine_a.dispose() await engine_b.dispose()
uv sync --frozen uv run pytest -q tests/server/cluster/test_stale_executor_effect_boundary.py \ -o junit_family=xunit1 --junitxml=stale-effect.xml
The JUnit properties record dispatches, accepted synthetic effects, SQL conflicts, and stored/remote final states. Fixture SHA-256:
2b5acc206f12c89def62dcc1d69ada06d04efc09332773f8a39127ef5f9c1da2.AI-assisted characterization and fixture by Mission (
gomission) with Codex; existing cluster test helpers credited to this repository.
Is your feature request related to a problem? Please describe.
Cluster mode (#1281) handles replicas that are alive and disagree: version checks on saves, cancel from any replica, resubscribe from any replica. It does nothing for a replica that dies while running
execute()(OOM kill, node loss, a rolling deploy without enough grace time).When that happens, the task stays non-terminal forever:
GetTaskreturnsTASK_STATE_WORKINGindefinitely.SubscribeToTaskon another replica keeps pollingtask_eventsand never ends, because no terminal event will ever be written.CancelTask.The client can't tell a dead replica from a slow agent.
status.timestamponly moves when the agent emits an event, so a long LLM or tool call looks exactly like a crash.Detecting abandoned tasks was a stated goal of the Go SDK's distributed mode (a2a-go discussion #139: "Supporting abandoned task detection (e.g. on process reboot)"), and a2a-go ships it (a2a-go #115,
a2asrv/workqueue). The Python design discussion (#1224) doesn't cover it.Describe the solution you'd like
An optional extension point for cluster mode, similar to a2a-go's
workqueue.Heartbeater/LeaseManager, along these lines:execute()runs, the replica renews a lease (task_id,owner,expires_at, plus the task's owner scope) on an interval.INPUT_REQUIRED/AUTH_REQUIREDis never treated as abandoned.FOR UPDATE SKIP LOCKED, so no leader is needed) and applies a policy:fail(default): writeTASK_STATE_FAILED;re-execute: run the turn again, with an attempt counter.VersionedTaskStore.save(..., event=...). This matters in cluster mode:SubscribeToTaskstream withFAILED;ConcurrentTaskModificationError, it re-reads, seesFAILED, and cancels its producer, the same way a remote cancel works today;FAILED.Pitfalls we hit running a heartbeat plus sweeper outside the SDK. They're probably worth designing in from the start:
execute()returns. The terminal or interrupted status event may still be on the event queue whenexecute()returns. A crash in that gap leaves a non-terminal task with no lease, which no sweep can ever find.INPUT_REQUIREDstate (update_with_messagedoesn't changestatus.state). Releasing on "saved an interrupted state" drops the lease at the start of the resumed turn.DatabaseTaskStore/VersionedDatabaseTaskStorefilter reads byowner_resolver(context). A sweeper callingget()with an emptyServerCallContextgetsNonefor authenticated tasks and silently does nothing. Store the owner scope with the lease and rebuild the context from it.Describe alternatives you've considered
TaskStorewrapper that releases the lease after the durable save. It works, but every deployment has to rebuild it. In cluster mode it also needs the versioned save path and the event table to end remote streams correctly, and those are SDK internals an application shouldn't have to coordinate with.Additional context
Code of Conduct