Skip to content

ant_ai.a2a.executor

A2AExecutor

Bases: AgentExecutor

A2AExecutor is responsible for processing an A2A request.

It handles the execution the workflow and propages the updates generated in it.

Source code in src/ant_ai/a2a/executor.py
 38
 39
 40
 41
 42
 43
 44
 45
 46
 47
 48
 49
 50
 51
 52
 53
 54
 55
 56
 57
 58
 59
 60
 61
 62
 63
 64
 65
 66
 67
 68
 69
 70
 71
 72
 73
 74
 75
 76
 77
 78
 79
 80
 81
 82
 83
 84
 85
 86
 87
 88
 89
 90
 91
 92
 93
 94
 95
 96
 97
 98
 99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
class A2AExecutor(AgentExecutor):
    """
    A2AExecutor is responsible for processing an A2A request.

    It handles the execution the workflow and propages the updates generated in it.
    """

    def __init__(
        self,
        agent: Agent,
        workflow: Workflow,
        *,
        stream_artifacts: bool = True,
        context_class: type[InvocationContext] = InvocationContext,
    ):
        """Initialize the A2AExecutor. A2AExecutor is a subclass of AgentExecutor. The AgentExecutor is the a2a-sdk class
        that is responsible for processing the request made to the agent.

        Args:
            agent: Agent that will be used to execute the workflow.
            workflow: Workflow that will be executed.
            stream_artifacts: Whether to translate ContentDeltaEvent into A2A
                artifact-update chunks. See `HVEventToA2A`.
            context_class: The `InvocationContext` (sub)class built for each
                request from the message metadata via `from_metadata`.
        """
        self.workflow: Workflow = workflow
        self.agent: Agent = agent
        self.context_class: type[InvocationContext] = context_class
        self._translator: HVEventToA2A = HVEventToA2A(stream_artifacts=stream_artifacts)
        self._a2a_to_hv: A2AToHVEvent = A2AToHVEvent()

    async def execute(
        self,
        context: RequestContext,
        event_queue: EventQueue,
    ) -> None:
        """Entry point of the A2A. This is what processes each request to the agent via a2a.

        Args:
            context: The object containing all info about the request.
            event_queue: The event queue to which events will be enqueued.
        """

        if not context.message:
            raise Exception("No message provided")

        task: Task = context.current_task or new_task_from_user_message(context.message)
        if not context.current_task:
            await event_queue.enqueue_event(task)

        updater = TaskUpdater(event_queue, task.id, task.context_id)
        with (
            obs.attach_propagation_context(context.metadata),
            obs.bind(
                agent_name=self.agent.name, task_id=task.id, context_id=task.context_id
            ),
        ):
            await obs.event("a2a.execute", task_id=task.id, context_id=task.context_id)
            try:
                await self._execute(context, updater, task)
            except asyncio.CancelledError:
                await obs.event(
                    "a2a.cancelled", task_id=task.id, context_id=task.context_id
                )
                raise
            except A2AError as exc:
                await obs.exception("a2a.error", exc)
                raise
            except ContextWindowExceededError as exc:
                # The message is the library's own, not the provider's, so it
                # can cross to the caller; the provider's text stays in the log.
                await obs.exception("a2a.error", exc)
                raise InternalError(str(exc)) from exc
            except Exception as exc:
                await obs.exception("a2a.error", exc)
                raise InternalError() from exc

    async def cancel(self, context: RequestContext, event_queue: EventQueue) -> None:
        """Cancel a running task by publishing the canceled status.

        Stopping the coroutine is the SDK's job: `DefaultRequestHandler.on_cancel_task`
        calls this and then cancels the producer task itself. Raising here (as this
        method used to) aborted the handler before it got that far, so `tasks/cancel`
        could not stop anything.

        What the handler needs from us is the canceled status: `consume_all` refuses
        any other final state, and a canceled status update is a final event, which
        is what closes the queue the cancelled producer never gets to close.
        """
        task: Task | None = context.current_task
        if task is None:
            raise Exception("No task to cancel.")

        await obs.event("a2a.cancel", task_id=task.id, context_id=task.context_id)
        await TaskUpdater(event_queue, task.id, task.context_id).cancel()

    async def _execute(
        self,
        context: RequestContext,
        updater: TaskUpdater,
        task: Task,
    ) -> None:
        ctx: InvocationContext = self.context_class.from_metadata(
            session_id=task.context_id, metadata=context.metadata
        )

        history: list[Message] = self._build_history(
            context.related_tasks, context.get_user_input()
        )

        await obs.event("a2a.history", history_messages=len(history))

        state: State = self.workflow.create_state(messages=history)
        token: Token[str] = current_session_id.set(ctx.session_id)

        await obs.event("a2a.workflow.start")
        try:
            async for event in self.workflow.stream(
                agent=self.agent, ctx=ctx, state=state
            ):
                if not isinstance(event, ContentDeltaEvent):
                    await obs.event(
                        "a2a.workflow.event",
                        workflow_event=getattr(event, "kind", type(event).__name__),
                        node=getattr(event.origin, "node", "-"),
                        step=getattr(event.origin, "run_step", "-"),
                    )
                if isinstance(event, CompletedEvent):
                    await persist_compression_checkpoint(state, updater)
                await self.process_event(event, updater)
        finally:
            current_session_id.reset(token)

        await obs.event("a2a.workflow.end")

    async def process_event(self, event: Event, updater: TaskUpdater) -> None:
        await self._translator.apply(event=event, updater=updater)

    def _build_history(
        self, related_tasks: list[Task], user_input: str
    ) -> list[Message]:
        checkpoint, cp_idx = find_compression_checkpoint(related_tasks)
        if checkpoint is not None:
            post_a2a = [
                m
                for r_task in related_tasks[cp_idx:]
                for m in r_task.history or []
                if not is_checkpoint_message(m)
            ]
            history = list(checkpoint) + self._convert_history(post_a2a)
        else:
            history = self._convert_history(
                [m for r_task in related_tasks for m in r_task.history or []]
            )
        history.append(Message(role="user", content=user_input))
        return history

    def _convert_history(self, a2a_history: list[A2AMessage]) -> list[Message]:
        return self._answer_dangling_tool_calls(
            [
                m
                for msg in a2a_history
                if (m := self._a2a_to_hv.to_history_message(msg)) is not None
            ]
        )

    @staticmethod
    def _answer_dangling_tool_calls(history: list[AnyMessage]) -> list[AnyMessage]:
        """Give every unanswered tool call a result, so the transcript replays.

        A task cancelled (or crashed) while its tools ran ends with the
        assistant's `tool_calls` turn and nothing after it: the run cannot
        emit results once the terminal status is on the queue. Chat-completion
        APIs reject a transcript where such a turn is followed by anything but
        its tool messages, which made the next turn in the same context fail.
        The synthetic results say what happened and are flagged as errors, so
        the model knows the work was not done.
        """
        repaired: list[AnyMessage] = []
        pending: dict[str, str] = {}  # call id -> tool name, unanswered so far

        def flush() -> None:
            repaired.extend(
                ToolCallResultMessage(
                    tool_call_id=call_id,
                    name=name,
                    content=_UNANSWERED_TOOL_CALL,
                    is_error=True,
                )
                for call_id, name in pending.items()
            )
            pending.clear()

        for m in history:
            if isinstance(m, ToolCallResultMessage):
                pending.pop(m.tool_call_id, None)
            elif pending:
                flush()
            if isinstance(m, ToolCallMessage):
                pending.update((tc.id, tc.function.name) for tc in m.tool_calls)
            repaired.append(m)
        flush()
        return repaired

