use std::cmp::Reverse;
use std::collections::hash_map::DefaultHasher;
use std::collections::VecDeque;
use std::hash::{Hash, Hasher};
use std::sync::atomic::{AtomicU64, Ordering};
use std::time::{Duration, Instant};
use dashmap::DashMap;
pub(super) const RECENT_TOOL_USE_IDS: usize = 8;
const RECENT_PROMPT_HASHES: usize = 8;
pub fn text_hash(s: &str) -> u64 {
let mut h = DefaultHasher::new();
s.trim().hash(&mut h);
h.finish()
}
pub type InstallId = String;
pub const SESSION_QUIET_WINDOW: Duration = Duration::from_secs(300);
const MAX_TRACKED_INSTALLS: usize = 1024;
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub enum Assurance {
Attested,
Inferred,
Unknown,
}
pub const ASSURANCE_KEY: &str = "ai.openlatch.session.assurance";
impl Assurance {
pub fn as_str(&self) -> &'static str {
match self {
Assurance::Attested => "attested",
Assurance::Inferred => "inferred",
Assurance::Unknown => "unknown",
}
}
}
#[derive(Clone, Debug)]
pub struct SessionActivity {
pub agent_id: String,
pub source: String,
pub session_id: String,
pub last_seen: Instant,
pub recent_tool_use_ids: VecDeque<String>,
pub recent_prompt_hashes: VecDeque<u64>,
}
#[derive(Clone, Debug, Default)]
pub struct SessionSignals<'a> {
pub tool_use_id: Option<&'a str>,
pub prompt: Option<&'a str>,
}
#[derive(Clone, Debug, Default)]
pub struct RequestSignals {
pub declared_session_id: Option<String>,
pub tool_use_ids: Vec<String>,
pub prompt_hashes: Vec<u64>,
}
impl RequestSignals {
fn is_empty(&self) -> bool {
self.declared_session_id.is_none()
&& self.tool_use_ids.is_empty()
&& self.prompt_hashes.is_empty()
}
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub enum Selector {
OnlyLive,
SessionId,
SessionIdUnregistered,
ToolResult,
Prompt,
MostRecentActive,
}
impl Selector {
pub fn as_str(self) -> &'static str {
match self {
Selector::OnlyLive => "only_live",
Selector::SessionId => "session_id",
Selector::SessionIdUnregistered => "session_id_unregistered",
Selector::ToolResult => "tool_result",
Selector::Prompt => "prompt",
Selector::MostRecentActive => "most_recent_active",
}
}
pub const ALL: [Selector; 6] = [
Selector::OnlyLive,
Selector::SessionId,
Selector::SessionIdUnregistered,
Selector::ToolResult,
Selector::Prompt,
Selector::MostRecentActive,
];
fn index(self) -> usize {
match self {
Selector::OnlyLive => 0,
Selector::SessionId => 1,
Selector::SessionIdUnregistered => 2,
Selector::ToolResult => 3,
Selector::Prompt => 4,
Selector::MostRecentActive => 5,
}
}
fn assurance(self) -> Assurance {
match self {
Selector::OnlyLive
| Selector::SessionId
| Selector::SessionIdUnregistered
| Selector::ToolResult
| Selector::Prompt => Assurance::Attested,
Selector::MostRecentActive => Assurance::Inferred,
}
}
}
impl SessionActivity {
fn apply(&mut self, signals: &SessionSignals<'_>) {
if let Some(id) = signals.tool_use_id.filter(|s| !s.is_empty()) {
if !self.recent_tool_use_ids.iter().any(|k| k == id) {
self.recent_tool_use_ids.push_front(id.to_string());
self.recent_tool_use_ids.truncate(RECENT_TOOL_USE_IDS);
}
}
if let Some(p) = signals.prompt.filter(|s| !s.trim().is_empty()) {
let h = text_hash(p);
if !self.recent_prompt_hashes.contains(&h) {
self.recent_prompt_hashes.push_front(h);
self.recent_prompt_hashes.truncate(RECENT_PROMPT_HASHES);
}
}
}
}
#[derive(Default)]
pub struct SessionRegistry {
active: DashMap<InstallId, Vec<SessionActivity>>,
}
impl SessionRegistry {
pub fn upsert(&self, install_id: &str, agent_id: &str, source: &str, session_id: &str) {
self.upsert_signals(
install_id,
agent_id,
source,
session_id,
&SessionSignals::default(),
);
}
pub fn upsert_signals(
&self,
install_id: &str,
agent_id: &str,
source: &str,
session_id: &str,
signals: &SessionSignals<'_>,
) {
if !self.active.contains_key(install_id) && self.active.len() >= MAX_TRACKED_INSTALLS {
self.evict_empty();
if self.active.len() >= MAX_TRACKED_INSTALLS {
return;
}
}
let now = Instant::now();
let mut entry = self.active.entry(install_id.to_string()).or_default();
entry.retain(|a| now.duration_since(a.last_seen) < SESSION_QUIET_WINDOW);
if let Some(existing) = entry.iter_mut().find(|a| a.session_id == session_id) {
existing.last_seen = now;
existing.agent_id = agent_id.to_string();
existing.source = source.to_string();
existing.apply(signals);
} else {
let mut fresh = SessionActivity {
agent_id: agent_id.to_string(),
source: source.to_string(),
session_id: session_id.to_string(),
last_seen: now,
recent_tool_use_ids: VecDeque::new(),
recent_prompt_hashes: VecDeque::new(),
};
fresh.apply(signals);
entry.push(fresh);
}
}
fn evict_empty(&self) {
let now = Instant::now();
self.active.retain(|_, v| {
v.retain(|a| now.duration_since(a.last_seen) < SESSION_QUIET_WINDOW);
!v.is_empty()
});
}
fn fresh(&self, install_id: &str) -> Vec<SessionActivity> {
let now = Instant::now();
match self.active.get(install_id) {
Some(v) => v
.iter()
.filter(|a| now.duration_since(a.last_seen) < SESSION_QUIET_WINDOW)
.cloned()
.collect(),
None => Vec::new(),
}
}
}
#[derive(Clone, Debug)]
pub struct Resolved {
pub agent_id: Option<String>,
pub source: Option<String>,
pub session_id: Option<String>,
pub assurance: Assurance,
}
impl Resolved {
pub fn unknown() -> Self {
Resolved {
agent_id: None,
source: None,
session_id: None,
assurance: Assurance::Unknown,
}
}
}
pub fn resolve_session(reg: &SessionRegistry, install: &str) -> Resolved {
resolve_session_with(reg, install, &RequestSignals::default())
}
pub fn resolve_session_with(
reg: &SessionRegistry,
install: &str,
signals: &RequestSignals,
) -> Resolved {
resolve_session_detailed(reg, install, signals).0
}
static SELECTOR_WINS: [AtomicU64; 6] = [
AtomicU64::new(0),
AtomicU64::new(0),
AtomicU64::new(0),
AtomicU64::new(0),
AtomicU64::new(0),
AtomicU64::new(0),
];
static NO_LIVE_SESSION: AtomicU64 = AtomicU64::new(0);
pub fn selector_wins() -> Vec<(&'static str, u64)> {
let mut out: Vec<(&'static str, u64)> = Selector::ALL
.iter()
.map(|s| (s.as_str(), SELECTOR_WINS[s.index()].load(Ordering::Relaxed)))
.collect();
out.push(("no_live_session", NO_LIVE_SESSION.load(Ordering::Relaxed)));
out
}
pub fn resolve_session_detailed(
reg: &SessionRegistry,
install: &str,
signals: &RequestSignals,
) -> (Resolved, Option<Selector>) {
let active = reg.fresh(install);
if let Some(sid) = signals
.declared_session_id
.as_deref()
.map(str::trim)
.filter(|s| !s.is_empty())
{
let selector = match active.iter().find(|a| a.session_id == sid) {
Some(hit) => {
SELECTOR_WINS[Selector::SessionId.index()].fetch_add(1, Ordering::Relaxed);
return (
resolved(hit, Selector::SessionId),
Some(Selector::SessionId),
);
}
None => Selector::SessionIdUnregistered,
};
SELECTOR_WINS[selector.index()].fetch_add(1, Ordering::Relaxed);
return (
Resolved {
agent_id: Some(install.to_string()),
source: None,
session_id: Some(sid.to_string()),
assurance: selector.assurance(),
},
Some(selector),
);
}
if active.is_empty() {
NO_LIVE_SESSION.fetch_add(1, Ordering::Relaxed);
return (Resolved::unknown(), None);
}
if active.len() == 1 {
SELECTOR_WINS[Selector::OnlyLive.index()].fetch_add(1, Ordering::Relaxed);
return (
resolved(&active[0], Selector::OnlyLive),
Some(Selector::OnlyLive),
);
}
let (winner, selector) = if signals.is_empty() {
(most_recent(active), Selector::MostRecentActive)
} else {
cascade_pick(active, signals)
};
SELECTOR_WINS[selector.index()].fetch_add(1, Ordering::Relaxed);
(resolved(&winner, selector), Some(selector))
}
fn cascade_pick(
active: Vec<SessionActivity>,
signals: &RequestSignals,
) -> (SessionActivity, Selector) {
let mut pool = active;
if !signals.tool_use_ids.is_empty() {
let matched = narrow(&pool, |a| {
signals
.tool_use_ids
.iter()
.any(|id| a.recent_tool_use_ids.iter().any(|k| k == id))
});
match matched.len() {
1 => {
return (
matched.into_iter().next().expect("len checked"),
Selector::ToolResult,
)
}
n if n > 1 => pool = matched,
_ => {}
}
}
if !signals.prompt_hashes.is_empty() {
let matched = narrow(&pool, |a| {
a.recent_prompt_hashes
.iter()
.any(|h| signals.prompt_hashes.contains(h))
});
match matched.len() {
1 => {
return (
matched.into_iter().next().expect("len checked"),
Selector::Prompt,
)
}
n if n > 1 => pool = matched,
_ => {}
}
}
(most_recent(pool), Selector::MostRecentActive)
}
fn narrow(
pool: &[SessionActivity],
pred: impl Fn(&SessionActivity) -> bool,
) -> Vec<SessionActivity> {
pool.iter().filter(|a| pred(a)).cloned().collect()
}
fn most_recent(mut pool: Vec<SessionActivity>) -> SessionActivity {
pool.sort_by_key(|a| Reverse(a.last_seen));
pool.remove(0)
}
fn resolved(a: &SessionActivity, selector: Selector) -> Resolved {
Resolved {
agent_id: Some(a.agent_id.clone()),
source: Some(a.source.clone()),
session_id: Some(a.session_id.clone()),
assurance: selector.assurance(),
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn zero_active_resolves_unknown() {
let reg = SessionRegistry::default();
let r = resolve_session(®, "agt_absent");
assert_eq!(r.assurance, Assurance::Unknown);
assert!(r.agent_id.is_none() && r.session_id.is_none() && r.source.is_none());
}
#[test]
fn one_active_resolves_attested_with_identifiers() {
let reg = SessionRegistry::default();
reg.upsert("agt_1", "agt_1", "claude-code", "sess_a");
let r = resolve_session(®, "agt_1");
assert_eq!(r.assurance, Assurance::Attested);
assert_eq!(r.session_id.as_deref(), Some("sess_a"));
assert_eq!(r.agent_id.as_deref(), Some("agt_1"));
assert_eq!(r.source.as_deref(), Some("claude-code"));
}
#[test]
fn two_concurrent_resolve_inferred_most_recent() {
let reg = SessionRegistry::default();
reg.upsert("agt_1", "agt_1", "claude-code", "sess_old");
std::thread::sleep(Duration::from_millis(5));
reg.upsert("agt_1", "agt_1", "claude-code", "sess_new");
let r = resolve_session(®, "agt_1");
assert_eq!(r.assurance, Assurance::Inferred);
assert_eq!(r.session_id.as_deref(), Some("sess_new"));
}
#[test]
fn refresh_updates_last_seen_not_count() {
let reg = SessionRegistry::default();
reg.upsert("agt_1", "agt_1", "claude-code", "sess_a");
reg.upsert("agt_1", "agt_1", "claude-code", "sess_a");
let r = resolve_session(®, "agt_1");
assert_eq!(r.assurance, Assurance::Attested);
}
#[test]
fn assurance_wire_strings_are_the_frozen_three() {
assert_eq!(Assurance::Attested.as_str(), "attested");
assert_eq!(Assurance::Inferred.as_str(), "inferred");
assert_eq!(Assurance::Unknown.as_str(), "unknown");
}
fn two_live() -> SessionRegistry {
let reg = SessionRegistry::default();
reg.upsert_signals(
"i",
"i",
"claude-code",
"sess_a",
&SessionSignals {
tool_use_id: Some("toolu_aaa"),
prompt: Some("build the parser"),
},
);
std::thread::sleep(Duration::from_millis(5));
reg.upsert_signals(
"i",
"i",
"claude-code",
"sess_b",
&SessionSignals {
tool_use_id: Some("toolu_bbb"),
prompt: Some("write the docs"),
},
);
reg
}
#[test]
fn no_signals_at_all_keeps_most_recent_active() {
let reg = two_live();
let r = resolve_session_with(®, "i", &RequestSignals::default());
assert_eq!(r.session_id.as_deref(), Some("sess_b"));
assert_eq!(r.assurance, Assurance::Inferred);
}
#[test]
fn tool_result_join_beats_recency_and_attests() {
let reg = two_live();
let signals = RequestSignals {
tool_use_ids: vec!["toolu_aaa".into()],
..Default::default()
};
let r = resolve_session_with(®, "i", &signals);
assert_eq!(r.session_id.as_deref(), Some("sess_a"));
assert_eq!(r.assurance, Assurance::Attested);
}
#[test]
fn prompt_match_attests() {
let reg = two_live();
let signals = RequestSignals {
prompt_hashes: vec![text_hash("build the parser")],
..Default::default()
};
let r = resolve_session_with(®, "i", &signals);
assert_eq!(r.session_id.as_deref(), Some("sess_a"));
assert_eq!(r.assurance, Assurance::Attested);
}
#[test]
fn signals_that_match_nothing_fall_through_to_most_recent() {
let reg = two_live();
let signals = RequestSignals {
declared_session_id: None,
tool_use_ids: vec!["toolu_zzz".into()],
prompt_hashes: vec![text_hash("something else entirely")],
};
let r = resolve_session_with(®, "i", &signals);
assert_eq!(r.session_id.as_deref(), Some("sess_b"));
assert_eq!(r.assurance, Assurance::Inferred);
}
#[test]
fn a_tie_narrows_the_pool_before_the_fallback_decides() {
let reg = SessionRegistry::default();
for sess in ["sess_a", "sess_b"] {
reg.upsert_signals(
"i",
"i",
"claude-code",
sess,
&SessionSignals {
prompt: Some("same opening prompt"),
..Default::default()
},
);
std::thread::sleep(Duration::from_millis(5));
}
reg.upsert_signals(
"i",
"i",
"claude-code",
"sess_outsider",
&SessionSignals {
prompt: Some("a different opening"),
..Default::default()
},
);
let signals = RequestSignals {
prompt_hashes: vec![text_hash("same opening prompt")],
..Default::default()
};
let r = resolve_session_with(®, "i", &signals);
assert_eq!(r.session_id.as_deref(), Some("sess_b"));
assert_eq!(r.assurance, Assurance::Inferred);
}
#[test]
fn empty_signals_with_cascade_on_are_the_metadata_only_path() {
let reg = two_live();
let r = resolve_session_with(®, "i", &RequestSignals::default());
assert_eq!(r.session_id.as_deref(), Some("sess_b"));
assert_eq!(r.assurance, Assurance::Inferred);
}
#[test]
fn one_live_session_attests_with_the_cascade_on() {
let reg = SessionRegistry::default();
reg.upsert("i", "i", "claude-code", "sess_only");
let signals = RequestSignals {
tool_use_ids: vec!["toolu_unrelated".into()],
..Default::default()
};
let r = resolve_session_with(®, "i", &signals);
assert_eq!(r.session_id.as_deref(), Some("sess_only"));
assert_eq!(r.assurance, Assurance::Attested);
}
#[test]
fn no_live_session_is_unknown_whatever_the_signals() {
let signals = RequestSignals {
tool_use_ids: vec!["toolu_aaa".into()],
..Default::default()
};
let r = resolve_session_with(&SessionRegistry::default(), "i", &signals);
assert_eq!(r.assurance, Assurance::Unknown);
assert!(r.session_id.is_none());
}
#[test]
fn a_request_that_names_its_session_beats_every_heuristic() {
let reg = two_live();
let signals = RequestSignals {
declared_session_id: Some("sess_a".into()),
tool_use_ids: vec!["toolu_bbb".into()],
..Default::default()
};
let (r, sel) = resolve_session_detailed(®, "i", &signals);
assert_eq!(r.session_id.as_deref(), Some("sess_a"));
assert_eq!(sel, Some(Selector::SessionId));
assert_eq!(r.assurance, Assurance::Attested);
assert_eq!(r.source.as_deref(), Some("claude-code"));
}
#[test]
fn a_named_session_the_registry_has_not_seen_yet_is_still_attributed() {
let reg = two_live();
let signals = RequestSignals {
declared_session_id: Some("sess_brand_new".into()),
..Default::default()
};
let (r, sel) = resolve_session_detailed(®, "i", &signals);
assert_eq!(r.session_id.as_deref(), Some("sess_brand_new"));
assert_eq!(sel, Some(Selector::SessionIdUnregistered));
assert_eq!(r.assurance, Assurance::Attested);
assert_eq!(r.agent_id.as_deref(), Some("i"));
assert!(
r.source.is_none(),
"a source we do not know must not be invented"
);
}
#[test]
fn a_named_session_resolves_even_with_an_empty_registry() {
let signals = RequestSignals {
declared_session_id: Some("sess_first".into()),
..Default::default()
};
let (r, sel) = resolve_session_detailed(&SessionRegistry::default(), "i", &signals);
assert_eq!(r.session_id.as_deref(), Some("sess_first"));
assert_eq!(sel, Some(Selector::SessionIdUnregistered));
assert_eq!(r.assurance, Assurance::Attested);
}
#[test]
fn the_zombie_session_misattribution_is_fixed() {
let reg = SessionRegistry::default();
for zombie in ["6693e5f6", "a5a39553", "92129c1b"] {
reg.upsert("i", "i", "claude-code", zombie);
}
let signals = RequestSignals {
declared_session_id: Some("5c7d9833".into()),
..Default::default()
};
let blind = resolve_session_with(®, "i", &RequestSignals::default());
assert_ne!(blind.session_id.as_deref(), Some("5c7d9833"));
assert_eq!(blind.assurance, Assurance::Inferred);
let seen = resolve_session_with(®, "i", &signals);
assert_eq!(seen.session_id.as_deref(), Some("5c7d9833"));
assert_eq!(seen.assurance, Assurance::Attested);
}
#[test]
fn an_empty_declared_session_id_falls_through_rather_than_labelling_nothing() {
let reg = two_live();
let signals = RequestSignals {
declared_session_id: Some(" ".into()),
tool_use_ids: vec!["toolu_aaa".into()],
..Default::default()
};
let (r, sel) = resolve_session_detailed(®, "i", &signals);
assert_eq!(r.session_id.as_deref(), Some("sess_a"));
assert_eq!(sel, Some(Selector::ToolResult));
}
#[test]
fn a_resolution_is_counted_under_the_selector_that_decided() {
let win = |sel: Selector| -> u64 {
selector_wins()
.into_iter()
.find(|(k, _)| *k == sel.as_str())
.map(|(_, v)| v)
.expect("every selector is reported")
};
let unattributed = || -> u64 {
selector_wins()
.into_iter()
.find(|(k, _)| *k == "no_live_session")
.map(|(_, v)| v)
.expect("no_live_session is reported")
};
let before = win(Selector::ToolResult);
let reg = two_live();
let signals = RequestSignals {
tool_use_ids: vec!["toolu_aaa".to_string()],
..Default::default()
};
let (_, sel) = resolve_session_detailed(®, "i", &signals);
assert_eq!(sel, Some(Selector::ToolResult));
assert!(
win(Selector::ToolResult) > before,
"a tool_result win must be counted"
);
let before_unattributed = unattributed();
let empty = SessionRegistry::default();
let (r, sel) = resolve_session_detailed(&empty, "nobody", &signals);
assert_eq!(sel, None);
assert_eq!(r.assurance, Assurance::Unknown);
assert!(
unattributed() > before_unattributed,
"a resolution with no live session must be counted as such, not as a selector win"
);
}
#[test]
fn prompt_hashes_accumulate_newest_first_and_are_bounded() {
let reg = SessionRegistry::default();
for p in ["same turn", "same turn"] {
reg.upsert_signals(
"i",
"i",
"claude-code",
"s",
&SessionSignals {
prompt: Some(p),
..Default::default()
},
);
}
assert_eq!(reg.fresh("i")[0].recent_prompt_hashes.len(), 1);
for n in 0..(RECENT_PROMPT_HASHES + 5) {
reg.upsert_signals(
"i",
"i",
"claude-code",
"s",
&SessionSignals {
prompt: Some(&format!("turn {n}")),
..Default::default()
},
);
}
let live = reg.fresh("i");
assert_eq!(live[0].recent_prompt_hashes.len(), RECENT_PROMPT_HASHES);
assert_eq!(
live[0].recent_prompt_hashes.front().copied(),
Some(text_hash(&format!("turn {}", RECENT_PROMPT_HASHES + 4)))
);
}
#[test]
fn an_entry_born_mid_conversation_still_joins() {
let reg = SessionRegistry::default();
for (sess, witnessed) in [
("sess_a", "waht is session Id"),
("sess_b", "write the docs"),
] {
reg.upsert_signals(
"i",
"i",
"claude-code",
sess,
&SessionSignals {
prompt: Some(witnessed),
..Default::default()
},
);
std::thread::sleep(Duration::from_millis(5));
}
let signals = RequestSignals {
declared_session_id: None,
tool_use_ids: Vec::new(),
prompt_hashes: vec![
text_hash("hello this a session without a cache"),
text_hash("waht is session Id"),
],
};
let r = resolve_session_with(®, "i", &signals);
assert_eq!(r.session_id.as_deref(), Some("sess_a"));
assert_eq!(r.assurance, Assurance::Attested);
let opening_only = RequestSignals {
declared_session_id: None,
tool_use_ids: Vec::new(),
prompt_hashes: vec![text_hash("hello this a session without a cache")],
};
let r = resolve_session_with(®, "i", &opening_only);
assert_eq!(r.session_id.as_deref(), Some("sess_b"));
assert_eq!(r.assurance, Assurance::Inferred);
}
#[test]
fn tool_use_ids_are_deduped_newest_first_and_bounded() {
let reg = SessionRegistry::default();
for id in ["toolu_1", "toolu_1"] {
reg.upsert_signals(
"i",
"i",
"claude-code",
"s",
&SessionSignals {
tool_use_id: Some(id),
..Default::default()
},
);
}
assert_eq!(reg.fresh("i")[0].recent_tool_use_ids.len(), 1);
for n in 2..(RECENT_TOOL_USE_IDS + 5) {
reg.upsert_signals(
"i",
"i",
"claude-code",
"s",
&SessionSignals {
tool_use_id: Some(&format!("toolu_{n}")),
..Default::default()
},
);
}
let live = reg.fresh("i");
assert_eq!(live[0].recent_tool_use_ids.len(), RECENT_TOOL_USE_IDS);
assert_eq!(
live[0].recent_tool_use_ids.front().map(String::as_str),
Some(format!("toolu_{}", RECENT_TOOL_USE_IDS + 4).as_str())
);
}
#[test]
fn a_signal_free_envelope_does_not_clear_what_an_earlier_one_recorded() {
let reg = SessionRegistry::default();
reg.upsert_signals(
"i",
"i",
"claude-code",
"s",
&SessionSignals {
tool_use_id: Some("toolu_keep"),
prompt: Some("keep me"),
},
);
reg.upsert("i", "i", "claude-code", "s");
let live = reg.fresh("i");
assert_eq!(
live[0].recent_tool_use_ids.front().map(String::as_str),
Some("toolu_keep")
);
assert_eq!(
live[0].recent_prompt_hashes.front().copied(),
Some(text_hash("keep me"))
);
}
}