pub struct PostgresTaskStore { /* private fields */ }Expand description
PostgreSQL-backed TaskStore.
Stores tasks as JSONB blobs in a tasks table. Suitable for multi-node
production deployments that need shared persistence and horizontal scaling.
§Schema
The store auto-creates the following table on first use:
CREATE TABLE IF NOT EXISTS tasks (
id TEXT PRIMARY KEY,
context_id TEXT NOT NULL,
state TEXT NOT NULL,
data JSONB NOT NULL,
created_at TIMESTAMPTZ NOT NULL DEFAULT now(),
updated_at TIMESTAMPTZ NOT NULL DEFAULT now()
);list() returns tasks most-recently-updated first (spec §3.1.4), ordered by
(updated_at DESC, id DESC) with a composite row-value cursor. The cursor
carries updated_at as a UTC-normalized microsecond string, so pagination
is stable regardless of the connection’s session time zone.
Implementations§
Source§impl PostgresTaskStore
impl PostgresTaskStore
Sourcepub const fn with_max_page_size(self, max: u32) -> Self
pub const fn with_max_page_size(self, max: u32) -> Self
Caps the page size list returns, however large a page is asked for.
Defaults to DEFAULT_MAX_PAGE_SIZE, which explains why this store
needs its own knob rather than reading TaskStoreConfig.
Sourcepub async fn new(url: &str) -> Result<Self, Error>
pub async fn new(url: &str) -> Result<Self, Error>
Opens a PostgreSQL connection pool and initializes the schema.
§Errors
Returns an error if the database cannot be opened or the schema migration fails.
Sourcepub async fn with_migrations(url: &str) -> Result<Self, Error>
pub async fn with_migrations(url: &str) -> Result<Self, Error>
Opens a PostgreSQL database with automatic schema migration.
Runs all pending migrations before returning the store. This is the recommended constructor for production deployments because it ensures the schema is always up to date without duplicating DDL statements.
§Errors
Returns an error if the database cannot be opened or any migration fails.
Sourcepub async fn from_pool(pool: PgPool) -> Result<Self, Error>
pub async fn from_pool(pool: PgPool) -> Result<Self, Error>
Creates a store from an existing connection pool.
§Errors
Returns an error if the schema migration fails.
Sourcepub async fn purge_expired(
&self,
policy: &RetentionPolicy,
) -> A2aResult<PurgeReport>
pub async fn purge_expired( &self, policy: &RetentionPolicy, ) -> A2aResult<PurgeReport>
Deletes terminal tasks that have outlived policy.
Nothing calls this for you. A persistent store keeps every task until
an operator says otherwise — see retention
for why that is the default and why the in-memory store does the
opposite — so this is the hook for whatever already schedules work: a
cron entry, a Kubernetes CronJob, a tokio interval in your own
binary.
Only Completed, Failed, Canceled and Rejected tasks are
eligible. A task still Working, or parked in InputRequired waiting
on a human, is never deleted however old it is.
Safe to run from several replicas at once: each batch is a single
DELETE whose subquery picks the rows, so two sweeps racing delete
disjoint sets rather than colliding.
§Errors
Returns an error if a delete fails. A sweep that fails partway has still committed its earlier batches; the counts in the returned report are lost in that case, but the deletions are not undone and the next sweep simply continues.
Trait Implementations§
Source§impl Clone for PostgresTaskStore
impl Clone for PostgresTaskStore
Source§fn clone(&self) -> PostgresTaskStore
fn clone(&self) -> PostgresTaskStore
1.0.0 (const: unstable) · Source§fn clone_from(&mut self, source: &Self)
fn clone_from(&mut self, source: &Self)
source. Read moreSource§impl Debug for PostgresTaskStore
impl Debug for PostgresTaskStore
Source§impl TaskStore for PostgresTaskStore
impl TaskStore for PostgresTaskStore
Source§fn save_artifact_delta<'a>(
&'a self,
task: &'a Task,
delta: ArtifactDelta,
) -> Pin<Box<dyn Future<Output = A2aResult<()>> + Send + 'a>>
fn save_artifact_delta<'a>( &'a self, task: &'a Task, delta: ArtifactDelta, ) -> Pin<Box<dyn Future<Output = A2aResult<()>> + Send + 'a>>
Appends into the stored JSONB document instead of rewriting it.
save serializes the whole task in Rust and ships it as a bind
parameter, so a streaming agent re-sends every artifact it has already
persisted on every subsequent event. This sends only what changed and
lets PostgreSQL splice it in with jsonb_set.
Unlike the SQLite implementation, which needs one path expression per
appended part, jsonb’s || concatenates two arrays — so any number of
parts lands in a single statement with constant SQL text.
§What this does and does not remove
Removed: the Rust-side serde_json::to_value of the whole task and the
transfer of the whole document. Both scale with the stream so far.
Not removed: PostgreSQL still rewrites the row. An UPDATE writes a
new tuple version under MVCC, and a JSONB document past the TOAST
threshold is rewritten out of line, so the statement stays linear in
document size. Only a normalized artifacts table could avoid that, and
the measurement in benches/benches/backpressure.rs puts the per-event
round trip well above the document-size term — so that surgery would buy
the smaller half. Recorded here rather than left implied.
updated_at is deliberately untouched: it carries the status
timestamp that orders list (§3.1.4), and appending an artifact does
not change a task’s status. Both other stores behave the same way, and a
divergence here would be invisible until someone paginated.
Falls back to save when the delta cannot be applied exactly: no
artifacts on the task, an index out of range, a Pushed that does not
name the last position, fewer parts present than claimed, or a stored
row whose document has no matching array. A store that is quietly wrong
is worse than one that is slower.
Source§fn save<'a>(
&'a self,
task: &'a Task,
) -> Pin<Box<dyn Future<Output = A2aResult<()>> + Send + 'a>>
fn save<'a>( &'a self, task: &'a Task, ) -> Pin<Box<dyn Future<Output = A2aResult<()>> + Send + 'a>>
Source§fn get<'a>(
&'a self,
id: &'a TaskId,
) -> Pin<Box<dyn Future<Output = A2aResult<Option<Task>>> + Send + 'a>>
fn get<'a>( &'a self, id: &'a TaskId, ) -> Pin<Box<dyn Future<Output = A2aResult<Option<Task>>> + Send + 'a>>
None if not found. Read moreSource§fn list<'a>(
&'a self,
params: &'a ListTasksParams,
) -> Pin<Box<dyn Future<Output = A2aResult<TaskListResponse>> + Send + 'a>>
fn list<'a>( &'a self, params: &'a ListTasksParams, ) -> Pin<Box<dyn Future<Output = A2aResult<TaskListResponse>> + Send + 'a>>
Source§fn insert_if_absent<'a>(
&'a self,
task: &'a Task,
) -> Pin<Box<dyn Future<Output = A2aResult<bool>> + Send + 'a>>
fn insert_if_absent<'a>( &'a self, task: &'a Task, ) -> Pin<Box<dyn Future<Output = A2aResult<bool>> + Send + 'a>>
Auto Trait Implementations§
impl !RefUnwindSafe for PostgresTaskStore
impl !UnwindSafe for PostgresTaskStore
impl Freeze for PostgresTaskStore
impl Send for PostgresTaskStore
impl Sync for PostgresTaskStore
impl Unpin for PostgresTaskStore
impl UnsafeUnpin for PostgresTaskStore
Blanket Implementations§
Source§impl<T> BorrowMut<T> for Twhere
T: ?Sized,
impl<T> BorrowMut<T> for Twhere
T: ?Sized,
Source§fn borrow_mut(&mut self) -> &mut T
fn borrow_mut(&mut self) -> &mut T
Source§impl<T> CloneToUninit for Twhere
T: Clone,
impl<T> CloneToUninit for Twhere
T: Clone,
Source§impl<T> FutureExt for T
impl<T> FutureExt for T
Source§fn with_context(self, otel_cx: Context) -> WithContext<Self> ⓘ
fn with_context(self, otel_cx: Context) -> WithContext<Self> ⓘ
Source§fn with_current_context(self) -> WithContext<Self> ⓘ
fn with_current_context(self) -> WithContext<Self> ⓘ
Source§impl<T> Instrument for T
impl<T> Instrument for T
Source§fn instrument(self, span: Span) -> Instrumented<Self> ⓘ
fn instrument(self, span: Span) -> Instrumented<Self> ⓘ
Source§fn in_current_span(self) -> Instrumented<Self> ⓘ
fn in_current_span(self) -> Instrumented<Self> ⓘ
Source§impl<T> IntoEither for T
impl<T> IntoEither for T
Source§fn into_either(self, into_left: bool) -> Either<Self, Self> ⓘ
fn into_either(self, into_left: bool) -> Either<Self, Self> ⓘ
self into a Left variant of Either<Self, Self>
if into_left is true.
Converts self into a Right variant of Either<Self, Self>
otherwise. Read moreSource§fn into_either_with<F>(self, into_left: F) -> Either<Self, Self> ⓘ
fn into_either_with<F>(self, into_left: F) -> Either<Self, Self> ⓘ
self into a Left variant of Either<Self, Self>
if into_left(&self) returns true.
Converts self into a Right variant of Either<Self, Self>
otherwise. Read moreSource§impl<T> IntoRequest<T> for T
impl<T> IntoRequest<T> for T
Source§fn into_request(self) -> Request<T>
fn into_request(self) -> Request<T>
T in a tonic::Request