Skip to main content

PostgresTaskStore

Struct PostgresTaskStore 

Source
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

Source

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.

Source

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.

Source

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.

Trait Implementations§

Source§

impl Clone for PostgresTaskStore

Source§

fn clone(&self) -> PostgresTaskStore

Returns a duplicate of the value. Read more
1.0.0 (const: unstable) · Source§

fn clone_from(&mut self, source: &Self)

Performs copy-assignment from source. Read more
Source§

impl Debug for PostgresTaskStore

Source§

fn fmt(&self, f: &mut Formatter<'_>) -> Result

Formats the value using the given formatter. Read more
Source§

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

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

Saves (creates or updates) a task. Read more
Source§

fn get<'a>( &'a self, id: &'a TaskId, ) -> Pin<Box<dyn Future<Output = A2aResult<Option<Task>>> + Send + 'a>>

Retrieves a task by its ID, returning None if not found. Read more
Source§

fn list<'a>( &'a self, params: &'a ListTasksParams, ) -> Pin<Box<dyn Future<Output = A2aResult<TaskListResponse>> + Send + 'a>>

Lists tasks matching the given filter parameters. Read more
Source§

fn insert_if_absent<'a>( &'a self, task: &'a Task, ) -> Pin<Box<dyn Future<Output = A2aResult<bool>> + Send + 'a>>

Atomically inserts a task only if no task with the same ID exists. Read more
Source§

fn delete<'a>( &'a self, id: &'a TaskId, ) -> Pin<Box<dyn Future<Output = A2aResult<()>> + Send + 'a>>

Deletes a task by its ID. Read more
Source§

fn count<'a>( &'a self, ) -> Pin<Box<dyn Future<Output = A2aResult<u64>> + Send + 'a>>

Returns the total number of tasks in the store. Read more

Auto Trait Implementations§

Blanket Implementations§

Source§

impl<T> Any for T
where T: 'static + ?Sized,

Source§

fn type_id(&self) -> TypeId

Gets the TypeId of self. Read more
Source§

impl<T> Borrow<T> for T
where T: ?Sized,

Source§

fn borrow(&self) -> &T

Immutably borrows from an owned value. Read more
Source§

impl<T> BorrowMut<T> for T
where T: ?Sized,

Source§

fn borrow_mut(&mut self) -> &mut T

Mutably borrows from an owned value. Read more
Source§

impl<T> CloneToUninit for T
where T: Clone,

Source§

unsafe fn clone_to_uninit(&self, dest: *mut u8)

🔬This is a nightly-only experimental API. (clone_to_uninit)
Performs copy-assignment from self to dest. Read more
Source§

impl<T> From<T> for T

Source§

fn from(t: T) -> T

Returns the argument unchanged.

Source§

impl<T> FromRef<T> for T
where T: Clone,

Source§

fn from_ref(input: &T) -> T

Converts to this type from a reference to the input type.
Source§

impl<T> FutureExt for T

Source§

fn with_context(self, otel_cx: Context) -> WithContext<Self>

Attaches the provided Context to this type, returning a WithContext wrapper. Read more
Source§

fn with_current_context(self) -> WithContext<Self>

Attaches the current Context to this type, returning a WithContext wrapper. Read more
Source§

impl<T> Instrument for T

Source§

fn instrument(self, span: Span) -> Instrumented<Self>

Instruments this type with the provided Span, returning an Instrumented wrapper. Read more
Source§

fn in_current_span(self) -> Instrumented<Self>

Instruments this type with the current Span, returning an Instrumented wrapper. Read more
Source§

impl<T, U> Into<U> for T
where U: From<T>,

Source§

fn into(self) -> U

Calls U::from(self).

That is, this conversion is whatever the implementation of From<T> for U chooses to do.

Source§

impl<T> IntoEither for T

Source§

fn into_either(self, into_left: bool) -> Either<Self, Self>

Converts 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 more
Source§

fn into_either_with<F>(self, into_left: F) -> Either<Self, Self>
where F: FnOnce(&Self) -> bool,

Converts 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 more
Source§

impl<T> IntoRequest<T> for T

Source§

fn into_request(self) -> Request<T>

Wrap the input message T in a tonic::Request
Source§

impl<L> LayerExt<L> for L

Source§

fn named_layer<S>(&self, service: S) -> Layered<<L as Layer<S>>::Service, S>
where L: Layer<S>,

Applies the layer to a service and wraps it in Layered.
Source§

impl<T> Same for T

Source§

type Output = T

Should always be Self
Source§

impl<T> ToOwned for T
where T: Clone,

Source§

type Owned = T

The resulting type after obtaining ownership.
Source§

fn to_owned(&self) -> T

Creates owned data from borrowed data, usually by cloning. Read more
Source§

fn clone_into(&self, target: &mut T)

Uses borrowed data to replace owned data, usually by cloning. Read more
Source§

impl<T, U> TryFrom<U> for T
where U: Into<T>,

Source§

type Error = Infallible

The type returned in the event of a conversion error.
Source§

fn try_from(value: U) -> Result<T, <T as TryFrom<U>>::Error>

Performs the conversion.
Source§

impl<T, U> TryInto<U> for T
where U: TryFrom<T>,

Source§

type Error = <U as TryFrom<T>>::Error

The type returned in the event of a conversion error.
Source§

fn try_into(self) -> Result<U, <U as TryFrom<T>>::Error>

Performs the conversion.
Source§

impl<V, T> VZip<V> for T
where V: MultiLane<T>,

Source§

fn vzip(self) -> V

Source§

impl<T> WithSubscriber for T

Source§

fn with_subscriber<S>(self, subscriber: S) -> WithDispatch<Self>
where S: Into<Dispatch>,

Attaches the provided Subscriber to this type, returning a WithDispatch wrapper. Read more
Source§

fn with_current_subscriber(self) -> WithDispatch<Self>

Attaches the current default Subscriber to this type, returning a WithDispatch wrapper. Read more