1//! The daemon: one process for every Claude Code session on the machine, 2//! serving HTTP on a Unix socket. It holds the Jev client, so the connection 3//! is opened once and stays warm, and the key is fetched once rather than at 4//! the start of every session. 5//! 6//! `POST /decide` a deciding event; answers with a `Decision` 7//! `POST /observe` an event the daemon needs but does not answer 8//! `GET /status` whether Jev can be asked, and the spend ledger's line 9//! `GET /recent` the last decisions, oldest first (`?count=20`) 10//! `POST /shutdown` stop, so a newer build can take the socket 11 12use std::collections::{HashMap, HashSet, VecDeque}; 13use std::convert::Infallible; 14use std::sync::atomic::{AtomicU64, Ordering}; 15use std::sync::{Arc, Mutex, PoisonError}; 16use std::time::Instant; 17 18use bytes::Bytes; 19use http_body_util::{BodyExt, Full}; 20use hyper::body::Incoming; 21use hyper::service::service_fn; 22use hyper::{Method, Request, Response, StatusCode}; 23use hyper_util::rt::TokioIo; 24use jev_http::{Guards, Jev}; 25use jev_protocol::ModelId; 26use jevhooks_events::{Decision, HookEvent, Verdict}; 27use serde::Deserialize; 28use serde_json::{Value, json}; 29use tokio::net::UnixListener; 30use tokio::sync::Notify; 31 32use crate::log::{DecisionLog, DecisionRecord, OutcomeRecord, Record, now}; 33use crate::headroom::{Config, Headroom}; 34use crate::{judge, paths}; 35 36/// Pinned, so an answer means the same thing from one week to the next; 37/// `jev-latest` pointed here on 2026-10-01. 38const MODEL: &str = "jev-1.13.0"; 39 40/// What the ledger records as the spender. 41const BINARY: &str = "jevhooks"; 42 43/// How many of a session's prompts a Stop is judged against. 44const REQUESTS_KEPT: usize = 5; 45 46/// How much of what was judged the decision log keeps. 47const SUBJECT_CHARS: usize = 500; 48 49/// Whether this process can ask Jev, decided once at start. 50enum Client { 51 Ready(Arc<Jev>), 52 /// No `TYPESAFE_API_KEY` in the environment the daemon was started with. 53 NoKey, 54 /// A key, but the client could not be built (the ledger, usually). 55 Broken(String), 56} 57 58/// What the daemon remembers about one Claude Code session. 59#[derive(Default)] 60struct Session { 61 /// The user's latest prompts, oldest first. 62 requests: VecDeque<String>, 63 /// What this session's questions have cost. 64 usd: f64, 65 /// Tool calls that were judged, so their outcomes are logged. 66 judged: HashSet<String>, 67} 68 69struct Daemon { 70 client: Client, 71 identity: String, 72 started: Instant, 73 log: DecisionLog, 74 headroom: Headroom, 75 /// Why the config file was not used, when it could not be. 76 config_error: Option<String>, 77 sessions: Mutex<HashMap<String, Session>>, 78 decided: AtomicU64, 79 observed: AtomicU64, 80 shutdown: Notify, 81} 82 83/// One event as the mod sends it. 84#[derive(Deserialize)] 85struct EventRequest { 86 session: String, 87 event: HookEvent, 88 /// The event's own fields: `tool`, `command` and `tool_use_id` for a tool 89 /// call; `prompt` for a submitted prompt; `last_assistant_message` and 90 /// the rest for a Stop. 91 #[serde(default)] 92 input: Value, 93 cwd: Option<String>, 94 root: Option<String>, 95} 96 97type Reply = Response<Full<Bytes>>; 98 99fn reply(status: StatusCode, body: &Value) -> Reply { 100 let mut response = Response::new(Full::new(Bytes::from(body.to_string()))); 101 *response.status_mut() = status; 102 response.headers_mut().insert(hyper::header::CONTENT_TYPE, hyper::header::HeaderValue::from_static("application/json")); 103 response 104} 105 106fn refuse(status: StatusCode, why: impl std::fmt::Display) -> Reply { 107 reply(status, &json!({ "error": why.to_string() })) 108} 109 110impl Daemon { 111 fn sessions(&self) -> std::sync::MutexGuard<'_, HashMap<String, Session>> { 112 self.sessions.lock().unwrap_or_else(PoisonError::into_inner) 113 } 114 115 fn session_usd(&self, session: &str) -> f64 { 116 self.sessions().get(session).map_or(0.0, |s| s.usd) 117 } 118 119 /// A verdict that needed no question, or could not get an answer. 120 fn unasked(&self, request: &EventRequest, started: Instant, line: String) -> Decision { 121 Decision { 122 verdict: Verdict::Pass, 123 line, 124 ms: started.elapsed().as_millis() as u64, 125 session_usd: self.session_usd(&request.session), 126 } 127 } 128 129 /// The Jev client, or why there is none to ask. 130 fn jev(&self) -> Result<Arc<Jev>, String> { 131 match &self.client { 132 Client::Ready(jev) => Ok(Arc::clone(jev)), 133 Client::NoKey => Err(format!("no {} where the daemon was started", jev_http::KEY)), 134 Client::Broken(why) => Err(format!("the Jev client could not be built: {why}")), 135 } 136 } 137 138 /// Books a judgment against its session and writes it to the log. 139 fn settle( 140 &self, 141 request: &EventRequest, 142 started: Instant, 143 kind: &'static str, 144 subject: &str, 145 tool_use_id: Option<&str>, 146 judged: Result<judge::Judged, String>, 147 ) -> Decision { 148 let (decision, judged) = match judged { 149 Ok(judged) => { 150 let session_usd = { 151 let mut sessions = self.sessions(); 152 let session = sessions.entry(request.session.clone()).or_default(); 153 session.usd += judged.meta.usage.dollars(); 154 if let Some(id) = tool_use_id { 155 session.judged.insert(id.to_owned()); 156 } 157 session.usd 158 }; 159 let decision = Decision { 160 verdict: judged.verdict, 161 line: judged.line.clone(), 162 ms: started.elapsed().as_millis() as u64, 163 session_usd, 164 }; 165 (decision, Some(judged)) 166 } 167 Err(why) => (self.unasked(request, started, format!("Jev not asked: {why}")), None), 168 }; 169 let meta = judged.as_ref().map(|judged| &judged.meta); 170 self.log.append(&Record::Decision(Box::new(DecisionRecord { 171 time: now(), 172 session: request.session.clone(), 173 kind: kind.to_owned(), 174 subject: subject.to_owned(), 175 tool_use_id: tool_use_id.map(str::to_owned), 176 verdict: decision.verdict, 177 line: decision.line.clone(), 178 answers: judged.as_ref().map(|judged| judged.answers.clone()), 179 ms: decision.ms, 180 jev_ms: meta.map(|m| m.jev_ms), 181 input_tokens: meta.map(|m| m.usage.input_tokens), 182 usd: meta.map(|m| m.usage.dollars()), 183 attempts: meta.map(|m| m.attempts), 184 request_id: meta.and_then(|m| m.request_id.clone()), 185 }))); 186 decision 187 } 188 189 async fn decide(&self, request: EventRequest) -> Result<Decision, String> { 190 let started = Instant::now(); 191 self.decided.fetch_add(1, Ordering::Relaxed); 192 let text = |field: &str| request.input.get(field).and_then(Value::as_str); 193 match request.event { 194 HookEvent::PreToolUse => { 195 let (Some("Bash"), Some(command)) = (text("tool"), text("command")) else { 196 return Ok(self.unasked(&request, started, "only Bash commands are judged".into())); 197 }; 198 let judged = match self.jev() { 199 Ok(jev) => { 200 let available = self.headroom.available_mb().await; 201 judge::bash(jev, command, request.cwd.as_deref(), request.root.as_deref(), available).await 202 } 203 Err(why) => Err(why), 204 }; 205 // A command is recognised by how it starts. 206 Ok(self.settle(&request, started, "bash", judge::head(command, SUBJECT_CHARS), text("tool_use_id"), judged)) 207 } 208 HookEvent::Stop => { 209 // A Stop that follows this plugin's own refusal is let through, 210 // so one refusal can never become a loop. 211 if request.input.get("stop_hook_active").and_then(Value::as_bool) == Some(true) { 212 return Ok(self.unasked(&request, started, "already refused once this turn".into())); 213 } 214 // Work still running in the background will wake the session; 215 // the turn has paused, not ended. 216 let waiting = |field: &str| request.input.get(field).and_then(Value::as_array).is_some_and(|a| !a.is_empty()); 217 if waiting("background_tasks") || waiting("session_crons") { 218 return Ok(self.unasked(&request, started, "background work is still running".into())); 219 } 220 let requests: Vec<String> = 221 self.sessions().get(&request.session).map(|s| s.requests.iter().cloned().collect()).unwrap_or_default(); 222 let (Some(last), false) = (text("last_assistant_message"), requests.is_empty()) else { 223 return Ok(self.unasked(&request, started, "no prompt or final message to judge".into())); 224 }; 225 let judged = match self.jev() { 226 Ok(jev) => judge::stop(jev, &requests, last).await, 227 Err(why) => Err(why), 228 }; 229 // A final message is recognised by how it ends. 230 Ok(self.settle(&request, started, "stop", judge::tail(last, SUBJECT_CHARS), None, judged)) 231 } 232 other => Err(format!("{other:?} is not a deciding event")), 233 } 234 } 235 236 fn observe(&self, request: EventRequest) -> Result<(), String> { 237 self.observed.fetch_add(1, Ordering::Relaxed); 238 let text = |field: &str| request.input.get(field).and_then(Value::as_str); 239 match request.event { 240 HookEvent::UserPromptSubmit => { 241 // Only what a person asked for is a request; wakeups and task 242 // notifications are the session talking to itself. 243 let from_person = matches!(text("source"), None | Some("user" | "sdk")); 244 if let (true, Some(prompt)) = (from_person, text("prompt")) { 245 let mut sessions = self.sessions(); 246 let session = sessions.entry(request.session).or_default(); 247 session.requests.push_back(prompt.to_owned()); 248 while session.requests.len() > REQUESTS_KEPT { 249 session.requests.pop_front(); 250 } 251 } 252 } 253 HookEvent::PostToolUse | HookEvent::PostToolUseFailure | HookEvent::PermissionDenied => { 254 let Some(id) = text("tool_use_id") else { return Ok(()) }; 255 let was_judged = self.sessions().get_mut(&request.session).is_some_and(|s| s.judged.remove(id)); 256 if was_judged { 257 let outcome = match request.event { 258 HookEvent::PostToolUse => "ran", 259 HookEvent::PostToolUseFailure => "failed", 260 _ => "denied", 261 }; 262 self.log.append(&Record::Outcome(OutcomeRecord { 263 time: now(), 264 session: request.session, 265 tool_use_id: id.to_owned(), 266 outcome: outcome.to_owned(), 267 })); 268 } 269 } 270 HookEvent::SessionEnd => { 271 self.sessions().remove(&request.session); 272 } 273 other => return Err(format!("{other:?} is {:?}, not observed", other.role())), 274 } 275 Ok(()) 276 } 277 278 async fn status(&self) -> Value { 279 let (jev, detail) = match &self.client { 280 Client::Ready(jev) => match jev.status() { 281 Ok(ledger) => ("ready", json!({ "line": ledger.line(), "ledger": ledger })), 282 Err(why) => ("ready", json!({ "line": format!("the spend ledger: {why}") })), 283 }, 284 Client::NoKey => ("no key", json!({ "line": format!("{} was not set where the daemon was started", jev_http::KEY) })), 285 Client::Broken(why) => ("broken", json!({ "line": why })), 286 }; 287 json!({ 288 "identity": self.identity, 289 "pid": std::process::id(), 290 "uptime_seconds": self.started.elapsed().as_secs(), 291 "jev": jev, 292 "model": MODEL, 293 "detail": detail, 294 "decided": self.decided.load(Ordering::Relaxed), 295 "observed": self.observed.load(Ordering::Relaxed), 296 "sessions": self.sessions().len(), 297 "decisions_log": self.log.path().display().to_string(), 298 "headroom": self.headroom.status().await, 299 "config": match &self.config_error { 300 Some(why) => json!({ "path": crate::headroom::config_path().display().to_string(), "error": why }), 301 None => json!({ "path": crate::headroom::config_path().display().to_string() }), 302 }, 303 "thresholds": { 304 "bash_ask_at_consequential": judge::BASH_ASK_AT_CONSEQUENTIAL, 305 "bash_ask_at_undo": judge::BASH_ASK_AT_UNDO, 306 "bash_allow_at_ordinary": judge::BASH_ALLOW_AT_ORDINARY, 307 "bash_allow_undo_at_most": judge::BASH_ALLOW_UNDO_AT_MOST, 308 "load_needs_mb": judge::LOAD_NEEDS_MB, 309 "stop_block_at_least": judge::STOP_BLOCK_AT_LEAST, 310 "answer_within_ms": judge::ANSWER_WITHIN.as_millis() as u64, 311 }, 312 }) 313 } 314 315 async fn route(self: Arc<Self>, request: Request<Incoming>) -> Reply { 316 let method = request.method().clone(); 317 let path = request.uri().path().to_owned(); 318 let query = request.uri().query().unwrap_or("").to_owned(); 319 let body = match request.into_body().collect().await { 320 Ok(body) => body.to_bytes(), 321 Err(e) => return refuse(StatusCode::BAD_REQUEST, format!("reading the body: {e}")), 322 }; 323 let event = || serde_json::from_slice::<EventRequest>(&body).map_err(|e| e.to_string()); 324 match (method, path.as_str()) { 325 (Method::POST, "/decide") => match event() { 326 Ok(request) => match self.decide(request).await { 327 Ok(decision) => reply(StatusCode::OK, &json!(decision)), 328 Err(why) => refuse(StatusCode::BAD_REQUEST, why), 329 }, 330 Err(why) => refuse(StatusCode::BAD_REQUEST, why), 331 }, 332 (Method::POST, "/observe") => match event().and_then(|request| self.observe(request)) { 333 Ok(()) => reply(StatusCode::OK, &json!({})), 334 Err(why) => refuse(StatusCode::BAD_REQUEST, why), 335 }, 336 (Method::GET, "/status") => reply(StatusCode::OK, &self.status().await), 337 (Method::GET, "/recent") => { 338 let count = query.strip_prefix("count=").and_then(|n| n.parse().ok()).unwrap_or(20); 339 match self.log.recent(count) { 340 Ok(records) => reply(StatusCode::OK, &json!(records)), 341 Err(e) => refuse(StatusCode::INTERNAL_SERVER_ERROR, format!("reading {}: {e}", self.log.path().display())), 342 } 343 } 344 (Method::POST, "/shutdown") => { 345 self.shutdown.notify_one(); 346 reply(StatusCode::OK, &json!({ "stopping": std::process::id() })) 347 } 348 _ => refuse(StatusCode::NOT_FOUND, format!("no route {path}")), 349 } 350 } 351} 352 353/// Runs the daemon until `/shutdown`. Returns at once, successfully, when 354/// another daemon already holds the lock. 355pub async fn serve() -> Result<(), Box<dyn std::error::Error>> { 356 std::fs::create_dir_all(paths::state_dir())?; 357 let lock = std::fs::OpenOptions::new().create(true).truncate(false).write(true).open(paths::lock())?; 358 if lock.try_lock().is_err() { 359 eprintln!("jevhooks: another daemon holds {}", paths::lock().display()); 360 return Ok(()); 361 } 362 // The lock is held, so a socket file left behind has no daemon behind it. 363 let socket = paths::socket(); 364 match std::fs::remove_file(&socket) { 365 Err(e) if e.kind() != std::io::ErrorKind::NotFound => return Err(e.into()), 366 _ => {} 367 } 368 let listener = UnixListener::bind(&socket)?; 369 370 let client = match ModelId::pinned(MODEL).map_err(|e| e.to_string()).map(|model| Jev::from_env(BINARY, model)) { 371 Ok(Some(Ok(jev))) => Client::Ready(Arc::new(jev.with_guards(Guards::from_env()))), 372 Ok(Some(Err(why))) | Err(why) => Client::Broken(why), 373 Ok(None) => Client::NoKey, 374 }; 375 // A config file that cannot be used is reported in `status`, and the 376 // daemon runs on the defaults: a typo must not switch the gate off. 377 let (config, config_error) = match Config::load() { 378 Ok(config) => (config, None), 379 Err(why) => { 380 eprintln!("jevhooks: config not used: {why}"); 381 (Config::default(), Some(why)) 382 } 383 }; 384 let daemon = Arc::new(Daemon { 385 client, 386 identity: paths::identity(), 387 started: Instant::now(), 388 log: DecisionLog::open(paths::decisions())?, 389 headroom: Headroom::new(&config), 390 config_error, 391 sessions: Mutex::default(), 392 decided: AtomicU64::new(0), 393 observed: AtomicU64::new(0), 394 shutdown: Notify::new(), 395 }); 396 eprintln!("jevhooks: serving on {} as {}", socket.display(), daemon.identity); 397 398 loop { 399 tokio::select! { 400 accepted = listener.accept() => { 401 let (stream, _) = accepted?; 402 let daemon = Arc::clone(&daemon); 403 tokio::spawn(async move { 404 let service = service_fn(move |request| { 405 let daemon = Arc::clone(&daemon); 406 async move { Ok::<_, Infallible>(daemon.route(request).await) } 407 }); 408 if let Err(e) = hyper::server::conn::http1::Builder::new().serve_connection(TokioIo::new(stream), service).await { 409 eprintln!("jevhooks: a connection ended badly: {e}"); 410 } 411 }); 412 } 413 () = daemon.shutdown.notified() => { 414 // Long enough for the reply to the shutdown request to leave. 415 tokio::time::sleep(std::time::Duration::from_millis(100)).await; 416 break; 417 } 418 } 419 } 420 drop(listener); 421 std::fs::remove_file(&socket).ok(); 422 eprintln!("jevhooks: stopped"); 423 Ok(()) 424}