1use 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
23fn 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
31fn 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 pub root: PathBuf,
47 pub sessions: Mutex<HashMap<MergeSessionId, ApiMergeSession>>,
54 pub policy_ceiling: Option<Policy>,
75 pub(crate) ops_since: Mutex<crate::ops_since_http::OpsSinceCache>,
79 pub blob_limits: Option<BlobLimits>,
87}
88
89#[derive(Debug, Clone, Copy, PartialEq, Eq)]
91pub struct BlobLimits {
92 pub max_blob_bytes: u64,
94 pub store_quota_bytes: u64,
96 pub max_manifest_entries: usize,
98}
99
100pub const CAPS: &[&str] = &[CAP_FILES_V1];
104
105pub const CAP_FILES_V1: &str = "files-v1";
107
108pub(crate) fn has_cap(header: Option<&str>, cap: &str) -> bool {
110 header.is_some_and(|h| h.split(',').any(|c| c.trim().eq_ignore_ascii_case(cap)))
111}
112
113pub struct ApiMergeSession {
120 pub inner: MergeSession,
121 pub src_branch: String,
122 pub dst_branch: String,
123}
124
125impl State {
126 pub fn open(root: PathBuf) -> anyhow::Result<Self> {
127 Self::open_with_ceiling(root, None)
128 }
129
130 pub fn open_with_ceiling(
135 root: PathBuf,
136 policy_ceiling: Option<Policy>,
137 ) -> anyhow::Result<Self> {
138 Ok(Self {
139 store: Mutex::new(Store::open(&root)?),
140 root,
141 sessions: Mutex::new(HashMap::new()),
142 policy_ceiling,
143 ops_since: Mutex::new(Default::default()),
144 blob_limits: None,
145 })
146 }
147
148 pub fn with_blob_limits(mut self, limits: Option<BlobLimits>) -> Self {
150 self.blob_limits = limits;
151 self
152 }
153
154 pub fn new_with_tenant(tenant_id: &str, store_root: PathBuf) -> anyhow::Result<Self> {
164 validate_tenant_id(tenant_id)?;
165 Self::open(store_root.join(tenant_id))
166 }
167
168 pub fn new_with_tenant_and_ceiling(
173 tenant_id: &str,
174 store_root: PathBuf,
175 policy_ceiling: Option<Policy>,
176 ) -> anyhow::Result<Self> {
177 validate_tenant_id(tenant_id)?;
178 Self::open_with_ceiling(store_root.join(tenant_id), policy_ceiling)
179 }
180}
181
182fn clamp_policy(requested: Policy, ceiling: &Policy) -> Policy {
198 let allow_effects: BTreeSet<String> = requested
199 .allow_effects
200 .intersection(&ceiling.allow_effects)
201 .cloned()
202 .collect();
203 let budget = match (requested.budget, ceiling.budget) {
204 (Some(r), Some(c)) => Some(r.min(c)),
205 (None, Some(c)) => Some(c),
206 (Some(r), None) => Some(r),
207 (None, None) => None,
208 };
209 Policy {
210 allow_effects,
211 allow_fs_read: ceiling.allow_fs_read.clone(),
212 allow_fs_write: ceiling.allow_fs_write.clone(),
213 allow_net_host: ceiling.allow_net_host.clone(),
214 allow_proc: ceiling.allow_proc.clone(),
215 allow_approval: ceiling.allow_approval.clone(),
216 budget,
217 }
218}
219
220fn validate_tenant_id(tenant_id: &str) -> anyhow::Result<()> {
221 if tenant_id.is_empty() {
222 anyhow::bail!("tenant_id must not be empty");
223 }
224 if tenant_id.len() > 64 {
225 anyhow::bail!("tenant_id must be at most 64 bytes");
226 }
227 if !tenant_id
228 .bytes()
229 .all(|b| b.is_ascii_alphanumeric() || b == b'_' || b == b'-')
230 {
231 anyhow::bail!(
232 "tenant_id {tenant_id:?} contains characters outside [A-Za-z0-9_-]"
233 );
234 }
235 Ok(())
236}
237
238#[derive(Debug, Serialize, Deserialize)]
239struct ErrorEnvelope {
240 error: String,
241 #[serde(skip_serializing_if = "Option::is_none")]
242 detail: Option<serde_json::Value>,
243}
244
245pub(crate) fn json_response(status: u16, body: &serde_json::Value) -> Response<std::io::Cursor<Vec<u8>>> {
246 let bytes = serde_json::to_vec(body).unwrap_or_else(|_| b"{}".to_vec());
247 Response::from_data(bytes)
248 .with_status_code(status)
249 .with_header(Header::from_bytes(&b"Content-Type"[..], &b"application/json"[..]).unwrap())
250}
251
252pub(crate) fn error_response(status: u16, msg: impl Into<String>) -> Response<std::io::Cursor<Vec<u8>>> {
253 json_response(status, &serde_json::to_value(ErrorEnvelope {
254 error: msg.into(), detail: None,
255 }).unwrap())
256}
257
258pub(crate) fn error_with_detail(status: u16, msg: impl Into<String>, detail: serde_json::Value)
259 -> Response<std::io::Cursor<Vec<u8>>>
260{
261 json_response(status, &serde_json::to_value(ErrorEnvelope {
262 error: msg.into(), detail: Some(detail),
263 }).unwrap())
264}
265
266pub const UNSATISFIABLE_PAIR_HINT: &str =
269 "republish from source to retire the stranded entry (#995)";
270
271pub(crate) fn unsatisfiable_pair_response(err: &lex_store::StoreError)
284 -> Option<Response<std::io::Cursor<Vec<u8>>>>
285{
286 let lex_store::StoreError::UnsatisfiablePair { sig_id, stage_id, filed_under } = err else {
287 return None;
288 };
289 Some(error_with_detail(422, "UnsatisfiablePair", serde_json::json!({
290 "sig_id": sig_id,
291 "stage_id": stage_id,
292 "filed_under": filed_under,
293 "message": err.to_string(),
294 "hint": UNSATISFIABLE_PAIR_HINT,
295 })))
296}
297
298fn write_error_response(prefix: &str, err: lex_store::StoreError)
304 -> Response<std::io::Cursor<Vec<u8>>>
305{
306 if let Some(resp) = unsatisfiable_pair_response(&err) {
307 return resp;
308 }
309 if let lex_store::StoreError::Contention { branch, attempts } = &err {
310 let body = serde_json::to_vec(&ErrorEnvelope {
311 error: format!("{prefix}: branch '{branch}' is contended (attempts={attempts})"),
312 detail: Some(serde_json::json!({
313 "kind": "contention",
314 "branch": branch,
315 "attempts": attempts,
316 })),
317 }).unwrap_or_else(|_| b"{}".to_vec());
318 return Response::from_data(body)
319 .with_status_code(503)
320 .with_header(Header::from_bytes(&b"Content-Type"[..], &b"application/json"[..]).unwrap())
321 .with_header(Header::from_bytes(&b"Retry-After"[..], &b"1"[..]).unwrap());
322 }
323 if let lex_store::StoreError::BudgetExceeded { session_id, cap, spent_after } = &err {
331 let body = serde_json::to_vec(&ErrorEnvelope {
332 error: format!(
333 "{prefix}: session `{session_id}` budget exceeded \
334 (spent_after={spent_after}, cap={cap})"
335 ),
336 detail: Some(serde_json::json!({
337 "kind": "budget_exceeded",
338 "session_id": session_id,
339 "cap": cap,
340 "spent_after": spent_after,
341 })),
342 }).unwrap_or_else(|_| b"{}".to_vec());
343 return Response::from_data(body)
344 .with_status_code(503)
345 .with_header(Header::from_bytes(&b"Content-Type"[..], &b"application/json"[..]).unwrap())
346 .with_header(Header::from_bytes(&b"Retry-After"[..], &b"0"[..]).unwrap());
347 }
348 error_response(500, format!("{prefix}: {err}"))
349}
350
351pub fn handle(state: Arc<State>, mut req: Request) -> std::io::Result<()> {
352 let method = req.method().clone();
353 let url = req.url().to_string();
354 let path = url.split('?').next().unwrap_or("").to_string();
355 let query = url.split_once('?').map(|(_, q)| q.to_string()).unwrap_or_default();
356
357 let x_lex_user = req.headers().iter()
362 .find(|h| h.field.equiv("x-lex-user"))
363 .map(|h| h.value.as_str().to_string());
364 let x_lex_caps = req.headers().iter()
367 .find(|h| h.field.equiv("x-lex-caps"))
368 .map(|h| h.value.as_str().to_string());
369
370 if matches!(method, Method::Post) && path == "/v1/pkg/publish" {
372 let mut body_bytes: Vec<u8> = Vec::new();
373 let _ = req.as_reader().read_to_end(&mut body_bytes);
374 let resp = pkg_publish_handler(&state, &body_bytes);
375 return req.respond(resp);
376 }
377
378 let mut body = String::new();
379 let _ = req.as_reader().read_to_string(&mut body);
380
381 let resp = route(&state, &method, &path, &query, &body, x_lex_user.as_deref(), x_lex_caps.as_deref());
382 req.respond(resp)
383}
384
385pub fn handle_with_auth<F>(state: Arc<State>, req: Request, auth: F) -> std::io::Result<()>
389where
390 F: FnOnce(&str, &[Header]) -> bool,
391{
392 let path = req.url().split('?').next().unwrap_or("").to_string();
393 if !auth(&path, req.headers()) {
394 return req.respond(
395 Response::from_data(br#"{"error":"unauthorized"}"#.to_vec())
396 .with_status_code(401)
397 .with_header(
398 Header::from_bytes(&b"Content-Type"[..], &b"application/json"[..]).unwrap(),
399 ),
400 );
401 }
402 handle(state, req)
403}
404
405fn route(
406 state: &State,
407 method: &Method,
408 path: &str,
409 query: &str,
410 body: &str,
411 x_lex_user: Option<&str>,
412 x_lex_caps: Option<&str>,
413) -> Response<std::io::Cursor<Vec<u8>>> {
414 match (method, path) {
415 (Method::Get, "/") => crate::web::activity_handler(state),
417 (Method::Get, "/web/branches") => crate::web::branches_handler(state),
418 (Method::Get, "/web/trust") => crate::web::trust_handler(state),
419 (Method::Get, "/web/attention") => crate::web::attention_handler(state),
420 (Method::Get, p) if p.starts_with("/web/branch/") => {
421 let name = &p["/web/branch/".len()..];
422 crate::web::branch_handler(state, name)
423 }
424 (Method::Get, p) if p.starts_with("/web/stage/") => {
425 let id = &p["/web/stage/".len()..];
426 crate::web::stage_html_handler(state, id)
427 }
428 (Method::Post, p) if p.starts_with("/web/stage/") && (
433 p.ends_with("/pin") || p.ends_with("/defer")
434 || p.ends_with("/block") || p.ends_with("/unblock")
435 ) => {
436 let prefix_len = "/web/stage/".len();
437 let last_slash = p.rfind('/').unwrap_or(p.len());
438 let id = &p[prefix_len..last_slash];
439 let verb = &p[last_slash + 1..];
440 let decision = match verb {
441 "pin" => crate::web::WebStageDecision::Pin,
442 "defer" => crate::web::WebStageDecision::Defer,
443 "block" => crate::web::WebStageDecision::Block,
444 "unblock" => crate::web::WebStageDecision::Unblock,
445 _ => unreachable!("matched in outer guard"),
446 };
447 crate::web::stage_decision_handler(state, id, body, decision, x_lex_user)
448 }
449 (Method::Get, "/v1/health") => json_response(200, &serde_json::json!({"ok": true, "caps": CAPS})),
451 (Method::Post, "/v1/parse") => parse_handler(body),
452 (Method::Post, "/v1/check") => check_handler(body),
453 (Method::Post, "/v1/publish") => publish_handler(state, body),
454 (Method::Post, "/v1/patch") => patch_handler(state, body),
455 (Method::Get, p) if p.starts_with("/v1/stage/") => {
456 let suffix = &p["/v1/stage/".len()..];
457 if let Some(id) = suffix.strip_suffix("/attestations") {
460 stage_attestations_handler(state, id)
461 } else {
462 stage_handler(state, suffix)
463 }
464 }
465 (Method::Post, "/v1/run") => run_handler(state, body, false),
466 (Method::Post, "/v1/replay") => run_handler(state, body, true),
467 (Method::Get, p) if p.starts_with("/v1/trace/") => {
468 let id = &p["/v1/trace/".len()..];
469 trace_handler(state, id)
470 }
471 (Method::Get, "/v1/diff") => diff_handler(state, query),
472 (Method::Post, "/v1/merge/start") => merge_start_handler(state, body),
473 (Method::Post, p) if p.starts_with("/v1/merge/") && p.ends_with("/resolve") => {
474 let id = &p["/v1/merge/".len()..p.len() - "/resolve".len()];
475 merge_resolve_handler(state, id, body)
476 }
477 (Method::Post, p) if p.starts_with("/v1/merge/") && p.ends_with("/commit") => {
478 let id = &p["/v1/merge/".len()..p.len() - "/commit".len()];
479 merge_commit_handler(state, id)
480 }
481 (Method::Post, "/v1/ops/batch") => ops_batch_handler(state, body),
483 (Method::Post, "/v1/attestations/batch") => attestations_batch_handler(state, body),
484 (Method::Post, "/v1/stages/batch") => crate::sync_http::stages_batch_handler(state, body),
487 (Method::Post, "/v1/stages/fetch") => crate::sync_http::stages_fetch_handler(state, body),
488 (Method::Post, "/v1/stages/missing") => crate::sync_http::stages_missing_handler(state, body),
489 (Method::Post, "/v1/intents/batch") => crate::sync_http::intents_batch_handler(state, body),
490 (Method::Post, "/v1/intents/fetch") => crate::sync_http::intents_fetch_handler(state, body),
491 (Method::Post, "/v1/locks/batch") => crate::sync_http::locks_batch_handler(state, body),
494 (Method::Post, "/v1/locks/fetch") => crate::sync_http::locks_fetch_handler(state, body),
495 (Method::Post, "/v1/blobs/missing") => crate::sync_http::blobs_missing_handler(state, body),
498 (Method::Post, "/v1/blobs/batch") => crate::sync_http::blobs_batch_handler(state, body),
499 (Method::Post, "/v1/blobs/fetch") => crate::sync_http::blobs_fetch_handler(state, body),
500 (Method::Post, "/v1/issues/batch") => crate::sync_http::issues_batch_handler(state, body),
503 (Method::Post, "/v1/issues/fetch") => crate::sync_http::issues_fetch_handler(state, body),
504 (Method::Get, "/v1/issues/list") => crate::sync_http::issues_list_handler(state),
505 (Method::Get, "/v1/issues") => crate::issues_http::issues_state_handler(state),
509 (Method::Get, "/v1/projects") => crate::issues_http::projects_handler(state),
510 (Method::Get, p) if p.starts_with("/v1/issues/") => {
511 crate::issues_http::issue_detail_handler(state, &p["/v1/issues/".len()..])
512 }
513 (Method::Get, "/v1/review/inbox") => crate::review_http::review_inbox_handler(state, query),
517 (Method::Post, "/v1/review/verdict") => crate::review_http::review_verdict_handler(state, body),
518 (Method::Get, "/v1/branches") => crate::branches_http::branches_list_handler(state),
519 (Method::Post, "/v1/branches") => crate::branches_http::branch_create_handler(state, body),
520 (Method::Post, p) if p.starts_with("/v1/branches/") && p.ends_with("/checkout") => {
521 let name = &p["/v1/branches/".len()..p.len() - "/checkout".len()];
522 crate::branches_http::branch_checkout_handler(state, name)
523 }
524 (Method::Get, p) if p.starts_with("/v1/branches/") && p.ends_with("/head") => {
527 let name = &p["/v1/branches/".len()..p.len() - "/head".len()];
528 crate::branches_http::branch_head_handler(state, name)
529 }
530 (Method::Post, p) if p.starts_with("/v1/branches/") && p.ends_with("/head") => {
531 let name = &p["/v1/branches/".len()..p.len() - "/head".len()];
532 crate::branches_http::branch_advance_head_handler(state, name, body)
533 }
534 (Method::Get, "/v1/ops/since") => crate::ops_since_http::ops_since_handler(state, query, x_lex_caps),
538 (Method::Get, "/v1/attestations/since") => attestations_since_handler(state, query),
539 (Method::Get, "/v1/pkg") => pkg_list_handler(state),
543 (Method::Put, p) if p.starts_with("/v1/pkg/") && p.ends_with("/visibility") => {
547 let name = &p["/v1/pkg/".len()..p.len() - "/visibility".len()];
548 pkg_set_visibility_handler(state, name, body)
549 }
550 (Method::Post, p) if p.starts_with("/v1/pkg/") && p.ends_with("/release") => {
553 let name = &p["/v1/pkg/".len()..p.len() - "/release".len()];
554 pkg_release_handler(state, name, body)
555 }
556 (Method::Get, p) if p.starts_with("/v1/pkg/") && p.ends_with("/head") => {
557 let name = &p["/v1/pkg/".len()..p.len() - "/head".len()];
558 pkg_head_handler(state, name)
559 }
560 (Method::Get, p) if p.starts_with("/v1/pkg/") && p.ends_with("/versions") => {
561 let name = &p["/v1/pkg/".len()..p.len() - "/versions".len()];
562 pkg_versions_handler(state, name)
563 }
564 (Method::Get, p) if p.starts_with("/v1/pkg/") && p.ends_with("/api-diff") => {
565 let name = &p["/v1/pkg/".len()..p.len() - "/api-diff".len()];
566 pkg_api_diff_handler(state, name, query)
567 }
568 (Method::Get, p) if p.starts_with("/v1/pkg/") && p.ends_with("/archive") => {
570 let inner = &p["/v1/pkg/".len()..p.len() - "/archive".len()];
571 if let Some((name, version)) = inner.split_once('/') {
573 pkg_archive_handler(state, name, version)
574 } else {
575 error_response(400, "expected /v1/pkg/{name}/{version}/archive")
576 }
577 }
578 (Method::Get, p) if p.starts_with("/v1/pkg/") && p["/v1/pkg/".len()..].contains('/') => {
580 let inner = &p["/v1/pkg/".len()..];
581 if let Some((name, version)) = inner.split_once('/') {
582 pkg_get_version_handler(state, name, version)
583 } else {
584 error_response(400, "expected /v1/pkg/{name}/{version}")
585 }
586 }
587 (Method::Get, p) if p.starts_with("/v1/pkg/") => {
588 let name = &p["/v1/pkg/".len()..];
589 pkg_get_handler(state, name)
590 }
591 (Method::Delete, p) if p.starts_with("/v1/pkg/") => {
592 let name = &p["/v1/pkg/".len()..];
593 pkg_delete_handler(state, name)
594 }
595 _ => error_response(404, format!("unknown route: {method:?} {path}")),
596 }
597}
598
599#[derive(Deserialize)]
600struct ParseReq { source: String }
601
602fn parse_handler(body: &str) -> Response<std::io::Cursor<Vec<u8>>> {
603 let req: ParseReq = match serde_json::from_str(body) {
604 Ok(r) => r, Err(e) => return error_response(400, format!("bad request: {e}")),
605 };
606 match load_program_from_str(&req.source) {
607 Ok(prog) => {
608 let stages = canonicalize_program(&prog);
609 json_response(200, &serde_json::to_value(&stages).unwrap())
610 }
611 Err(e) => error_response(400, format!("syntax error: {e}")),
612 }
613}
614
615pub(crate) fn check_handler(body: &str) -> Response<std::io::Cursor<Vec<u8>>> {
616 let req: ParseReq = match serde_json::from_str(body) {
617 Ok(r) => r, Err(e) => return error_response(400, format!("bad request: {e}")),
618 };
619 let prog = match load_program_from_str(&req.source) {
620 Ok(p) => p, Err(e) => return error_response(400, format!("syntax error: {e}")),
621 };
622 let stages = canonicalize_program(&prog);
623 match lex_types::check_program(&stages) {
624 Ok(_) => json_response(200, &serde_json::json!({"ok": true})),
625 Err(errs) => json_response(422, &serde_json::to_value(&errs).unwrap()),
626 }
627}
628
629#[derive(Deserialize)]
630struct PublishReq { source: String, #[serde(default)] activate: bool }
631
632pub(crate) fn publish_handler(state: &State, body: &str) -> Response<std::io::Cursor<Vec<u8>>> {
633 let req: PublishReq = match serde_json::from_str(body) {
634 Ok(r) => r, Err(e) => return error_response(400, format!("bad request: {e}")),
635 };
636 let prog = match load_program_from_str(&req.source) {
637 Ok(p) => p, Err(e) => return error_response(400, format!("syntax error: {e}")),
638 };
639 let mut stages = canonicalize_program(&prog);
643 if let Err(errs) = lex_types::check_and_rewrite_program(&mut stages) {
644 return error_with_detail(422, "type errors", serde_json::to_value(&errs).unwrap());
645 }
646 let example_errors = lex_runtime::evaluate_examples(&stages);
652 if !example_errors.is_empty() {
653 return error_with_detail(422, "example mismatch",
654 serde_json::to_value(&example_errors).unwrap_or_default());
655 }
656
657 let store = state.store.lock().unwrap();
658 let branch = store.current_branch();
659
660 let old_head = match store.branch_head(&branch) {
662 Ok(h) => h,
663 Err(e) => return error_response(500, format!("branch_head: {e}")),
664 };
665 let old_pairs: Vec<(String, String)> =
672 old_head.iter().map(|(sig, stg)| (sig.clone(), stg.clone())).collect();
673 let old_head_stages: Vec<lex_ast::Stage> =
674 store.get_asts_for_sigs_bulk(&old_pairs).into_iter().filter_map(Result::ok).collect();
675 let old_fns = stage_fns(&old_head_stages);
676 let new_fns = stage_fns(&stages);
677 let old_types = stage_types(&old_head_stages);
678 let new_types = stage_types(&stages);
679 let report =
680 lex_vcs::compute_diff_with_types(&old_fns, &new_fns, &old_types, &new_types, false);
681
682 let mut new_imports: lex_vcs::ImportMap = lex_vcs::ImportMap::new();
684 {
685 let entry = new_imports.entry("<source>".into()).or_default();
686 for s in &stages {
687 if let lex_ast::Stage::Import(im) = s {
688 entry.insert(lex_vcs::ImportRef {
689 reference: im.reference.clone(),
690 alias: im.alias.clone(),
691 });
692 }
693 }
694 }
695
696 match store.publish_program(&branch, &stages, &report, &new_imports, req.activate) {
697 Ok(outcome) => {
698 record_examples_for_publish(&store, &stages, &outcome);
702 json_response(200, &serde_json::json!({
703 "ops": outcome.ops,
704 "head_op": outcome.head_op,
705 }))
706 }
707 Err(lex_store::StoreError::TypeError(errs)) => {
715 error_with_detail(422, "type errors", serde_json::to_value(&errs).unwrap())
716 }
717 Err(e) => write_error_response("publish_program", e),
718 }
719}
720
721#[derive(Deserialize)]
722struct PatchReq {
723 stage_id: String,
724 patch: lex_ast::Patch,
725 #[serde(default)] activate: bool,
726}
727
728fn patch_handler(state: &State, body: &str) -> Response<std::io::Cursor<Vec<u8>>> {
731 let req: PatchReq = match serde_json::from_str(body) {
732 Ok(r) => r, Err(e) => return error_response(400, format!("bad request: {e}")),
733 };
734 let store = state.store.lock().unwrap();
735
736 let original = match store.get_ast(&req.stage_id) {
738 Ok(s) => s, Err(e) => return error_response(404, format!("stage: {e}")),
739 };
740
741 let patched = match lex_ast::apply_patch(&original, &req.patch) {
743 Ok(s) => s,
744 Err(e) => return error_with_detail(422, "patch failed",
745 serde_json::to_value(&e).unwrap_or_default()),
746 };
747
748 let branch = store.current_branch();
759
760 let sig = match lex_ast::sig_id(&patched) {
762 Some(s) => s,
763 None => return error_response(500, "patched stage has no sig_id"),
764 };
765
766 let new_id = match store.publish(&patched) {
769 Ok(id) => id, Err(e) => return error_response(500, format!("publish: {e}")),
770 };
771
772 let original_effects: std::collections::BTreeSet<String> = match &original {
774 lex_ast::Stage::FnDecl(fd) => fd.effects.iter().map(|e| e.name.clone()).collect(),
775 _ => std::collections::BTreeSet::new(),
776 };
777 let patched_effects: std::collections::BTreeSet<String> = match &patched {
778 lex_ast::Stage::FnDecl(fd) => fd.effects.iter().map(|e| e.name.clone()).collect(),
779 _ => std::collections::BTreeSet::new(),
780 };
781 let head_now = match store.get_branch(&branch) {
782 Ok(b) => b.and_then(|b| b.head_op),
783 Err(e) => return error_response(500, format!("get_branch: {e}")),
784 };
785 let kind = if original_effects != patched_effects {
786 let from_budget = lex_vcs::operation_budget_from_effects(&original_effects);
793 let to_budget = lex_vcs::operation_budget_from_effects(&patched_effects);
794 let to_sig_id = lex_ast::sig_id(&patched).filter(|s| *s != sig);
798 lex_vcs::OperationKind::ChangeEffectSig {
799 sig_id: sig.clone(),
800 from_stage_id: req.stage_id.clone(),
801 to_stage_id: new_id.clone(),
802 from_effects: original_effects,
803 to_effects: patched_effects,
804 from_budget,
805 to_budget,
806 to_sig_id,
807 }
808 } else {
809 let budget = lex_vcs::operation_budget_from_effects(&original_effects);
810 lex_vcs::OperationKind::ModifyBody {
811 sig_id: sig.clone(),
812 from_stage_id: req.stage_id.clone(),
813 to_stage_id: new_id.clone(),
814 from_budget: budget,
815 to_budget: budget,
816 to_sig_id: lex_ast::sig_id(&patched).filter(|s| *s != sig),
819 }
820 };
821 let transition = lex_store::transition_for_kind(&kind);
825 let op = lex_vcs::Operation::new(
826 kind,
827 head_now.into_iter().collect::<Vec<_>>(),
828 );
829 let op_id = match store.apply_operation_gated(&branch, op, transition) {
830 Ok(id) => id,
831 Err(lex_store::StoreError::TypeError(errs)) => return error_with_detail(
832 422, "type errors after patch", serde_json::to_value(&errs).unwrap_or_default()),
833 Err(e) => return write_error_response("apply_operation_gated", e),
834 };
835 if req.activate {
836 if let Err(e) = store.activate(&new_id) {
837 return error_response(500, format!("activate: {e}"));
838 }
839 }
840
841 let status = format!("{:?}",
842 store.get_status(&new_id).unwrap_or(lex_store::StageStatus::Draft)).to_lowercase();
843 json_response(200, &serde_json::json!({
844 "old_stage_id": req.stage_id,
845 "new_stage_id": new_id,
846 "sig_id": sig,
847 "status": status,
848 "op_id": op_id,
849 }))
850}
851
852pub(crate) fn stage_handler(state: &State, id: &str) -> Response<std::io::Cursor<Vec<u8>>> {
853 let store = state.store.lock().unwrap();
854 let meta = match store.get_metadata(id) {
855 Ok(m) => m, Err(e) => return error_response(404, format!("{e}")),
856 };
857 let ast = match store.get_ast(id) {
858 Ok(a) => a, Err(e) => return error_response(404, format!("{e}")),
859 };
860 let status = format!("{:?}", store.get_status(id).unwrap_or(lex_store::StageStatus::Draft)).to_lowercase();
861 json_response(200, &serde_json::json!({
862 "metadata": meta,
863 "ast": ast,
864 "status": status,
865 }))
866}
867
868pub(crate) fn stage_attestations_handler(state: &State, id: &str) -> Response<std::io::Cursor<Vec<u8>>> {
877 let store = state.store.lock().unwrap();
878 if let Err(e) = store.get_metadata(id) {
879 return error_response(404, format!("{e}"));
880 }
881 let log = match store.attestation_log() {
882 Ok(l) => l,
883 Err(e) => return error_response(500, format!("attestation log: {e}")),
884 };
885 let mut listing = match log.list_for_stage(&id.to_string()) {
886 Ok(v) => v,
887 Err(e) => return error_response(500, format!("list_for_stage: {e}")),
888 };
889 listing.sort_by_key(|a| std::cmp::Reverse(a.timestamp));
890 json_response(200, &serde_json::json!({"attestations": listing}))
891}
892
893#[derive(Deserialize, Default)]
894struct PolicyJson {
895 #[serde(default)] allow_effects: Vec<String>,
896 #[serde(default)] allow_fs_read: Vec<String>,
897 #[serde(default)] allow_fs_write: Vec<String>,
898 #[serde(default)] budget: Option<u64>,
899}
900
901impl PolicyJson {
902 fn into_policy(self) -> Policy {
903 Policy {
904 allow_effects: self.allow_effects.into_iter().collect::<BTreeSet<_>>(),
905 allow_fs_read: self.allow_fs_read.into_iter().map(PathBuf::from).collect(),
906 allow_fs_write: self.allow_fs_write.into_iter().map(PathBuf::from).collect(),
907 allow_net_host: Vec::new(),
908 allow_proc: Vec::new(),
909 allow_approval: Vec::new(),
910 budget: self.budget,
911 }
912 }
913}
914
915#[derive(Deserialize)]
916struct RunReq {
917 source: String,
918 #[serde(rename = "fn")] func: String,
919 #[serde(default)] args: Vec<serde_json::Value>,
920 #[serde(default)] policy: PolicyJson,
921 #[serde(default)] overrides: IndexMap<String, serde_json::Value>,
922}
923
924pub(crate) fn run_handler(state: &State, body: &str, with_overrides: bool) -> Response<std::io::Cursor<Vec<u8>>> {
925 let req: RunReq = match serde_json::from_str(body) {
926 Ok(r) => r, Err(e) => return error_response(400, format!("bad request: {e}")),
927 };
928 let prog = match load_program_from_str(&req.source) {
929 Ok(p) => p, Err(e) => return error_response(400, format!("syntax error: {e}")),
930 };
931 let stages = canonicalize_program(&prog);
932 if let Err(errs) = lex_types::check_program(&stages) {
933 return error_with_detail(422, "type errors", serde_json::to_value(&errs).unwrap());
934 }
935 let bc = compile_program(&stages);
936 let mut policy = req.policy.into_policy();
937 if let Some(ceiling) = &state.policy_ceiling {
943 policy = clamp_policy(policy, ceiling);
944 }
945 if let Err(violations) = check_policy(&bc, &policy) {
946 return error_with_detail(403, "policy violation", serde_json::to_value(&violations).unwrap());
947 }
948
949 let mut recorder = lex_trace::Recorder::new();
950 if with_overrides && !req.overrides.is_empty() {
951 recorder = recorder.with_overrides(req.overrides);
952 }
953 let handle = recorder.handle();
954 let handler = DefaultHandler::new(policy);
955 let mut vm = Vm::with_handler(&bc, Box::new(handler));
956 vm.set_tracer(Box::new(recorder));
957
958 let vargs: Vec<Value> = req.args.iter().map(json_to_value).collect();
959 let started = std::time::SystemTime::now().duration_since(std::time::UNIX_EPOCH).unwrap().as_secs();
960 let result = vm.call(&req.func, vargs);
961 let ended = std::time::SystemTime::now().duration_since(std::time::UNIX_EPOCH).unwrap().as_secs();
962
963 let store = state.store.lock().unwrap();
964 let (root_out, root_err, status) = match &result {
965 Ok(v) => (Some(value_to_json(v)), None, 200u16),
966 Err(e) => (None, Some(format!("{e}")), 200u16),
967 };
968 let tree = handle.finalize(req.func.clone(), serde_json::Value::Null,
969 root_out.clone(), root_err.clone(), started, ended);
970 let run_id = match store.save_trace(&tree) {
971 Ok(id) => id,
972 Err(e) => return error_response(500, format!("save_trace: {e}")),
973 };
974
975 let mut body = serde_json::json!({
976 "run_id": run_id,
977 "output": root_out,
978 });
979 if let Some(err) = root_err {
980 body["error"] = serde_json::Value::String(err);
981 }
982 json_response(status, &body)
983}
984
985fn trace_handler(state: &State, id: &str) -> Response<std::io::Cursor<Vec<u8>>> {
986 let store = state.store.lock().unwrap();
987 match store.load_trace(id) {
988 Ok(t) => json_response(200, &serde_json::to_value(&t).unwrap()),
989 Err(e) => error_response(404, format!("{e}")),
990 }
991}
992
993fn diff_handler(state: &State, query: &str) -> Response<std::io::Cursor<Vec<u8>>> {
994 let mut a = None;
995 let mut b = None;
996 for kv in query.split('&') {
997 if let Some((k, v)) = kv.split_once('=') {
998 match k { "a" => a = Some(v.to_string()), "b" => b = Some(v.to_string()), _ => {} }
999 }
1000 }
1001 let (Some(a), Some(b)) = (a, b) else {
1002 return error_response(400, "missing a or b query params");
1003 };
1004 let store = state.store.lock().unwrap();
1005 let ta = match store.load_trace(&a) { Ok(t) => t, Err(e) => return error_response(404, format!("a: {e}")) };
1006 let tb = match store.load_trace(&b) { Ok(t) => t, Err(e) => return error_response(404, format!("b: {e}")) };
1007 match lex_trace::diff_runs(&ta, &tb) {
1008 Some(d) => json_response(200, &serde_json::to_value(&d).unwrap()),
1009 None => json_response(200, &serde_json::json!({"divergence": null})),
1010 }
1011}
1012
1013fn json_to_value(v: &serde_json::Value) -> Value { Value::from_json(v) }
1014
1015fn value_to_json(v: &Value) -> serde_json::Value { v.to_json() }
1016
1017#[derive(Deserialize)]
1018struct MergeStartReq {
1019 src_branch: String,
1020 dst_branch: String,
1021}
1022
1023fn merge_start_handler(state: &State, body: &str) -> Response<std::io::Cursor<Vec<u8>>> {
1034 let req: MergeStartReq = match serde_json::from_str(body) {
1035 Ok(r) => r, Err(e) => return error_response(400, format!("bad request: {e}")),
1036 };
1037 let store = state.store.lock().unwrap();
1038 let src_head = match store.get_branch(&req.src_branch) {
1039 Ok(Some(b)) => b.head_op,
1040 Ok(None) => return error_response(404, format!("unknown src branch `{}`", req.src_branch)),
1041 Err(e) => return error_response(500, format!("src branch read: {e}")),
1042 };
1043 let dst_head = match store.get_branch(&req.dst_branch) {
1044 Ok(Some(b)) => b.head_op,
1045 Ok(None) => return error_response(404, format!("unknown dst branch `{}`", req.dst_branch)),
1046 Err(e) => return error_response(500, format!("dst branch read: {e}")),
1047 };
1048 let log = match lex_vcs::OpLog::open(store.root()) {
1049 Ok(l) => l,
1050 Err(e) => return error_response(500, format!("op log: {e}")),
1051 };
1052 let merge_id = mint_merge_id();
1056 let session = match MergeSession::start(
1057 merge_id.clone(),
1058 &log,
1059 src_head.as_ref(),
1060 dst_head.as_ref(),
1061 ) {
1062 Ok(s) => s,
1063 Err(e) => return error_response(500, format!("merge start: {e}")),
1064 };
1065 let conflicts: Vec<&lex_vcs::ConflictRecord> = session.remaining_conflicts();
1066 let auto_resolved_count = session.auto_resolved.len();
1067 let body = serde_json::json!({
1068 "merge_id": merge_id,
1069 "src_head": session.src_head,
1070 "dst_head": session.dst_head,
1071 "lca": session.lca,
1072 "conflicts": conflicts,
1073 "auto_resolved_count": auto_resolved_count,
1074 });
1075 drop(conflicts);
1076 drop(store);
1077 let wrapped = ApiMergeSession {
1078 inner: session,
1079 src_branch: req.src_branch,
1080 dst_branch: req.dst_branch,
1081 };
1082 state.sessions.lock().unwrap().insert(merge_id, wrapped);
1083 json_response(200, &body)
1084}
1085
1086#[derive(Deserialize)]
1087struct MergeResolveReq {
1088 resolutions: Vec<MergeResolveEntry>,
1093}
1094
1095#[derive(Deserialize)]
1096struct MergeResolveEntry {
1097 conflict_id: String,
1098 resolution: lex_vcs::Resolution,
1099}
1100
1101fn merge_resolve_handler(
1112 state: &State,
1113 merge_id: &str,
1114 body: &str,
1115) -> Response<std::io::Cursor<Vec<u8>>> {
1116 let req: MergeResolveReq = match serde_json::from_str(body) {
1117 Ok(r) => r, Err(e) => return error_response(400, format!("bad request: {e}")),
1118 };
1119 let mut sessions = state.sessions.lock().unwrap();
1120 let Some(wrapped) = sessions.get_mut(merge_id) else {
1121 return error_response(404, format!("unknown merge_id `{merge_id}`"));
1122 };
1123 let pairs: Vec<(String, lex_vcs::Resolution)> = req.resolutions.into_iter()
1124 .map(|e| (e.conflict_id, e.resolution))
1125 .collect();
1126 let store = state.store.lock().unwrap();
1131 let checker = lex_store::MergeResolutionChecker::new(&store, wrapped.dst_branch.clone());
1132 let verdicts = wrapped.inner.resolve_checked(pairs, &checker);
1133 drop(store);
1134 let remaining: Vec<&lex_vcs::ConflictRecord> = wrapped.inner.remaining_conflicts();
1135 let body = serde_json::json!({
1136 "verdicts": verdicts,
1137 "remaining_conflicts": remaining,
1138 });
1139 json_response(200, &body)
1140}
1141
1142fn merge_commit_handler(
1164 state: &State,
1165 merge_id: &str,
1166) -> Response<std::io::Cursor<Vec<u8>>> {
1167 use std::collections::BTreeMap;
1168 let wrapped = match state.sessions.lock().unwrap().remove(merge_id) {
1169 Some(w) => w,
1170 None => return error_response(404, format!("unknown merge_id `{merge_id}`")),
1171 };
1172 let dst_branch = wrapped.dst_branch.clone();
1173 let src_head = wrapped.inner.src_head.clone();
1174 let dst_head = wrapped.inner.dst_head.clone();
1175 let auto_resolved = wrapped.inner.auto_resolved.clone();
1176
1177 let mut entries: BTreeMap<lex_vcs::SigId, Option<lex_vcs::StageId>> = BTreeMap::new();
1180
1181 for outcome in &auto_resolved {
1183 if let lex_vcs::MergeOutcome::Src { sig_id, stage_id } = outcome {
1184 entries.insert(sig_id.clone(), stage_id.clone());
1185 }
1186 }
1187
1188 let resolved = match wrapped.inner.commit() {
1190 Ok(r) => r,
1191 Err(lex_vcs::CommitError::ConflictsRemaining(ids)) => {
1192 return error_with_detail(
1196 422,
1197 "conflicts remaining",
1198 serde_json::json!({"unresolved": ids}),
1199 );
1200 }
1201 };
1202
1203 for (conflict_id, resolution) in resolved {
1204 match resolution {
1205 lex_vcs::Resolution::TakeOurs => {
1206 }
1208 lex_vcs::Resolution::TakeTheirs => {
1209 match resolve_take_theirs(state, &src_head, &conflict_id) {
1219 Ok(stage_id) => {
1220 entries.insert(conflict_id.clone(), stage_id);
1221 }
1222 Err(e) => return error_response(500, format!("resolve take_theirs: {e}")),
1223 }
1224 }
1225 lex_vcs::Resolution::Custom { op } => {
1226 match op.kind.merge_target() {
1235 Some((sig, stage)) => {
1236 if sig != conflict_id {
1237 return error_with_detail(
1238 422,
1239 "custom op targets a different sig than the conflict",
1240 serde_json::json!({
1241 "conflict_id": conflict_id,
1242 "op_targets": sig,
1243 }),
1244 );
1245 }
1246 entries.insert(conflict_id, stage);
1247 }
1248 None => {
1249 return error_with_detail(
1250 422,
1251 "custom op kind doesn't yield a single sig→stage delta",
1252 serde_json::json!({
1253 "conflict_id": conflict_id,
1254 "kind": serde_json::to_value(&op.kind).unwrap_or(serde_json::Value::Null),
1255 }),
1256 );
1257 }
1258 }
1259 }
1260 lex_vcs::Resolution::Defer => {
1261 return error_response(500, "internal: Defer slipped past commit gate");
1263 }
1264 }
1265 }
1266
1267 let resolved_count = entries.len();
1268 let mut parents: Vec<lex_vcs::OpId> = Vec::new();
1269 if let Some(d) = dst_head { parents.push(d); }
1270 if let Some(s) = src_head { parents.push(s); }
1271 let op = lex_vcs::Operation::new(
1272 lex_vcs::OperationKind::Merge { resolved: resolved_count },
1273 parents,
1274 );
1275 let transition = lex_vcs::StageTransition::Merge { entries };
1276 let store = state.store.lock().unwrap();
1277 match store.apply_merge_op_gated(&dst_branch, op, transition) {
1280 Ok(new_head_op) => json_response(200, &serde_json::json!({
1281 "new_head_op": new_head_op,
1282 "dst_branch": dst_branch,
1283 })),
1284 Err(lex_store::StoreError::TypeError(errs)) => error_with_detail(
1285 422, "merged program has type errors", serde_json::to_value(&errs).unwrap_or_default()),
1286 Err(e @ lex_store::StoreError::DependencyConflict { .. }) => {
1292 let detail = match &e {
1293 lex_store::StoreError::DependencyConflict { package, dst_version, src_version } => {
1294 serde_json::json!({
1295 "kind": "dependency_conflict",
1296 "package": package,
1297 "dst_branch": dst_branch,
1298 "dst_version": dst_version,
1299 "src_branch": wrapped.src_branch,
1300 "src_version": src_version,
1301 })
1302 }
1303 _ => serde_json::Value::Null,
1304 };
1305 error_with_detail(409, e.to_string(), detail)
1306 }
1307 Err(e) => write_error_response("apply merge op", e),
1308 }
1309}
1310
1311fn resolve_take_theirs(
1316 state: &State,
1317 src_head: &Option<lex_vcs::OpId>,
1318 sig: &lex_vcs::SigId,
1319) -> std::io::Result<Option<lex_vcs::StageId>> {
1320 let store = state.store.lock().unwrap();
1321 let log = lex_vcs::OpLog::open(store.root())?;
1322 let Some(head) = src_head.as_ref() else { return Ok(None); };
1323 let mut current: Option<lex_vcs::StageId> = None;
1326 for record in log.walk_forward(head, None)? {
1327 match &record.produces {
1328 lex_vcs::StageTransition::Create { sig_id, stage_id }
1329 if sig_id == sig => { current = Some(stage_id.clone()); }
1330 lex_vcs::StageTransition::Replace { sig_id, to, .. }
1331 if sig_id == sig => { current = Some(to.clone()); }
1332 lex_vcs::StageTransition::Remove { sig_id, .. }
1333 if sig_id == sig => { current = None; }
1334 lex_vcs::StageTransition::Rename { from, to, body_stage_id }
1335 if from == sig || to == sig => {
1336 if from == sig { current = None; }
1337 if to == sig { current = Some(body_stage_id.clone()); }
1338 }
1339 lex_vcs::StageTransition::Merge { entries } => {
1340 if let Some(opt) = entries.get(sig) {
1341 current = opt.clone();
1342 }
1343 }
1344 _ => {}
1345 }
1346 }
1347 Ok(current)
1348}
1349
1350fn mint_merge_id() -> MergeSessionId {
1351 use std::sync::atomic::{AtomicU64, Ordering};
1352 static COUNTER: AtomicU64 = AtomicU64::new(0);
1353 let nanos = SystemTime::now()
1354 .duration_since(UNIX_EPOCH)
1355 .map(|d| d.as_nanos())
1356 .unwrap_or(0);
1357 let n = COUNTER.fetch_add(1, Ordering::Relaxed);
1358 format!("merge_{nanos:x}_{n:x}")
1359}
1360
1361pub(crate) fn ops_batch_handler(state: &State, body: &str)
1397 -> Response<std::io::Cursor<Vec<u8>>>
1398{
1399 let records: Vec<lex_vcs::OperationRecord> = match serde_json::from_str(body) {
1400 Ok(r) => r,
1401 Err(e) => return error_response(400,
1402 format!("body must be a JSON array of OperationRecord: {e}")),
1403 };
1404 let store = state.store.lock().unwrap();
1405 let log = match lex_vcs::OpLog::open(store.root()) {
1406 Ok(l) => l,
1407 Err(e) => return error_response(500, format!("opening op log: {e}")),
1408 };
1409
1410 let mut batch_ids: std::collections::BTreeSet<lex_vcs::OpId> =
1418 std::collections::BTreeSet::new();
1419 for rec in &records {
1420 let expected = rec.op.op_id();
1421 if expected != rec.op_id {
1422 return error_with_detail(409, "OpIdMismatch", serde_json::json!({
1423 "supplied": rec.op_id,
1424 "expected": expected,
1425 }));
1426 }
1427 for parent in &rec.op.parents {
1428 let known = match log.get(parent) {
1429 Ok(Some(_)) => true,
1430 Ok(None) => false,
1431 Err(e) => return error_response(500, format!("op log read: {e}")),
1432 };
1433 if !known && !batch_ids.contains(parent) {
1434 return error_with_detail(422, "MissingParent", serde_json::json!({
1435 "op_id": rec.op_id,
1436 "missing_parent": parent,
1437 }));
1438 }
1439 }
1440 if let Err(resp) = crate::sync_http::check_set_files(&store, state.blob_limits, rec) {
1441 return resp;
1442 }
1443 batch_ids.insert(rec.op_id.clone());
1444 }
1445
1446 let mut added = 0usize;
1449 let mut added_ids: Vec<&lex_vcs::OpId> = Vec::new();
1450 for rec in &records {
1451 let already_present = matches!(log.get(&rec.op_id), Ok(Some(_)));
1452 match log.put(rec) {
1453 Ok(()) => {
1454 if !already_present {
1455 added += 1;
1456 added_ids.push(&rec.op_id);
1457 }
1458 }
1459 Err(e) => return error_response(500, format!("op log write: {e}")),
1460 }
1461 }
1462
1463 json_response(200, &serde_json::json!({
1464 "received": records.len(),
1465 "added": added,
1466 "skipped": records.len() - added,
1467 "added_ids": added_ids,
1468 }))
1469}
1470
1471pub(crate) fn attestations_batch_handler(state: &State, body: &str)
1493 -> Response<std::io::Cursor<Vec<u8>>>
1494{
1495 let attestations: Vec<lex_vcs::Attestation> = match serde_json::from_str(body) {
1496 Ok(a) => a,
1497 Err(e) => return error_response(400,
1498 format!("body must be a JSON array of Attestation: {e}")),
1499 };
1500 let store = state.store.lock().unwrap();
1501 let log = match store.attestation_log() {
1502 Ok(l) => l,
1503 Err(e) => return error_response(500, format!("opening attestation log: {e}")),
1504 };
1505 let op_log = match lex_vcs::OpLog::open(store.root()) {
1506 Ok(l) => l,
1507 Err(e) => return error_response(500, format!("opening op log: {e}")),
1508 };
1509
1510 for att in &attestations {
1512 let expected = lex_vcs::Attestation::with_timestamp(
1515 att.stage_id.clone(),
1516 att.op_id.clone(),
1517 att.intent_id.clone(),
1518 att.kind.clone(),
1519 att.result.clone(),
1520 att.produced_by.clone(),
1521 att.cost.clone(),
1522 att.timestamp,
1523 ).attestation_id;
1524 if expected != att.attestation_id {
1525 return error_with_detail(409, "AttestationIdMismatch", serde_json::json!({
1526 "supplied": att.attestation_id,
1527 "expected": expected,
1528 }));
1529 }
1530 if let Some(op_id) = &att.op_id {
1534 match op_log.get(op_id) {
1535 Ok(Some(_)) => {}
1536 Ok(None) => return error_with_detail(422, "UnknownOp", serde_json::json!({
1537 "attestation_id": att.attestation_id,
1538 "op_id": op_id,
1539 })),
1540 Err(e) => return error_response(500, format!("op log read: {e}")),
1541 }
1542 }
1543 }
1544
1545 let mut added = 0usize;
1549 let mut added_ids: Vec<&lex_vcs::AttestationId> = Vec::new();
1550 for att in &attestations {
1551 let already_present = matches!(log.get(&att.attestation_id), Ok(Some(_)));
1552 match log.put(att) {
1553 Ok(()) => {
1554 if !already_present {
1555 added += 1;
1556 added_ids.push(&att.attestation_id);
1557 }
1558 }
1559 Err(e) => return error_response(500, format!("attestation log write: {e}")),
1560 }
1561 }
1562
1563 json_response(200, &serde_json::json!({
1564 "received": attestations.len(),
1565 "added": added,
1566 "skipped": attestations.len() - added,
1567 "added_ids": added_ids,
1568 }))
1569}
1570
1571pub(crate) fn attestations_since_handler(state: &State, query: &str)
1584 -> Response<std::io::Cursor<Vec<u8>>>
1585{
1586 let mut after_op: Option<String> = None;
1587 let mut limit: Option<usize> = None;
1588 for kv in query.split('&') {
1589 let Some((k, v)) = kv.split_once('=') else { continue };
1590 match k {
1591 "after-op" => after_op = Some(v.to_string()),
1592 "limit" => {
1593 limit = Some(match v.parse::<usize>() {
1594 Ok(n) => n,
1595 Err(_) => return error_response(400,
1596 format!("limit must be a positive integer, got `{v}`")),
1597 });
1598 }
1599 _ => {}
1600 }
1601 }
1602
1603 let store = state.store.lock().unwrap();
1604 let log = match store.attestation_log() {
1605 Ok(l) => l,
1606 Err(e) => return error_response(500, format!("opening attestation log: {e}")),
1607 };
1608
1609 let exclude: std::collections::BTreeSet<String> = match &after_op {
1613 None => std::collections::BTreeSet::new(),
1614 Some(cutoff) => {
1615 let op_log = match lex_vcs::OpLog::open(store.root()) {
1616 Ok(l) => l,
1617 Err(e) => return error_response(500, format!("opening op log: {e}")),
1618 };
1619 match op_log.walk_back(cutoff, None) {
1620 Ok(records) => records.into_iter().map(|r| r.op_id).collect(),
1621 Err(_) => {
1622 std::collections::BTreeSet::new()
1626 }
1627 }
1628 }
1629 };
1630
1631 let all = match log.list_all() {
1632 Ok(v) => v,
1633 Err(e) => return error_response(500, format!("listing attestations: {e}")),
1634 };
1635 let mut filtered: Vec<lex_vcs::Attestation> = all
1636 .into_iter()
1637 .filter(|a| match &a.op_id {
1638 Some(op_id) => !exclude.contains(op_id),
1639 None => true,
1643 })
1644 .collect();
1645 filtered.sort_by(|a, b| {
1649 a.timestamp.cmp(&b.timestamp)
1650 .then_with(|| a.attestation_id.cmp(&b.attestation_id))
1651 });
1652 if let Some(n) = limit {
1653 filtered.truncate(n);
1654 }
1655
1656 json_response(200, &serde_json::to_value(&filtered).unwrap_or_default())
1657}
1658
1659#[derive(Debug, Clone, Default, serde::Serialize, serde::Deserialize)]
1668struct DepSpec {
1669 #[serde(default, skip_serializing_if = "Option::is_none")]
1670 registry: Option<String>,
1671 #[serde(default, skip_serializing_if = "Option::is_none")]
1672 version: Option<String>,
1673 #[serde(default, skip_serializing_if = "Option::is_none")]
1674 git: Option<String>,
1675 #[serde(default, skip_serializing_if = "Option::is_none")]
1676 branch: Option<String>,
1677 #[serde(default, skip_serializing_if = "Option::is_none")]
1678 tag: Option<String>,
1679 #[serde(default, skip_serializing_if = "Option::is_none")]
1680 rev: Option<String>,
1681 #[serde(default, skip_serializing_if = "Option::is_none")]
1682 path: Option<String>,
1683}
1684
1685impl DepSpec {
1686 fn to_toml_inline(&self) -> Option<String> {
1691 let mut parts: Vec<String> = Vec::new();
1692 let mut push = |k: &str, v: &Option<String>| {
1693 if let Some(val) = v {
1694 parts.push(format!("{k} = {}", toml_str(val)));
1695 }
1696 };
1697 push("registry", &self.registry);
1698 push("version", &self.version);
1699 push("git", &self.git);
1700 push("branch", &self.branch);
1701 push("tag", &self.tag);
1702 push("rev", &self.rev);
1703 push("path", &self.path);
1704 if parts.is_empty() {
1705 None
1706 } else {
1707 Some(format!("{{ {} }}", parts.join(", ")))
1708 }
1709 }
1710}
1711
1712fn toml_str(s: &str) -> String {
1714 format!("\"{}\"", s.replace('\\', "\\\\").replace('"', "\\\""))
1715}
1716
1717#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
1719struct PkgRecord {
1720 name: String,
1721 version: String,
1722 head_op: Option<String>,
1723 published_at: u64,
1724 function_names: Vec<String>,
1726 #[serde(default)]
1731 dependencies: Vec<String>,
1732 #[serde(default, skip_serializing_if = "std::collections::BTreeMap::is_empty")]
1738 dependency_specs: std::collections::BTreeMap<String, DepSpec>,
1739 ops: Vec<serde_json::Value>,
1741}
1742
1743#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize, serde::Deserialize, Default)]
1750#[serde(rename_all = "lowercase")]
1751pub enum Visibility {
1752 #[default]
1753 Private,
1754 Public,
1755}
1756
1757#[derive(Debug, Clone, serde::Serialize, serde::Deserialize, Default)]
1762struct PkgIndex {
1763 latest: Option<String>,
1765 versions: Vec<PkgVersionSummary>,
1767 #[serde(default)]
1771 visibility: Visibility,
1772}
1773
1774#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
1775struct PkgVersionSummary {
1776 version: String,
1777 head_op: Option<String>,
1778 published_at: u64,
1779}
1780
1781fn pkg_name_dir(root: &std::path::Path, name: &str) -> PathBuf {
1782 root.join("packages").join(name)
1783}
1784
1785fn pkg_index_path(root: &std::path::Path, name: &str) -> PathBuf {
1786 pkg_name_dir(root, name).join("index.json")
1787}
1788
1789fn pkg_version_path(root: &std::path::Path, name: &str, version: &str) -> PathBuf {
1790 pkg_name_dir(root, name).join(format!("{version}.json"))
1791}
1792
1793fn pkg_archive_path(root: &std::path::Path, name: &str, version: &str) -> PathBuf {
1794 pkg_name_dir(root, name).join(format!("{version}.tar.gz"))
1795}
1796
1797fn load_pkg_index(root: &std::path::Path, name: &str) -> Option<PkgIndex> {
1798 let bytes = std::fs::read(pkg_index_path(root, name)).ok()?;
1799 serde_json::from_slice(&bytes).ok()
1800}
1801
1802fn load_pkg_record(root: &std::path::Path, name: &str, version: &str) -> Option<PkgRecord> {
1803 let bytes = std::fs::read(pkg_version_path(root, name, version)).ok()?;
1804 serde_json::from_slice(&bytes).ok()
1805}
1806
1807fn load_latest_pkg_record(root: &std::path::Path, name: &str) -> Option<PkgRecord> {
1808 let index = load_pkg_index(root, name)?;
1809 let latest = index.latest.clone()?;
1810 load_pkg_record(root, name, &latest)
1811}
1812
1813fn pkg_is_public(root: &std::path::Path, name: &str) -> bool {
1817 load_pkg_index(root, name).map(|i| i.visibility) == Some(Visibility::Public)
1818}
1819
1820fn valid_pkg_segment(s: &str) -> bool {
1825 !s.is_empty()
1826 && s.len() <= 128
1827 && s != "."
1828 && s != ".."
1829 && s.chars().all(|c| c.is_ascii_alphanumeric() || matches!(c, '.' | '_' | '-'))
1830}
1831
1832#[derive(Deserialize)]
1833struct VisibilityReq {
1834 visibility: Visibility,
1835}
1836
1837fn pkg_set_visibility_handler(
1844 state: &State,
1845 name: &str,
1846 body: &str,
1847) -> Response<std::io::Cursor<Vec<u8>>> {
1848 if !valid_pkg_segment(name) {
1849 return error_response(400, format!("invalid package name {name:?}"));
1850 }
1851 let req: VisibilityReq = match serde_json::from_str(body) {
1852 Ok(r) => r,
1853 Err(e) => return error_response(400, format!("bad request: {e}")),
1854 };
1855 let mut index = match load_pkg_index(&state.root, name) {
1856 Some(i) => i,
1857 None => return error_response(404, format!("package {name:?} not found")),
1858 };
1859 index.visibility = req.visibility;
1860 let bytes = serde_json::to_vec_pretty(&index).unwrap_or_default();
1861 match std::fs::write(pkg_index_path(&state.root, name), bytes) {
1862 Ok(()) => json_response(
1863 200,
1864 &serde_json::json!({ "name": name, "visibility": index.visibility }),
1865 ),
1866 Err(e) => error_response(500, format!("write index: {e}")),
1867 }
1868}
1869
1870#[derive(serde::Deserialize)]
1871struct ReleaseReq {
1872 version: String,
1873 #[serde(default)]
1874 branch: Option<String>,
1875 #[serde(default)]
1878 dependencies: Vec<String>,
1879 #[serde(default)]
1884 dependency_specs: std::collections::BTreeMap<String, DepSpec>,
1885}
1886
1887fn pkg_release_handler(state: &State, name: &str, body: &str) -> Response<std::io::Cursor<Vec<u8>>> {
1896 if !valid_pkg_segment(name) {
1897 return error_response(400, format!("invalid package name {name:?}"));
1898 }
1899 let req: ReleaseReq = match serde_json::from_str(body) {
1900 Ok(r) => r,
1901 Err(e) => return error_response(400, format!("bad request: {e}")),
1902 };
1903 let version = req.version.trim().to_string();
1904 if version.is_empty() || !valid_pkg_segment(&version) {
1905 return error_response(400, "version must be a non-empty, path-safe string (e.g. 1.2.0)");
1906 }
1907 if load_pkg_record(&state.root, name, &version).is_some() {
1909 return error_response(
1910 409,
1911 format!("{name}@{version} already released; releases are immutable — bump the version"),
1912 );
1913 }
1914
1915 let store = state.store.lock().unwrap();
1916 let branch = req.branch.unwrap_or_else(|| store.current_branch());
1917 let head_op = match store.get_branch(&branch) {
1918 Ok(Some(b)) => b.head_op,
1919 Ok(None) => return error_response(404, format!("unknown branch {branch:?}")),
1920 Err(e) => return error_response(500, format!("get_branch: {e}")),
1921 };
1922 let Some(head_op) = head_op else {
1923 return error_response(400, format!("branch {branch:?} has no commits to release"));
1924 };
1925
1926 let predecessor = load_pkg_index(&state.root, name)
1937 .map(|i| i.versions)
1938 .unwrap_or_default()
1939 .into_iter()
1940 .filter_map(|v| lex_syntax::semver::parse_exact(&v.version).map(|p| (p, v)))
1941 .filter(|(p, _)| lex_syntax::semver::parse_exact(&version).map(|n| *p < n).unwrap_or(false))
1942 .max_by_key(|(p, _)| *p)
1943 .map(|(_, v)| v);
1944 if let Some(prev) = predecessor {
1945 if let (Some(prev_head), Some(declared)) = (
1946 prev.head_op.clone(),
1947 lex_syntax::semver::bump_between(&prev.version, &version),
1948 ) {
1949 if let (Ok(prev_api), Ok(new_api)) = (
1950 lex_store::api::public_api_at_op(&store, &prev_head),
1951 lex_store::api::public_api_at_op(&store, &head_op),
1952 ) {
1953 use lex_store::api::ApiChange;
1954 use lex_syntax::semver::Bump;
1955 let (required, why) = match lex_store::api::classify_api_change(&prev_api, &new_api) {
1956 ApiChange::Breaking(d) => (Bump::Major, d),
1957 ApiChange::Additive(d) => (Bump::Minor, d),
1958 ApiChange::None => (Bump::Patch, String::new()),
1959 };
1960 if declared < required {
1961 let need = match required {
1962 Bump::Major => "major",
1963 Bump::Minor => "minor",
1964 Bump::Patch => "patch",
1965 };
1966 return error_response(
1967 422,
1968 format!(
1969 "version bump too small: {} → {version} is a {declared:?} bump, \
1970 but the API change ({why}) requires a {need} bump",
1971 prev.version
1972 ),
1973 );
1974 }
1975 }
1976 }
1977 }
1978
1979 let head = store.branch_head(&branch).unwrap_or_default();
1982 let pairs: Vec<(String, String)> = head.iter().map(|(s, st)| (s.clone(), st.clone())).collect();
1983 let function_names: Vec<String> = store
1984 .get_asts_for_sigs_bulk(&pairs)
1985 .into_iter()
1986 .filter_map(|r| r.ok())
1987 .filter_map(|s| match s {
1988 lex_ast::Stage::FnDecl(fd) => Some(fd.name),
1989 _ => None,
1990 })
1991 .collect();
1992 let mut deps: std::collections::BTreeSet<String> = req.dependencies.into_iter().collect();
1996 if let Ok(extracted) = lex_store::api::external_dependencies_at_op(&store, &head_op) {
1997 deps.extend(extracted);
1998 }
1999 deps.extend(req.dependency_specs.keys().cloned());
2002 let dependencies: Vec<String> = deps.into_iter().collect();
2003 drop(store);
2004
2005 let published_at = std::time::SystemTime::now()
2006 .duration_since(std::time::UNIX_EPOCH)
2007 .map(|d| d.as_secs())
2008 .unwrap_or(0);
2009 let record = PkgRecord {
2010 name: name.to_string(),
2011 version: version.clone(),
2012 head_op: Some(head_op.clone()),
2013 published_at,
2014 function_names,
2015 dependencies,
2016 dependency_specs: req.dependency_specs,
2017 ops: Vec::new(),
2018 };
2019 if let Err(e) = save_pkg_record(&state.root, &record, None) {
2020 return error_response(500, format!("write release: {e}"));
2021 }
2022 json_response(
2023 201,
2024 &serde_json::json!({
2025 "name": name,
2026 "version": version,
2027 "head_op": head_op,
2028 "branch": branch,
2029 }),
2030 )
2031}
2032
2033fn public_pkg_names(root: &std::path::Path) -> Vec<String> {
2036 list_pkg_names(root)
2037 .into_iter()
2038 .filter(|name| pkg_is_public(root, name))
2039 .collect()
2040}
2041
2042fn public_pkg_list_handler(state: &State) -> Response<std::io::Cursor<Vec<u8>>> {
2046 let packages: Vec<serde_json::Value> = public_pkg_names(&state.root)
2047 .iter()
2048 .filter_map(|name| {
2049 let r = load_latest_pkg_record(&state.root, name)?;
2050 Some(serde_json::json!({
2051 "name": r.name,
2052 "version": r.version,
2053 "head_op": r.head_op,
2054 "published_at": r.published_at,
2055 }))
2056 })
2057 .collect();
2058 json_response(200, &serde_json::json!({ "packages": packages }))
2059}
2060
2061#[derive(Debug, PartialEq, Eq)]
2066enum PublicTarget {
2067 List,
2068 Latest(String),
2069 Versions(String),
2070 ApiDiff(String),
2071 Head(String),
2072 Version(String, String),
2073 Archive(String, String),
2074}
2075
2076impl PublicTarget {
2077 fn pkg_name(&self) -> Option<&str> {
2079 match self {
2080 PublicTarget::List => None,
2081 PublicTarget::Latest(n)
2082 | PublicTarget::Versions(n)
2083 | PublicTarget::ApiDiff(n)
2084 | PublicTarget::Head(n)
2085 | PublicTarget::Version(n, _)
2086 | PublicTarget::Archive(n, _) => Some(n),
2087 }
2088 }
2089}
2090
2091fn resolve_public(method: &Method, path: &str) -> Result<PublicTarget, u16> {
2096 if !matches!(method, Method::Get) {
2097 return Err(405);
2098 }
2099 let rest = path.trim_matches('/');
2100 if rest.is_empty() {
2101 return Ok(PublicTarget::List);
2102 }
2103 let segs: Vec<&str> = rest.split('/').collect();
2104 if !segs.iter().all(|s| valid_pkg_segment(s)) {
2105 return Err(404);
2106 }
2107 match segs.as_slice() {
2108 [n] => Ok(PublicTarget::Latest(n.to_string())),
2109 [n, "versions"] => Ok(PublicTarget::Versions(n.to_string())),
2110 [n, "api-diff"] => Ok(PublicTarget::ApiDiff(n.to_string())),
2111 [n, "head"] => Ok(PublicTarget::Head(n.to_string())),
2112 [n, v, "archive"] => Ok(PublicTarget::Archive(n.to_string(), v.to_string())),
2113 [n, v] => Ok(PublicTarget::Version(n.to_string(), v.to_string())),
2114 _ => Err(404),
2115 }
2116}
2117
2118pub fn route_public(
2129 state: &State,
2130 method: &Method,
2131 path: &str,
2132 query: &str,
2133) -> Response<std::io::Cursor<Vec<u8>>> {
2134 let target = match resolve_public(method, path) {
2135 Ok(t) => t,
2136 Err(405) => return error_response(405, "public read is GET-only"),
2137 Err(_) => return error_response(404, "not found"),
2138 };
2139 if let PublicTarget::List = target {
2141 return public_pkg_list_handler(state);
2142 }
2143 if let Some(name) = target.pkg_name() {
2145 if !pkg_is_public(&state.root, name) {
2146 return error_response(404, format!("package {name:?} not found"));
2147 }
2148 }
2149 match target {
2150 PublicTarget::List => unreachable!("handled above"),
2151 PublicTarget::Latest(n) => pkg_get_handler(state, &n),
2152 PublicTarget::Versions(n) => pkg_versions_handler(state, &n),
2153 PublicTarget::ApiDiff(n) => pkg_api_diff_handler(state, &n, query),
2154 PublicTarget::Head(n) => pkg_head_handler(state, &n),
2155 PublicTarget::Version(n, v) => pkg_get_version_handler(state, &n, &v),
2156 PublicTarget::Archive(n, v) => pkg_archive_handler(state, &n, &v),
2157 }
2158}
2159
2160fn save_pkg_record(
2161 root: &std::path::Path,
2162 record: &PkgRecord,
2163 archive: Option<&[u8]>,
2166) -> std::io::Result<()> {
2167 let dir = pkg_name_dir(root, &record.name);
2168 std::fs::create_dir_all(&dir)?;
2169
2170 let rec_bytes = serde_json::to_vec_pretty(record).unwrap_or_default();
2172 std::fs::write(pkg_version_path(root, &record.name, &record.version), rec_bytes)?;
2173
2174 if let Some(archive) = archive {
2176 std::fs::write(pkg_archive_path(root, &record.name, &record.version), archive)?;
2177 }
2178
2179 let mut index = load_pkg_index(root, &record.name).unwrap_or_default();
2181 index.latest = Some(record.version.clone());
2182 if !index.versions.iter().any(|v| v.version == record.version) {
2183 index.versions.push(PkgVersionSummary {
2184 version: record.version.clone(),
2185 head_op: record.head_op.clone(),
2186 published_at: record.published_at,
2187 });
2188 }
2189 let idx_bytes = serde_json::to_vec_pretty(&index).unwrap_or_default();
2190 std::fs::write(pkg_index_path(root, &record.name), idx_bytes)
2191}
2192
2193fn list_pkg_names(root: &std::path::Path) -> Vec<String> {
2194 let dir = root.join("packages");
2195 let Ok(entries) = std::fs::read_dir(&dir) else {
2196 return Vec::new();
2197 };
2198 let mut names: Vec<String> = entries
2199 .filter_map(|e| e.ok())
2200 .filter(|e| e.path().is_dir())
2201 .filter_map(|e| e.file_name().into_string().ok())
2202 .collect();
2203 names.sort();
2204 names
2205}
2206
2207fn collect_lex_files(dir: &std::path::Path, out: &mut Vec<PathBuf>) {
2208 let Ok(entries) = std::fs::read_dir(dir) else { return };
2209 let mut entries: Vec<_> = entries.filter_map(|e| e.ok()).collect();
2210 entries.sort_by_key(|e| e.path());
2211 for entry in entries {
2212 let path = entry.path();
2213 if path.is_dir() {
2214 collect_lex_files(&path, out);
2215 } else if path.extension().and_then(|x| x.to_str()) == Some("lex") {
2216 out.push(path);
2217 }
2218 }
2219}
2220
2221fn pkg_publish_handler(state: &State, body: &[u8]) -> Response<std::io::Cursor<Vec<u8>>> {
2224 let tmp = match tempfile::TempDir::new() {
2225 Ok(t) => t,
2226 Err(e) => return error_response(500, format!("create temp dir: {e}")),
2227 };
2228 {
2229 let gz = flate2::read::GzDecoder::new(std::io::Cursor::new(body));
2230 let mut ar = tar::Archive::new(gz);
2231 if let Err(e) = ar.unpack(tmp.path()) {
2232 return error_response(400, format!("unpack archive: {e}"));
2233 }
2234 }
2235
2236 let toml_path = tmp.path().join("lex.toml");
2237 if !toml_path.exists() {
2238 return error_response(400, "archive must contain lex.toml at root");
2239 }
2240 let manifest = match Manifest::load(&toml_path) {
2241 Ok(m) => m,
2242 Err(e) => return error_response(400, format!("lex.toml: {e}")),
2243 };
2244 let (pkg_name, pkg_version) = match &manifest.package {
2245 Some(m) => (m.name.clone(), m.version.clone()),
2246 None => return error_response(400, "lex.toml must have a [package] section"),
2247 };
2248
2249 if load_pkg_record(&state.root, &pkg_name, &pkg_version).is_some() {
2253 return error_response(
2254 409,
2255 format!(
2256 "package {pkg_name}@{pkg_version} already published; \
2257 bump the version in lex.toml to publish a new release"
2258 ),
2259 );
2260 }
2261
2262 let src_dir = tmp.path().join("src");
2263 if !src_dir.exists() {
2264 return error_response(400, "archive must contain a src/ directory");
2265 }
2266 let mut lex_files: Vec<PathBuf> = Vec::new();
2267 collect_lex_files(&src_dir, &mut lex_files);
2268 if lex_files.is_empty() {
2269 return error_response(400, "no .lex files found in src/");
2270 }
2271
2272 let store = state.store.lock().unwrap();
2273 let branch = store.current_branch();
2274
2275 let old_head = match store.branch_head(&branch) {
2295 Ok(h) => h,
2296 Err(e) => return error_response(500, format!("branch_head: {e}")),
2297 };
2298 let old_pairs: Vec<(String, String)> =
2305 old_head.iter().map(|(sig, stage)| (sig.clone(), stage.clone())).collect();
2306 let mut old_fns_by_name: BTreeMap<String, Vec<lex_ast::FnDecl>> = BTreeMap::new();
2307 for fd in store.get_asts_for_sigs_bulk(&old_pairs)
2308 .into_iter()
2309 .filter_map(|r| r.ok())
2310 .filter_map(|s| match s { lex_ast::Stage::FnDecl(fd) => Some(fd), _ => None })
2311 {
2312 old_fns_by_name.entry(fd.name.clone()).or_default().push(fd);
2313 }
2314 let mut old_types_by_name: BTreeMap<String, lex_ast::TypeDecl> = BTreeMap::new();
2317 for td in store.get_asts_for_sigs_bulk(&old_pairs)
2318 .into_iter()
2319 .filter_map(|r| r.ok())
2320 .filter_map(|s| match s { lex_ast::Stage::TypeDecl(td) => Some(td), _ => None })
2321 {
2322 old_types_by_name.insert(td.name.clone(), td);
2323 }
2324 fn structural_key(fd: &lex_ast::FnDecl) -> Option<String> {
2332 let mut anon = fd.clone();
2333 anon.name = String::new();
2334 lex_ast::sig_id(&lex_ast::Stage::FnDecl(anon))
2335 }
2336
2337 fn take_matching(
2347 map: &mut BTreeMap<String, Vec<lex_ast::FnDecl>>,
2348 name: &str,
2349 new_fd: &lex_ast::FnDecl,
2350 ) -> Option<lex_ast::FnDecl> {
2351 let candidates = map.get_mut(name)?;
2352 let idx = match candidates.len() {
2353 0 => return None,
2354 1 => 0,
2355 _ => {
2356 let want = structural_key(new_fd);
2357 candidates.iter().position(|c| structural_key(c) == want)?
2358 }
2359 };
2360 let matched = candidates.remove(idx);
2361 if candidates.is_empty() {
2362 map.remove(name);
2363 }
2364 Some(matched)
2365 }
2366
2367 let loaded = match load_package(&lex_files, tmp.path(), &pkg_name, true) {
2391 Ok(p) => p,
2392 Err(e) => return error_response(400, format!("load package: {e}")),
2393 };
2394 let mut stages = canonicalize_program(&loaded.program);
2395 if let Err(errs) = lex_types::check_and_rewrite_program(&mut stages) {
2401 return error_with_detail(
2402 422,
2403 format!("type errors in package {pkg_name}"),
2404 serde_json::to_value(&errs).unwrap(),
2405 );
2406 }
2407 let new_fns = stage_fns(&stages);
2408 let all_function_names: Vec<String> = new_fns.keys().cloned().collect();
2409
2410 let mut old_fns: BTreeMap<String, lex_ast::FnDecl> = BTreeMap::new();
2415 for (name, new_fd) in &new_fns {
2416 if let Some(fd) = take_matching(&mut old_fns_by_name, name, new_fd) {
2417 old_fns.insert(name.clone(), fd);
2418 }
2419 }
2420 let new_types = stage_types(&stages);
2421 let old_types: BTreeMap<String, lex_ast::TypeDecl> = new_types
2425 .keys()
2426 .filter_map(|n| old_types_by_name.get(n).map(|td| (n.clone(), td.clone())))
2427 .collect();
2428 let report =
2429 lex_vcs::compute_diff_with_types(&old_fns, &new_fns, &old_types, &new_types, false);
2430
2431 let mut new_imports = lex_vcs::ImportMap::new();
2439 for (file, modules) in &loaded.imports_by_file {
2440 let entry = new_imports.entry(file.clone()).or_default();
2441 for (reference, alias) in modules {
2442 entry.insert(lex_vcs::ImportRef {
2443 reference: reference.clone(),
2444 alias: alias.clone(),
2445 });
2446 }
2447 }
2448
2449 let outcome = match store.publish_program_with_intent(
2452 &branch,
2453 &stages,
2454 &report,
2455 &new_imports,
2456 false,
2457 None,
2458 None,
2459 &loaded.module_prefixes,
2460 ) {
2461 Ok(outcome) => outcome,
2462 Err(lex_store::StoreError::TypeError(errs)) => {
2463 return error_with_detail(422, "type errors", serde_json::to_value(&errs).unwrap());
2464 }
2465 Err(e) => return write_error_response("publish_program", e),
2466 };
2467 let all_ops: Vec<serde_json::Value> = match serde_json::to_value(&outcome.ops) {
2468 Ok(serde_json::Value::Array(arr)) => arr,
2469 _ => Vec::new(),
2470 };
2471 let final_head_op = outcome.head_op;
2472
2473 let now = SystemTime::now()
2498 .duration_since(UNIX_EPOCH)
2499 .map(|d| d.as_secs())
2500 .unwrap_or(0);
2501 let dependencies: Vec<String> = final_head_op
2504 .as_ref()
2505 .and_then(|h| lex_store::api::external_dependencies_at_op(&store, h).ok())
2506 .unwrap_or_default();
2507 let record = PkgRecord {
2508 name: pkg_name.clone(),
2509 version: pkg_version,
2510 head_op: final_head_op.clone(),
2511 published_at: now,
2512 function_names: all_function_names,
2513 dependencies,
2514 dependency_specs: Default::default(),
2518 ops: all_ops.clone(),
2519 };
2520 if let Err(e) = save_pkg_record(&state.root, &record, Some(body)) {
2521 return error_response(500, format!("save package index: {e}"));
2522 }
2523
2524 json_response(200, &serde_json::json!({
2525 "package": pkg_name,
2526 "ops": all_ops,
2527 "head_op": final_head_op,
2528 }))
2529}
2530
2531fn pkg_list_handler(state: &State) -> Response<std::io::Cursor<Vec<u8>>> {
2533 let names = list_pkg_names(&state.root);
2534 let packages: Vec<serde_json::Value> = names.iter()
2535 .filter_map(|name| {
2536 let idx = load_pkg_index(&state.root, name)?;
2537 let latest = idx.latest.as_deref()?;
2538 let r = load_pkg_record(&state.root, name, latest)?;
2539 Some(serde_json::json!({
2540 "name": r.name,
2541 "version": r.version,
2542 "head_op": r.head_op,
2543 "published_at": r.published_at,
2544 }))
2545 })
2546 .collect();
2547 json_response(200, &serde_json::json!({ "packages": packages }))
2548}
2549
2550fn pkg_get_handler(state: &State, name: &str) -> Response<std::io::Cursor<Vec<u8>>> {
2552 match load_latest_pkg_record(&state.root, name) {
2553 Some(r) => json_response(200, &serde_json::json!({
2554 "name": r.name,
2555 "version": r.version,
2556 "head_op": r.head_op,
2557 "published_at": r.published_at,
2558 "function_names": r.function_names,
2559 "ops": r.ops,
2560 })),
2561 None => error_response(404, format!("package {name:?} not found")),
2562 }
2563}
2564
2565fn pkg_versions_handler(state: &State, name: &str) -> Response<std::io::Cursor<Vec<u8>>> {
2567 match load_pkg_index(&state.root, name) {
2568 Some(idx) => json_response(200, &serde_json::json!({
2569 "name": name,
2570 "latest": idx.latest,
2571 "versions": idx.versions,
2572 })),
2573 None => error_response(404, format!("package {name:?} not found")),
2574 }
2575}
2576
2577fn pkg_api_diff_handler(state: &State, name: &str, query: &str) -> Response<std::io::Cursor<Vec<u8>>> {
2582 let mut from: Option<String> = None;
2583 let mut to: Option<String> = None;
2584 for kv in query.split('&') {
2585 match kv.split_once('=') {
2586 Some(("from", v)) => from = Some(v.to_string()),
2587 Some(("to", v)) => to = Some(v.to_string()),
2588 _ => {}
2589 }
2590 }
2591 let (Some(from), Some(to)) = (from, to) else {
2592 return error_response(400, "api-diff requires ?from=<version>&to=<version>");
2593 };
2594 let head_of = |v: &str| load_pkg_record(&state.root, name, v).and_then(|r| r.head_op);
2595 let (Some(from_head), Some(to_head)) = (head_of(&from), head_of(&to)) else {
2596 return error_response(404, format!("{name}: unknown release in {from}..{to}"));
2597 };
2598
2599 let store = state.store.lock().unwrap();
2600 let (prev_api, new_api) = match (
2601 lex_store::api::public_api_at_op(&store, &from_head),
2602 lex_store::api::public_api_at_op(&store, &to_head),
2603 ) {
2604 (Ok(a), Ok(b)) => (a, b),
2605 _ => return error_response(500, "could not read package APIs for the given releases"),
2606 };
2607 let (change, detail) = match lex_store::api::classify_api_change(&prev_api, &new_api) {
2608 lex_store::api::ApiChange::Breaking(d) => ("breaking", d),
2609 lex_store::api::ApiChange::Additive(d) => ("additive", d),
2610 lex_store::api::ApiChange::None => ("none", String::new()),
2611 };
2612 let renames = lex_store::api::detect_renames(&prev_api, &new_api);
2613 json_response(200, &serde_json::json!({
2614 "name": name, "from": from, "to": to,
2615 "change": change, "detail": detail, "renames": renames,
2616 }))
2617}
2618
2619fn pkg_get_version_handler(state: &State, name: &str, version: &str) -> Response<std::io::Cursor<Vec<u8>>> {
2621 match load_pkg_record(&state.root, name, version) {
2622 Some(r) => json_response(200, &serde_json::json!({
2623 "name": r.name,
2624 "version": r.version,
2625 "head_op": r.head_op,
2626 "published_at": r.published_at,
2627 "function_names": r.function_names,
2628 "dependencies": r.dependencies,
2629 "ops": r.ops,
2630 })),
2631 None => error_response(404, format!("package {name:?}@{version:?} not found")),
2632 }
2633}
2634
2635fn pkg_archive_handler(state: &State, name: &str, version: &str) -> Response<std::io::Cursor<Vec<u8>>> {
2637 let gzip = |bytes: Vec<u8>| {
2638 Response::from_data(bytes).with_status_code(200).with_header(
2639 tiny_http::Header::from_bytes(&b"Content-Type"[..], &b"application/gzip"[..]).unwrap(),
2640 )
2641 };
2642
2643 if let Ok(bytes) = std::fs::read(pkg_archive_path(&state.root, name, version)) {
2645 return gzip(bytes);
2646 }
2647
2648 if let Some(record) = load_pkg_record(&state.root, name, version) {
2652 if let Some(head_op) = record.head_op.clone() {
2653 match render_op_log_archive(state, name, version, &head_op, &record.dependency_specs) {
2654 Ok(bytes) => return gzip(bytes),
2655 Err(e) => {
2656 return error_response(500, format!("rendering archive for {name:?}@{version:?}: {e}"));
2657 }
2658 }
2659 }
2660 }
2661
2662 error_response(404, format!("archive for {name:?}@{version:?} not found"))
2663}
2664
2665const DEFAULT_ARCHIVE_MODULE: &str = "src/lib.lex";
2670
2671fn render_op_log_archive(
2676 state: &State,
2677 name: &str,
2678 version: &str,
2679 head_op: &str,
2680 dependency_specs: &std::collections::BTreeMap<String, DepSpec>,
2681) -> Result<Vec<u8>, String> {
2682 let (files, lock): (Vec<(String, String)>, Option<String>) = {
2688 let store = state.store.lock().unwrap();
2689 let lock = store.committed_lock_inherited(head_op).ok().flatten();
2695 let head = lex_store::render::package_head_at_op(&store, head_op)
2696 .map_err(|e| format!("reading head {head_op}: {e}"))?;
2697 let files = match lex_store::render::render_source(&store, &head)
2698 .map_err(|e| format!("rendering source at {head_op}: {e}"))?
2699 {
2700 lex_store::render::RenderedSource::Single { path, src } => {
2701 vec![(path.unwrap_or_else(|| DEFAULT_ARCHIVE_MODULE.to_string()), src)]
2704 }
2705 lex_store::render::RenderedSource::Multi(tree) => tree.into_iter().collect(),
2706 };
2707 (files, lock)
2708 };
2709
2710 let mut manifest = format!("[package]\nname = \"{name}\"\nversion = \"{version}\"\n");
2715 let dep_lines: Vec<String> = dependency_specs
2716 .iter()
2717 .filter_map(|(dep_name, spec)| spec.to_toml_inline().map(|inline| format!("{dep_name} = {inline}")))
2718 .collect();
2719 if !dep_lines.is_empty() {
2720 manifest.push_str("\n[dependencies]\n");
2721 for line in dep_lines {
2722 manifest.push_str(&line);
2723 manifest.push('\n');
2724 }
2725 }
2726 let mut enc = flate2::write::GzEncoder::new(Vec::new(), flate2::Compression::default());
2727 {
2728 let mut ar = tar::Builder::new(&mut enc);
2729 let mut append = |p: &str, data: &[u8]| -> std::io::Result<()> {
2730 let mut h = tar::Header::new_gnu();
2731 h.set_size(data.len() as u64);
2732 h.set_mode(0o644);
2733 h.set_cksum();
2734 ar.append_data(&mut h, p, data)
2735 };
2736 append("lex.toml", manifest.as_bytes()).map_err(|e| e.to_string())?;
2737 if let Some(lock) = &lock {
2738 append("lex.lock", lock.as_bytes()).map_err(|e| e.to_string())?;
2739 }
2740 for (path, src) in &files {
2741 append(path, src.as_bytes()).map_err(|e| e.to_string())?;
2742 }
2743 ar.finish().map_err(|e| e.to_string())?;
2744 }
2745 enc.finish().map_err(|e| e.to_string())
2746}
2747
2748fn pkg_head_handler(state: &State, name: &str) -> Response<std::io::Cursor<Vec<u8>>> {
2750 match load_latest_pkg_record(&state.root, name) {
2751 Some(r) => json_response(200, &serde_json::json!({
2752 "name": r.name,
2753 "version": r.version,
2754 "head_op": r.head_op,
2755 })),
2756 None => error_response(404, format!("package {name:?} not found")),
2757 }
2758}
2759
2760fn pkg_delete_handler(state: &State, name: &str) -> Response<std::io::Cursor<Vec<u8>>> {
2762 let record = match load_latest_pkg_record(&state.root, name) {
2763 Some(r) => r,
2764 None => return error_response(404, format!("package {name:?} not found")),
2765 };
2766
2767 let store = state.store.lock().unwrap();
2768 let branch = store.current_branch();
2769
2770 let head = match store.branch_head(&branch) {
2771 Ok(h) => h,
2772 Err(e) => return error_response(500, format!("branch_head: {e}")),
2773 };
2774
2775 let head_pairs: Vec<(String, String)> = head
2782 .iter()
2783 .map(|(sig, stage)| (sig.clone(), stage.clone()))
2784 .collect();
2785 let old_fns: BTreeMap<String, lex_ast::FnDecl> = store
2786 .get_asts_for_sigs_bulk(&head_pairs)
2787 .into_iter()
2788 .filter_map(|r| r.ok())
2789 .filter_map(|s| match s {
2790 lex_ast::Stage::FnDecl(fd)
2791 if record.function_names.contains(&fd.name) => Some((fd.name.clone(), fd)),
2792 _ => None,
2793 })
2794 .collect();
2795
2796 let new_fns: BTreeMap<String, lex_ast::FnDecl> = BTreeMap::new();
2797 let report = lex_vcs::compute_diff(&old_fns, &new_fns, false);
2798 let empty_imports = lex_vcs::ImportMap::new();
2799
2800 match store.publish_program(&branch, &[], &report, &empty_imports, false) {
2801 Ok(outcome) => {
2802 let ver = record.version.clone();
2804 let _ = std::fs::remove_file(pkg_version_path(&state.root, name, &ver));
2805 let _ = std::fs::remove_file(pkg_archive_path(&state.root, name, &ver));
2806 if let Some(mut idx) = load_pkg_index(&state.root, name) {
2808 idx.versions.retain(|v| v.version != ver);
2809 idx.latest = idx.versions.last().map(|v| v.version.clone());
2810 if idx.versions.is_empty() {
2811 let _ = std::fs::remove_dir_all(pkg_name_dir(&state.root, name));
2812 } else {
2813 let bytes = serde_json::to_vec_pretty(&idx).unwrap_or_default();
2814 let _ = std::fs::write(pkg_index_path(&state.root, name), bytes);
2815 }
2816 }
2817 json_response(200, &serde_json::json!({
2818 "deleted": name,
2819 "version": ver,
2820 "ops": outcome.ops,
2821 "head_op": outcome.head_op,
2822 }))
2823 }
2824 Err(lex_store::StoreError::TypeError(errs)) => {
2825 error_with_detail(422, "type errors", serde_json::to_value(&errs).unwrap())
2826 }
2827 Err(e) => write_error_response("retract package", e),
2828 }
2829}
2830
2831#[cfg(test)]
2832mod dep_spec_tests {
2833 use super::DepSpec;
2834
2835 #[test]
2836 fn a_dual_spec_renders_both_git_and_vcs_refs_vcs_first() {
2837 let spec = DepSpec {
2838 registry: Some("vcs.lexlang.org/lex-official/lex-schema".into()),
2839 version: Some("^0.9".into()),
2840 git: Some("https://github.com/alpibrusl/lex-schema".into()),
2841 ..Default::default()
2842 };
2843 assert_eq!(
2844 spec.to_toml_inline().as_deref(),
2845 Some("{ registry = \"vcs.lexlang.org/lex-official/lex-schema\", version = \"^0.9\", git = \"https://github.com/alpibrusl/lex-schema\" }"),
2846 );
2847 }
2848
2849 #[test]
2850 fn bare_git_and_bare_registry_specs_render_their_own_keys() {
2851 let git = DepSpec { git: Some("https://x/g".into()), tag: Some("v1".into()), ..Default::default() };
2852 assert_eq!(git.to_toml_inline().as_deref(), Some("{ git = \"https://x/g\", tag = \"v1\" }"));
2853 let reg = DepSpec { registry: Some("vcs/r".into()), version: Some("1.0.0".into()), ..Default::default() };
2854 assert_eq!(reg.to_toml_inline().as_deref(), Some("{ registry = \"vcs/r\", version = \"1.0.0\" }"));
2855 }
2856
2857 #[test]
2858 fn an_empty_spec_renders_nothing() {
2859 assert_eq!(DepSpec::default().to_toml_inline(), None);
2860 }
2861}
2862
2863#[cfg(test)]
2864mod policy_ceiling_tests {
2865 use super::*;
2866 use lex_runtime::Policy;
2867 use std::path::PathBuf;
2868
2869 fn permissive_request() -> Policy {
2873 Policy {
2874 allow_effects: ["io", "fs_read", "fs_write", "net", "proc"]
2875 .iter()
2876 .map(|s| s.to_string())
2877 .collect(),
2878 allow_fs_read: vec![PathBuf::from("/")],
2879 allow_fs_write: vec![PathBuf::from("/")],
2880 allow_net_host: Vec::new(),
2881 allow_proc: Vec::new(),
2882 allow_approval: Vec::new(),
2883 budget: None,
2884 }
2885 }
2886
2887 #[test]
2888 fn ceiling_drops_effects_the_caller_was_not_granted() {
2889 let ceiling = Policy {
2890 allow_effects: ["io", "time"].iter().map(|s| s.to_string()).collect(),
2891 ..Policy::default()
2892 };
2893 let got = clamp_policy(permissive_request(), &ceiling);
2894 assert!(got.allow_effects.contains("io"));
2895 assert!(!got.allow_effects.contains("proc"), "proc must not survive a ceiling without it");
2896 assert!(!got.allow_effects.contains("fs_write"));
2897 assert!(!got.allow_effects.contains("net"));
2898 assert!(!got.allow_effects.contains("time"));
2900 }
2901
2902 #[test]
2903 fn ceiling_scopes_override_caller_scopes() {
2904 let ceiling = Policy {
2905 allow_effects: ["fs_read"].iter().map(|s| s.to_string()).collect(),
2906 allow_fs_read: vec![PathBuf::from("/srv/tenant")],
2907 ..Policy::default()
2908 };
2909 let got = clamp_policy(permissive_request(), &ceiling);
2910 assert_eq!(got.allow_fs_read, vec![PathBuf::from("/srv/tenant")]);
2913 assert!(got.allow_fs_write.is_empty());
2914 assert!(got.allow_proc.is_empty());
2915 assert!(got.allow_net_host.is_empty());
2916 }
2917
2918 #[test]
2919 fn ceiling_caps_budget_and_prefers_the_smaller() {
2920 let mut req = permissive_request();
2922 req.budget = None;
2923 let ceiling = Policy { budget: Some(1_000), ..Policy::default() };
2924 assert_eq!(clamp_policy(req, &ceiling).budget, Some(1_000));
2925
2926 let mut req2 = permissive_request();
2928 req2.budget = Some(50);
2929 let ceiling2 = Policy { budget: Some(1_000), ..Policy::default() };
2930 assert_eq!(clamp_policy(req2, &ceiling2).budget, Some(50));
2931 }
2932
2933 #[test]
2934 fn empty_ceiling_is_pure_only() {
2935 let got = clamp_policy(permissive_request(), &Policy::default());
2936 assert!(got.allow_effects.is_empty(), "an empty ceiling grants nothing");
2937 assert!(got.allow_proc.is_empty());
2938 assert!(got.allow_fs_write.is_empty());
2939 }
2940}
2941
2942#[cfg(test)]
2943mod public_read_tests {
2944 use super::*;
2945
2946 fn seed_pkg(root: &std::path::Path, name: &str, version: &str) {
2949 let record = PkgRecord {
2950 name: name.to_string(),
2951 version: version.to_string(),
2952 head_op: Some(format!("op-{name}")),
2953 published_at: 1,
2954 function_names: vec![format!("{name}.f")],
2955 dependencies: vec![],
2956 dependency_specs: Default::default(),
2957 ops: vec![],
2958 };
2959 save_pkg_record(root, &record, Some(format!("ARCHIVE:{name}@{version}").as_bytes()))
2960 .expect("seed package");
2961 }
2962
2963 #[test]
2964 fn new_package_defaults_to_private() {
2965 let tmp = tempfile::TempDir::new().unwrap();
2966 seed_pkg(tmp.path(), "lex-schema", "0.9.2");
2967 assert!(!pkg_is_public(tmp.path(), "lex-schema"));
2968 assert!(!pkg_is_public(tmp.path(), "does-not-exist"));
2970 }
2971
2972 #[test]
2973 fn set_visibility_round_trips_and_index_persists() {
2974 let tmp = tempfile::TempDir::new().unwrap();
2975 let state = State::open(tmp.path().to_path_buf()).unwrap();
2976 seed_pkg(tmp.path(), "lex-schema", "0.9.2");
2977
2978 let _ = pkg_set_visibility_handler(&state, "lex-schema", r#"{"visibility":"public"}"#);
2979 assert!(pkg_is_public(tmp.path(), "lex-schema"));
2980 let idx = load_pkg_index(tmp.path(), "lex-schema").unwrap();
2982 assert_eq!(idx.latest.as_deref(), Some("0.9.2"));
2983 assert_eq!(idx.versions.len(), 1);
2984
2985 let _ = pkg_set_visibility_handler(&state, "lex-schema", r#"{"visibility":"private"}"#);
2986 assert!(!pkg_is_public(tmp.path(), "lex-schema"));
2987 }
2988
2989 #[test]
2990 fn set_visibility_on_unknown_package_is_a_noop() {
2991 let tmp = tempfile::TempDir::new().unwrap();
2992 let state = State::open(tmp.path().to_path_buf()).unwrap();
2993 let _ = pkg_set_visibility_handler(&state, "ghost", r#"{"visibility":"public"}"#);
2995 assert!(load_pkg_index(tmp.path(), "ghost").is_none());
2996 }
2997
2998 #[test]
2999 fn public_listing_omits_private_packages() {
3000 let tmp = tempfile::TempDir::new().unwrap();
3001 let state = State::open(tmp.path().to_path_buf()).unwrap();
3002 seed_pkg(tmp.path(), "pub-pkg", "1.0.0");
3003 seed_pkg(tmp.path(), "priv-pkg", "1.0.0");
3004 let _ = pkg_set_visibility_handler(&state, "pub-pkg", r#"{"visibility":"public"}"#);
3005
3006 let names = public_pkg_names(tmp.path());
3007 assert_eq!(names, vec!["pub-pkg".to_string()]);
3008 }
3009
3010 #[test]
3011 fn resolve_public_maps_routes() {
3012 let get = Method::Get;
3013 assert_eq!(resolve_public(&get, "").unwrap(), PublicTarget::List);
3014 assert_eq!(resolve_public(&get, "/").unwrap(), PublicTarget::List);
3015 assert_eq!(
3016 resolve_public(&get, "/lex-schema").unwrap(),
3017 PublicTarget::Latest("lex-schema".into())
3018 );
3019 assert_eq!(
3020 resolve_public(&get, "/lex-schema/versions").unwrap(),
3021 PublicTarget::Versions("lex-schema".into())
3022 );
3023 assert_eq!(
3024 resolve_public(&get, "/lex-schema/head").unwrap(),
3025 PublicTarget::Head("lex-schema".into())
3026 );
3027 assert_eq!(
3028 resolve_public(&get, "/lex-schema/0.9.2").unwrap(),
3029 PublicTarget::Version("lex-schema".into(), "0.9.2".into())
3030 );
3031 assert_eq!(
3032 resolve_public(&get, "/lex-schema/0.9.2/archive").unwrap(),
3033 PublicTarget::Archive("lex-schema".into(), "0.9.2".into())
3034 );
3035 }
3036
3037 #[test]
3038 fn resolve_public_rejects_bad_method_and_traversal() {
3039 assert_eq!(resolve_public(&Method::Put, "/lex-schema"), Err(405));
3041 assert_eq!(resolve_public(&Method::Post, "").err(), Some(405));
3042 assert_eq!(resolve_public(&Method::Get, "/.."), Err(404));
3044 assert_eq!(resolve_public(&Method::Get, "/lex-schema/../etc"), Err(404));
3045 assert_eq!(resolve_public(&Method::Get, "/a/b/c/d"), Err(404));
3046 assert!(resolve_public(&Method::Get, "/lex schema").is_err());
3048 }
3049}