__init__

__init__(
    agent: Agent,
    workflow: Workflow,
    *,
    stream_artifacts: bool = True,
    context_class: type[
        InvocationContext
    ] = InvocationContext,
)

Initialize the A2AExecutor. A2AExecutor is a subclass of AgentExecutor. The AgentExecutor is the a2a-sdk class that is responsible for processing the request made to the agent.

Parameters:

Name Type Description Default
agent Agent

Agent that will be used to execute the workflow.

required
workflow Workflow

Workflow that will be executed.

required
stream_artifacts bool

Whether to translate ContentDeltaEvent into A2A artifact-update chunks. See HVEventToA2A.

True
context_class type[InvocationContext]

The InvocationContext (sub)class built for each request from the message metadata via from_metadata.

InvocationContext
Source code in src/ant_ai/a2a/executor.py
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
def __init__(
    self,
    agent: Agent,
    workflow: Workflow,
    *,
    stream_artifacts: bool = True,
    context_class: type[InvocationContext] = InvocationContext,
):
    """Initialize the A2AExecutor. A2AExecutor is a subclass of AgentExecutor. The AgentExecutor is the a2a-sdk class
    that is responsible for processing the request made to the agent.

    Args:
        agent: Agent that will be used to execute the workflow.
        workflow: Workflow that will be executed.
        stream_artifacts: Whether to translate ContentDeltaEvent into A2A
            artifact-update chunks. See `HVEventToA2A`.
        context_class: The `InvocationContext` (sub)class built for each
            request from the message metadata via `from_metadata`.
    """
    self.workflow: Workflow = workflow
    self.agent: Agent = agent
    self.context_class: type[InvocationContext] = context_class
    self._translator: HVEventToA2A = HVEventToA2A(stream_artifacts=stream_artifacts)
    self._a2a_to_hv: A2AToHVEvent = A2AToHVEvent()

