[Bug]: Cluster mode: SubscribeToTask on a replica holding a paused ActiveTask returns a stale INPUT_REQUIRED snapshot and closes
I maintainer di solito rispondono entro 2 giorni
Valutazione
Questa issue non è ancora stata valutata.
Descrizione
What happened?
Version: a2a-sdk 1.2.2 (code unchanged on main), cluster mode (VersionedDatabaseTaskStore + DatabaseTaskEventStream).
In cluster mode, SubscribeToTask on a replica that ran an earlier turn of a multi-turn task returns a stale snapshot and closes the stream right away, while the current turn is still running on another replica. A replica that never saw the task streams the same turn correctly.
Steps
- Replica A runs turn 1. The agent pauses with
TASK_STATE_INPUT_REQUIRED. A keeps itsActiveTaskinActiveTaskRegistry, because only terminal states trigger cleanup. - The client's reply lands on replica B, which runs turn 2 (
WORKING…COMPLETED). - While turn 2 is running, the client resubscribes. The load balancer sends it to replica A.
Expected: a snapshot of the current state (WORKING), then turn 2's events from the shared event stream, ending at COMPLETED. This is what replica C returns.
Actual: replica A returns Task(TASK_STATE_INPUT_REQUIRED), the state from turn 1, and closes the stream. The client concludes the agent is waiting for input, even though the task is running and later completes.
resubscribe on C (never saw the task): Task(WORKING) -> StatusUpdate(WORKING) x3 -> StatusUpdate(COMPLETED) -> <stream closed>
resubscribe on A (ran turn 1): Task(INPUT_REQUIRED) -> <stream closed>
GetTask afterwards: TASK_STATE_COMPLETED
This is consistent: 5 out of 5 runs.
Cause
on_subscribe_to_task (default_request_handler_v2.py, around line 575) takes the local fast path whenever the replica's registry has an entry for the task:
# Shared-stream mode. Fast path: this replica runs the agent -> tap it.
local = await self._active_task_registry.get(task_id)
if local is not None:
async for event in local.subscribe(include_initial_task=True):
yield event
return
Once a turn has ended in an interrupted state, having a registry entry no longer means "this replica runs the agent". The paused ActiveTask on A is idle:
local.subscribe(include_initial_task=True)returns itsTaskManager's cached snapshot, which still saysINPUT_REQUIRED. Unlike the producer loop, this path doesn't callinvalidate()first.- A then follows its own in-memory queue, which never receives turn 2's events because they are produced on B.
The handler already read the current snapshot and version from _versioned_store a few lines earlier, so the remote path (_subscribe_remote) would return the right answer.
Possible fixes
- In shared-stream mode, take the local fast path only while the local
ActiveTaskhas a request in flight. Otherwise use_subscribe_remote(task_id, task, snapshot_version, stream). - Or evict an
ActiveTaskfrom the registry when its turn ends in an interrupted state, when cluster mode is on. That would also stop paused tasks from staying in memory indefinitely on every replica that ever ran a turn for them.
Related: #1188 and #1281. This is the same root cause (a paused ActiveTask outliving its turn), but on the subscribe path, which #1281's per-request invalidate() doesn't cover.
Repro (standalone; needs only a2a-sdk[sqlite] or aiosqlite):
repro_resubscribe_stale.py
"""Resubscribe on a replica that holds a paused ActiveTask returns a stale snapshot.
a2a-sdk 1.2.2, cluster mode (VersionedDatabaseTaskStore + DatabaseTaskEventStream),
three replicas sharing one SQLite database.
"""
import asyncio
import os
import tempfile
from sqlalchemy.ext.asyncio import create_async_engine
from a2a.helpers.proto_helpers import get_message_text
from a2a.server.agent_execution import AgentExecutor, RequestContext
from a2a.server.cluster import DatabaseTaskEventStream, VersionedDatabaseTaskStore
from a2a.server.context import ServerCallContext
from a2a.server.events.event_queue import EventQueue
from a2a.server.request_handlers import DefaultRequestHandler
from a2a.server.tasks.task_updater import TaskUpdater
from a2a.types import a2a_pb2 as pb
DB = os.path.join(tempfile.mkdtemp(), "cluster.db")
class TwoTurnAgent(AgentExecutor):
"""Turn 1 asks for input; turn 2 works for ~1.6s, then completes."""
async def execute(self, context: RequestContext, queue: EventQueue) -> None:
if context.current_task is None:
await queue.enqueue_event(pb.Task(
id=context.task_id, context_id=context.context_id,
status=pb.TaskStatus(state=pb.TASK_STATE_SUBMITTED),
history=[context.message]))
updater = TaskUpdater(queue, context.task_id, context.context_id)
if get_message_text(context.message) == "hi":
await updater.requires_input()
return
for _ in range(4):
await updater.start_work()
await asyncio.sleep(0.4)
await updater.complete()
async def cancel(self, context: RequestContext, queue: EventQueue) -> None:
await TaskUpdater(queue, context.task_id, context.context_id).cancel()
CARD = pb.AgentCard(name="repro", capabilities=pb.AgentCapabilities(streaming=True))
STORES = []
def replica() -> DefaultRequestHandler:
engine = create_async_engine(f"sqlite+aiosqlite:///{DB}")
store, stream = VersionedDatabaseTaskStore(engine), DatabaseTaskEventStream(engine)
STORES.extend([store, stream])
return DefaultRequestHandler(agent_executor=TwoTurnAgent(), task_store=store,
event_stream=stream, agent_card=CARD)
def send(text: str, task_id: str | None = None) -> pb.SendMessageRequest:
msg = pb.Message(message_id=text, role=pb.ROLE_USER, context_id="ctx",
parts=[pb.Part(text=text)])
if task_id:
msg.task_id = task_id
return pb.SendMessageRequest(
message=msg, configuration=pb.SendMessageConfiguration(return_immediately=True))
def label(event) -> str:
status = getattr(event, "status", None)
return f"{type(event).__name__}({pb.TaskState.Name(status.state)})" if status else type(event).__name__
async def resubscribe(name: str, handler: DefaultRequestHandler, task_id: str) -> None:
events = []
async def tail() -> None:
async for event in handler.on_subscribe_to_task(
pb.SubscribeToTaskRequest(id=task_id), ServerCallContext()
):
events.append(label(event))
try:
await asyncio.wait_for(tail(), timeout=5)
events.append("<stream closed>")
except asyncio.TimeoutError:
events.append("<still open after 5s>")
print(f"resubscribe on {name}: {' -> '.join(events)}")
async def main() -> None:
a, b, c = replica(), replica(), replica()
for s in STORES:
await s.initialize()
ctx = ServerCallContext()
task = await a.on_message_send(send("hi"), ctx) # turn 1 on A -> INPUT_REQUIRED
await asyncio.sleep(0.3)
await b.on_message_send(send("go", task.id), ctx) # turn 2 runs on B
await asyncio.sleep(0.1)
await asyncio.gather(
resubscribe("A (ran turn 1)", a, task.id),
resubscribe("C (never saw the task)", c, task.id),
)
final = await c.on_get_task(pb.GetTaskRequest(id=task.id), ctx)
print("GetTask afterwards:", pb.TaskState.Name(final.status.state))
asyncio.run(main())
Relevant log output
resubscribe on C (never saw the task): Task(TASK_STATE_WORKING) -> TaskStatusUpdateEvent(TASK_STATE_WORKING) -> TaskStatusUpdateEvent(TASK_STATE_WORKING) -> TaskStatusUpdateEvent(TASK_STATE_WORKING) -> TaskStatusUpdateEvent(TASK_STATE_COMPLETED) -> <stream closed>
resubscribe on A (ran turn 1): Task(TASK_STATE_INPUT_REQUIRED) -> <stream closed>
GetTask afterwards: TASK_STATE_COMPLETED
Code of Conduct
- I agree to follow this project's Code of Conduct
- Lingua principale
- Python
- Stelle
- 2.2k
- Fork
- 509
- Merge medio
- 3g 18h
- PR unite (30g)
- 44
Preparare l'ambiente
- Nessun Dockerfile né file Docker Compose
- Ha un modello di pull request
- Leggi la guida per i contributori
Come iniziare
- Leggi tutta la issue e poi la guida ai contributi del progetto.
- Commenta sulla issue per dire che te ne occupi tu — evita che due persone facciano lo stesso lavoro.
- Fai un fork del repository e lavora su un branch.
- Apri una pull request che faccia riferimento al numero della issue.
Altre issue di a2aproject/a2a-python
-
[Bug]: REST task/request id sanitizationForse già presa @Linux2010 l’ha presa 103 giorni fa. Apertamaintainers-only
Difficoltà 2/5 1-3 ore Idoneità per principianti 65/100
a2aproject/a2a-python#805 · 1 commento ·
I maintainer di solito rispondono entro 2 giorni
-
[Feat]: Cluster mode: detect and recover tasks abandoned by a crashed replica (heartbeat/lease)Forse già presa @rohityan l’ha presa 2 giorni fa. Aperta
a2aproject/a2a-python#1324 · 1 assegnatario ·
I maintainer di solito rispondono entro 2 giorni
-
[Bug]: After a streamed task is cancelled, its background producer never finishesForse già presa @rohityan l’ha presa 2 giorni fa. Apertacomponent: server status:awaiting response
a2aproject/a2a-python#1322 · 1 commento · 1 assegnatario ·
I maintainer di solito rispondono entro 2 giorni
-
v0.3 gRPC and REST SendMessage without configuration run non-blockingForse già presa @rohityan l’ha presa 2 giorni fa. Apertaquestion status:awaiting response
a2aproject/a2a-python#1321 · 1 commento · 1 assegnatario ·
I maintainer di solito rispondono entro 2 giorni
-
[Bug]: Push notification store failure rewrites a completed task as FAILED (DefaultRequestHandlerV2)Forse già presa @rohityan l’ha presa 4 giorni fa. Apertacomponent: server status:awaiting response
a2aproject/a2a-python#1313 · 1 commento · 1 assegnatario ·
I maintainer di solito rispondono entro 2 giorni
Tutte le issue di a2aproject/a2a-python
Issue simili
-
HTML: <template> content is extracted as document textForse già presa @ryanmeowy l’ha presa oggi. Apertabug html
Difficoltà 1/5 Meno di un'ora Idoneità per principianti 82/100
docling-project/docling#4714 · 2 commenti ·
I maintainer di solito rispondono entro 1 giorno
-
[BUG] Qdrant RAG client applies score_threshold to raw cosine similarity, not the 0-1 score it returnsForse già presa @roydonsequeira l’ha presa oggi. Apertabug
Difficoltà 2/5 1-3 ore Idoneità per principianti 70/100
I maintainer di solito rispondono entro 1 giorno
-
Host test failure in core/direct_io.zig on Linux kernel 6.17: O_DIRECT open succeeds on procfs, so the test's 'plain' fd is not plainForse già presa Una pull request collegata a questa issue è aperta o già unita. Aperta
Difficoltà 2/5 1-3 ore Idoneità per principianti 72/100
ashhart/TensorFold#536 ·
I maintainer di solito rispondono entro 1 giorno
-
area/install-update comp/cli duplicate P2 python:uv sweeper:risk-compatibility type/bug
Difficoltà 1/5 Meno di un'ora Idoneità per principianti 62/100
NousResearch/hermes-agent#135440 · 1 commento ·
I maintainer di solito rispondono entro 1 giorno
-
bug
Difficoltà 2/5 1-3 ore Idoneità per principianti 85/100
Deepak3699/Ai_Mentor#244 ·
I maintainer di solito rispondono entro 1 giorno