a2a.server.cluster.database_event_stream module¶
- class a2a.server.cluster.database_event_stream.DatabaseTaskEventStream(engine: AsyncEngine, create_table: bool = True, table_name: str = 'task_events', poll_interval_s: float = 0.5)¶
Bases:
TaskEventStreamTaskEventStream backed by polling the shared
task_eventstable.- async destroy(task_id: str) None¶
No-op: the append-only log is retained; subscribers stop on their own.
- async initialize() None¶
Creates the
task_eventstable if requested.
- async publish(task_id: str, event: VersionedEvent) None¶
No-op: events are persisted transactionally by the task store.
- async subscribe(task_id: str, *, after: TaskVersion) AsyncGenerator[VersionedEvent, None]¶
Polls the log for events of task_id newer than after.
- a2a.server.cluster.database_event_stream.stream_response_to_event(response: StreamResponse) Message | Task | TaskStatusUpdateEvent | TaskArtifactUpdateEvent¶
Converts a StreamResponse proto back to an internal Event.