Skip to main content

MemoryTaskStore

Struct MemoryTaskStore 

Source
pub struct MemoryTaskStore { /* private fields */ }
Expand description

In-memory TaskStore backed by a HashMap.

This is the default store. Suitable for single-instance deployments. For horizontal scaling, use an external store that shares state across instances. A lazy background worker signals task expiry at its exact TTL deadline and physically removes expired records at the configured cleanup interval. It uses a standard thread rather than assuming construction happens inside a Tokio runtime, and holds only weak task state while sleeping.

Completion wakeups for wait_for_completion use a per-task tokio::sync::Notify, which is an implementation detail of this store.

Implementations§

Source§

impl MemoryTaskStore

Source

pub fn new() -> Self

Create a task store with the default lifecycle and retention policy.

That is a five-minute TTL, one-minute cleanup interval, and the finite TaskRetentionLimits::default byte/count bounds.

Source

pub fn with_config(config: MemoryTaskStoreConfig) -> Self

Create a task store with an explicit lifecycle policy and the default finite retention limits.

Source

pub fn with_retention_limits(retention_limits: TaskRetentionLimits) -> Self

Create a task store with the default lifecycle policy and explicit record and encoded-payload limits.

Source

pub fn with_config_and_retention( config: MemoryTaskStoreConfig, retention_limits: TaskRetentionLimits, ) -> Self

Create a task store with explicit lifecycle and retention policies.

Source

pub fn cleanup_expired(&self) -> usize

Remove expired tasks immediately.

Returns the number removed. Not part of the TaskStore trait; external backends typically expire entries natively (e.g. Redis TTL).

The configured worker already calls this retirement path periodically; applications may call it to reclaim memory sooner. Calling it is an optimization, not a correctness requirement: expiry has already cancelled work, woken waiters, and made the task read as absent.

§Example
use tower_mcp::async_task::{MemoryTaskStore, TaskStore};

let store = MemoryTaskStore::new();
// A one millisecond retention window, and no terminal state: the
// clock runs from creation, so the task retires while still working.
let (id, _cancel) = store
    .create_task("deploy", serde_json::json!({}), Some(1), None)
    .await
    .unwrap();
tokio::time::sleep(std::time::Duration::from_millis(10)).await;

// Already invisible, before anything has been reclaimed.
assert!(store.get_task(&id).await.unwrap().is_none());
assert!(store.list_tasks(None).await.unwrap().is_empty());

// Cleanup only frees the memory the entry was still holding.
assert_eq!(store.cleanup_expired(), 1);
assert_eq!(store.cleanup_expired(), 0);
Source

pub fn usage(&self) -> TaskStoreUsage

Return content-free count and encoded-byte gauges.

The snapshot is taken under the same lock as task mutations, so its count and both byte totals always describe one committed store state.

Trait Implementations§

Source§

impl Clone for MemoryTaskStore

Source§

fn clone(&self) -> MemoryTaskStore

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 MemoryTaskStore

Source§

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

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

impl Default for MemoryTaskStore

Source§

fn default() -> Self

Returns the “default value” for a type. Read more
Source§

impl TaskStore for MemoryTaskStore

Source§

fn task_presence<'life0, 'life1, 'async_trait>( &'life0 self, task_id: &'life1 str, ) -> Pin<Box<dyn Future<Output = Result<TaskPresence>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

