Skip to content

[Feat]: Cluster mode: detect and recover tasks abandoned by a crashed replica (heartbeat/lease) #1324

Description

@KashyapNasit

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:

  • GetTask returns TASK_STATE_WORKING indefinitely.
  • Every SubscribeToTask on another replica keeps polling task_events and never ends, because no terminal event will ever be written.
  • Nothing ever re-runs or fails the task. The only way out is a client-issued CancelTask.
  • A push-only client never hears anything again.

The client can't tell a dead replica from a slow agent. status.timestamp only 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:

  1. Lease while executing. While execute() runs, the replica renews a lease (task_id, owner, expires_at, plus the task's owner scope) on an interval.
  2. Release when the turn ends. Release the lease when the task moves into a terminal or interrupted state, so a task waiting on INPUT_REQUIRED / AUTH_REQUIRED is never treated as abandoned.
  3. Sweep. On every replica, a sweeper atomically claims leases that have expired (for example FOR UPDATE SKIP LOCKED, so no leader is needed) and applies a policy:
    • fail (default): write TASK_STATE_FAILED;
    • re-execute: run the turn again, with an attempt counter.
  4. Write the failure through VersionedTaskStore.save(..., event=...). This matters in cluster mode:
    • the event row ends every remote SubscribeToTask stream with FAILED;
    • the version bump fences a replica that wasn't really dead (a long GC pause or a network partition). Its next save gets ConcurrentTaskModificationError, it re-reads, sees FAILED, and cancels its producer, the same way a remote cancel works today;
    • push senders can then notify FAILED.

Pitfalls we hit running a heartbeat plus sweeper outside the SDK. They're probably worth designing in from the start:

  • Release after the durable save, not when execute() returns. The terminal or interrupted status event may still be on the event queue when execute() returns. A crash in that gap leaves a non-terminal task with no lease, which no sweep can ever find.
  • Release on a state change, not on whatever state is being saved. On a resumed turn, the first save is the held user-message flush, and it re-saves the stale INPUT_REQUIRED state (update_with_message doesn't change status.state). Releasing on "saved an interrupted state" drops the lease at the start of the resumed turn.
  • A renewal must not recreate a released lease. If renewal is an upsert, a heartbeat that lands just after the release brings the lease back. The sweeper later claims it and fails a task that was correctly paused.
  • The sweeper needs the owner scope. DatabaseTaskStore / VersionedDatabaseTaskStore filter reads by owner_resolver(context). A sweeper calling get() with an empty ServerCallContext gets None for authenticated tasks and silently does nothing. Store the owner scope with the lease and rebuild the context from it.

Describe alternatives you've considered

  • Handle it in the application (what we do today): heartbeat rows, an orphan sweeper, and a TaskStore wrapper 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.
  • Leave it out of scope, as a2a-rs did (a2a-rs#149). That's reasonable for single-process servers. But cluster mode is aimed specifically at multi-replica deployments, and replicas dying is the normal failure in that setting.
  • A client-side overall timeout is a useful backstop, but the client can't tell which tasks are actually dead.
  • The spec-level "Expectations" proposal (A2A#2209) is about business deadlines, not process liveness.

Additional context

  • Defaults that worked for us: heartbeat every 30s, a task counts as abandoned after 120s without one, and the sweep runs every 30s. That's roughly 2–2.5 minutes to detect a dead replica.
  • Happy to share more detail or help with an implementation if the maintainers think this belongs in the SDK.

Code of Conduct

  • I agree to follow this project's Code of Conduct

Activity

  1. added theissue type on Oct 9, 2026
  2. added
    component: serverIssues related to frameworks for agent execution, HTTP/event handling, database persistence logic.
    on Oct 9, 2026
  3. rohityan commented on Oct 9, 2026

    @rohityan

    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.

  4. assigned and unassigned on Oct 9, 2026
  5. rohityan commented on Oct 9, 2026

    @rohityan

    Hey @ishymko , PTAL

  6. gomission commented on Oct 10, 2026

    @gomission

    The default fail policy looks useful for ending orphaned client streams. Two additional fault-injection cases may help specify the optional re-execute contract:

    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.

  7. gomission commented on Oct 10, 2026

    @gomission

    Follow-up on the stale-owner test vector above: I executed a three-case characterization on unchanged SDK 494a8ece0ad9815afd8ebd59e9a818e286eb0d80 using DefaultRequestHandlerV2 / ActiveTask, the real VersionedDatabaseTaskStore, and DatabaseTaskEventStream.

    Two handler objects and two SQLAlchemy engines share a SQLite file in one Python process. The executor pauses after WORKING is durably stored. A direct versioned FAILED save, 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
    FAILED write, task-store fence only 1 1 1 FAILED / FAILED
    FAILED write plus a synthetic resource generation fence 1 0 1 FAILED / FAILED

    Before resuming the executor, both failure cases assert that the remote SubscribeToTask stream has ended as FAILED, 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 durable FAILED state 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.sh and explicit fixture ty check exit 0. Full uv run pytest --timeout=45 reports 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.py at the pinned revision. It reuses the repository's existing context/request/handler helpers from that directory's conftest.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.

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

Labels

component: serverIssues related to frameworks for agent execution, HTTP/event handling, database persistence logic.status: needs review

Projects

No projects

    Milestone

    No milestone

    Relationships

    None yet

    Development

    No branches or pull requests

    Issue actions