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: TaskEventStream

TaskEventStream backed by polling the shared task_events table.

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_events table 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.