execute async

execute(
    context: RequestContext, event_queue: EventQueue
) -> None

Entry point of the A2A. This is what processes each request to the agent via a2a.

Parameters:

Name Type Description Default
context RequestContext

The object containing all info about the request.

required
event_queue EventQueue

The event queue to which events will be enqueued.

required
Source code in src/ant_ai/a2a/executor.py
 70
 71
 72
 73
 74
 75
 76
 77
 78
 79
 80
 81
 82
 83
 84
 85
 86
 87
 88
 89
 90
 91
 92
 93
 94
 95
 96
 97
 98
 99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
async def execute(
    self,
    context: RequestContext,
    event_queue: EventQueue,
) -> None:
    """Entry point of the A2A. This is what processes each request to the agent via a2a.

    Args:
        context: The object containing all info about the request.
        event_queue: The event queue to which events will be enqueued.
    """

    if not context.message:
        raise Exception("No message provided")

    task: Task = context.current_task or new_task_from_user_message(context.message)
    if not context.current_task:
        await event_queue.enqueue_event(task)

    updater = TaskUpdater(event_queue, task.id, task.context_id)
    with (
        obs.attach_propagation_context(context.metadata),
        obs.bind(
            agent_name=self.agent.name, task_id=task.id, context_id=task.context_id
        ),
    ):
        await obs.event("a2a.execute", task_id=task.id, context_id=task.context_id)
        try:
            await self._execute(context, updater, task)
        except asyncio.CancelledError:
            await obs.event(
                "a2a.cancelled", task_id=task.id, context_id=task.context_id
            )
            raise
        except A2AError as exc:
            await obs.exception("a2a.error", exc)
            raise
        except ContextWindowExceededError as exc:
            # The message is the library's own, not the provider's, so it
            # can cross to the caller; the provider's text stays in the log.
            await obs.exception("a2a.error", exc)
            raise InternalError(str(exc)) from exc
        except Exception as exc:
            await obs.exception("a2a.error", exc)
            raise InternalError() from exc

cancel async

cancel(
    context: RequestContext, event_queue: EventQueue
) -> None

Cancel a running task by publishing the canceled status.

Stopping the coroutine is the SDK's job: DefaultRequestHandler.on_cancel_task calls this and then cancels the producer task itself. Raising here (as this method used to) aborted the handler before it got that far, so tasks/cancel could not stop anything.

What the handler needs from us is the canceled status: consume_all refuses any other final state, and a canceled status update is a final event, which is what closes the queue the cancelled producer never gets to close.

Source code in src/ant_ai/a2a/executor.py
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
async def cancel(self, context: RequestContext, event_queue: EventQueue) -> None:
    """Cancel a running task by publishing the canceled status.

    Stopping the coroutine is the SDK's job: `DefaultRequestHandler.on_cancel_task`
    calls this and then cancels the producer task itself. Raising here (as this
    method used to) aborted the handler before it got that far, so `tasks/cancel`
    could not stop anything.

    What the handler needs from us is the canceled status: `consume_all` refuses
    any other final state, and a canceled status update is a final event, which
    is what closes the queue the cancelled producer never gets to close.
    """
    task: Task | None = context.current_task
    if task is None:
        raise Exception("No task to cancel.")

    await obs.event("a2a.cancel", task_id=task.id, context_id=task.context_id)
    await TaskUpdater(event_queue, task.id, task.context_id).cancel()