pub struct Deduplicator { /* private fields */ }Expand description
In-process request deduplicator that coalesces identical concurrent requests and caches recently completed results.
§How it works
- The caller derives a stable cache key (see
dedup_key) from the prompt and session. Deduplicator::check_and_registeratomically checks the shared state map and returns one of three outcomes:DeduplicationResult::New— the caller is the first to see this key; it receives aDeduplicationTokenand must process the request.DeduplicationResult::InProgress— another task is already working; the caller should callDeduplicator::wait_for_resultto block.DeduplicationResult::Cached— a prior result is still within thecache_durationTTL; the caller can return it immediately.
- On success the worker calls
Deduplicator::complete; on failure it callsDeduplicator::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
impl Deduplicator
Sourcepub fn new(cache_duration: Duration) -> Self
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.
Sourcepub fn signal_shutdown(&self)
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.
Sourcepub async fn shutdown(&self)
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.
Sourcepub async fn check_and_register(&self, key: &str) -> DeduplicationResult
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 withdedup_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;
}Sourcepub async fn wait_for_result(&self, key: &str) -> Option<String>
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
key— The same key passed toDeduplicator::check_and_register.
§Returns
Some(result) when the pending request completes, or None if the
key is not tracked (e.g. the worker called Deduplicator::fail).
Sourcepub async fn complete(&self, token: DeduplicationToken, result: String)
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
token— TheDeduplicationTokenreturned byDeduplicator::check_and_register.result— The serialised response to cache and broadcast.
Sourcepub async fn fail(&self, token: DeduplicationToken)
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
token— TheDeduplicationTokenreturned byDeduplicator::check_and_register.
Sourcepub fn stats(&self) -> DeduplicationStats
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.
Sourcepub fn with_semantic(self, threshold: f32) -> Self
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.0requires exact vector match (default);0.95catches 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);Sourcepub async fn check_and_register_with_embedding(
&self,
key: &str,
embedding: Option<Vec<f32>>,
) -> DeduplicationResult
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
impl Clone for Deduplicator
Source§fn clone(&self) -> Deduplicator
fn clone(&self) -> Deduplicator
1.0.0 (const: unstable) · Source§fn clone_from(&mut self, source: &Self)
fn clone_from(&mut self, source: &Self)
source. Read moreSource§impl Drop for Deduplicator
impl Drop for Deduplicator
Auto Trait Implementations§
impl !RefUnwindSafe for Deduplicator
impl !UnwindSafe for Deduplicator
impl Freeze for Deduplicator
impl Send for Deduplicator
impl Sync for Deduplicator
impl Unpin for Deduplicator
impl UnsafeUnpin for Deduplicator
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
impl<ST, DT> CastableFrom<ST, Initialized, Initialized> for DT
impl<ST, DT> CastableFrom<ST, Uninit, Uninit> for DT
Source§impl<T> CloneToUninit for Twhere
T: Clone,
impl<T> CloneToUninit for Twhere
T: Clone,
§impl<T> FutureExt for T
impl<T> FutureExt for T
§fn with_context(self, otel_cx: Context) -> WithContext<Self>
fn with_context(self, otel_cx: Context) -> WithContext<Self>
§fn with_current_context(self) -> WithContext<Self>
fn with_current_context(self) -> WithContext<Self>
§impl<T> Instrument for T
impl<T> Instrument for T
§fn instrument(self, span: Span) -> Instrumented<Self>
fn instrument(self, span: Span) -> Instrumented<Self>
§fn in_current_span(self) -> Instrumented<Self>
fn in_current_span(self) -> Instrumented<Self>
Source§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