Skip to main content

server/
http.rs

1//! Thin HTTP wrapper over [`SharedDb`]. Every endpoint is a lock, a public
2//! core-api call, then a response — no business logic.
3//!
4//! # Single-sink design
5//!
6//! The HTTP router is the designated broadcast producer. [`router`] installs
7//! one `broadcast::Sender` as the [`core_api::GraphDb`] event sink. MCP and
8//! CLI mutations on the same [`SharedDb`] fire into that same sink (one
9//! producer, many `/watch` subscribers). A second [`router`] call replaces
10//! the sink and terminates every existing subscriber with
11//! [`tokio::sync::broadcast::error::RecvError::Closed`].
12
13use crate::json::{
14    edge_history_result_json, node_edges_json, node_history_json, node_info_json, params_from_json,
15    parse_ingest_edges, result_set_json, rule_def_from_json,
16};
17use crate::{AppState, AuthIdentity};
18use arrow_bridge::to_ipc_bytes;
19use axum::extract::{Extension, Path, Query, Request, State};
20use axum::http::{header, HeaderValue, Method, StatusCode};
21use axum::middleware::{self, Next};
22use axum::response::{IntoResponse, Response};
23use axum::routing::{get, post};
24use axum::{Json, Router};
25use core_api::{
26    is_write_query, json_to_rows, json_to_value, AutoFk, BackupReport, BatchOp, DegreeConfig, Dir,
27    GraphError, IngestOptions, MaskMode, NodeMask, PageRankConfig, ResultSet, SharedDb,
28    SuggestConfig, Value, WccConfig, SUGGEST_DEFAULT_SEED,
29};
30use serde_json::{json, Value as Js};
31use std::collections::{BTreeMap, HashMap};
32use std::net::SocketAddr;
33use std::path::PathBuf;
34use tower_http::services::ServeDir;
35
36/// Build the HTTP router over `db`. Read endpoints take the read lock;
37/// `/ingest` takes the write lock. Guards are dropped before any `.await`.
38/// `GET /watch` upgrades to a WebSocket fed by the post-commit sink.
39///
40/// **Call at most once per [`SharedDb`].** A second call replaces the sink
41/// and terminates all existing `/watch` subscribers with
42/// [`tokio::sync::broadcast::error::RecvError::Closed`].
43///
44/// Installing the watch sink replaces any previously installed
45/// [`core_api::GraphDb::set_event_sink`]. The sink only
46/// `broadcast::Sender::send`s (non-blocking) and never re-enters `db`.
47///
48/// # Blocking
49///
50/// [`SharedDb`] uses a std [`std::sync::RwLock`]. Write handlers
51/// (`query_write`, `/ingest`, `create_rule`) run in
52/// `tokio::task::spawn_blocking` and drop the write guard before `.await`.
53/// Reads stay on the worker (neighborhood is µs). `suggest` and `algo`
54/// already use the blocking pool.
55pub fn router(db: SharedDb) -> Router {
56    router_with_auth(db, None)
57}
58
59/// [`router`] with an optional bearer/`?token=` requirement on every route
60/// except unauthenticated `GET /health`.
61///
62/// Role enforcement is not active on this entry point; use
63/// [`router_with_role_tokens`] when role-bound tokens are required.
64pub fn router_with_auth(db: SharedDb, token: Option<String>) -> Router {
65    build_app(
66        db,
67        token,
68        HashMap::new(),
69        UiFallback::None,
70        default_advertise_addr(),
71    )
72}
73
74/// [`router`] with a full-access token and a map of role-bound tokens.
75///
76/// Role-bound tokens (`role_tokens`: bearer → role name) receive masked reads
77/// on `/query` and `/node/*`, and 403 on all write or subscription endpoints.
78/// Unknown token: 401.  Token bound to a role not in the DB at request time: 401.
79pub fn router_with_role_tokens(
80    db: SharedDb,
81    token: Option<String>,
82    role_tokens: HashMap<String, String>,
83) -> Router {
84    build_app(
85        db,
86        token,
87        role_tokens,
88        UiFallback::None,
89        default_advertise_addr(),
90    )
91}
92
93/// Same as [`router_with_auth`], then `ServeDir` as the fallback so API routes win.
94///
95/// Role enforcement is not active on this entry point; use
96/// [`router_with_role_tokens`] when role-bound tokens are required.
97pub fn router_with_ui(
98    db: SharedDb,
99    ui_dir: impl AsRef<std::path::Path>,
100    token: Option<String>,
101) -> Router {
102    build_app(
103        db,
104        token,
105        HashMap::new(),
106        UiFallback::Dir(ui_dir.as_ref().to_path_buf()),
107        default_advertise_addr(),
108    )
109}
110
111#[cfg(feature = "embed-ui")]
112static EMBEDDED_UI: include_dir::Dir<'_> =
113    include_dir::include_dir!("$CARGO_MANIFEST_DIR/../../ui/dist");
114
115/// [`router`] plus the `embed-ui` static tree as fallback.
116#[cfg(feature = "embed-ui")]
117pub fn router_with_embedded_ui(db: SharedDb) -> Router {
118    build_app(
119        db,
120        None,
121        HashMap::new(),
122        UiFallback::Embedded,
123        default_advertise_addr(),
124    )
125}
126
127#[cfg(feature = "embed-ui")]
128async fn embedded_fallback(uri: axum::http::Uri) -> Response {
129    let rel = if uri.path() == "/" || uri.path().is_empty() {
130        "index.html"
131    } else {
132        uri.path().trim_start_matches('/')
133    };
134    if rel.split('/').any(|seg| seg == "..") {
135        return StatusCode::NOT_FOUND.into_response();
136    }
137    match EMBEDDED_UI.get_file(rel) {
138        Some(file) => (
139            StatusCode::OK,
140            [(header::CONTENT_TYPE, embedded_ctype(rel))],
141            file.contents(),
142        )
143            .into_response(),
144        None => StatusCode::NOT_FOUND.into_response(),
145    }
146}
147
148#[cfg(feature = "embed-ui")]
149fn embedded_ctype(path: &str) -> &'static str {
150    if path.ends_with(".html") {
151        "text/html; charset=utf-8"
152    } else if path.ends_with(".js") {
153        "application/javascript; charset=utf-8"
154    } else if path.ends_with(".css") {
155        "text/css; charset=utf-8"
156    } else if path.ends_with(".woff2") {
157        "font/woff2"
158    } else if path.ends_with(".svg") {
159        "image/svg+xml"
160    } else if path.ends_with(".ico") {
161        "image/x-icon"
162    } else if path.ends_with(".txt") {
163        "text/plain; charset=utf-8"
164    } else {
165        "application/octet-stream"
166    }
167}
168
169/// Bind `addr` (port 0 is ephemeral) and serve.
170///
171/// Sends the resolved local address on `ready` once the listener is accepting.
172/// Does not hold a database lock.
173///
174/// Role enforcement is not active on this entry point; use
175/// [`serve_with_role_tokens`] when role-bound tokens are required.
176#[deprecated(
177    since = "0.2.0",
178    note = "Use `serve_with_role_tokens` instead; this variant silently ignores role-token configuration."
179)]
180#[doc(hidden)]
181pub async fn serve(
182    db: SharedDb,
183    addr: SocketAddr,
184    ready: tokio::sync::oneshot::Sender<SocketAddr>,
185    token: Option<String>,
186) -> std::io::Result<()> {
187    serve_inner(db, addr, ready, UiFallback::None, token, HashMap::new()).await
188}
189
190/// [`serve`] with role-bound tokens in addition to the optional full-access token.
191///
192/// `role_tokens` maps bearer values to role names.  See [`router_with_role_tokens`]
193/// for the enforcement semantics.
194pub async fn serve_with_role_tokens(
195    db: SharedDb,
196    addr: SocketAddr,
197    ready: tokio::sync::oneshot::Sender<SocketAddr>,
198    token: Option<String>,
199    role_tokens: HashMap<String, String>,
200) -> std::io::Result<()> {
201    serve_inner(db, addr, ready, UiFallback::None, token, role_tokens).await
202}
203
204/// [`serve`] plus a UI dist directory mounted behind the API routes.
205///
206/// Role enforcement is not active on this entry point; use
207/// [`serve_with_ui_and_role_tokens`] when role-bound tokens are required.
208#[deprecated(
209    since = "0.2.0",
210    note = "Use `serve_with_ui_and_role_tokens` instead; this variant silently ignores role-token configuration."
211)]
212#[doc(hidden)]
213pub async fn serve_with_ui(
214    db: SharedDb,
215    addr: SocketAddr,
216    ready: tokio::sync::oneshot::Sender<SocketAddr>,
217    ui_dir: PathBuf,
218    token: Option<String>,
219) -> std::io::Result<()> {
220    serve_inner(
221        db,
222        addr,
223        ready,
224        UiFallback::Dir(ui_dir),
225        token,
226        HashMap::new(),
227    )
228    .await
229}
230
231/// [`serve_with_ui`] with role-bound tokens.
232pub async fn serve_with_ui_and_role_tokens(
233    db: SharedDb,
234    addr: SocketAddr,
235    ready: tokio::sync::oneshot::Sender<SocketAddr>,
236    ui_dir: PathBuf,
237    token: Option<String>,
238    role_tokens: HashMap<String, String>,
239) -> std::io::Result<()> {
240    serve_inner(db, addr, ready, UiFallback::Dir(ui_dir), token, role_tokens).await
241}
242
243/// [`serve`] plus the compiled-in UI (no-op fallback if `embed-ui` is off).
244///
245/// `role_tokens` maps bearer values to role names; see [`router_with_role_tokens`]
246/// for enforcement semantics.  Pass `HashMap::new()` when no role tokens are needed.
247#[cfg(feature = "embed-ui")]
248pub async fn serve_with_embedded_ui(
249    db: SharedDb,
250    addr: SocketAddr,
251    ready: tokio::sync::oneshot::Sender<SocketAddr>,
252    token: Option<String>,
253    role_tokens: HashMap<String, String>,
254) -> std::io::Result<()> {
255    serve_inner(db, addr, ready, UiFallback::Embedded, token, role_tokens).await
256}
257
258enum UiFallback {
259    None,
260    Dir(PathBuf),
261    #[cfg(feature = "embed-ui")]
262    Embedded,
263}
264
265async fn serve_inner(
266    db: SharedDb,
267    addr: SocketAddr,
268    ready: tokio::sync::oneshot::Sender<SocketAddr>,
269    ui: UiFallback,
270    token: Option<String>,
271    role_tokens: HashMap<String, String>,
272) -> std::io::Result<()> {
273    let listener = tokio::net::TcpListener::bind(addr).await?;
274    let local = listener.local_addr()?;
275    if ready.send(local).is_err() {
276        // Caller dropped the readiness receiver; still serve.
277        eprintln!("serve: readiness receiver dropped before bind notify");
278    }
279    let app = build_app(db, token, role_tokens, ui, local);
280    axum::serve(listener, app).await
281}
282
283fn default_advertise_addr() -> SocketAddr {
284    SocketAddr::from(([127, 0, 0, 1], 8080))
285}
286
287fn build_app(
288    db: SharedDb,
289    token: Option<String>,
290    role_tokens: HashMap<String, String>,
291    ui: UiFallback,
292    addr: SocketAddr,
293) -> Router {
294    debug_assert!(
295        !db.read().has_event_sink(),
296        "router() must be called at most once per SharedDb; a second call \
297         replaces the sink and terminates all existing /watch subscribers \
298         with RecvError::Closed"
299    );
300    let (tx, _) = tokio::sync::broadcast::channel(1024);
301    {
302        let tx = tx.clone();
303        db.write().set_event_sink(Box::new(move |ev| {
304            let _ = tx.send(ev);
305        }));
306    }
307    let state = AppState {
308        db,
309        watch: tx,
310        token,
311        role_tokens,
312        addr,
313    };
314    let app = Router::new()
315        .route("/health", get(health))
316        .route("/query", post(query))
317        .route("/stats", get(stats))
318        .route("/ingest", post(ingest))
319        .route("/rules", post(create_rule))
320        .route("/suggest", get(suggest))
321        .route("/explain", get(explain))
322        .route("/node/{key}", get(node_info))
323        .route("/node/{key}", axum::routing::delete(delete_node))
324        .route("/node/{key}/edges", get(node_edges))
325        .route("/node/{key}/neighborhood", get(neighborhood))
326        .route("/node/{key}/history", get(node_history_handler))
327        .route("/history/edge", get(edge_history_handler))
328        .route("/history/was_linked", get(was_linked_handler))
329        .route(
330            "/node/{key}/prop/{field}",
331            axum::routing::put(set_node_prop),
332        )
333        .route(
334            "/node/{key}/prop/{field}",
335            axum::routing::delete(remove_node_prop),
336        )
337        // Simple BatchOp-mapped endpoints — routed through the group-commit
338        // queue (submit_batch) so concurrent node/edge CRUD does not hold
339        // the write lock during fsync.
340        .route("/nodes", post(create_node))
341        .route("/nodes/{key}/rename", post(rename_node))
342        .route("/edges", post(create_edge))
343        .route("/edges/upsert", post(upsert_edge))
344        .route(
345            "/edges/{etype}/{src}/{dst}",
346            axum::routing::delete(delete_edge),
347        )
348        .route("/algo/pagerank", post(algo_pagerank))
349        .route("/algo/wcc", post(algo_wcc))
350        .route("/algo/degree", post(algo_degree))
351        .route("/backup", post(backup))
352        .route("/watch", get(crate::ws::watch))
353        .route("/subscribe", get(crate::subscribe::subscribe))
354        .with_state(state.clone());
355    let app = match ui {
356        UiFallback::None => app,
357        UiFallback::Dir(dir) => app.fallback_service(ServeDir::new(dir)),
358        #[cfg(feature = "embed-ui")]
359        UiFallback::Embedded => app.fallback(embedded_fallback),
360    };
361    app.layer(middleware::from_fn_with_state(state, auth_middleware))
362}
363
364async fn health(State(state): State<AppState>) -> Response {
365    let (nodes, edges) = {
366        let g = state.db.read();
367        let s = g.stats();
368        (s.nodes_live, s.edges)
369    };
370    json_ok(json!({
371        "ok": true,
372        "nodes": nodes,
373        "edges": edges,
374        "addr": state.addr.to_string(),
375    }))
376}
377
378/// Run a GraphDb write on the blocking pool. The write guard lives only
379/// inside `f` and is dropped before this future awaits.
380async fn blocking_write<T, F>(f: F) -> std::result::Result<T, Response>
381where
382    T: Send + 'static,
383    F: FnOnce() -> core_api::Result<T> + Send + 'static,
384{
385    match tokio::task::spawn_blocking(f).await {
386        Ok(Ok(v)) => Ok(v),
387        Ok(Err(e)) => Err(graph_err(e)),
388        Err(_) => Err(err_response("write task panicked")),
389    }
390}
391
392const TOKEN_COOKIE: &str = "mushroomdb_token";
393
394async fn auth_middleware(State(state): State<AppState>, mut req: Request, next: Next) -> Response {
395    // No auth configured at all: every request is Full, no restriction.
396    if state.token.is_none() && state.role_tokens.is_empty() {
397        req.extensions_mut().insert(AuthIdentity::Full);
398        return next.run(req).await;
399    }
400
401    // Health is always open regardless of token configuration.
402    if req.method() == Method::GET && req.uri().path() == "/health" {
403        req.extensions_mut().insert(AuthIdentity::Full);
404        return next.run(req).await;
405    }
406
407    let presented = request_token(&req);
408
409    // Check full-access token first.
410    if let Some(ref full_tok) = state.token.clone().filter(|s| !s.is_empty()) {
411        if presented.as_deref() == Some(full_tok.as_str()) {
412            let set_cookie = presented_bearer_or_query(&req).as_deref() == Some(full_tok.as_str());
413            req.extensions_mut().insert(AuthIdentity::Full);
414            let mut res = next.run(req).await;
415            if set_cookie && is_html_response(&res) {
416                attach_token_cookie(&mut res, full_tok);
417            }
418            return res;
419        }
420    }
421
422    // Check role-bound tokens.
423    if let Some(tok) = presented.as_deref() {
424        if let Some(role_name) = state.role_tokens.get(tok) {
425            // Early-deny paths whose handlers use WebSocket upgrade: the WS
426            // extraction consumes the request before the handler body runs, so
427            // the handler-level identity check is unreachable in those cases.
428            // Returning 403 here also removes the need for the handler to
429            // hold a read lock just to enforce this.
430            let path = req.uri().path();
431            if path == "/subscribe" || path == "/watch" {
432                return forbidden("role-bound token: this endpoint is not permitted");
433            }
434            req.extensions_mut()
435                .insert(AuthIdentity::Role(role_name.clone()));
436            return next.run(req).await;
437        }
438    }
439
440    // Nothing matched → 401.
441    unauthorized()
442}
443
444fn request_token(req: &Request) -> Option<String> {
445    presented_bearer_or_query(req).or_else(|| presented_cookie(req))
446}
447
448fn presented_bearer_or_query(req: &Request) -> Option<String> {
449    if let Some(header) = req
450        .headers()
451        .get(header::AUTHORIZATION)
452        .and_then(|v| v.to_str().ok())
453    {
454        if let Some(value) = bearer_token(header) {
455            return Some(value.to_string());
456        }
457    }
458    query_param(req.uri().query().unwrap_or(""), "token")
459}
460
461fn presented_cookie(req: &Request) -> Option<String> {
462    let header = req.headers().get(header::COOKIE)?.to_str().ok()?;
463    cookie_named(header, TOKEN_COOKIE).map(str::to_string)
464}
465
466fn cookie_named<'a>(header: &'a str, name: &str) -> Option<&'a str> {
467    for part in header.split(';') {
468        let part = part.trim();
469        let Some((k, v)) = part.split_once('=') else {
470            continue;
471        };
472        if k.trim() == name {
473            return Some(v.trim());
474        }
475    }
476    None
477}
478
479fn is_html_response(res: &Response) -> bool {
480    res.headers()
481        .get(header::CONTENT_TYPE)
482        .and_then(|v| v.to_str().ok())
483        .is_some_and(|ct| {
484            ct.split(';')
485                .next()
486                .unwrap_or("")
487                .trim()
488                .eq_ignore_ascii_case("text/html")
489        })
490}
491
492fn attach_token_cookie(res: &mut Response, token: &str) {
493    let value = format!("{TOKEN_COOKIE}={token}; Path=/; SameSite=Lax; HttpOnly");
494    if let Ok(hv) = HeaderValue::from_str(&value) {
495        res.headers_mut().insert(header::SET_COOKIE, hv);
496    }
497}
498
499fn bearer_token(header: &str) -> Option<&str> {
500    let (scheme, value) = header.split_once(' ')?;
501    if scheme.eq_ignore_ascii_case("Bearer") {
502        Some(value.trim())
503    } else {
504        None
505    }
506}
507
508fn query_param(query: &str, key: &str) -> Option<String> {
509    for pair in query.split('&') {
510        if pair.is_empty() {
511            continue;
512        }
513        match pair.split_once('=') {
514            Some((k, v)) if k == key => return percent_decode_plus(v),
515            None if pair == key => return Some(String::new()),
516            _ => {}
517        }
518    }
519    None
520}
521
522/// `application/x-www-form-urlencoded`: `+` is space, `%HH` is a byte.
523fn percent_decode_plus(s: &str) -> Option<String> {
524    let bytes = s.as_bytes();
525    let mut out = Vec::with_capacity(bytes.len());
526    let mut i = 0;
527    while i < bytes.len() {
528        match bytes[i] {
529            b'+' => {
530                out.push(b' ');
531                i += 1;
532            }
533            b'%' => {
534                if i + 2 >= bytes.len() {
535                    return None;
536                }
537                let hi = from_hex(bytes[i + 1])?;
538                let lo = from_hex(bytes[i + 2])?;
539                out.push((hi << 4) | lo);
540                i += 3;
541            }
542            c => {
543                out.push(c);
544                i += 1;
545            }
546        }
547    }
548    String::from_utf8(out).ok()
549}
550
551fn from_hex(b: u8) -> Option<u8> {
552    match b {
553        b'0'..=b'9' => Some(b - b'0'),
554        b'a'..=b'f' => Some(b - b'a' + 10),
555        b'A'..=b'F' => Some(b - b'A' + 10),
556        _ => None,
557    }
558}
559
560fn unauthorized() -> Response {
561    (
562        StatusCode::UNAUTHORIZED,
563        Json(json!({"error": "unauthorized"})),
564    )
565        .into_response()
566}
567
568/// 403 response for role-bound token operations that are denied in v1.
569fn forbidden(detail: &str) -> Response {
570    (StatusCode::FORBIDDEN, Json(json!({"error": detail}))).into_response()
571}
572
573/// Convert a `mask_for_role` error into an HTTP response.
574///
575/// - `GraphError::Corrupt` → 500: the roles sidecar is poisoned; the server
576///   refuses all role-token requests until the file is fixed and the DB is
577///   re-opened.  Full-access tokens are unaffected.
578/// - `GraphError::KeyNotFound` with a `role:` prefix → 401: the token is bound
579///   to a role that does not exist in the DB at this request time.
580fn role_mask_err(e: GraphError) -> Response {
581    match e {
582        GraphError::Corrupt { detail } => (
583            StatusCode::INTERNAL_SERVER_ERROR,
584            Json(json!({"error": format!("roles misconfigured: {detail}")})),
585        )
586            .into_response(),
587        GraphError::KeyNotFound { key } if key.starts_with("role:") => unauthorized(),
588        other => graph_err(other),
589    }
590}
591
592fn err_response(detail: impl Into<String>) -> Response {
593    (
594        StatusCode::BAD_REQUEST,
595        Json(json!({"error": detail.into()})),
596    )
597        .into_response()
598}
599
600fn graph_err(e: GraphError) -> Response {
601    let detail = match e {
602        GraphError::QueryError { detail } | GraphError::IngestError { detail } => detail,
603        other => other.to_string(),
604    };
605    err_response(detail)
606}
607
608fn key_not_found(key: String) -> Response {
609    (
610        StatusCode::NOT_FOUND,
611        Json(json!({"error": GraphError::KeyNotFound { key }.to_string()})),
612    )
613        .into_response()
614}
615
616fn conflict_response(key: String) -> Response {
617    (
618        StatusCode::CONFLICT,
619        Json(json!({"error": GraphError::DuplicateKey { key }.to_string()})),
620    )
621        .into_response()
622}
623
624fn json_ok(value: Js) -> Response {
625    (StatusCode::OK, Json(value)).into_response()
626}
627
628fn ingest_options(v: Option<&Js>) -> Result<IngestOptions, String> {
629    let Some(v) = v else {
630        return Ok(IngestOptions::default());
631    };
632    if v.is_null() {
633        return Ok(IngestOptions::default());
634    }
635    let obj = v
636        .as_object()
637        .ok_or_else(|| "options must be an object".to_string())?;
638    let mut opts = IngestOptions::default();
639    if let Some(kf) = obj.get("key_field") {
640        opts.key_field = kf
641            .as_str()
642            .ok_or_else(|| "options.key_field must be a string".to_string())?
643            .to_string();
644    }
645    if let Some(fk) = obj.get("auto_fk") {
646        if fk == &Js::Bool(false) || fk.as_str() == Some("off") {
647            opts.auto_fk = AutoFk::Off;
648        } else if let Some(m) = fk.as_object() {
649            let suf = m
650                .get("suffix")
651                .and_then(Js::as_str)
652                .ok_or_else(|| "options.auto_fk.suffix must be a string".to_string())?;
653            opts.auto_fk = AutoFk::Auto {
654                suffix: suf.to_string(),
655            };
656        } else {
657            return Err("options.auto_fk must be false, \"off\", or {suffix}".into());
658        }
659    }
660    Ok(opts)
661}
662
663/// Format a `ResultSet` as Arrow IPC or JSON depending on `format`.
664fn format_query_result(rs: ResultSet, format: &str) -> Response {
665    match format {
666        "" => match to_ipc_bytes(&rs) {
667            Ok(bytes) => (
668                StatusCode::OK,
669                [(header::CONTENT_TYPE, "application/vnd.apache.arrow.stream")],
670                bytes,
671            )
672                .into_response(),
673            Err(e) => err_response(e),
674        },
675        "json" => json_ok(result_set_json(&rs)),
676        other => err_response(format!("unknown format: {other}")),
677    }
678}
679
680async fn query(
681    State(state): State<AppState>,
682    Extension(identity): Extension<AuthIdentity>,
683    Query(qs): Query<BTreeMap<String, String>>,
684    Json(body): Json<Js>,
685) -> Response {
686    let cypher = match body.get("cypher").and_then(Js::as_str) {
687        Some(s) => s.to_string(),
688        None => return err_response("missing cypher"),
689    };
690    let params = match params_from_json(body.get("params")) {
691        Ok(p) => p,
692        Err(e) => return err_response(e),
693    };
694    let format = qs.get("format").map(String::as_str).unwrap_or("");
695
696    // Parse client-supplied mask (optional array of node keys).
697    let mask_keys: Option<Vec<String>> = match body.get("mask") {
698        None | Some(Js::Null) => None,
699        Some(Js::Array(arr)) => {
700            let mut keys = Vec::with_capacity(arr.len());
701            for v in arr {
702                match v.as_str() {
703                    Some(s) => keys.push(s.to_string()),
704                    None => return err_response("mask must be an array of strings"),
705                }
706            }
707            Some(keys)
708        }
709        Some(_) => return err_response("mask must be an array of strings"),
710    };
711
712    // Role token: writes are denied; reads are auto-masked by the role's
713    // visibility mask.  A client-supplied mask can only narrow, never widen.
714    // Lock-free epoch snapshot: mask and query execute on the same frozen state
715    // (constraint 2 — RBAC mask coherence), without holding the read lock.
716    if let AuthIdentity::Role(ref role_name) = identity {
717        let is_write = match is_write_query(&cypher) {
718            Ok(b) => b,
719            Err(e) => return err_response(e),
720        };
721        if is_write {
722            return forbidden("role-bound token: writes are not permitted");
723        }
724        let snap = state.db.reader();
725        let role_mask = match snap.mask_for_role(role_name) {
726            Ok(m) => m,
727            Err(e) => return role_mask_err(e),
728        };
729        let effective_mask = if let Some(ref keys) = mask_keys {
730            // Client mask intersects role mask — never widens visibility.
731            let client_mask = NodeMask::from_ids(keys.iter().filter_map(|k| snap.resolve_key(k)));
732            role_mask.intersect(&client_mask)
733        } else {
734            role_mask
735        };
736        return match snap.query_masked(&cypher, &params, &effective_mask) {
737            Ok(rs) => format_query_result(rs, format),
738            Err(GraphError::MaskedReadOnly) => (
739                StatusCode::BAD_REQUEST,
740                Json(json!({"error": "masked queries are read-only"})),
741            )
742                .into_response(),
743            Err(e) => graph_err(e),
744        };
745    }
746
747    // Full token: existing paths below.
748
749    // When a client-supplied mask is present, route to query_masked (read-only).
750    // Hold a single read guard for both from_keys and query_masked so the mask
751    // and the query execute on the same database snapshot.
752    //
753    // `stub_hidden: true` opts into MaskMode::Stub for the mask; Cypher query
754    // behaviour is identical in both modes (hidden nodes are excluded from
755    // query results regardless of mode).
756    if let Some(ref keys) = mask_keys {
757        let stub_hidden = body
758            .get("stub_hidden")
759            .and_then(|v| v.as_bool())
760            .unwrap_or(false);
761        let db = state.db.read();
762        let mask = {
763            let m = NodeMask::from_keys(&*db, keys.iter().map(String::as_str));
764            if stub_hidden {
765                m.with_mode(MaskMode::Stub)
766            } else {
767                m
768            }
769        };
770        return match db.query_masked(&cypher, &params, &mask) {
771            Ok(rs) => format_query_result(rs, format),
772            Err(GraphError::MaskedReadOnly) => (
773                StatusCode::BAD_REQUEST,
774                Json(json!({"error": "masked queries are read-only"})),
775            )
776                .into_response(),
777            Err(e) => graph_err(e),
778        };
779    }
780
781    // Detect write statements to dispatch to the correct lock.
782    // Write statements (CREATE / MATCH…SET / MATCH…DELETE / MERGE) need the
783    // write lock so mutations flow through WAL + rule engine with fsync before
784    // the response is sent.  Read queries (MATCH … RETURN …) use the read lock.
785    let is_write = match is_write_query(&cypher) {
786        Ok(b) => b,
787        Err(e) => return err_response(e),
788    };
789
790    let rs = if is_write {
791        let db = state.db.clone();
792        match blocking_write(move || db.write().query_write(&cypher, &params)).await {
793            Ok(rs) => rs,
794            Err(resp) => return resp,
795        }
796    } else {
797        match state.db.read().query(&cypher, &params) {
798            Ok(rs) => rs,
799            Err(e) => return graph_err(e),
800        }
801    };
802
803    format_query_result(rs, format)
804}
805
806async fn stats(
807    State(state): State<AppState>,
808    Extension(identity): Extension<AuthIdentity>,
809) -> Response {
810    // v1: deny role tokens — raw counts leak graph size beyond the role's subgraph.
811    if let AuthIdentity::Role(_) = identity {
812        return forbidden("role-bound token: /stats requires a full-access token");
813    }
814    let snap = {
815        let g = state.db.read();
816        g.stats()
817    };
818    match serde_json::to_value(&snap) {
819        Ok(v) => json_ok(v),
820        Err(e) => err_response(e.to_string()),
821    }
822}
823
824async fn ingest(
825    State(state): State<AppState>,
826    Extension(identity): Extension<AuthIdentity>,
827    Json(body): Json<Js>,
828) -> Response {
829    if let AuthIdentity::Role(_) = identity {
830        return forbidden("role-bound token: writes are not permitted");
831    }
832    let label = match body.get("label").and_then(Js::as_str) {
833        Some(s) => s.to_string(),
834        None => return err_response("missing label"),
835    };
836    let rows = match body.get("rows") {
837        Some(r) => r,
838        None => return err_response("missing rows"),
839    };
840    let mut converted = match json_to_rows(rows) {
841        Ok(c) => c,
842        Err(e) => return graph_err(e),
843    };
844    let opts = match ingest_options(body.get("options")) {
845        Ok(o) => o,
846        Err(e) => return err_response(e),
847    };
848    let taken = std::mem::take(&mut converted.rows);
849    let edges = match body.get("edges") {
850        None | Some(Js::Null) => Vec::new(),
851        Some(raw) => match parse_ingest_edges(raw) {
852            Ok(e) => e,
853            Err(e) => return err_response(e),
854        },
855    };
856    let db = state.db.clone();
857    let report =
858        match blocking_write(move || db.write().ingest_with_edges(&label, taken, &opts, &edges))
859            .await
860        {
861            Ok(r) => converted.into_report(r),
862            Err(resp) => return resp,
863        };
864    match serde_json::to_value(&report) {
865        Ok(v) => json_ok(v),
866        Err(e) => err_response(e.to_string()),
867    }
868}
869
870/// `GET /suggest` — profile the database and return rule suggestions.
871///
872/// # Locking and blocking strategy
873///
874/// `suggest_rules_with_config` is CPU-intensive and synchronous. Running it on a
875/// Tokio worker thread would starve the executor. This handler offloads the work to
876/// `tokio::task::spawn_blocking`, which uses the blocking thread-pool. The
877/// `std::sync::RwLock` read guard is acquired and held inside the blocking task —
878/// reads don't block other reads; writes wait for the guard to drop. The global
879/// budget (`SuggestConfig::global_budget_ms`, default 5 s) caps lock-hold time.
880async fn suggest(
881    State(state): State<AppState>,
882    Extension(identity): Extension<AuthIdentity>,
883) -> Response {
884    // v1: deny role tokens — suggest scans the full graph and would reveal
885    // existence of nodes outside the role's mask.
886    if let AuthIdentity::Role(_) = identity {
887        return forbidden("role-bound token: /suggest requires a full-access token");
888    }
889    let db = state.db.clone();
890    match tokio::task::spawn_blocking(move || {
891        let config = SuggestConfig::default();
892        db.read()
893            .suggest_rules_with_config(&config, SUGGEST_DEFAULT_SEED)
894    })
895    .await
896    {
897        Ok(report) => json_ok(serde_json::to_value(&report).unwrap_or_else(|_| json!({}))),
898        Err(_) => err_response("suggest task panicked"),
899    }
900}
901
902async fn create_rule(
903    State(state): State<AppState>,
904    Extension(identity): Extension<AuthIdentity>,
905    Json(body): Json<Js>,
906) -> Response {
907    if let AuthIdentity::Role(_) = identity {
908        return forbidden("role-bound token: writes are not permitted");
909    }
910    let def = match rule_def_from_json(body) {
911        Ok(d) => d,
912        Err(e) => return err_response(e),
913    };
914    let name = def.name.clone();
915    let db = state.db.clone();
916    match blocking_write(move || db.write().create_rule(def)).await {
917        Ok(()) => json_ok(json!({"ok": true, "name": name})),
918        Err(resp) => resp,
919    }
920}
921
922async fn explain(
923    State(state): State<AppState>,
924    Extension(identity): Extension<AuthIdentity>,
925    Query(qs): Query<BTreeMap<String, String>>,
926) -> Response {
927    // v1: deny role tokens — explain reveals hidden-node linkage through rules.
928    if let AuthIdentity::Role(_) = identity {
929        return forbidden(
930            "role-bound token: /explain requires a full-access token \
931             (v1: explain may reveal hidden-node linkage; revisit when stubs land)",
932        );
933    }
934    let a = match qs.get("a") {
935        Some(s) if !s.is_empty() => s.clone(),
936        _ => return err_response("missing query param a"),
937    };
938    let b = match qs.get("b") {
939        Some(s) if !s.is_empty() => s.clone(),
940        _ => return err_response("missing query param b"),
941    };
942    let out = {
943        let g = state.db.read();
944        g.explain(&a, &b)
945    };
946    match out {
947        Ok(v) => match serde_json::to_value(&v) {
948            Ok(j) => json_ok(j),
949            Err(e) => err_response(e.to_string()),
950        },
951        Err(e) => graph_err(e),
952    }
953}
954
955async fn node_info(
956    State(state): State<AppState>,
957    Extension(identity): Extension<AuthIdentity>,
958    Path(key): Path<String>,
959    Query(qs): Query<BTreeMap<String, String>>,
960) -> Response {
961    if let AuthIdentity::Role(ref role_name) = identity {
962        // Role-token path: hard-coded Omit mode; hidden keys are indistinguishable
963        // from absent keys.  `stub_hidden` query param is silently ignored here —
964        // role paths must NEVER produce stubs (RBAC invariant).
965        let snap = state.db.reader();
966        let role_mask = match snap.mask_for_role(role_name) {
967            Ok(m) => m,
968            Err(e) => return role_mask_err(e),
969        };
970        if !snap
971            .resolve_key(&key)
972            .is_some_and(|id| role_mask.contains_id(id))
973        {
974            return key_not_found(key);
975        }
976        return match snap.node_info(&key) {
977            Some(info) => json_ok(node_info_json(&info)),
978            None => key_not_found(key),
979        };
980    }
981
982    // Full-token path: optional client mask + stub_hidden via query params.
983    // `mask=key1,key2` — comma-separated visible keys (empty string = no mask).
984    // `stub_hidden=true` — opt into MaskMode::Stub for this request.
985    let mask_param = qs.get("mask").map(String::as_str).unwrap_or("").trim();
986    if !mask_param.is_empty() {
987        let stub_hidden = qs
988            .get("stub_hidden")
989            .map(|v| v == "true" || v == "1")
990            .unwrap_or(false);
991        let g = state.db.read();
992        let mask = {
993            let keys = mask_param
994                .split(',')
995                .map(str::trim)
996                .filter(|s| !s.is_empty());
997            let m = NodeMask::from_keys(&*g, keys);
998            if stub_hidden {
999                m.with_mode(MaskMode::Stub)
1000            } else {
1001                m
1002            }
1003        };
1004        return match g.node_info_masked(&key, &mask) {
1005            Some(core_api::MaskedNodeResult::Visible(info)) => json_ok(node_info_json(&info)),
1006            Some(core_api::MaskedNodeResult::Restricted) => {
1007                json_ok(crate::json::stub_node_json(&key))
1008            }
1009            None => key_not_found(key),
1010        };
1011    }
1012
1013    let info = {
1014        let g = state.db.read();
1015        g.node_info(&key)
1016    };
1017    match info {
1018        Some(info) => json_ok(node_info_json(&info)),
1019        None => key_not_found(key),
1020    }
1021}
1022
1023async fn node_edges(
1024    State(state): State<AppState>,
1025    Extension(identity): Extension<AuthIdentity>,
1026    Path(key): Path<String>,
1027    Query(qs): Query<BTreeMap<String, String>>,
1028) -> Response {
1029    if let AuthIdentity::Role(ref role_name) = identity {
1030        // Role-token path: hard-coded Omit mode.  `stub_hidden` query param is
1031        // silently ignored — role paths must NEVER produce stubs (RBAC invariant).
1032        let snap = state.db.reader();
1033        let role_mask = match snap.mask_for_role(role_name) {
1034            Ok(m) => m,
1035            Err(e) => return role_mask_err(e),
1036        };
1037        if !snap
1038            .resolve_key(&key)
1039            .is_some_and(|id| role_mask.contains_id(id))
1040        {
1041            return key_not_found(key);
1042        }
1043        return match snap.node_edges(&key) {
1044            Ok(edges) => {
1045                // Filter out edges whose OTHER endpoint is hidden in the role
1046                // mask.  A role token must not learn about hidden neighbors via
1047                // the edge list even when the entry key itself is visible.
1048                let visible: Vec<_> = edges
1049                    .into_iter()
1050                    .filter(|e| {
1051                        let other = if e.src_key == key {
1052                            &e.dst_key
1053                        } else {
1054                            &e.src_key
1055                        };
1056                        snap.resolve_key(other)
1057                            .is_some_and(|id| role_mask.contains_id(id))
1058                    })
1059                    .collect();
1060                json_ok(node_edges_json(&visible))
1061            }
1062            Err(GraphError::KeyNotFound { key }) => key_not_found(key),
1063            Err(e) => graph_err(e),
1064        };
1065    }
1066
1067    // Full-token path: optional client mask + stub_hidden via query params.
1068    let mask_param = qs.get("mask").map(String::as_str).unwrap_or("").trim();
1069    if !mask_param.is_empty() {
1070        let stub_hidden = qs
1071            .get("stub_hidden")
1072            .map(|v| v == "true" || v == "1")
1073            .unwrap_or(false);
1074        let g = state.db.read();
1075        let mask = {
1076            let keys = mask_param
1077                .split(',')
1078                .map(str::trim)
1079                .filter(|s| !s.is_empty());
1080            let m = NodeMask::from_keys(&*g, keys);
1081            if stub_hidden {
1082                m.with_mode(MaskMode::Stub)
1083            } else {
1084                m
1085            }
1086        };
1087        return match g.node_edges_masked(&key, &mask) {
1088            Ok(edges) => json_ok(crate::json::masked_edges_json(&edges)),
1089            Err(GraphError::KeyNotFound { key }) => key_not_found(key),
1090            Err(e) => graph_err(e),
1091        };
1092    }
1093
1094    let out = {
1095        let g = state.db.read();
1096        g.node_edges(&key)
1097    };
1098    match out {
1099        Ok(edges) => json_ok(node_edges_json(&edges)),
1100        Err(GraphError::KeyNotFound { key }) => key_not_found(key),
1101        Err(e) => graph_err(e),
1102    }
1103}
1104
1105async fn neighborhood(
1106    State(state): State<AppState>,
1107    Extension(identity): Extension<AuthIdentity>,
1108    Path(key): Path<String>,
1109    Query(qs): Query<BTreeMap<String, String>>,
1110) -> Response {
1111    let depth = match qs.get("depth") {
1112        None => 1u32,
1113        Some(s) => match s.parse() {
1114            Ok(d) => d,
1115            Err(_) => return err_response("depth must be an integer"),
1116        },
1117    };
1118    let dir = match qs.get("dir").map(String::as_str).unwrap_or("both") {
1119        s if s.eq_ignore_ascii_case("out") => Dir::Out,
1120        s if s.eq_ignore_ascii_case("in") => Dir::In,
1121        s if s.eq_ignore_ascii_case("both") => Dir::Both,
1122        other => return err_response(format!("unknown dir: {other}")),
1123    };
1124    let edge_type_names: Option<Vec<String>> = qs.get("edge_types").map(|s| {
1125        s.split(',')
1126            .map(str::trim)
1127            .filter(|t| !t.is_empty())
1128            .map(str::to_string)
1129            .collect()
1130    });
1131    let etype_refs: Option<Vec<&str>> = edge_type_names
1132        .as_ref()
1133        .map(|v| v.iter().map(String::as_str).collect());
1134    if let AuthIdentity::Role(ref role_name) = identity {
1135        // Lock-free epoch snapshot: mask and neighborhood BFS on same frozen state
1136        // (constraint 2 — RBAC mask coherence).
1137        let snap = state.db.reader();
1138        let role_mask = match snap.mask_for_role(role_name) {
1139            Ok(m) => m,
1140            Err(e) => return role_mask_err(e),
1141        };
1142        if !snap
1143            .resolve_key(&key)
1144            .is_some_and(|id| role_mask.contains_id(id))
1145        {
1146            return key_not_found(key);
1147        }
1148        // Use the mask-aware BFS: hidden nodes are excluded from results AND
1149        // cannot be used as traversal intermediaries (never-leak invariant).
1150        let rs = match snap.neighborhood_masked(&key, depth, etype_refs.as_deref(), dir, &role_mask)
1151        {
1152            Some(rs) => rs,
1153            None => return key_not_found(key),
1154        };
1155        return json_ok(result_set_json(&rs));
1156    }
1157    // Full-token path: optional client mask + stub_hidden via query params.
1158    // `mask=key1,key2` — comma-separated visible keys.
1159    // `stub_hidden=true` — opt into MaskMode::Stub; hidden direct neighbours
1160    // appear as stub rows (label: null) in the result; BFS does not expand
1161    // through them in either mode.
1162    let mask_param = qs.get("mask").map(String::as_str).unwrap_or("").trim();
1163    if !mask_param.is_empty() {
1164        let stub_hidden = qs
1165            .get("stub_hidden")
1166            .map(|v| v == "true" || v == "1")
1167            .unwrap_or(false);
1168        let g = state.db.read();
1169        let mask = {
1170            let keys = mask_param
1171                .split(',')
1172                .map(str::trim)
1173                .filter(|s| !s.is_empty());
1174            let m = NodeMask::from_keys(&*g, keys);
1175            if stub_hidden {
1176                m.with_mode(MaskMode::Stub)
1177            } else {
1178                m
1179            }
1180        };
1181        return match g.neighborhood_masked(&key, depth, etype_refs.as_deref(), dir, &mask) {
1182            Some(rs) => json_ok(result_set_json(&rs)),
1183            None => graph_err(GraphError::KeyNotFound { key: key.clone() }),
1184        };
1185    }
1186
1187    // Unmasked full-token path (no mask param).
1188    let rs = {
1189        let g = state.db.read();
1190        match g.node_ref(&key) {
1191            Some(n) => Ok(n.neighborhood(depth, etype_refs.as_deref(), dir)),
1192            None => Err(GraphError::KeyNotFound { key: key.clone() }),
1193        }
1194    };
1195    match rs {
1196        Ok(rs) => json_ok(result_set_json(&rs)),
1197        Err(e) => graph_err(e),
1198    }
1199}
1200
1201/// `POST /algo/pagerank` — run PageRank over the unified topology.
1202///
1203/// # Locking and blocking strategy
1204///
1205/// PageRank is CPU-intensive and synchronous. This handler offloads the work
1206/// to `tokio::task::spawn_blocking` (blocking thread-pool). The read guard is
1207/// acquired and held inside the blocking task — reads don't block other reads.
1208/// The `budget_ms` field in [`PageRankConfig`] caps lock-hold time.
1209async fn algo_pagerank(
1210    State(state): State<AppState>,
1211    Extension(identity): Extension<AuthIdentity>,
1212    Json(body): Json<serde_json::Value>,
1213) -> Response {
1214    // v1: deny role tokens — algo endpoints scan the full graph.
1215    if let AuthIdentity::Role(_) = identity {
1216        return forbidden("role-bound token: /algo/* requires a full-access token");
1217    }
1218    let config: PageRankConfig = match serde_json::from_value(body) {
1219        Ok(c) => c,
1220        Err(e) => return err_response(format!("invalid pagerank config: {e}")),
1221    };
1222    let db = state.db.clone();
1223    match tokio::task::spawn_blocking(move || db.read().pagerank(&config)).await {
1224        Ok(report) => json_ok(serde_json::to_value(&report).unwrap_or_else(|_| json!({}))),
1225        Err(_) => err_response("pagerank task panicked"),
1226    }
1227}
1228
1229/// `POST /algo/wcc` — weakly-connected components over the unified topology.
1230///
1231/// Mirrors the `suggest` locking and blocking pattern exactly: spawn_blocking,
1232/// read guard inside, `budget_ms` in config caps lock-hold time.
1233async fn algo_wcc(
1234    State(state): State<AppState>,
1235    Extension(identity): Extension<AuthIdentity>,
1236    Json(body): Json<serde_json::Value>,
1237) -> Response {
1238    if let AuthIdentity::Role(_) = identity {
1239        return forbidden("role-bound token: /algo/* requires a full-access token");
1240    }
1241    let config: WccConfig = match serde_json::from_value(body) {
1242        Ok(c) => c,
1243        Err(e) => return err_response(format!("invalid wcc config: {e}")),
1244    };
1245    let db = state.db.clone();
1246    match tokio::task::spawn_blocking(move || db.read().connected_components(&config)).await {
1247        Ok(report) => json_ok(serde_json::to_value(&report).unwrap_or_else(|_| json!({}))),
1248        Err(_) => err_response("wcc task panicked"),
1249    }
1250}
1251
1252/// `POST /algo/degree` — degree centrality over the unified topology.
1253///
1254/// Mirrors the `suggest` locking and blocking pattern exactly.
1255async fn algo_degree(
1256    State(state): State<AppState>,
1257    Extension(identity): Extension<AuthIdentity>,
1258    Json(body): Json<serde_json::Value>,
1259) -> Response {
1260    if let AuthIdentity::Role(_) = identity {
1261        return forbidden("role-bound token: /algo/* requires a full-access token");
1262    }
1263    let config: DegreeConfig = match serde_json::from_value(body) {
1264        Ok(c) => c,
1265        Err(e) => return err_response(format!("invalid degree config: {e}")),
1266    };
1267    let db = state.db.clone();
1268    match tokio::task::spawn_blocking(move || db.read().degree_centrality(&config)).await {
1269        Ok(report) => json_ok(serde_json::to_value(&report).unwrap_or_else(|_| json!({}))),
1270        Err(_) => err_response("degree task panicked"),
1271    }
1272}
1273
1274// ── Simple BatchOp-mapped endpoints ──────────────────────────────────────────
1275//
1276// These endpoints map 1:1 to BatchOp variants and route through submit_batch
1277// so concurrent writes share one WAL fsync per drain group, keeping reader
1278// p95 latency low under write bursts.  Auth checks (role tokens → 403) run
1279// before enqueue so RBAC enforcement is unchanged.
1280//
1281// Complex paths (/query Cypher writes, /ingest bulk JSON) stay on db.write()
1282// because they are multi-step operations that cannot be pre-expressed as a
1283// Vec<BatchOp> without redesigning the query executor.
1284
1285/// Parse a JSON object into a `Vec<(String, Value)>` prop list.
1286fn props_from_json_obj(v: &serde_json::Value) -> Result<Vec<(String, Value)>, String> {
1287    let obj = match v.as_object() {
1288        Some(o) => o,
1289        None => return Err("props must be a JSON object".into()),
1290    };
1291    let mut out = Vec::with_capacity(obj.len());
1292    for (k, val) in obj {
1293        if let Some(v) = json_to_value(val.clone()) {
1294            out.push((k.clone(), v));
1295        }
1296    }
1297    Ok(out)
1298}
1299
1300/// `POST /nodes` — create or upsert a node via the group-commit queue.
1301///
1302/// Body: `{"label": "Person", "key": "alice", "props": {"age": 30}}`
1303async fn create_node(
1304    State(state): State<AppState>,
1305    Extension(identity): Extension<AuthIdentity>,
1306    Json(body): Json<Js>,
1307) -> Response {
1308    if let AuthIdentity::Role(_) = identity {
1309        return forbidden("role-bound token: writes are not permitted");
1310    }
1311    let label = match body.get("label").and_then(Js::as_str) {
1312        Some(s) => s.to_string(),
1313        None => return err_response("missing label"),
1314    };
1315    let key = match body.get("key").and_then(Js::as_str) {
1316        Some(s) => s.to_string(),
1317        None => return err_response("missing key"),
1318    };
1319    let props = match body.get("props") {
1320        None | Some(Js::Null) => vec![],
1321        Some(v) => match props_from_json_obj(v) {
1322            Ok(p) => p,
1323            Err(e) => return err_response(e),
1324        },
1325    };
1326    let db = state.db.clone();
1327    match blocking_write(move || db.submit_batch(vec![BatchOp::InsertNode { label, key, props }]))
1328        .await
1329    {
1330        Ok((nodes, edges)) => json_ok(json!({"ok": true, "nodes": nodes, "edges": edges})),
1331        Err(resp) => resp,
1332    }
1333}
1334
1335/// `DELETE /node/{key}` — delete a node via the group-commit queue.
1336async fn delete_node(
1337    State(state): State<AppState>,
1338    Extension(identity): Extension<AuthIdentity>,
1339    Path(key): Path<String>,
1340) -> Response {
1341    if let AuthIdentity::Role(_) = identity {
1342        return forbidden("role-bound token: writes are not permitted");
1343    }
1344    let db = state.db.clone();
1345    match blocking_write(move || db.submit_batch(vec![BatchOp::DeleteNode { key }])).await {
1346        Ok(_) => json_ok(json!({"ok": true})),
1347        Err(resp) => resp,
1348    }
1349}
1350
1351/// `POST /edges` — create an edge via the group-commit queue.
1352///
1353/// Body: `{"type": "KNOWS", "src": "alice", "dst": "bob"}`
1354async fn create_edge(
1355    State(state): State<AppState>,
1356    Extension(identity): Extension<AuthIdentity>,
1357    Json(body): Json<Js>,
1358) -> Response {
1359    if let AuthIdentity::Role(_) = identity {
1360        return forbidden("role-bound token: writes are not permitted");
1361    }
1362    let edge_type = match body.get("type").and_then(Js::as_str) {
1363        Some(s) => s.to_string(),
1364        None => return err_response("missing type"),
1365    };
1366    let src = match body.get("src").and_then(Js::as_str) {
1367        Some(s) => s.to_string(),
1368        None => return err_response("missing src"),
1369    };
1370    let dst = match body.get("dst").and_then(Js::as_str) {
1371        Some(s) => s.to_string(),
1372        None => return err_response("missing dst"),
1373    };
1374    let db = state.db.clone();
1375    match blocking_write(move || {
1376        db.submit_batch(vec![BatchOp::InsertEdge {
1377            edge_type,
1378            src_key: src,
1379            dst_key: dst,
1380        }])
1381    })
1382    .await
1383    {
1384        Ok(_) => json_ok(json!({"ok": true})),
1385        Err(resp) => resp,
1386    }
1387}
1388
1389/// `DELETE /edges/{etype}/{src}/{dst}` — delete an edge via the group-commit queue.
1390async fn delete_edge(
1391    State(state): State<AppState>,
1392    Extension(identity): Extension<AuthIdentity>,
1393    Path((etype, src, dst)): Path<(String, String, String)>,
1394) -> Response {
1395    if let AuthIdentity::Role(_) = identity {
1396        return forbidden("role-bound token: writes are not permitted");
1397    }
1398    let db = state.db.clone();
1399    match blocking_write(move || {
1400        db.submit_batch(vec![BatchOp::DeleteEdge {
1401            edge_type: etype,
1402            src_key: src,
1403            dst_key: dst,
1404        }])
1405    })
1406    .await
1407    {
1408        Ok(_) => json_ok(json!({"ok": true})),
1409        Err(resp) => resp,
1410    }
1411}
1412
1413/// `POST /nodes/{key}/rename` — rename a node's key via the group-commit queue.
1414///
1415/// Body: `{"new_key": "alice2"}`
1416/// Returns 404 on KeyNotFound, 409 on DuplicateKey, 200 on success.
1417async fn rename_node(
1418    State(state): State<AppState>,
1419    Extension(identity): Extension<AuthIdentity>,
1420    Path(key): Path<String>,
1421    Json(body): Json<Js>,
1422) -> Response {
1423    if let AuthIdentity::Role(_) = identity {
1424        return forbidden("role-bound token: writes are not permitted");
1425    }
1426    let new_key = match body.get("new_key").and_then(Js::as_str) {
1427        Some(s) => s.to_string(),
1428        None => return err_response("missing new_key"),
1429    };
1430    let db = state.db.clone();
1431    match tokio::task::spawn_blocking(move || {
1432        db.submit_batch(vec![BatchOp::RenameNode {
1433            old_key: key,
1434            new_key,
1435        }])
1436    })
1437    .await
1438    {
1439        Ok(Ok(_)) => json_ok(json!({"ok": true})),
1440        Ok(Err(GraphError::KeyNotFound { key })) => key_not_found(key),
1441        Ok(Err(GraphError::DuplicateKey { key })) => conflict_response(key),
1442        Ok(Err(e)) => graph_err(e),
1443        Err(_) => err_response("write task panicked"),
1444    }
1445}
1446
1447/// `POST /edges/upsert` — insert an edge, auto-creating missing endpoints.
1448///
1449/// Body: `{"edge_type":"KNOWS","src_key":"alice","dst_key":"bob","placeholder_label":"Person"}`
1450/// Returns `{"nodes_created": N, "edge_inserted": bool}`.
1451async fn upsert_edge(
1452    State(state): State<AppState>,
1453    Extension(identity): Extension<AuthIdentity>,
1454    Json(body): Json<Js>,
1455) -> Response {
1456    if let AuthIdentity::Role(_) = identity {
1457        return forbidden("role-bound token: writes are not permitted");
1458    }
1459    let edge_type = match body.get("edge_type").and_then(Js::as_str) {
1460        Some(s) => s.to_string(),
1461        None => return err_response("missing edge_type"),
1462    };
1463    let src_key = match body.get("src_key").and_then(Js::as_str) {
1464        Some(s) => s.to_string(),
1465        None => return err_response("missing src_key"),
1466    };
1467    let dst_key = match body.get("dst_key").and_then(Js::as_str) {
1468        Some(s) => s.to_string(),
1469        None => return err_response("missing dst_key"),
1470    };
1471    let placeholder_label = match body.get("placeholder_label").and_then(Js::as_str) {
1472        Some(s) => s.to_string(),
1473        None => return err_response("missing placeholder_label"),
1474    };
1475    let db = state.db.clone();
1476    match blocking_write(move || {
1477        db.submit_batch(vec![BatchOp::InsertEdgeUpsert {
1478            edge_type,
1479            src_key,
1480            dst_key,
1481            placeholder_label,
1482        }])
1483    })
1484    .await
1485    {
1486        Ok((nodes, edges)) => json_ok(json!({
1487            "nodes_created": nodes,
1488            "edge_inserted": edges > 0,
1489        })),
1490        Err(resp) => resp,
1491    }
1492}
1493
1494/// `PUT /node/{key}/prop/{field}` — set a property via the group-commit queue.
1495///
1496/// Body: `{"value": 42}`
1497async fn set_node_prop(
1498    State(state): State<AppState>,
1499    Extension(identity): Extension<AuthIdentity>,
1500    Path((key, field)): Path<(String, String)>,
1501    Json(body): Json<Js>,
1502) -> Response {
1503    if let AuthIdentity::Role(_) = identity {
1504        return forbidden("role-bound token: writes are not permitted");
1505    }
1506    let value = match body.get("value").and_then(|v| json_to_value(v.clone())) {
1507        Some(v) => v,
1508        None => return err_response("missing or null value"),
1509    };
1510    let db = state.db.clone();
1511    match blocking_write(move || db.submit_batch(vec![BatchOp::SetProp { key, field, value }]))
1512        .await
1513    {
1514        Ok(_) => json_ok(json!({"ok": true})),
1515        Err(resp) => resp,
1516    }
1517}
1518
1519// ── History endpoints ─────────────────────────────────────────────────────────
1520//
1521// These are cold-path diagnostic endpoints. They scan the on-disk WAL and
1522// must NOT extend ReaderSnapshot — they use db.read() directly per the
1523// controller ruling (see module doc and task-2-brief.md).
1524//
1525// Role masking: node visibility is checked under the SAME read guard as the
1526// history call (coherent snapshot). A node outside the role's mask responds
1527// identically to an absent node — no existence oracle.
1528
1529/// `GET /node/{key}/history` — return the WAL change history for `key`.
1530///
1531/// Response: `{ key, history: [{commit, change}], total_commits }`.
1532/// Role tokens: if `key` is hidden by the role mask, responds with 404
1533/// (same shape as querying an absent key — no existence oracle).
1534async fn node_history_handler(
1535    State(state): State<AppState>,
1536    Extension(identity): Extension<AuthIdentity>,
1537    Path(key): Path<String>,
1538) -> Response {
1539    if let AuthIdentity::Role(ref role_name) = identity {
1540        let g = state.db.read();
1541        let role_mask = match g.mask_for_role(role_name) {
1542            Ok(m) => m,
1543            Err(e) => return role_mask_err(e),
1544        };
1545        // Hidden keys must respond identically to absent keys (no oracle).
1546        if !role_mask.contains_node(&*g, &key) {
1547            return key_not_found(key);
1548        }
1549        let entries = match g.node_history(&key) {
1550            Ok(e) => e,
1551            Err(e) => return graph_err(e),
1552        };
1553        let total_commits = match g.wal_total_commits() {
1554            Ok(n) => n,
1555            Err(e) => return graph_err(e),
1556        };
1557        // Filter EdgeAdded/EdgeRemoved entries whose `other` endpoint is hidden.
1558        // A role token must not learn about hidden nodes via edge history events —
1559        // mirrors the same protection in `node_edges` (http.rs ~978-989).
1560        use core_api::HistoryChange;
1561        let visible: Vec<_> = entries
1562            .into_iter()
1563            .filter(|entry| match &entry.change {
1564                HistoryChange::EdgeAdded { other, .. }
1565                | HistoryChange::EdgeRemoved { other, .. } => role_mask.contains_node(&*g, other),
1566                _ => true,
1567            })
1568            .collect();
1569        return json_ok(node_history_json(&key, &visible, total_commits));
1570    }
1571    // Full identity: no masking. Return 404 for absent keys (consistent with
1572    // GET /node/{key} and the Role branch above).
1573    let g = state.db.read();
1574    if !g.has_node(&key) {
1575        return key_not_found(key);
1576    }
1577    let entries = match g.node_history(&key) {
1578        Ok(e) => e,
1579        Err(e) => return graph_err(e),
1580    };
1581    let total_commits = match g.wal_total_commits() {
1582        Ok(n) => n,
1583        Err(e) => return graph_err(e),
1584    };
1585    json_ok(node_history_json(&key, &entries, total_commits))
1586}
1587
1588/// `GET /history/edge?a=&b=` — return the edge lifecycle between two nodes.
1589///
1590/// Response: `{ a, b, events: [{edge_type, commit, event, rule}], total_commits }`.
1591/// Role tokens: BOTH `a` AND `b` must be visible in the role mask, otherwise
1592/// responds with 404 for the first invisible key (no existence oracle).
1593async fn edge_history_handler(
1594    State(state): State<AppState>,
1595    Extension(identity): Extension<AuthIdentity>,
1596    Query(qs): Query<BTreeMap<String, String>>,
1597) -> Response {
1598    let a = match qs.get("a").filter(|s| !s.is_empty()) {
1599        Some(s) => s.clone(),
1600        None => return err_response("missing query param a"),
1601    };
1602    let b = match qs.get("b").filter(|s| !s.is_empty()) {
1603        Some(s) => s.clone(),
1604        None => return err_response("missing query param b"),
1605    };
1606    if let AuthIdentity::Role(ref role_name) = identity {
1607        let g = state.db.read();
1608        let role_mask = match g.mask_for_role(role_name) {
1609            Ok(m) => m,
1610            Err(e) => return role_mask_err(e),
1611        };
1612        // BOTH endpoints must be visible (no oracle for either).
1613        if !role_mask.contains_node(&*g, &a) {
1614            return key_not_found(a);
1615        }
1616        if !role_mask.contains_node(&*g, &b) {
1617            return key_not_found(b);
1618        }
1619        let result = match g.edge_history(&a, &b) {
1620            Ok(r) => r,
1621            Err(e) => return graph_err(e),
1622        };
1623        return json_ok(edge_history_result_json(&a, &b, &result));
1624    }
1625    // Full identity: no masking.
1626    let g = state.db.read();
1627    let result = match g.edge_history(&a, &b) {
1628        Ok(r) => r,
1629        Err(e) => return graph_err(e),
1630    };
1631    json_ok(edge_history_result_json(&a, &b, &result))
1632}
1633
1634/// `GET /history/was_linked?a=&b=&edge_type=&at_commit=` — point-in-time edge check.
1635///
1636/// Response: `{ a, b, edge_type, at_commit, linked }`.
1637/// Returns 400 (not 500) when `at_commit` is outside the visible horizon.
1638/// Role tokens: BOTH `a` AND `b` must be visible, same-as-absent otherwise.
1639async fn was_linked_handler(
1640    State(state): State<AppState>,
1641    Extension(identity): Extension<AuthIdentity>,
1642    Query(qs): Query<BTreeMap<String, String>>,
1643) -> Response {
1644    let a = match qs.get("a").filter(|s| !s.is_empty()) {
1645        Some(s) => s.clone(),
1646        None => return err_response("missing query param a"),
1647    };
1648    let b = match qs.get("b").filter(|s| !s.is_empty()) {
1649        Some(s) => s.clone(),
1650        None => return err_response("missing query param b"),
1651    };
1652    let edge_type = match qs.get("edge_type").filter(|s| !s.is_empty()) {
1653        Some(s) => s.clone(),
1654        None => return err_response("missing query param edge_type"),
1655    };
1656    let at_commit: u64 = match qs.get("at_commit") {
1657        Some(s) => match s.parse() {
1658            Ok(n) => n,
1659            Err(_) => return err_response("at_commit must be a non-negative integer"),
1660        },
1661        None => return err_response("missing query param at_commit"),
1662    };
1663
1664    if let AuthIdentity::Role(ref role_name) = identity {
1665        let g = state.db.read();
1666        let role_mask = match g.mask_for_role(role_name) {
1667            Ok(m) => m,
1668            Err(e) => return role_mask_err(e),
1669        };
1670        if !role_mask.contains_node(&*g, &a) {
1671            return key_not_found(a);
1672        }
1673        if !role_mask.contains_node(&*g, &b) {
1674            return key_not_found(b);
1675        }
1676        return match g.was_linked(&a, &b, &edge_type, at_commit) {
1677            Ok(linked) => json_ok(json!({
1678                "a": a, "b": b, "edge_type": edge_type,
1679                "at_commit": at_commit, "linked": linked,
1680            })),
1681            Err(GraphError::CommitOutOfRange { .. }) => (
1682                StatusCode::BAD_REQUEST,
1683                Json(json!({"error": format!("commit {at_commit} is out of range")})),
1684            )
1685                .into_response(),
1686            Err(e) => graph_err(e),
1687        };
1688    }
1689
1690    // Full identity: no masking.
1691    let g = state.db.read();
1692    match g.was_linked(&a, &b, &edge_type, at_commit) {
1693        Ok(linked) => json_ok(json!({
1694            "a": a, "b": b, "edge_type": edge_type,
1695            "at_commit": at_commit, "linked": linked,
1696        })),
1697        Err(GraphError::CommitOutOfRange { .. }) => (
1698            StatusCode::BAD_REQUEST,
1699            Json(json!({"error": format!("commit {at_commit} is out of range")})),
1700        )
1701            .into_response(),
1702        Err(e) => graph_err(e),
1703    }
1704}
1705
1706/// `DELETE /node/{key}/prop/{field}` — remove a property via the group-commit queue.
1707async fn remove_node_prop(
1708    State(state): State<AppState>,
1709    Extension(identity): Extension<AuthIdentity>,
1710    Path((key, field)): Path<(String, String)>,
1711) -> Response {
1712    if let AuthIdentity::Role(_) = identity {
1713        return forbidden("role-bound token: writes are not permitted");
1714    }
1715    let db = state.db.clone();
1716    match blocking_write(move || db.submit_batch(vec![BatchOp::RemoveProp { key, field }])).await {
1717        Ok(_) => json_ok(json!({"ok": true})),
1718        Err(resp) => resp,
1719    }
1720}
1721
1722/// `POST /backup` — take a consistent backup of the database to `dest`.
1723///
1724/// The read guard is held for the duration of the file copies, which is the
1725/// correct cross-process synchronisation point: the server is the single
1726/// process touching the files, so holding the read lock excludes concurrent
1727/// in-process writers.  This is the safe alternative to running the
1728/// `mushroomdb backup` CLI against a live-served store.
1729///
1730/// Request body: `{"dest": "/absolute/path/to/backup-dir"}`
1731///
1732/// Responses:
1733/// - `200 OK` — backup completed; body is a `BackupReport` JSON object.
1734/// - `500 Internal Server Error` — backup succeeded but verification failed;
1735///   body is the `BackupReport` JSON object (examine `files` and `bytes`).
1736/// - `400 Bad Request` — missing or invalid `dest`.
1737/// - `403 Forbidden` — role-bound token; this endpoint requires a full-access token.
1738async fn backup(
1739    State(state): State<AppState>,
1740    Extension(identity): Extension<AuthIdentity>,
1741    Json(body): Json<Js>,
1742) -> Response {
1743    if let AuthIdentity::Role(_) = identity {
1744        return forbidden("role-bound token: /backup requires a full-access token");
1745    }
1746    let dest = match body.get("dest").and_then(Js::as_str) {
1747        Some(s) if !s.is_empty() => std::path::PathBuf::from(s),
1748        _ => return err_response("missing or empty \"dest\" field"),
1749    };
1750    let db = state.db.clone();
1751    let report: BackupReport = match tokio::task::spawn_blocking(move || {
1752        // The read guard is held inside spawn_blocking so file copies happen
1753        // with the write lock excluded.  The guard drops at end of closure.
1754        let g = db.read();
1755        g.backup_to(&dest)
1756    })
1757    .await
1758    {
1759        Ok(Ok(r)) => r,
1760        Ok(Err(e)) => return graph_err(e),
1761        Err(_) => return err_response("backup task panicked"),
1762    };
1763
1764    let body = match serde_json::to_value(BackupReportJson::from(&report)) {
1765        Ok(v) => v,
1766        Err(e) => return err_response(e.to_string()),
1767    };
1768
1769    if report.verified {
1770        json_ok(body)
1771    } else {
1772        (StatusCode::INTERNAL_SERVER_ERROR, Json(body)).into_response()
1773    }
1774}
1775
1776/// JSON-serialisable projection of [`BackupReport`].
1777#[derive(serde::Serialize)]
1778struct BackupReportJson<'a> {
1779    files: &'a [String],
1780    bytes: u64,
1781    verified: bool,
1782}
1783
1784impl<'a> From<&'a BackupReport> for BackupReportJson<'a> {
1785    fn from(r: &'a BackupReport) -> Self {
1786        Self {
1787            files: &r.files,
1788            bytes: r.bytes,
1789            verified: r.verified,
1790        }
1791    }
1792}
1793
1794#[cfg(test)]
1795mod tests {
1796    use super::*;
1797    use crate::json::result_set_json;
1798    use core_api::{DegreeConfig, PageRankConfig, ResultSet, Value, WccConfig};
1799
1800    #[test]
1801    fn nan_float_cell_serializes_as_null() {
1802        let mut rs = ResultSet::new(vec!["n".into()]);
1803        rs.push_row(vec![Some(Value::Float(f64::NAN))]);
1804        let j = result_set_json(&rs);
1805        assert_eq!(j["rows"][0][0], Js::Null);
1806    }
1807
1808    /// Verify that `POST /algo/pagerank` accepts an empty JSON body `{}` and
1809    /// applies server defaults (regression guard for `#[serde(default)]`).
1810    #[test]
1811    fn pagerank_config_empty_body_uses_defaults() {
1812        let config: PageRankConfig = serde_json::from_str("{}").unwrap();
1813        let default = PageRankConfig::default();
1814        assert_eq!(config.damping, default.damping);
1815        assert_eq!(config.max_iters, default.max_iters);
1816        assert_eq!(config.tol, default.tol);
1817        assert_eq!(config.budget_ms, default.budget_ms);
1818        assert_eq!(config.edge_type, default.edge_type);
1819    }
1820
1821    /// Same guard for `POST /algo/wcc`.
1822    #[test]
1823    fn wcc_config_empty_body_uses_defaults() {
1824        let config: WccConfig = serde_json::from_str("{}").unwrap();
1825        let default = WccConfig::default();
1826        assert_eq!(config.budget_ms, default.budget_ms);
1827        assert_eq!(config.edge_type, default.edge_type);
1828    }
1829
1830    /// Same guard for `POST /algo/degree`.
1831    #[test]
1832    fn degree_config_empty_body_uses_defaults() {
1833        let config: DegreeConfig = serde_json::from_str("{}").unwrap();
1834        let default = DegreeConfig::default();
1835        assert_eq!(config.budget_ms, default.budget_ms);
1836        assert_eq!(config.edge_type, default.edge_type);
1837    }
1838}