use indexmap::IndexMap;
use lex_ast::canonicalize_program;
use lex_bytecode::{compile_program, vm::Vm, Value};
use lex_runtime::{check_program as check_policy, DefaultHandler, Policy};
use lex_store::Store;
use crate::publish_examples::record_examples_for_publish;
use lex_syntax::{load_package, load_program_from_str, Manifest};
use lex_vcs::{MergeSession, MergeSessionId};
use serde::{Deserialize, Serialize};
use std::collections::{BTreeMap, BTreeSet, HashMap};
use std::path::PathBuf;
use std::sync::{Arc, Mutex};
use std::time::{SystemTime, UNIX_EPOCH};
use tiny_http::{Header, Method, Request, Response};
pub struct State {
pub store: Mutex<Store>,
pub root: PathBuf,
pub sessions: Mutex<HashMap<MergeSessionId, ApiMergeSession>>,
pub policy_ceiling: Option<Policy>,
}
pub struct ApiMergeSession {
pub inner: MergeSession,
pub src_branch: String,
pub dst_branch: String,
}
impl State {
pub fn open(root: PathBuf) -> anyhow::Result<Self> {
Self::open_with_ceiling(root, None)
}
pub fn open_with_ceiling(
root: PathBuf,
policy_ceiling: Option<Policy>,
) -> anyhow::Result<Self> {
Ok(Self {
store: Mutex::new(Store::open(&root)?),
root,
sessions: Mutex::new(HashMap::new()),
policy_ceiling,
})
}
pub fn new_with_tenant(tenant_id: &str, store_root: PathBuf) -> anyhow::Result<Self> {
validate_tenant_id(tenant_id)?;
Self::open(store_root.join(tenant_id))
}
pub fn new_with_tenant_and_ceiling(
tenant_id: &str,
store_root: PathBuf,
policy_ceiling: Option<Policy>,
) -> anyhow::Result<Self> {
validate_tenant_id(tenant_id)?;
Self::open_with_ceiling(store_root.join(tenant_id), policy_ceiling)
}
}
fn clamp_policy(requested: Policy, ceiling: &Policy) -> Policy {
let allow_effects: BTreeSet<String> = requested
.allow_effects
.intersection(&ceiling.allow_effects)
.cloned()
.collect();
let budget = match (requested.budget, ceiling.budget) {
(Some(r), Some(c)) => Some(r.min(c)),
(None, Some(c)) => Some(c),
(Some(r), None) => Some(r),
(None, None) => None,
};
Policy {
allow_effects,
allow_fs_read: ceiling.allow_fs_read.clone(),
allow_fs_write: ceiling.allow_fs_write.clone(),
allow_net_host: ceiling.allow_net_host.clone(),
allow_proc: ceiling.allow_proc.clone(),
allow_approval: ceiling.allow_approval.clone(),
budget,
}
}
fn validate_tenant_id(tenant_id: &str) -> anyhow::Result<()> {
if tenant_id.is_empty() {
anyhow::bail!("tenant_id must not be empty");
}
if tenant_id.len() > 64 {
anyhow::bail!("tenant_id must be at most 64 bytes");
}
if !tenant_id
.bytes()
.all(|b| b.is_ascii_alphanumeric() || b == b'_' || b == b'-')
{
anyhow::bail!(
"tenant_id {tenant_id:?} contains characters outside [A-Za-z0-9_-]"
);
}
Ok(())
}
#[derive(Debug, Serialize, Deserialize)]
struct ErrorEnvelope {
error: String,
#[serde(skip_serializing_if = "Option::is_none")]
detail: Option<serde_json::Value>,
}
fn json_response(status: u16, body: &serde_json::Value) -> Response<std::io::Cursor<Vec<u8>>> {
let bytes = serde_json::to_vec(body).unwrap_or_else(|_| b"{}".to_vec());
Response::from_data(bytes)
.with_status_code(status)
.with_header(Header::from_bytes(&b"Content-Type"[..], &b"application/json"[..]).unwrap())
}
fn error_response(status: u16, msg: impl Into<String>) -> Response<std::io::Cursor<Vec<u8>>> {
json_response(status, &serde_json::to_value(ErrorEnvelope {
error: msg.into(), detail: None,
}).unwrap())
}
fn error_with_detail(status: u16, msg: impl Into<String>, detail: serde_json::Value)
-> Response<std::io::Cursor<Vec<u8>>>
{
json_response(status, &serde_json::to_value(ErrorEnvelope {
error: msg.into(), detail: Some(detail),
}).unwrap())
}
fn write_error_response(prefix: &str, err: lex_store::StoreError)
-> Response<std::io::Cursor<Vec<u8>>>
{
if let lex_store::StoreError::Contention { branch, attempts } = &err {
let body = serde_json::to_vec(&ErrorEnvelope {
error: format!("{prefix}: branch '{branch}' is contended (attempts={attempts})"),
detail: Some(serde_json::json!({
"kind": "contention",
"branch": branch,
"attempts": attempts,
})),
}).unwrap_or_else(|_| b"{}".to_vec());
return Response::from_data(body)
.with_status_code(503)
.with_header(Header::from_bytes(&b"Content-Type"[..], &b"application/json"[..]).unwrap())
.with_header(Header::from_bytes(&b"Retry-After"[..], &b"1"[..]).unwrap());
}
if let lex_store::StoreError::BudgetExceeded { session_id, cap, spent_after } = &err {
let body = serde_json::to_vec(&ErrorEnvelope {
error: format!(
"{prefix}: session `{session_id}` budget exceeded \
(spent_after={spent_after}, cap={cap})"
),
detail: Some(serde_json::json!({
"kind": "budget_exceeded",
"session_id": session_id,
"cap": cap,
"spent_after": spent_after,
})),
}).unwrap_or_else(|_| b"{}".to_vec());
return Response::from_data(body)
.with_status_code(503)
.with_header(Header::from_bytes(&b"Content-Type"[..], &b"application/json"[..]).unwrap())
.with_header(Header::from_bytes(&b"Retry-After"[..], &b"0"[..]).unwrap());
}
error_response(500, format!("{prefix}: {err}"))
}
pub fn handle(state: Arc<State>, mut req: Request) -> std::io::Result<()> {
let method = req.method().clone();
let url = req.url().to_string();
let path = url.split('?').next().unwrap_or("").to_string();
let query = url.split_once('?').map(|(_, q)| q.to_string()).unwrap_or_default();
let x_lex_user = req.headers().iter()
.find(|h| h.field.equiv("x-lex-user"))
.map(|h| h.value.as_str().to_string());
if matches!(method, Method::Post) && path == "/v1/pkg/publish" {
let mut body_bytes: Vec<u8> = Vec::new();
let _ = req.as_reader().read_to_end(&mut body_bytes);
let resp = pkg_publish_handler(&state, &body_bytes);
return req.respond(resp);
}
let mut body = String::new();
let _ = req.as_reader().read_to_string(&mut body);
let resp = route(&state, &method, &path, &query, &body, x_lex_user.as_deref());
req.respond(resp)
}
pub fn handle_with_auth<F>(state: Arc<State>, req: Request, auth: F) -> std::io::Result<()>
where
F: FnOnce(&str, &[Header]) -> bool,
{
let path = req.url().split('?').next().unwrap_or("").to_string();
if !auth(&path, req.headers()) {
return req.respond(
Response::from_data(br#"{"error":"unauthorized"}"#.to_vec())
.with_status_code(401)
.with_header(
Header::from_bytes(&b"Content-Type"[..], &b"application/json"[..]).unwrap(),
),
);
}
handle(state, req)
}
fn route(
state: &State,
method: &Method,
path: &str,
query: &str,
body: &str,
x_lex_user: Option<&str>,
) -> Response<std::io::Cursor<Vec<u8>>> {
match (method, path) {
(Method::Get, "/") => crate::web::activity_handler(state),
(Method::Get, "/web/branches") => crate::web::branches_handler(state),
(Method::Get, "/web/trust") => crate::web::trust_handler(state),
(Method::Get, "/web/attention") => crate::web::attention_handler(state),
(Method::Get, p) if p.starts_with("/web/branch/") => {
let name = &p["/web/branch/".len()..];
crate::web::branch_handler(state, name)
}
(Method::Get, p) if p.starts_with("/web/stage/") => {
let id = &p["/web/stage/".len()..];
crate::web::stage_html_handler(state, id)
}
(Method::Post, p) if p.starts_with("/web/stage/") && (
p.ends_with("/pin") || p.ends_with("/defer")
|| p.ends_with("/block") || p.ends_with("/unblock")
) => {
let prefix_len = "/web/stage/".len();
let last_slash = p.rfind('/').unwrap_or(p.len());
let id = &p[prefix_len..last_slash];
let verb = &p[last_slash + 1..];
let decision = match verb {
"pin" => crate::web::WebStageDecision::Pin,
"defer" => crate::web::WebStageDecision::Defer,
"block" => crate::web::WebStageDecision::Block,
"unblock" => crate::web::WebStageDecision::Unblock,
_ => unreachable!("matched in outer guard"),
};
crate::web::stage_decision_handler(state, id, body, decision, x_lex_user)
}
(Method::Get, "/v1/health") => json_response(200, &serde_json::json!({"ok": true})),
(Method::Post, "/v1/parse") => parse_handler(body),
(Method::Post, "/v1/check") => check_handler(body),
(Method::Post, "/v1/publish") => publish_handler(state, body),
(Method::Post, "/v1/patch") => patch_handler(state, body),
(Method::Get, p) if p.starts_with("/v1/stage/") => {
let suffix = &p["/v1/stage/".len()..];
if let Some(id) = suffix.strip_suffix("/attestations") {
stage_attestations_handler(state, id)
} else {
stage_handler(state, suffix)
}
}
(Method::Post, "/v1/run") => run_handler(state, body, false),
(Method::Post, "/v1/replay") => run_handler(state, body, true),
(Method::Get, p) if p.starts_with("/v1/trace/") => {
let id = &p["/v1/trace/".len()..];
trace_handler(state, id)
}
(Method::Get, "/v1/diff") => diff_handler(state, query),
(Method::Post, "/v1/merge/start") => merge_start_handler(state, body),
(Method::Post, p) if p.starts_with("/v1/merge/") && p.ends_with("/resolve") => {
let id = &p["/v1/merge/".len()..p.len() - "/resolve".len()];
merge_resolve_handler(state, id, body)
}
(Method::Post, p) if p.starts_with("/v1/merge/") && p.ends_with("/commit") => {
let id = &p["/v1/merge/".len()..p.len() - "/commit".len()];
merge_commit_handler(state, id)
}
(Method::Post, "/v1/ops/batch") => ops_batch_handler(state, body),
(Method::Post, "/v1/attestations/batch") => attestations_batch_handler(state, body),
(Method::Get, p) if p.starts_with("/v1/branches/") && p.ends_with("/head") => {
let name = &p["/v1/branches/".len()..p.len() - "/head".len()];
branch_head_handler(state, name)
}
(Method::Get, "/v1/ops/since") => ops_since_handler(state, query),
(Method::Get, "/v1/attestations/since") => attestations_since_handler(state, query),
(Method::Get, "/v1/pkg") => pkg_list_handler(state),
(Method::Put, p) if p.starts_with("/v1/pkg/") && p.ends_with("/visibility") => {
let name = &p["/v1/pkg/".len()..p.len() - "/visibility".len()];
pkg_set_visibility_handler(state, name, body)
}
(Method::Get, p) if p.starts_with("/v1/pkg/") && p.ends_with("/head") => {
let name = &p["/v1/pkg/".len()..p.len() - "/head".len()];
pkg_head_handler(state, name)
}
(Method::Get, p) if p.starts_with("/v1/pkg/") && p.ends_with("/versions") => {
let name = &p["/v1/pkg/".len()..p.len() - "/versions".len()];
pkg_versions_handler(state, name)
}
(Method::Get, p) if p.starts_with("/v1/pkg/") && p.ends_with("/archive") => {
let inner = &p["/v1/pkg/".len()..p.len() - "/archive".len()];
if let Some((name, version)) = inner.split_once('/') {
pkg_archive_handler(state, name, version)
} else {
error_response(400, "expected /v1/pkg/{name}/{version}/archive")
}
}
(Method::Get, p) if p.starts_with("/v1/pkg/") && p["/v1/pkg/".len()..].contains('/') => {
let inner = &p["/v1/pkg/".len()..];
if let Some((name, version)) = inner.split_once('/') {
pkg_get_version_handler(state, name, version)
} else {
error_response(400, "expected /v1/pkg/{name}/{version}")
}
}
(Method::Get, p) if p.starts_with("/v1/pkg/") => {
let name = &p["/v1/pkg/".len()..];
pkg_get_handler(state, name)
}
(Method::Delete, p) if p.starts_with("/v1/pkg/") => {
let name = &p["/v1/pkg/".len()..];
pkg_delete_handler(state, name)
}
_ => error_response(404, format!("unknown route: {method:?} {path}")),
}
}
#[derive(Deserialize)]
struct ParseReq { source: String }
fn parse_handler(body: &str) -> Response<std::io::Cursor<Vec<u8>>> {
let req: ParseReq = match serde_json::from_str(body) {
Ok(r) => r, Err(e) => return error_response(400, format!("bad request: {e}")),
};
match load_program_from_str(&req.source) {
Ok(prog) => {
let stages = canonicalize_program(&prog);
json_response(200, &serde_json::to_value(&stages).unwrap())
}
Err(e) => error_response(400, format!("syntax error: {e}")),
}
}
pub(crate) fn check_handler(body: &str) -> Response<std::io::Cursor<Vec<u8>>> {
let req: ParseReq = match serde_json::from_str(body) {
Ok(r) => r, Err(e) => return error_response(400, format!("bad request: {e}")),
};
let prog = match load_program_from_str(&req.source) {
Ok(p) => p, Err(e) => return error_response(400, format!("syntax error: {e}")),
};
let stages = canonicalize_program(&prog);
match lex_types::check_program(&stages) {
Ok(_) => json_response(200, &serde_json::json!({"ok": true})),
Err(errs) => json_response(422, &serde_json::to_value(&errs).unwrap()),
}
}
#[derive(Deserialize)]
struct PublishReq { source: String, #[serde(default)] activate: bool }
pub(crate) fn publish_handler(state: &State, body: &str) -> Response<std::io::Cursor<Vec<u8>>> {
let req: PublishReq = match serde_json::from_str(body) {
Ok(r) => r, Err(e) => return error_response(400, format!("bad request: {e}")),
};
let prog = match load_program_from_str(&req.source) {
Ok(p) => p, Err(e) => return error_response(400, format!("syntax error: {e}")),
};
let mut stages = canonicalize_program(&prog);
if let Err(errs) = lex_types::check_and_rewrite_program(&mut stages) {
return error_with_detail(422, "type errors", serde_json::to_value(&errs).unwrap());
}
let example_errors = lex_runtime::evaluate_examples(&stages);
if !example_errors.is_empty() {
return error_with_detail(422, "example mismatch",
serde_json::to_value(&example_errors).unwrap_or_default());
}
let store = state.store.lock().unwrap();
let branch = store.current_branch();
let old_head = match store.branch_head(&branch) {
Ok(h) => h,
Err(e) => return error_response(500, format!("branch_head: {e}")),
};
let old_fns: std::collections::BTreeMap<String, lex_ast::FnDecl> = old_head.values()
.filter_map(|stg| store.get_ast(stg).ok())
.filter_map(|s| match s {
lex_ast::Stage::FnDecl(fd) => Some((fd.name.clone(), fd)),
_ => None,
})
.collect();
let new_fns: std::collections::BTreeMap<String, lex_ast::FnDecl> = stages.iter()
.filter_map(|s| match s {
lex_ast::Stage::FnDecl(fd) => Some((fd.name.clone(), fd.clone())),
_ => None,
})
.collect();
let report = lex_vcs::compute_diff(&old_fns, &new_fns, false);
let mut new_imports: lex_vcs::ImportMap = lex_vcs::ImportMap::new();
{
let entry = new_imports.entry("<source>".into()).or_default();
for s in &stages {
if let lex_ast::Stage::Import(im) = s {
entry.insert(im.reference.clone());
}
}
}
match store.publish_program(&branch, &stages, &report, &new_imports, req.activate) {
Ok(outcome) => {
record_examples_for_publish(&store, &stages, &outcome);
json_response(200, &serde_json::json!({
"ops": outcome.ops,
"head_op": outcome.head_op,
}))
}
Err(lex_store::StoreError::TypeError(errs)) => {
error_with_detail(422, "type errors", serde_json::to_value(&errs).unwrap())
}
Err(e) => write_error_response("publish_program", e),
}
}
#[derive(Deserialize)]
struct PatchReq {
stage_id: String,
patch: lex_ast::Patch,
#[serde(default)] activate: bool,
}
fn patch_handler(state: &State, body: &str) -> Response<std::io::Cursor<Vec<u8>>> {
let req: PatchReq = match serde_json::from_str(body) {
Ok(r) => r, Err(e) => return error_response(400, format!("bad request: {e}")),
};
let store = state.store.lock().unwrap();
let original = match store.get_ast(&req.stage_id) {
Ok(s) => s, Err(e) => return error_response(404, format!("stage: {e}")),
};
let patched = match lex_ast::apply_patch(&original, &req.patch) {
Ok(s) => s,
Err(e) => return error_with_detail(422, "patch failed",
serde_json::to_value(&e).unwrap_or_default()),
};
let branch = store.current_branch();
let sig = match lex_ast::sig_id(&patched) {
Some(s) => s,
None => return error_response(500, "patched stage has no sig_id"),
};
let new_id = match store.publish(&patched) {
Ok(id) => id, Err(e) => return error_response(500, format!("publish: {e}")),
};
let original_effects: std::collections::BTreeSet<String> = match &original {
lex_ast::Stage::FnDecl(fd) => fd.effects.iter().map(|e| e.name.clone()).collect(),
_ => std::collections::BTreeSet::new(),
};
let patched_effects: std::collections::BTreeSet<String> = match &patched {
lex_ast::Stage::FnDecl(fd) => fd.effects.iter().map(|e| e.name.clone()).collect(),
_ => std::collections::BTreeSet::new(),
};
let head_now = match store.get_branch(&branch) {
Ok(b) => b.and_then(|b| b.head_op),
Err(e) => return error_response(500, format!("get_branch: {e}")),
};
let kind = if original_effects != patched_effects {
let from_budget = lex_vcs::operation_budget_from_effects(&original_effects);
let to_budget = lex_vcs::operation_budget_from_effects(&patched_effects);
lex_vcs::OperationKind::ChangeEffectSig {
sig_id: sig.clone(),
from_stage_id: req.stage_id.clone(),
to_stage_id: new_id.clone(),
from_effects: original_effects,
to_effects: patched_effects,
from_budget,
to_budget,
}
} else {
let budget = lex_vcs::operation_budget_from_effects(&original_effects);
lex_vcs::OperationKind::ModifyBody {
sig_id: sig.clone(),
from_stage_id: req.stage_id.clone(),
to_stage_id: new_id.clone(),
from_budget: budget,
to_budget: budget,
}
};
let transition = lex_vcs::StageTransition::Replace {
sig_id: sig.clone(),
from: req.stage_id.clone(),
to: new_id.clone(),
};
let op = lex_vcs::Operation::new(
kind,
head_now.into_iter().collect::<Vec<_>>(),
);
let op_id = match store.apply_operation_gated(&branch, op, transition) {
Ok(id) => id,
Err(lex_store::StoreError::TypeError(errs)) => return error_with_detail(
422, "type errors after patch", serde_json::to_value(&errs).unwrap_or_default()),
Err(e) => return write_error_response("apply_operation_gated", e),
};
if req.activate {
if let Err(e) = store.activate(&new_id) {
return error_response(500, format!("activate: {e}"));
}
}
let status = format!("{:?}",
store.get_status(&new_id).unwrap_or(lex_store::StageStatus::Draft)).to_lowercase();
json_response(200, &serde_json::json!({
"old_stage_id": req.stage_id,
"new_stage_id": new_id,
"sig_id": sig,
"status": status,
"op_id": op_id,
}))
}
pub(crate) fn stage_handler(state: &State, id: &str) -> Response<std::io::Cursor<Vec<u8>>> {
let store = state.store.lock().unwrap();
let meta = match store.get_metadata(id) {
Ok(m) => m, Err(e) => return error_response(404, format!("{e}")),
};
let ast = match store.get_ast(id) {
Ok(a) => a, Err(e) => return error_response(404, format!("{e}")),
};
let status = format!("{:?}", store.get_status(id).unwrap_or(lex_store::StageStatus::Draft)).to_lowercase();
json_response(200, &serde_json::json!({
"metadata": meta,
"ast": ast,
"status": status,
}))
}
pub(crate) fn stage_attestations_handler(state: &State, id: &str) -> Response<std::io::Cursor<Vec<u8>>> {
let store = state.store.lock().unwrap();
if let Err(e) = store.get_metadata(id) {
return error_response(404, format!("{e}"));
}
let log = match store.attestation_log() {
Ok(l) => l,
Err(e) => return error_response(500, format!("attestation log: {e}")),
};
let mut listing = match log.list_for_stage(&id.to_string()) {
Ok(v) => v,
Err(e) => return error_response(500, format!("list_for_stage: {e}")),
};
listing.sort_by_key(|a| std::cmp::Reverse(a.timestamp));
json_response(200, &serde_json::json!({"attestations": listing}))
}
#[derive(Deserialize, Default)]
struct PolicyJson {
#[serde(default)] allow_effects: Vec<String>,
#[serde(default)] allow_fs_read: Vec<String>,
#[serde(default)] allow_fs_write: Vec<String>,
#[serde(default)] budget: Option<u64>,
}
impl PolicyJson {
fn into_policy(self) -> Policy {
Policy {
allow_effects: self.allow_effects.into_iter().collect::<BTreeSet<_>>(),
allow_fs_read: self.allow_fs_read.into_iter().map(PathBuf::from).collect(),
allow_fs_write: self.allow_fs_write.into_iter().map(PathBuf::from).collect(),
allow_net_host: Vec::new(),
allow_proc: Vec::new(),
allow_approval: Vec::new(),
budget: self.budget,
}
}
}
#[derive(Deserialize)]
struct RunReq {
source: String,
#[serde(rename = "fn")] func: String,
#[serde(default)] args: Vec<serde_json::Value>,
#[serde(default)] policy: PolicyJson,
#[serde(default)] overrides: IndexMap<String, serde_json::Value>,
}
pub(crate) fn run_handler(state: &State, body: &str, with_overrides: bool) -> Response<std::io::Cursor<Vec<u8>>> {
let req: RunReq = match serde_json::from_str(body) {
Ok(r) => r, Err(e) => return error_response(400, format!("bad request: {e}")),
};
let prog = match load_program_from_str(&req.source) {
Ok(p) => p, Err(e) => return error_response(400, format!("syntax error: {e}")),
};
let stages = canonicalize_program(&prog);
if let Err(errs) = lex_types::check_program(&stages) {
return error_with_detail(422, "type errors", serde_json::to_value(&errs).unwrap());
}
let bc = compile_program(&stages);
let mut policy = req.policy.into_policy();
if let Some(ceiling) = &state.policy_ceiling {
policy = clamp_policy(policy, ceiling);
}
if let Err(violations) = check_policy(&bc, &policy) {
return error_with_detail(403, "policy violation", serde_json::to_value(&violations).unwrap());
}
let mut recorder = lex_trace::Recorder::new();
if with_overrides && !req.overrides.is_empty() {
recorder = recorder.with_overrides(req.overrides);
}
let handle = recorder.handle();
let handler = DefaultHandler::new(policy);
let mut vm = Vm::with_handler(&bc, Box::new(handler));
vm.set_tracer(Box::new(recorder));
let vargs: Vec<Value> = req.args.iter().map(json_to_value).collect();
let started = std::time::SystemTime::now().duration_since(std::time::UNIX_EPOCH).unwrap().as_secs();
let result = vm.call(&req.func, vargs);
let ended = std::time::SystemTime::now().duration_since(std::time::UNIX_EPOCH).unwrap().as_secs();
let store = state.store.lock().unwrap();
let (root_out, root_err, status) = match &result {
Ok(v) => (Some(value_to_json(v)), None, 200u16),
Err(e) => (None, Some(format!("{e}")), 200u16),
};
let tree = handle.finalize(req.func.clone(), serde_json::Value::Null,
root_out.clone(), root_err.clone(), started, ended);
let run_id = match store.save_trace(&tree) {
Ok(id) => id,
Err(e) => return error_response(500, format!("save_trace: {e}")),
};
let mut body = serde_json::json!({
"run_id": run_id,
"output": root_out,
});
if let Some(err) = root_err {
body["error"] = serde_json::Value::String(err);
}
json_response(status, &body)
}
fn trace_handler(state: &State, id: &str) -> Response<std::io::Cursor<Vec<u8>>> {
let store = state.store.lock().unwrap();
match store.load_trace(id) {
Ok(t) => json_response(200, &serde_json::to_value(&t).unwrap()),
Err(e) => error_response(404, format!("{e}")),
}
}
fn diff_handler(state: &State, query: &str) -> Response<std::io::Cursor<Vec<u8>>> {
let mut a = None;
let mut b = None;
for kv in query.split('&') {
if let Some((k, v)) = kv.split_once('=') {
match k { "a" => a = Some(v.to_string()), "b" => b = Some(v.to_string()), _ => {} }
}
}
let (Some(a), Some(b)) = (a, b) else {
return error_response(400, "missing a or b query params");
};
let store = state.store.lock().unwrap();
let ta = match store.load_trace(&a) { Ok(t) => t, Err(e) => return error_response(404, format!("a: {e}")) };
let tb = match store.load_trace(&b) { Ok(t) => t, Err(e) => return error_response(404, format!("b: {e}")) };
match lex_trace::diff_runs(&ta, &tb) {
Some(d) => json_response(200, &serde_json::to_value(&d).unwrap()),
None => json_response(200, &serde_json::json!({"divergence": null})),
}
}
fn json_to_value(v: &serde_json::Value) -> Value { Value::from_json(v) }
fn value_to_json(v: &Value) -> serde_json::Value { v.to_json() }
#[derive(Deserialize)]
struct MergeStartReq {
src_branch: String,
dst_branch: String,
}
fn merge_start_handler(state: &State, body: &str) -> Response<std::io::Cursor<Vec<u8>>> {
let req: MergeStartReq = match serde_json::from_str(body) {
Ok(r) => r, Err(e) => return error_response(400, format!("bad request: {e}")),
};
let store = state.store.lock().unwrap();
let src_head = match store.get_branch(&req.src_branch) {
Ok(Some(b)) => b.head_op,
Ok(None) => return error_response(404, format!("unknown src branch `{}`", req.src_branch)),
Err(e) => return error_response(500, format!("src branch read: {e}")),
};
let dst_head = match store.get_branch(&req.dst_branch) {
Ok(Some(b)) => b.head_op,
Ok(None) => return error_response(404, format!("unknown dst branch `{}`", req.dst_branch)),
Err(e) => return error_response(500, format!("dst branch read: {e}")),
};
let log = match lex_vcs::OpLog::open(store.root()) {
Ok(l) => l,
Err(e) => return error_response(500, format!("op log: {e}")),
};
let merge_id = mint_merge_id();
let session = match MergeSession::start(
merge_id.clone(),
&log,
src_head.as_ref(),
dst_head.as_ref(),
) {
Ok(s) => s,
Err(e) => return error_response(500, format!("merge start: {e}")),
};
let conflicts: Vec<&lex_vcs::ConflictRecord> = session.remaining_conflicts();
let auto_resolved_count = session.auto_resolved.len();
let body = serde_json::json!({
"merge_id": merge_id,
"src_head": session.src_head,
"dst_head": session.dst_head,
"lca": session.lca,
"conflicts": conflicts,
"auto_resolved_count": auto_resolved_count,
});
drop(conflicts);
drop(store);
let wrapped = ApiMergeSession {
inner: session,
src_branch: req.src_branch,
dst_branch: req.dst_branch,
};
state.sessions.lock().unwrap().insert(merge_id, wrapped);
json_response(200, &body)
}
#[derive(Deserialize)]
struct MergeResolveReq {
resolutions: Vec<MergeResolveEntry>,
}
#[derive(Deserialize)]
struct MergeResolveEntry {
conflict_id: String,
resolution: lex_vcs::Resolution,
}
fn merge_resolve_handler(
state: &State,
merge_id: &str,
body: &str,
) -> Response<std::io::Cursor<Vec<u8>>> {
let req: MergeResolveReq = match serde_json::from_str(body) {
Ok(r) => r, Err(e) => return error_response(400, format!("bad request: {e}")),
};
let mut sessions = state.sessions.lock().unwrap();
let Some(wrapped) = sessions.get_mut(merge_id) else {
return error_response(404, format!("unknown merge_id `{merge_id}`"));
};
let pairs: Vec<(String, lex_vcs::Resolution)> = req.resolutions.into_iter()
.map(|e| (e.conflict_id, e.resolution))
.collect();
let store = state.store.lock().unwrap();
let checker = lex_store::MergeResolutionChecker::new(&store, wrapped.dst_branch.clone());
let verdicts = wrapped.inner.resolve_checked(pairs, &checker);
drop(store);
let remaining: Vec<&lex_vcs::ConflictRecord> = wrapped.inner.remaining_conflicts();
let body = serde_json::json!({
"verdicts": verdicts,
"remaining_conflicts": remaining,
});
json_response(200, &body)
}
fn merge_commit_handler(
state: &State,
merge_id: &str,
) -> Response<std::io::Cursor<Vec<u8>>> {
use std::collections::BTreeMap;
let wrapped = match state.sessions.lock().unwrap().remove(merge_id) {
Some(w) => w,
None => return error_response(404, format!("unknown merge_id `{merge_id}`")),
};
let dst_branch = wrapped.dst_branch.clone();
let src_head = wrapped.inner.src_head.clone();
let dst_head = wrapped.inner.dst_head.clone();
let auto_resolved = wrapped.inner.auto_resolved.clone();
let mut entries: BTreeMap<lex_vcs::SigId, Option<lex_vcs::StageId>> = BTreeMap::new();
for outcome in &auto_resolved {
if let lex_vcs::MergeOutcome::Src { sig_id, stage_id } = outcome {
entries.insert(sig_id.clone(), stage_id.clone());
}
}
let resolved = match wrapped.inner.commit() {
Ok(r) => r,
Err(lex_vcs::CommitError::ConflictsRemaining(ids)) => {
return error_with_detail(
422,
"conflicts remaining",
serde_json::json!({"unresolved": ids}),
);
}
};
for (conflict_id, resolution) in resolved {
match resolution {
lex_vcs::Resolution::TakeOurs => {
}
lex_vcs::Resolution::TakeTheirs => {
match resolve_take_theirs(state, &src_head, &conflict_id) {
Ok(stage_id) => {
entries.insert(conflict_id.clone(), stage_id);
}
Err(e) => return error_response(500, format!("resolve take_theirs: {e}")),
}
}
lex_vcs::Resolution::Custom { op } => {
match op.kind.merge_target() {
Some((sig, stage)) => {
if sig != conflict_id {
return error_with_detail(
422,
"custom op targets a different sig than the conflict",
serde_json::json!({
"conflict_id": conflict_id,
"op_targets": sig,
}),
);
}
entries.insert(conflict_id, stage);
}
None => {
return error_with_detail(
422,
"custom op kind doesn't yield a single sig→stage delta",
serde_json::json!({
"conflict_id": conflict_id,
"kind": serde_json::to_value(&op.kind).unwrap_or(serde_json::Value::Null),
}),
);
}
}
}
lex_vcs::Resolution::Defer => {
return error_response(500, "internal: Defer slipped past commit gate");
}
}
}
let resolved_count = entries.len();
let mut parents: Vec<lex_vcs::OpId> = Vec::new();
if let Some(d) = dst_head { parents.push(d); }
if let Some(s) = src_head { parents.push(s); }
let op = lex_vcs::Operation::new(
lex_vcs::OperationKind::Merge { resolved: resolved_count },
parents,
);
let transition = lex_vcs::StageTransition::Merge { entries };
let store = state.store.lock().unwrap();
match store.apply_merge_op_gated(&dst_branch, op, transition) {
Ok(new_head_op) => json_response(200, &serde_json::json!({
"new_head_op": new_head_op,
"dst_branch": dst_branch,
})),
Err(lex_store::StoreError::TypeError(errs)) => error_with_detail(
422, "merged program has type errors", serde_json::to_value(&errs).unwrap_or_default()),
Err(e) => write_error_response("apply merge op", e),
}
}
fn resolve_take_theirs(
state: &State,
src_head: &Option<lex_vcs::OpId>,
sig: &lex_vcs::SigId,
) -> std::io::Result<Option<lex_vcs::StageId>> {
let store = state.store.lock().unwrap();
let log = lex_vcs::OpLog::open(store.root())?;
let Some(head) = src_head.as_ref() else { return Ok(None); };
let mut current: Option<lex_vcs::StageId> = None;
for record in log.walk_forward(head, None)? {
match &record.produces {
lex_vcs::StageTransition::Create { sig_id, stage_id }
if sig_id == sig => { current = Some(stage_id.clone()); }
lex_vcs::StageTransition::Replace { sig_id, to, .. }
if sig_id == sig => { current = Some(to.clone()); }
lex_vcs::StageTransition::Remove { sig_id, .. }
if sig_id == sig => { current = None; }
lex_vcs::StageTransition::Rename { from, to, body_stage_id }
if from == sig || to == sig => {
if from == sig { current = None; }
if to == sig { current = Some(body_stage_id.clone()); }
}
lex_vcs::StageTransition::Merge { entries } => {
if let Some(opt) = entries.get(sig) {
current = opt.clone();
}
}
_ => {}
}
}
Ok(current)
}
fn mint_merge_id() -> MergeSessionId {
use std::sync::atomic::{AtomicU64, Ordering};
static COUNTER: AtomicU64 = AtomicU64::new(0);
let nanos = SystemTime::now()
.duration_since(UNIX_EPOCH)
.map(|d| d.as_nanos())
.unwrap_or(0);
let n = COUNTER.fetch_add(1, Ordering::Relaxed);
format!("merge_{nanos:x}_{n:x}")
}
pub(crate) fn ops_batch_handler(state: &State, body: &str)
-> Response<std::io::Cursor<Vec<u8>>>
{
let records: Vec<lex_vcs::OperationRecord> = match serde_json::from_str(body) {
Ok(r) => r,
Err(e) => return error_response(400,
format!("body must be a JSON array of OperationRecord: {e}")),
};
let store = state.store.lock().unwrap();
let log = match lex_vcs::OpLog::open(store.root()) {
Ok(l) => l,
Err(e) => return error_response(500, format!("opening op log: {e}")),
};
let mut batch_ids: std::collections::BTreeSet<lex_vcs::OpId> =
std::collections::BTreeSet::new();
for rec in &records {
let expected = rec.op.op_id();
if expected != rec.op_id {
return error_with_detail(409, "OpIdMismatch", serde_json::json!({
"supplied": rec.op_id,
"expected": expected,
}));
}
for parent in &rec.op.parents {
let known = match log.get(parent) {
Ok(Some(_)) => true,
Ok(None) => false,
Err(e) => return error_response(500, format!("op log read: {e}")),
};
if !known && !batch_ids.contains(parent) {
return error_with_detail(422, "MissingParent", serde_json::json!({
"op_id": rec.op_id,
"missing_parent": parent,
}));
}
}
batch_ids.insert(rec.op_id.clone());
}
let mut added = 0usize;
let mut added_ids: Vec<&lex_vcs::OpId> = Vec::new();
for rec in &records {
let already_present = matches!(log.get(&rec.op_id), Ok(Some(_)));
match log.put(rec) {
Ok(()) => {
if !already_present {
added += 1;
added_ids.push(&rec.op_id);
}
}
Err(e) => return error_response(500, format!("op log write: {e}")),
}
}
json_response(200, &serde_json::json!({
"received": records.len(),
"added": added,
"skipped": records.len() - added,
"added_ids": added_ids,
}))
}
pub(crate) fn attestations_batch_handler(state: &State, body: &str)
-> Response<std::io::Cursor<Vec<u8>>>
{
let attestations: Vec<lex_vcs::Attestation> = match serde_json::from_str(body) {
Ok(a) => a,
Err(e) => return error_response(400,
format!("body must be a JSON array of Attestation: {e}")),
};
let store = state.store.lock().unwrap();
let log = match store.attestation_log() {
Ok(l) => l,
Err(e) => return error_response(500, format!("opening attestation log: {e}")),
};
let op_log = match lex_vcs::OpLog::open(store.root()) {
Ok(l) => l,
Err(e) => return error_response(500, format!("opening op log: {e}")),
};
for att in &attestations {
let expected = lex_vcs::Attestation::with_timestamp(
att.stage_id.clone(),
att.op_id.clone(),
att.intent_id.clone(),
att.kind.clone(),
att.result.clone(),
att.produced_by.clone(),
att.cost.clone(),
att.timestamp,
).attestation_id;
if expected != att.attestation_id {
return error_with_detail(409, "AttestationIdMismatch", serde_json::json!({
"supplied": att.attestation_id,
"expected": expected,
}));
}
if let Some(op_id) = &att.op_id {
match op_log.get(op_id) {
Ok(Some(_)) => {}
Ok(None) => return error_with_detail(422, "UnknownOp", serde_json::json!({
"attestation_id": att.attestation_id,
"op_id": op_id,
})),
Err(e) => return error_response(500, format!("op log read: {e}")),
}
}
}
let mut added = 0usize;
let mut added_ids: Vec<&lex_vcs::AttestationId> = Vec::new();
for att in &attestations {
let already_present = matches!(log.get(&att.attestation_id), Ok(Some(_)));
match log.put(att) {
Ok(()) => {
if !already_present {
added += 1;
added_ids.push(&att.attestation_id);
}
}
Err(e) => return error_response(500, format!("attestation log write: {e}")),
}
}
json_response(200, &serde_json::json!({
"received": attestations.len(),
"added": added,
"skipped": attestations.len() - added,
"added_ids": added_ids,
}))
}
pub(crate) fn branch_head_handler(state: &State, name: &str)
-> Response<std::io::Cursor<Vec<u8>>>
{
let store = state.store.lock().unwrap();
let head = match store.get_branch(name) {
Ok(Some(b)) => b.head_op,
Ok(None) => None,
Err(e) => return error_response(500, format!("get_branch: {e}")),
};
json_response(200, &serde_json::json!({
"branch": name,
"head_op": head,
}))
}
pub(crate) fn ops_since_handler(state: &State, query: &str)
-> Response<std::io::Cursor<Vec<u8>>>
{
let mut after: Option<String> = None;
let mut branch = String::from("main");
let mut limit: Option<usize> = None;
for kv in query.split('&') {
let Some((k, v)) = kv.split_once('=') else { continue };
match k {
"after" => after = Some(v.to_string()),
"branch" => branch = v.to_string(),
"limit" => {
limit = Some(match v.parse::<usize>() {
Ok(n) => n,
Err(_) => return error_response(400,
format!("limit must be a positive integer, got `{v}`")),
});
}
_ => {}
}
}
let store = state.store.lock().unwrap();
let log = match lex_vcs::OpLog::open(store.root()) {
Ok(l) => l,
Err(e) => return error_response(500, format!("opening op log: {e}")),
};
let head = match store.get_branch(&branch) {
Ok(Some(b)) => b.head_op,
Ok(None) => None,
Err(e) => return error_response(500, format!("get_branch: {e}")),
};
let Some(head) = head else {
return json_response(200, &serde_json::json!([]));
};
let ops_since = match log.ops_since(&head, after.as_ref()) {
Ok(o) => o,
Err(e) => return error_response(500, format!("ops_since: {e}")),
};
let mut ops = ops_since;
ops.reverse();
if let Some(n) = limit {
ops.truncate(n);
}
json_response(200, &serde_json::to_value(&ops).unwrap_or_default())
}
pub(crate) fn attestations_since_handler(state: &State, query: &str)
-> Response<std::io::Cursor<Vec<u8>>>
{
let mut after_op: Option<String> = None;
let mut limit: Option<usize> = None;
for kv in query.split('&') {
let Some((k, v)) = kv.split_once('=') else { continue };
match k {
"after-op" => after_op = Some(v.to_string()),
"limit" => {
limit = Some(match v.parse::<usize>() {
Ok(n) => n,
Err(_) => return error_response(400,
format!("limit must be a positive integer, got `{v}`")),
});
}
_ => {}
}
}
let store = state.store.lock().unwrap();
let log = match store.attestation_log() {
Ok(l) => l,
Err(e) => return error_response(500, format!("opening attestation log: {e}")),
};
let exclude: std::collections::BTreeSet<String> = match &after_op {
None => std::collections::BTreeSet::new(),
Some(cutoff) => {
let op_log = match lex_vcs::OpLog::open(store.root()) {
Ok(l) => l,
Err(e) => return error_response(500, format!("opening op log: {e}")),
};
match op_log.walk_back(cutoff, None) {
Ok(records) => records.into_iter().map(|r| r.op_id).collect(),
Err(_) => {
std::collections::BTreeSet::new()
}
}
}
};
let all = match log.list_all() {
Ok(v) => v,
Err(e) => return error_response(500, format!("listing attestations: {e}")),
};
let mut filtered: Vec<lex_vcs::Attestation> = all
.into_iter()
.filter(|a| match &a.op_id {
Some(op_id) => !exclude.contains(op_id),
None => true,
})
.collect();
filtered.sort_by(|a, b| {
a.timestamp.cmp(&b.timestamp)
.then_with(|| a.attestation_id.cmp(&b.attestation_id))
});
if let Some(n) = limit {
filtered.truncate(n);
}
json_response(200, &serde_json::to_value(&filtered).unwrap_or_default())
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
struct PkgRecord {
name: String,
version: String,
head_op: Option<String>,
published_at: u64,
function_names: Vec<String>,
ops: Vec<serde_json::Value>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize, serde::Deserialize, Default)]
#[serde(rename_all = "lowercase")]
pub enum Visibility {
#[default]
Private,
Public,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize, Default)]
struct PkgIndex {
latest: Option<String>,
versions: Vec<PkgVersionSummary>,
#[serde(default)]
visibility: Visibility,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
struct PkgVersionSummary {
version: String,
head_op: Option<String>,
published_at: u64,
}
fn pkg_name_dir(root: &std::path::Path, name: &str) -> PathBuf {
root.join("packages").join(name)
}
fn pkg_index_path(root: &std::path::Path, name: &str) -> PathBuf {
pkg_name_dir(root, name).join("index.json")
}
fn pkg_version_path(root: &std::path::Path, name: &str, version: &str) -> PathBuf {
pkg_name_dir(root, name).join(format!("{version}.json"))
}
fn pkg_archive_path(root: &std::path::Path, name: &str, version: &str) -> PathBuf {
pkg_name_dir(root, name).join(format!("{version}.tar.gz"))
}
fn load_pkg_index(root: &std::path::Path, name: &str) -> Option<PkgIndex> {
let bytes = std::fs::read(pkg_index_path(root, name)).ok()?;
serde_json::from_slice(&bytes).ok()
}
fn load_pkg_record(root: &std::path::Path, name: &str, version: &str) -> Option<PkgRecord> {
let bytes = std::fs::read(pkg_version_path(root, name, version)).ok()?;
serde_json::from_slice(&bytes).ok()
}
fn load_latest_pkg_record(root: &std::path::Path, name: &str) -> Option<PkgRecord> {
let index = load_pkg_index(root, name)?;
let latest = index.latest.clone()?;
load_pkg_record(root, name, &latest)
}
fn pkg_is_public(root: &std::path::Path, name: &str) -> bool {
load_pkg_index(root, name).map(|i| i.visibility) == Some(Visibility::Public)
}
fn valid_pkg_segment(s: &str) -> bool {
!s.is_empty()
&& s.len() <= 128
&& s != "."
&& s != ".."
&& s.chars().all(|c| c.is_ascii_alphanumeric() || matches!(c, '.' | '_' | '-'))
}
#[derive(Deserialize)]
struct VisibilityReq {
visibility: Visibility,
}
fn pkg_set_visibility_handler(
state: &State,
name: &str,
body: &str,
) -> Response<std::io::Cursor<Vec<u8>>> {
if !valid_pkg_segment(name) {
return error_response(400, format!("invalid package name {name:?}"));
}
let req: VisibilityReq = match serde_json::from_str(body) {
Ok(r) => r,
Err(e) => return error_response(400, format!("bad request: {e}")),
};
let mut index = match load_pkg_index(&state.root, name) {
Some(i) => i,
None => return error_response(404, format!("package {name:?} not found")),
};
index.visibility = req.visibility;
let bytes = serde_json::to_vec_pretty(&index).unwrap_or_default();
match std::fs::write(pkg_index_path(&state.root, name), bytes) {
Ok(()) => json_response(
200,
&serde_json::json!({ "name": name, "visibility": index.visibility }),
),
Err(e) => error_response(500, format!("write index: {e}")),
}
}
fn public_pkg_names(root: &std::path::Path) -> Vec<String> {
list_pkg_names(root)
.into_iter()
.filter(|name| pkg_is_public(root, name))
.collect()
}
fn public_pkg_list_handler(state: &State) -> Response<std::io::Cursor<Vec<u8>>> {
let packages: Vec<serde_json::Value> = public_pkg_names(&state.root)
.iter()
.filter_map(|name| {
let r = load_latest_pkg_record(&state.root, name)?;
Some(serde_json::json!({
"name": r.name,
"version": r.version,
"head_op": r.head_op,
"published_at": r.published_at,
}))
})
.collect();
json_response(200, &serde_json::json!({ "packages": packages }))
}
#[derive(Debug, PartialEq, Eq)]
enum PublicTarget {
List,
Latest(String),
Versions(String),
Head(String),
Version(String, String),
Archive(String, String),
}
impl PublicTarget {
fn pkg_name(&self) -> Option<&str> {
match self {
PublicTarget::List => None,
PublicTarget::Latest(n)
| PublicTarget::Versions(n)
| PublicTarget::Head(n)
| PublicTarget::Version(n, _)
| PublicTarget::Archive(n, _) => Some(n),
}
}
}
fn resolve_public(method: &Method, path: &str) -> Result<PublicTarget, u16> {
if !matches!(method, Method::Get) {
return Err(405);
}
let rest = path.trim_matches('/');
if rest.is_empty() {
return Ok(PublicTarget::List);
}
let segs: Vec<&str> = rest.split('/').collect();
if !segs.iter().all(|s| valid_pkg_segment(s)) {
return Err(404);
}
match segs.as_slice() {
[n] => Ok(PublicTarget::Latest(n.to_string())),
[n, "versions"] => Ok(PublicTarget::Versions(n.to_string())),
[n, "head"] => Ok(PublicTarget::Head(n.to_string())),
[n, v, "archive"] => Ok(PublicTarget::Archive(n.to_string(), v.to_string())),
[n, v] => Ok(PublicTarget::Version(n.to_string(), v.to_string())),
_ => Err(404),
}
}
pub fn route_public(
state: &State,
method: &Method,
path: &str,
_query: &str,
) -> Response<std::io::Cursor<Vec<u8>>> {
let target = match resolve_public(method, path) {
Ok(t) => t,
Err(405) => return error_response(405, "public read is GET-only"),
Err(_) => return error_response(404, "not found"),
};
if let PublicTarget::List = target {
return public_pkg_list_handler(state);
}
if let Some(name) = target.pkg_name() {
if !pkg_is_public(&state.root, name) {
return error_response(404, format!("package {name:?} not found"));
}
}
match target {
PublicTarget::List => unreachable!("handled above"),
PublicTarget::Latest(n) => pkg_get_handler(state, &n),
PublicTarget::Versions(n) => pkg_versions_handler(state, &n),
PublicTarget::Head(n) => pkg_head_handler(state, &n),
PublicTarget::Version(n, v) => pkg_get_version_handler(state, &n, &v),
PublicTarget::Archive(n, v) => pkg_archive_handler(state, &n, &v),
}
}
fn save_pkg_record(
root: &std::path::Path,
record: &PkgRecord,
archive: &[u8],
) -> std::io::Result<()> {
let dir = pkg_name_dir(root, &record.name);
std::fs::create_dir_all(&dir)?;
let rec_bytes = serde_json::to_vec_pretty(record).unwrap_or_default();
std::fs::write(pkg_version_path(root, &record.name, &record.version), rec_bytes)?;
std::fs::write(pkg_archive_path(root, &record.name, &record.version), archive)?;
let mut index = load_pkg_index(root, &record.name).unwrap_or_default();
index.latest = Some(record.version.clone());
if !index.versions.iter().any(|v| v.version == record.version) {
index.versions.push(PkgVersionSummary {
version: record.version.clone(),
head_op: record.head_op.clone(),
published_at: record.published_at,
});
}
let idx_bytes = serde_json::to_vec_pretty(&index).unwrap_or_default();
std::fs::write(pkg_index_path(root, &record.name), idx_bytes)
}
fn list_pkg_names(root: &std::path::Path) -> Vec<String> {
let dir = root.join("packages");
let Ok(entries) = std::fs::read_dir(&dir) else {
return Vec::new();
};
let mut names: Vec<String> = entries
.filter_map(|e| e.ok())
.filter(|e| e.path().is_dir())
.filter_map(|e| e.file_name().into_string().ok())
.collect();
names.sort();
names
}
fn collect_lex_files(dir: &std::path::Path, out: &mut Vec<PathBuf>) {
let Ok(entries) = std::fs::read_dir(dir) else { return };
let mut entries: Vec<_> = entries.filter_map(|e| e.ok()).collect();
entries.sort_by_key(|e| e.path());
for entry in entries {
let path = entry.path();
if path.is_dir() {
collect_lex_files(&path, out);
} else if path.extension().and_then(|x| x.to_str()) == Some("lex") {
out.push(path);
}
}
}
fn pkg_publish_handler(state: &State, body: &[u8]) -> Response<std::io::Cursor<Vec<u8>>> {
let tmp = match tempfile::TempDir::new() {
Ok(t) => t,
Err(e) => return error_response(500, format!("create temp dir: {e}")),
};
{
let gz = flate2::read::GzDecoder::new(std::io::Cursor::new(body));
let mut ar = tar::Archive::new(gz);
if let Err(e) = ar.unpack(tmp.path()) {
return error_response(400, format!("unpack archive: {e}"));
}
}
let toml_path = tmp.path().join("lex.toml");
if !toml_path.exists() {
return error_response(400, "archive must contain lex.toml at root");
}
let manifest = match Manifest::load(&toml_path) {
Ok(m) => m,
Err(e) => return error_response(400, format!("lex.toml: {e}")),
};
let (pkg_name, pkg_version) = match &manifest.package {
Some(m) => (m.name.clone(), m.version.clone()),
None => return error_response(400, "lex.toml must have a [package] section"),
};
if load_pkg_record(&state.root, &pkg_name, &pkg_version).is_some() {
return error_response(
409,
format!(
"package {pkg_name}@{pkg_version} already published; \
bump the version in lex.toml to publish a new release"
),
);
}
let src_dir = tmp.path().join("src");
if !src_dir.exists() {
return error_response(400, "archive must contain a src/ directory");
}
let mut lex_files: Vec<PathBuf> = Vec::new();
collect_lex_files(&src_dir, &mut lex_files);
if lex_files.is_empty() {
return error_response(400, "no .lex files found in src/");
}
let store = state.store.lock().unwrap();
let branch = store.current_branch();
let old_head = match store.branch_head(&branch) {
Ok(h) => h,
Err(e) => return error_response(500, format!("branch_head: {e}")),
};
let old_pairs: Vec<(String, String)> =
old_head.iter().map(|(sig, stage)| (sig.clone(), stage.clone())).collect();
let mut old_fns_by_name: BTreeMap<String, Vec<lex_ast::FnDecl>> = BTreeMap::new();
for fd in store.get_asts_for_sigs_bulk(&old_pairs)
.into_iter()
.filter_map(|r| r.ok())
.filter_map(|s| match s { lex_ast::Stage::FnDecl(fd) => Some(fd), _ => None })
{
old_fns_by_name.entry(fd.name.clone()).or_default().push(fd);
}
fn structural_key(fd: &lex_ast::FnDecl) -> Option<String> {
let mut anon = fd.clone();
anon.name = String::new();
lex_ast::sig_id(&lex_ast::Stage::FnDecl(anon))
}
fn take_matching(
map: &mut BTreeMap<String, Vec<lex_ast::FnDecl>>,
name: &str,
new_fd: &lex_ast::FnDecl,
) -> Option<lex_ast::FnDecl> {
let candidates = map.get_mut(name)?;
let idx = match candidates.len() {
0 => return None,
1 => 0,
_ => {
let want = structural_key(new_fd);
candidates.iter().position(|c| structural_key(c) == want)?
}
};
let matched = candidates.remove(idx);
if candidates.is_empty() {
map.remove(name);
}
Some(matched)
}
let loaded = match load_package(&lex_files, tmp.path(), &pkg_name) {
Ok(p) => p,
Err(e) => return error_response(400, format!("load package: {e}")),
};
let mut stages = canonicalize_program(&loaded.program);
if let Err(errs) = lex_types::check_and_rewrite_program(&mut stages) {
return error_with_detail(
422,
format!("type errors in package {pkg_name}"),
serde_json::to_value(&errs).unwrap(),
);
}
let new_fns: BTreeMap<String, lex_ast::FnDecl> = stages
.iter()
.filter_map(|s| match s {
lex_ast::Stage::FnDecl(fd) => Some((fd.name.clone(), fd.clone())),
_ => None,
})
.collect();
let all_function_names: Vec<String> = new_fns.keys().cloned().collect();
let mut old_fns: BTreeMap<String, lex_ast::FnDecl> = BTreeMap::new();
for (name, new_fd) in &new_fns {
if let Some(fd) = take_matching(&mut old_fns_by_name, name, new_fd) {
old_fns.insert(name.clone(), fd);
}
}
let report = lex_vcs::compute_diff(&old_fns, &new_fns, false);
let mut new_imports = lex_vcs::ImportMap::new();
for (file, modules) in &loaded.imports_by_file {
let entry = new_imports.entry(file.clone()).or_default();
for m in modules {
entry.insert(m.clone());
}
}
let outcome = match store.publish_program(&branch, &stages, &report, &new_imports, false) {
Ok(outcome) => outcome,
Err(lex_store::StoreError::TypeError(errs)) => {
return error_with_detail(422, "type errors", serde_json::to_value(&errs).unwrap());
}
Err(e) => return write_error_response("publish_program", e),
};
let all_ops: Vec<serde_json::Value> = match serde_json::to_value(&outcome.ops) {
Ok(serde_json::Value::Array(arr)) => arr,
_ => Vec::new(),
};
let final_head_op = outcome.head_op;
let now = SystemTime::now()
.duration_since(UNIX_EPOCH)
.map(|d| d.as_secs())
.unwrap_or(0);
let record = PkgRecord {
name: pkg_name.clone(),
version: pkg_version,
head_op: final_head_op.clone(),
published_at: now,
function_names: all_function_names,
ops: all_ops.clone(),
};
if let Err(e) = save_pkg_record(&state.root, &record, body) {
return error_response(500, format!("save package index: {e}"));
}
json_response(200, &serde_json::json!({
"package": pkg_name,
"ops": all_ops,
"head_op": final_head_op,
}))
}
fn pkg_list_handler(state: &State) -> Response<std::io::Cursor<Vec<u8>>> {
let names = list_pkg_names(&state.root);
let packages: Vec<serde_json::Value> = names.iter()
.filter_map(|name| {
let idx = load_pkg_index(&state.root, name)?;
let latest = idx.latest.as_deref()?;
let r = load_pkg_record(&state.root, name, latest)?;
Some(serde_json::json!({
"name": r.name,
"version": r.version,
"head_op": r.head_op,
"published_at": r.published_at,
}))
})
.collect();
json_response(200, &serde_json::json!({ "packages": packages }))
}
fn pkg_get_handler(state: &State, name: &str) -> Response<std::io::Cursor<Vec<u8>>> {
match load_latest_pkg_record(&state.root, name) {
Some(r) => json_response(200, &serde_json::json!({
"name": r.name,
"version": r.version,
"head_op": r.head_op,
"published_at": r.published_at,
"function_names": r.function_names,
"ops": r.ops,
})),
None => error_response(404, format!("package {name:?} not found")),
}
}
fn pkg_versions_handler(state: &State, name: &str) -> Response<std::io::Cursor<Vec<u8>>> {
match load_pkg_index(&state.root, name) {
Some(idx) => json_response(200, &serde_json::json!({
"name": name,
"latest": idx.latest,
"versions": idx.versions,
})),
None => error_response(404, format!("package {name:?} not found")),
}
}
fn pkg_get_version_handler(state: &State, name: &str, version: &str) -> Response<std::io::Cursor<Vec<u8>>> {
match load_pkg_record(&state.root, name, version) {
Some(r) => json_response(200, &serde_json::json!({
"name": r.name,
"version": r.version,
"head_op": r.head_op,
"published_at": r.published_at,
"function_names": r.function_names,
"ops": r.ops,
})),
None => error_response(404, format!("package {name:?}@{version:?} not found")),
}
}
fn pkg_archive_handler(state: &State, name: &str, version: &str) -> Response<std::io::Cursor<Vec<u8>>> {
let path = pkg_archive_path(&state.root, name, version);
match std::fs::read(&path) {
Ok(bytes) => Response::from_data(bytes)
.with_status_code(200)
.with_header(
tiny_http::Header::from_bytes(
&b"Content-Type"[..],
&b"application/gzip"[..],
)
.unwrap(),
),
Err(_) => error_response(404, format!("archive for {name:?}@{version:?} not found")),
}
}
fn pkg_head_handler(state: &State, name: &str) -> Response<std::io::Cursor<Vec<u8>>> {
match load_latest_pkg_record(&state.root, name) {
Some(r) => json_response(200, &serde_json::json!({
"name": r.name,
"version": r.version,
"head_op": r.head_op,
})),
None => error_response(404, format!("package {name:?} not found")),
}
}
fn pkg_delete_handler(state: &State, name: &str) -> Response<std::io::Cursor<Vec<u8>>> {
let record = match load_latest_pkg_record(&state.root, name) {
Some(r) => r,
None => return error_response(404, format!("package {name:?} not found")),
};
let store = state.store.lock().unwrap();
let branch = store.current_branch();
let head = match store.branch_head(&branch) {
Ok(h) => h,
Err(e) => return error_response(500, format!("branch_head: {e}")),
};
let head_pairs: Vec<(String, String)> = head
.iter()
.map(|(sig, stage)| (sig.clone(), stage.clone()))
.collect();
let old_fns: BTreeMap<String, lex_ast::FnDecl> = store
.get_asts_for_sigs_bulk(&head_pairs)
.into_iter()
.filter_map(|r| r.ok())
.filter_map(|s| match s {
lex_ast::Stage::FnDecl(fd)
if record.function_names.contains(&fd.name) => Some((fd.name.clone(), fd)),
_ => None,
})
.collect();
let new_fns: BTreeMap<String, lex_ast::FnDecl> = BTreeMap::new();
let report = lex_vcs::compute_diff(&old_fns, &new_fns, false);
let empty_imports = lex_vcs::ImportMap::new();
match store.publish_program(&branch, &[], &report, &empty_imports, false) {
Ok(outcome) => {
let ver = record.version.clone();
let _ = std::fs::remove_file(pkg_version_path(&state.root, name, &ver));
let _ = std::fs::remove_file(pkg_archive_path(&state.root, name, &ver));
if let Some(mut idx) = load_pkg_index(&state.root, name) {
idx.versions.retain(|v| v.version != ver);
idx.latest = idx.versions.last().map(|v| v.version.clone());
if idx.versions.is_empty() {
let _ = std::fs::remove_dir_all(pkg_name_dir(&state.root, name));
} else {
let bytes = serde_json::to_vec_pretty(&idx).unwrap_or_default();
let _ = std::fs::write(pkg_index_path(&state.root, name), bytes);
}
}
json_response(200, &serde_json::json!({
"deleted": name,
"version": ver,
"ops": outcome.ops,
"head_op": outcome.head_op,
}))
}
Err(lex_store::StoreError::TypeError(errs)) => {
error_with_detail(422, "type errors", serde_json::to_value(&errs).unwrap())
}
Err(e) => write_error_response("retract package", e),
}
}
#[cfg(test)]
mod policy_ceiling_tests {
use super::*;
use lex_runtime::Policy;
use std::path::PathBuf;
fn permissive_request() -> Policy {
Policy {
allow_effects: ["io", "fs_read", "fs_write", "net", "proc"]
.iter()
.map(|s| s.to_string())
.collect(),
allow_fs_read: vec![PathBuf::from("/")],
allow_fs_write: vec![PathBuf::from("/")],
allow_net_host: Vec::new(),
allow_proc: Vec::new(),
allow_approval: Vec::new(),
budget: None,
}
}
#[test]
fn ceiling_drops_effects_the_caller_was_not_granted() {
let ceiling = Policy {
allow_effects: ["io", "time"].iter().map(|s| s.to_string()).collect(),
..Policy::default()
};
let got = clamp_policy(permissive_request(), &ceiling);
assert!(got.allow_effects.contains("io"));
assert!(!got.allow_effects.contains("proc"), "proc must not survive a ceiling without it");
assert!(!got.allow_effects.contains("fs_write"));
assert!(!got.allow_effects.contains("net"));
assert!(!got.allow_effects.contains("time"));
}
#[test]
fn ceiling_scopes_override_caller_scopes() {
let ceiling = Policy {
allow_effects: ["fs_read"].iter().map(|s| s.to_string()).collect(),
allow_fs_read: vec![PathBuf::from("/srv/tenant")],
..Policy::default()
};
let got = clamp_policy(permissive_request(), &ceiling);
assert_eq!(got.allow_fs_read, vec![PathBuf::from("/srv/tenant")]);
assert!(got.allow_fs_write.is_empty());
assert!(got.allow_proc.is_empty());
assert!(got.allow_net_host.is_empty());
}
#[test]
fn ceiling_caps_budget_and_prefers_the_smaller() {
let mut req = permissive_request();
req.budget = None;
let ceiling = Policy { budget: Some(1_000), ..Policy::default() };
assert_eq!(clamp_policy(req, &ceiling).budget, Some(1_000));
let mut req2 = permissive_request();
req2.budget = Some(50);
let ceiling2 = Policy { budget: Some(1_000), ..Policy::default() };
assert_eq!(clamp_policy(req2, &ceiling2).budget, Some(50));
}
#[test]
fn empty_ceiling_is_pure_only() {
let got = clamp_policy(permissive_request(), &Policy::default());
assert!(got.allow_effects.is_empty(), "an empty ceiling grants nothing");
assert!(got.allow_proc.is_empty());
assert!(got.allow_fs_write.is_empty());
}
}
#[cfg(test)]
mod public_read_tests {
use super::*;
fn seed_pkg(root: &std::path::Path, name: &str, version: &str) {
let record = PkgRecord {
name: name.to_string(),
version: version.to_string(),
head_op: Some(format!("op-{name}")),
published_at: 1,
function_names: vec![format!("{name}.f")],
ops: vec![],
};
save_pkg_record(root, &record, format!("ARCHIVE:{name}@{version}").as_bytes())
.expect("seed package");
}
#[test]
fn new_package_defaults_to_private() {
let tmp = tempfile::TempDir::new().unwrap();
seed_pkg(tmp.path(), "lex-schema", "0.9.2");
assert!(!pkg_is_public(tmp.path(), "lex-schema"));
assert!(!pkg_is_public(tmp.path(), "does-not-exist"));
}
#[test]
fn set_visibility_round_trips_and_index_persists() {
let tmp = tempfile::TempDir::new().unwrap();
let state = State::open(tmp.path().to_path_buf()).unwrap();
seed_pkg(tmp.path(), "lex-schema", "0.9.2");
let _ = pkg_set_visibility_handler(&state, "lex-schema", r#"{"visibility":"public"}"#);
assert!(pkg_is_public(tmp.path(), "lex-schema"));
let idx = load_pkg_index(tmp.path(), "lex-schema").unwrap();
assert_eq!(idx.latest.as_deref(), Some("0.9.2"));
assert_eq!(idx.versions.len(), 1);
let _ = pkg_set_visibility_handler(&state, "lex-schema", r#"{"visibility":"private"}"#);
assert!(!pkg_is_public(tmp.path(), "lex-schema"));
}
#[test]
fn set_visibility_on_unknown_package_is_a_noop() {
let tmp = tempfile::TempDir::new().unwrap();
let state = State::open(tmp.path().to_path_buf()).unwrap();
let _ = pkg_set_visibility_handler(&state, "ghost", r#"{"visibility":"public"}"#);
assert!(load_pkg_index(tmp.path(), "ghost").is_none());
}
#[test]
fn public_listing_omits_private_packages() {
let tmp = tempfile::TempDir::new().unwrap();
let state = State::open(tmp.path().to_path_buf()).unwrap();
seed_pkg(tmp.path(), "pub-pkg", "1.0.0");
seed_pkg(tmp.path(), "priv-pkg", "1.0.0");
let _ = pkg_set_visibility_handler(&state, "pub-pkg", r#"{"visibility":"public"}"#);
let names = public_pkg_names(tmp.path());
assert_eq!(names, vec!["pub-pkg".to_string()]);
}
#[test]
fn resolve_public_maps_routes() {
let get = Method::Get;
assert_eq!(resolve_public(&get, "").unwrap(), PublicTarget::List);
assert_eq!(resolve_public(&get, "/").unwrap(), PublicTarget::List);
assert_eq!(
resolve_public(&get, "/lex-schema").unwrap(),
PublicTarget::Latest("lex-schema".into())
);
assert_eq!(
resolve_public(&get, "/lex-schema/versions").unwrap(),
PublicTarget::Versions("lex-schema".into())
);
assert_eq!(
resolve_public(&get, "/lex-schema/head").unwrap(),
PublicTarget::Head("lex-schema".into())
);
assert_eq!(
resolve_public(&get, "/lex-schema/0.9.2").unwrap(),
PublicTarget::Version("lex-schema".into(), "0.9.2".into())
);
assert_eq!(
resolve_public(&get, "/lex-schema/0.9.2/archive").unwrap(),
PublicTarget::Archive("lex-schema".into(), "0.9.2".into())
);
}
#[test]
fn resolve_public_rejects_bad_method_and_traversal() {
assert_eq!(resolve_public(&Method::Put, "/lex-schema"), Err(405));
assert_eq!(resolve_public(&Method::Post, "").err(), Some(405));
assert_eq!(resolve_public(&Method::Get, "/.."), Err(404));
assert_eq!(resolve_public(&Method::Get, "/lex-schema/../etc"), Err(404));
assert_eq!(resolve_public(&Method::Get, "/a/b/c/d"), Err(404));
assert!(resolve_public(&Method::Get, "/lex schema").is_err());
}
}