Skip to main content

lex_api/
handlers.rs

1//! Request routing for the agent API.
2//!
3//! Each handler is a synchronous function that returns
4//! `Result<serde_json::Value, ApiError>`. The dispatcher wraps the result
5//! in an HTTP response — successes as 200 with the JSON body, structured
6//! errors as 4xx/5xx with a JSON envelope.
7
8use indexmap::IndexMap;
9use lex_ast::canonicalize_program;
10use lex_bytecode::{compile_program, vm::Vm, Value};
11use lex_runtime::{check_program as check_policy, DefaultHandler, Policy};
12use lex_store::Store;
13use crate::publish_examples::record_examples_for_publish;
14use lex_syntax::{load_package, load_program_from_str, Manifest};
15use lex_vcs::{MergeSession, MergeSessionId};
16use serde::{Deserialize, Serialize};
17use std::collections::{BTreeMap, BTreeSet, HashMap};
18use std::path::PathBuf;
19use std::sync::{Arc, Mutex};
20use std::time::{SystemTime, UNIX_EPOCH};
21use tiny_http::{Header, Method, Request, Response};
22
23/// The function declarations in a canonicalized program, by name.
24fn stage_fns(stages: &[lex_ast::Stage]) -> BTreeMap<String, lex_ast::FnDecl> {
25    stages.iter().filter_map(|s| match s {
26        lex_ast::Stage::FnDecl(fd) => Some((fd.name.clone(), fd.clone())),
27        _ => None,
28    }).collect()
29}
30
31/// The type declarations in a canonicalized program, by name — so the
32/// publish diff captures `type`s alongside functions (#895).
33fn stage_types(stages: &[lex_ast::Stage]) -> BTreeMap<String, lex_ast::TypeDecl> {
34    stages.iter().filter_map(|s| match s {
35        lex_ast::Stage::TypeDecl(td) => Some((td.name.clone(), td.clone())),
36        _ => None,
37    }).collect()
38}
39
40pub struct State {
41    pub store: Mutex<Store>,
42    /// Filesystem root of the store. Held alongside the `Store`
43    /// itself so handlers that need to read store-level files
44    /// (e.g. `users.json` for actor auth) don't have to round-
45    /// trip through the lock.
46    pub root: PathBuf,
47    /// In-memory merge sessions, keyed by MergeSessionId. Sessions
48    /// are ephemeral by design (#134 foundation): they live for the
49    /// lifetime of the server process and are GC'd on commit. A
50    /// future slice can persist them to disk so a session survives
51    /// process restarts. For now an agent that gets unlucky with a
52    /// restart re-runs `merge/start` and gets a fresh session.
53    pub sessions: Mutex<HashMap<MergeSessionId, ApiMergeSession>>,
54    /// Optional server-imposed ceiling on the effect policy honored
55    /// by `/v1/run` and `/v1/replay`. `None` (the default, used by
56    /// single-tenant `lex serve`) runs the caller's request policy
57    /// as-is — the operator *is* the caller there, so that's
58    /// intended. When `Some`, the request policy is clamped via
59    /// [`clamp_policy`] so it can only *narrow* the ceiling, never
60    /// widen it.
61    ///
62    /// Any embedder that exposes this API to untrusted callers — a
63    /// hosted, multi-tenant gateway like lex-hub — MUST set this.
64    /// Without it the request body can grant itself `[proc]`
65    /// (arbitrary subprocess spawn), `[fs_*]` over `/`, and
66    /// unrestricted `[net]`: arbitrary code execution as the server
67    /// process. See lex-hub#6.
68    ///
69    /// NOTE: an empty scope list means "any path/host" in the
70    /// runtime, so a ceiling that puts `fs_read`/`fs_write`/`net` in
71    /// `allow_effects` MUST also populate the matching scope list
72    /// (`allow_fs_read`, …) or it re-opens the wildcard. Granting
73    /// none of those kinds is the safe default.
74    pub policy_ceiling: Option<Policy>,
75    /// Per-head op-history indexes and paged deltas behind
76    /// `/v1/ops/since` (#971), so a paged pull walks the op log once
77    /// instead of once per page. See [`crate::ops_since_http`].
78    pub(crate) ops_since: Mutex<crate::ops_since_http::OpsSinceCache>,
79    /// Optional server-imposed limits on the files-beside-the-op-log blob
80    /// space (#1007): `/v1/blobs/batch` refuses an oversize blob (413) or
81    /// one that would take the store past its quota (507), and
82    /// `/v1/ops/batch` refuses a `SetFiles` whose manifest has too many
83    /// entries. `None` (the default, single-tenant `lex serve`) is
84    /// unlimited. A hosted, multi-tenant embedder such as lex-hub should
85    /// set it — the same shape as [`policy_ceiling`](State::policy_ceiling).
86    pub blob_limits: Option<BlobLimits>,
87    /// `produced_by.tool` names that clients may NOT claim. Defaults to
88    /// empty (single-tenant `lex serve` and existing embedders behave
89    /// exactly as before). When non-empty, `POST /v1/attestations/batch`
90    /// refuses — `403` `ReservedProducer`, whole batch, nothing written —
91    /// any attestation whose `produced_by.tool` matches an entry. An entry
92    /// is an exact name, or — when it ends in `*` — a prefix
93    /// (`"lex-store::review:*"` reserves every `lex-store::review:<who>`
94    /// producer, whose suffix is variable). Matching trims and ASCII
95    /// case-folds both sides; blank entries are ignored.
96    ///
97    /// The point is to keep a name that only the server writes (the hub's
98    /// own `lex-hub-ci`, see `lex_store::HUB_CI_PRODUCER_TOOL`, and the
99    /// `lex-store::review:*` family, see `lex_store::REVIEW_PRODUCER_RESERVATION`)
100    /// from being minted by a tenant key holder. Server-internal writers
101    /// (`Store::verify_head_and_attest`, `record_review`, …) call the
102    /// store directly and are unaffected. A hosted embedder such as
103    /// lex-hub should set it; there is deliberately no default name here.
104    pub reserved_producers: Vec<String>,
105    /// Attestation KINDS that clients may NOT file (#1066). Defaults to
106    /// empty (single-tenant `lex serve` and existing embedders behave
107    /// exactly as before). When non-empty, `POST /v1/attestations/batch`
108    /// refuses — `403` `ReservedKind`, whole batch, nothing written — any
109    /// attestation whose kind matches an entry.
110    ///
111    /// This is the counterpart of [`reserved_producers`](State::reserved_producers)
112    /// for readers that key on the KIND rather than on who produced it:
113    /// `Store::latest_review_verdict` (the review inbox, `promote`'s
114    /// "standing Reject") takes the latest `Review` on a stage whatever its
115    /// `produced_by.tool`, so reserving the `lex-store::review:*` producer
116    /// family alone does not protect verdicts — a client files a `Review`
117    /// under `evil-tool`. An embedder that stamps reviewer identity
118    /// server-side reserves [`REVIEW_KIND`] as well.
119    ///
120    /// An entry is the serde tag of `lex_vcs::AttestationKind`
121    /// (`"review"`, `"type_check"`, …; see the `*_KIND` consts). Matching
122    /// is against the kind of the PARSED attestation — exactly what the
123    /// store would persist — so no body shape can store a reserved kind
124    /// without being judged as it (a body with duplicate `"kind"` keys does
125    /// not parse at all: `400`).
126    /// Matching trims and ASCII case-folds both sides, and additionally
127    /// ignores `_`, so `"Review"`, `" REVIEW "`, `"TypeCheck"` and
128    /// `"type_check"` all mean what they say; an entry ending in `*` is a
129    /// prefix (as for producers); blank entries are ignored.
130    ///
131    /// Server-internal writers (`Store::verify_head_and_attest`,
132    /// `record_review`, `POST /v1/review/verdict`, …) call the store
133    /// directly and are unaffected. Reserving a kind is not authenticity
134    /// (a signature is): reserving `TypeCheck` should wait for
135    /// signature-checked gates.
136    pub reserved_kinds: Vec<String>,
137}
138
139/// Serde tag of `AttestationKind::Review` — the value an embedder passes to
140/// [`State::with_reserved_kinds`] to keep verdicts server-stamped (#1066).
141/// Pair it with `lex_store::REVIEW_PRODUCER_RESERVATION`.
142pub const REVIEW_KIND: &str = "review";
143
144/// Serde tag of `AttestationKind::TypeCheck` (#1066). Reserving it should
145/// wait for signature-checked gates; exported so the name is not retyped.
146pub const TYPE_CHECK_KIND: &str = "type_check";
147
148/// Limits on one store's blob space (#1007). See [`State::blob_limits`].
149#[derive(Debug, Clone, Copy, PartialEq, Eq)]
150pub struct BlobLimits {
151    /// Largest single blob accepted, in decoded bytes.
152    pub max_blob_bytes: u64,
153    /// Total decoded bytes the store's blob space may hold.
154    pub store_quota_bytes: u64,
155    /// Most entries a `SetFiles` manifest may name.
156    pub max_manifest_entries: usize,
157}
158
159/// Capabilities this server advertises on `/v1/health` (#1007). A client
160/// states the ones it speaks in the `X-Lex-Caps` request header
161/// (comma-separated).
162pub const CAPS: &[&str] = &[CAP_FILES_V1, CAP_INTENT_ORIGIN_V1];
163
164/// The server stores and serves `SetFiles` ops and their blobs (#1007).
165pub const CAP_FILES_V1: &str = "files-v1";
166
167/// The server stores and serves `Intent.origin` (#892) — the external-VCS
168/// provenance of an imported intent. A server without it deserializes an
169/// origin-bearing intent, silently drops the unknown field and re-stores an
170/// intent whose bytes no longer match its id, so `lex op push` refuses to
171/// send one to a hub that doesn't advertise this.
172pub const CAP_INTENT_ORIGIN_V1: &str = "intent-origin-v1";
173
174/// Whether an `X-Lex-Caps` header value names `cap`.
175pub(crate) fn has_cap(header: Option<&str>, cap: &str) -> bool {
176    header.is_some_and(|h| h.split(',').any(|c| c.trim().eq_ignore_ascii_case(cap)))
177}
178
179/// Server-side wrapper around [`MergeSession`] carrying the
180/// branch names that started the merge. The lex-vcs session
181/// itself only tracks `OpId` heads; commit needs the dst branch
182/// name to advance the right head, and the src branch name is
183/// kept for round-trip auditability ("which branch did we merge
184/// from?").
185pub struct ApiMergeSession {
186    pub inner: MergeSession,
187    pub src_branch: String,
188    pub dst_branch: String,
189}
190
191/// The entry of `entries` that `claimed` matches: trim + ASCII case-fold
192/// (`fold`) both sides; an entry ending in `*` is a prefix; blank entries
193/// are ignored. Shared by `reserved_producers` and `reserved_kinds`.
194fn find_reserved<'a>(entries: &'a [String], claimed: &str, fold: fn(&str) -> String) -> Option<&'a str> {
195    let claimed = fold(claimed);
196    entries.iter().map(String::as_str).find(|entry| {
197        let entry = fold(entry);
198        match entry.strip_suffix('*') {
199            Some(prefix) => !prefix.is_empty() && claimed.starts_with(prefix),
200            None => !entry.is_empty() && claimed == entry,
201        }
202    })
203}
204
205fn fold_name(s: &str) -> String {
206    s.trim().to_ascii_lowercase()
207}
208
209/// Like [`fold_name`], and `_` is dropped so `TypeCheck` and `type_check`
210/// (the serde tag) are the same kind.
211fn fold_kind(s: &str) -> String {
212    s.trim().chars().filter(|c| *c != '_').map(|c| c.to_ascii_lowercase()).collect()
213}
214
215/// The serde tag of `kind` — the exact `"kind"` string the store persists.
216fn attestation_kind_tag(kind: &lex_vcs::AttestationKind) -> String {
217    serde_json::to_value(kind)
218        .ok()
219        .and_then(|v| v.get("kind").and_then(|t| t.as_str().map(str::to_string)))
220        .unwrap_or_default()
221}
222
223impl State {
224    pub fn open(root: PathBuf) -> anyhow::Result<Self> {
225        Self::open_with_ceiling(root, None)
226    }
227
228    /// Like [`State::open`] but installs a [`policy_ceiling`](State::policy_ceiling)
229    /// that `/v1/run` and `/v1/replay` clamp the caller's request
230    /// policy against. Embedders exposing this API to untrusted
231    /// callers must use this constructor (or set the field directly).
232    pub fn open_with_ceiling(
233        root: PathBuf,
234        policy_ceiling: Option<Policy>,
235    ) -> anyhow::Result<Self> {
236        Ok(Self {
237            store: Mutex::new(Store::open(&root)?),
238            root,
239            sessions: Mutex::new(HashMap::new()),
240            policy_ceiling,
241            ops_since: Mutex::new(Default::default()),
242            blob_limits: None,
243            reserved_producers: Vec::new(),
244            reserved_kinds: Vec::new(),
245        })
246    }
247
248    /// Install [`blob_limits`](State::blob_limits) (#1007).
249    pub fn with_blob_limits(mut self, limits: Option<BlobLimits>) -> Self {
250        self.blob_limits = limits;
251        self
252    }
253
254    /// Install [`reserved_producers`](State::reserved_producers): the
255    /// `produced_by.tool` names clients may not claim through the
256    /// attestation-writing HTTP endpoints.
257    pub fn with_reserved_producers(mut self, tools: Vec<String>) -> Self {
258        self.reserved_producers = tools;
259        self
260    }
261
262    /// Install [`reserved_kinds`](State::reserved_kinds): the attestation
263    /// kinds clients may not file through `POST /v1/attestations/batch`.
264    pub fn with_reserved_kinds(mut self, kinds: Vec<String>) -> Self {
265        self.reserved_kinds = kinds;
266        self
267    }
268
269    /// The reserved producer name `att` claims, if any.
270    fn reserved_producer_claimed(&self, att: &lex_vcs::Attestation) -> Option<&str> {
271        find_reserved(&self.reserved_producers, &att.produced_by.tool, fold_name)
272    }
273
274    /// The reserved kind entry `att`'s kind matches, if any, with the
275    /// kind's serde tag. Judges the PARSED kind (what the store persists).
276    fn reserved_kind_claimed(&self, att: &lex_vcs::Attestation) -> Option<(&str, String)> {
277        if self.reserved_kinds.is_empty() {
278            return None;
279        }
280        let tag = attestation_kind_tag(&att.kind);
281        find_reserved(&self.reserved_kinds, &tag, fold_kind).map(|entry| (entry, tag))
282    }
283
284    /// Construct a per-tenant `State` by prefixing `store_root` with the
285    /// tenant id. Single-tenant `lex serve` is unaffected — it calls
286    /// `State::open` directly.
287    ///
288    /// `tenant_id` is restricted to `[A-Za-z0-9_-]{1,64}`: anything else
289    /// (path separators, `..`, NUL, absolute paths, dotfiles, empty
290    /// string) is rejected before touching the filesystem. Without this
291    /// `PathBuf::join("/etc")` would silently replace `store_root`, and
292    /// `PathBuf::join("../foo")` would escape the tenant root.
293    pub fn new_with_tenant(tenant_id: &str, store_root: PathBuf) -> anyhow::Result<Self> {
294        validate_tenant_id(tenant_id)?;
295        Self::open(store_root.join(tenant_id))
296    }
297
298    /// Multi-tenant constructor that also installs a policy ceiling
299    /// for `/v1/run` / `/v1/replay`. The path-traversal guard from
300    /// [`new_with_tenant`](State::new_with_tenant) and the effect
301    /// ceiling are the two halves a hosted gateway needs.
302    pub fn new_with_tenant_and_ceiling(
303        tenant_id: &str,
304        store_root: PathBuf,
305        policy_ceiling: Option<Policy>,
306    ) -> anyhow::Result<Self> {
307        validate_tenant_id(tenant_id)?;
308        Self::open_with_ceiling(store_root.join(tenant_id), policy_ceiling)
309    }
310}
311
312/// Clamp a caller-supplied [`Policy`] to a server-imposed `ceiling`
313/// so it can only *narrow* the granted capabilities, never widen
314/// them. Used by [`run_handler`] when [`State::policy_ceiling`] is
315/// set — i.e. when an embedder exposes `/v1/run` to untrusted
316/// callers and must not let the request body grant itself `[proc]`,
317/// arbitrary `[fs_*]` paths, or unrestricted `[net]`.
318///
319/// - **Effects**: set-intersection of request and ceiling. The
320///   caller may drop effects but never add one the ceiling withheld.
321/// - **Scopes** (fs paths, proc binaries, net hosts): taken from the
322///   ceiling outright. The caller cannot widen them, and — because an
323///   empty scope list means "any" in the runtime — we must not let a
324///   caller's empty list collapse the ceiling's restriction back to a
325///   wildcard.
326/// - **Budget**: the more restrictive (smaller) of the two.
327fn clamp_policy(requested: Policy, ceiling: &Policy) -> Policy {
328    let allow_effects: BTreeSet<String> = requested
329        .allow_effects
330        .intersection(&ceiling.allow_effects)
331        .cloned()
332        .collect();
333    let budget = match (requested.budget, ceiling.budget) {
334        (Some(r), Some(c)) => Some(r.min(c)),
335        (None, Some(c)) => Some(c),
336        (Some(r), None) => Some(r),
337        (None, None) => None,
338    };
339    Policy {
340        allow_effects,
341        allow_fs_read: ceiling.allow_fs_read.clone(),
342        allow_fs_write: ceiling.allow_fs_write.clone(),
343        allow_net_host: ceiling.allow_net_host.clone(),
344        allow_proc: ceiling.allow_proc.clone(),
345        allow_approval: ceiling.allow_approval.clone(),
346        budget,
347    }
348}
349
350fn validate_tenant_id(tenant_id: &str) -> anyhow::Result<()> {
351    if tenant_id.is_empty() {
352        anyhow::bail!("tenant_id must not be empty");
353    }
354    if tenant_id.len() > 64 {
355        anyhow::bail!("tenant_id must be at most 64 bytes");
356    }
357    if !tenant_id
358        .bytes()
359        .all(|b| b.is_ascii_alphanumeric() || b == b'_' || b == b'-')
360    {
361        anyhow::bail!(
362            "tenant_id {tenant_id:?} contains characters outside [A-Za-z0-9_-]"
363        );
364    }
365    Ok(())
366}
367
368#[derive(Debug, Serialize, Deserialize)]
369struct ErrorEnvelope {
370    error: String,
371    #[serde(skip_serializing_if = "Option::is_none")]
372    detail: Option<serde_json::Value>,
373}
374
375pub(crate) fn json_response(status: u16, body: &serde_json::Value) -> Response<std::io::Cursor<Vec<u8>>> {
376    let bytes = serde_json::to_vec(body).unwrap_or_else(|_| b"{}".to_vec());
377    Response::from_data(bytes)
378        .with_status_code(status)
379        .with_header(Header::from_bytes(&b"Content-Type"[..], &b"application/json"[..]).unwrap())
380}
381
382pub(crate) fn error_response(status: u16, msg: impl Into<String>) -> Response<std::io::Cursor<Vec<u8>>> {
383    json_response(status, &serde_json::to_value(ErrorEnvelope {
384        error: msg.into(), detail: None,
385    }).unwrap())
386}
387
388pub(crate) fn error_with_detail(status: u16, msg: impl Into<String>, detail: serde_json::Value)
389    -> Response<std::io::Cursor<Vec<u8>>>
390{
391    json_response(status, &serde_json::to_value(ErrorEnvelope {
392        error: msg.into(), detail: Some(detail),
393    }).unwrap())
394}
395
396/// The actionable next step for a refused unsatisfiable pair (#992). Shared
397/// with the `lex op push` client, which prints it verbatim.
398pub const UNSATISFIABLE_PAIR_HINT: &str =
399    "republish from source to retire the stranded entry (#995)";
400
401/// `StoreError::UnsatisfiablePair` → **422**, never 500 (#992). The request
402/// asked to move a head onto a `(sig, stage)` pair no store can hold — a
403/// problem with the client's data, not the server. The body names the pair,
404/// the sig that actually owns the stage, and what to do about it:
405///
406/// ```json
407/// { "error": "UnsatisfiablePair",
408///   "detail": { "sig_id": "...", "stage_id": "...", "filed_under": "...",
409///               "hint": "republish from source to retire the stranded entry (#995)" } }
410/// ```
411///
412/// `None` for any other error, so callers fall through to their own mapping.
413pub(crate) fn unsatisfiable_pair_response(err: &lex_store::StoreError)
414    -> Option<Response<std::io::Cursor<Vec<u8>>>>
415{
416    let lex_store::StoreError::UnsatisfiablePair { sig_id, stage_id, filed_under } = err else {
417        return None;
418    };
419    Some(error_with_detail(422, "UnsatisfiablePair", serde_json::json!({
420        "sig_id": sig_id,
421        "stage_id": stage_id,
422        "filed_under": filed_under,
423        "message": err.to_string(),
424        "hint": UNSATISFIABLE_PAIR_HINT,
425    })))
426}
427
428/// Map a `StoreError` from a write path (`apply_operation` /
429/// `apply_operation_checked`) to an HTTP response. The only special
430/// case today is `Contention` (#262 multi-writer CAS retries
431/// exhausted), which maps to 503 with a `Retry-After` header so
432/// clients back off rather than hammering the same branch tip.
433pub(crate) fn write_error_response(prefix: &str, err: lex_store::StoreError)
434    -> Response<std::io::Cursor<Vec<u8>>>
435{
436    if let Some(resp) = unsatisfiable_pair_response(&err) {
437        return resp;
438    }
439    if let lex_store::StoreError::Contention { branch, attempts } = &err {
440        let body = serde_json::to_vec(&ErrorEnvelope {
441            error: format!("{prefix}: branch '{branch}' is contended (attempts={attempts})"),
442            detail: Some(serde_json::json!({
443                "kind": "contention",
444                "branch": branch,
445                "attempts": attempts,
446            })),
447        }).unwrap_or_else(|_| b"{}".to_vec());
448        return Response::from_data(body)
449            .with_status_code(503)
450            .with_header(Header::from_bytes(&b"Content-Type"[..], &b"application/json"[..]).unwrap())
451            .with_header(Header::from_bytes(&b"Retry-After"[..], &b"1"[..]).unwrap());
452    }
453    // #292 slice 3: budget overflow → 503 with `Retry-After: 0`.
454    // Unlike Contention (where a retry might land after another
455    // writer finishes), there's no point retrying a budget-
456    // exceeded op — the caller needs to raise the cap, switch
457    // sessions, or refactor the work. The `Retry-After: 0`
458    // signals "don't bother retrying as-is" while still using
459    // the canonical "service refused" status code.
460    if let lex_store::StoreError::BudgetExceeded { session_id, cap, spent_after } = &err {
461        let body = serde_json::to_vec(&ErrorEnvelope {
462            error: format!(
463                "{prefix}: session `{session_id}` budget exceeded \
464                 (spent_after={spent_after}, cap={cap})"
465            ),
466            detail: Some(serde_json::json!({
467                "kind": "budget_exceeded",
468                "session_id": session_id,
469                "cap": cap,
470                "spent_after": spent_after,
471            })),
472        }).unwrap_or_else(|_| b"{}".to_vec());
473        return Response::from_data(body)
474            .with_status_code(503)
475            .with_header(Header::from_bytes(&b"Content-Type"[..], &b"application/json"[..]).unwrap())
476            .with_header(Header::from_bytes(&b"Retry-After"[..], &b"0"[..]).unwrap());
477    }
478    error_response(500, format!("{prefix}: {err}"))
479}
480
481pub fn handle(state: Arc<State>, mut req: Request) -> std::io::Result<()> {
482    let method = req.method().clone();
483    let url = req.url().to_string();
484    let path = url.split('?').next().unwrap_or("").to_string();
485    let query = url.split_once('?').map(|(_, q)| q.to_string()).unwrap_or_default();
486
487    // `X-Lex-User` is the v3d session identifier — set by humans
488    // operating the web UI through whatever proxy fronts auth, or
489    // by AI agents calling the JSON API. We pluck it once here so
490    // every handler can take it as a borrowed string.
491    let x_lex_user = req.headers().iter()
492        .find(|h| h.field.equiv("x-lex-user"))
493        .map(|h| h.value.as_str().to_string());
494    // `X-Lex-Caps` (#1007): the capabilities the client speaks, so a route
495    // can refuse to hand an old client history it would mis-store.
496    let x_lex_caps = req.headers().iter()
497        .find(|h| h.field.equiv("x-lex-caps"))
498        .map(|h| h.value.as_str().to_string());
499
500    // POST /v1/pkg/publish sends a raw tar.gz body — read bytes before routing.
501    if matches!(method, Method::Post) && path == "/v1/pkg/publish" {
502        let mut body_bytes: Vec<u8> = Vec::new();
503        let _ = req.as_reader().read_to_end(&mut body_bytes);
504        let resp = pkg_publish_handler(&state, &body_bytes);
505        return req.respond(resp);
506    }
507
508    let mut body = String::new();
509    let _ = req.as_reader().read_to_string(&mut body);
510
511    let resp = route(&state, &method, &path, &query, &body, x_lex_user.as_deref(), x_lex_caps.as_deref());
512    req.respond(resp)
513}
514
515/// Auth-gated entry point. Calls `auth(path, headers)` before routing;
516/// returns 401 JSON when it returns false. Keeps auth logic out of the
517/// product-agnostic core.
518pub fn handle_with_auth<F>(state: Arc<State>, req: Request, auth: F) -> std::io::Result<()>
519where
520    F: FnOnce(&str, &[Header]) -> bool,
521{
522    let path = req.url().split('?').next().unwrap_or("").to_string();
523    if !auth(&path, req.headers()) {
524        return req.respond(
525            Response::from_data(br#"{"error":"unauthorized"}"#.to_vec())
526                .with_status_code(401)
527                .with_header(
528                    Header::from_bytes(&b"Content-Type"[..], &b"application/json"[..]).unwrap(),
529                ),
530        );
531    }
532    handle(state, req)
533}
534
535fn route(
536    state: &State,
537    method: &Method,
538    path: &str,
539    query: &str,
540    body: &str,
541    x_lex_user: Option<&str>,
542    x_lex_caps: Option<&str>,
543) -> Response<std::io::Cursor<Vec<u8>>> {
544    match (method, path) {
545        // ---- lex-tea v2 (HTML browser) ------------------------
546        (Method::Get, "/") => crate::web::activity_handler(state),
547        (Method::Get, "/web/branches") => crate::web::branches_handler(state),
548        (Method::Get, "/web/trust") => crate::web::trust_handler(state),
549        (Method::Get, "/web/attention") => crate::web::attention_handler(state),
550        (Method::Get, p) if p.starts_with("/web/branch/") => {
551            let name = &p["/web/branch/".len()..];
552            crate::web::branch_handler(state, name)
553        }
554        (Method::Get, p) if p.starts_with("/web/stage/") => {
555            let id = &p["/web/stage/".len()..];
556            crate::web::stage_html_handler(state, id)
557        }
558        // lex-tea v3 human-triage actions (#172). HTML forms post
559        // to /web/stage/<id>/{pin,defer,block,unblock} with a
560        // `reason` body. All four share one handler; the verb in
561        // the path picks the AttestationKind.
562        (Method::Post, p) if p.starts_with("/web/stage/") && (
563            p.ends_with("/pin") || p.ends_with("/defer")
564            || p.ends_with("/block") || p.ends_with("/unblock")
565        ) => {
566            let prefix_len = "/web/stage/".len();
567            let last_slash = p.rfind('/').unwrap_or(p.len());
568            let id = &p[prefix_len..last_slash];
569            let verb = &p[last_slash + 1..];
570            let decision = match verb {
571                "pin"     => crate::web::WebStageDecision::Pin,
572                "defer"   => crate::web::WebStageDecision::Defer,
573                "block"   => crate::web::WebStageDecision::Block,
574                "unblock" => crate::web::WebStageDecision::Unblock,
575                _ => unreachable!("matched in outer guard"),
576            };
577            crate::web::stage_decision_handler(state, id, body, decision, x_lex_user)
578        }
579        // ---- JSON API -----------------------------------------
580        (Method::Get, "/v1/health") => json_response(200, &serde_json::json!({"ok": true, "caps": CAPS})),
581        (Method::Post, "/v1/parse") => parse_handler(body),
582        (Method::Post, "/v1/check") => check_handler(body),
583        (Method::Post, "/v1/publish") => publish_handler(state, body),
584        (Method::Post, "/v1/patch") => patch_handler(state, body),
585        (Method::Post, "/v1/transform") => crate::transform_http::transform_handler(state, body),
586        (Method::Get, p) if p.starts_with("/v1/stage/") => {
587            let suffix = &p["/v1/stage/".len()..];
588            // Match `/v1/stage/<id>/attestations` first so a literal
589            // stage_id of "attestations" can't be misrouted.
590            if let Some(id) = suffix.strip_suffix("/attestations") {
591                stage_attestations_handler(state, id)
592            } else {
593                stage_handler(state, suffix)
594            }
595        }
596        (Method::Post, "/v1/run") => run_handler(state, body, false),
597        (Method::Post, "/v1/replay") => run_handler(state, body, true),
598        (Method::Get, p) if p.starts_with("/v1/trace/") => {
599            let id = &p["/v1/trace/".len()..];
600            trace_handler(state, id)
601        }
602        (Method::Get, "/v1/diff") => diff_handler(state, query),
603        (Method::Post, "/v1/merge/start") => merge_start_handler(state, body),
604        (Method::Post, p) if p.starts_with("/v1/merge/") && p.ends_with("/resolve") => {
605            let id = &p["/v1/merge/".len()..p.len() - "/resolve".len()];
606            merge_resolve_handler(state, id, body)
607        }
608        (Method::Post, p) if p.starts_with("/v1/merge/") && p.ends_with("/commit") => {
609            let id = &p["/v1/merge/".len()..p.len() - "/commit".len()];
610            merge_commit_handler(state, id)
611        }
612        // ---- #242: append-only sync of op log + attestation log
613        (Method::Post, "/v1/ops/batch") => ops_batch_handler(state, body),
614        (Method::Post, "/v1/attestations/batch") => attestations_batch_handler(state, body),
615        // Content half of push/pull: the stage (code) and intent blobs the
616        // op records reference. `batch` receives, `fetch` returns by id.
617        (Method::Post, "/v1/stages/batch") => crate::sync_http::stages_batch_handler(state, body),
618        (Method::Post, "/v1/stages/fetch") => crate::sync_http::stages_fetch_handler(state, body),
619        (Method::Post, "/v1/stages/missing") => crate::sync_http::stages_missing_handler(state, body),
620        (Method::Post, "/v1/intents/batch") => crate::sync_http::intents_batch_handler(state, body),
621        (Method::Post, "/v1/intents/fetch") => crate::sync_http::intents_fetch_handler(state, body),
622        // #930 P2b-1: committed lockfiles travel with the package so the
623        // write-time gate can resolve a head's pinned dependencies.
624        (Method::Post, "/v1/locks/batch") => crate::sync_http::locks_batch_handler(state, body),
625        (Method::Post, "/v1/locks/fetch") => crate::sync_http::locks_fetch_handler(state, body),
626        // #1007: files beside the op-log — content-addressed blobs a
627        // `SetFiles` manifest names. Pushed before the ops that need them.
628        (Method::Post, "/v1/blobs/missing") => crate::sync_http::blobs_missing_handler(state, body),
629        (Method::Post, "/v1/blobs/batch") => crate::sync_http::blobs_batch_handler(state, body),
630        (Method::Post, "/v1/blobs/fetch") => crate::sync_http::blobs_fetch_handler(state, body),
631        // #949 phase 1: typed issues travel with the package (content-
632        // addressed, like intents) so a pulled op-log carries its work items.
633        (Method::Post, "/v1/issues/batch") => crate::sync_http::issues_batch_handler(state, body),
634        (Method::Post, "/v1/issues/fetch") => crate::sync_http::issues_fetch_handler(state, body),
635        (Method::Get, "/v1/issues/list") => crate::sync_http::issues_list_handler(state),
636        // #949 phase 3: derived issue/project state, computed from the log —
637        // a board is a view, nobody drags cards. The `/list` literal above
638        // wins over the `/v1/issues/<id>` prefix arm below.
639        (Method::Get, "/v1/issues") => crate::issues_http::issues_state_handler(state),
640        (Method::Get, "/v1/projects") => crate::issues_http::projects_handler(state),
641        (Method::Get, p) if p.starts_with("/v1/issues/") => {
642            crate::issues_http::issue_detail_handler(state, &p["/v1/issues/".len()..])
643        }
644        // ---- #839 follow-up: branch management over HTTP so a remote
645        // client can create/switch branches (and thus drive the merge
646        // gates end to end), not just probe heads.
647        (Method::Get, "/v1/review/inbox") => crate::review_http::review_inbox_handler(state, query),
648        (Method::Post, "/v1/review/verdict") => crate::review_http::review_verdict_handler(state, body),
649        (Method::Get, "/v1/branches") => crate::branches_http::branches_list_handler(state),
650        (Method::Post, "/v1/branches") => crate::branches_http::branch_create_handler(state, body),
651        (Method::Post, p) if p.starts_with("/v1/branches/") && p.ends_with("/checkout") => {
652            let name = &p["/v1/branches/".len()..p.len() - "/checkout".len()];
653            crate::branches_http::branch_checkout_handler(state, name)
654        }
655        // Branch head: GET probes it (for `op push`'s delta), POST advances
656        // it (the ref half of push, fast-forward-only). Both in branches_http.
657        (Method::Get, p) if p.starts_with("/v1/branches/") && p.ends_with("/head") => {
658            let name = &p["/v1/branches/".len()..p.len() - "/head".len()];
659            crate::branches_http::branch_head_handler(state, name)
660        }
661        (Method::Post, p) if p.starts_with("/v1/branches/") && p.ends_with("/head") => {
662            let name = &p["/v1/branches/".len()..p.len() - "/head".len()];
663            crate::branches_http::branch_advance_head_handler(state, name, body)
664        }
665        // ---- #260: append-only fetch (inverse of #242 push)
666        // Body is a JSON array of OperationRecords reachable from
667        // `branch.head_op` but not from `after`, oldest-first.
668        (Method::Get, "/v1/ops/since") => crate::ops_since_http::ops_since_handler(state, query, x_lex_caps),
669        (Method::Get, "/v1/attestations/since") => attestations_since_handler(state, query),
670        // ---- #4: package concept ----------------------------------
671        // POST /v1/pkg/publish is handled in handle() before route()
672        // (binary body), so it doesn't appear here.
673        (Method::Get, "/v1/pkg") => pkg_list_handler(state),
674        // Owner-only visibility toggle (authed via the front door). Must
675        // precede the generic `/v1/pkg/{name}` arms; it's a PUT, so it
676        // can't collide with the GET/DELETE arms regardless.
677        (Method::Put, p) if p.starts_with("/v1/pkg/") && p.ends_with("/visibility") => {
678            let name = &p["/v1/pkg/".len()..p.len() - "/visibility".len()];
679            pkg_set_visibility_handler(state, name, body)
680        }
681        // Cut an immutable versioned release of an op-log-hosted package
682        // (#893). POST, before the generic /v1/pkg/{name} arms.
683        (Method::Post, p) if p.starts_with("/v1/pkg/") && p.ends_with("/release") => {
684            let name = &p["/v1/pkg/".len()..p.len() - "/release".len()];
685            pkg_release_handler(state, name, body)
686        }
687        (Method::Get, p) if p.starts_with("/v1/pkg/") && p.ends_with("/head") => {
688            let name = &p["/v1/pkg/".len()..p.len() - "/head".len()];
689            pkg_head_handler(state, name)
690        }
691        (Method::Get, p) if p.starts_with("/v1/pkg/") && p.ends_with("/versions") => {
692            let name = &p["/v1/pkg/".len()..p.len() - "/versions".len()];
693            pkg_versions_handler(state, name)
694        }
695        (Method::Get, p) if p.starts_with("/v1/pkg/") && p.ends_with("/api-diff") => {
696            let name = &p["/v1/pkg/".len()..p.len() - "/api-diff".len()];
697            pkg_api_diff_handler(state, name, query)
698        }
699        // /v1/pkg/{name}/{version}/archive — must match before the generic /{name}/{version}
700        (Method::Get, p) if p.starts_with("/v1/pkg/") && p.ends_with("/archive") => {
701            let inner = &p["/v1/pkg/".len()..p.len() - "/archive".len()];
702            // inner = "{name}/{version}"
703            if let Some((name, version)) = inner.split_once('/') {
704                pkg_archive_handler(state, name, version)
705            } else {
706                error_response(400, "expected /v1/pkg/{name}/{version}/archive")
707            }
708        }
709        // /v1/pkg/{name}/{version}
710        (Method::Get, p) if p.starts_with("/v1/pkg/") && p["/v1/pkg/".len()..].contains('/') => {
711            let inner = &p["/v1/pkg/".len()..];
712            if let Some((name, version)) = inner.split_once('/') {
713                pkg_get_version_handler(state, name, version)
714            } else {
715                error_response(400, "expected /v1/pkg/{name}/{version}")
716            }
717        }
718        (Method::Get, p) if p.starts_with("/v1/pkg/") => {
719            let name = &p["/v1/pkg/".len()..];
720            pkg_get_handler(state, name)
721        }
722        (Method::Delete, p) if p.starts_with("/v1/pkg/") => {
723            let name = &p["/v1/pkg/".len()..];
724            pkg_delete_handler(state, name)
725        }
726        _ => error_response(404, format!("unknown route: {method:?} {path}")),
727    }
728}
729
730#[derive(Deserialize)]
731struct ParseReq { source: String }
732
733fn parse_handler(body: &str) -> Response<std::io::Cursor<Vec<u8>>> {
734    let req: ParseReq = match serde_json::from_str(body) {
735        Ok(r) => r, Err(e) => return error_response(400, format!("bad request: {e}")),
736    };
737    match load_program_from_str(&req.source) {
738        Ok(prog) => {
739            let stages = canonicalize_program(&prog);
740            json_response(200, &serde_json::to_value(&stages).unwrap())
741        }
742        Err(e) => error_response(400, format!("syntax error: {e}")),
743    }
744}
745
746pub(crate) fn check_handler(body: &str) -> Response<std::io::Cursor<Vec<u8>>> {
747    let req: ParseReq = match serde_json::from_str(body) {
748        Ok(r) => r, Err(e) => return error_response(400, format!("bad request: {e}")),
749    };
750    let prog = match load_program_from_str(&req.source) {
751        Ok(p) => p, Err(e) => return error_response(400, format!("syntax error: {e}")),
752    };
753    let stages = canonicalize_program(&prog);
754    match lex_types::check_program(&stages) {
755        Ok(_) => json_response(200, &serde_json::json!({"ok": true})),
756        Err(errs) => json_response(422, &serde_json::to_value(&errs).unwrap()),
757    }
758}
759
760#[derive(Deserialize)]
761struct PublishReq { source: String, #[serde(default)] activate: bool }
762
763pub(crate) fn publish_handler(state: &State, body: &str) -> Response<std::io::Cursor<Vec<u8>>> {
764    let req: PublishReq = match serde_json::from_str(body) {
765        Ok(r) => r, Err(e) => return error_response(400, format!("bad request: {e}")),
766    };
767    let prog = match load_program_from_str(&req.source) {
768        Ok(p) => p, Err(e) => return error_response(400, format!("syntax error: {e}")),
769    };
770    // #168: rewrite stdlib parse calls to parse_strict so the
771    // bytecode emitted from these stages enforces required-field
772    // checks at runtime.
773    let mut stages = canonicalize_program(&prog);
774    if let Err(errs) = lex_types::check_and_rewrite_program(&mut stages) {
775        return error_with_detail(422, "type errors", serde_json::to_value(&errs).unwrap());
776    }
777    // #835 Tier 1: behavioral example gate. check_and_rewrite_program
778    // only type-checks `examples {}`; run them and refuse the publish
779    // if any declared example evaluates to the wrong value. Non-breaking:
780    // functions without examples (and effectful ones, which can't have
781    // them) produce no cases.
782    let example_errors = lex_runtime::evaluate_examples(&stages);
783    if !example_errors.is_empty() {
784        return error_with_detail(422, "example mismatch",
785            serde_json::to_value(&example_errors).unwrap_or_default());
786    }
787
788    let store = state.store.lock().unwrap();
789    let branch = store.current_branch();
790
791    // Compute diff between what's already on the branch and the new program.
792    let old_head = match store.branch_head(&branch) {
793        Ok(h) => h,
794        Err(e) => return error_response(500, format!("branch_head: {e}")),
795    };
796    // Fns + types (#895) on both sides. Old side is the branch head, read
797    // in one pass through the SigId the head names each stage by (#971):
798    // `get_ast` per entry re-read the whole stage index once per live
799    // declaration, and — StageIds being name-independent — resolved two
800    // functions differing only in name to one of the two, so the other was
801    // re-reported as an Add on every republish (#826).
802    let old_pairs: Vec<(String, String)> =
803        old_head.iter().map(|(sig, stg)| (sig.clone(), stg.clone())).collect();
804    let old_head_stages: Vec<lex_ast::Stage> =
805        store.get_asts_for_sigs_bulk(&old_pairs).into_iter().filter_map(Result::ok).collect();
806    let old_fns = stage_fns(&old_head_stages);
807    let new_fns = stage_fns(&stages);
808    let old_types = stage_types(&old_head_stages);
809    let new_types = stage_types(&stages);
810    let report =
811        lex_vcs::compute_diff_with_types(&old_fns, &new_fns, &old_types, &new_types, false);
812
813    // Build new imports map from any Import stages in the source.
814    let mut new_imports: lex_vcs::ImportMap = lex_vcs::ImportMap::new();
815    {
816        let entry = new_imports.entry("<source>".into()).or_default();
817        for s in &stages {
818            if let lex_ast::Stage::Import(im) = s {
819                entry.insert(lex_vcs::ImportRef {
820                    reference: im.reference.clone(),
821                    alias: im.alias.clone(),
822                });
823            }
824        }
825    }
826
827    match store.publish_program(&branch, &stages, &report, &new_imports, req.activate) {
828        Ok(outcome) => {
829            // #835 Tier 1: record the behavioral-examples verdict for each
830            // published fn-stage that declares examples. Best-effort — a
831            // failure to record must not fail an otherwise-good publish.
832            record_examples_for_publish(&store, &stages, &outcome);
833            json_response(200, &serde_json::json!({
834                "ops": outcome.ops,
835                "head_op": outcome.head_op,
836            }))
837        }
838        // The store-write gate (#130) also type-checks at the top
839        // of `publish_program`. The handler above already pre-checks,
840        // so this branch is reached only on a race or a state we
841        // didn't see at handler time. Surface the structured
842        // envelope (422) instead of a generic 500 — same shape the
843        // initial pre-check uses, so a client only has one error
844        // contract to handle.
845        Err(lex_store::StoreError::TypeError(errs)) => {
846            error_with_detail(422, "type errors", serde_json::to_value(&errs).unwrap())
847        }
848        Err(e) => write_error_response("publish_program", e),
849    }
850}
851
852#[derive(Deserialize)]
853struct PatchReq {
854    stage_id: String,
855    patch: lex_ast::Patch,
856    #[serde(default)] activate: bool,
857    /// #837 piece A: the branch to write to. Absent keeps the historical
858    /// behaviour (the server's global current branch); a supplied name is
859    /// honoured, and an unknown one is a 404.
860    #[serde(default)] branch: Option<String>,
861    /// #837 piece A: attribute the write. Absent records no intent, as before.
862    #[serde(default)] intent: Option<crate::transform_http::IntentSpec>,
863}
864
865/// POST /v1/patch — apply a structured edit to a stored stage's
866/// canonical AST, type-check the result, and publish a new stage.
867///
868/// Optional `branch` / `intent` (#837 piece A) name the branch explicitly and
869/// attribute the op; without them the request behaves exactly as it always did.
870/// For the typed transforms, and a request that *requires* a branch, see
871/// `POST /v1/transform`.
872fn patch_handler(state: &State, body: &str) -> Response<std::io::Cursor<Vec<u8>>> {
873    let req: PatchReq = match serde_json::from_str(body) {
874        Ok(r) => r, Err(e) => return error_response(400, format!("bad request: {e}")),
875    };
876    let intent = match req.intent.clone() {
877        Some(spec) => match spec.into_intent(crate::transform_http::default_http_session) {
878            Ok(i) => Some(i),
879            Err(e) => return error_response(400, format!("bad request: {e}")),
880        },
881        None => None,
882    };
883    let store = state.store.lock().unwrap();
884    let explicit_branch = req.branch.is_some();
885    if let Some(b) = &req.branch {
886        match store.list_branches() {
887            Ok(bs) if bs.iter().any(|x| x == b) => {}
888            Ok(_) => return error_response(404, format!("unknown branch `{b}`")),
889            Err(e) => return error_response(500, format!("list_branches: {e}")),
890        }
891    }
892
893    // 1. Load.
894    let original = match store.get_ast(&req.stage_id) {
895        Ok(s) => s, Err(e) => return error_response(404, format!("stage: {e}")),
896    };
897
898    // 2. Apply.
899    let patched = match lex_ast::apply_patch(&original, &req.patch) {
900        Ok(s) => s,
901        Err(e) => return error_with_detail(422, "patch failed",
902            serde_json::to_value(&e).unwrap_or_default()),
903    };
904
905    // 3. No isolated check (#833): the gated apply below type-checks
906    // the *composed* program — the branch head with the patched stage
907    // swapped in. Stricter where it matters (a body that no longer
908    // composes with its callers is refused) and correct where the old
909    // isolated check was wrong (a body calling a sibling was rejected
910    // as an unknown identifier a one-stage program couldn't see).
911
912    // Routing through the gated apply so /v1/patch participates in the
913    // op DAG. We know this op is always a body change on the existing
914    // sig (a patch can't add a brand-new fn).
915    let branch = req.branch.clone().unwrap_or_else(|| store.current_branch());
916
917    // Find the sig — patched stage's sig must match the original's.
918    let sig = match lex_ast::sig_id(&patched) {
919        Some(s) => s,
920        None => return error_response(500, "patched stage has no sig_id"),
921    };
922
923    // Persist before the gate (its RepairHint on rejection is
924    // addressed to this stage); activate only once the head moved.
925    let new_id = match store.publish(&patched) {
926        Ok(id) => id, Err(e) => return error_response(500, format!("publish: {e}")),
927    };
928
929    // Determine op kind: ChangeEffectSig if effects differ, ModifyBody otherwise.
930    let original_effects: std::collections::BTreeSet<String> = match &original {
931        lex_ast::Stage::FnDecl(fd) => fd.effects.iter().map(|e| e.name.clone()).collect(),
932        _ => std::collections::BTreeSet::new(),
933    };
934    let patched_effects: std::collections::BTreeSet<String> = match &patched {
935        lex_ast::Stage::FnDecl(fd) => fd.effects.iter().map(|e| e.name.clone()).collect(),
936        _ => std::collections::BTreeSet::new(),
937    };
938    let head_now = match store.get_branch(&branch) {
939        Ok(b) => b.and_then(|b| b.head_op),
940        Err(e) => return error_response(500, format!("get_branch: {e}")),
941    };
942    let kind = if original_effects != patched_effects {
943        // #247: budget delta is part of the canonical payload now.
944        // Patch endpoints don't currently rehydrate the AST to
945        // recompute budgets, so leave them None — clients that
946        // need budget tracking should publish through the diff
947        // pipeline (`lex publish`) where `compute_diff` populates
948        // them.
949        let from_budget = lex_vcs::operation_budget_from_effects(&original_effects);
950        let to_budget = lex_vcs::operation_budget_from_effects(&patched_effects);
951        // #992: the patched stage declares the new effects, so it hashes to a
952        // different SigId and the store files it under that one. Record it, or
953        // the head keeps the old sig pointing at a stage no store holds there.
954        let to_sig_id = lex_ast::sig_id(&patched).filter(|s| *s != sig);
955        lex_vcs::OperationKind::ChangeEffectSig {
956            sig_id: sig.clone(),
957            from_stage_id: req.stage_id.clone(),
958            to_stage_id: new_id.clone(),
959            from_effects: original_effects,
960            to_effects: patched_effects,
961            from_budget,
962            to_budget,
963            to_sig_id,
964        }
965    } else {
966        let budget = lex_vcs::operation_budget_from_effects(&original_effects);
967        lex_vcs::OperationKind::ModifyBody {
968            sig_id: sig.clone(),
969            from_stage_id: req.stage_id.clone(),
970            to_stage_id: new_id.clone(),
971            from_budget: budget,
972            to_budget: budget,
973            // #992: a patch that touches the signature (types, examples)
974            // moves the sig just as an effect change does.
975            to_sig_id: lex_ast::sig_id(&patched).filter(|s| *s != sig),
976        }
977    };
978    // #992: derive the transition from the op, never hard-code `Replace`. A
979    // sig-moving op must retire the old sig and bind the new one; a `Replace`
980    // here left `(old_sig, new_stage)` at the head — a pair no store holds.
981    let transition = lex_store::transition_for_kind(&kind);
982    let op = lex_vcs::Operation::new(
983        kind,
984        head_now.into_iter().collect::<Vec<_>>(),
985    );
986    let op = match &intent {
987        Some(i) => op.with_intent(i.intent_id.clone()),
988        None => op,
989    };
990    let op_id = match store.apply_operation_gated_with_intent(&branch, op, transition, intent.as_ref()) {
991        Ok(id) => id,
992        Err(lex_store::StoreError::TypeError(errs)) => return error_with_detail(
993            422, "type errors after patch", serde_json::to_value(&errs).unwrap_or_default()),
994        Err(e) => return write_error_response("apply_operation_gated", e),
995    };
996    if req.activate {
997        if let Err(e) = store.activate(&new_id) {
998            return error_response(500, format!("activate: {e}"));
999        }
1000    }
1001
1002    let status = format!("{:?}",
1003        store.get_status(&new_id).unwrap_or(lex_store::StageStatus::Draft)).to_lowercase();
1004    let mut resp = serde_json::json!({
1005        "old_stage_id": req.stage_id,
1006        "new_stage_id": new_id,
1007        "sig_id": sig,
1008        "status": status,
1009        "op_id": op_id,
1010    });
1011    // Only echoed when the caller opted in, so a request without `branch` /
1012    // `intent` gets the response it always got, byte for byte.
1013    if explicit_branch {
1014        resp["branch"] = serde_json::json!(branch);
1015    }
1016    if let Some(i) = &intent {
1017        resp["intent_id"] = serde_json::json!(i.intent_id);
1018    }
1019    json_response(200, &resp)
1020}
1021
1022pub(crate) fn stage_handler(state: &State, id: &str) -> Response<std::io::Cursor<Vec<u8>>> {
1023    let store = state.store.lock().unwrap();
1024    let meta = match store.get_metadata(id) {
1025        Ok(m) => m, Err(e) => return error_response(404, format!("{e}")),
1026    };
1027    let ast = match store.get_ast(id) {
1028        Ok(a) => a, Err(e) => return error_response(404, format!("{e}")),
1029    };
1030    let status = format!("{:?}", store.get_status(id).unwrap_or(lex_store::StageStatus::Draft)).to_lowercase();
1031    json_response(200, &serde_json::json!({
1032        "metadata": meta,
1033        "ast": ast,
1034        "status": status,
1035    }))
1036}
1037
1038/// `GET /v1/stage/<id>/attestations` — every persisted attestation
1039/// for this stage, newest-first by arrival order. Issue #132's
1040/// queryable-evidence consumer surface.
1041///
1042/// 404s on unknown stage_id (matches `/v1/stage/<id>`'s shape so a
1043/// caller round-tripping both endpoints sees consistent errors).
1044/// Empty list (200) is *evidence of absence*: the stage exists but
1045/// no producer has attested it.
1046pub(crate) fn stage_attestations_handler(state: &State, id: &str) -> Response<std::io::Cursor<Vec<u8>>> {
1047    let store = state.store.lock().unwrap();
1048    if let Err(e) = store.get_metadata(id) {
1049        return error_response(404, format!("{e}"));
1050    }
1051    let log = match store.attestation_log() {
1052        Ok(l) => l,
1053        Err(e) => return error_response(500, format!("attestation log: {e}")),
1054    };
1055    // Newest-first by ARRIVAL order (server-assigned), not by the
1056    // writer-supplied timestamp; legacy unstamped entries sort below
1057    // stamped ones. `lex op pull` re-puts these in reverse (oldest
1058    // first) so a puller's local arrival order matches this store's.
1059    let mut listing = match log.list_for_stage_by_arrival(&id.to_string()) {
1060        Ok(v) => v,
1061        Err(e) => return error_response(500, format!("list_for_stage: {e}")),
1062    };
1063    listing.reverse();
1064    json_response(200, &serde_json::json!({"attestations": listing}))
1065}
1066
1067#[derive(Deserialize, Default)]
1068struct PolicyJson {
1069    #[serde(default)] allow_effects: Vec<String>,
1070    #[serde(default)] allow_fs_read: Vec<String>,
1071    #[serde(default)] allow_fs_write: Vec<String>,
1072    #[serde(default)] budget: Option<u64>,
1073}
1074
1075impl PolicyJson {
1076    fn into_policy(self) -> Policy {
1077        Policy {
1078            allow_effects: self.allow_effects.into_iter().collect::<BTreeSet<_>>(),
1079            allow_fs_read: self.allow_fs_read.into_iter().map(PathBuf::from).collect(),
1080            allow_fs_write: self.allow_fs_write.into_iter().map(PathBuf::from).collect(),
1081            allow_net_host: Vec::new(),
1082            allow_proc: Vec::new(),
1083            allow_approval: Vec::new(),
1084            budget: self.budget,
1085        }
1086    }
1087}
1088
1089#[derive(Deserialize)]
1090struct RunReq {
1091    source: String,
1092    #[serde(rename = "fn")] func: String,
1093    #[serde(default)] args: Vec<serde_json::Value>,
1094    #[serde(default)] policy: PolicyJson,
1095    #[serde(default)] overrides: IndexMap<String, serde_json::Value>,
1096}
1097
1098pub(crate) fn run_handler(state: &State, body: &str, with_overrides: bool) -> Response<std::io::Cursor<Vec<u8>>> {
1099    let req: RunReq = match serde_json::from_str(body) {
1100        Ok(r) => r, Err(e) => return error_response(400, format!("bad request: {e}")),
1101    };
1102    let prog = match load_program_from_str(&req.source) {
1103        Ok(p) => p, Err(e) => return error_response(400, format!("syntax error: {e}")),
1104    };
1105    let stages = canonicalize_program(&prog);
1106    if let Err(errs) = lex_types::check_program(&stages) {
1107        return error_with_detail(422, "type errors", serde_json::to_value(&errs).unwrap());
1108    }
1109    let bc = compile_program(&stages);
1110    let mut policy = req.policy.into_policy();
1111    // When a server-imposed ceiling is present (multi-tenant
1112    // embedders like lex-hub), the request policy can only narrow
1113    // it — never grant itself proc/fs/net beyond what the operator
1114    // allowed. Single-tenant `lex serve` leaves the ceiling unset
1115    // and runs the caller's policy verbatim.
1116    if let Some(ceiling) = &state.policy_ceiling {
1117        policy = clamp_policy(policy, ceiling);
1118    }
1119    if let Err(violations) = check_policy(&bc, &policy) {
1120        return error_with_detail(403, "policy violation", serde_json::to_value(&violations).unwrap());
1121    }
1122
1123    let mut recorder = lex_trace::Recorder::new();
1124    if with_overrides && !req.overrides.is_empty() {
1125        recorder = recorder.with_overrides(req.overrides);
1126    }
1127    let handle = recorder.handle();
1128    let handler = DefaultHandler::new(policy);
1129    let mut vm = Vm::with_handler(&bc, Box::new(handler));
1130    vm.set_tracer(Box::new(recorder));
1131
1132    let vargs: Vec<Value> = req.args.iter().map(json_to_value).collect();
1133    let started = std::time::SystemTime::now().duration_since(std::time::UNIX_EPOCH).unwrap().as_secs();
1134    let result = vm.call(&req.func, vargs);
1135    let ended = std::time::SystemTime::now().duration_since(std::time::UNIX_EPOCH).unwrap().as_secs();
1136
1137    let store = state.store.lock().unwrap();
1138    let (root_out, root_err, status) = match &result {
1139        Ok(v) => (Some(value_to_json(v)), None, 200u16),
1140        Err(e) => (None, Some(format!("{e}")), 200u16),
1141    };
1142    let tree = handle.finalize(req.func.clone(), serde_json::Value::Null,
1143        root_out.clone(), root_err.clone(), started, ended);
1144    let run_id = match store.save_trace(&tree) {
1145        Ok(id) => id,
1146        Err(e) => return error_response(500, format!("save_trace: {e}")),
1147    };
1148
1149    let mut body = serde_json::json!({
1150        "run_id": run_id,
1151        "output": root_out,
1152    });
1153    if let Some(err) = root_err {
1154        body["error"] = serde_json::Value::String(err);
1155    }
1156    json_response(status, &body)
1157}
1158
1159fn trace_handler(state: &State, id: &str) -> Response<std::io::Cursor<Vec<u8>>> {
1160    let store = state.store.lock().unwrap();
1161    match store.load_trace(id) {
1162        Ok(t) => json_response(200, &serde_json::to_value(&t).unwrap()),
1163        Err(e) => error_response(404, format!("{e}")),
1164    }
1165}
1166
1167fn diff_handler(state: &State, query: &str) -> Response<std::io::Cursor<Vec<u8>>> {
1168    let mut a = None;
1169    let mut b = None;
1170    for kv in query.split('&') {
1171        if let Some((k, v)) = kv.split_once('=') {
1172            match k { "a" => a = Some(v.to_string()), "b" => b = Some(v.to_string()), _ => {} }
1173        }
1174    }
1175    let (Some(a), Some(b)) = (a, b) else {
1176        return error_response(400, "missing a or b query params");
1177    };
1178    let store = state.store.lock().unwrap();
1179    let ta = match store.load_trace(&a) { Ok(t) => t, Err(e) => return error_response(404, format!("a: {e}")) };
1180    let tb = match store.load_trace(&b) { Ok(t) => t, Err(e) => return error_response(404, format!("b: {e}")) };
1181    match lex_trace::diff_runs(&ta, &tb) {
1182        Some(d) => json_response(200, &serde_json::to_value(&d).unwrap()),
1183        None => json_response(200, &serde_json::json!({"divergence": null})),
1184    }
1185}
1186
1187fn json_to_value(v: &serde_json::Value) -> Value { Value::from_json(v) }
1188
1189fn value_to_json(v: &Value) -> serde_json::Value { v.to_json() }
1190
1191#[derive(Deserialize)]
1192struct MergeStartReq {
1193    src_branch: String,
1194    dst_branch: String,
1195}
1196
1197/// `POST /v1/merge/start` (#134) — open a stateful merge between two
1198/// branch heads and return the conflicts the agent needs to
1199/// resolve. Auto-resolved sigs (one-sided changes, identical
1200/// changes both sides) are returned for audit but don't block
1201/// commit.
1202///
1203/// Response: `{ merge_id, src_head, dst_head, lca, conflicts,
1204/// auto_resolved_count }`. The session is held in process memory
1205/// keyed by `merge_id` for subsequent `resolve` / `commit` calls
1206/// (next slices).
1207fn merge_start_handler(state: &State, body: &str) -> Response<std::io::Cursor<Vec<u8>>> {
1208    let req: MergeStartReq = match serde_json::from_str(body) {
1209        Ok(r) => r, Err(e) => return error_response(400, format!("bad request: {e}")),
1210    };
1211    let store = state.store.lock().unwrap();
1212    let src_head = match store.get_branch(&req.src_branch) {
1213        Ok(Some(b)) => b.head_op,
1214        Ok(None) => return error_response(404, format!("unknown src branch `{}`", req.src_branch)),
1215        Err(e) => return error_response(500, format!("src branch read: {e}")),
1216    };
1217    let dst_head = match store.get_branch(&req.dst_branch) {
1218        Ok(Some(b)) => b.head_op,
1219        Ok(None) => return error_response(404, format!("unknown dst branch `{}`", req.dst_branch)),
1220        Err(e) => return error_response(500, format!("dst branch read: {e}")),
1221    };
1222    let log = match lex_vcs::OpLog::open(store.root()) {
1223        Ok(l) => l,
1224        Err(e) => return error_response(500, format!("op log: {e}")),
1225    };
1226    // Caller doesn't choose merge_ids — minted server-side from
1227    // wall clock + a per-process counter avoids leaking session
1228    // ids' shape into the public surface.
1229    let merge_id = mint_merge_id();
1230    let mut session = match MergeSession::start(
1231        merge_id.clone(),
1232        &log,
1233        src_head.as_ref(),
1234        dst_head.as_ref(),
1235    ) {
1236        Ok(s) => s,
1237        Err(e) => return error_response(500, format!("merge start: {e}")),
1238    };
1239    // #1007 PR 7: the files dimension of the merge. `lex-vcs` doesn't
1240    // know about manifests (see `merge_session.rs`'s module docs), so
1241    // the store computes the 3-way diff and the session just tracks the
1242    // conflicts it surfaced — the same layering `MergeResolutionChecker`
1243    // already uses for sig-conflict type-checking.
1244    let manifest_outcome = match store.manifest_merge(
1245        session.lca.as_deref(),
1246        dst_head.as_deref(),
1247        src_head.as_deref(),
1248    ) {
1249        Ok(o) => o,
1250        Err(lex_store::StoreError::AmbiguousManifest { op_id }) => {
1251            return error_with_detail(422, "AmbiguousManifest", serde_json::json!({
1252                "op_id": op_id,
1253                "reason": "a merge ancestor has an ambiguous files manifest; \
1254                           append a set_files op resolving it before merging",
1255            }));
1256        }
1257        Err(e) => return error_response(500, format!("manifest merge: {e}")),
1258    };
1259    let (file_conflicts, needs_setfiles): (Vec<lex_vcs::FileConflict>, bool) =
1260        match manifest_outcome {
1261            lex_store::ManifestMergeOutcome::NoChange => (Vec::new(), false),
1262            lex_store::ManifestMergeOutcome::Needed { conflicts, .. } => (conflicts, true),
1263        };
1264    session.attach_file_conflicts(file_conflicts, needs_setfiles);
1265
1266    let conflicts: Vec<&lex_vcs::ConflictRecord> = session.remaining_conflicts();
1267    let auto_resolved_count = session.auto_resolved.len();
1268    let remaining_file_conflicts: Vec<&lex_vcs::FileConflict> = session.remaining_file_conflicts();
1269    let body = serde_json::json!({
1270        "merge_id": merge_id,
1271        "src_head": session.src_head,
1272        "dst_head": session.dst_head,
1273        "lca":      session.lca,
1274        "conflicts": conflicts,
1275        "auto_resolved_count": auto_resolved_count,
1276        "file_conflicts": remaining_file_conflicts,
1277        "needs_setfiles": session.needs_setfiles(),
1278    });
1279    drop(conflicts);
1280    drop(remaining_file_conflicts);
1281    drop(store);
1282    let wrapped = ApiMergeSession {
1283        inner: session,
1284        src_branch: req.src_branch,
1285        dst_branch: req.dst_branch,
1286    };
1287    state.sessions.lock().unwrap().insert(merge_id, wrapped);
1288    json_response(200, &body)
1289}
1290
1291#[derive(Deserialize)]
1292struct MergeResolveReq {
1293    /// Each entry is `(conflict_id, resolution)`. The resolution is
1294    /// the same shape as `lex_vcs::Resolution`'s tagged JSON form
1295    /// — `{"kind":"take_ours"}`, `{"kind":"take_theirs"}`,
1296    /// `{"kind":"defer"}`, or `{"kind":"custom","op":{...}}`.
1297    #[serde(default)]
1298    resolutions: Vec<MergeResolveEntry>,
1299    /// File conflicts (#1007 PR 7), keyed by path instead of sig id.
1300    /// Resolution shape is `lex_vcs::FileResolution`'s tagged JSON form
1301    /// — `{"kind":"take_ours"}`, `{"kind":"take_theirs"}`, or
1302    /// `{"kind":"defer"}` (no `custom`; see the module docs on
1303    /// `merge_session::FileResolution` for why).
1304    #[serde(default)]
1305    file_resolutions: Vec<MergeFileResolveEntry>,
1306}
1307
1308#[derive(Deserialize)]
1309struct MergeResolveEntry {
1310    conflict_id: String,
1311    resolution: lex_vcs::Resolution,
1312}
1313
1314#[derive(Deserialize)]
1315struct MergeFileResolveEntry {
1316    path: lex_vcs::FilePath,
1317    resolution: lex_vcs::FileResolution,
1318}
1319
1320/// `POST /v1/merge/<id>/resolve` (#134) — submit batched
1321/// resolutions against the conflicts surfaced by `merge/start`.
1322/// Returns one verdict per input: accepted (recorded against the
1323/// session) or rejected (with structured reason). The session
1324/// stays alive across calls so an agent can iterate.
1325///
1326/// Errors:
1327/// - 404 if `merge_id` doesn't refer to a live session (a typo
1328///   or a session GC'd by a server restart).
1329/// - 400 on malformed body.
1330fn merge_resolve_handler(
1331    state: &State,
1332    merge_id: &str,
1333    body: &str,
1334) -> Response<std::io::Cursor<Vec<u8>>> {
1335    let req: MergeResolveReq = match serde_json::from_str(body) {
1336        Ok(r) => r, Err(e) => return error_response(400, format!("bad request: {e}")),
1337    };
1338    let mut sessions = state.sessions.lock().unwrap();
1339    let Some(wrapped) = sessions.get_mut(merge_id) else {
1340        return error_response(404, format!("unknown merge_id `{merge_id}`"));
1341    };
1342    let pairs: Vec<(String, lex_vcs::Resolution)> = req.resolutions.into_iter()
1343        .map(|e| (e.conflict_id, e.resolution))
1344        .collect();
1345    // #834: type-check each resolution against dst's head at submission
1346    // time (not only at commit) so the agent's "submit N, see which
1347    // broke, retry" loop gets its feedback here. The store composes the
1348    // projected program; the session owns the merge→delta semantics.
1349    let store = state.store.lock().unwrap();
1350    let checker = lex_store::MergeResolutionChecker::new(&store, wrapped.dst_branch.clone());
1351    let verdicts = wrapped.inner.resolve_checked(pairs, &checker);
1352    drop(store);
1353    let file_pairs: Vec<(lex_vcs::FilePath, lex_vcs::FileResolution)> = req.file_resolutions
1354        .into_iter()
1355        .map(|e| (e.path, e.resolution))
1356        .collect();
1357    let file_verdicts = wrapped.inner.resolve_files(file_pairs);
1358    let remaining: Vec<&lex_vcs::ConflictRecord> = wrapped.inner.remaining_conflicts();
1359    let remaining_files: Vec<&lex_vcs::FileConflict> = wrapped.inner.remaining_file_conflicts();
1360    let body = serde_json::json!({
1361        "verdicts": verdicts,
1362        "remaining_conflicts": remaining,
1363        "file_verdicts": file_verdicts,
1364        "remaining_file_conflicts": remaining_files,
1365    });
1366    json_response(200, &body)
1367}
1368
1369/// `POST /v1/merge/<id>/commit` (#134) — finalize a merge
1370/// session. Builds a `Merge` op from the auto-resolved sigs +
1371/// the conflict resolutions, applies it to the dst branch, and
1372/// returns the new head op id. The session is dropped on
1373/// success; the caller would re-run `merge/start` to land
1374/// further changes.
1375///
1376/// Errors:
1377/// - 404: unknown `merge_id`.
1378/// - 422: conflicts remaining (pass `Defer` or just don't
1379///   resolve a conflict and you land here). Body carries the
1380///   list so the caller knows which still need attention.
1381/// - 422: a `Custom` resolution was used. The data layer
1382///   supports them but landing them via HTTP needs an extra
1383///   pass to apply the custom op against the dst branch
1384///   first; deferred to a follow-up slice. Use TakeOurs /
1385///   TakeTheirs for now.
1386/// - 409: dependency conflict (#977) — the branches pin the same
1387///   package at different versions, so no merged lock exists. Body
1388///   `detail` names the package and both versions; nothing is written.
1389/// - 500: filesystem error while landing the merge op.
1390fn merge_commit_handler(
1391    state: &State,
1392    merge_id: &str,
1393) -> Response<std::io::Cursor<Vec<u8>>> {
1394    use std::collections::BTreeMap;
1395    let wrapped = match state.sessions.lock().unwrap().remove(merge_id) {
1396        Some(w) => w,
1397        None => return error_response(404, format!("unknown merge_id `{merge_id}`")),
1398    };
1399    let dst_branch = wrapped.dst_branch.clone();
1400    let src_head = wrapped.inner.src_head.clone();
1401    let dst_head = wrapped.inner.dst_head.clone();
1402    let lca = wrapped.inner.lca.clone();
1403    let auto_resolved = wrapped.inner.auto_resolved.clone();
1404
1405    // Conflict resolutions (sig conflicts, then file conflicts — #1007 PR 7).
1406    let commit_out = match wrapped.inner.commit() {
1407        Ok(r) => r,
1408        Err(lex_vcs::CommitError::ConflictsRemaining(ids)) => {
1409            // Re-insert isn't possible since we removed above; the
1410            // caller will need to re-start. That's acceptable: a
1411            // commit-with-unresolved-conflicts is operator error.
1412            return error_with_detail(
1413                422,
1414                "conflicts remaining",
1415                serde_json::json!({"unresolved": ids}),
1416            );
1417        }
1418        Err(lex_vcs::CommitError::FileConflictsRemaining(paths)) => {
1419            return error_with_detail(
1420                422,
1421                "file conflicts remaining",
1422                serde_json::json!({"unresolved_files": paths}),
1423            );
1424        }
1425    };
1426
1427    // Translate auto-resolved + resolutions into the StageTransition::Merge
1428    // entries map. Every sig the merge decided goes in — including the ones
1429    // dst already has (`take_ours`, dst-only changes) — so the merged head
1430    // does not depend on the order a replay applies the two sides in
1431    // (#1062); see `Store::merge_pins`. `take_theirs` is read from src's
1432    // head, not reconstructed from the session (the commit consumed it).
1433    let store = state.store.lock().unwrap();
1434    let mut entries: BTreeMap<lex_vcs::SigId, Option<lex_vcs::StageId>> =
1435        match store.merge_pins(dst_head.as_ref(), src_head.as_ref(), &auto_resolved, &commit_out.resolved) {
1436            Ok(e) => e,
1437            Err(e) => return write_error_response("resolve merge entries", e),
1438        };
1439
1440    for (conflict_id, resolution) in commit_out.resolved {
1441        match resolution {
1442            // Pinned by `merge_pins` above (dst already has its head; the
1443            // `theirs` stage is src's).
1444            lex_vcs::Resolution::TakeOurs | lex_vcs::Resolution::TakeTheirs => {}
1445            lex_vcs::Resolution::Custom { op } => {
1446                // The agent's brand-new op carries the merge target
1447                // in its kind (e.g. ModifyBody.to_stage_id). The op
1448                // itself isn't separately recorded in the log here
1449                // — its head-map effect is folded into the merge
1450                // op's entries map. Callers that want the op as a
1451                // first-class history entry should publish it via
1452                // /v1/publish first and submit a TakeTheirs/TakeOurs
1453                // resolution against the resulting head.
1454                match op.kind.merge_target() {
1455                    Some((sig, stage)) => {
1456                        if sig != conflict_id {
1457                            return error_with_detail(
1458                                422,
1459                                "custom op targets a different sig than the conflict",
1460                                serde_json::json!({
1461                                    "conflict_id": conflict_id,
1462                                    "op_targets": sig,
1463                                }),
1464                            );
1465                        }
1466                        entries.insert(conflict_id, stage);
1467                    }
1468                    None => {
1469                        return error_with_detail(
1470                            422,
1471                            "custom op kind doesn't yield a single sig→stage delta",
1472                            serde_json::json!({
1473                                "conflict_id": conflict_id,
1474                                "kind": serde_json::to_value(&op.kind).unwrap_or(serde_json::Value::Null),
1475                            }),
1476                        );
1477                    }
1478                }
1479            }
1480            lex_vcs::Resolution::Defer => {
1481                // Unreachable: commit() rejects Defer above.
1482                return error_response(500, "internal: Defer slipped past commit gate");
1483            }
1484        }
1485    }
1486
1487    let resolved_count = entries.len();
1488    let mut parents: Vec<lex_vcs::OpId> = Vec::new();
1489    if let Some(d) = dst_head.clone() { parents.push(d); }
1490    if let Some(s) = src_head.clone() { parents.push(s); }
1491    let op = lex_vcs::Operation::new(
1492        lex_vcs::OperationKind::Merge { resolved: resolved_count },
1493        parents,
1494    );
1495    let transition = lex_vcs::StageTransition::Merge { entries };
1496
1497    // #1007 PR 7: if the manifests disagreed at all, build the merged
1498    // manifest now — auto-resolved paths plus every resolved file
1499    // conflict — so it can be landed as a `SetFiles` right after the
1500    // merge op. Recomputes the same 3-way diff `merge/start` already
1501    // showed the caller (over the same `lca`/`dst_head`/`src_head`, so
1502    // it reproduces byte-for-byte unless the store's blobs changed
1503    // underneath the session, which the manifest-closure check below
1504    // would catch anyway).
1505    let manifest_blob: Option<String> = if commit_out.needs_setfiles {
1506        match store.manifest_merge(lca.as_deref(), dst_head.as_deref(), src_head.as_deref()) {
1507            Ok(lex_store::ManifestMergeOutcome::NoChange) => None,
1508            Ok(lex_store::ManifestMergeOutcome::Needed { auto_entries, conflicts }) => {
1509                let file_resolutions: BTreeMap<lex_vcs::FilePath, lex_vcs::FileResolution> =
1510                    commit_out.resolved_files.into_iter().collect();
1511                match store.build_merged_manifest(
1512                    auto_entries,
1513                    &conflicts,
1514                    &file_resolutions,
1515                    "pending-merge",
1516                ) {
1517                    Ok(blob_id) => Some(blob_id),
1518                    Err(e) => return write_error_response("build merged manifest", e),
1519                }
1520            }
1521            Err(lex_store::StoreError::AmbiguousManifest { op_id }) => {
1522                return error_with_detail(422, "AmbiguousManifest", serde_json::json!({
1523                    "op_id": op_id,
1524                }));
1525            }
1526            Err(e) => return write_error_response("manifest merge", e),
1527        }
1528    } else {
1529        None
1530    };
1531
1532    // Gated (#833): lands the merge op, type-checks the real post-merge
1533    // head, rolls the head back on a TypeError — and (#1007 PR 7) also
1534    // appends the `SetFiles` above when `manifest_blob` is `Some`,
1535    // rolling all the way back if that fails too (see
1536    // `apply_merge_op_gated_with_manifest`'s doc comment).
1537    match store.apply_merge_op_gated_with_manifest(&dst_branch, op, transition, manifest_blob.as_deref(), None) {
1538        Ok(new_head_op) => json_response(200, &serde_json::json!({
1539            "new_head_op": new_head_op,
1540            "dst_branch": dst_branch,
1541        })),
1542        Err(lex_store::StoreError::TypeError(errs)) => error_with_detail(
1543            422, "merged program has type errors", serde_json::to_value(&errs).unwrap_or_default()),
1544        // #977: the two branches pin the same dependency at different
1545        // versions. The merge commits the union of both parents' locks and
1546        // refuses to pick a side silently; nothing was written, so the dst
1547        // branch is unchanged. 409: the request is fine, the branches'
1548        // states conflict — align the version on one branch and retry.
1549        Err(e @ lex_store::StoreError::DependencyConflict { .. }) => {
1550            let detail = match &e {
1551                lex_store::StoreError::DependencyConflict { package, dst_version, src_version } => {
1552                    serde_json::json!({
1553                        "kind": "dependency_conflict",
1554                        "package": package,
1555                        "dst_branch": dst_branch,
1556                        "dst_version": dst_version,
1557                        "src_branch": wrapped.src_branch,
1558                        "src_version": src_version,
1559                    })
1560                }
1561                _ => serde_json::Value::Null,
1562            };
1563            error_with_detail(409, e.to_string(), detail)
1564        }
1565        Err(e @ lex_store::StoreError::AmbiguousManifest { .. }) => {
1566            error_with_detail(422, "AmbiguousManifest", serde_json::json!({"detail": e.to_string()}))
1567        }
1568        Err(e @ lex_store::StoreError::InvalidManifest(_))
1569        | Err(e @ lex_store::StoreError::MissingBlobs(_)) => {
1570            error_with_detail(422, e.to_string(), serde_json::Value::Null)
1571        }
1572        Err(e) => write_error_response("apply merge op", e),
1573    }
1574}
1575
1576fn mint_merge_id() -> MergeSessionId {
1577    use std::sync::atomic::{AtomicU64, Ordering};
1578    static COUNTER: AtomicU64 = AtomicU64::new(0);
1579    let nanos = SystemTime::now()
1580        .duration_since(UNIX_EPOCH)
1581        .map(|d| d.as_nanos())
1582        .unwrap_or(0);
1583    let n = COUNTER.fetch_add(1, Ordering::Relaxed);
1584    format!("merge_{nanos:x}_{n:x}")
1585}
1586
1587// ---- #242: append-only sync ---------------------------------------
1588
1589/// `POST /v1/ops/batch` (#242). Server endpoint for `lex op push`.
1590///
1591/// Body: a JSON array of `OperationRecord`s. The handler validates
1592/// DAG integrity by checking that every op's `parents` either
1593/// already exist on the remote *or* appear earlier in the same
1594/// batch. This lets a client send a topologically-ordered slice
1595/// without first probing for what's already there.
1596///
1597/// Response shape:
1598///
1599/// ```json
1600/// { "received": N, "added": M, "skipped": (N-M), "added_ids": [...] }
1601/// ```
1602///
1603/// Failure modes:
1604///
1605/// * `400` — body isn't a JSON array of op records.
1606/// * `422` with `{ "error": "MissingParent", "detail": { "op_id":
1607///   ..., "missing_parent": ... } }` if a parent is unreachable.
1608///   The whole batch is rejected; nothing is persisted. The client
1609///   should backfill the missing op and retry.
1610/// * `409` if the supplied `op_id` doesn't match the canonical
1611///   hash of the record's payload — content addressing must hold
1612///   over the wire.
1613/// * `422` `MissingBlobs` `{ op_id, ids }` / `InvalidManifest`
1614///   `{ op_id, manifest, reason }` for a `SetFiles` op (#1007) whose
1615///   manifest or entry blobs this store lacks, or whose manifest is not
1616///   a valid one — so a head can never name a file set it can't serve.
1617///   Blobs go first (`/v1/blobs/batch`).
1618///
1619/// Idempotency: a record whose `op_id` already exists is silently
1620/// skipped (not added, not rejected). Pushing the same payload
1621/// twice is `received == N, added == 0` on the second call.
1622pub(crate) fn ops_batch_handler(state: &State, body: &str)
1623    -> Response<std::io::Cursor<Vec<u8>>>
1624{
1625    let records: Vec<lex_vcs::OperationRecord> = match serde_json::from_str(body) {
1626        Ok(r) => r,
1627        Err(e) => return error_response(400,
1628            format!("body must be a JSON array of OperationRecord: {e}")),
1629    };
1630    let store = state.store.lock().unwrap();
1631    let log = match lex_vcs::OpLog::open(store.root()) {
1632        Ok(l) => l,
1633        Err(e) => return error_response(500, format!("opening op log: {e}")),
1634    };
1635
1636    // Validate every record before persisting any of them.
1637    //
1638    // 1. Content-addressing: the supplied `op_id` must match the
1639    //    canonical hash of `record.op`. Otherwise the client is
1640    //    sending a forged or corrupted record.
1641    // 2. DAG integrity: every parent must either already exist in
1642    //    the local log OR appear earlier in this batch.
1643    let mut batch_ids: std::collections::BTreeSet<lex_vcs::OpId> =
1644        std::collections::BTreeSet::new();
1645    for rec in &records {
1646        let expected = rec.op.op_id();
1647        if expected != rec.op_id {
1648            return error_with_detail(409, "OpIdMismatch", serde_json::json!({
1649                "supplied": rec.op_id,
1650                "expected": expected,
1651            }));
1652        }
1653        for parent in &rec.op.parents {
1654            let known = match log.get(parent) {
1655                Ok(Some(_)) => true,
1656                Ok(None) => false,
1657                Err(e) => return error_response(500, format!("op log read: {e}")),
1658            };
1659            if !known && !batch_ids.contains(parent) {
1660                return error_with_detail(422, "MissingParent", serde_json::json!({
1661                    "op_id": rec.op_id,
1662                    "missing_parent": parent,
1663                }));
1664            }
1665        }
1666        if let Err(resp) = crate::sync_http::check_set_files(&store, state.blob_limits, rec) {
1667            return resp;
1668        }
1669        batch_ids.insert(rec.op_id.clone());
1670    }
1671
1672    // Persist. `OpLog::put` is idempotent so a re-push is a no-op
1673    // for already-present records.
1674    let mut added = 0usize;
1675    let mut added_ids: Vec<&lex_vcs::OpId> = Vec::new();
1676    for rec in &records {
1677        let already_present = matches!(log.get(&rec.op_id), Ok(Some(_)));
1678        match log.put(rec) {
1679            Ok(()) => {
1680                if !already_present {
1681                    added += 1;
1682                    added_ids.push(&rec.op_id);
1683                }
1684            }
1685            Err(e) => return error_response(500, format!("op log write: {e}")),
1686        }
1687    }
1688
1689    json_response(200, &serde_json::json!({
1690        "received": records.len(),
1691        "added": added,
1692        "skipped": records.len() - added,
1693        "added_ids": added_ids,
1694    }))
1695}
1696
1697/// `POST /v1/attestations/batch` (#242). Server endpoint for `lex
1698/// attest push`.
1699///
1700/// Body: a JSON array of `Attestation`s. Validates that each
1701/// attestation's `op_id` (when set) refers to an op that already
1702/// exists on the remote — `attestation_id` is then re-derivable
1703/// from the canonical form, so cross-store dedup just works.
1704///
1705/// Response: same shape as `ops_batch_handler` but `added_ids` is
1706/// the list of accepted `attestation_id`s.
1707///
1708/// Failure modes:
1709///
1710/// * `400` for malformed JSON.
1711/// * `422` with `{ "error": "UnknownOp", "detail": { ... } }` if
1712///   an attestation's `op_id` references an op the remote doesn't
1713///   know about. Whole batch rejected.
1714/// * `409` `AttestationIdMismatch` if the supplied id doesn't
1715///   match the canonical hash.
1716/// * `403` `ReservedKind` if any attestation's kind is in
1717///   [`State::reserved_kinds`] (#1066). Whole batch rejected, nothing
1718///   written.
1719/// * `403` `ReservedProducer` if any attestation claims a
1720///   `produced_by.tool` in [`State::reserved_producers`]. Whole batch
1721///   rejected, nothing written.
1722///
1723/// Check order (deliberate; a hosted embedder's dispatcher-level guard
1724/// answers the same way): malformed JSON `400` first — the kind is judged
1725/// on the parsed `Attestation`, so nothing unparsable is ever judged —
1726/// then `ReservedKind`, then `ReservedProducer` (a batch violating both
1727/// answers `ReservedKind`), then `409`/`422` validation, then persist.
1728///
1729/// Idempotency: same as the ops endpoint — content-addressed dedup.
1730pub(crate) fn attestations_batch_handler(state: &State, body: &str)
1731    -> Response<std::io::Cursor<Vec<u8>>>
1732{
1733    let attestations: Vec<lex_vcs::Attestation> = match serde_json::from_str(body) {
1734        Ok(a) => a,
1735        Err(e) => return error_response(400,
1736            format!("body must be a JSON array of Attestation: {e}")),
1737    };
1738    let store = state.store.lock().unwrap();
1739    let log = match store.attestation_log() {
1740        Ok(l) => l,
1741        Err(e) => return error_response(500, format!("opening attestation log: {e}")),
1742    };
1743    let op_log = match lex_vcs::OpLog::open(store.root()) {
1744        Ok(l) => l,
1745        Err(e) => return error_response(500, format!("opening op log: {e}")),
1746    };
1747
1748    // Reserved kinds first (#1066), over the WHOLE batch and on the PARSED
1749    // kind — the very value `log.put` serializes below — so a body that
1750    // smuggles a second `"kind"` tag is judged as the one that would be
1751    // stored. Readers such as `latest_review_verdict` key on kind, not on
1752    // the producer name, so the producer reservation alone is not enough.
1753    for att in &attestations {
1754        if let Some((reserved, kind)) = state.reserved_kind_claimed(att) {
1755            return error_with_detail(403, "ReservedKind", serde_json::json!({
1756                "attestation_id": att.attestation_id,
1757                "kind": kind,
1758                "reserved_kind": reserved,
1759                "message": "this attestation kind is reserved for the server; \
1760                            clients may not submit attestations of it \
1761                            (review verdicts go through POST /v1/review/verdict)",
1762            }));
1763        }
1764    }
1765
1766    // Reserved producers next, over the WHOLE batch, so a mixed batch
1767    // (one legitimate + one reserved) is refused outright and nothing is
1768    // written. Server-internal writers (hosted CI, review, examples) call
1769    // the store directly and never come through here.
1770    for att in &attestations {
1771        if let Some(reserved) = state.reserved_producer_claimed(att) {
1772            return error_with_detail(403, "ReservedProducer", serde_json::json!({
1773                "attestation_id": att.attestation_id,
1774                "produced_by_tool": att.produced_by.tool,
1775                "reserved_tool": reserved,
1776                "message": "this producer name is reserved for the server; \
1777                            clients may not submit attestations claiming it",
1778            }));
1779        }
1780    }
1781
1782    // Validate before persisting any record.
1783    for att in &attestations {
1784        // Content-addressing: re-derive attestation_id from the
1785        // payload and reject mismatches.
1786        let expected = lex_vcs::Attestation::with_timestamp(
1787            att.stage_id.clone(),
1788            att.op_id.clone(),
1789            att.intent_id.clone(),
1790            att.kind.clone(),
1791            att.result.clone(),
1792            att.produced_by.clone(),
1793            att.cost.clone(),
1794            att.timestamp,
1795        ).attestation_id;
1796        if expected != att.attestation_id {
1797            return error_with_detail(409, "AttestationIdMismatch", serde_json::json!({
1798                "supplied": att.attestation_id,
1799                "expected": expected,
1800            }));
1801        }
1802        // The op_id field, if set, must point at an op the remote
1803        // knows about. Without this check, attestations would
1804        // dangle into a future sync that never lands the op.
1805        if let Some(op_id) = &att.op_id {
1806            match op_log.get(op_id) {
1807                Ok(Some(_)) => {}
1808                Ok(None) => return error_with_detail(422, "UnknownOp", serde_json::json!({
1809                    "attestation_id": att.attestation_id,
1810                    "op_id": op_id,
1811                })),
1812                Err(e) => return error_response(500, format!("op log read: {e}")),
1813            }
1814        }
1815    }
1816
1817    // Persist. `AttestationLog::put` is idempotent on
1818    // `attestation_id` and the by-stage index is rewritten as a
1819    // marker file, also idempotent.
1820    let mut added = 0usize;
1821    let mut added_ids: Vec<&lex_vcs::AttestationId> = Vec::new();
1822    for att in &attestations {
1823        let already_present = matches!(log.get(&att.attestation_id), Ok(Some(_)));
1824        match log.put(att) {
1825            Ok(()) => {
1826                if !already_present {
1827                    added += 1;
1828                    added_ids.push(&att.attestation_id);
1829                }
1830            }
1831            Err(e) => return error_response(500, format!("attestation log write: {e}")),
1832        }
1833    }
1834
1835    json_response(200, &serde_json::json!({
1836        "received": attestations.len(),
1837        "added": added,
1838        "skipped": attestations.len() - added,
1839        "added_ids": added_ids,
1840    }))
1841}
1842
1843/// `GET /v1/attestations/since?after-op=<op_id>&limit=<n>` (#260).
1844/// Mirror of `ops_since_handler` for the attestation log.
1845///
1846/// Returns attestations whose `op_id` field is reachable from
1847/// **any** branch's head — not just one — and not in `after_op`'s
1848/// ancestry. The cross-branch fan-out matches the push side:
1849/// attestations are stage-keyed, not branch-keyed, so a single
1850/// "since this op" filter is the right shape.
1851///
1852/// Attestations with `op_id: None` (e.g. `Override`,
1853/// `ProducerBlock`) are always included — the cutoff doesn't apply.
1854/// `--limit` caps the response.
1855pub(crate) fn attestations_since_handler(state: &State, query: &str)
1856    -> Response<std::io::Cursor<Vec<u8>>>
1857{
1858    let mut after_op: Option<String> = None;
1859    let mut limit: Option<usize> = None;
1860    for kv in query.split('&') {
1861        let Some((k, v)) = kv.split_once('=') else { continue };
1862        match k {
1863            "after-op" => after_op = Some(v.to_string()),
1864            "limit" => {
1865                limit = Some(match v.parse::<usize>() {
1866                    Ok(n) => n,
1867                    Err(_) => return error_response(400,
1868                        format!("limit must be a positive integer, got `{v}`")),
1869                });
1870            }
1871            _ => {}
1872        }
1873    }
1874
1875    let store = state.store.lock().unwrap();
1876    let log = match store.attestation_log() {
1877        Ok(l) => l,
1878        Err(e) => return error_response(500, format!("opening attestation log: {e}")),
1879    };
1880
1881    // Build the exclude set: every op_id reachable from `after_op`,
1882    // inclusive. Attestations whose op_id is in this set were
1883    // already known to the caller.
1884    let exclude: std::collections::BTreeSet<String> = match &after_op {
1885        None => std::collections::BTreeSet::new(),
1886        Some(cutoff) => {
1887            let op_log = match lex_vcs::OpLog::open(store.root()) {
1888                Ok(l) => l,
1889                Err(e) => return error_response(500, format!("opening op log: {e}")),
1890            };
1891            match op_log.walk_back(cutoff, None) {
1892                Ok(records) => records.into_iter().map(|r| r.op_id).collect(),
1893                Err(_) => {
1894                    // Cutoff op doesn't exist on this remote. Treat
1895                    // as "no exclude" — caller will get every
1896                    // attestation. They may dedup client-side.
1897                    std::collections::BTreeSet::new()
1898                }
1899            }
1900        }
1901    };
1902
1903    let all = match log.list_all() {
1904        Ok(v) => v,
1905        Err(e) => return error_response(500, format!("listing attestations: {e}")),
1906    };
1907    let mut filtered: Vec<lex_vcs::Attestation> = all
1908        .into_iter()
1909        .filter(|a| match &a.op_id {
1910            Some(op_id) => !exclude.contains(op_id),
1911            // No op_id = doesn't participate in the cutoff; always
1912            // ship it on the first pull, server-side idempotency
1913            // dedupes on the client.
1914            None => true,
1915        })
1916        .collect();
1917    // Stable order: oldest-first by `timestamp`, then by
1918    // `attestation_id` for ties. Lets the client land them
1919    // deterministically.
1920    filtered.sort_by(|a, b| {
1921        a.timestamp.cmp(&b.timestamp)
1922            .then_with(|| a.attestation_id.cmp(&b.attestation_id))
1923    });
1924    if let Some(n) = limit {
1925        filtered.truncate(n);
1926    }
1927
1928    json_response(200, &serde_json::to_value(&filtered).unwrap_or_default())
1929}
1930
1931// ── Package concept (#4) ────────────────────────────────────────────────────
1932
1933/// Full coordinates of one declared dependency, captured from the releaser's
1934/// `lex.toml` so the rendered archive can reproduce a faithful `[dependencies]`
1935/// table — carrying BOTH a vcs (`registry` + `version`) reference and a `git`
1936/// mirror when the source declared both. Every field is optional so a bare
1937/// registry dep, a bare git dep, or a dual dep all round-trip. `#[serde(default,
1938/// skip_serializing_if)]` keeps the record compact and back-compatible.
1939#[derive(Debug, Clone, Default, serde::Serialize, serde::Deserialize)]
1940struct DepSpec {
1941    #[serde(default, skip_serializing_if = "Option::is_none")]
1942    registry: Option<String>,
1943    #[serde(default, skip_serializing_if = "Option::is_none")]
1944    version: Option<String>,
1945    #[serde(default, skip_serializing_if = "Option::is_none")]
1946    git: Option<String>,
1947    #[serde(default, skip_serializing_if = "Option::is_none")]
1948    branch: Option<String>,
1949    #[serde(default, skip_serializing_if = "Option::is_none")]
1950    tag: Option<String>,
1951    #[serde(default, skip_serializing_if = "Option::is_none")]
1952    rev: Option<String>,
1953    #[serde(default, skip_serializing_if = "Option::is_none")]
1954    path: Option<String>,
1955}
1956
1957impl DepSpec {
1958    /// Render this spec as the body of a `lex.toml` inline dependency table,
1959    /// e.g. `{ registry = "…", version = "^0.9", git = "…" }`. Keys are emitted
1960    /// in a stable order (vcs primary first, then the git mirror) so the
1961    /// rendered manifest is deterministic.
1962    fn to_toml_inline(&self) -> Option<String> {
1963        let mut parts: Vec<String> = Vec::new();
1964        let mut push = |k: &str, v: &Option<String>| {
1965            if let Some(val) = v {
1966                parts.push(format!("{k} = {}", toml_str(val)));
1967            }
1968        };
1969        push("registry", &self.registry);
1970        push("version", &self.version);
1971        push("git", &self.git);
1972        push("branch", &self.branch);
1973        push("tag", &self.tag);
1974        push("rev", &self.rev);
1975        push("path", &self.path);
1976        if parts.is_empty() {
1977            None
1978        } else {
1979            Some(format!("{{ {} }}", parts.join(", ")))
1980        }
1981    }
1982}
1983
1984/// Quote a string as a TOML basic string (escaping `\` and `"`).
1985fn toml_str(s: &str) -> String {
1986    format!("\"{}\"", s.replace('\\', "\\\\").replace('"', "\\\""))
1987}
1988
1989/// Per-version record stored at `{store_root}/packages/{name}/{version}.json`.
1990#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
1991struct PkgRecord {
1992    name: String,
1993    version: String,
1994    head_op: Option<String>,
1995    published_at: u64,
1996    /// Function names introduced or updated by this version (for retract).
1997    function_names: Vec<String>,
1998    /// External package dependencies of this release — declared by the
1999    /// releaser (from `lex.toml`) and/or extracted from the head's
2000    /// non-inlined imports. The edge set of the dependency graph (#893
2001    /// propagation). `#[serde(default)]` keeps pre-existing records readable.
2002    #[serde(default)]
2003    dependencies: Vec<String>,
2004    /// Full dependency coordinates (name → git+vcs refs) captured from the
2005    /// releaser's `lex.toml`, so the rendered install archive reproduces a
2006    /// faithful `[dependencies]` table (with both git and vcs references).
2007    /// `#[serde(default)]` keeps pre-existing records (which lack it) readable;
2008    /// when empty, the archive falls back to a bare `[package]` manifest.
2009    #[serde(default, skip_serializing_if = "std::collections::BTreeMap::is_empty")]
2010    dependency_specs: std::collections::BTreeMap<String, DepSpec>,
2011    /// Raw op JSON from each file in this publish.
2012    ops: Vec<serde_json::Value>,
2013}
2014
2015/// Whether a package is reachable without authentication.
2016///
2017/// Per-package (applies across all versions), GitHub-style: an org's
2018/// store can hold a mix of public and private packages. Defaults to
2019/// `Private`, so an index written before this field existed (or any
2020/// freshly published package) is private until an owner opts in.
2021#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize, serde::Deserialize, Default)]
2022#[serde(rename_all = "lowercase")]
2023pub enum Visibility {
2024    #[default]
2025    Private,
2026    Public,
2027}
2028
2029/// Index stored at `{store_root}/packages/{name}/index.json`.
2030///
2031/// Tracks which versions have been published and which is "latest", so
2032/// consumers can resolve `{name}@latest` without listing every file.
2033#[derive(Debug, Clone, serde::Serialize, serde::Deserialize, Default)]
2034struct PkgIndex {
2035    /// The most recently published version string (human label, not OpId).
2036    latest: Option<String>,
2037    /// All published versions, newest-last.
2038    versions: Vec<PkgVersionSummary>,
2039    /// Public/private flag. `#[serde(default)]` keeps pre-existing
2040    /// `index.json` files (which lack the field) deserializing as
2041    /// `Private`.
2042    #[serde(default)]
2043    visibility: Visibility,
2044}
2045
2046#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
2047struct PkgVersionSummary {
2048    version: String,
2049    head_op: Option<String>,
2050    published_at: u64,
2051}
2052
2053fn pkg_name_dir(root: &std::path::Path, name: &str) -> PathBuf {
2054    root.join("packages").join(name)
2055}
2056
2057fn pkg_index_path(root: &std::path::Path, name: &str) -> PathBuf {
2058    pkg_name_dir(root, name).join("index.json")
2059}
2060
2061fn pkg_version_path(root: &std::path::Path, name: &str, version: &str) -> PathBuf {
2062    pkg_name_dir(root, name).join(format!("{version}.json"))
2063}
2064
2065fn pkg_archive_path(root: &std::path::Path, name: &str, version: &str) -> PathBuf {
2066    pkg_name_dir(root, name).join(format!("{version}.tar.gz"))
2067}
2068
2069fn load_pkg_index(root: &std::path::Path, name: &str) -> Option<PkgIndex> {
2070    let bytes = std::fs::read(pkg_index_path(root, name)).ok()?;
2071    serde_json::from_slice(&bytes).ok()
2072}
2073
2074fn load_pkg_record(root: &std::path::Path, name: &str, version: &str) -> Option<PkgRecord> {
2075    let bytes = std::fs::read(pkg_version_path(root, name, version)).ok()?;
2076    serde_json::from_slice(&bytes).ok()
2077}
2078
2079fn load_latest_pkg_record(root: &std::path::Path, name: &str) -> Option<PkgRecord> {
2080    let index = load_pkg_index(root, name)?;
2081    let latest = index.latest.clone()?;
2082    load_pkg_record(root, name, &latest)
2083}
2084
2085/// A package is public iff its index exists and is marked `Public`.
2086/// A missing index (unknown package) is treated as private, so the
2087/// public surface never distinguishes "private" from "does not exist".
2088fn pkg_is_public(root: &std::path::Path, name: &str) -> bool {
2089    load_pkg_index(root, name).map(|i| i.visibility) == Some(Visibility::Public)
2090}
2091
2092/// Reject package/version path segments that could escape the
2093/// `packages/` directory or otherwise aren't valid names. Mirrors the
2094/// tenant-id guard's spirit (defense in depth — lex-hub validates the
2095/// tenant, this validates the package/version).
2096fn valid_pkg_segment(s: &str) -> bool {
2097    !s.is_empty()
2098        && s.len() <= 128
2099        && s != "."
2100        && s != ".."
2101        && s.chars().all(|c| c.is_ascii_alphanumeric() || matches!(c, '.' | '_' | '-'))
2102}
2103
2104#[derive(Deserialize)]
2105struct VisibilityReq {
2106    visibility: Visibility,
2107}
2108
2109/// `PUT /v1/pkg/{name}/visibility` — set a package public or private.
2110///
2111/// Authorization is implicit: this runs against a single tenant's store,
2112/// and the caller only reaches *this* store because the front door
2113/// (lex-hub) authenticated their token and selected it. A caller can
2114/// therefore only change visibility of packages they own.
2115fn pkg_set_visibility_handler(
2116    state: &State,
2117    name: &str,
2118    body: &str,
2119) -> Response<std::io::Cursor<Vec<u8>>> {
2120    if !valid_pkg_segment(name) {
2121        return error_response(400, format!("invalid package name {name:?}"));
2122    }
2123    let req: VisibilityReq = match serde_json::from_str(body) {
2124        Ok(r) => r,
2125        Err(e) => return error_response(400, format!("bad request: {e}")),
2126    };
2127    let mut index = match load_pkg_index(&state.root, name) {
2128        Some(i) => i,
2129        None => return error_response(404, format!("package {name:?} not found")),
2130    };
2131    index.visibility = req.visibility;
2132    let bytes = serde_json::to_vec_pretty(&index).unwrap_or_default();
2133    match std::fs::write(pkg_index_path(&state.root, name), bytes) {
2134        Ok(()) => json_response(
2135            200,
2136            &serde_json::json!({ "name": name, "visibility": index.visibility }),
2137        ),
2138        Err(e) => error_response(500, format!("write index: {e}")),
2139    }
2140}
2141
2142#[derive(serde::Deserialize)]
2143struct ReleaseReq {
2144    version: String,
2145    #[serde(default)]
2146    branch: Option<String>,
2147    /// External package dependencies (from the releaser's `lex.toml`).
2148    /// Unioned with any non-inlined external imports found in the head.
2149    #[serde(default)]
2150    dependencies: Vec<String>,
2151    /// Full dependency coordinates (name → git+vcs refs) from the releaser's
2152    /// `lex.toml`, so the rendered install archive reproduces a faithful
2153    /// `[dependencies]` table. Optional (back-compatible); when omitted, the
2154    /// archive manifest carries only `[package]`.
2155    #[serde(default)]
2156    dependency_specs: std::collections::BTreeMap<String, DepSpec>,
2157}
2158
2159/// `POST /v1/pkg/{name}/release` — cut an immutable versioned release of
2160/// an **op-log-hosted** package: snapshot the current branch head as
2161/// `name@version` in the registry (#893). Unlike `POST /v1/pkg/publish`
2162/// (archive upload), this records only the op-log ref (`head_op`) — a
2163/// consumer resolves the version, `op pull`s that head, and renders
2164/// source. A published version is immutable: re-releasing an existing
2165/// version is a 409, so a resolved+locked dependency can never change
2166/// under a consumer.
2167fn pkg_release_handler(state: &State, name: &str, body: &str) -> Response<std::io::Cursor<Vec<u8>>> {
2168    if !valid_pkg_segment(name) {
2169        return error_response(400, format!("invalid package name {name:?}"));
2170    }
2171    let req: ReleaseReq = match serde_json::from_str(body) {
2172        Ok(r) => r,
2173        Err(e) => return error_response(400, format!("bad request: {e}")),
2174    };
2175    let version = req.version.trim().to_string();
2176    if version.is_empty() || !valid_pkg_segment(&version) {
2177        return error_response(400, "version must be a non-empty, path-safe string (e.g. 1.2.0)");
2178    }
2179    // Immutable: a published version can never be overwritten.
2180    if load_pkg_record(&state.root, name, &version).is_some() {
2181        return error_response(
2182            409,
2183            format!("{name}@{version} already released; releases are immutable — bump the version"),
2184        );
2185    }
2186
2187    let store = state.store.lock().unwrap();
2188    let branch = req.branch.unwrap_or_else(|| store.current_branch());
2189    let head_op = match store.get_branch(&branch) {
2190        Ok(Some(b)) => b.head_op,
2191        Ok(None) => return error_response(404, format!("unknown branch {branch:?}")),
2192        Err(e) => return error_response(500, format!("get_branch: {e}")),
2193    };
2194    let Some(head_op) = head_op else {
2195        return error_response(400, format!("branch {branch:?} has no commits to release"));
2196    };
2197
2198    // Version-bump gate (#893): the version increment must be at least what
2199    // the public-API change requires — a breaking change (a removed or
2200    // re-signatured public declaration) needs a *major* bump, an addition at
2201    // least a *minor*. This is what makes `^`/`~` resolution safe: a caret
2202    // update can't silently pull a breaking change mislabeled as a patch.
2203    //
2204    // Compare against the **semver predecessor** — the highest already-published
2205    // version strictly less than the new one (not merely `latest`, so releasing
2206    // 2.0.0 after 1.5.0 diffs against 1.5.0). No predecessor (first version, or
2207    // a back-port below everything) means nothing to gate.
2208    let predecessor = load_pkg_index(&state.root, name)
2209        .map(|i| i.versions)
2210        .unwrap_or_default()
2211        .into_iter()
2212        .filter_map(|v| lex_syntax::semver::parse_exact(&v.version).map(|p| (p, v)))
2213        .filter(|(p, _)| lex_syntax::semver::parse_exact(&version).map(|n| *p < n).unwrap_or(false))
2214        .max_by_key(|(p, _)| *p)
2215        .map(|(_, v)| v);
2216    if let Some(prev) = predecessor {
2217        if let (Some(prev_head), Some(declared)) = (
2218            prev.head_op.clone(),
2219            lex_syntax::semver::bump_between(&prev.version, &version),
2220        ) {
2221            if let (Ok(prev_api), Ok(new_api)) = (
2222                lex_store::api::public_api_at_op(&store, &prev_head),
2223                lex_store::api::public_api_at_op(&store, &head_op),
2224            ) {
2225                use lex_store::api::ApiChange;
2226                use lex_syntax::semver::Bump;
2227                let (required, why) = match lex_store::api::classify_api_change(&prev_api, &new_api) {
2228                    ApiChange::Breaking(d) => (Bump::Major, d),
2229                    ApiChange::Additive(d) => (Bump::Minor, d),
2230                    ApiChange::None => (Bump::Patch, String::new()),
2231                };
2232                if declared < required {
2233                    let need = match required {
2234                        Bump::Major => "major",
2235                        Bump::Minor => "minor",
2236                        Bump::Patch => "patch",
2237                    };
2238                    return error_response(
2239                        422,
2240                        format!(
2241                            "version bump too small: {} → {version} is a {declared:?} bump, \
2242                             but the API change ({why}) requires a {need} bump",
2243                            prev.version
2244                        ),
2245                    );
2246                }
2247            }
2248        }
2249    }
2250
2251    // The package's exported function names at this head (for retract /
2252    // the catalog), read through the SigId the head names each by.
2253    let head = store.branch_head(&branch).unwrap_or_default();
2254    let pairs: Vec<(String, String)> = head.iter().map(|(s, st)| (s.clone(), st.clone())).collect();
2255    let function_names: Vec<String> = store
2256        .get_asts_for_sigs_bulk(&pairs)
2257        .into_iter()
2258        .filter_map(|r| r.ok())
2259        .filter_map(|s| match s {
2260            lex_ast::Stage::FnDecl(fd) => Some(fd.name),
2261            _ => None,
2262        })
2263        .collect();
2264    // Dependency-graph edges (#893 propagation): the releaser's declared deps
2265    // (from lex.toml) unioned with any external imports still visible in the
2266    // head (most are inlined at publish, so the declared list is primary).
2267    let mut deps: std::collections::BTreeSet<String> = req.dependencies.into_iter().collect();
2268    if let Ok(extracted) = lex_store::api::external_dependencies_at_op(&store, &head_op) {
2269        deps.extend(extracted);
2270    }
2271    // Names captured as full coordinates count as dependency-graph edges too,
2272    // so a release that sends only `dependency_specs` still records the edges.
2273    deps.extend(req.dependency_specs.keys().cloned());
2274    let dependencies: Vec<String> = deps.into_iter().collect();
2275    drop(store);
2276
2277    let published_at = std::time::SystemTime::now()
2278        .duration_since(std::time::UNIX_EPOCH)
2279        .map(|d| d.as_secs())
2280        .unwrap_or(0);
2281    let record = PkgRecord {
2282        name: name.to_string(),
2283        version: version.clone(),
2284        head_op: Some(head_op.clone()),
2285        published_at,
2286        function_names,
2287        dependencies,
2288        dependency_specs: req.dependency_specs,
2289        ops: Vec::new(),
2290    };
2291    if let Err(e) = save_pkg_record(&state.root, &record, None) {
2292        return error_response(500, format!("write release: {e}"));
2293    }
2294    json_response(
2295        201,
2296        &serde_json::json!({
2297            "name": name,
2298            "version": version,
2299            "head_op": head_op,
2300            "branch": branch,
2301        }),
2302    )
2303}
2304
2305/// Names of the tenant's PUBLIC packages, sorted. Pure (filesystem in,
2306/// names out) so the visibility filter is unit-testable.
2307fn public_pkg_names(root: &std::path::Path) -> Vec<String> {
2308    list_pkg_names(root)
2309        .into_iter()
2310        .filter(|name| pkg_is_public(root, name))
2311        .collect()
2312}
2313
2314/// `GET /v1/public/<tenant>` (org page) — list only the tenant's PUBLIC
2315/// packages (latest version of each). Private packages are omitted, so
2316/// their existence is not revealed.
2317fn public_pkg_list_handler(state: &State) -> Response<std::io::Cursor<Vec<u8>>> {
2318    let packages: Vec<serde_json::Value> = public_pkg_names(&state.root)
2319        .iter()
2320        .filter_map(|name| {
2321            let r = load_latest_pkg_record(&state.root, name)?;
2322            Some(serde_json::json!({
2323                "name": r.name,
2324                "version": r.version,
2325                "head_op": r.head_op,
2326                "published_at": r.published_at,
2327            }))
2328        })
2329        .collect();
2330    json_response(200, &serde_json::json!({ "packages": packages }))
2331}
2332
2333/// A resolved public read target. Separated from response formatting so
2334/// the routing/guard logic is unit-testable without constructing HTTP
2335/// responses. `Err(status)` is a guard failure (405 = non-GET,
2336/// 404 = invalid/unknown route).
2337#[derive(Debug, PartialEq, Eq)]
2338enum PublicTarget {
2339    List,
2340    Latest(String),
2341    Versions(String),
2342    ApiDiff(String),
2343    Head(String),
2344    Version(String, String),
2345    Archive(String, String),
2346}
2347
2348impl PublicTarget {
2349    /// The package name a target refers to, if any (`List` has none).
2350    fn pkg_name(&self) -> Option<&str> {
2351        match self {
2352            PublicTarget::List => None,
2353            PublicTarget::Latest(n)
2354            | PublicTarget::Versions(n)
2355            | PublicTarget::ApiDiff(n)
2356            | PublicTarget::Head(n)
2357            | PublicTarget::Version(n, _)
2358            | PublicTarget::Archive(n, _) => Some(n),
2359        }
2360    }
2361}
2362
2363/// Pure routing decision for `/v1/public/<tenant>` reads. `path` is the
2364/// portion after the tenant (leading `/` ok). Enforces GET-only and
2365/// per-segment name validation; does NOT consult the store (visibility
2366/// is checked by the caller, which has store access).
2367fn resolve_public(method: &Method, path: &str) -> Result<PublicTarget, u16> {
2368    if !matches!(method, Method::Get) {
2369        return Err(405);
2370    }
2371    let rest = path.trim_matches('/');
2372    if rest.is_empty() {
2373        return Ok(PublicTarget::List);
2374    }
2375    let segs: Vec<&str> = rest.split('/').collect();
2376    if !segs.iter().all(|s| valid_pkg_segment(s)) {
2377        return Err(404);
2378    }
2379    match segs.as_slice() {
2380        [n] => Ok(PublicTarget::Latest(n.to_string())),
2381        [n, "versions"] => Ok(PublicTarget::Versions(n.to_string())),
2382        [n, "api-diff"] => Ok(PublicTarget::ApiDiff(n.to_string())),
2383        [n, "head"] => Ok(PublicTarget::Head(n.to_string())),
2384        [n, v, "archive"] => Ok(PublicTarget::Archive(n.to_string(), v.to_string())),
2385        [n, v] => Ok(PublicTarget::Version(n.to_string(), v.to_string())),
2386        _ => Err(404),
2387    }
2388}
2389
2390/// Unauthenticated, read-only access to **public** packages in `state`'s
2391/// store. `path` is the portion of the URL after `/v1/public/<tenant>`
2392/// (with a leading `/`); lex-hub resolves `<tenant>` → store and calls
2393/// this. Visibility gating, GET-only enforcement, and segment validation
2394/// all live here so the whole public surface is auditable in one place.
2395///
2396/// Everything served here is package-scoped — manifests and the source
2397/// archive of a public package's own publish — so it cannot leak code
2398/// from a private package that happens to share content-addressed stages
2399/// in the same store.
2400pub fn route_public(
2401    state: &State,
2402    method: &Method,
2403    path: &str,
2404    query: &str,
2405) -> Response<std::io::Cursor<Vec<u8>>> {
2406    let target = match resolve_public(method, path) {
2407        Ok(t) => t,
2408        Err(405) => return error_response(405, "public read is GET-only"),
2409        Err(_) => return error_response(404, "not found"),
2410    };
2411    // The org listing already filters to public packages itself.
2412    if let PublicTarget::List = target {
2413        return public_pkg_list_handler(state);
2414    }
2415    // Single 404 for both "private" and "absent" — never reveal which.
2416    if let Some(name) = target.pkg_name() {
2417        if !pkg_is_public(&state.root, name) {
2418            return error_response(404, format!("package {name:?} not found"));
2419        }
2420    }
2421    match target {
2422        PublicTarget::List => unreachable!("handled above"),
2423        PublicTarget::Latest(n) => pkg_get_handler(state, &n),
2424        PublicTarget::Versions(n) => pkg_versions_handler(state, &n),
2425        PublicTarget::ApiDiff(n) => pkg_api_diff_handler(state, &n, query),
2426        PublicTarget::Head(n) => pkg_head_handler(state, &n),
2427        PublicTarget::Version(n, v) => pkg_get_version_handler(state, &n, &v),
2428        PublicTarget::Archive(n, v) => pkg_archive_handler(state, &n, &v),
2429    }
2430}
2431
2432fn save_pkg_record(
2433    root: &std::path::Path,
2434    record: &PkgRecord,
2435    // `None` for an op-log release, whose source of truth is the op-log at
2436    // `head_op` (pull + render); `Some` for an archive-upload publish.
2437    archive: Option<&[u8]>,
2438) -> std::io::Result<()> {
2439    let dir = pkg_name_dir(root, &record.name);
2440    std::fs::create_dir_all(&dir)?;
2441
2442    // Per-version record.
2443    let rec_bytes = serde_json::to_vec_pretty(record).unwrap_or_default();
2444    std::fs::write(pkg_version_path(root, &record.name, &record.version), rec_bytes)?;
2445
2446    // Archive (tar.gz) for the download endpoint, when one was uploaded.
2447    if let Some(archive) = archive {
2448        std::fs::write(pkg_archive_path(root, &record.name, &record.version), archive)?;
2449    }
2450
2451    // Update the index.
2452    let mut index = load_pkg_index(root, &record.name).unwrap_or_default();
2453    index.latest = Some(record.version.clone());
2454    if !index.versions.iter().any(|v| v.version == record.version) {
2455        index.versions.push(PkgVersionSummary {
2456            version: record.version.clone(),
2457            head_op: record.head_op.clone(),
2458            published_at: record.published_at,
2459        });
2460    }
2461    let idx_bytes = serde_json::to_vec_pretty(&index).unwrap_or_default();
2462    std::fs::write(pkg_index_path(root, &record.name), idx_bytes)
2463}
2464
2465fn list_pkg_names(root: &std::path::Path) -> Vec<String> {
2466    let dir = root.join("packages");
2467    let Ok(entries) = std::fs::read_dir(&dir) else {
2468        return Vec::new();
2469    };
2470    let mut names: Vec<String> = entries
2471        .filter_map(|e| e.ok())
2472        .filter(|e| e.path().is_dir())
2473        .filter_map(|e| e.file_name().into_string().ok())
2474        .collect();
2475    names.sort();
2476    names
2477}
2478
2479fn collect_lex_files(dir: &std::path::Path, out: &mut Vec<PathBuf>) {
2480    let Ok(entries) = std::fs::read_dir(dir) else { return };
2481    let mut entries: Vec<_> = entries.filter_map(|e| e.ok()).collect();
2482    entries.sort_by_key(|e| e.path());
2483    for entry in entries {
2484        let path = entry.path();
2485        if path.is_dir() {
2486            collect_lex_files(&path, out);
2487        } else if path.extension().and_then(|x| x.to_str()) == Some("lex") {
2488            out.push(path);
2489        }
2490    }
2491}
2492
2493/// `POST /v1/pkg/publish` — publish a multi-file package from a `.tar.gz`
2494/// archive containing `lex.toml` and `src/**/*.lex`.
2495fn pkg_publish_handler(state: &State, body: &[u8]) -> Response<std::io::Cursor<Vec<u8>>> {
2496    let tmp = match tempfile::TempDir::new() {
2497        Ok(t) => t,
2498        Err(e) => return error_response(500, format!("create temp dir: {e}")),
2499    };
2500    {
2501        let gz = flate2::read::GzDecoder::new(std::io::Cursor::new(body));
2502        let mut ar = tar::Archive::new(gz);
2503        if let Err(e) = ar.unpack(tmp.path()) {
2504            return error_response(400, format!("unpack archive: {e}"));
2505        }
2506    }
2507
2508    let toml_path = tmp.path().join("lex.toml");
2509    if !toml_path.exists() {
2510        return error_response(400, "archive must contain lex.toml at root");
2511    }
2512    let manifest = match Manifest::load(&toml_path) {
2513        Ok(m) => m,
2514        Err(e) => return error_response(400, format!("lex.toml: {e}")),
2515    };
2516    let (pkg_name, pkg_version) = match &manifest.package {
2517        Some(m) => (m.name.clone(), m.version.clone()),
2518        None => return error_response(400, "lex.toml must have a [package] section"),
2519    };
2520
2521    // Reject a same-(name, version) re-publish BEFORE any store write:
2522    // this ran after the publish loop, so a 409'd duplicate still
2523    // appended its ops to the tenant's op log (#826).
2524    if load_pkg_record(&state.root, &pkg_name, &pkg_version).is_some() {
2525        return error_response(
2526            409,
2527            format!(
2528                "package {pkg_name}@{pkg_version} already published; \
2529                 bump the version in lex.toml to publish a new release"
2530            ),
2531        );
2532    }
2533
2534    let src_dir = tmp.path().join("src");
2535    if !src_dir.exists() {
2536        return error_response(400, "archive must contain a src/ directory");
2537    }
2538    let mut lex_files: Vec<PathBuf> = Vec::new();
2539    collect_lex_files(&src_dir, &mut lex_files);
2540    if lex_files.is_empty() {
2541        return error_response(400, "no .lex files found in src/");
2542    }
2543
2544    let store = state.store.lock().unwrap();
2545    let branch = store.current_branch();
2546
2547    // `old_fns_by_name` mirrors the branch's current function set,
2548    // GROUPED by name rather than collapsed to one entry per name. The
2549    // branch is tenant-wide and carries history: several live functions
2550    // can share a bare name (#818 — three unrelated `validate`s across
2551    // `field.lex`/`schema.lex`/`validator.lex` in a real package, all
2552    // published before names were file-prefixed). `SigId` disambiguates
2553    // them correctly — it hashes the full signature, not just the name —
2554    // so the bug was ever collapsing multiple SigIds sharing a name down
2555    // to one `FnDecl`, silently discarding the others and corrupting
2556    // later diffs against them. The package side can no longer produce
2557    // such a collision: one pass, one prefix-mangled name per
2558    // declaration (#828).
2559    //
2560    // Reading it once, outside any per-file loop, is what fixed #813:
2561    // `branch_head` walks the whole branch op history and `get_ast` is a
2562    // disk fetch per live function, so re-deriving this per file was
2563    // O(files * (history_size + live_fn_count)) — tens of minutes on a
2564    // tenant with 110k+ accumulated ops. Since #828 there is one pass, so
2565    // the shape is structural rather than a discipline to maintain.
2566    let old_head = match store.branch_head(&branch) {
2567        Ok(h) => h,
2568        Err(e) => return error_response(500, format!("branch_head: {e}")),
2569    };
2570    // Resolve each AST through the SigId the branch head names, never its
2571    // StageId: StageIds are name-independent, so two live functions
2572    // differing only in name share one, `stage_index` maps it to just one
2573    // of their sigs, and the name that lookup missed got re-reported as an
2574    // Add on every publish of unchanged source (#826). It also skips that
2575    // index (#825) — see `get_asts_for_sigs_bulk`.
2576    let old_pairs: Vec<(String, String)> =
2577        old_head.iter().map(|(sig, stage)| (sig.clone(), stage.clone())).collect();
2578    let mut old_fns_by_name: BTreeMap<String, Vec<lex_ast::FnDecl>> = BTreeMap::new();
2579    for fd in store.get_asts_for_sigs_bulk(&old_pairs)
2580        .into_iter()
2581        .filter_map(|r| r.ok())
2582        .filter_map(|s| match s { lex_ast::Stage::FnDecl(fd) => Some(fd), _ => None })
2583    {
2584        old_fns_by_name.entry(fd.name.clone()).or_default().push(fd);
2585    }
2586    // Old `type`s on the branch, by (mangled) name — captured too (#895).
2587    // Types aren't overloaded, so a plain name map needs no take_matching.
2588    let mut old_types_by_name: BTreeMap<String, lex_ast::TypeDecl> = BTreeMap::new();
2589    for td in store.get_asts_for_sigs_bulk(&old_pairs)
2590        .into_iter()
2591        .filter_map(|r| r.ok())
2592        .filter_map(|s| match s { lex_ast::Stage::TypeDecl(td) => Some(td), _ => None })
2593    {
2594        old_types_by_name.insert(td.name.clone(), td);
2595    }
2596    // A name-independent fingerprint of a function's *contract*
2597    // (effects, param types, return type, examples — everything SigId
2598    // hashes except the name). Used only to disambiguate when multiple
2599    // candidates share a bare name: if this file's own declaration
2600    // structurally matches exactly one of them, that's unambiguously
2601    // the same evolving function; otherwise it's a new, unrelated
2602    // declaration that happens to reuse a name used elsewhere.
2603    fn structural_key(fd: &lex_ast::FnDecl) -> Option<String> {
2604        let mut anon = fd.clone();
2605        anon.name = String::new();
2606        lex_ast::sig_id(&lex_ast::Stage::FnDecl(anon))
2607    }
2608
2609    // Look up (and consume) the candidate in `map[name]` that matches
2610    // `new_fd`'s identity. One candidate is taken as unambiguous (so a
2611    // signature-changing edit still reads as a modification of the same
2612    // function); several always require an exact structural match to
2613    // disambiguate, and no match means this is a distinct, unrelated
2614    // declaration reusing a name used elsewhere — `None` rather than a
2615    // guess that mis-attributes history. Safe because the candidate pool
2616    // is complete and never grows: it is built once, up front, from the
2617    // whole branch.
2618    fn take_matching(
2619        map: &mut BTreeMap<String, Vec<lex_ast::FnDecl>>,
2620        name: &str,
2621        new_fd: &lex_ast::FnDecl,
2622    ) -> Option<lex_ast::FnDecl> {
2623        let candidates = map.get_mut(name)?;
2624        let idx = match candidates.len() {
2625            0 => return None,
2626            1 => 0,
2627            _ => {
2628                let want = structural_key(new_fd);
2629                candidates.iter().position(|c| structural_key(c) == want)?
2630            }
2631        };
2632        let matched = candidates.remove(idx);
2633        if candidates.is_empty() {
2634            map.remove(name);
2635        }
2636        Some(matched)
2637    }
2638
2639    // ---- one pass over the whole package -------------------------
2640    // `load_package` merges every file in the archive into ONE program
2641    // through a single shared loader pass. The flattening per-file loads
2642    // this replaced gave each top-level file its own copy of everything
2643    // it imported, so a shared dependency was canonicalized,
2644    // type-checked, diffed and published once per importing file: the
2645    // real 21-file `lex-schema` package, whose `error.lex` is imported by
2646    // 17 of its files, produced 2,239 `FnDecl`s for 693 distinct names
2647    // and paid for all 2,239 (#828). Collapsing 21 `publish_program`
2648    // calls into one matters even more than the 3.2x itself, because each
2649    // call independently reads every live function on the branch and
2650    // walks the op log for `old_imports`.
2651    //
2652    // The cost is that nothing is published under its bare source name
2653    // any more: `fn validate` in `src/field.lex` is
2654    // `field_<hash>.validate`. That is what makes one program safe to
2655    // check as a unit — the checker's global scope is keyed by name, and
2656    // two files may each declare their own `validate` (#818) — and it is
2657    // the naming change #828 asks for in exchange for the single pass.
2658    // The archive-publish path type-checks the loaded program with no
2659    // dependency resolver, so it still inlines registry/git deps to stay
2660    // self-contained (#930). The op-log-native `op push` path publishes
2661    // without inlining and resolves at the gate instead.
2662    let loaded = match load_package(&lex_files, tmp.path(), &pkg_name, /*inline_packages=*/ true) {
2663        Ok(p) => p,
2664        Err(e) => return error_response(400, format!("load package: {e}")),
2665    };
2666    let mut stages = canonicalize_program(&loaded.program);
2667    // Type errors are reported for the package, not per file: every name
2668    // in them carries its file's mangling prefix, so the offending file
2669    // is still named. Nothing is published unless the whole package
2670    // checks, where before each file published as it was processed and a
2671    // later failure left the earlier files' ops applied.
2672    if let Err(errs) = lex_types::check_and_rewrite_program(&mut stages) {
2673        return error_with_detail(
2674            422,
2675            format!("type errors in package {pkg_name}"),
2676            serde_json::to_value(&errs).unwrap(),
2677        );
2678    }
2679    let new_fns = stage_fns(&stages);
2680    let all_function_names: Vec<String> = new_fns.keys().cloned().collect();
2681
2682    // Resolve each declaration against the branch's current state. No
2683    // in-request bookkeeping is needed now: one pass means each name is
2684    // declared once, so there is no earlier file in this request whose
2685    // just-published version a later one has to diff against.
2686    let mut old_fns: BTreeMap<String, lex_ast::FnDecl> = BTreeMap::new();
2687    for (name, new_fd) in &new_fns {
2688        if let Some(fd) = take_matching(&mut old_fns_by_name, name, new_fd) {
2689            old_fns.insert(name.clone(), fd);
2690        }
2691    }
2692    let new_types = stage_types(&stages);
2693    // Only diff old types this publish also declares — one on the branch
2694    // but absent from this file-set is left alone (as unclaimed old fns
2695    // are), not read as removed.
2696    let old_types: BTreeMap<String, lex_ast::TypeDecl> = new_types
2697        .keys()
2698        .filter_map(|n| old_types_by_name.get(n).map(|td| (n.clone(), td.clone())))
2699        .collect();
2700    let report =
2701        lex_vcs::compute_diff_with_types(&old_fns, &new_fns, &old_types, &new_types, false);
2702
2703    // Imports stay attributed per file — `AddImport`/`RemoveImport` carry
2704    // an `in_file`, and history records these same root-relative keys.
2705    // Each file now gets only the modules it imports itself; a flattened
2706    // per-file load could not tell those from its children's.
2707    //
2708    // `imports_by_file` carries each import's real `as` alias (#909/#930), so
2709    // a non-default `import "..." as x` round-trips as `x`.
2710    let mut new_imports = lex_vcs::ImportMap::new();
2711    for (file, modules) in &loaded.imports_by_file {
2712        let entry = new_imports.entry(file.clone()).or_default();
2713        for (reference, alias) in modules {
2714            entry.insert(lex_vcs::ImportRef {
2715                reference: reference.clone(),
2716                alias: alias.clone(),
2717            });
2718        }
2719    }
2720
2721    // Record each declaration's source file (#894) so an HTTP-published
2722    // package de-flattens in `export-git` too, not just a CLI-published one.
2723    let outcome = match store.publish_program_with_intent(
2724        &branch,
2725        &stages,
2726        &report,
2727        &new_imports,
2728        false,
2729        None,
2730        None,
2731        &loaded.module_prefixes,
2732    ) {
2733        Ok(outcome) => outcome,
2734        Err(lex_store::StoreError::TypeError(errs)) => {
2735            return error_with_detail(422, "type errors", serde_json::to_value(&errs).unwrap());
2736        }
2737        Err(e) => return write_error_response("publish_program", e),
2738    };
2739    let all_ops: Vec<serde_json::Value> = match serde_json::to_value(&outcome.ops) {
2740        Ok(serde_json::Value::Array(arr)) => arr,
2741        _ => Vec::new(),
2742    };
2743    let final_head_op = outcome.head_op;
2744
2745    // Deliberately no "genuinely removed" cleanup pass here. Whatever
2746    // remains in `old_fns_by_name` was never claimed by any file in
2747    // THIS archive — but the branch this walks is scoped to the whole
2748    // TENANT, not to this one package: a tenant that has ever published
2749    // more than one package (confirmed in production — `lex-schema` and
2750    // `lex-ocpi` share a tenant) has every other package's functions
2751    // sitting in `old_fns_by_name` too, forever unclaimed by any file in
2752    // *this* package's own archive. An earlier version of this handler
2753    // treated all such leftovers as "removed" and would have emitted
2754    // RemoveFunction ops for a completely unrelated package's functions
2755    // on every single publish. Caught before it shipped (`diff_to_ops`
2756    // failed atomically on a stale SigId before applying anything, so
2757    // no data was actually lost) — see alpibrusl/lex-lang#818's
2758    // follow-up. Nothing here currently tracks which package "owns" a
2759    // given branch function, so there's no reliable way to tell a
2760    // genuine same-package removal from another package's untouched
2761    // function; leaving a deleted function's stage un-removed (it just
2762    // sits there, unreferenced) is the safe default until package-scoped
2763    // ownership is tracked, not silently deleting a stranger's data.
2764    // Since #828 the same applies to a package's own previous names: the
2765    // first publish after file-prefixed naming landed leaves the bare
2766    // names it used to publish under sitting unreferenced, for the same
2767    // reason — this view cannot tell them from another package's.
2768
2769    let now = SystemTime::now()
2770        .duration_since(UNIX_EPOCH)
2771        .map(|d| d.as_secs())
2772        .unwrap_or(0);
2773    // External deps from the head's non-inlined imports (archive publish; the
2774    // op-log release route also accepts declared deps from lex.toml).
2775    let dependencies: Vec<String> = final_head_op
2776        .as_ref()
2777        .and_then(|h| lex_store::api::external_dependencies_at_op(&store, h).ok())
2778        .unwrap_or_default();
2779    let record = PkgRecord {
2780        name: pkg_name.clone(),
2781        version: pkg_version,
2782        head_op: final_head_op.clone(),
2783        published_at: now,
2784        function_names: all_function_names,
2785        dependencies,
2786        // The archive-upload path stores the uploaded archive verbatim (its
2787        // lex.toml already carries the full dependency table), so it is served
2788        // directly rather than re-rendered — no captured specs needed here.
2789        dependency_specs: Default::default(),
2790        ops: all_ops.clone(),
2791    };
2792    if let Err(e) = save_pkg_record(&state.root, &record, Some(body)) {
2793        return error_response(500, format!("save package index: {e}"));
2794    }
2795
2796    json_response(200, &serde_json::json!({
2797        "package": pkg_name,
2798        "ops": all_ops,
2799        "head_op": final_head_op,
2800    }))
2801}
2802
2803/// `GET /v1/pkg` — list packages published by this tenant (latest version of each).
2804fn pkg_list_handler(state: &State) -> Response<std::io::Cursor<Vec<u8>>> {
2805    let names = list_pkg_names(&state.root);
2806    let packages: Vec<serde_json::Value> = names.iter()
2807        .filter_map(|name| {
2808            let idx = load_pkg_index(&state.root, name)?;
2809            let latest = idx.latest.as_deref()?;
2810            let r = load_pkg_record(&state.root, name, latest)?;
2811            Some(serde_json::json!({
2812                "name": r.name,
2813                "version": r.version,
2814                "head_op": r.head_op,
2815                "published_at": r.published_at,
2816            }))
2817        })
2818        .collect();
2819    json_response(200, &serde_json::json!({ "packages": packages }))
2820}
2821
2822/// `GET /v1/pkg/{name}` — latest version details for a package.
2823fn pkg_get_handler(state: &State, name: &str) -> Response<std::io::Cursor<Vec<u8>>> {
2824    match load_latest_pkg_record(&state.root, name) {
2825        Some(r) => json_response(200, &serde_json::json!({
2826            "name": r.name,
2827            "version": r.version,
2828            "head_op": r.head_op,
2829            "published_at": r.published_at,
2830            "function_names": r.function_names,
2831            "ops": r.ops,
2832        })),
2833        None => error_response(404, format!("package {name:?} not found")),
2834    }
2835}
2836
2837/// `GET /v1/pkg/{name}/versions` — all published versions for a package.
2838fn pkg_versions_handler(state: &State, name: &str) -> Response<std::io::Cursor<Vec<u8>>> {
2839    match load_pkg_index(&state.root, name) {
2840        Some(idx) => json_response(200, &serde_json::json!({
2841            "name": name,
2842            "latest": idx.latest,
2843            "versions": idx.versions,
2844        })),
2845        None => error_response(404, format!("package {name:?} not found")),
2846    }
2847}
2848
2849/// `GET /v1/pkg/{name}/api-diff?from=<v>&to=<v>` — how the public API changed
2850/// between two releases (#893 propagation): the classification and the
2851/// mechanically-propagatable renames, so `lex propagate` can auto-derive its
2852/// `--rename` edits from a hosted release pair rather than have them restated.
2853fn pkg_api_diff_handler(state: &State, name: &str, query: &str) -> Response<std::io::Cursor<Vec<u8>>> {
2854    let mut from: Option<String> = None;
2855    let mut to: Option<String> = None;
2856    for kv in query.split('&') {
2857        match kv.split_once('=') {
2858            Some(("from", v)) => from = Some(v.to_string()),
2859            Some(("to", v)) => to = Some(v.to_string()),
2860            _ => {}
2861        }
2862    }
2863    let (Some(from), Some(to)) = (from, to) else {
2864        return error_response(400, "api-diff requires ?from=<version>&to=<version>");
2865    };
2866    let head_of = |v: &str| load_pkg_record(&state.root, name, v).and_then(|r| r.head_op);
2867    let (Some(from_head), Some(to_head)) = (head_of(&from), head_of(&to)) else {
2868        return error_response(404, format!("{name}: unknown release in {from}..{to}"));
2869    };
2870
2871    let store = state.store.lock().unwrap();
2872    let (prev_api, new_api) = match (
2873        lex_store::api::public_api_at_op(&store, &from_head),
2874        lex_store::api::public_api_at_op(&store, &to_head),
2875    ) {
2876        (Ok(a), Ok(b)) => (a, b),
2877        _ => return error_response(500, "could not read package APIs for the given releases"),
2878    };
2879    let (change, detail) = match lex_store::api::classify_api_change(&prev_api, &new_api) {
2880        lex_store::api::ApiChange::Breaking(d) => ("breaking", d),
2881        lex_store::api::ApiChange::Additive(d) => ("additive", d),
2882        lex_store::api::ApiChange::None => ("none", String::new()),
2883    };
2884    let renames = lex_store::api::detect_renames(&prev_api, &new_api);
2885    json_response(200, &serde_json::json!({
2886        "name": name, "from": from, "to": to,
2887        "change": change, "detail": detail, "renames": renames,
2888    }))
2889}
2890
2891/// `GET /v1/pkg/{name}/{version}` — specific version details.
2892fn pkg_get_version_handler(state: &State, name: &str, version: &str) -> Response<std::io::Cursor<Vec<u8>>> {
2893    match load_pkg_record(&state.root, name, version) {
2894        Some(r) => json_response(200, &serde_json::json!({
2895            "name": r.name,
2896            "version": r.version,
2897            "head_op": r.head_op,
2898            "published_at": r.published_at,
2899            "function_names": r.function_names,
2900            "dependencies": r.dependencies,
2901            "ops": r.ops,
2902        })),
2903        None => error_response(404, format!("package {name:?}@{version:?} not found")),
2904    }
2905}
2906
2907/// `GET /v1/pkg/{name}/{version}/archive` — download the source tar.gz.
2908fn pkg_archive_handler(state: &State, name: &str, version: &str) -> Response<std::io::Cursor<Vec<u8>>> {
2909    let gzip = |bytes: Vec<u8>| {
2910        Response::from_data(bytes).with_status_code(200).with_header(
2911            tiny_http::Header::from_bytes(&b"Content-Type"[..], &b"application/gzip"[..]).unwrap(),
2912        )
2913    };
2914
2915    // 1. A stored archive (the `lex pkg publish --registry` upload path).
2916    if let Ok(bytes) = std::fs::read(pkg_archive_path(&state.root, name, version)) {
2917        return gzip(bytes);
2918    }
2919
2920    // 2. An op-log-native release (op push + `POST …/release`, #911) has no
2921    //    stored archive — render one from the pinned op-log head so the
2922    //    package installs like any other (#920).
2923    if let Some(record) = load_pkg_record(&state.root, name, version) {
2924        if let Some(head_op) = record.head_op.clone() {
2925            match render_op_log_archive(state, name, version, &head_op, &record.dependency_specs) {
2926                Ok(bytes) => return gzip(bytes),
2927                Err(e) => {
2928                    return error_response(500, format!("rendering archive for {name:?}@{version:?}: {e}"));
2929                }
2930            }
2931        }
2932    }
2933
2934    error_response(404, format!("archive for {name:?}@{version:?} not found"))
2935}
2936
2937/// The archive path for a single-module head that records no source file of
2938/// its own — every op predating `in_file`. Historical, and kept so already
2939/// released packages in that state do not change name underneath their
2940/// dependents; the cure for them is a re-push, not a rename here (#988).
2941const DEFAULT_ARCHIVE_MODULE: &str = "src/lib.lex";
2942
2943/// Build a gzip-tar package archive (`lex.toml` + the package's `src/*.lex`)
2944/// by rendering the op-log head `head_op` to source — the composition of a
2945/// registry release (#911) with the op-log, so an `op push`-hosted package is
2946/// installable without a separately-uploaded archive (#920).
2947fn render_op_log_archive(
2948    state: &State,
2949    name: &str,
2950    version: &str,
2951    head_op: &str,
2952    dependency_specs: &std::collections::BTreeMap<String, DepSpec>,
2953) -> Result<Vec<u8>, String> {
2954    // De-flatten the head into its source tree — one file for a single-module
2955    // package, or the full `src/*.lex` layout for a multi-module one (#894).
2956    // The same renderer `lex export-git` uses, so the installed source matches
2957    // the git mirror. Either way the *names* come from the head itself; this
2958    // used to invent `src/lib.lex` for the single-module arm (#988).
2959    let (files, lock): (Vec<(String, String)>, Option<String>) = {
2960        let store = state.store.lock().unwrap();
2961        // #943: ship the lock this head was published against, so a consumer
2962        // that installs this package resolves ITS dependencies against the
2963        // pins it was built with. Without it a cached package directory has
2964        // no `lex.lock`, and any caret registry dependency of a dependency
2965        // was unresolvable (`UnlockedRegistryDep`).
2966        let lock = store.committed_lock_inherited(head_op).ok().flatten();
2967        let head = lex_store::render::package_head_at_op(&store, head_op)
2968            .map_err(|e| format!("reading head {head_op}: {e}"))?;
2969        let files = match lex_store::render::render_source(&store, &head)
2970            .map_err(|e| format!("rendering source at {head_op}: {e}"))?
2971        {
2972            lex_store::render::RenderedSource::Single { path, src } => {
2973                // A head that records no file keeps the registry's historical
2974                // name; one that does keeps its own (#988).
2975                vec![(path.unwrap_or_else(|| DEFAULT_ARCHIVE_MODULE.to_string()), src)]
2976            }
2977            lex_store::render::RenderedSource::Multi(tree) => tree.into_iter().collect(),
2978        };
2979        (files, lock)
2980    };
2981
2982    // Reconstruct the manifest from the release record. When the release
2983    // captured full dependency coordinates, emit a faithful `[dependencies]`
2984    // table (carrying both git and vcs references where the source declared
2985    // both); otherwise a bare `[package]` (pre-dependency-capture releases).
2986    let mut manifest = format!("[package]\nname = \"{name}\"\nversion = \"{version}\"\n");
2987    let dep_lines: Vec<String> = dependency_specs
2988        .iter()
2989        .filter_map(|(dep_name, spec)| spec.to_toml_inline().map(|inline| format!("{dep_name} = {inline}")))
2990        .collect();
2991    if !dep_lines.is_empty() {
2992        manifest.push_str("\n[dependencies]\n");
2993        for line in dep_lines {
2994            manifest.push_str(&line);
2995            manifest.push('\n');
2996        }
2997    }
2998    let mut enc = flate2::write::GzEncoder::new(Vec::new(), flate2::Compression::default());
2999    {
3000        let mut ar = tar::Builder::new(&mut enc);
3001        let mut append = |p: &str, data: &[u8]| -> std::io::Result<()> {
3002            let mut h = tar::Header::new_gnu();
3003            h.set_size(data.len() as u64);
3004            h.set_mode(0o644);
3005            h.set_cksum();
3006            ar.append_data(&mut h, p, data)
3007        };
3008        append("lex.toml", manifest.as_bytes()).map_err(|e| e.to_string())?;
3009        if let Some(lock) = &lock {
3010            append("lex.lock", lock.as_bytes()).map_err(|e| e.to_string())?;
3011        }
3012        for (path, src) in &files {
3013            append(path, src.as_bytes()).map_err(|e| e.to_string())?;
3014        }
3015        ar.finish().map_err(|e| e.to_string())?;
3016    }
3017    enc.finish().map_err(|e| e.to_string())
3018}
3019
3020/// `GET /v1/pkg/{name}/head` — head op for a package's latest version.
3021fn pkg_head_handler(state: &State, name: &str) -> Response<std::io::Cursor<Vec<u8>>> {
3022    match load_latest_pkg_record(&state.root, name) {
3023        Some(r) => json_response(200, &serde_json::json!({
3024            "name": r.name,
3025            "version": r.version,
3026            "head_op": r.head_op,
3027        })),
3028        None => error_response(404, format!("package {name:?} not found")),
3029    }
3030}
3031
3032/// `DELETE /v1/pkg/{name}` — retract the latest version of a package.
3033fn pkg_delete_handler(state: &State, name: &str) -> Response<std::io::Cursor<Vec<u8>>> {
3034    let record = match load_latest_pkg_record(&state.root, name) {
3035        Some(r) => r,
3036        None => return error_response(404, format!("package {name:?} not found")),
3037    };
3038
3039    let store = state.store.lock().unwrap();
3040    let branch = store.current_branch();
3041
3042    let head = match store.branch_head(&branch) {
3043        Ok(h) => h,
3044        Err(e) => return error_response(500, format!("branch_head: {e}")),
3045    };
3046
3047    // Build old_fns from this package's function names that are still on
3048    // the branch, reading each AST through the SigId the head names. A
3049    // StageId-keyed read is ambiguous when two live functions differ only
3050    // in name (#826), and here that ambiguity decides what gets REMOVED:
3051    // it could both miss one of this package's functions and match a name
3052    // belonging to another package sharing the stage.
3053    let head_pairs: Vec<(String, String)> = head
3054        .iter()
3055        .map(|(sig, stage)| (sig.clone(), stage.clone()))
3056        .collect();
3057    let old_fns: BTreeMap<String, lex_ast::FnDecl> = store
3058        .get_asts_for_sigs_bulk(&head_pairs)
3059        .into_iter()
3060        .filter_map(|r| r.ok())
3061        .filter_map(|s| match s {
3062            lex_ast::Stage::FnDecl(fd)
3063                if record.function_names.contains(&fd.name) => Some((fd.name.clone(), fd)),
3064            _ => None,
3065        })
3066        .collect();
3067
3068    let new_fns: BTreeMap<String, lex_ast::FnDecl> = BTreeMap::new();
3069    let report = lex_vcs::compute_diff(&old_fns, &new_fns, false);
3070    let empty_imports = lex_vcs::ImportMap::new();
3071
3072    match store.publish_program(&branch, &[], &report, &empty_imports, false) {
3073        Ok(outcome) => {
3074            // Remove the version record and archive, then update the index.
3075            let ver = record.version.clone();
3076            let _ = std::fs::remove_file(pkg_version_path(&state.root, name, &ver));
3077            let _ = std::fs::remove_file(pkg_archive_path(&state.root, name, &ver));
3078            // Update index: remove this version, set latest to previous if any.
3079            if let Some(mut idx) = load_pkg_index(&state.root, name) {
3080                idx.versions.retain(|v| v.version != ver);
3081                idx.latest = idx.versions.last().map(|v| v.version.clone());
3082                if idx.versions.is_empty() {
3083                    let _ = std::fs::remove_dir_all(pkg_name_dir(&state.root, name));
3084                } else {
3085                    let bytes = serde_json::to_vec_pretty(&idx).unwrap_or_default();
3086                    let _ = std::fs::write(pkg_index_path(&state.root, name), bytes);
3087                }
3088            }
3089            json_response(200, &serde_json::json!({
3090                "deleted": name,
3091                "version": ver,
3092                "ops": outcome.ops,
3093                "head_op": outcome.head_op,
3094            }))
3095        }
3096        Err(lex_store::StoreError::TypeError(errs)) => {
3097            error_with_detail(422, "type errors", serde_json::to_value(&errs).unwrap())
3098        }
3099        Err(e) => write_error_response("retract package", e),
3100    }
3101}
3102
3103#[cfg(test)]
3104mod dep_spec_tests {
3105    use super::DepSpec;
3106
3107    #[test]
3108    fn a_dual_spec_renders_both_git_and_vcs_refs_vcs_first() {
3109        let spec = DepSpec {
3110            registry: Some("vcs.lexlang.org/lex-official/lex-schema".into()),
3111            version: Some("^0.9".into()),
3112            git: Some("https://github.com/alpibrusl/lex-schema".into()),
3113            ..Default::default()
3114        };
3115        assert_eq!(
3116            spec.to_toml_inline().as_deref(),
3117            Some("{ registry = \"vcs.lexlang.org/lex-official/lex-schema\", version = \"^0.9\", git = \"https://github.com/alpibrusl/lex-schema\" }"),
3118        );
3119    }
3120
3121    #[test]
3122    fn bare_git_and_bare_registry_specs_render_their_own_keys() {
3123        let git = DepSpec { git: Some("https://x/g".into()), tag: Some("v1".into()), ..Default::default() };
3124        assert_eq!(git.to_toml_inline().as_deref(), Some("{ git = \"https://x/g\", tag = \"v1\" }"));
3125        let reg = DepSpec { registry: Some("vcs/r".into()), version: Some("1.0.0".into()), ..Default::default() };
3126        assert_eq!(reg.to_toml_inline().as_deref(), Some("{ registry = \"vcs/r\", version = \"1.0.0\" }"));
3127    }
3128
3129    #[test]
3130    fn an_empty_spec_renders_nothing() {
3131        assert_eq!(DepSpec::default().to_toml_inline(), None);
3132    }
3133}
3134
3135#[cfg(test)]
3136mod policy_ceiling_tests {
3137    use super::*;
3138    use lex_runtime::Policy;
3139    use std::path::PathBuf;
3140
3141    /// A maximally-permissive policy of the kind a malicious caller
3142    /// would put in a `/v1/run` body: every dangerous effect plus
3143    /// fs over `/`.
3144    fn permissive_request() -> Policy {
3145        Policy {
3146            allow_effects: ["io", "fs_read", "fs_write", "net", "proc"]
3147                .iter()
3148                .map(|s| s.to_string())
3149                .collect(),
3150            allow_fs_read: vec![PathBuf::from("/")],
3151            allow_fs_write: vec![PathBuf::from("/")],
3152            allow_net_host: Vec::new(),
3153            allow_proc: Vec::new(),
3154            allow_approval: Vec::new(),
3155            budget: None,
3156        }
3157    }
3158
3159    #[test]
3160    fn ceiling_drops_effects_the_caller_was_not_granted() {
3161        let ceiling = Policy {
3162            allow_effects: ["io", "time"].iter().map(|s| s.to_string()).collect(),
3163            ..Policy::default()
3164        };
3165        let got = clamp_policy(permissive_request(), &ceiling);
3166        assert!(got.allow_effects.contains("io"));
3167        assert!(!got.allow_effects.contains("proc"), "proc must not survive a ceiling without it");
3168        assert!(!got.allow_effects.contains("fs_write"));
3169        assert!(!got.allow_effects.contains("net"));
3170        // `time` is in the ceiling but not the request → intersection drops it.
3171        assert!(!got.allow_effects.contains("time"));
3172    }
3173
3174    #[test]
3175    fn ceiling_scopes_override_caller_scopes() {
3176        let ceiling = Policy {
3177            allow_effects: ["fs_read"].iter().map(|s| s.to_string()).collect(),
3178            allow_fs_read: vec![PathBuf::from("/srv/tenant")],
3179            ..Policy::default()
3180        };
3181        let got = clamp_policy(permissive_request(), &ceiling);
3182        // Caller asked for "/" but only the ceiling's scope survives —
3183        // an empty/wider caller list can never widen the ceiling.
3184        assert_eq!(got.allow_fs_read, vec![PathBuf::from("/srv/tenant")]);
3185        assert!(got.allow_fs_write.is_empty());
3186        assert!(got.allow_proc.is_empty());
3187        assert!(got.allow_net_host.is_empty());
3188    }
3189
3190    #[test]
3191    fn ceiling_caps_budget_and_prefers_the_smaller() {
3192        // Caller wants unlimited; ceiling caps it.
3193        let mut req = permissive_request();
3194        req.budget = None;
3195        let ceiling = Policy { budget: Some(1_000), ..Policy::default() };
3196        assert_eq!(clamp_policy(req, &ceiling).budget, Some(1_000));
3197
3198        // Caller asks for less than the ceiling → keep the caller's.
3199        let mut req2 = permissive_request();
3200        req2.budget = Some(50);
3201        let ceiling2 = Policy { budget: Some(1_000), ..Policy::default() };
3202        assert_eq!(clamp_policy(req2, &ceiling2).budget, Some(50));
3203    }
3204
3205    #[test]
3206    fn empty_ceiling_is_pure_only() {
3207        let got = clamp_policy(permissive_request(), &Policy::default());
3208        assert!(got.allow_effects.is_empty(), "an empty ceiling grants nothing");
3209        assert!(got.allow_proc.is_empty());
3210        assert!(got.allow_fs_write.is_empty());
3211    }
3212}
3213
3214#[cfg(test)]
3215mod public_read_tests {
3216    use super::*;
3217
3218    /// Write a minimal package (index + per-version record + an archive
3219    /// blob) straight into a temp store, bypassing the publish pipeline.
3220    fn seed_pkg(root: &std::path::Path, name: &str, version: &str) {
3221        let record = PkgRecord {
3222            name: name.to_string(),
3223            version: version.to_string(),
3224            head_op: Some(format!("op-{name}")),
3225            published_at: 1,
3226            function_names: vec![format!("{name}.f")],
3227            dependencies: vec![],
3228            dependency_specs: Default::default(),
3229            ops: vec![],
3230        };
3231        save_pkg_record(root, &record, Some(format!("ARCHIVE:{name}@{version}").as_bytes()))
3232            .expect("seed package");
3233    }
3234
3235    #[test]
3236    fn new_package_defaults_to_private() {
3237        let tmp = tempfile::TempDir::new().unwrap();
3238        seed_pkg(tmp.path(), "lex-schema", "0.9.2");
3239        assert!(!pkg_is_public(tmp.path(), "lex-schema"));
3240        // Unknown packages are also "not public" — never distinguished.
3241        assert!(!pkg_is_public(tmp.path(), "does-not-exist"));
3242    }
3243
3244    #[test]
3245    fn set_visibility_round_trips_and_index_persists() {
3246        let tmp = tempfile::TempDir::new().unwrap();
3247        let state = State::open(tmp.path().to_path_buf()).unwrap();
3248        seed_pkg(tmp.path(), "lex-schema", "0.9.2");
3249
3250        let _ = pkg_set_visibility_handler(&state, "lex-schema", r#"{"visibility":"public"}"#);
3251        assert!(pkg_is_public(tmp.path(), "lex-schema"));
3252        // The version list survives the index rewrite (we don't clobber it).
3253        let idx = load_pkg_index(tmp.path(), "lex-schema").unwrap();
3254        assert_eq!(idx.latest.as_deref(), Some("0.9.2"));
3255        assert_eq!(idx.versions.len(), 1);
3256
3257        let _ = pkg_set_visibility_handler(&state, "lex-schema", r#"{"visibility":"private"}"#);
3258        assert!(!pkg_is_public(tmp.path(), "lex-schema"));
3259    }
3260
3261    #[test]
3262    fn set_visibility_on_unknown_package_is_a_noop() {
3263        let tmp = tempfile::TempDir::new().unwrap();
3264        let state = State::open(tmp.path().to_path_buf()).unwrap();
3265        // No package seeded → handler returns 404 and writes nothing.
3266        let _ = pkg_set_visibility_handler(&state, "ghost", r#"{"visibility":"public"}"#);
3267        assert!(load_pkg_index(tmp.path(), "ghost").is_none());
3268    }
3269
3270    #[test]
3271    fn public_listing_omits_private_packages() {
3272        let tmp = tempfile::TempDir::new().unwrap();
3273        let state = State::open(tmp.path().to_path_buf()).unwrap();
3274        seed_pkg(tmp.path(), "pub-pkg", "1.0.0");
3275        seed_pkg(tmp.path(), "priv-pkg", "1.0.0");
3276        let _ = pkg_set_visibility_handler(&state, "pub-pkg", r#"{"visibility":"public"}"#);
3277
3278        let names = public_pkg_names(tmp.path());
3279        assert_eq!(names, vec!["pub-pkg".to_string()]);
3280    }
3281
3282    #[test]
3283    fn resolve_public_maps_routes() {
3284        let get = Method::Get;
3285        assert_eq!(resolve_public(&get, "").unwrap(), PublicTarget::List);
3286        assert_eq!(resolve_public(&get, "/").unwrap(), PublicTarget::List);
3287        assert_eq!(
3288            resolve_public(&get, "/lex-schema").unwrap(),
3289            PublicTarget::Latest("lex-schema".into())
3290        );
3291        assert_eq!(
3292            resolve_public(&get, "/lex-schema/versions").unwrap(),
3293            PublicTarget::Versions("lex-schema".into())
3294        );
3295        assert_eq!(
3296            resolve_public(&get, "/lex-schema/head").unwrap(),
3297            PublicTarget::Head("lex-schema".into())
3298        );
3299        assert_eq!(
3300            resolve_public(&get, "/lex-schema/0.9.2").unwrap(),
3301            PublicTarget::Version("lex-schema".into(), "0.9.2".into())
3302        );
3303        assert_eq!(
3304            resolve_public(&get, "/lex-schema/0.9.2/archive").unwrap(),
3305            PublicTarget::Archive("lex-schema".into(), "0.9.2".into())
3306        );
3307    }
3308
3309    #[test]
3310    fn resolve_public_rejects_bad_method_and_traversal() {
3311        // Non-GET → 405.
3312        assert_eq!(resolve_public(&Method::Put, "/lex-schema"), Err(405));
3313        assert_eq!(resolve_public(&Method::Post, "").err(), Some(405));
3314        // Path traversal / invalid segments → 404, never a filesystem touch.
3315        assert_eq!(resolve_public(&Method::Get, "/.."), Err(404));
3316        assert_eq!(resolve_public(&Method::Get, "/lex-schema/../etc"), Err(404));
3317        assert_eq!(resolve_public(&Method::Get, "/a/b/c/d"), Err(404));
3318        // Slashes elsewhere can't smuggle a deep path: each segment is checked.
3319        assert!(resolve_public(&Method::Get, "/lex schema").is_err());
3320    }
3321}