serve.rsannotatedserve.rssource424 lines · 18.1 KB · raw

The daemon: one process for every Claude Code session on the machine, serving HTTP on a Unix socket. It holds the Jev client, so the connection is opened once and stays warm, and the key is fetched once rather than at the start of every session.

POST /decide a deciding event; answers with a Decision POST /observe an event the daemon needs but does not answer GET /status whether Jev can be asked, and the spend ledger's line GET /recent the last decisions, oldest first (?count=20) POST /shutdown stop, so a newer build can take the socket

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;
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};

Pinned, so an answer means the same thing from one week to the next; jev-latest pointed here on 2026-10-01.

38const MODEL: &str = "jev-1.13.0";

What the ledger records as the spender.

41const BINARY: &str = "jevhooks";

How many of a session's prompts a Stop is judged against.

44const REQUESTS_KEPT: usize = 5;

How much of what was judged the decision log keeps.

47const SUBJECT_CHARS: usize = 500;

Whether this process can ask Jev, decided once at start.

50enum Client {
51    Ready(Arc<Jev>),

No TYPESAFE_API_KEY in the environment the daemon was started with.

53    NoKey,

A key, but the client could not be built (the ledger, usually).

55    Broken(String),
56}

What the daemon remembers about one Claude Code session.

59#[derive(Default)]
60struct Session {

The user's latest prompts, oldest first.

62    requests: VecDeque<String>,

What this session's questions have cost.

64    usd: f64,

Tool calls that were judged, so their outcomes are logged.

66    judged: HashSet<String>,
67}
69struct Daemon {
70    client: Client,
71    identity: String,
72    started: Instant,
73    log: DecisionLog,
74    headroom: Headroom,

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}

One event as the mod sends it.

84#[derive(Deserialize)]
85struct EventRequest {
86    session: String,
87    event: HookEvent,

The event's own fields: tool, command and tool_use_id for a tool call; prompt for a submitted prompt; last_assistant_message and the rest for a Stop.

91    #[serde(default)]
92    input: Value,
93    cwd: Option<String>,
94    root: Option<String>,
95}
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    }

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    }

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    }

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    }
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}

Runs the daemon until /shutdown. Returns at once, successfully, when 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}