Skip to content

[FEATURE]: A2A tasks run on the runtime #783

Description

@justintime4tea

Summary

AuraAgentExecutor runs the agent inside the execution stream, keeps its own cancel map (TaskCancelState and its guard), passes its own token to RunOptions::cancelled_by, and folds StreamItems into task status by hand. BusBridgedExecutor then republishes every execution event over the session-store bus with its own sequence numbers so another instance can relay a subscribe. That is a second registry and a second journal, both A2A-private, and the task's state is a fold over stream items rather than a reading of the run — which is how #737 happens: a cancel racing the fold records Canceled over a task the fold already completed.

Move the executor onto the runtime. A task starts a run; task status is a projection of the run's lifecycle; subscribe is a subscription; and the run's status, not a fold, decides whether a cancel still applies.

Goals

Data structures

Proposed. The projection is a function over the subscription; the ids it needs are the task's.

pub struct A2aTaskIds { pub task_id: String, pub context_id: String }

/// Projects one run event onto the A2A execution stream. `None` for events A2A
/// has no frame for.
pub fn project_a2a(event: &SessionEvent, ids: &A2aTaskIds) -> Option<StreamResponse>;   // None for another run's events

// Lifecycle(Started)                       -> StatusUpdate { state: Working }
// Agent(TextDelta | Completed)             -> ArtifactUpdate (one part per chunk, as today)
// Agent(ApprovalPending) | Lifecycle(Parked) -> StatusUpdate { state: InputRequired }  (the #209 wiring)
// Lifecycle(Finished)                      -> StatusUpdate { state: Completed, final: true }
// Lifecycle(Cancelled)                     -> StatusUpdate { state: Canceled,  final: true }
// Lifecycle(Failed)                        -> StatusUpdate { state: Failed,    final: true }
// Lifecycle(ObserverAttached | ...)        -> None

Additional Context

Task-store writes stay where they are. The store is the durable record of a task; the runtime is the live record of a run; the executor keeps them in step, as it does today. The task record gains the run's RunId and SessionId (in its metadata, as the upstream Task type allows), so an instance serving tasks/{id}:subscribe for a task it did not start can find the session and the cursor without a fold over the execution stream.

A worker's tool calls within an orchestrated run are attributed to the run, not the worker, until #732 gives MCP a stream identity. The projection carries whatever AgentContext the event has; it does not repair it.

a2a-server's own ActiveExecution — a broadcast per task that its request handler subscribes clients to, with a "subscription fell behind active execution" error on lag — is upstream and stays. The executor's stream feeds it; the runtime's journal is what feeds the executor. The UnboundedSender<StreamResponse> the executor writes into is the one channel end this seam keeps, because the upstream executor interface is a stream; it is one hop at the seam, not a reach into the run. Because the task holds the claim for its whole life, ClaimsExhausted never fires for an A2A run before it is terminal, which is why #784 has no a2a liveness entry.

Searched Issues

  • No similar issues found

Code of Conduct

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

Metadata

Metadata

Assignees

No one assigned

    Labels

    No labels
    No labels

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions