serve.rsannotatedserve.rssource424 lines · 18.1 KB · raw
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}