use std::collections::{HashMap, HashSet};
use std::io::{Read, Seek, SeekFrom, Write};
use std::path::PathBuf;
use serde_json::{json, Value};
use supercode_interchange::sidecar::ms_to_rfc3339;
use crate::mail_route::{LiveSession, LiveSessions};
use crate::mailbox::MailAddress;
use crate::HarnessHomes;
const LAUNCHER_PROGRAMS: &[&str] = &[
"codex",
"claude",
"node",
"bun",
"deno",
"sh",
"bash",
"zsh",
"fish",
"dash",
"nohup",
"env",
"setsid",
"timeout",
"npx",
"npm",
"caffeinate",
];
const HEARTBEAT_MS: i64 = 9 * 60 * 1000;
const JOURNAL_MAX_BYTES: u64 = 16 * 1024 * 1024;
fn now_ms() -> i64 {
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|elapsed| elapsed.as_millis() as i64)
.unwrap_or_default()
}
pub fn journal_path() -> PathBuf {
crate::teams::teams_home().join("sessions-tree.jsonl")
}
#[derive(Debug, Clone, Default)]
struct Launch {
key: String,
by: Option<String>,
at: i64,
resumes: Option<String>,
kind: Option<String>,
}
#[derive(Debug, Default)]
struct Journal {
origins: HashMap<String, Launch>,
panes: HashMap<String, Launch>,
}
fn resumed_by_key(key: &str) -> Option<String> {
let rest = key.strip_prefix("open:")?;
let rest = rest.strip_prefix("native:").unwrap_or(rest);
let (_, id) = rest.split_once(':')?;
(!id.is_empty()).then(|| id.to_string())
}
fn read_journal() -> Journal {
let mut journal = Journal::default();
let Ok(text) = std::fs::read_to_string(crate::teams::teams_home().join("launches.jsonl"))
else {
return journal;
};
let mut issued: HashMap<(String, String), i64> = HashMap::new();
for line in text.lines() {
let Ok(record) = serde_json::from_str::<Value>(line) else {
continue;
};
let (Some(key), Some(t)) = (record["key"].as_str(), record["t"].as_i64()) else {
continue;
};
let pane = record["pane"].as_str().map(str::to_string);
let at = *issued
.entry((key.to_string(), pane.clone().unwrap_or_default()))
.or_insert(t);
let resumes = record["resumes"]
.as_str()
.map(str::to_string)
.or_else(|| resumed_by_key(key));
let launch = Launch {
key: key.to_string(),
by: record["by"].as_str().map(str::to_string),
at,
resumes: resumes.clone(),
kind: record["kind"].as_str().map(str::to_string),
};
if let (Some(native), None) = (record["native_session"].as_str(), &resumes) {
journal
.origins
.entry(native.to_string())
.or_insert_with(|| launch.clone());
}
if let Some(pane) = pane {
journal.panes.insert(pane, launch);
}
}
journal
}
#[derive(Debug, Clone)]
struct Process {
ppid: u32,
program: String,
}
#[cfg(unix)]
fn processes() -> HashMap<u32, Process> {
let Ok(output) = std::process::Command::new("ps")
.args(["-axo", "pid=,ppid=,comm="])
.output()
else {
return HashMap::new();
};
String::from_utf8_lossy(&output.stdout)
.lines()
.filter_map(|line| {
let mut words = line.split_whitespace();
let pid = words.next()?.parse().ok()?;
let ppid = words.next()?.parse().ok()?;
let command = words.collect::<Vec<_>>().join(" ");
let program = command
.rsplit('/')
.next()
.unwrap_or_default()
.trim_start_matches('-')
.to_string();
Some((pid, Process { ppid, program }))
})
.collect()
}
#[cfg(not(unix))]
fn processes() -> HashMap<u32, Process> {
HashMap::new()
}
fn ancestry(pid: u32, table: &HashMap<u32, Process>) -> Vec<u32> {
let mut chain = vec![pid];
let mut current = pid;
while let Some(process) = table.get(¤t) {
if process.ppid <= 1 || chain.contains(&process.ppid) || chain.len() > 64 {
break;
}
chain.push(process.ppid);
current = process.ppid;
}
chain
}
fn orphaned(pid: u32, table: &HashMap<u32, Process>) -> bool {
let chain = ancestry(pid, table);
let top = chain.last().and_then(|top| table.get(top));
top.is_some_and(|top| top.ppid == 1)
&& chain.iter().all(|pid| {
table
.get(pid)
.is_some_and(|process| LAUNCHER_PROGRAMS.contains(&process.program.as_str()))
})
}
#[cfg(target_os = "macos")]
fn starter_environment(pid: u32) -> HashMap<String, String> {
let Ok(output) = std::process::Command::new("ps")
.args(["-E", "-o", "command=", "-p", &pid.to_string()])
.output()
else {
return HashMap::new();
};
let text = String::from_utf8_lossy(&output.stdout);
let variables = starter_variables(text.split_whitespace());
variables
}
#[cfg(target_os = "linux")]
fn starter_environment(pid: u32) -> HashMap<String, String> {
let Ok(environ) = std::fs::read(format!("/proc/{pid}/environ")) else {
return HashMap::new();
};
let variables = starter_variables(
environ
.split(|byte| *byte == 0)
.filter_map(|entry| std::str::from_utf8(entry).ok()),
);
variables
}
#[cfg(not(any(target_os = "macos", target_os = "linux")))]
fn starter_environment(_pid: u32) -> HashMap<String, String> {
HashMap::new()
}
#[cfg(any(target_os = "macos", target_os = "linux"))]
fn starter_variables<'a>(entries: impl Iterator<Item = &'a str>) -> HashMap<String, String> {
entries
.filter_map(|entry| entry.split_once('='))
.filter(|(name, value)| {
matches!(
*name,
"CLAUDE_CODE_SESSION_ID" | "CODEX_THREAD_ID" | "SUPERCODE_TEAMS_PANE"
) && !value.is_empty()
})
.map(|(name, value)| (name.to_string(), value.to_string()))
.collect()
}
fn is_daemon_pane(name: &str) -> bool {
name.strip_prefix("p_").is_some_and(|rest| {
!rest.is_empty()
&& rest
.chars()
.all(|c| c.is_ascii_alphanumeric() || c == '_' || c == '-')
})
}
#[derive(Debug, Clone)]
struct StartedBy {
parent: Option<String>,
evidence: &'static str,
recorded: bool,
detail: Option<String>,
launch_key: Option<String>,
at: Option<i64>,
}
impl StartedBy {
fn none() -> Self {
Self {
parent: None,
evidence: "none",
recorded: false,
detail: None,
launch_key: None,
at: None,
}
}
fn from_launch(launch: &Launch, detail: &str) -> Self {
Self {
parent: launch.by.clone(),
evidence: "launch",
recorded: true,
detail: Some(detail.to_string()),
launch_key: Some(launch.key.clone()),
at: Some(launch.at),
}
}
fn json(&self) -> Value {
let kind = match (&self.parent, self.evidence) {
(Some(_), _) => "session",
(None, "launch" | "pane_launch") => "person",
_ => "unrecorded",
};
let mut value = json!({"kind": kind, "evidence": self.evidence, "recorded": self.recorded});
if let Some(detail) = &self.detail {
value["detail"] = json!(detail);
}
if let Some(key) = &self.launch_key {
value["launch_key"] = json!(key);
}
if let Some(at) = self.at {
value["at"] = json!(ms_to_rfc3339(at));
}
value
}
}
pub fn live_tree(homes: &HarnessHomes, ask_other_machines: bool) -> Value {
let roots_read = std::thread::spawn(crate::mail_route::daemon_pane_roots);
let codex_read = std::thread::spawn(crate::codex_peer::session_processes);
let live = LiveSessions::read(homes);
let pane_roots = roots_read.join().unwrap_or_default();
let codex_pids = codex_read.join().unwrap_or_default();
let table = processes();
let journal = read_journal();
let machine = crate::mailbox::local_machine_name();
let pid_of = |session: &LiveSession| -> Option<u32> {
session.pid.or_else(|| {
(session.address.harness == "codex")
.then(|| codex_pids.get(&session.address.session_id).copied())
.flatten()
})
};
let root_panes: HashMap<u32, &String> = pane_roots
.iter()
.map(|(pane, root)| (*root, pane))
.collect();
let pane_of = |session: &LiveSession| -> Option<String> {
if let Some(pane) = session
.tmux
.as_deref()
.and_then(|tmux| tmux.split(':').next())
.filter(|name| is_daemon_pane(name))
{
return Some(pane.to_string());
}
let pid = pid_of(session)?;
ancestry(pid, &table)
.iter()
.find_map(|ancestor| root_panes.get(ancestor).map(|pane| pane.to_string()))
};
let sessions = live.all();
let pids: Vec<Option<u32>> = sessions.iter().map(pid_of).collect();
let panes: Vec<Option<String>> = sessions.iter().map(pane_of).collect();
let by_pid: HashMap<u32, String> = sessions
.iter()
.zip(&pids)
.filter_map(|(session, pid)| pid.map(|pid| (pid, session.address.to_string())))
.collect();
let process_parents: Vec<Option<String>> = pids
.iter()
.map(|pid| {
ancestry((*pid)?, &table)
.into_iter()
.skip(1)
.find_map(|ancestor| by_pid.get(&ancestor).cloned())
})
.collect();
let by_pane: HashMap<&str, String> = sessions
.iter()
.zip(&panes)
.zip(&process_parents)
.filter(|(_, parent)| parent.is_none())
.filter_map(|((session, pane), _)| {
pane.as_deref()
.map(|pane| (pane, session.address.to_string()))
})
.collect();
let address_of = |harness: &str, id: &str| {
MailAddress::new(&machine, harness, id)
.ok()
.map(|a| a.to_string())
};
let mut rows: Vec<Value> = Vec::new();
for (((session, pid), pane), process_parent) in
sessions.iter().zip(&pids).zip(&panes).zip(&process_parents)
{
let id = session.address.session_id.as_str();
let address = session.address.to_string();
let pane_launch = pane.as_deref().and_then(|pane| journal.panes.get(pane));
let resume = pane_launch.filter(|launch| launch.resumes.as_deref() == Some(id));
let started = if let Some(origin) = journal.origins.get(id) {
StartedBy::from_launch(origin, "its launch recorded this session")
} else if let Some(first) = resume.and_then(|_| first_recorded_start(&address)) {
first
} else if let Some(parent) = process_parent {
StartedBy {
parent: Some(parent.clone()),
evidence: "process_parent",
recorded: false,
detail: Some("a parent process of it is that session's".into()),
launch_key: None,
at: None,
}
} else if let Some(launch) = pane_launch.filter(|_| resume.is_none()) {
let root = pane
.as_deref()
.and_then(|pane| pane_roots.get(pane))
.copied();
let launched = pid.is_some_and(|pid| {
Some(pid) == root
|| (table.get(&pid).map(|process| process.ppid) == root
&& launch.kind.as_deref() == Some(session.address.harness.as_str()))
});
let mut started = StartedBy::from_launch(launch, "its pane's launch ran it");
if !launched {
started.evidence = "pane_launch";
started.recorded = false;
started.detail = Some(
"it runs in that launch's pane; the launch ran another program there".into(),
);
}
started
} else if let Some(started) = pid.filter(|pid| orphaned(*pid, &table)).and_then(|pid| {
environment_starter(
pid,
&session.address,
pane.as_deref(),
&by_pane,
&by_pid,
&address_of,
)
}) {
started
} else {
StartedBy::none()
};
let mut row = json!({
"kind": "session",
"address": address,
"harness": session.address.harness,
"session_id": id,
"name": session.name.rsplit_once('@').map_or(session.name.as_str(), |(name, _)| name),
"folder": session.cwd,
"state": session.status,
"pane": pane,
"pid": pid,
"parent": started.parent,
"started_by": started.json(),
"resumed_by": Value::Null,
});
if let Some(resume) = resume {
row["resumed_by"] =
json!({"by": resume.by, "launch_key": resume.key, "at": ms_to_rfc3339(resume.at)});
}
rows.push(row);
}
for thread in live.codex_threads() {
let short: String = thread.address.session_id.chars().take(8).collect();
rows.push(json!({
"kind": "codex_thread",
"address": thread.address.to_string(),
"harness": "codex",
"session_id": thread.address.session_id,
"name": format!("codex-{short}"),
"folder": thread.cwd,
"state": thread.status,
"pane": Value::Null,
"pid": Value::Null,
"parent": thread.parent.to_string(),
"started_by": {"kind": "session", "evidence": "codex_thread", "recorded": true,
"detail": "Codex recorded it as a subagent thread of that conversation"},
"resumed_by": Value::Null,
}));
}
let tree = tree_of(rows, &machine, ask_other_machines);
json!({"machine": machine, "at": ms_to_rfc3339(now_ms()), "sessions": tree.0, "parents_elsewhere": tree.1})
}
fn environment_starter(
pid: u32,
own: &MailAddress,
own_pane: Option<&str>,
by_pane: &HashMap<&str, String>,
by_pid: &HashMap<u32, String>,
address_of: &dyn Fn(&str, &str) -> Option<String>,
) -> Option<StartedBy> {
let environment = starter_environment(pid);
let mut named: Vec<(String, &'static str)> = Vec::new();
for (variable, harness) in [
("CLAUDE_CODE_SESSION_ID", "claude-code"),
("CODEX_THREAD_ID", "codex"),
] {
if let Some(id) = environment
.get(variable)
.filter(|id| **id != own.session_id)
{
if let Some(address) = address_of(harness, id) {
named.push((address, variable));
}
}
}
if let Some(address) = environment
.get("SUPERCODE_TEAMS_PANE")
.filter(|pane| Some(pane.as_str()) != own_pane)
.and_then(|pane| by_pane.get(pane.as_str()))
.filter(|address| **address != own.to_string())
{
if !named.iter().any(|(named, _)| named == address) {
named.push((address.clone(), "SUPERCODE_TEAMS_PANE"));
}
}
let pid_of: HashMap<&String, u32> = by_pid
.iter()
.map(|(pid, address)| (address, *pid))
.collect();
let nearest = if named.len() > 1 {
named.iter().find(|(address, _)| {
pid_of.get(address).is_some_and(|pid| {
let theirs = starter_environment(*pid);
let theirs: HashSet<String> = theirs
.iter()
.filter_map(|(variable, id)| match variable.as_str() {
"CLAUDE_CODE_SESSION_ID" => address_of("claude-code", id),
"CODEX_THREAD_ID" => address_of("codex", id),
_ => None,
})
.collect();
named
.iter()
.filter(|(other, _)| other != address)
.all(|(other, _)| theirs.contains(other))
})
})
} else {
named.first()
};
let (parent, variable) = nearest?.clone();
Some(StartedBy {
parent: Some(parent),
evidence: "process_environment",
recorded: false,
detail: Some(format!(
"its parent process has exited; its environment's {variable} names that session"
)),
launch_key: None,
at: None,
})
}
fn tree_of(
mut rows: Vec<Value>,
machine: &str,
ask_other_machines: bool,
) -> (Vec<Value>, Vec<Value>) {
let index: HashMap<String, usize> = rows
.iter()
.enumerate()
.filter_map(|(i, row)| Some((row["address"].as_str()?.to_string(), i)))
.collect();
for i in 0..rows.len() {
let mut seen = HashSet::from([i]);
let mut current = rows[i]["parent"]
.as_str()
.and_then(|parent| index.get(parent))
.copied();
while let Some(next) = current {
if next == i {
rows[i]["parent"] = Value::Null;
rows[i]["started_by"]["cycle"] = json!(true);
break;
}
if !seen.insert(next) {
break;
}
current = rows[next]["parent"]
.as_str()
.and_then(|parent| index.get(parent))
.copied();
}
}
let mut children: HashMap<String, Vec<String>> = HashMap::new();
let mut elsewhere: Vec<String> = Vec::new();
for row in &rows {
if let Some(parent) = row["parent"].as_str() {
children
.entry(parent.to_string())
.or_default()
.push(row["address"].as_str().unwrap_or_default().to_string());
if !index.contains_key(parent) && !elsewhere.iter().any(|known| known == parent) {
elsewhere.push(parent.to_string());
}
}
}
for row in &mut rows {
let address = row["address"].as_str().unwrap_or_default().to_string();
row["children"] = json!(children.remove(&address).unwrap_or_default());
}
let machines: HashSet<String> = elsewhere
.iter()
.filter_map(|address| MailAddress::parse(address).ok())
.map(|address| address.machine)
.filter(|other| other != machine && ask_other_machines)
.collect();
let asking: Vec<_> = machines
.into_iter()
.map(|other| {
std::thread::spawn(move || {
let listed = crate::mailbox::teams_mail(&other, &json!({"op": "list"}))
.ok()
.and_then(|answer| answer["sessions"].as_array().cloned());
(other, listed)
})
})
.collect();
let asked: Vec<(String, Option<Vec<Value>>)> = asking
.into_iter()
.filter_map(|handle| handle.join().ok())
.collect();
let elsewhere = elsewhere
.into_iter()
.map(|address| {
let parsed = MailAddress::parse(&address).ok();
let listing = parsed.as_ref().and_then(|parsed| {
asked
.iter()
.find(|(other, _)| *other == parsed.machine)
.and_then(|(_, listed)| listed.as_ref())
});
let row =
listing.and_then(|rows| rows.iter().find(|row| row["address"] == address.as_str()));
let name = row
.and_then(|row| row["name"].as_str())
.map(|name| {
name.rsplit_once('@')
.map_or(name, |(name, _)| name)
.to_string()
})
.or_else(|| parsed.as_ref().and_then(crate::mail_route::remembered_name));
let here = parsed.as_ref().is_some_and(|a| a.machine == machine);
json!({
"address": address,
"machine": parsed.as_ref().map(|a| a.machine.clone()),
"name": name,
"here": here,
"running": if here { Some(false) } else { listing.map(|_| row.is_some()) },
"state": row.and_then(|row| row["status"].as_str()),
"children": children.get(&address).cloned().unwrap_or_default(),
})
})
.collect();
(rows, elsewhere)
}
fn first_recorded_start(address: &str) -> Option<StartedBy> {
let path = journal_path();
for file in [path.with_extension("jsonl.1"), path] {
let Ok(text) = std::fs::read_to_string(&file) else {
continue;
};
for line in text.lines().filter(|line| line.contains(address)) {
let Ok(record) = serde_json::from_str::<Value>(line) else {
continue;
};
let found = record["sessions"]
.as_array()
.into_iter()
.flatten()
.find(|row| {
row["address"] == address
&& row["started_by"]["recorded"] == true
&& row["started_by"]["evidence"] == "launch"
&& row["resumed_by"].is_null()
});
if let Some(row) = found {
let started = &row["started_by"];
return Some(StartedBy {
parent: row["parent"].as_str().map(str::to_string),
evidence: "launch",
recorded: true,
detail: Some(format!(
"its first launch, as the session tree recorded it at {}",
record["t"].as_i64().map(ms_to_rfc3339).unwrap_or_default()
)),
launch_key: started["launch_key"].as_str().map(str::to_string),
at: started["at"]
.as_str()
.and_then(supercode_interchange::sidecar::rfc3339_to_ms),
});
}
}
}
None
}
fn shape(tree: &Value) -> Vec<(String, String, String)> {
let mut shape: Vec<(String, String, String)> = tree["sessions"]
.as_array()
.into_iter()
.flatten()
.map(|row| {
let text = |key: &str| row[key].as_str().unwrap_or_default().to_string();
(text("address"), text("parent"), text("pane"))
})
.collect();
shape.sort();
shape
}
fn journal_tail() -> (Option<Value>, Option<i64>) {
let Ok(mut file) = std::fs::File::open(journal_path()) else {
return (None, None);
};
let length = file.metadata().map(|meta| meta.len()).unwrap_or_default();
let start = length.saturating_sub(4 * 1024 * 1024);
if file.seek(SeekFrom::Start(start)).is_err() {
return (None, None);
}
let mut text = String::new();
if file.read_to_string(&mut text).is_err() {
return (None, None);
}
let mut last_at = None;
for line in text.lines().rev() {
let Ok(record) = serde_json::from_str::<Value>(line) else {
continue;
};
last_at = last_at.or_else(|| record["t"].as_i64());
if record.get("sessions").is_some() {
return (Some(record), last_at);
}
}
(None, last_at)
}
pub fn record(tree: &Value) -> std::io::Result<&'static str> {
let path = journal_path();
if let Some(parent) = path.parent() {
std::fs::create_dir_all(parent)?;
}
if std::fs::metadata(&path).is_ok_and(|meta| meta.len() > JOURNAL_MAX_BYTES) {
std::fs::rename(&path, path.with_extension("jsonl.1"))?;
}
let now = now_ms();
let (last, last_at) = journal_tail();
let line = match &last {
Some(last) if shape(last) == shape(tree) => {
if last_at.is_some_and(|at| now - at < HEARTBEAT_MS) {
return Ok("nothing");
}
json!({"t": now, "unchanged_since": last["t"]})
}
_ => json!({"t": now, "machine": tree["machine"], "sessions": tree["sessions"],
"parents_elsewhere": tree["parents_elsewhere"]}),
};
let written = if line.get("sessions").is_some() {
"tree"
} else {
"heartbeat"
};
let mut file = std::fs::OpenOptions::new()
.create(true)
.append(true)
.open(&path)?;
file.write_all(format!("{line}\n").as_bytes())?;
Ok(written)
}
pub fn recorded_tree(at_ms: i64) -> Option<Value> {
let path = journal_path();
let mut text = std::fs::read_to_string(path.with_extension("jsonl.1")).unwrap_or_default();
text.push_str(&std::fs::read_to_string(&path).unwrap_or_default());
let mut found: Option<Value> = None;
let mut confirmed: Option<i64> = None;
let mut next_change: Option<i64> = None;
for line in text.lines() {
let Ok(record) = serde_json::from_str::<Value>(line) else {
continue;
};
let Some(t) = record["t"].as_i64() else {
continue;
};
let full = record.get("sessions").is_some();
if t <= at_ms {
if full {
found = Some(record);
}
confirmed = Some(t);
} else if full {
next_change = Some(t);
break;
} else if found.is_some() && next_change.is_none() {
confirmed = Some(t);
}
}
let mut tree = found?;
let recorded_at = tree["t"].as_i64().unwrap_or_default();
tree["at"] = json!(ms_to_rfc3339(at_ms));
tree["recorded_at"] = json!(ms_to_rfc3339(recorded_at));
tree["confirmed_through"] = json!(confirmed.map(ms_to_rfc3339));
tree["next_change_at"] = json!(next_change.map(ms_to_rfc3339));
if let Some(object) = tree.as_object_mut() {
object.remove("t");
}
Some(tree)
}