Skip to content

ant_ai.a2a.translator

HVEventToA2A

Translator that converts internal HV Events to A2A updates by applying the appropriate handler based on the event class. Each handler is responsible for taking an Event and using the TaskUpdater to propagate the corresponding update to A2A.

Instantiated once per A2AExecutor and shared across every concurrent task it serves, so this class must stay stateless: is_first/stream_id uniqueness comes from the emitting LLMStep (a fresh id per generation), never from per-artifact bookkeeping kept here.

Source code in src/ant_ai/a2a/translator.py
 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
class HVEventToA2A:
    """
    Translator that converts internal HV Events to A2A updates by applying the appropriate handler based on the event class. Each handler is responsible for taking an Event and using the TaskUpdater to propagate the corresponding update to A2A.

    Instantiated once per A2AExecutor and shared across every concurrent task
    it serves, so this class must stay stateless: `is_first`/`stream_id`
    uniqueness comes from the emitting `LLMStep` (a fresh id per generation),
    never from per-artifact bookkeeping kept here.
    """

    def __init__(self, *, stream_artifacts: bool = True) -> None:
        """
        Initializes the translator and registers handlers. Adopting a single-dispatch like approach for translating Events to A2A updates, where handlers are registered via a decorator and stored in a mapping of event class to handler method.

        Args:
            stream_artifacts: Whether to translate ContentDeltaEvent into A2A
                TaskArtifactUpdateEvent chunks via `add_artifact`. When False,
                deltas are dropped and only the terminal whole-event message
                is sent. The terminal message is always sent either way, so
                this is purely additive for peers that read artifacts.
        """
        self._stream_artifacts = stream_artifacts
        self._handlers: dict[type[Event], Handler] = {}
        self._register_handlers()

    def _register_handlers(self) -> None:
        """
        Scan instance methods and register decorated handlers.
        """
        for attr_name in dir(self):
            method = getattr(self, attr_name)
            types = getattr(method, "_event_types", None)
            if types:
                for t in types:
                    self._handlers[t] = method

    async def apply(self, event: Event, updater: TaskUpdater) -> None:
        """Applies the appropriate handler for the given Event based on its class, using the TaskUpdater to propagate updates to A2A. This method serves as the main entry point for translating Events to A2A updates, abstracting away the specific handling logic into separate methods for each event class.

        Args:
            event: The internal HV Event to be translated.
            updater: The A2A TaskUpdater instance used to propagate the translated event.

        Raises:
            ValueError: If no handler is registered for the class of the given event.
        """
        event_handler: Handler | None = self._handlers.get(type(event))
        if not event_handler:
            raise ValueError(
                f"No handler registered for event type: {type(event).__name__}"
            )

        await event_handler(event, updater)

    @handler(StartEvent)
    async def _start(self, event: Event, updater: TaskUpdater) -> None:
        await updater.start_work()

    @handler(UpdateEvent, MaxStepsReachedEvent)
    async def _update(self, event: Event, updater: TaskUpdater) -> None:
        metadata: dict[str, Any] = A2AMetadata(event=event).model_dump()
        await updater.update_status(
            state=TaskState.TASK_STATE_WORKING,
            metadata=metadata,
        )

    @handler(
        ToolCallingEvent,
        ToolResultEvent,
        FinalAnswerEvent,
        ReasoningEvent,
    )
    async def _agent_message(self, event: Event, updater: TaskUpdater) -> None:
        metadata: dict[str, Any] = A2AMetadata(event=event).model_dump()
        msg = updater.new_agent_message(parts=[Part(text=event.content)])
        msg.metadata.update(metadata)

        await updater.update_status(
            state=TaskState.TASK_STATE_WORKING,
            message=msg,
            metadata=metadata,
        )

        stream_id = getattr(event, "stream_id", None)
        if not self._stream_artifacts or not stream_id:
            return

        artifact_ids: list[str] = []
        if isinstance(event, ReasoningEvent):
            artifact_ids.append(f"{stream_id}:reasoning")
        else:
            if event.content:
                artifact_ids.append(f"{stream_id}:content")
            tool_calls = getattr(event, "tool_calls", None) or []
            artifact_ids += [f"{stream_id}:tool:{i}" for i in range(len(tool_calls))]

        for artifact_id in artifact_ids:
            await updater.add_artifact(
                parts=[], artifact_id=artifact_id, append=True, last_chunk=True
            )

    @handler(ContentDeltaEvent)
    async def _content_delta(
        self, event: ContentDeltaEvent, updater: TaskUpdater
    ) -> None:
        if not self._stream_artifacts:
            return
        artifact_id = (
            f"{event.stream_id}:tool:{event.tool_call_index}"
            if event.tool_call_index is not None
            else f"{event.stream_id}:{event.target_kind}"
        )
        metadata: dict[str, Any] = A2AMetadata(event=event).model_dump()
        await updater.add_artifact(
            parts=[Part(text=event.delta)],
            artifact_id=artifact_id,
            append=not event.is_first,
            metadata=metadata,
        )

    @handler(ClarificationNeededEvent)
    async def _input_required(self, event: Event, updater: TaskUpdater) -> None:
        metadata: dict[str, Any] = A2AMetadata(event=event).model_dump()
        msg = updater.new_agent_message(parts=[Part(text=event.content)])
        msg.metadata.update(metadata)
        await updater.update_status(
            state=TaskState.TASK_STATE_INPUT_REQUIRED,
            message=msg,
            metadata=metadata,
        )

    @handler(CompletedEvent)
    async def _completed(self, event: Event, updater: TaskUpdater) -> None:
        await updater.complete()

