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;
const TOPICAL_WORDS: usize = 4;
const TOPICAL_FLOOR: f32 = 0.22;
const _: () = assert!(TOPICAL_FLOOR > 0.10 && TOPICAL_FLOOR < 0.35);
const TOPICAL_LIMIT: usize = 10;
fn worth_widening(no_hits: bool, needle: &str) -> bool {
no_hits || needle.split_whitespace().count() >= TOPICAL_WORDS
}
#[derive(Default)]
struct Topics {
model: Option<crate::embed::Model>,
index: Option<crate::embed::index::Index>,
absent: bool,
}
fn topical(
topics: &mut Topics,
needle: &str,
targets: &[crate::session::search::Target],
hits: &mut HashMap<String, String>,
) {
if !worth_widening(hits.is_empty(), needle) {
return;
}
if topics.absent {
return;
}
if topics.model.is_none() {
if !crate::embed::fetch::present() {
topics.absent = true;
return;
}
match crate::embed::Model::load(&crate::embed::fetch::model_dir()) {
Ok(m) => topics.model = Some(m),
Err(_) => {
topics.absent = true;
return;
}
}
}
let Some(model) = topics.model.as_ref() else {
return;
};
let index = topics.index.get_or_insert_with(|| {
crate::embed::index::Index::load(&crate::config::EMBEDDING_INDEX_FILE).unwrap_or_default()
});
if index.refresh(model, targets) > 0 {
let _ = index.save(&crate::config::EMBEDDING_INDEX_FILE);
}
let query = model.embed(&crate::embed::topic_of(needle));
for (key, snippet, score) in index.search(&query, TOPICAL_FLOOR, TOPICAL_LIMIT) {
hits.entry(key)
.or_insert_with(|| format!("~{:.0}% {snippet}", score * 100.0));
}
}
fn scan(
cache: &mut ScanCache,
targets: &[crate::session::search::Target],
needle: &str,
) -> HashMap<String, String> {
use rayon::prelude::*;
let query = crate::session::search::Query::parse(needle);
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_query(target, &query).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();
let mut topics = Topics::default();
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::Devin => crate::session::devin::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 mut hits = loader.gently(|| scan(&mut scans, &targets, &needle));
loader.gently(|| topical(&mut topics, &needle, &targets, &mut hits));
if tx.send(Response::Scanned { query, hits }).is_err() {
break;
}
}
Request::Shutdown => break,
}
}
loader.store().save();
})
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn only_an_empty_or_sentence_like_query_reaches_the_index() {
assert!(worth_widening(true, "nemotron"));
assert!(worth_widening(true, "a"));
assert!(!worth_widening(false, "nemotron"));
assert!(!worth_widening(false, "vast.ai nemotron"));
assert!(!worth_widening(false, "musl static build"));
assert!(worth_widening(false, "where did I price out GPUs"));
assert!(worth_widening(
false,
"the conversation about running a language model"
));
}
}