Skip to main content

Deduplicator

Struct Deduplicator 

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

In-process request deduplicator that coalesces identical concurrent requests and caches recently completed results.

§How it works

  1. The caller derives a stable cache key (see dedup_key) from the prompt and session.
  2. Deduplicator::check_and_register atomically checks the shared state map and returns one of three outcomes:
  3. On success the worker calls Deduplicator::complete; on failure it calls Deduplicator::fail, which removes the entry.

§Thread safety

Deduplicator is Clone + Send + Sync. All clones share the same underlying Arc<DashMap>. A background task cleans up expired entries every 60 seconds; it stops when the last Deduplicator clone is dropped.

§Examples

use std::time::Duration;
use tokio_prompt_orchestrator::enhanced::{Deduplicator, DeduplicationResult};

let dedup = Deduplicator::new(Duration::from_secs(300));
let key = "dedup:g:abc123";

match dedup.check_and_register(key).await {
    DeduplicationResult::New(token) => {
        let result = "hello".to_string();
        dedup.complete(token, result).await;
    }
    DeduplicationResult::InProgress => {
        let _ = dedup.wait_for_result(key).await;
    }
    DeduplicationResult::Cached(result) => println!("{result}"),
}

Implementations§

Source§

impl Deduplicator

Source

pub fn new(cache_duration: Duration) -> Self

Create a new Deduplicator with the given cache TTL.

A background cleanup task is spawned immediately. It wakes up every 60 seconds to evict expired entries and exits when the last Deduplicator clone is dropped.

§Arguments
  • cache_duration — How long a completed result remains cached before being treated as a fresh request. Common choices: 5 minutes for interactive use, 1 hour for batch/idempotent workloads.
§Examples
use std::time::Duration;
use tokio_prompt_orchestrator::enhanced::Deduplicator;

let dedup = Deduplicator::new(Duration::from_secs(300));
§Panics

Spawns a background cleanup task, so it must be called from within a Tokio runtime.

Source

pub fn signal_shutdown(&self)

Signal the background cleanup task to stop without waiting for it.

Sets the shutdown AtomicBool to true. The cleanup loop exits on its next wake-up. Call shutdown if you need to await completion.

Source

pub async fn shutdown(&self)

Gracefully shut down the background cleanup task.

Sets the shutdown flag so the cleanup loop exits on its next wake-up, then waits for the task to finish. Safe to call multiple times.

Source

pub async fn check_and_register(&self, key: &str) -> DeduplicationResult

Atomically check whether a request is new, in-progress, or cached, and register it as in-progress if it is new.

Uses DashMap::entry() for a compare-and-insert that prevents multiple concurrent callers from each receiving DeduplicationResult::New for the same key — only one will win the race.

§Arguments
  • key — Stable cache key; derive one with dedup_key.
§Returns

A DeduplicationResult indicating whether the caller should process the request, wait for another task, or reuse a cached result.

§Examples
use std::time::Duration;
use tokio_prompt_orchestrator::enhanced::{Deduplicator, DeduplicationResult};

let dedup = Deduplicator::new(Duration::from_secs(60));
if let DeduplicationResult::New(token) = dedup.check_and_register("key").await {
    dedup.complete(token, "result".to_string()).await;
}
Source

pub async fn wait_for_result(&self, key: &str) -> Option<String>

Wait for an in-progress request to complete and return its result.

Subscribes to the internal broadcast channel for the given key. If the request has already completed by the time this is called, the cached result is returned immediately without waiting.

§Arguments
§Returns

Some(result) when the pending request completes, or None if the key is not tracked (e.g. the worker called Deduplicator::fail).

Source

pub async fn complete(&self, token: DeduplicationToken, result: String)

Mark a request as successfully completed and cache its result.

Notifies all tasks currently blocked in Deduplicator::wait_for_result for the same key. The result is retained in the cache for cache_duration so subsequent callers receive DeduplicationResult::Cached.

§Arguments
Source

pub async fn fail(&self, token: DeduplicationToken)

Mark a request as failed and remove it from tracking.

After this call, the next Deduplicator::check_and_register for the same key will receive DeduplicationResult::New so the request can be retried. Any tasks waiting in Deduplicator::wait_for_result will receive None on their next recv() after the sender is dropped.

§Arguments
Source

pub fn stats(&self) -> DeduplicationStats

Return a snapshot of current deduplication statistics.

The counts are computed by iterating the internal map in O(n). Use sparingly on hot paths; prefer Prometheus counters for high-frequency monitoring.

Source

pub fn clear(&self)

Clear all cached results

Source

pub fn with_semantic(self, threshold: f32) -> Self

Enable semantic (embedding-based) deduplication.

When enabled, check_and_register_with_embedding compares new embeddings against all stored embeddings using cosine similarity. Any stored embedding with similarity ≥ threshold is treated as a cache hit.

§Arguments
  • threshold — Cosine similarity score in [0.0, 1.0]. 1.0 requires exact vector match (default); 0.95 catches near-paraphrases.
§Example
use std::time::Duration;
use tokio_prompt_orchestrator::enhanced::Deduplicator;

let dedup = Deduplicator::new(Duration::from_secs(300))
    .with_semantic(0.95);
Source

pub async fn check_and_register_with_embedding( &self, key: &str, embedding: Option<Vec<f32>>, ) -> DeduplicationResult

Like check_and_register but also performs a semantic similarity scan against previously registered embeddings.

If embedding is Some and semantic deduplication is enabled (threshold < 1.0), all stored embeddings are scanned. The first match whose cosine similarity meets the threshold is returned as DeduplicationResult::Cached with an empty string (the caller should use wait_for_result with the matched key to obtain the actual cached value).

Falls back to exact-key lookup when embedding is None or the threshold equals 1.0.

§Arguments
  • key — Exact cache key for this request.
  • embedding — Optional dense vector embedding of the prompt.

Trait Implementations§

Source§

impl Clone for Deduplicator

Source§

fn clone(&self) -> Deduplicator

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 Drop for Deduplicator

Source§

fn drop(&mut self)

Executes the destructor for this type. Read more
Source§

fn pin_drop(self: Pin<&mut Self>)

🔬This is a nightly-only experimental API. (pin_ergonomics)
Execute the destructor for this type, but different to Drop::drop, it requires self to be pinned. 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
§

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

§

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> From<T> for T

Source§

fn from(t: T) -> T

Returns the argument unchanged.

§

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

§

fn from_ref(input: &T) -> T

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

impl<T> FutureExt for T

§

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

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

fn with_current_context(self) -> WithContext<Self>

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

impl<T> Instrument for T

§

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

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

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> IntoRequest<T> for T

Source§

fn into_request(self) -> Request<T>

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

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

§

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
§

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
§

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 = 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.
§

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

§

fn vzip(self) -> V

§

impl<T> WithSubscriber for T

§

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
§

fn with_current_subscriber(self) -> WithDispatch<Self>

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