__init__

__init__(*, stream_artifacts: bool = True) -> None

Initializes the translator and registers handlers. Adopting a single-dispatch like approach for translating Events to A2A updates, where handlers are registered via a decorator and stored in a mapping of event class to handler method.

Parameters:

Name Type Description Default
stream_artifacts bool

Whether to translate ContentDeltaEvent into A2A TaskArtifactUpdateEvent chunks via add_artifact. When False, deltas are dropped and only the terminal whole-event message is sent. The terminal message is always sent either way, so this is purely additive for peers that read artifacts.

True
Source code in src/ant_ai/a2a/translator.py
69
70
71
72
73
74
75
76
77
78
79
80
81
82
def __init__(self, *, stream_artifacts: bool = True) -> None:
    """
    Initializes the translator and registers handlers. Adopting a single-dispatch like approach for translating Events to A2A updates, where handlers are registered via a decorator and stored in a mapping of event class to handler method.

    Args:
        stream_artifacts: Whether to translate ContentDeltaEvent into A2A
            TaskArtifactUpdateEvent chunks via `add_artifact`. When False,
            deltas are dropped and only the terminal whole-event message
            is sent. The terminal message is always sent either way, so
            this is purely additive for peers that read artifacts.
    """
    self._stream_artifacts = stream_artifacts
    self._handlers: dict[type[Event], Handler] = {}
    self._register_handlers()

apply async

apply(event: Event, updater: TaskUpdater) -> None

Applies the appropriate handler for the given Event based on its class, using the TaskUpdater to propagate updates to A2A. This method serves as the main entry point for translating Events to A2A updates, abstracting away the specific handling logic into separate methods for each event class.

Parameters:

Name Type Description Default
event Event

The internal HV Event to be translated.

required
updater TaskUpdater

The A2A TaskUpdater instance used to propagate the translated event.

required

Raises:

Type Description
ValueError

If no handler is registered for the class of the given event.

Source code in src/ant_ai/a2a/translator.py
 95
 96
 97
 98
 99
100
101
102
103
104
105
106
107
108
109
110
111
async def apply(self, event: Event, updater: TaskUpdater) -> None:
    """Applies the appropriate handler for the given Event based on its class, using the TaskUpdater to propagate updates to A2A. This method serves as the main entry point for translating Events to A2A updates, abstracting away the specific handling logic into separate methods for each event class.

    Args:
        event: The internal HV Event to be translated.
        updater: The A2A TaskUpdater instance used to propagate the translated event.

    Raises:
        ValueError: If no handler is registered for the class of the given event.
    """
    event_handler: Handler | None = self._handlers.get(type(event))
    if not event_handler:
        raise ValueError(
            f"No handler registered for event type: {type(event).__name__}"
        )

    await event_handler(event, updater)

A2AToHVEvent

Translator that converts A2A messages and events into internal HV Events. Uses singledispatchmethod to define translation logic for different input types, allowing for flexible handling of various A2A message and event formats.

