use std::collections::HashMap;
use std::path::{Path, PathBuf};
use std::sync::{Mutex, MutexGuard, PoisonError};
use std::time::Duration;
use chrono::{DateTime, Utc};
use serde::{Deserialize, Serialize};
use tokio::sync::watch;
pub mod stream;
pub mod watcher;
const DEFAULT_SESSION_TTL: Duration = Duration::from_secs(300);
const ENDED_SESSION_TTL: Duration = Duration::from_secs(10);
const DEFAULT_WINDOW_TTL: Duration = Duration::from_secs(30);
const MAX_SESSIONS: usize = 512;
const MAX_WINDOWS: usize = 256;
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum SessionState {
Starting,
Working,
Idle,
WaitingForInput,
WaitingForPermission,
Ended,
}
impl SessionState {
#[must_use]
pub fn for_event(event: &SessionEvent, current: Option<Self>) -> Self {
match event {
SessionEvent::SessionStart => Self::Starting,
SessionEvent::UserPromptSubmit
| SessionEvent::PreToolUse
| SessionEvent::PostToolUse => Self::Working,
SessionEvent::TranscriptGrew => match current {
Some(held @ (Self::WaitingForInput | Self::WaitingForPermission | Self::Ended)) => {
held
}
_ => Self::Working,
},
SessionEvent::Stop => Self::Idle,
SessionEvent::StreamState(state) => *state,
SessionEvent::Notification(NotificationKind::PermissionPrompt) => {
Self::WaitingForPermission
}
SessionEvent::Notification(
NotificationKind::IdlePrompt | NotificationKind::AgentNeedsInput,
) => Self::WaitingForInput,
SessionEvent::Notification(NotificationKind::Other)
| SessionEvent::TranscriptDiscovered => current.unwrap_or(Self::Idle),
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum NotificationKind {
PermissionPrompt,
IdlePrompt,
AgentNeedsInput,
Other,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum SessionEvent {
SessionStart,
UserPromptSubmit,
PreToolUse,
PostToolUse,
Stop,
Notification(NotificationKind),
TranscriptGrew,
TranscriptDiscovered,
StreamState(SessionState),
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(tag = "kind", rename_all = "snake_case")]
pub enum Source {
Terminal,
VsCode {
window_key: String,
},
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct ObserveRequest {
pub session_id: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub cwd: Option<PathBuf>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub transcript_path: Option<PathBuf>,
pub event: SessionEvent,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub repo: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub model: Option<String>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct WindowReport {
pub key: String,
#[serde(default)]
pub folders: Vec<PathBuf>,
#[serde(default)]
pub tabs: usize,
#[serde(default)]
pub terminals: usize,
}
impl WindowReport {
#[must_use]
fn has_embedding(&self) -> bool {
self.tabs > 0 || self.terminals > 0
}
}
#[derive(Debug, Clone, Serialize)]
pub struct SessionEntry {
pub session_id: String,
#[serde(skip_serializing_if = "Option::is_none")]
pub cwd: Option<PathBuf>,
#[serde(skip_serializing_if = "Option::is_none")]
pub transcript_path: Option<PathBuf>,
#[serde(skip_serializing_if = "Option::is_none")]
pub repo: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub model: Option<String>,
pub state: SessionState,
pub source: Source,
pub last_event: SessionEvent,
pub started_at: DateTime<Utc>,
pub last_seen: DateTime<Utc>,
}
#[derive(Debug, Clone)]
struct WindowEntry {
report: WindowReport,
last_seen: DateTime<Utc>,
}
pub struct SessionsRegistry {
sessions: Mutex<HashMap<String, SessionEntry>>,
windows: Mutex<HashMap<String, WindowEntry>>,
session_ttl: Duration,
ended_ttl: Duration,
window_ttl: Duration,
changes: watch::Sender<u64>,
}
impl SessionsRegistry {
#[must_use]
pub fn new() -> Self {
Self {
sessions: Mutex::new(HashMap::new()),
windows: Mutex::new(HashMap::new()),
session_ttl: DEFAULT_SESSION_TTL,
ended_ttl: ENDED_SESSION_TTL,
window_ttl: DEFAULT_WINDOW_TTL,
changes: watch::channel(0).0,
}
}
#[must_use]
pub fn subscribe_changes(&self) -> watch::Receiver<u64> {
self.changes.subscribe()
}
pub(crate) fn bump(&self) {
self.changes.send_modify(|v| *v = v.wrapping_add(1));
}
fn lock_sessions(&self) -> MutexGuard<'_, HashMap<String, SessionEntry>> {
self.sessions.lock().unwrap_or_else(PoisonError::into_inner)
}
fn lock_windows(&self) -> MutexGuard<'_, HashMap<String, WindowEntry>> {
self.windows.lock().unwrap_or_else(PoisonError::into_inner)
}
pub fn observe(&self, req: ObserveRequest) {
let now = Utc::now();
let changed = {
let mut sessions = self.lock_sessions();
let reaped = reap_sessions(&mut sessions, self.session_ttl, self.ended_ttl, now);
let mutated = if let Some(entry) = sessions.get_mut(&req.session_id) {
let next = SessionState::for_event(&req.event, Some(entry.state));
let state_changed = next != entry.state;
entry.state = next;
entry.last_event = req.event;
entry.last_seen = now;
let filled_cwd = fill(&mut entry.cwd, req.cwd);
let filled_transcript = fill(&mut entry.transcript_path, req.transcript_path);
let filled_repo = fill(&mut entry.repo, req.repo);
let filled_model = fill(&mut entry.model, req.model);
state_changed || filled_cwd || filled_transcript || filled_repo || filled_model
} else {
if sessions.len() >= MAX_SESSIONS {
evict_oldest_session(&mut sessions);
}
let state = SessionState::for_event(&req.event, None);
sessions.insert(
req.session_id.clone(),
SessionEntry {
session_id: req.session_id,
cwd: req.cwd,
transcript_path: req.transcript_path,
repo: req.repo,
model: req.model,
state,
source: Source::Terminal,
last_event: req.event,
started_at: now,
last_seen: now,
},
);
true
};
mutated || reaped > 0
};
if changed {
self.bump();
}
}
pub fn end(&self, session_id: &str, _reason: Option<&str>) -> bool {
let now = Utc::now();
let (known, reaped) = {
let mut sessions = self.lock_sessions();
let reaped = reap_sessions(&mut sessions, self.session_ttl, self.ended_ttl, now);
let known = match sessions.get_mut(session_id) {
Some(entry) => {
entry.state = SessionState::Ended;
entry.last_event = SessionEvent::Stop;
entry.last_seen = now;
true
}
None => false,
};
(known, reaped)
};
if known || reaped > 0 {
self.bump();
}
known
}
pub fn report_window(&self, report: WindowReport) {
let now = Utc::now();
let changed = {
let mut windows = self.lock_windows();
let reaped = reap_windows(&mut windows, self.window_ttl, now);
let mutated = if let Some(previous) = windows.get(&report.key) {
previous.report.folders != report.folders
|| previous.report.has_embedding() != report.has_embedding()
} else {
if windows.len() >= MAX_WINDOWS {
evict_oldest_window(&mut windows);
}
true
};
windows.insert(
report.key.clone(),
WindowEntry {
report,
last_seen: now,
},
);
mutated || reaped > 0
};
if changed {
self.bump();
}
}
pub fn unregister_window(&self, key: &str) -> bool {
let removed = {
let mut windows = self.lock_windows();
windows.remove(key).is_some()
};
if removed {
self.bump();
}
removed
}
pub fn list(&self) -> Vec<SessionEntry> {
let now = Utc::now();
let mut sessions: Vec<SessionEntry> = {
let mut guard = self.lock_sessions();
reap_sessions(&mut guard, self.session_ttl, self.ended_ttl, now);
guard.values().cloned().collect()
};
let windows: Vec<WindowReport> = {
let mut guard = self.lock_windows();
reap_windows(&mut guard, self.window_ttl, now);
guard
.values()
.map(|e| e.report.clone())
.filter(WindowReport::has_embedding)
.collect()
};
for session in &mut sessions {
session.source = resolve_source(session.cwd.as_deref(), &windows);
}
sessions.sort_by(|a, b| {
a.repo
.cmp(&b.repo)
.then_with(|| a.session_id.cmp(&b.session_id))
});
sessions
}
pub fn focus_folder(&self, session_id: &str) -> Option<PathBuf> {
let cwd = {
let sessions = self.lock_sessions();
sessions.get(session_id).and_then(|e| e.cwd.clone())
}?;
let now = Utc::now();
let mut windows = self.lock_windows();
reap_windows(&mut windows, self.window_ttl, now);
windows
.values()
.map(|e| &e.report)
.filter(|w| w.has_embedding())
.filter(|w| w.folders.iter().any(|f| cwd.starts_with(f)))
.find_map(|w| w.folders.first().cloned())
}
}
impl Default for SessionsRegistry {
fn default() -> Self {
Self::new()
}
}
fn fill<T: PartialEq>(slot: &mut Option<T>, incoming: Option<T>) -> bool {
match incoming {
Some(value) if slot.as_ref() != Some(&value) => {
*slot = Some(value);
true
}
_ => false,
}
}
fn resolve_source(cwd: Option<&Path>, windows: &[WindowReport]) -> Source {
let Some(cwd) = cwd else {
return Source::Terminal;
};
let matched = windows
.iter()
.filter(|w| w.folders.iter().any(|f| cwd.starts_with(f)))
.min_by(|a, b| a.key.cmp(&b.key));
match matched {
Some(window) => Source::VsCode {
window_key: window.key.clone(),
},
None => Source::Terminal,
}
}
fn reap_sessions(
sessions: &mut HashMap<String, SessionEntry>,
session_ttl: Duration,
ended_ttl: Duration,
now: DateTime<Utc>,
) -> usize {
let session_max = session_ttl.as_secs() as i64;
let ended_max = ended_ttl.as_secs() as i64;
let before = sessions.len();
sessions.retain(|_, e| {
let max_age = if e.state == SessionState::Ended {
ended_max
} else {
session_max
};
(now - e.last_seen).num_seconds() <= max_age
});
before - sessions.len()
}
fn reap_windows(
windows: &mut HashMap<String, WindowEntry>,
ttl: Duration,
now: DateTime<Utc>,
) -> usize {
let max_age = ttl.as_secs() as i64;
let before = windows.len();
windows.retain(|_, e| (now - e.last_seen).num_seconds() <= max_age);
before - windows.len()
}
fn evict_oldest_session(sessions: &mut HashMap<String, SessionEntry>) {
let oldest = sessions
.values()
.min_by(|a, b| {
a.last_seen
.cmp(&b.last_seen)
.then_with(|| a.session_id.cmp(&b.session_id))
})
.map(|e| e.session_id.clone());
if let Some(key) = oldest {
sessions.remove(&key);
}
}
fn evict_oldest_window(windows: &mut HashMap<String, WindowEntry>) {
let oldest = windows
.iter()
.min_by(|a, b| a.1.last_seen.cmp(&b.1.last_seen).then_with(|| a.0.cmp(b.0)))
.map(|(k, _)| k.clone());
if let Some(key) = oldest {
windows.remove(&key);
}
}
#[cfg(test)]
#[allow(clippy::unwrap_used, clippy::expect_used)]
mod tests {
use super::*;
fn observe_request(session_id: &str, event: SessionEvent, cwd: Option<&str>) -> ObserveRequest {
ObserveRequest {
session_id: session_id.to_string(),
cwd: cwd.map(PathBuf::from),
transcript_path: None,
event,
repo: None,
model: None,
}
}
#[test]
fn list_is_empty_initially() {
let reg = SessionsRegistry::new();
assert!(reg.list().is_empty());
}
#[test]
fn observe_then_list_round_trips_and_infers_state() {
let reg = SessionsRegistry::new();
reg.observe(observe_request(
"s1",
SessionEvent::SessionStart,
Some("/tmp/a"),
));
let sessions = reg.list();
assert_eq!(sessions.len(), 1);
assert_eq!(sessions[0].session_id, "s1");
assert_eq!(sessions[0].state, SessionState::Starting);
assert_eq!(sessions[0].source, Source::Terminal);
}
#[test]
fn observe_is_idempotent_upsert_advancing_state() {
let reg = SessionsRegistry::new();
reg.observe(observe_request(
"s1",
SessionEvent::SessionStart,
Some("/tmp/a"),
));
reg.observe(observe_request("s1", SessionEvent::PreToolUse, None));
let sessions = reg.list();
assert_eq!(sessions.len(), 1, "same session_id upserts, not duplicates");
assert_eq!(sessions[0].state, SessionState::Working);
assert_eq!(sessions[0].cwd.as_deref(), Some(Path::new("/tmp/a")));
}
#[test]
fn state_machine_covers_every_event() {
use NotificationKind::*;
use SessionEvent::*;
let cases = [
(SessionStart, SessionState::Starting),
(UserPromptSubmit, SessionState::Working),
(PreToolUse, SessionState::Working),
(PostToolUse, SessionState::Working),
(Stop, SessionState::Idle),
(
Notification(PermissionPrompt),
SessionState::WaitingForPermission,
),
(Notification(IdlePrompt), SessionState::WaitingForInput),
(Notification(AgentNeedsInput), SessionState::WaitingForInput),
(TranscriptGrew, SessionState::Working),
(TranscriptDiscovered, SessionState::Idle),
(
StreamState(SessionState::WaitingForPermission),
SessionState::WaitingForPermission,
),
(StreamState(SessionState::Idle), SessionState::Idle),
];
for (event, expected) in cases {
assert_eq!(
SessionState::for_event(&event, None),
expected,
"event {event:?}"
);
}
assert_eq!(
SessionState::for_event(&Notification(Other), Some(SessionState::Working)),
SessionState::Working
);
assert_eq!(
SessionState::for_event(&TranscriptDiscovered, Some(SessionState::Working)),
SessionState::Working
);
assert_eq!(
SessionState::for_event(
&StreamState(SessionState::Idle),
Some(SessionState::Working)
),
SessionState::Idle
);
for held in [
SessionState::WaitingForInput,
SessionState::WaitingForPermission,
SessionState::Ended,
] {
assert_eq!(
SessionState::for_event(&TranscriptGrew, Some(held)),
held,
"growth must not overwrite {held:?}"
);
}
for other in [
SessionState::Working,
SessionState::Idle,
SessionState::Starting,
] {
assert_eq!(
SessionState::for_event(&TranscriptGrew, Some(other)),
SessionState::Working,
"growth from {other:?}"
);
}
}
#[test]
fn end_marks_ended_and_reaps_quickly() {
let reg = SessionsRegistry::new();
reg.observe(observe_request(
"s1",
SessionEvent::PreToolUse,
Some("/tmp/a"),
));
assert!(reg.end("s1", Some("clear")));
assert!(!reg.end("ghost", None));
let sessions = reg.list();
assert_eq!(sessions.len(), 1);
assert_eq!(sessions[0].state, SessionState::Ended);
{
let mut guard = reg.lock_sessions();
guard.get_mut("s1").unwrap().last_seen = Utc::now() - chrono::Duration::seconds(30);
}
assert!(reg.list().is_empty(), "ended entry reaps after ended TTL");
}
#[test]
fn stale_working_session_reaps_but_recent_survives() {
let reg = SessionsRegistry::new();
reg.observe(observe_request("fresh", SessionEvent::PreToolUse, None));
reg.observe(observe_request("stale", SessionEvent::PreToolUse, None));
{
let mut guard = reg.lock_sessions();
guard.get_mut("stale").unwrap().last_seen =
Utc::now() - chrono::Duration::seconds(1000);
}
let ids: Vec<String> = reg.list().into_iter().map(|s| s.session_id).collect();
assert_eq!(ids, vec!["fresh".to_string()]);
}
#[test]
fn source_is_vscode_when_cwd_is_under_a_reporting_window() {
let reg = SessionsRegistry::new();
reg.observe(observe_request(
"s1",
SessionEvent::PreToolUse,
Some("/home/me/proj/sub"),
));
reg.report_window(WindowReport {
key: "w1".to_string(),
folders: vec![PathBuf::from("/home/me/proj")],
tabs: 1,
terminals: 0,
});
let sessions = reg.list();
assert_eq!(
sessions[0].source,
Source::VsCode {
window_key: "w1".to_string()
}
);
}
#[test]
fn source_is_terminal_when_window_has_no_embedding() {
let reg = SessionsRegistry::new();
reg.observe(observe_request(
"s1",
SessionEvent::PreToolUse,
Some("/home/me/proj"),
));
reg.report_window(WindowReport {
key: "w1".to_string(),
folders: vec![PathBuf::from("/home/me/proj")],
tabs: 0,
terminals: 0,
});
assert_eq!(reg.list()[0].source, Source::Terminal);
}
#[test]
fn window_report_is_upsert_and_unregister_removes() {
let reg = SessionsRegistry::new();
reg.report_window(WindowReport {
key: "w1".to_string(),
folders: vec![PathBuf::from("/p")],
tabs: 1,
terminals: 0,
});
reg.report_window(WindowReport {
key: "w1".to_string(),
folders: vec![PathBuf::from("/p")],
tabs: 2,
terminals: 1,
});
assert!(reg.unregister_window("w1"));
assert!(!reg.unregister_window("w1"));
}
#[test]
fn stale_window_stops_tagging_source() {
let reg = SessionsRegistry::new();
reg.observe(observe_request(
"s1",
SessionEvent::PreToolUse,
Some("/p/sub"),
));
reg.report_window(WindowReport {
key: "w1".to_string(),
folders: vec![PathBuf::from("/p")],
tabs: 1,
terminals: 0,
});
{
let mut guard = reg.lock_windows();
guard.get_mut("w1").unwrap().last_seen = Utc::now() - chrono::Duration::seconds(120);
}
assert_eq!(reg.list()[0].source, Source::Terminal);
}
#[test]
fn resolve_source_prefers_lowest_key_on_overlap() {
let windows = vec![
WindowReport {
key: "w2".to_string(),
folders: vec![PathBuf::from("/p")],
tabs: 1,
terminals: 0,
},
WindowReport {
key: "w1".to_string(),
folders: vec![PathBuf::from("/p")],
tabs: 1,
terminals: 0,
},
];
assert_eq!(
resolve_source(Some(Path::new("/p/x")), &windows),
Source::VsCode {
window_key: "w1".to_string()
}
);
assert_eq!(resolve_source(None, &windows), Source::Terminal);
}
#[test]
fn focus_folder_resolves_matching_window_folder() {
let reg = SessionsRegistry::new();
reg.observe(observe_request(
"s1",
SessionEvent::PreToolUse,
Some("/home/me/proj/sub"),
));
assert!(reg.focus_folder("s1").is_none(), "no window yet");
reg.report_window(WindowReport {
key: "w1".to_string(),
folders: vec![PathBuf::from("/home/me/proj")],
tabs: 1,
terminals: 0,
});
assert_eq!(reg.focus_folder("s1"), Some(PathBuf::from("/home/me/proj")));
assert!(reg.focus_folder("ghost").is_none());
}
#[test]
fn evict_oldest_session_drops_the_longest_silent() {
let now = Utc::now();
let mut sessions = HashMap::new();
for (id, age) in [("young", 0), ("old", 100), ("older", 200)] {
sessions.insert(
id.to_string(),
SessionEntry {
session_id: id.to_string(),
cwd: None,
transcript_path: None,
repo: None,
model: None,
state: SessionState::Working,
source: Source::Terminal,
last_event: SessionEvent::PreToolUse,
started_at: now,
last_seen: now - chrono::Duration::seconds(age),
},
);
}
evict_oldest_session(&mut sessions);
assert!(!sessions.contains_key("older"));
assert!(sessions.contains_key("young"));
assert!(sessions.contains_key("old"));
}
#[test]
fn list_sorts_by_repo_then_session_id() {
let reg = SessionsRegistry::new();
for (id, repo) in [("z", "repo-a"), ("a", "repo-b"), ("m", "repo-a")] {
reg.observe(ObserveRequest {
session_id: id.to_string(),
cwd: None,
transcript_path: None,
event: SessionEvent::PreToolUse,
repo: Some(repo.to_string()),
model: None,
});
}
let ordered: Vec<(String, String)> = reg
.list()
.into_iter()
.map(|s| (s.session_id, s.repo.unwrap()))
.collect();
assert_eq!(
ordered,
vec![
("m".to_string(), "repo-a".to_string()),
("z".to_string(), "repo-a".to_string()),
("a".to_string(), "repo-b".to_string()),
]
);
}
#[test]
fn serialized_session_shapes_are_stable() {
let reg = SessionsRegistry::new();
reg.observe(ObserveRequest {
session_id: "s1".to_string(),
cwd: Some(PathBuf::from("/p")),
transcript_path: None,
event: SessionEvent::Notification(NotificationKind::PermissionPrompt),
repo: Some("proj".to_string()),
model: None,
});
let value = serde_json::to_value(®.list()[0]).unwrap();
assert_eq!(value["state"], "waiting_for_permission");
assert_eq!(value["source"]["kind"], "terminal");
assert_eq!(value["repo"], "proj");
assert!(value.get("model").is_none());
assert!(value.get("transcript_path").is_none());
}
#[test]
fn stream_state_is_authoritative_and_round_trips() {
let reg = SessionsRegistry::new();
reg.observe(observe_request("s1", SessionEvent::PreToolUse, Some("/p")));
assert_eq!(reg.list()[0].state, SessionState::Working);
reg.observe(observe_request(
"s1",
SessionEvent::StreamState(SessionState::WaitingForPermission),
None,
));
assert_eq!(reg.list()[0].state, SessionState::WaitingForPermission);
let value = serde_json::to_value(SessionEvent::StreamState(
SessionState::WaitingForPermission,
))
.unwrap();
assert_eq!(
value,
serde_json::json!({ "stream_state": "waiting_for_permission" })
);
let event: SessionEvent = serde_json::from_value(value).unwrap();
assert_eq!(
event,
SessionEvent::StreamState(SessionState::WaitingForPermission)
);
}
#[test]
fn default_constructs_an_empty_registry() {
let reg = SessionsRegistry::default();
assert!(reg.list().is_empty());
}
#[test]
fn fill_only_overwrites_with_a_present_value() {
let mut slot = Some("keep");
fill(&mut slot, None);
assert_eq!(slot, Some("keep"));
fill(&mut slot, Some("new"));
assert_eq!(slot, Some("new"));
let mut empty: Option<&str> = None;
fill(&mut empty, Some("filled"));
assert_eq!(empty, Some("filled"));
}
#[test]
fn observe_at_session_cap_evicts_the_longest_silent() {
let reg = SessionsRegistry::new();
{
let mut sessions = reg.lock_sessions();
let base = Utc::now();
for i in 0..MAX_SESSIONS {
let id = format!("s{i:04}");
sessions.insert(
id.clone(),
SessionEntry {
session_id: id.clone(),
cwd: None,
transcript_path: None,
repo: None,
model: None,
state: SessionState::Working,
source: Source::Terminal,
last_event: SessionEvent::PreToolUse,
started_at: base,
last_seen: base - chrono::Duration::milliseconds(i as i64),
},
);
}
}
reg.observe(observe_request("fresh", SessionEvent::PreToolUse, None));
let sessions = reg.lock_sessions();
assert_eq!(sessions.len(), MAX_SESSIONS);
assert!(sessions.contains_key("fresh"));
assert!(!sessions.contains_key(&format!("s{:04}", MAX_SESSIONS - 1)));
assert!(sessions.contains_key("s0000"));
}
#[test]
fn report_window_at_cap_evicts_the_longest_silent() {
let reg = SessionsRegistry::new();
{
let mut windows = reg.lock_windows();
let base = Utc::now();
for i in 0..MAX_WINDOWS {
let key = format!("w{i:04}");
windows.insert(
key.clone(),
WindowEntry {
report: WindowReport {
key: key.clone(),
folders: vec![],
tabs: 1,
terminals: 0,
},
last_seen: base - chrono::Duration::milliseconds(i as i64),
},
);
}
}
reg.report_window(WindowReport {
key: "fresh".to_string(),
folders: vec![],
tabs: 1,
terminals: 0,
});
let windows = reg.lock_windows();
assert_eq!(windows.len(), MAX_WINDOWS);
assert!(windows.contains_key("fresh"));
assert!(!windows.contains_key(&format!("w{:04}", MAX_WINDOWS - 1)));
assert!(windows.contains_key("w0000"));
}
#[test]
fn evict_oldest_window_breaks_ties_by_key() {
let now = Utc::now();
let mut windows = HashMap::new();
let at = |key: &str, secs: i64| WindowEntry {
report: WindowReport {
key: key.to_string(),
folders: vec![],
tabs: 1,
terminals: 0,
},
last_seen: now - chrono::Duration::seconds(secs),
};
windows.insert("young".to_string(), at("young", 0));
windows.insert("old-b".to_string(), at("old-b", 10));
windows.insert("old-a".to_string(), at("old-a", 10));
evict_oldest_window(&mut windows);
assert!(!windows.contains_key("old-a"));
assert!(windows.contains_key("old-b"));
assert!(windows.contains_key("young"));
let mut empty: HashMap<String, WindowEntry> = HashMap::new();
evict_oldest_window(&mut empty);
assert!(empty.is_empty());
}
fn window_report(key: &str, folder: &str, embedded: bool) -> WindowReport {
WindowReport {
key: key.to_string(),
folders: vec![PathBuf::from(folder)],
tabs: usize::from(embedded),
terminals: 0,
}
}
#[test]
fn subscribe_changes_starts_seen_and_a_new_session_bumps() {
let reg = SessionsRegistry::new();
let mut rx = reg.subscribe_changes();
assert!(!rx.has_changed().unwrap());
reg.observe(observe_request("s1", SessionEvent::SessionStart, None));
assert!(rx.has_changed().unwrap(), "a new session should bump");
rx.borrow_and_update();
assert!(!rx.has_changed().unwrap());
}
#[test]
fn observe_bumps_on_a_state_transition_but_not_on_a_repeat_sighting() {
let reg = SessionsRegistry::new();
reg.observe(observe_request(
"s1",
SessionEvent::PreToolUse,
Some("/tmp/a"),
));
let mut rx = reg.subscribe_changes();
reg.observe(observe_request(
"s1",
SessionEvent::PreToolUse,
Some("/tmp/a"),
));
assert!(
!rx.has_changed().unwrap(),
"a repeat sighting with no visible change must not bump"
);
reg.observe(observe_request("s1", SessionEvent::Stop, None));
assert!(rx.has_changed().unwrap(), "a state transition should bump");
rx.borrow_and_update();
reg.observe(observe_request("s1", SessionEvent::Stop, Some("/tmp/b")));
assert!(
rx.has_changed().unwrap(),
"a newly-filled `cwd` should bump"
);
}
#[test]
fn transcript_growth_does_not_clobber_a_waiting_session() {
let reg = SessionsRegistry::new();
reg.observe(observe_request(
"s1",
SessionEvent::Notification(NotificationKind::PermissionPrompt),
Some("/tmp/a"),
));
let rx = reg.subscribe_changes();
reg.observe(observe_request(
"s1",
SessionEvent::TranscriptGrew,
Some("/tmp/a"),
));
assert_eq!(reg.list()[0].state, SessionState::WaitingForPermission);
assert!(
!rx.has_changed().unwrap(),
"state did not change, so nothing a consumer renders did either (#1414)"
);
reg.observe(observe_request("s1", SessionEvent::PostToolUse, None));
assert_eq!(reg.list()[0].state, SessionState::Working);
assert!(
rx.has_changed().unwrap(),
"the release is a real transition"
);
}
#[test]
fn transcript_growth_does_not_revive_an_ended_session() {
let reg = SessionsRegistry::new();
reg.observe(observe_request(
"s1",
SessionEvent::PreToolUse,
Some("/tmp/a"),
));
assert!(reg.end("s1", Some("clear")));
let rx = reg.subscribe_changes();
reg.observe(observe_request(
"s1",
SessionEvent::TranscriptGrew,
Some("/tmp/a"),
));
assert_eq!(reg.list()[0].state, SessionState::Ended);
assert!(
!rx.has_changed().unwrap(),
"state did not change, so nothing a consumer renders did either (#1414)"
);
}
#[test]
fn end_bumps_only_for_a_known_session() {
let reg = SessionsRegistry::new();
reg.observe(observe_request("s1", SessionEvent::PreToolUse, None));
let mut rx = reg.subscribe_changes();
assert!(!reg.end("ghost", None), "an unknown session is a no-op");
assert!(
!rx.has_changed().unwrap(),
"ending an unknown session must not bump"
);
assert!(reg.end("s1", None));
assert!(rx.has_changed().unwrap(), "a real end should bump");
rx.borrow_and_update();
}
#[test]
fn window_report_bumps_only_when_it_changes_the_source_join() {
let reg = SessionsRegistry::new();
reg.report_window(window_report("w1", "/p", true));
let mut rx = reg.subscribe_changes();
reg.report_window(window_report("w1", "/p", true));
assert!(
!rx.has_changed().unwrap(),
"an unchanged window refresh must not bump"
);
reg.report_window(window_report("w1", "/p", false));
assert!(
rx.has_changed().unwrap(),
"an embedding that vanished should bump"
);
rx.borrow_and_update();
reg.report_window(window_report("w1", "/q", false));
assert!(rx.has_changed().unwrap(), "changed folders should bump");
rx.borrow_and_update();
reg.report_window(window_report("w2", "/r", true));
assert!(rx.has_changed().unwrap(), "a new window should bump");
}
#[test]
fn unregister_window_bumps_only_when_it_removes() {
let reg = SessionsRegistry::new();
reg.report_window(window_report("w1", "/p", true));
let rx = reg.subscribe_changes();
assert!(!reg.unregister_window("ghost"));
assert!(
!rx.has_changed().unwrap(),
"a no-op unregister must not bump"
);
assert!(reg.unregister_window("w1"));
assert!(
rx.has_changed().unwrap(),
"a removing unregister should bump"
);
}
#[tokio::test]
async fn a_burst_of_bumps_coalesces_into_one_wakeup() {
let reg = SessionsRegistry::new();
let mut rx = reg.subscribe_changes();
reg.observe(observe_request("s1", SessionEvent::SessionStart, None));
reg.observe(observe_request("s2", SessionEvent::SessionStart, None));
reg.observe(observe_request("s3", SessionEvent::SessionStart, None));
rx.changed().await.unwrap();
assert!(
!rx.has_changed().unwrap(),
"a burst should collapse into a single wakeup"
);
}
#[test]
fn list_does_not_bump_even_when_it_reaps() {
let reg = SessionsRegistry::new();
reg.observe(observe_request("s1", SessionEvent::PreToolUse, None));
{
let mut sessions = reg.lock_sessions();
let entry = sessions.get_mut("s1").unwrap();
entry.last_seen = Utc::now() - chrono::Duration::seconds(600);
}
let rx = reg.subscribe_changes();
assert!(reg.list().is_empty(), "the stale session should be reaped");
assert!(!rx.has_changed().unwrap(), "`list` must never bump");
}
}