use super::*;
use crate::loader::Loader;
use crate::loader::Stats;
use std::sync::mpsc::Receiver;
pub(super) enum Request {
Refresh,
RefreshLive,
Data(Box<Session>),
Delete(Box<Session>),
Terminate {
session_key: String,
pid: u32,
},
SendKeys {
pid: u32,
text: String,
},
Insight {
which: &'static str,
},
HookClaims(HashMap<String, Vec<u32>>),
Scan {
query: String,
targets: Vec<crate::session::search::Target>,
},
Shutdown,
}
pub(super) enum Response {
Discovered(Vec<Session>),
Annotated(Box<Session>),
Sessions(Box<(Vec<Session>, Stats)>),
LiveRows(Box<(Vec<Session>, Stats)>),
Data(String, Box<SessionData>),
Quota(Box<Quota>),
PricingReady,
UpdateAvailable(String),
Terminated {
session_key: String,
result: Result<(), String>,
},
Deleted {
session_key: String,
result: Result<(), String>,
},
KeysSent {
result: Result<(), String>,
},
Remote {
host: String,
snapshot: crate::fleet::Snapshot,
},
Scanned {
query: String,
hits: HashMap<String, String>,
},
Insight(String),
}
type ScanCache = HashMap<(String, String), Option<String>>;
const MAX_SCAN_CACHE: usize = 20_000;
fn scan(
cache: &mut ScanCache,
targets: &[crate::session::search::Target],
needle: &str,
) -> HashMap<String, String> {
use rayon::prelude::*;
let found: Vec<(&crate::session::search::Target, Option<String>)> = targets
.par_iter()
.map(|target| {
let memo = (!target.running)
.then(|| cache.get(&(target.key.clone(), needle.to_string())))
.flatten();
match memo {
Some(remembered) => (target, remembered.clone()),
None => (
target,
crate::session::search::find(target, needle).map(|hit| hit.snippet),
),
}
})
.collect();
if cache.len() + found.len() > MAX_SCAN_CACHE {
cache.clear();
}
let mut hits = HashMap::new();
for (target, snippet) in found {
if !target.running {
cache.insert((target.key.clone(), needle.to_string()), snippet.clone());
}
if let Some(snippet) = snippet {
hits.insert(target.key.clone(), snippet);
}
}
hits
}
pub(super) fn spawn_worker(
plan: Plan,
rx: Receiver<Request>,
tx: Sender<Response>,
) -> std::thread::JoinHandle<()> {
std::thread::spawn(move || {
let mut loader = Loader::new();
let mut sent_initial_discovery = false;
let mut live_rows: Vec<Session> = Vec::new();
let mut scans: ScanCache = HashMap::new();
while let Ok(req) = rx.recv() {
match req {
Request::Refresh => {
let first_load = !sent_initial_discovery;
sent_initial_discovery = true;
let sessions = loader.load_progressive(
plan,
first_load,
|sessions| {
if first_load {
let _ = tx.send(Response::Discovered(sessions.to_vec()));
}
},
|session| {
if first_load {
let _ = tx.send(Response::Annotated(Box::new(session.clone())));
}
},
);
let stats = crate::loader::compute_stats(&sessions);
live_rows = sessions.clone();
if tx
.send(Response::Sessions(Box::new((sessions, stats))))
.is_err()
{
break;
}
}
Request::HookClaims(claims) => {
loader.set_hook_claims(claims);
}
Request::RefreshLive => {
if live_rows.is_empty() {
continue;
}
let moved = loader.refresh_live(plan, &mut live_rows);
let stats = crate::loader::compute_stats(&live_rows);
if tx
.send(Response::LiveRows(Box::new((moved, stats))))
.is_err()
{
break;
}
}
Request::Data(session) => {
let data = loader.store().session_data_fresh(&session);
if tx
.send(Response::Data(session.key(), Box::new(data)))
.is_err()
{
break;
}
}
Request::Delete(session) => {
let result = match session.provider {
Provider::Claude => crate::session::claude::delete(&session)
.map_err(|error| error.to_string()),
Provider::Codex => crate::session::codex::delete(&session)
.map_err(|error| error.to_string()),
Provider::Cursor => crate::session::cursor::delete(&session)
.map_err(|error| error.to_string()),
Provider::Gemini => crate::session::gemini::delete(&session)
.map_err(|error| error.to_string()),
Provider::OpenCode => crate::session::opencode::delete(&session)
.map_err(|error| error.to_string()),
Provider::Pi => {
crate::session::pi::delete(&session).map_err(|error| error.to_string())
}
Provider::Windsurf => crate::session::windsurf::delete(&session)
.map_err(|error| error.to_string()),
};
if result.is_ok() {
loader.store().evict(&session);
}
if tx
.send(Response::Deleted {
session_key: session.key(),
result,
})
.is_err()
{
break;
}
}
Request::Terminate { session_key, pid } => {
let result = crate::proc::terminate(pid);
if tx
.send(Response::Terminated {
session_key,
result,
})
.is_err()
{
break;
}
}
Request::SendKeys { pid, text } => {
let result = crate::inject::send_line(pid, &text);
if tx.send(Response::KeysSent { result }).is_err() {
break;
}
}
Request::Insight { which } => {
let sessions = loader.load(plan);
let store = loader.store();
let text = loader.gently(|| {
let found = crate::insight::from_store(&sessions, store);
let all: Vec<&crate::insight::Analysis> = found.iter().collect();
match which {
"optimize" => crate::insight::optimize::report(&all),
_ => crate::insight::compare::report(&all),
}
});
if tx.send(Response::Insight(text)).is_err() {
break;
}
}
Request::Scan { query, targets } => {
let needle = query.to_ascii_lowercase();
let hits = loader.gently(|| scan(&mut scans, &targets, &needle));
if tx.send(Response::Scanned { query, hits }).is_err() {
break;
}
}
Request::Shutdown => break,
}
}
loader.store().save();
})
}