Source code in src/ant_ai/a2a/translator.py
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
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
class A2AToHVEvent:
    """
    Translator that converts A2A messages and events into internal HV Events. Uses singledispatchmethod to define translation logic for different input types, allowing for flexible handling of various A2A message and event formats.
    """

    @singledispatchmethod
    def translate(self, raw: Any) -> Event | None:
        return None

    @translate.register
    def _(self, raw: A2AMessage) -> Event | None:
        if not raw.metadata:
            return None
        md: dict[str, Any] = _json_format.MessageToDict(raw.metadata)
        event = md.get("event")
        if not event:
            return None
        event["task_id"] = raw.task_id
        event["session_id"] = raw.context_id
        return _any_event_adapter.validate_python(event)

    @translate.register
    def _(self, raw: TaskStatusUpdateEvent) -> Event | None:
        if not raw.metadata:
            return None
        md: dict[str, Any] = _json_format.MessageToDict(raw.metadata)
        event = md.get("event")
        if not event:
            return None
        event["task_id"] = raw.task_id
        event["session_id"] = raw.context_id
        return _any_event_adapter.validate_python(event)

    @translate.register
    def _(self, raw: TaskArtifactUpdateEvent) -> Event | None:
        artifact = raw.artifact
        md: dict[str, Any] = (
            _json_format.MessageToDict(artifact.metadata) if artifact.metadata else {}
        )
        event = md.get("event")
        if event:
            event["task_id"] = raw.task_id
            event["session_id"] = raw.context_id
            return _any_event_adapter.validate_python(event)

        text = "".join(
            p.text for p in artifact.parts if p.WhichOneof("content") == "text"
        )
        if not text:
            return None
        # No structured `event` metadata key, e.g. a spec-compliant peer that
        # isn't this codebase — best-effort reconstruction as a content delta.
        return ContentDeltaEvent(
            delta=text,
            stream_id=artifact.artifact_id,
            task_id=raw.task_id,
            session_id=raw.context_id,
        )

    @translate.register
    def _(self, raw: Task) -> Event | None:
        if not raw.metadata:
            return None
        md: dict[str, Any] = _json_format.MessageToDict(raw.metadata)
        event = md.get("event")
        if not event:
            return None
        event["task_id"] = raw.id
        event["session_id"] = raw.context_id
        return _any_event_adapter.validate_python(event)

    def to_history_message(self, raw: A2AMessage) -> AnyMessage | None:
        """Convert an A2A history message to the appropriate internal message type.

        Uses embedded event metadata when available to reconstruct structured
        messages (ToolCallMessage, ToolCallResultMessage). Falls back to plain
        text when no metadata is present. Returns None for non-conversation
        events (ReasoningEvent, UpdateEvent, etc.) that carry no LLM context.
        """
        from a2a.helpers import get_message_text
        from a2a.types import Role

        event = self.translate(raw)
        if event is None:
            text = get_message_text(raw)
            if raw.role != Role.ROLE_AGENT:
                return Message(role="user", content=text)
            return Message(role="assistant", content=text) if text else None
        if isinstance(event, ToolCallingEvent):
            return ToolCallMessage(tool_calls=list(event.tool_calls))
        if isinstance(event, ToolResultEvent):
            return ToolCallResultMessage(
                name=event.name,
                tool_call_id=event.tool_call_id,
                content=event.content,
                is_error=event.is_error,
            )
        if isinstance(event, FinalAnswerEvent):
            return Message(role="assistant", content=event.content)
        return None

to_history_message

to_history_message(raw: Message) -> AnyMessage | None

Convert an A2A history message to the appropriate internal message type.

Uses embedded event metadata when available to reconstruct structured messages (ToolCallMessage, ToolCallResultMessage). Falls back to plain text when no metadata is present. Returns None for non-conversation events (ReasoningEvent, UpdateEvent, etc.) that carry no LLM context.

Source code in src/ant_ai/a2a/translator.py
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
def to_history_message(self, raw: A2AMessage) -> AnyMessage | None:
    """Convert an A2A history message to the appropriate internal message type.

    Uses embedded event metadata when available to reconstruct structured
    messages (ToolCallMessage, ToolCallResultMessage). Falls back to plain
    text when no metadata is present. Returns None for non-conversation
    events (ReasoningEvent, UpdateEvent, etc.) that carry no LLM context.
    """
    from a2a.helpers import get_message_text
    from a2a.types import Role

    event = self.translate(raw)
    if event is None:
        text = get_message_text(raw)
        if raw.role != Role.ROLE_AGENT:
            return Message(role="user", content=text)
        return Message(role="assistant", content=text) if text else None
    if isinstance(event, ToolCallingEvent):
        return ToolCallMessage(tool_calls=list(event.tool_calls))
    if isinstance(event, ToolResultEvent):
        return ToolCallResultMessage(
            name=event.name,
            tool_call_id=event.tool_call_id,
            content=event.content,
            is_error=event.is_error,
        )
    if isinstance(event, FinalAnswerEvent):
        return Message(role="assistant", content=event.content)
    return None

handler

handler(*event_types: type[Event])

Decorator used to mark methods as handlers for specific Event classes.

Source code in src/ant_ai/a2a/translator.py
47
48
49
50
51
52
53
54
55
56
def handler(*event_types: type[Event]):
    """
    Decorator used to mark methods as handlers for specific Event classes.
    """

    def decorator(fn: UnboundHandler):
        fn._event_types = event_types  # ty:ignore[unresolved-attribute]
        return fn

    return decorator