When one agent asks another agent to do something that takes longer than a single request, both sides need a shared handle for that piece of work: something to ask about, cancel, stream from and come back to after a network drop. In the Agent2Agent (A2A) protocol that handle is the Task. A client agent sends a message; the remote agent either answers directly with a message or creates a task and reports its progress through a small set of states until it finishes.
This article is about building the task layer on the remote side, and using it well from the client side. It works from the A2A 1.0 specification, whose data model is defined in the a2a.proto file of the a2aproject/A2A repository and published at a2a-protocol.org. It covers the Task object on the wire, a state machine you can enforce, a five-component server design with working Python, a worked delegation, and the failure modes that show up once tasks run for minutes instead of milliseconds. The lifecycle story is told in A2A task state; here the focus is the machinery underneath it.
What a task is, and when you do not need one
A task is a stateful unit of work with an identity that outlives the request that created it. It carries an id, the contextId of the conversation it belongs to, a current status, the artifacts it has produced, an optional history of messages and free-form metadata. The status is itself a small object: a state, an optional agent message explaining it, and a timestamp.
Not every exchange needs one. A2A lets the remote agent reply to SendMessage with either a Task or a plain Message. If the answer is quick, deterministic and has nothing to cancel, return a Message and keep no state. Create a task when any of these is true: the work may outlast an HTTP timeout, the caller may want to cancel it, it produces files or structured results in pieces, or it may pause to ask the caller for more input. Creating tasks for everything fills your store with one-second tasks nobody ever queries.
Two identifiers do different jobs. The task id names one unit of work. The context id groups related tasks and messages into one conversation, so a follow-up request after a task has finished becomes a new task in the same context. A finished task is never reopened: the specification states that tasks in a terminal state cannot accept further messages.
The Task on the wire
A2A 1.0 defines its operations as SendMessage, SendStreamingMessage, GetTask, ListTasks, CancelTask and SubscribeToTask, plus create, get, list and delete operations for task push-notification configs and GetExtendedAgentCard. Pre-1.0 SDKs used slash-style method names such as message/send and lower-case state strings; if you interoperate with agents pinned to 0.3, expect to translate. In JSON the proto fields appear in lowerCamelCase. A non-blocking request that starts a reconciliation job looks like this:
{
"jsonrpc": "2.0",
"id": "req-41",
"method": "SendMessage",
"params": {
"message": {
"messageId": "9f2c6a3e-0d51-4c1e-9a7b-2f0f3c1d8e11",
"role": "ROLE_USER",
"parts": [{"text": "Reconcile September invoices against the ledger export."},
{"url": "https://files.example.com/ledger-2026-09.csv",
"mediaType": "text/csv", "filename": "ledger-2026-09.csv"}]
},
"configuration": {
"acceptedOutputModes": ["application/json", "text/plain"],
"historyLength": 0,
"returnImmediately": true
}
}
}The messageId is required and is generated by the client; it is your idempotency key. A part holds exactly one of text, raw bytes, a url or structured data, with optional mediaType and filename. historyLength of 0 asks the server to omit history from the response, and returnImmediately asks it not to wait for the work. The response wraps either a task or a message:
{
"jsonrpc": "2.0",
"id": "req-41",
"result": {
"task": {
"id": "tsk_01JB7Q",
"contextId": "ctx_01JB7P",
"status": {"state": "TASK_STATE_SUBMITTED", "timestamp": "2026-10-01T05:58:03Z"}
}
}
}| State | Kind | Meaning for the client |
|---|---|---|
TASK_STATE_SUBMITTED | Active | Accepted and queued, nothing done yet |
TASK_STATE_WORKING | Active | An executor is on it |
TASK_STATE_INPUT_REQUIRED | Interrupted | Paused; reply with a message carrying the same taskId |
TASK_STATE_AUTH_REQUIRED | Interrupted | Paused until the caller supplies authorization |
TASK_STATE_COMPLETED | Terminal | Artifacts are final |
TASK_STATE_FAILED | Terminal | Ended with an error |
TASK_STATE_CANCELED | Terminal | Ended by a cancel request |
TASK_STATE_REJECTED | Terminal | The agent declined the work |
A server built from five parts
A production task layer separates five responsibilities. The API handler validates requests, authenticates the caller and checks the messageId against recent ones. The task store holds the authoritative row for each task, with a version number. A work queue decouples accepting a task from running it, so a burst of requests becomes a backlog instead of a crash. Executors claim task ids and run the actual agent logic. Every state change and artifact chunk is appended to a per-task event log, which feeds both streaming subscribers and the push notifier.
The single rule that keeps this design correct is that the store is the only source of truth. Streams and webhooks are notifications that the store changed. A client that misses an event calls GetTask and reads the row.
Code: the task manager
The sketch below is an in-memory version of the store, dedupe map and event log. Every transition is a compare-and-set on the version: if a cancel and a completion race, exactly one of them wins and the other gets a conflict instead of silently overwriting a terminal state.
import asyncio, time, uuid
from dataclasses import dataclass, field
TERMINAL = {"TASK_STATE_COMPLETED", "TASK_STATE_FAILED",
"TASK_STATE_CANCELED", "TASK_STATE_REJECTED"}
INTERRUPTED = {"TASK_STATE_INPUT_REQUIRED", "TASK_STATE_AUTH_REQUIRED"}
# Server policy: which moves this agent allows. The spec defines the states and which
# are terminal or interrupted; the exact edges are an implementation decision.
ALLOWED = {
"TASK_STATE_SUBMITTED": {"TASK_STATE_WORKING", "TASK_STATE_REJECTED",
"TASK_STATE_CANCELED", "TASK_STATE_FAILED"},
"TASK_STATE_WORKING": {"TASK_STATE_INPUT_REQUIRED", "TASK_STATE_AUTH_REQUIRED",
"TASK_STATE_COMPLETED", "TASK_STATE_FAILED",
"TASK_STATE_CANCELED"},
"TASK_STATE_INPUT_REQUIRED": {"TASK_STATE_WORKING", "TASK_STATE_CANCELED",
"TASK_STATE_FAILED"},
"TASK_STATE_AUTH_REQUIRED": {"TASK_STATE_WORKING", "TASK_STATE_CANCELED",
"TASK_STATE_FAILED"},
}
@dataclass
class TaskRow:
id: str
context_id: str
state: str = "TASK_STATE_SUBMITTED"
version: int = 0
history: list = field(default_factory=list)
artifacts: dict = field(default_factory=dict)
events: list = field(default_factory=list) # (seq, event) per task
class Conflict(Exception): pass
class TaskManager:
def __init__(self):
self.tasks, self.by_message = {}, {}
self.queue = asyncio.Queue()
self.subs = {} # task id -> set of asyncio.Queue
self.changed = {} # task id -> asyncio.Condition
def _publish(self, row, event):
row.events.append((len(row.events), event))
for q in self.subs.get(row.id, ()):
q.put_nowait(event)
cond = self.changed.setdefault(row.id, asyncio.Condition())
async def wake():
async with cond:
cond.notify_all()
asyncio.create_task(wake())
def transition(self, task_id, expected_version, new_state, status_msg=None):
row = self.tasks[task_id]
if row.version != expected_version:
raise Conflict(f"{task_id}: version {row.version}, expected {expected_version}")
if new_state not in ALLOWED.get(row.state, set()):
raise Conflict(f"{task_id}: {row.state} -> {new_state} not allowed")
row.state, row.version = new_state, row.version + 1
self._publish(row, {"statusUpdate": {
"taskId": row.id, "contextId": row.context_id,
"status": {"state": new_state, "message": status_msg,
"timestamp": time.strftime("%Y-%m-%dT%H:%M:%SZ", time.gmtime())}}})
return row.version
async def send_message(self, msg, return_immediately=False):
if msg["messageId"] in self.by_message: # client retry: same answer
return self.tasks[self.by_message[msg["messageId"]]]
if msg.get("taskId"): # continuing an existing task
row = self.tasks.get(msg["taskId"])
if row is None:
raise KeyError("TaskNotFoundError")
if row.state in TERMINAL:
raise Conflict("terminal tasks accept no further messages; start a new "
"task in the same contextId")
if row.state not in INTERRUPTED:
raise Conflict("task is not waiting for input")
row.history.append(msg)
self.by_message[msg["messageId"]] = row.id
self.transition(row.id, row.version, "TASK_STATE_WORKING")
else: # a new unit of work
row = TaskRow(id=f"tsk_{uuid.uuid4().hex[:12]}",
context_id=msg.get("contextId") or f"ctx_{uuid.uuid4().hex[:12]}",
history=[msg])
self.tasks[row.id] = row
self.by_message[msg["messageId"]] = row.id
await self.queue.put(row.id)
if not return_immediately: # block until it settles
cond = self.changed.setdefault(row.id, asyncio.Condition())
async with cond:
await cond.wait_for(lambda: row.state in TERMINAL | INTERRUPTED)
return row
def emit_artifact(self, task_id, artifact_id, parts, append=False, last_chunk=False):
row = self.tasks[task_id]
art = row.artifacts.setdefault(artifact_id, {"artifactId": artifact_id, "parts": []})
art["parts"] = art["parts"] + parts if append else list(parts)
self._publish(row, {"artifactUpdate": {
"taskId": row.id, "contextId": row.context_id,
"artifact": {"artifactId": artifact_id, "parts": parts},
"append": append, "lastChunk": last_chunk}})The executor loop claims ids and drives the transitions. Note what it does on errors: it turns them into TASK_STATE_FAILED with a short, safe message, never a stack trace, because status messages are visible to the calling agent and, through it, possibly to a model prompt.
async def worker(tm: TaskManager, agent):
while True:
task_id = await tm.queue.get()
row = tm.tasks[task_id]
v = row.version # version seen at claim time
try:
if row.state == "TASK_STATE_SUBMITTED":
v = tm.transition(task_id, v, "TASK_STATE_WORKING")
outcome = await agent.run(row, tm) # may emit artifacts as it goes
if outcome.needs_input:
v = tm.transition(task_id, v, "TASK_STATE_INPUT_REQUIRED",
status_msg=outcome.question)
else:
v = tm.transition(task_id, v, "TASK_STATE_COMPLETED")
except Conflict:
pass # a cancel or a newer claim moved the task first
except Exception as exc:
try:
tm.transition(task_id, v, "TASK_STATE_FAILED", status_msg={
"messageId": str(uuid.uuid4()), "role": "ROLE_AGENT",
"parts": [{"text": f"internal error: {type(exc).__name__}"}]})
except Conflict:
pass # someone else already settled it
Worked example: delegating a reconciliation
A finance assistant agent delegates invoice reconciliation to a ledger agent. Follow one task through the system.
- The assistant sends the message shown above with
returnImmediatelyset to true. The handler has not seen the messageId, creates tasktsk_01JB7QinTASK_STATE_SUBMITTED, enqueues it and returns the task within a few milliseconds. - The assistant opens
SubscribeToTaskfor that id. Our server sends the current Task snapshot first, then live events. - An executor claims the task and moves it to
TASK_STATE_WORKINGat version 1. It streams a partial artifactmatches.jsonin three chunks: the first withappendfalse, the next two withappendtrue, the last withlastChunktrue. - Forty invoices have no ledger entry. The executor moves the task to
TASK_STATE_INPUT_REQUIREDwith a status message asking whether to treat unmatched invoices under 50 EUR as rounding. The stream reports it; the assistant asks its user. - The network drops while the user thinks. Nothing is lost: the task sits in the store. On reconnect the assistant calls
GetTaskto confirm the state, then sends a new message withtaskIdset totsk_01JB7Qand a fresh messageId. - The handler appends the message to history, moves the task back to working and re-enqueues it. The executor finishes, emits the final artifact and sets
TASK_STATE_COMPLETED. - Next week the user asks for October. That is a new task with the same
contextId, so the ledger agent can reuse the rounding answer.
Client patterns: block, stream, poll or push
A client has four ways to learn the outcome, and robust clients combine two of them.
| Pattern | How | Use when | Watch out for |
|---|---|---|---|
| Blocking | SendMessage with returnImmediately false | Work finishes well inside your HTTP timeout | Proxies cutting long requests; a cut does not cancel the task |
| Streaming | SendStreamingMessage, or SubscribeToTask later | You want progress and partial artifacts | Reconnects; always re-read with GetTask after one |
| Polling | GetTask with backoff | Simple clients, firewalled callers | Load from tight loops; cap and jitter the interval |
| Push | A task push-notification config with a webhook | Tasks that run for hours | Verifying the sender; see push notifications |
The safe default for anything over a few seconds is: send non-blocking, stream for progress, and fall back to GetTask polling with exponential backoff whenever the stream breaks. If you set historyLength on GetTask, the server returns at most that many recent messages, which keeps poll responses small on long multi-turn tasks. The details of the event stream are in A2A streaming.
Failure modes
- Duplicate tasks from retries. A client times out on a blocking send and retries. Without a dedupe map keyed by messageId the server starts the work twice. Keep the map for at least the client's retry horizon, and persist it with the task row so a restart does not forget it. Idempotency in A2A covers the general pattern.
- Lost wakeups. A worker crashes after claiming a task but before changing its state, and the task sits in working forever. Use a lease: the claim records an expiry, a reaper re-enqueues tasks whose lease lapsed, and the agent logic must therefore be safe to run twice.
- Terminal-state overwrite. A late completion overwrites a cancel, or a retry flips a failed task back to working. The compare-and-set and the transition table prevent both; an unconditional
UPDATEdoes not. - Event gaps on reconnect. A subscriber reconnects and misses the chunk that carried
lastChunk, so it waits forever. Treat the stream as a hint: on every reconnect, read the task and its artifacts from the store. - Unbounded growth. History and event logs grow with every turn. Set retention per state, for example terminal tasks kept 30 days and their event logs 24 hours, and document how long ids stay valid.
- Cancel that does not stop work. Marking a row canceled does nothing to a running model call. The executor must check for cancellation between steps; task cancellation covers the cooperative unwind.
Design trade-offs
Queue in front of executors. A queue gives back-pressure and lets you scale executors separately, but adds a hop and a second place where work can be lost. Using the task store itself as the queue, with a claim column and a lease, removes that second place at the cost of polling the store.
Multi-tenancy. The 1.0 request messages carry a tenant field. Partition store, queues and rate limits by tenant from the start; retrofitting it means migrating every task id.
What to do next
- Decide which of your agent's skills return a Message and which create a Task, and write the rule down in the agent's documentation.
- Model the task row with id, contextId, state, version, lease expiry and a messageId index, and enforce transitions with compare-and-set.
- Write your transition table as code and unit-test every illegal edge, especially transitions out of terminal states.
- Put a queue or a claimable column between the handler and the executors, with a lease and a reaper.
- Append every status and artifact change to a per-task event log, and drive streaming and push from it.
- In your client, send non-blocking, subscribe for progress and fall back to GetTask with jittered backoff after any disconnect.
- Set retention for terminal tasks and event logs, and load-test with retries enabled to confirm no duplicate work.