use dashmap::DashMap;
use std::collections::hash_map::DefaultHasher;
use std::hash::{Hash, Hasher};
use std::time::{Duration, Instant};
const DEDUP_TTL: Duration = Duration::from_millis(100);
const CALL_ID_TTL: Duration = Duration::from_secs(600);
pub struct DedupStore {
inner: DashMap<u64, Instant>,
ttl: Duration,
}
impl DedupStore {
pub fn new() -> Self {
Self {
inner: DashMap::new(),
ttl: DEDUP_TTL,
}
}
pub fn check_and_insert(
&self,
session_id: &str,
tool_name: &str,
tool_input: &serde_json::Value,
) -> bool {
let key = self.compute_hash(session_id, tool_name, tool_input);
self.claim(key, self.ttl)
}
pub fn check_and_insert_call(&self, session_id: &str, event_type: &str, call_id: &str) -> bool {
let key = self.compute_hash(
session_id,
event_type,
&serde_json::json!({ "tool_use_id": call_id }),
);
self.claim(key, CALL_ID_TTL)
}
fn claim(&self, key: u64, ttl: Duration) -> bool {
let now = Instant::now();
if let Some(deadline) = self.inner.get(&key) {
if now < *deadline.value() {
return true; }
}
self.inner.insert(key, now + ttl);
false
}
fn compute_hash(
&self,
session_id: &str,
tool_name: &str,
tool_input: &serde_json::Value,
) -> u64 {
let mut hasher = DefaultHasher::new();
session_id.hash(&mut hasher);
tool_name.hash(&mut hasher);
hash_value(tool_input, &mut hasher);
hasher.finish()
}
pub fn evict_expired(&self) {
let now = Instant::now();
self.inner.retain(|_, deadline| now < *deadline);
}
}
impl Default for DedupStore {
fn default() -> Self {
Self::new()
}
}
const PLUGIN_CLAIM_TTL: Duration = Duration::from_secs(10);
const PLUGIN_CLAIMS_PRUNE_EVERY: Duration = Duration::from_secs(1);
const PLUGIN_CLAIMS_PRUNE_ABOVE: usize = 4096;
type ClaimKey = (String, String);
#[derive(Default)]
pub struct PluginClaims {
inner: DashMap<ClaimKey, (u32, Instant)>,
last_prune: std::sync::Mutex<Option<Instant>>,
}
impl PluginClaims {
fn key(session_id: &str, identity: &serde_json::Value) -> ClaimKey {
let canonical = serde_json_canonicalizer::to_string(identity).unwrap_or_default();
(session_id.to_string(), canonical)
}
fn prune(&self, now: Instant) {
{
let mut last = self.last_prune.lock().unwrap_or_else(|e| e.into_inner());
let due = last.is_none_or(|at| now.duration_since(at) >= PLUGIN_CLAIMS_PRUNE_EVERY);
if !due && self.inner.len() <= PLUGIN_CLAIMS_PRUNE_ABOVE {
return;
}
*last = Some(now);
}
self.inner.retain(|_, (_, deadline)| now < *deadline);
}
pub fn insert(&self, session_id: &str, identity: &serde_json::Value) {
let now = Instant::now();
self.prune(now);
let mut entry = self
.inner
.entry(Self::key(session_id, identity))
.or_insert((0, now));
let (count, deadline) = entry.value_mut();
if now >= *deadline {
*count = 0;
}
*count += 1;
*deadline = now + PLUGIN_CLAIM_TTL;
}
pub fn consume(&self, session_id: &str, identity: &serde_json::Value) -> bool {
let now = Instant::now();
let key = Self::key(session_id, identity);
let Some(mut entry) = self.inner.get_mut(&key) else {
return false;
};
let (count, deadline) = entry.value_mut();
let hit = now < *deadline && *count > 0;
if hit {
*count -= 1;
}
let spent = *count == 0 || now >= *deadline;
drop(entry);
if spent {
self.inner
.remove_if(&key, |_, (count, deadline)| *count == 0 || now >= *deadline);
}
hit
}
#[cfg(test)]
fn len(&self) -> usize {
self.inner.len()
}
}
fn hash_value(value: &serde_json::Value, hasher: &mut impl Hasher) {
match value {
serde_json::Value::Null => 0u8.hash(hasher),
serde_json::Value::Bool(b) => {
1u8.hash(hasher);
b.hash(hasher);
}
serde_json::Value::Number(n) => {
2u8.hash(hasher);
format!("{n}").hash(hasher);
}
serde_json::Value::String(s) => {
3u8.hash(hasher);
s.hash(hasher);
}
serde_json::Value::Array(arr) => {
4u8.hash(hasher);
arr.len().hash(hasher);
for item in arr {
hash_value(item, hasher);
}
}
serde_json::Value::Object(map) => {
5u8.hash(hasher);
map.len().hash(hasher);
let mut keys: Vec<&String> = map.keys().collect();
keys.sort();
for key in keys {
key.hash(hasher);
hash_value(&map[key], hasher);
}
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use serde_json::json;
#[test]
fn test_first_event_is_not_duplicate() {
let store = DedupStore::new();
let input = json!({"command": "ls -la"});
let is_dup = store.check_and_insert("session-1", "bash", &input);
assert!(!is_dup, "first occurrence must not be a duplicate");
}
#[test]
fn test_same_event_within_ttl_is_duplicate() {
let store = DedupStore::new();
let input = json!({"command": "ls -la"});
let first = store.check_and_insert("session-1", "bash", &input);
let second = store.check_and_insert("session-1", "bash", &input);
assert!(!first, "first occurrence must not be a duplicate");
assert!(second, "immediate repeat must be detected as duplicate");
}
#[test]
fn test_same_event_after_ttl_is_not_duplicate() {
let store = DedupStore::new();
let input = json!({"command": "ls -la"});
store.check_and_insert("session-1", "bash", &input);
std::thread::sleep(Duration::from_millis(150));
let after_ttl = store.check_and_insert("session-1", "bash", &input);
assert!(!after_ttl, "event after TTL expiry must not be a duplicate");
}
#[test]
fn test_different_events_are_not_duplicates() {
let store = DedupStore::new();
let first = store.check_and_insert("session-1", "bash", &json!({"cmd": "ls"}));
let second =
store.check_and_insert("session-1", "read_file", &json!({"path": "/etc/hosts"}));
assert!(!first, "first event must not be a duplicate");
assert!(!second, "different event must not be a duplicate");
}
#[test]
fn test_different_sessions_same_tool_not_duplicate() {
let store = DedupStore::new();
let input = json!({"command": "ls"});
store.check_and_insert("session-A", "bash", &input);
let second = store.check_and_insert("session-B", "bash", &input);
assert!(
!second,
"same tool call in different session must not be a duplicate"
);
}
#[test]
fn test_evict_expired_removes_old_entries() {
let store = DedupStore::new();
let input = json!({"key": "value"});
store.check_and_insert("session-1", "bash", &input);
assert!(
store.check_and_insert("session-1", "bash", &input),
"entry must exist before eviction"
);
std::thread::sleep(Duration::from_millis(150));
store.evict_expired();
let after_evict = store.check_and_insert("session-1", "bash", &input);
assert!(
!after_evict,
"evicted entry must not be a duplicate on next check"
);
}
#[test]
fn a_call_id_is_claimed_past_the_content_window() {
let store = DedupStore::new();
assert!(!store.check_and_insert_call("s", "pre_tool_use", "call_1"));
std::thread::sleep(Duration::from_millis(150));
store.evict_expired();
assert!(store.check_and_insert_call("s", "pre_tool_use", "call_1"));
assert!(!store.check_and_insert_call("s", "pre_tool_use", "call_2"));
assert!(!store.check_and_insert_call("s", "post_tool_use", "call_1"));
assert!(!store.check_and_insert_call("other", "pre_tool_use", "call_1"));
}
#[test]
fn test_dedup_key_order_independent() {
let store = DedupStore::new();
let mut map_a = serde_json::Map::new();
map_a.insert("command".to_string(), json!("ls -la"));
map_a.insert("path".to_string(), json!("/tmp"));
let input_a = serde_json::Value::Object(map_a);
let mut map_b = serde_json::Map::new();
map_b.insert("path".to_string(), json!("/tmp"));
map_b.insert("command".to_string(), json!("ls -la"));
let input_b = serde_json::Value::Object(map_b);
let first = store.check_and_insert("session-1", "bash", &input_a);
let second = store.check_and_insert("session-1", "bash", &input_b);
assert!(!first, "first occurrence must not be a duplicate");
assert!(
second,
"same logical event with different key order must be detected as duplicate"
);
}
#[test]
fn test_hash_value_sorts_nested_keys() {
let input_a = json!({"z": {"b": 2, "a": 1}, "a": [{"y": 1, "x": 2}]});
let input_b = json!({"a": [{"x": 2, "y": 1}], "z": {"a": 1, "b": 2}});
let mut hasher_a = DefaultHasher::new();
let mut hasher_b = DefaultHasher::new();
super::hash_value(&input_a, &mut hasher_a);
super::hash_value(&input_b, &mut hasher_b);
assert_eq!(
hasher_a.finish(),
hasher_b.finish(),
"nested keys in different order must produce the same hash"
);
}
#[test]
fn plugin_claims_count_and_expire() {
let claims = PluginClaims::default();
let call = json!({"tool": "read_file", "input": {"path": "/a"}});
claims.insert("s", &call);
claims.insert("s", &call);
assert!(claims.consume("s", &call));
assert!(claims.consume("s", &call));
assert!(!claims.consume("s", &call), "two claims, two folds");
assert_eq!(claims.len(), 0);
claims.insert("cli", &call);
claims.insert("cli", &json!({"prompt": "hi"}));
assert_eq!(claims.len(), 2);
claims.prune(Instant::now() + PLUGIN_CLAIM_TTL + Duration::from_secs(1));
assert_eq!(
claims.len(),
0,
"expired claims are dropped by time, not by count"
);
}
}