a2a.server.cluster package

Submodules

Module contents

exception a2a.server.cluster.ConcurrentTaskModificationError(task_id: str)

Bases: Exception

Raised 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: 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.

class a2a.server.cluster.LegacyTaskStoreAdapter(store: TaskStore)

Bases: VersionedTaskStore

Runs 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.

property store: TaskStore

The wrapped task store.

class a2a.server.cluster.StoredTask(task: Task, version: TaskVersion)

Bases: object

A task together with the version it was read at.

task: Task
version: TaskVersion
class a2a.server.cluster.TaskEventStream

Bases: ABC

Delivers 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: object

A 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: VersionedTaskStore

VersionedTaskStore backed by SQLAlchemy.

property as_task_store: TaskStore

The underlying non-versioned TaskStore.

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 to task_events in the same transaction for cross-replica replay.

class a2a.server.cluster.VersionedEvent(event: Message | Task | TaskStatusUpdateEvent | TaskArtifactUpdateEvent, version: TaskVersion)

Bases: object

An event together with the task version produced by applying it.

event: Message | Task | TaskStatusUpdateEvent | TaskArtifactUpdateEvent
version: TaskVersion
class a2a.server.cluster.VersionedTaskStore

Bases: ABC

A 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.