a2a.server.cluster package¶
Submodules¶
Module contents¶
- exception a2a.server.cluster.ConcurrentTaskModificationError(task_id: str)¶
Bases:
ExceptionRaised by VersionedTaskStore.save when prev_version is stale.
- class a2a.server.cluster.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.
- class a2a.server.cluster.LegacyTaskStoreAdapter(store: TaskStore)¶
Bases:
VersionedTaskStoreRuns an unversioned TaskStore under the VersionedTaskStore interface.
- async delete(task_id: str, context: ServerCallContext) None¶
Deletes a task via the wrapped store.
- async get(task_id: str, context: ServerCallContext) StoredTask | None¶
Gets from the wrapped store, pairing the result with MISSING.
- async list(params: ListTasksRequest, context: ServerCallContext) ListTasksResponse¶
Lists tasks via the wrapped store.
- async save(task: Task, *, event: Message | Task | TaskStatusUpdateEvent | TaskArtifactUpdateEvent | None, prev: Task | None, prev_version: TaskVersion, context: ServerCallContext) TaskVersion¶
Saves via the wrapped store; ignores version args, returns MISSING.
- class a2a.server.cluster.StoredTask(task: Task, version: TaskVersion)¶
Bases:
objectA task together with the version it was read at.
- version: TaskVersion¶
- class a2a.server.cluster.TaskEventStream¶
Bases:
ABCDelivers task events across replicas.
- abstractmethod async destroy(task_id: str) None¶
Releases resources for a task that has reached a terminal state.
- abstractmethod async publish(task_id: str, event: VersionedEvent) None¶
Publishes one event for task_id to all replicas.
- abstractmethod subscribe(task_id: str, *, after: TaskVersion) AsyncGenerator[VersionedEvent, None]¶
Yields events for task_id newer than after.
- class a2a.server.cluster.TaskVersion(value: int | str)¶
Bases:
objectA version marker a VersionedTaskStore assigns to a stored Task.
Prevents concurrent state re-writes. The wrapped value is the store’s choice - a counter, a commit timestamp, etc. Callers do not read it or do arithmetic on it; they only pass it back to save and order two versions with is_after.
- MISSING: ClassVar[TaskVersion] = TaskVersion.MISSING¶
- is_after(other: TaskVersion) bool¶
Whether self is a strictly later version than other.
- property is_missing: bool¶
Whether this token means “versioning is not tracked”.
- class a2a.server.cluster.VersionedDatabaseTaskStore(engine: ~sqlalchemy.ext.asyncio.engine.AsyncEngine, create_table: bool = True, table_name: str = 'tasks', owner_resolver: ~collections.abc.Callable[[~a2a.server.context.ServerCallContext], str] = <function resolve_user_scope>, core_to_model_conversion: ~collections.abc.Callable[[~a2a_pb2.Task, str], ~a2a.server.models.TaskModel] | None = None, model_to_core_conversion: ~collections.abc.Callable[[~a2a.server.models.TaskModel], ~a2a_pb2.Task] | None = None, event_table_name: str = 'task_events', version_table_name: str = 'task_versions', max_attempts: int = 5, retry_delay_s: float = 0.02)¶
Bases:
VersionedTaskStoreVersionedTaskStore backed by SQLAlchemy.
- async delete(task_id: str, context: ServerCallContext) None¶
Deletes a task and its version row via the underlying store.
- async get(task_id: str, context: ServerCallContext) StoredTask | None¶
Returns the task with its stored version, or None if absent.
Retries transient contention (e.g. a lock held while another writer commits) rather than failing the read.
- async initialize() None¶
Initializes the database schema (task, version and event tables).
- async list(params: ListTasksRequest, context: ServerCallContext) ListTasksResponse¶
Lists tasks via the underlying store.
- async save(task: Task, *, event: Message | Task | TaskStatusUpdateEvent | TaskArtifactUpdateEvent | None = None, prev: Task | None = None, prev_version: TaskVersion, context: ServerCallContext) TaskVersion¶
Persists task with a compare-and-swap on the version side table.
The version lives in
task_versions, not on the tasks row, and counts writes to the task: the first write sets 1 and every later write adds 1. Updates CAS on it; cancel overwrites a non-terminal task. When event is provided it is appended totask_eventsin the same transaction for cross-replica replay.
- class a2a.server.cluster.VersionedEvent(event: Message | Task | TaskStatusUpdateEvent | TaskArtifactUpdateEvent, version: TaskVersion)¶
Bases:
objectAn event together with the task version produced by applying it.
- event: Message | Task | TaskStatusUpdateEvent | TaskArtifactUpdateEvent¶
- version: TaskVersion¶
- class a2a.server.cluster.VersionedTaskStore¶
Bases:
ABCA TaskStore variant with snapshot versioning to prevent concurrent re-writes.
- abstractmethod async delete(task_id: str, context: ServerCallContext) None¶
Deletes a task from the store by ID.
- abstractmethod async get(task_id: str, context: ServerCallContext) StoredTask | None¶
Retrieves a task with its version, or None if it does not exist.
- abstractmethod async list(params: ListTasksRequest, context: ServerCallContext) ListTasksResponse¶
Retrieves a list of tasks from the store.
- abstractmethod async save(task: Task, *, event: Message | Task | TaskStatusUpdateEvent | TaskArtifactUpdateEvent | None, prev: Task | None, prev_version: TaskVersion, context: ServerCallContext) TaskVersion¶
Persists task and returns its new version.
- Parameters:
task – The task state to persist.
event – The event that produced this state, or None for a direct write.
prev – The task as previously read, for implementations that diff.
prev_version – The version task was derived from. Implementations MUST raise ConcurrentTaskModificationError if the currently stored version differs. TaskVersion.MISSING marks a first write. A write moving task to CANCELED overwrites a non-terminal stored task without a version check, and raises ConcurrentTaskModificationError if the stored task is already terminal or absent.
context – The server call context (used to resolve the owner).
- Returns:
The new TaskVersion for the persisted task.
- Raises:
ConcurrentTaskModificationError – If prev_version is stale.