This store keeps expired records until cleanup_expired runs, so it can tell an owner that a task expired rather than that it never existed (#1249). One read resolves both, so expiry cannot change between deciding presence and reading the owner.

Source§

fn create_task<'life0, 'life1, 'async_trait>( &'life0 self, tool_name: &'life1 str, arguments: Value, ttl: Option<u64>, owner: TaskOwner, ) -> Pin<Box<dyn Future<Output = Result<(String, CancellationToken)>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

Create and store a new task owned by owner. Read more
Source§

fn get_task<'life0, 'life1, 'async_trait>( &'life0 self, task_id: &'life1 str, ) -> Pin<Box<dyn Future<Output = Result<Option<TaskObject>>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

Get task object by ID. Returns None if unknown.
Source§

fn set_task_meta<'life0, 'life1, 'async_trait>( &'life0 self, task_id: &'life1 str, meta: Value, ) -> Pin<Box<dyn Future<Output = Result<bool>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

Persist protocol _meta for a task. Read more
Source§

fn discard_task<'life0, 'life1, 'async_trait>( &'life0 self, task_id: &'life1 str, ) -> Pin<Box<dyn Future<Output = Result<bool>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

Remove a task that could not finish initialization. Read more
Source§

fn task_owner<'life0, 'life1, 'async_trait>( &'life0 self, task_id: &'life1 str, ) -> Pin<Box<dyn Future<Output = Result<Option<TaskOwner>>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

Read a task’s owner. Read more
Source§

fn get_task_result<'life0, 'life1, 'async_trait>( &'life0 self, task_id: &'life1 str, ) -> Pin<Box<dyn Future<Output = Result<Option<TaskSnapshot>>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

Get a task’s full snapshot (task object, result, error) by ID.
Source§

fn wait_for_completion<'life0, 'life1, 'async_trait>( &'life0 self, task_id: &'life1 str, ) -> Pin<Box<dyn Future<Output = Result<Option<TaskSnapshot>>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

Wait for a task to reach a terminal state, then return its snapshot. Read more
Source§

fn list_tasks<'life0, 'async_trait>( &'life0 self, status_filter: Option<TaskStatus>, ) -> Pin<Box<dyn Future<Output = Result<Vec<TaskObject>>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait,

List all tasks, optionally filtered by status.
Source§

fn require_input<'life0, 'life1, 'life2, 'async_trait>( &'life0 self, task_id: &'life1 str, requests: InputRequests, message: Option<&'life2 str>, ) -> Pin<Box<dyn Future<Output = Result<bool>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait, 'life2: 'async_trait,

Mark a task as requiring input, recording the requests to be answered. Read more
Source§

fn outstanding_input_requests<'life0, 'life1, 'async_trait>( &'life0 self, task_id: &'life1 str, ) -> Pin<Box<dyn Future<Output = Result<Option<InputRequests>>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

Read the requests a task is currently waiting on. Read more
Source§

fn apply_input_responses<'life0, 'life1, 'async_trait>( &'life0 self, task_id: &'life1 str, responses: InputResponses, ) -> Pin<Box<dyn Future<Output = Result<Option<AppliedInputResponses>>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

Apply tasks/update.inputResponses to a task. Read more
Source§

fn set_status<'life0, 'life1, 'life2, 'async_trait>( &'life0 self, task_id: &'life1 str, status: TaskStatus, message: Option<&'life2 str>, ) -> Pin<Box<dyn Future<Output = Result<bool>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait, 'life2: 'async_trait,

Set a task’s non-terminal status and message. Read more
Source§

fn resume_context<'life0, 'life1, 'async_trait>( &'life0 self, task_id: &'life1 str, ) -> Pin<Box<dyn Future<Output = Result<Option<TaskResumeContext>>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

Everything needed to re-invoke a task’s handler after its input requests were answered. Read more
Source§

fn set_ttl<'life0, 'life1, 'async_trait>( &'life0 self, task_id: &'life1 str, ttl_ms: u64, ) -> Pin<Box<dyn Future<Output = Result<bool>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

Update a task’s time-to-live, measured from creation. Read more
Source§

fn complete_task<'life0, 'life1, 'async_trait>( &'life0 self, task_id: &'life1 str, result: CallToolResult, ) -> Pin<Box<dyn Future<Output = Result<bool>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

Mark a task as completed with a result. Read more
Source§

fn fail_task<'life0, 'life1, 'async_trait>( &'life0 self, task_id: &'life1 str, error: JsonRpcError, ) -> Pin<Box<dyn Future<Output = Result<bool>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

Mark a task as failed with a structured execution error. Read more
Source§

fn cancel_task<'life0, 'life1, 'life2, 'async_trait>( &'life0 self, task_id: &'life1 str, reason: Option<&'life2 str>, ) -> Pin<Box<dyn Future<Output = Result<Option<TaskObject>>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait, 'life2: 'async_trait,

Cancel a task. Read more
Source§

fn input_responses<'life0, 'life1, 'async_trait>( &'life0 self, task_id: &'life1 str, ) -> Pin<Box<dyn Future<Output = Result<Option<InputResponses>>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

Every answer accumulated for a task so far, keyed as issued. 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<ST, DT> CastableFrom<ST, Initialized, Initialized> for DT
where ST: ?Sized, DT: ?Sized,

Source§

impl<ST, DT> CastableFrom<ST, Uninit, Uninit> for DT
where ST: ?Sized, DT: ?Sized,

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> DynClone for T
where T: Clone,

Source§

fn __clone_box(&self, _: Private) -> *mut ()

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<A, B, T> HttpServerConnExec<A, B> for T
where B: Body,

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> PolicyExt for T
where T: ?Sized,

Source§

fn and<P, B, E>(self, other: P) -> And<T, P>
where T: Sized + Policy<B, E>, P: Policy<B, E>,

Create a new Policy that returns Action::Follow only if self and other return Action::Follow. Read more
Source§

fn or<P, B, E>(self, other: P) -> Or<T, P>
where T: Sized + Policy<B, E>, P: Policy<B, E>,

Create a new Policy that returns Action::Follow if either self or other returns Action::Follow. Read more
Source§

impl<T> Read<Exclusive, BecauseExclusive> for T
where T: ?Sized,

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

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