use std::path::Path;
use std::sync::{Arc, Mutex, PoisonError};
use anyhow::{anyhow, bail, Context, Result};
use async_trait::async_trait;
use git2::Repository;
use serde_json::{json, Value};
use tokio::task::JoinHandle;
use tokio_util::sync::CancellationToken;
use crate::daemon::service::{
DaemonService, MenuAction, MenuItem, MenuSnapshot, ServiceStatus, ServiceStream,
};
use crate::daemon::services::worktrees::focus_window;
use crate::sessions::{ObserveRequest, SessionEntry, SessionState, SessionsRegistry, WindowReport};
pub const SERVICE_NAME: &str = "sessions";
const SUBMENU_TITLE: &str = "Claude Sessions";
struct WatcherTask {
token: CancellationToken,
handle: JoinHandle<()>,
}
pub struct SessionsService {
registry: Arc<SessionsRegistry>,
watcher: Mutex<Option<WatcherTask>>,
}
impl SessionsService {
#[must_use]
pub fn new() -> Self {
Self {
registry: Arc::new(SessionsRegistry::new()),
watcher: Mutex::new(None),
}
}
pub fn start_watcher(&self) {
if tokio::runtime::Handle::try_current().is_err() {
tracing::debug!("no tokio runtime; sessions transcript watcher not started");
return;
}
let mut guard = self.watcher.lock().unwrap_or_else(PoisonError::into_inner);
if guard.is_some() {
return;
}
let token = CancellationToken::new();
let handle = crate::sessions::watcher::spawn(self.registry.clone(), token.clone());
*guard = Some(WatcherTask { token, handle });
}
#[cfg(test)]
pub(crate) fn registry(&self) -> &Arc<SessionsRegistry> {
&self.registry
}
}
impl Default for SessionsService {
fn default() -> Self {
Self::new()
}
}
#[async_trait]
impl DaemonService for SessionsService {
fn name(&self) -> &'static str {
SERVICE_NAME
}
async fn handle(&self, op: &str, payload: Value) -> Result<Value> {
match op {
"observe" => {
let mut req: ObserveRequest =
serde_json::from_value(payload).context("invalid `observe` payload")?;
if req.session_id.trim().is_empty() {
bail!("`observe` requires a non-empty `session_id`");
}
if req.repo.is_none() {
if let Some(cwd) = req.cwd.clone() {
req.repo = tokio::task::spawn_blocking(move || repo_name_for(&cwd))
.await
.unwrap_or_default();
}
}
self.registry.observe(req);
Ok(json!({ "ok": true }))
}
"end" => {
let session_id = require_str(&payload, "session_id", "end")?;
let reason = payload.get("reason").and_then(Value::as_str);
Ok(json!({ "ended": self.registry.end(session_id, reason) }))
}
"window" => {
let req: WindowReport =
serde_json::from_value(payload).context("invalid `window` payload")?;
if req.key.trim().is_empty() {
bail!("`window` requires a non-empty `key`");
}
self.registry.report_window(req);
Ok(json!({ "ok": true }))
}
"window-unregister" => {
let key = require_str(&payload, "key", "window-unregister")?;
Ok(json!({ "removed": self.registry.unregister_window(key) }))
}
"list" => Ok(sessions_payload(&self.registry)),
other => bail!("unknown sessions op: {other}"),
}
}
fn subscribe(&self, op: &str, _payload: &Value) -> Option<Box<dyn ServiceStream>> {
if op != "subscribe" {
return None;
}
Some(Box::new(SessionsStream {
registry: self.registry.clone(),
changes: self.registry.subscribe_changes(),
}))
}
fn menu(&self) -> MenuSnapshot {
MenuSnapshot {
title: SUBMENU_TITLE.to_string(),
items: menu_items_for(&self.registry.list()),
}
}
async fn menu_action(&self, action_id: &str) -> Result<()> {
if let Some(session_id) = action_id.strip_prefix("focus:") {
let folder = self.registry.focus_folder(session_id).ok_or_else(|| {
anyhow!("session {session_id} is not open in a known VS Code window")
})?;
focus_window(&folder)?;
return Ok(());
}
bail!("unknown sessions menu action: {action_id}")
}
async fn status(&self) -> ServiceStatus {
let sessions = self.registry.list();
let summary = status_summary(&sessions);
ServiceStatus {
name: SERVICE_NAME.to_string(),
healthy: true,
summary,
detail: json!({ "sessions": sessions }),
}
}
async fn shutdown(&self) {
let task = self
.watcher
.lock()
.unwrap_or_else(PoisonError::into_inner)
.take();
if let Some(task) = task {
task.token.cancel();
let _ = task.handle.await;
}
}
}
fn sessions_payload(registry: &SessionsRegistry) -> Value {
json!({ "sessions": registry.list() })
}
struct SessionsStream {
registry: Arc<SessionsRegistry>,
changes: tokio::sync::watch::Receiver<u64>,
}
#[async_trait]
impl ServiceStream for SessionsStream {
async fn changed(&mut self) {
if self.changes.changed().await.is_err() {
std::future::pending::<()>().await;
}
}
async fn snapshot(&self) -> Value {
sessions_payload(&self.registry)
}
}
fn require_str<'a>(payload: &'a Value, field: &str, op: &str) -> Result<&'a str> {
payload
.get(field)
.and_then(Value::as_str)
.ok_or_else(|| anyhow!("`{op}` requires `{field}`"))
}
fn repo_name_for(cwd: &Path) -> Option<String> {
let repo = Repository::discover(cwd).ok()?;
main_repo_name(repo.commondir())
}
fn main_repo_name(commondir: &Path) -> Option<String> {
let file_name = commondir.file_name()?.to_string_lossy().into_owned();
if file_name == ".git" {
commondir
.parent()
.and_then(Path::file_name)
.map(|n| n.to_string_lossy().into_owned())
} else {
Some(
file_name
.strip_suffix(".git")
.unwrap_or(&file_name)
.to_string(),
)
}
}
fn status_summary(sessions: &[SessionEntry]) -> String {
if sessions.is_empty() {
return "0 session(s)".to_string();
}
let mut working = 0;
let mut waiting = 0;
let mut idle = 0;
for s in sessions {
match s.state {
SessionState::Working | SessionState::Starting => working += 1,
SessionState::WaitingForInput | SessionState::WaitingForPermission => waiting += 1,
SessionState::Idle | SessionState::Ended => idle += 1,
}
}
format!(
"{} session(s): {working} working, {waiting} waiting, {idle} idle",
sessions.len()
)
}
fn state_glyph(state: SessionState) -> &'static str {
match state {
SessionState::Starting => "…",
SessionState::Working => "⚙",
SessionState::Idle => "◦",
SessionState::WaitingForInput => "?",
SessionState::WaitingForPermission => "!",
SessionState::Ended => "×",
}
}
fn display_name(entry: &SessionEntry) -> String {
if let Some(repo) = &entry.repo {
return repo.clone();
}
if let Some(cwd) = &entry.cwd {
if let Some(name) = cwd.file_name() {
return name.to_string_lossy().into_owned();
}
}
entry.session_id.chars().take(8).collect()
}
fn menu_items_for(sessions: &[SessionEntry]) -> Vec<MenuItem> {
use crate::sessions::Source;
if sessions.is_empty() {
return vec![MenuItem::Label("No active sessions".to_string())];
}
sessions
.iter()
.map(|entry| {
let label = format!(
"{} {} {}",
display_name(entry),
state_glyph(entry.state),
state_label(entry.state),
);
match &entry.source {
Source::VsCode { .. } => MenuItem::Action(MenuAction {
id: format!("focus:{}", entry.session_id),
label,
enabled: true,
}),
Source::Terminal => MenuItem::Label(label),
}
})
.collect()
}
fn state_label(state: SessionState) -> &'static str {
match state {
SessionState::Starting => "starting",
SessionState::Working => "working",
SessionState::Idle => "idle",
SessionState::WaitingForInput => "waiting",
SessionState::WaitingForPermission => "permission",
SessionState::Ended => "ended",
}
}
#[cfg(test)]
#[allow(clippy::unwrap_used, clippy::expect_used)]
mod tests {
use super::*;
use crate::sessions::{NotificationKind, SessionEvent, Source};
use std::path::PathBuf;
fn service() -> SessionsService {
SessionsService::new()
}
#[tokio::test]
async fn observe_then_list_round_trips() {
let svc = service();
let ok = svc
.handle(
"observe",
json!({ "session_id": "s1", "event": "session_start", "cwd": "/tmp/x" }),
)
.await
.unwrap();
assert_eq!(ok, json!({ "ok": true }));
let listed = svc.handle("list", Value::Null).await.unwrap();
let sessions = listed["sessions"].as_array().unwrap();
assert_eq!(sessions.len(), 1);
assert_eq!(sessions[0]["session_id"], "s1");
assert_eq!(sessions[0]["state"], "starting");
}
#[tokio::test]
async fn observe_rejects_blank_session_id() {
let svc = service();
let err = svc
.handle("observe", json!({ "session_id": " ", "event": "stop" }))
.await
.unwrap_err();
assert!(err.to_string().contains("session_id"), "{err}");
}
#[tokio::test]
async fn end_marks_ended() {
let svc = service();
svc.handle(
"observe",
json!({ "session_id": "s1", "event": "pre_tool_use" }),
)
.await
.unwrap();
let reply = svc
.handle("end", json!({ "session_id": "s1", "reason": "clear" }))
.await
.unwrap();
assert_eq!(reply, json!({ "ended": true }));
let reply = svc
.handle("end", json!({ "session_id": "ghost" }))
.await
.unwrap();
assert_eq!(reply, json!({ "ended": false }));
}
#[tokio::test]
async fn window_report_tags_source_and_unregister_removes() {
let svc = service();
svc.handle(
"observe",
json!({ "session_id": "s1", "event": "pre_tool_use", "cwd": "/home/me/proj/sub" }),
)
.await
.unwrap();
svc.handle(
"window",
json!({ "key": "w1", "folders": ["/home/me/proj"], "tabs": 1, "terminals": 0 }),
)
.await
.unwrap();
let listed = svc.handle("list", Value::Null).await.unwrap();
assert_eq!(listed["sessions"][0]["source"]["kind"], "vs_code");
assert_eq!(listed["sessions"][0]["source"]["window_key"], "w1");
let removed = svc
.handle("window-unregister", json!({ "key": "w1" }))
.await
.unwrap();
assert_eq!(removed, json!({ "removed": true }));
let listed = svc.handle("list", Value::Null).await.unwrap();
assert_eq!(listed["sessions"][0]["source"]["kind"], "terminal");
}
#[tokio::test]
async fn unknown_op_errors() {
let svc = service();
let err = svc.handle("frobnicate", Value::Null).await.unwrap_err();
assert!(err.to_string().contains("unknown sessions op"), "{err}");
}
#[tokio::test]
async fn status_summarizes_states() {
let svc = service();
svc.registry().observe(ObserveRequest {
session_id: "w".to_string(),
cwd: None,
transcript_path: None,
event: SessionEvent::PreToolUse,
repo: None,
model: None,
});
svc.registry().observe(ObserveRequest {
session_id: "p".to_string(),
cwd: None,
transcript_path: None,
event: SessionEvent::Notification(NotificationKind::PermissionPrompt),
repo: None,
model: None,
});
let status = svc.status().await;
assert!(status.healthy);
assert!(
status.summary.contains("2 session(s)"),
"{}",
status.summary
);
assert!(status.summary.contains("1 working"), "{}", status.summary);
assert!(status.summary.contains("1 waiting"), "{}", status.summary);
}
#[test]
fn menu_items_placeholder_when_empty() {
let items = menu_items_for(&[]);
assert_eq!(items.len(), 1);
assert!(matches!(&items[0], MenuItem::Label(l) if l.contains("No active")));
}
#[test]
fn menu_item_is_clickable_only_for_vscode_sessions() {
let now = chrono::Utc::now();
let base = |source: Source| SessionEntry {
session_id: "sid-12345678".to_string(),
cwd: Some(PathBuf::from("/p")),
transcript_path: None,
repo: Some("proj".to_string()),
model: None,
state: SessionState::Working,
source,
last_event: SessionEvent::PreToolUse,
started_at: now,
last_seen: now,
};
let terminal = menu_items_for(&[base(Source::Terminal)]);
assert!(
matches!(&terminal[0], MenuItem::Label(l) if l.contains("proj") && l.contains("working"))
);
let vscode = menu_items_for(&[base(Source::VsCode {
window_key: "w1".to_string(),
})]);
match &vscode[0] {
MenuItem::Action(a) => assert_eq!(a.id, "focus:sid-12345678"),
other => panic!("expected an action, got {other:?}"),
}
}
#[test]
fn repo_name_for_non_repo_is_none() {
assert_eq!(repo_name_for(Path::new("/nonexistent/xyz")), None);
}
#[test]
fn main_repo_name_handles_layouts() {
assert_eq!(
main_repo_name(Path::new("/home/me/proj/.git")).as_deref(),
Some("proj")
);
assert_eq!(
main_repo_name(Path::new("/home/me/bare.git")).as_deref(),
Some("bare")
);
}
fn entry(id: &str, state: SessionState, repo: Option<&str>, cwd: Option<&str>) -> SessionEntry {
let now = chrono::Utc::now();
SessionEntry {
session_id: id.to_string(),
cwd: cwd.map(PathBuf::from),
transcript_path: None,
repo: repo.map(str::to_string),
model: None,
state,
source: Source::Terminal,
last_event: SessionEvent::PreToolUse,
started_at: now,
last_seen: now,
}
}
#[test]
fn default_constructs_an_empty_service() {
let svc = SessionsService::default();
assert!(svc.registry().list().is_empty());
}
#[test]
fn start_watcher_is_a_noop_outside_a_runtime() {
let svc = SessionsService::new();
svc.start_watcher();
assert!(svc
.watcher
.lock()
.unwrap_or_else(PoisonError::into_inner)
.is_none());
}
#[tokio::test]
async fn start_watcher_is_idempotent_and_shutdown_stops_it() {
let svc = SessionsService::new();
svc.start_watcher();
svc.start_watcher();
assert!(svc
.watcher
.lock()
.unwrap_or_else(PoisonError::into_inner)
.is_some());
svc.shutdown().await;
assert!(svc
.watcher
.lock()
.unwrap_or_else(PoisonError::into_inner)
.is_none());
}
#[test]
fn menu_renders_every_state_and_name_fallback() {
let sessions = vec![
entry("s1", SessionState::Starting, Some("repo-a"), Some("/a")),
entry("s2", SessionState::Working, None, Some("/home/me/proj")),
entry("s3", SessionState::Idle, None, None),
entry("s4", SessionState::WaitingForInput, Some("r"), Some("/b")),
entry(
"s5",
SessionState::WaitingForPermission,
Some("r"),
Some("/c"),
),
entry("s6", SessionState::Ended, Some("r"), Some("/d")),
];
let items = menu_items_for(&sessions);
assert_eq!(items.len(), 6);
let labels: Vec<&str> = items
.iter()
.map(|i| match i {
MenuItem::Label(l) => l.as_str(),
_ => panic!("terminal sessions render as labels"),
})
.collect();
assert!(labels[0].contains("repo-a") && labels[0].contains("starting"));
assert!(labels[1].contains("proj") && labels[1].contains("working")); assert!(labels[2].contains("s3") && labels[2].contains("idle")); assert!(labels[3].contains("waiting"));
assert!(labels[4].contains("permission"));
assert!(labels[5].contains("ended"));
}
#[test]
fn menu_serves_the_snapshot_title() {
let svc = SessionsService::new();
svc.registry()
.observe(observe_req("s1", SessionEvent::Stop, None));
let snapshot = svc.menu();
assert_eq!(snapshot.title, SUBMENU_TITLE);
assert_eq!(snapshot.items.len(), 1);
}
fn observe_req(id: &str, event: SessionEvent, cwd: Option<&str>) -> ObserveRequest {
ObserveRequest {
session_id: id.to_string(),
cwd: cwd.map(PathBuf::from),
transcript_path: None,
event,
repo: None,
model: None,
}
}
#[tokio::test]
async fn menu_action_errors_on_unknown_and_missing_window() {
let svc = SessionsService::new();
let err = svc.menu_action("frobnicate").await.unwrap_err();
assert!(
err.to_string().contains("unknown sessions menu action"),
"{err}"
);
let err = svc.menu_action("focus:nope").await.unwrap_err();
assert!(
err.to_string()
.contains("not open in a known VS Code window"),
"{err}"
);
}
#[test]
fn repo_name_for_resolves_a_real_repo() {
let tmp = tempfile::tempdir().unwrap();
let repo_dir = tmp.path().join("myrepo");
std::fs::create_dir(&repo_dir).unwrap();
git2::Repository::init(&repo_dir).unwrap();
assert_eq!(repo_name_for(&repo_dir).as_deref(), Some("myrepo"));
}
#[tokio::test]
async fn status_summary_counts_idle_and_ended() {
let svc = SessionsService::new();
svc.registry()
.observe(observe_req("i", SessionEvent::Stop, None)); svc.registry().end("i2", None); svc.registry()
.observe(observe_req("i2", SessionEvent::PreToolUse, None));
svc.registry().end("i2", Some("done")); let status = svc.status().await;
assert!(status.summary.contains("2 idle"), "{}", status.summary);
}
#[tokio::test]
async fn subscribe_streams_only_for_the_subscribe_op() {
let svc = service();
assert!(svc.subscribe("subscribe", &Value::Null).is_some());
assert!(svc.subscribe("list", &Value::Null).is_none());
assert!(svc.subscribe("observe", &Value::Null).is_none());
assert!(svc.subscribe("window", &Value::Null).is_none());
assert!(svc.subscribe("bogus", &Value::Null).is_none());
}
#[tokio::test]
async fn subscribe_snapshot_matches_the_list_op() {
let svc = service();
let stream = svc
.subscribe("subscribe", &Value::Null)
.expect("subscribe stream");
assert_eq!(stream.snapshot().await, json!({ "sessions": [] }));
svc.handle(
"observe",
json!({ "session_id": "s1", "event": "pre_tool_use", "cwd": "/tmp/x" }),
)
.await
.unwrap();
let snap = stream.snapshot().await;
let listed = svc.handle("list", Value::Null).await.unwrap();
assert_eq!(snap, listed);
assert_eq!(snap["sessions"][0]["state"], json!("working"));
}
#[tokio::test]
async fn subscribe_changed_wakes_on_a_visible_change() {
let svc = service();
let mut stream = svc
.subscribe("subscribe", &Value::Null)
.expect("subscribe stream");
tokio::select! {
() = stream.changed() => panic!("changed resolved with no registry change"),
() = tokio::time::sleep(std::time::Duration::from_millis(50)) => {}
}
svc.handle(
"observe",
json!({ "session_id": "s1", "event": "session_start" }),
)
.await
.unwrap();
tokio::time::timeout(std::time::Duration::from_secs(2), stream.changed())
.await
.expect("changed should resolve after a visible change");
}
#[tokio::test]
async fn changed_parks_when_every_sender_is_gone() {
let (tx, rx) = tokio::sync::watch::channel(0u64);
let mut stream = SessionsStream {
registry: Arc::new(SessionsRegistry::new()),
changes: rx,
};
drop(tx);
tokio::select! {
() = stream.changed() => panic!("changed resolved after the sender was dropped"),
() = tokio::time::sleep(std::time::Duration::from_millis(50)) => {}
}
}
}