a2a.server.cluster.database_task_store module¶
- class a2a.server.cluster.database_task_store.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.