lex_store/store.rs
1//! `Store` — content-addressed code repository.
2//!
3//! The filesystem is the source of truth. All operations read/write JSON
4//! files under `<root>/stages/<SigId>/`. There is no SQLite cache: every
5//! query walks the directory and parses what's needed. `cargo test`
6//! runs aren't perf-critical and the §4.6 acceptance requires the
7//! rebuild-from-filesystem property anyway.
8
9use crate::branches::DEFAULT_BRANCH;
10use crate::model::*;
11use lex_ast::{sig_id, stage_id, Stage};
12use serde::de::DeserializeOwned;
13use serde::Serialize;
14use std::collections::{BTreeMap, BTreeSet};
15use std::fs;
16use std::path::{Path, PathBuf};
17use std::time::{SystemTime, UNIX_EPOCH};
18
19#[derive(Debug, thiserror::Error)]
20pub enum StoreError {
21 #[error("io error: {0}")]
22 Io(#[from] std::io::Error),
23 #[error("serialization error: {0}")]
24 Serde(#[from] serde_json::Error),
25 #[error("imports cannot be published as stages")]
26 CannotPublishImport,
27 #[error("unknown stage_id `{0}`")]
28 UnknownStage(String),
29 #[error("unknown sig_id `{0}`")]
30 UnknownSig(String),
31 /// A head names `(sig_id, stage_id)`, but the AST at `stage_id` is filed
32 /// under `filed_under` — the sig it actually hashes to — so no store can
33 /// ever hold the pair (#992). Distinct from [`Self::UnknownStage`]: the
34 /// content is present, it is just bound to the wrong identity. The shape
35 /// a pre-#992 `ChangeEffectSig` left behind (old sig, new stage).
36 #[error(
37 "unsatisfiable head entry: sig `{sig_id}` is bound to stage `{stage_id}`, \
38 but that stage is filed under sig `{filed_under}` (the sig its AST hashes to), \
39 so no store can hold the pair (#992)"
40 )]
41 UnsatisfiablePair { sig_id: String, stage_id: String, filed_under: String },
42 #[error("invalid lifecycle transition: {0}")]
43 InvalidTransition(String),
44 #[error("unknown branch `{0}`")]
45 UnknownBranch(String),
46 /// A branch-head advance (e.g. the ref half of `op push`) was
47 /// asked to move `branch` to `attempted`, but `attempted` is not a
48 /// descendant of the branch's `current` head — a non-fast-forward
49 /// that would orphan history. Refused, git-style, so a disjoint or
50 /// diverged push can't silently clobber a shared branch. The op
51 /// objects may already be present; only the ref is left unchanged.
52 #[error("non-fast-forward on `{branch}`: {attempted} is not a descendant of current head {current}")]
53 NonFastForward { branch: String, current: lex_vcs::OpId, attempted: lex_vcs::OpId },
54 #[error("unknown blob `{0}`")]
55 UnknownBlob(String),
56 #[error("unknown blob ref `{namespace}/{key}`")]
57 UnknownBlobRef { namespace: String, key: String },
58 /// A files manifest (#1007) failed to decode or validate.
59 #[error("invalid files manifest: {0}")]
60 InvalidManifest(crate::files::ManifestError),
61 /// A `SetFiles` (#1007) names blobs — its manifest or entries — this
62 /// store does not hold. Nothing was written; the head is unchanged.
63 #[error("files manifest references {} missing blob(s)", .0.len())]
64 MissingBlobs(Vec<crate::files::BlobId>),
65 /// The op exists but has nothing a regenerator could reproduce, by
66 /// kind (#1007) — a typed refusal, not a replay miss.
67 #[error("op {op_id} is not replayable: {why}")]
68 NotReplayable { op_id: lex_vcs::OpId, why: crate::files::NotReplayable },
69 #[error("unknown op_id `{0}`")]
70 UnknownOp(lex_vcs::OpId),
71 /// `replay_request` was asked to replay an op whose own produced stage
72 /// can't be loaded (#868). Unloadable *context* in the parent program is
73 /// skipped and reported instead (`ReplayRequest::skipped`), but the
74 /// target is what a regeneration is compared against — without it there
75 /// is nothing to replay.
76 #[error("cannot replay op {op_id}: its target stage {stage_id} cannot be loaded ({reason})")]
77 ReplayTargetUnloadable { op_id: String, stage_id: String, reason: String },
78 /// A head stage exists but could not be read — corrupt bytes, a
79 /// permission error, any I/O or parse failure other than "absent".
80 /// Lenient reconstruction (#868) skips only *absent* stages; this is
81 /// surfaced instead so store corruption is never hidden behind an
82 /// "incomplete context" report.
83 #[error("stage {stage_id} (sig {sig_id}) is present but unreadable: {reason}")]
84 StageUnreadable { sig_id: String, stage_id: String, reason: String },
85 /// A typed AST transform (#280) — e.g. `ReplaceMatchArm` — was
86 /// asked to operate on a node it couldn't address (wrong kind,
87 /// out-of-range arm index, unknown NodeId, etc.). Distinct from
88 /// `TypeError` (which means the transform succeeded but its
89 /// output didn't typecheck) so callers can render the right
90 /// error message.
91 #[error("transform failed: {0}")]
92 TransformError(lex_ast::TransformError),
93 #[error(transparent)]
94 Apply(#[from] lex_vcs::ApplyError),
95 /// The candidate program — i.e. the source the caller is
96 /// publishing — doesn't typecheck. The branch head is unchanged
97 /// and no op records are persisted. Issue #130's "always-valid
98 /// HEAD" invariant: the gate runs before any side effect, so a
99 /// type-broken publish leaves no footprint.
100 #[error("type errors in published program: {} error(s)", .0.len())]
101 TypeError(Vec<lex_types::TypeError>),
102 /// A dependency being resolved for the write-time gate (#930) is a
103 /// multi-module package. Per-module signature extraction (picking the
104 /// imported module's file out of the de-flattened tree) is not yet
105 /// implemented; single-module (leaf) dependencies resolve today. A
106 /// caller can treat this as "cannot resolve here" rather than a hard
107 /// failure.
108 #[error("multi-module dependency resolution is not yet supported")]
109 UnsupportedMultiModuleDependency,
110 /// A merge was refused because its two parents pin the **same**
111 /// dependency at **different** versions (#977). The merge commit records
112 /// the union of both parents' lock entries as its own lock; picking one of
113 /// two conflicting pins silently would change what the merged code was
114 /// tested against, so the merge is refused instead. Nothing moves —
115 /// neither branch advances. Resolve by aligning the version on one branch
116 /// (e.g. `lex pkg update` there, then push) and merging again.
117 #[error(
118 "dependency conflict: `{package}` is pinned at {dst_version} on the \
119 destination branch but {src_version} on the source branch; align the \
120 version on one branch and merge again"
121 )]
122 DependencyConflict { package: String, dst_version: String, src_version: String },
123 /// A merge's file-manifest diff (#1007 PR 7) was asked to 3-way-merge
124 /// a side whose own `manifest_at` is already `Ambiguous` — a merge of
125 /// merges where an earlier merge landed without the `SetFiles` it
126 /// should have appended. There is no well-defined manifest on that
127 /// side to diff against, so this merge is refused too rather than
128 /// guessing. Resolve by appending a `SetFiles` on `op_id` first (the
129 /// same fix `check_head_files` names for a plain head advance).
130 #[error(
131 "ambiguous manifest at `{op_id}`: append a set_files op resolving it before merging"
132 )]
133 AmbiguousManifest { op_id: lex_vcs::OpId },
134 /// A typed issue's example targets a function that isn't declared at
135 /// the head being evaluated (#949) — the example cannot be attached, so
136 /// the issue cannot be judged there.
137 #[error("issue example targets `{0}`, which is not declared at this head")]
138 IssueTarget(String),
139 /// A proposed acceptance was refused for the issue it names (#956) —
140 /// e.g. the issue is already typed, so it is not open to refinement.
141 #[error("{0}")]
142 IssueRefinement(String),
143 /// The op was persisted but a `required_attestations` rule in
144 /// `policy.json` (#245) refused to advance the branch head past
145 /// it. The op record is durable — re-running with the missing
146 /// attestations recorded will succeed without re-persisting —
147 /// but the branch is unchanged. Surfaced as a structured JSON
148 /// envelope on the HTTP API.
149 #[error(
150 "branch advance blocked: op {} missing attestations: {}",
151 .0.op_id, .0.missing.join(", ")
152 )]
153 BranchAdvanceBlocked(crate::policy::BranchAdvanceBlocked),
154 /// All retry attempts of the CAS branch-head advance failed
155 /// because another writer kept advancing the same branch
156 /// (#262). The op record itself is durable in the op log
157 /// (orphaned), so re-running with backoff would eventually
158 /// land — return `503 Contention { retry_after }` from the
159 /// HTTP API and let the client back off.
160 #[error("branch advance contention on `{branch}`: {attempts} retries exhausted")]
161 Contention { branch: String, attempts: u32 },
162 /// The op was persisted but its stage carries an attestation
163 /// produced by a retroactively quarantined tool (#248). The
164 /// branch head is unchanged. The op record stays in the log
165 /// (audit trail intact); re-running with the producer
166 /// unblocked, or with un-contaminated attestations, succeeds
167 /// without re-persisting the op.
168 #[error(
169 "branch advance blocked: op {} touches stage {} with an attestation from \
170 quarantined producer `{}` (blocked at {}, attestation at {})",
171 .0.op_id, .0.stage_id, .0.tool_id, .0.blocked_at, .0.attestation_at
172 )]
173 ProducerBlocked(crate::policy::ProducerBlocked),
174 /// The op would push its session's monotonic budget over the
175 /// cap configured in `policy.session_budgets` (#292 slice 3).
176 /// The op is *not* persisted; the branch head is unchanged.
177 /// The caller should either start a new session, raise the
178 /// cap, or refactor to fit the budget. HTTP API maps to 503.
179 #[error("session `{session_id}` budget exceeded: spent_after={spent_after} > cap={cap}")]
180 BudgetExceeded {
181 session_id: String,
182 cap: u64,
183 spent_after: u64,
184 },
185}
186
187/// The outcome returned by [`Store::publish_program`].
188#[derive(Debug, Clone, serde::Serialize)]
189pub struct PublishOutcome {
190 pub ops: Vec<PublishOp>,
191 pub head_op: Option<lex_vcs::OpId>,
192}
193
194/// Everything a regenerator needs to *replay* an op (#836 G3), produced
195/// by [`Store::replay_request`]. The model call is external: a harness
196/// feeds `prompt` + `parent_program` to `model`, then hands the
197/// regenerated stage to [`Store::replay_compare`].
198#[derive(Debug, Clone, serde::Serialize)]
199pub struct ReplayRequest {
200 pub op_id: String,
201 /// The sig the op changed — the function to regenerate.
202 pub target_sig: String,
203 /// The target function's name (the recorded stage is a function
204 /// for every replayable op).
205 #[serde(default, skip_serializing_if = "Option::is_none")]
206 pub target_name: Option<String>,
207 /// The target function's rendered signature (`fn name(...) -> T`),
208 /// so a regenerator knows the interface to implement without
209 /// re-deriving it from the sig hash.
210 #[serde(default, skip_serializing_if = "Option::is_none")]
211 pub target_signature: Option<String>,
212 /// The stage id a faithful regeneration should reproduce.
213 pub expected_stage_id: String,
214 /// The recorded intent prompt (`None` if the op carried no intent).
215 pub prompt: Option<String>,
216 /// The recorded model (`provider/name[@version]`), if any.
217 pub model: Option<String>,
218 /// The recorded session id, if any.
219 pub session_id: Option<String>,
220 /// The program the change was made against — the parent state
221 /// rendered to source — the context a regenerator needs.
222 pub parent_program: String,
223 /// Declarations the parent state names but the store could not load
224 /// (#868) — GC'd, superseded or never-persisted stages. They are left out
225 /// of `parent_program` rather than failing the replay; an entry with
226 /// `called_by_target` means the regenerator is looking at a program with
227 /// a dangling reference. Omitted from JSON when empty.
228 #[serde(default, skip_serializing_if = "Vec::is_empty")]
229 pub skipped: Vec<SkippedStage>,
230}
231
232impl ReplayRequest {
233 /// Whether the parent program shown to the regenerator is missing
234 /// declarations (#868) — any skip at all makes the context incomplete.
235 pub fn context_incomplete(&self) -> bool {
236 !self.skipped.is_empty()
237 }
238
239 /// A one-line human note describing the incomplete context, suitable
240 /// for appending to a replay verdict's detail. `None` when nothing was
241 /// skipped.
242 pub fn context_note(&self) -> Option<String> {
243 if self.skipped.is_empty() {
244 return None;
245 }
246 let called: Vec<&str> = self
247 .skipped
248 .iter()
249 .filter(|s| s.called_by_target)
250 .map(|s| s.name.as_deref().unwrap_or(s.sig_id.as_str()))
251 .collect();
252 let mut note = format!(
253 "replay context incomplete: {} declaration(s) of the parent program could not be loaded",
254 self.skipped.len()
255 );
256 if !called.is_empty() {
257 note.push_str(&format!(", including {} called by the target", called.join(", ")));
258 }
259 Some(note)
260 }
261}
262
263/// A head declaration left out of a reconstructed program because its stage
264/// could not be loaded (#868).
265#[derive(Debug, Clone, PartialEq, serde::Serialize, serde::Deserialize)]
266pub struct SkippedStage {
267 pub sig_id: String,
268 pub stage_id: String,
269 /// The declaration's (de-mangled) name, when recoverable from the
270 /// stage's metadata or another stage of the same sig.
271 #[serde(default, skip_serializing_if = "Option::is_none")]
272 pub name: Option<String>,
273 /// Why it could not be loaded (the underlying store error).
274 pub reason: String,
275 /// Whether the replay target references this declaration by name — i.e.
276 /// the regenerator's context has a dangling reference. Conservative (a
277 /// name match anywhere in the target counts) and `false` when the name
278 /// isn't recoverable. Only set by `replay_request`.
279 #[serde(default)]
280 pub called_by_target: bool,
281}
282
283/// A program reconstructed from the op-log with unloadable declarations
284/// skipped rather than failing the whole reconstruction (#868).
285#[derive(Debug, Clone)]
286pub struct ReconstructedProgram {
287 pub stages: Vec<Stage>,
288 pub skipped: Vec<SkippedStage>,
289}
290
291/// The result of comparing a regenerated candidate against an op's
292/// recorded output (#836 G3), returned by [`Store::replay_compare`]
293/// after it emits the `Replay` attestation.
294#[derive(Debug, Clone, serde::Serialize)]
295pub struct ReplayOutcome {
296 pub op_id: String,
297 pub expected_stage_id: String,
298 /// The candidate's stage id when it regenerated the same sig, else
299 /// `None`.
300 pub produced_stage_id: Option<String>,
301 /// Whether the regeneration reproduced the recorded change (exact or
302 /// behavioral).
303 pub reproduced: bool,
304 /// Set when reproduction was behavioral (same values over N sampled
305 /// inputs) rather than an exact stage-id match. `None` for an exact
306 /// match or a genuine miss.
307 #[serde(default, skip_serializing_if = "Option::is_none")]
308 pub behavioral_samples: Option<usize>,
309 /// The id of the `Replay` attestation this comparison emitted.
310 pub attestation_id: String,
311 /// Set when the replay ran against an incomplete parent program —
312 /// declarations the store could not load were skipped (#868). A
313 /// not-reproduced verdict under incomplete context is weaker evidence
314 /// than one under full context. Omitted from JSON when false.
315 #[serde(default, skip_serializing_if = "std::ops::Not::not")]
316 pub context_incomplete: bool,
317 /// The skipped declarations (see [`ReplayRequest::skipped`]). Omitted
318 /// from JSON when empty.
319 #[serde(default, skip_serializing_if = "Vec::is_empty")]
320 pub skipped: Vec<SkippedStage>,
321}
322
323impl ReplayOutcome {
324 /// Carry a request's skipped-context report onto its verdict (#868).
325 pub fn with_context_of(mut self, req: &ReplayRequest) -> Self {
326 self.context_incomplete = req.context_incomplete();
327 self.skipped = req.skipped.clone();
328 self
329 }
330}
331
332/// One applied operation within a [`PublishOutcome`].
333#[derive(Debug, Clone, serde::Serialize)]
334pub struct PublishOp {
335 pub op_id: lex_vcs::OpId,
336 pub kind: serde_json::Value,
337}
338
339/// One entry in the per-`SigId` stage history surfaced by
340/// `Store::sig_history`. Newest-first ordering is the responsibility
341/// of the producer.
342#[derive(Debug, Clone, serde::Serialize, PartialEq)]
343pub struct StageHistoryEntry {
344 pub stage_id: String,
345 pub status: StageStatus,
346 /// Wall-clock seconds of the most recent transition.
347 pub last_at: u64,
348 /// Wall-clock seconds when this stage was first written to the
349 /// store (its initial Draft transition). `None` for stages
350 /// whose lifecycle log doesn't include an explicit Draft entry
351 /// — shouldn't happen for stages published via `Store::publish`,
352 /// but the type allows hand-edited stores.
353 #[serde(skip_serializing_if = "Option::is_none")]
354 pub published_at: Option<u64>,
355}
356
357/// Per-candidate metadata surfaced by [`Store::list_candidates`]
358/// (#294). Returned sorted by `op_id` for deterministic output.
359#[derive(Debug, Clone, serde::Serialize, serde::Deserialize, PartialEq, Eq)]
360pub struct CandidateInfo {
361 pub op_id: lex_vcs::OpId,
362 pub stage_id: lex_vcs::StageId,
363 /// Author intent. Always set for `Candidate` ops emitted via
364 /// [`Store::propose_candidate`]; `None` only if a
365 /// hand-written raw op skipped the intent tag.
366 pub intent_id: Option<lex_vcs::IntentId>,
367}
368
369/// One line of `stage_index.jsonl`. See `Store::lookup_lifecycle`.
370#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
371struct StageIndexEntry {
372 stage_id: String,
373 sig_id: String,
374}
375
376/// Sentinel `sig_id` value recording "a full scan already established
377/// this stage_id exists nowhere in the store" (#825). Never a real
378/// sig — sig directory names are never empty.
379const MISSING_STAGE_MARKER: &str = "";
380
381/// Resolves a head's external (registry/git) dependencies to their public
382/// module signatures, so the write-time gate can type-check a head that keeps
383/// `import "<pkg>/mod" as <alias>` edges instead of inlining the dependency
384/// (#930). The returned map is keyed by import *reference* (`"lex-nt/lib"`) →
385/// that dependency module's record type (as [`crate::render::module_record_at_op`]
386/// or a source-based equivalent produces), exactly the shape
387/// [`lex_types::check_program_with_modules`] consumes.
388///
389/// Implementations differ by context and live in the crate that has the
390/// resolution machinery: the client (`lex publish`) resolves from the
391/// working-copy `lex.lock` + local package cache; the hub resolves from the
392/// committed lock + its own hosted stores (cross-tenant), which a single
393/// [`Store`] cannot reach on its own. When no resolver is installed the gate
394/// resolves an empty map — only stdlib binds and any external reference is an
395/// unbound-name error, exactly as before #930 (so inlined heads, which carry
396/// no external edges, are unaffected).
397pub trait DepResolver: Send + Sync {
398 /// `head_op` is the op the head is known by when the gate has one (merge,
399 /// patch and hub-verify reconstruct a committed head); `None` for a
400 /// candidate not yet committed (`publish`), where the client resolver
401 /// falls back to the working-copy lock.
402 fn resolve_modules(
403 &self,
404 stages: &[Stage],
405 head_op: Option<&str>,
406 ) -> BTreeMap<String, lex_types::Ty>;
407
408 /// The resolved dependencies' exported **type declarations**, keyed by the
409 /// same import reference as [`Self::resolve_modules`] (#930 completeness).
410 /// The write-time gate registers these under each import's alias so a head
411 /// that references a dependency's *type* — not just its functions — checks.
412 ///
413 /// Defaulted to empty so existing resolvers keep compiling and simply
414 /// resolve function signatures only (their previous behavior); a resolver
415 /// that can supply type declarations overrides this.
416 fn resolve_module_types(
417 &self,
418 _stages: &[Stage],
419 _head_op: Option<&str>,
420 ) -> BTreeMap<String, Vec<lex_ast::TypeDecl>> {
421 BTreeMap::new()
422 }
423
424 /// Each resolved dependency import's **module mangle prefix**, keyed by the
425 /// same import reference (#963). When a dependency is resolved as a whole
426 /// package, its exported types are named by their path-derived prefix
427 /// (`error_<hash>.DbErr`); this lets the write-time gate map the import
428 /// alias to that prefix so an alias-qualified type reference (`e.DbErr`)
429 /// resolves to the same type the dependency's own signatures name.
430 ///
431 /// Defaulted to empty: a resolver that returns bare (alias-qualified) type
432 /// names supplies none, and the gate keeps the #930 alias-qualification
433 /// path.
434 fn resolve_module_prefixes(
435 &self,
436 _stages: &[Stage],
437 _head_op: Option<&str>,
438 ) -> BTreeMap<String, String> {
439 BTreeMap::new()
440 }
441
442 /// Everything the gate needs in one call (#943/#944): the three maps
443 /// above plus structured **diagnostics** for imports that could not be
444 /// resolved (`unpinned_dependency`, `unresolved_dependency`,
445 /// `dependency_conflict`). The write-time gate calls only this.
446 ///
447 /// Defaulted to the three methods above with no diagnostics, so existing
448 /// resolvers behave exactly as before. A resolver that resolves
449 /// recursively overrides this (doing the work once) and may implement the
450 /// three methods in terms of it.
451 fn resolve(&self, stages: &[Stage], head_op: Option<&str>) -> crate::deps::ResolvedDeps {
452 crate::deps::ResolvedDeps {
453 modules: self.resolve_modules(stages, head_op),
454 types: self.resolve_module_types(stages, head_op),
455 prefixes: self.resolve_module_prefixes(stages, head_op),
456 diagnostics: Vec::new(),
457 }
458 }
459}
460
461/// Whether a dependency diagnostic is about `import "<reference>"`.
462fn names_import(d: &lex_types::TypeError, reference: &str) -> bool {
463 match d {
464 lex_types::TypeError::UnpinnedDependency { reference: r, .. }
465 | lex_types::TypeError::UnresolvedDependency { reference: r, .. } => r == reference,
466 _ => false,
467 }
468}
469
470pub struct Store {
471 root: PathBuf,
472 /// Optional dependency resolver (#930). Injected by the client or the hub;
473 /// `None` in a bare store means the gate resolves an empty module map.
474 resolver: Option<std::sync::Arc<dyn DepResolver>>,
475}
476
477impl Store {
478 /// Open or create a store rooted at `root`.
479 pub fn open(root: impl AsRef<Path>) -> Result<Self, StoreError> {
480 let root = root.as_ref().to_path_buf();
481 fs::create_dir_all(root.join("stages"))?;
482 fs::create_dir_all(root.join("traces"))?;
483 let store = Self { root, resolver: None };
484 store.ensure_stage_index();
485 Ok(store)
486 }
487
488 /// Install the dependency resolver the write-time gate uses to type-check
489 /// heads that keep external `import` edges (#930). Builder-style so a
490 /// caller can write `Store::open(p)?.with_dep_resolver(r)`.
491 pub fn with_dep_resolver(mut self, resolver: std::sync::Arc<dyn DepResolver>) -> Self {
492 self.resolver = Some(resolver);
493 self
494 }
495
496 /// Set the dependency resolver in place (for a `Store` already owned, e.g.
497 /// behind a `Mutex` in the HTTP `State`).
498 pub fn set_dep_resolver(&mut self, resolver: std::sync::Arc<dyn DepResolver>) {
499 self.resolver = Some(resolver);
500 }
501
502 /// The resolved dependencies for a head being gated — the installed
503 /// resolver's answer, or nothing when none is installed (#930): then only
504 /// stdlib binds and an external reference is an unbound-name error.
505 fn resolved_deps(&self, stages: &[Stage], head_op: Option<&str>) -> crate::deps::ResolvedDeps {
506 match &self.resolver {
507 Some(r) => r.resolve(stages, head_op),
508 None => crate::deps::ResolvedDeps::default(),
509 }
510 }
511
512 /// Type-check a gated head against its resolved dependencies (#930). A
513 /// dependency that did not resolve is reported by its structured
514 /// diagnostic (#943/#944) — ahead of, and alongside, whatever the checker
515 /// then finds — and fails the gate even if the checker would not: an
516 /// import the gate cannot verify is not a verified head.
517 fn check_with_resolved_deps(
518 &self,
519 stages: &[Stage],
520 head_op: Option<&str>,
521 ) -> Result<(), Vec<lex_types::TypeError>> {
522 let deps = self.resolved_deps(stages, head_op);
523 let checked = lex_types::check_program_with_deps(stages, &deps.modules, &deps.types, &deps.prefixes);
524 match (deps.diagnostics.is_empty(), checked) {
525 (true, Ok(_)) => Ok(()),
526 (true, Err(errors)) => Err(errors),
527 (false, result) => {
528 // The checker's `unknown_identifier` on an unresolved import's
529 // alias is the vague symptom the diagnostic replaces; keep
530 // every other checker error.
531 let failed_aliases: std::collections::BTreeSet<&str> = stages
532 .iter()
533 .filter_map(|s| match s {
534 Stage::Import(i) if deps.diagnostics.iter().any(|d| names_import(d, &i.reference)) => {
535 Some(i.alias.as_str())
536 }
537 _ => None,
538 })
539 .collect();
540 let mut more: Vec<lex_types::TypeError> = result.err().unwrap_or_default();
541 more.retain(|e| match e {
542 lex_types::TypeError::UnknownIdentifier { name, .. } => {
543 let head = name.split('.').next().unwrap_or(name);
544 !failed_aliases.contains(head)
545 }
546 _ => true,
547 });
548 let mut errors = deps.diagnostics;
549 errors.extend(more);
550 Err(errors)
551 }
552 }
553 }
554
555 /// The `import` edges of the head at `head_op`, as `Stage::Import`s (#930).
556 ///
557 /// A SigId→stage map holds only fn/type declarations; a head's imports are
558 /// `AddImport` ops and are absent from it. Every gate that type-checks a
559 /// reconstructed head must prepend these or a non-inlined head's
560 /// `<alias>.name` references fail as `unknown_identifier` even when the
561 /// dependency resolved correctly — the resolver scans these imports to find
562 /// each dependency, and the checker's Pass 1 binds the alias to the
563 /// resolved module. Returns empty when the head can't be read (an inlined
564 /// or import-free head needs nothing).
565 ///
566 /// Shared by all the write-path gates so they agree (#945).
567 fn head_import_stages(&self, head_op: &str) -> Vec<Stage> {
568 crate::render::package_head_at_op(self, head_op)
569 .map(|ph| {
570 ph.flat_imports
571 .into_iter()
572 .map(|(reference, alias)| Stage::Import(lex_ast::Import { reference, alias }))
573 .collect::<Vec<_>>()
574 })
575 .unwrap_or_default()
576 }
577
578 /// Prepend the head's import edges to `decls`, skipping any import the
579 /// program already carries (so a caller that already has its own imports —
580 /// a publish of full file contents, say — doesn't get duplicates).
581 ///
582 /// The gate helper: `stages = store.with_head_imports(head_op, decls)`.
583 fn with_head_imports(&self, head_op: Option<&str>, decls: Vec<Stage>) -> Vec<Stage> {
584 let Some(head) = head_op else { return decls };
585 let present: std::collections::BTreeSet<(String, String)> = decls
586 .iter()
587 .filter_map(|s| match s {
588 Stage::Import(i) => Some((i.reference.clone(), i.alias.clone())),
589 _ => None,
590 })
591 .collect();
592 let mut stages: Vec<Stage> = self
593 .head_import_stages(head)
594 .into_iter()
595 .filter(|s| match s {
596 Stage::Import(i) => !present.contains(&(i.reference.clone(), i.alias.clone())),
597 _ => true,
598 })
599 .collect();
600 stages.extend(decls);
601 stages
602 }
603
604 /// One-time migration for a store that predates the reverse
605 /// index (#822), or whose previous rebuild pass never finished
606 /// (e.g. the process was killed or its client disconnected
607 /// mid-request — server-side work keeps running either way, but
608 /// a *restart* genuinely stops it): build the index in a single
609 /// pass instead of leaving every subsequent `lookup_lifecycle`
610 /// call to discover its own entry via the slow per-call scan-
611 /// and-backfill fallback.
612 ///
613 /// That per-call fallback is fine for the rare individual miss
614 /// it was designed for, but pathological as a *bulk* cold-start
615 /// strategy: on a tenant with a few thousand functions it means
616 /// redoing an O(total sigs) scan from scratch for *each* of a
617 /// few thousand cold entries — O(total sigs²) — which measured
618 /// as a near-stall (page-cache thrashing) on a memory-
619 /// constrained host. A single pass over `list_sigs()` is
620 /// O(total sigs) total.
621 ///
622 /// Gated on a dedicated completion marker
623 /// (`stage_index.complete`), NOT on `stage_index.jsonl`'s mere
624 /// existence — a partially-built index file (left behind by an
625 /// interrupted rebuild, lazy or bulk) must still trigger a
626 /// re-run so the remaining entries get backfilled in one more
627 /// cheap O(total sigs) pass, not silently be mistaken for
628 /// "already done" and fall back to the slow per-call path for
629 /// whatever's left. `rebuild_stage_index` already skips entries
630 /// it finds present, so re-running it against a partial index
631 /// only does the work that remains. The marker is written only
632 /// after a full pass returns `Ok`, so a failed pass (e.g. an I/O
633 /// error partway through `list_sigs`) is retried on the next
634 /// open rather than being marked done.
635 ///
636 /// Runs once per `Store::open` call — which, in a long-lived
637 /// server (lex-hub caches one `Store` per tenant for the life of
638 /// the process), means once per tenant per process lifetime, not
639 /// once per request. Once the marker exists (the steady state
640 /// after the first successful run on any given host) this is a
641 /// single cheap file-existence check. Best-effort like the rest
642 /// of the index: any failure here just leaves the slower per-call
643 /// fallback as the only path, never breaks correctness.
644 fn ensure_stage_index(&self) {
645 if self.stage_index_complete_marker_path().exists() {
646 return;
647 }
648 if self.rebuild_stage_index().is_ok() {
649 let _ = fs::write(self.stage_index_complete_marker_path(), "");
650 }
651 }
652
653 fn stage_index_complete_marker_path(&self) -> PathBuf {
654 self.root.join("stage_index.complete")
655 }
656
657 /// Build (or top up) the reverse index in one pass over every
658 /// SigId in the store, rather than relying on `lookup_lifecycle`
659 /// to discover entries one at a time. Safe to call at any time,
660 /// including on a partially-built index (e.g. one left behind by
661 /// an interrupted request that was populating it lazily): already-
662 /// indexed stage_ids are skipped, so this only does the work that
663 /// remains. Returns the number of newly-added entries.
664 pub fn rebuild_stage_index(&self) -> Result<usize, StoreError> {
665 // A sig's lifecycle can list the same stage_id more than once
666 // (Draft, then later Active, then Deprecated all carry the
667 // same stage_id with a different status) -- track newly-seen
668 // keys locally too, not just what was already on disk at the
669 // start, so a repeated stage_id within one sig's transitions
670 // doesn't get appended to the index more than once.
671 let mut existing = self.load_stage_index();
672 let mut added = 0usize;
673 for sig in self.list_sigs()? {
674 let Ok(life) = self.read_lifecycle(&sig) else { continue };
675 for t in &life.transitions {
676 if !existing.contains_key(&t.stage_id) {
677 self.append_stage_index_entry(&t.stage_id, &sig);
678 existing.insert(t.stage_id.clone(), sig.clone());
679 added += 1;
680 }
681 }
682 }
683 Ok(added)
684 }
685
686 pub fn root(&self) -> &Path {
687 &self.root
688 }
689
690 // ── Generic content-addressed blobs (#5 / M6.1) ──────────────────────────
691 //
692 // The stage store holds typed Lex ASTs; loom-style artifacts (generated
693 // code, JSON, prose) are opaque text. These blob methods give the store a
694 // generic content-addressed object alongside stages, plus a lightweight
695 // ref namespace so callers can bind names (e.g. a sprint's node ids) to
696 // blob shas without touching the operation-log branch machinery.
697 //
698 // The sha is the lowercase hex SHA-256 of the content's UTF-8 bytes —
699 // identical to Lex's `crypto.sha256_str`, so a blob written here and an
700 // artifact content-addressed in loom's SQLite store share the same id and
701 // are interchangeable by reference. Store-scoped, so under lex-hub each
702 // tenant store gets its own blob space for free.
703
704 fn blobs_dir(&self) -> PathBuf {
705 self.root.join("blobs")
706 }
707
708 fn blob_refs_dir(&self) -> PathBuf {
709 self.root.join("blobrefs")
710 }
711
712 /// Content-address `content` and persist it under `<root>/blobs/<sha>`.
713 /// Returns the sha. Text wrapper over [`Self::put_blob_bytes`]; the id is
714 /// the sha of the UTF-8 bytes, so it is unchanged from before #1007.
715 pub fn put_blob(&self, content: &str) -> Result<String, StoreError> {
716 self.put_blob_bytes(content.as_bytes())
717 }
718
719 /// Content-address arbitrary `bytes` (#1007: files beside the op-log may
720 /// be binary) and persist them under `<root>/blobs/<sha>`. Returns the
721 /// lowercase hex SHA-256 — the same value `sha256sum` prints. Idempotent:
722 /// re-putting identical content is a no-op. Concurrency-safe — writes to a
723 /// unique temp file then atomically renames onto the content-addressed
724 /// path, so parallel writers of the same content can't corrupt it.
725 pub fn put_blob_bytes(&self, bytes: &[u8]) -> Result<String, StoreError> {
726 use sha2::{Digest, Sha256};
727 let sha = hex::encode(Sha256::digest(bytes));
728 let dir = self.blobs_dir();
729 let path = dir.join(&sha);
730 if path.exists() {
731 return Ok(sha);
732 }
733 fs::create_dir_all(&dir)?;
734 static SEQ: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(0);
735 let n = SEQ.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
736 let tmp = dir.join(format!(".{sha}.{}.{n}.tmp", std::process::id()));
737 fs::write(&tmp, bytes)?;
738 // rename is atomic on the same filesystem; identical content makes a
739 // last-writer-wins race harmless.
740 fs::rename(&tmp, &path)?;
741 Ok(sha)
742 }
743
744 /// Read a blob by its sha as text. `UnknownBlob` if absent; an
745 /// `InvalidData` I/O error if the blob is not UTF-8 (use
746 /// [`Self::get_blob_bytes`] for binary content).
747 pub fn get_blob(&self, sha: &str) -> Result<String, StoreError> {
748 let bytes = self.get_blob_bytes(sha)?;
749 String::from_utf8(bytes).map_err(|e| {
750 StoreError::Io(std::io::Error::new(std::io::ErrorKind::InvalidData, e))
751 })
752 }
753
754 /// Read a blob's exact bytes by its sha. `UnknownBlob` if absent.
755 pub fn get_blob_bytes(&self, sha: &str) -> Result<Vec<u8>, StoreError> {
756 // A blob id is 64 lowercase hex chars; anything else (e.g. `../x`)
757 // cannot name a blob and must never be joined onto the blobs dir.
758 if !crate::files::is_blob_id(sha) {
759 return Err(StoreError::UnknownBlob(sha.to_string()));
760 }
761 match fs::read(self.blobs_dir().join(sha)) {
762 Ok(b) => Ok(b),
763 Err(e) if e.kind() == std::io::ErrorKind::NotFound => {
764 Err(StoreError::UnknownBlob(sha.to_string()))
765 }
766 Err(e) => Err(StoreError::Io(e)),
767 }
768 }
769
770 /// Byte length of a blob without reading it, or `None` if absent.
771 pub fn blob_len(&self, sha: &str) -> Option<u64> {
772 if !crate::files::is_blob_id(sha) {
773 return None;
774 }
775 fs::metadata(self.blobs_dir().join(sha)).ok().map(|m| m.len())
776 }
777
778 /// Total bytes held in this store's blob space (#1007): the sum of every
779 /// stored blob's length, ignoring in-flight temp files. What a per-store
780 /// blob quota is measured against. O(number of blobs).
781 pub fn blob_bytes_used(&self) -> Result<u64, StoreError> {
782 let rd = match fs::read_dir(self.blobs_dir()) {
783 Ok(rd) => rd,
784 Err(e) if e.kind() == std::io::ErrorKind::NotFound => return Ok(0),
785 Err(e) => return Err(e.into()),
786 };
787 let mut total = 0u64;
788 for ent in rd {
789 let ent = ent?;
790 if !crate::files::is_blob_id(&ent.file_name().to_string_lossy()) {
791 continue;
792 }
793 total = total.saturating_add(ent.metadata()?.len());
794 }
795 Ok(total)
796 }
797
798 /// Whether a blob with this sha exists.
799 pub fn has_blob(&self, sha: &str) -> bool {
800 crate::files::is_blob_id(sha) && self.blobs_dir().join(sha).exists()
801 }
802
803 /// Every blob id currently on disk, paired with its file's last
804 /// modification time (#1007 PR 7 blob GC) — the age a blob's grace
805 /// period is measured against. `.tmp` write-in-progress files (see
806 /// [`Self::put_blob_bytes`]) never match `is_blob_id` and are
807 /// skipped, so a concurrent writer can't have its temp file swept.
808 pub(crate) fn list_blob_ids_with_mtime(&self) -> Result<Vec<(String, SystemTime)>, StoreError> {
809 let rd = match fs::read_dir(self.blobs_dir()) {
810 Ok(rd) => rd,
811 Err(e) if e.kind() == std::io::ErrorKind::NotFound => return Ok(Vec::new()),
812 Err(e) => return Err(e.into()),
813 };
814 let mut out = Vec::new();
815 for ent in rd {
816 let ent = ent?;
817 let name = ent.file_name().to_string_lossy().to_string();
818 if !crate::files::is_blob_id(&name) {
819 continue;
820 }
821 let mtime = ent.metadata()?.modified()?;
822 out.push((name, mtime));
823 }
824 Ok(out)
825 }
826
827 /// Delete a blob by id (#1007 PR 7 blob GC). No-op if the id doesn't
828 /// look like a blob id or the file is already gone — GC is expected
829 /// to run concurrently with itself across replicas without erroring.
830 pub(crate) fn delete_blob(&self, id: &str) -> Result<(), StoreError> {
831 if !crate::files::is_blob_id(id) {
832 return Ok(());
833 }
834 match fs::remove_file(self.blobs_dir().join(id)) {
835 Ok(()) => Ok(()),
836 Err(e) if e.kind() == std::io::ErrorKind::NotFound => Ok(()),
837 Err(e) => Err(e.into()),
838 }
839 }
840
841 /// Every blob sha bound anywhere under `blobrefs/**` (#1007 §3 blob
842 /// GC: "PLUS everything under the existing `blobrefs/**` namespace").
843 /// Walks the whole tree regardless of namespace depth or shape (locks
844 /// keyed by op id, loom artifacts keyed by node id, …) so a new
845 /// namespace never needs a GC update to stay safe — the cost of a
846 /// false "still live" is a few retained blobs, not a correctness bug.
847 pub(crate) fn all_blob_ref_shas(&self) -> Result<BTreeSet<String>, StoreError> {
848 let mut out = BTreeSet::new();
849 let mut stack = vec![self.blob_refs_dir()];
850 while let Some(dir) = stack.pop() {
851 let rd = match fs::read_dir(&dir) {
852 Ok(rd) => rd,
853 Err(e) if e.kind() == std::io::ErrorKind::NotFound => continue,
854 Err(e) => return Err(e.into()),
855 };
856 for ent in rd {
857 let ent = ent?;
858 let ft = ent.file_type()?;
859 if ft.is_dir() {
860 stack.push(ent.path());
861 } else if ft.is_file() {
862 if let Ok(sha) = fs::read_to_string(ent.path()) {
863 out.insert(sha.trim().to_string());
864 }
865 }
866 }
867 }
868 Ok(out)
869 }
870
871 /// Bind `key` to a blob `sha` within `namespace` (e.g. namespace
872 /// `"loom/sprint-abc"`, key `"build-node"`). Overwrites an existing
873 /// binding. The namespace may contain `/`; neither namespace nor key may
874 /// contain a `..` path component.
875 pub fn set_blob_ref(&self, namespace: &str, key: &str, sha: &str) -> Result<(), StoreError> {
876 let dir = self.blob_ref_namespace_dir(namespace, key)?;
877 fs::create_dir_all(&dir)?;
878 fs::write(dir.join(key), sha.as_bytes())?;
879 Ok(())
880 }
881
882 /// Resolve `namespace`/`key` to a blob sha. `UnknownBlobRef` if unbound.
883 pub fn get_blob_ref(&self, namespace: &str, key: &str) -> Result<String, StoreError> {
884 let dir = self.blob_ref_namespace_dir(namespace, key)?;
885 match fs::read_to_string(dir.join(key)) {
886 Ok(s) => Ok(s.trim().to_string()),
887 Err(e) if e.kind() == std::io::ErrorKind::NotFound => Err(StoreError::UnknownBlobRef {
888 namespace: namespace.to_string(),
889 key: key.to_string(),
890 }),
891 Err(e) => Err(StoreError::Io(e)),
892 }
893 }
894
895 /// Blob-ref namespace for committed lockfiles (#930 phase 2b-1).
896 const LOCK_NS: &'static str = "lock";
897
898 /// Record the `lex.lock` committed with the package head `head_op` — the
899 /// exact dependency versions and op-log heads that head was built and
900 /// type-checks against (#930 phase 2b-1: "HEAD + its committed lex.lock
901 /// always type-checks"). Content-addressed via [`Self::put_blob`] and
902 /// bound under the `lock` namespace keyed by the head op, so it is
903 /// idempotent (a re-push converges) and travels with the package through
904 /// the same object-sync path as stages and intents. Keyed by head op
905 /// rather than by branch so re-verifying a *historical* head resolves it
906 /// against the lock that head actually committed, not whatever the branch
907 /// points at now.
908 pub fn set_committed_lock(&self, head_op: &str, lock_toml: &str) -> Result<(), StoreError> {
909 let sha = self.put_blob(lock_toml)?;
910 self.set_blob_ref(Self::LOCK_NS, head_op, &sha)
911 }
912
913 /// The `lex.lock` committed with `head_op`, or `None` when the head
914 /// carries no committed lock — a dependency-free package, or one
915 /// published before locks were committed (the write-time gate then has no
916 /// registry/git dependencies to resolve, exactly as today).
917 pub fn committed_lock(&self, head_op: &str) -> Result<Option<String>, StoreError> {
918 match self.get_blob_ref(Self::LOCK_NS, head_op) {
919 Ok(sha) => Ok(Some(self.get_blob(&sha)?)),
920 Err(StoreError::UnknownBlobRef { .. }) => Ok(None),
921 Err(e) => Err(e),
922 }
923 }
924
925 /// The lock **governing** `head_op`: its own committed lock if it has one,
926 /// otherwise the nearest ancestor's (#975).
927 ///
928 /// [`Self::committed_lock`] is an exact-key lookup, and a lock is only ever
929 /// committed for a head a *client pushes* — or, since #977, for a merge op,
930 /// which commits the union of its parents' locks
931 /// ([`Self::merged_lock_for_parents`]). A head landed through `/v1/patch`
932 /// has none of its own, so the exact lookup returned `None`, the dependency
933 /// resolver got no pins, and a non-inlined head was rejected as
934 /// `unknown_identifier "<alias>"` even though every dependency was
935 /// resolvable. A head inherits its ancestors' pins until a new lock is
936 /// committed: the correct model for a patch head, which genuinely has no
937 /// new pins to record.
938 ///
939 /// Breadth-first, so the *nearest* ancestor wins. That returns one lock
940 /// whole, which is only right along a single line of history — which is
941 /// why a merge op does not rely on it and commits its own merged lock.
942 pub fn committed_lock_inherited(&self, head_op: &str) -> Result<Option<String>, StoreError> {
943 use std::collections::{BTreeSet, VecDeque};
944 let log = lex_vcs::OpLog::open(self.root())?;
945 let mut queue: VecDeque<String> = VecDeque::new();
946 let mut seen: BTreeSet<String> = BTreeSet::new();
947 queue.push_back(head_op.to_string());
948 while let Some(id) = queue.pop_front() {
949 if !seen.insert(id.clone()) {
950 continue;
951 }
952 if let Some(lock) = self.committed_lock(&id)? {
953 return Ok(Some(lock));
954 }
955 // A missing or unreadable op just ends that branch of the walk: an
956 // incomplete local op-log must not fail dependency resolution.
957 if let Ok(Some(rec)) = log.get(&id) {
958 for p in rec.op.parents {
959 queue.push_back(p);
960 }
961 }
962 }
963 Ok(None)
964 }
965
966 /// The lock a merge op with `parents` commits as its own (#977): the
967 /// **union** of every parent's governing lock
968 /// ([`Self::committed_lock_inherited`]).
969 ///
970 /// Inheriting one parent's lock whole is wrong for a merge: if the *source*
971 /// branch introduced a dependency, only the source's lock pins it, and a
972 /// merge that inherited the destination's lock failed its gate as
973 /// `unknown_identifier "<alias>"`. Unioning fixes that.
974 ///
975 /// `parents` must be in merge order — the first is the branch being merged
976 /// INTO (`dst`), later ones are merged in. Note an `Operation`'s own
977 /// `parents` are *sorted* (so the op id is order-independent) and do not
978 /// carry that order; the caller supplies it. It only matters for naming
979 /// the sides of a conflict and for which entry is kept on a tie. When two parents pin the **same** package at
980 /// **different** versions the merge is refused with
981 /// [`StoreError::DependencyConflict`] rather than silently picking one:
982 /// either choice changes what one side's code was tested against. When the
983 /// versions agree, the first parent's entry is kept (filling in a
984 /// `head_op` it lacks from a later parent).
985 ///
986 /// `Ok(None)` when no parent is governed by any lock (a dependency-free
987 /// package), and also when a parent's lock can't be parsed — the merge then
988 /// carries no lock of its own and falls back to inheritance exactly as
989 /// before #977, rather than refusing a merge over an unreadable lock it
990 /// did not write. Read-only: nothing is written.
991 pub fn merged_lock_for_parents(
992 &self,
993 parents: &[lex_vcs::OpId],
994 ) -> Result<Option<String>, StoreError> {
995 use lex_syntax::lock::{LockEntry, LockFile, LOCK_FORMAT_VERSION};
996 let mut raws: Vec<String> = Vec::new();
997 for p in parents {
998 if let Some(raw) = self.committed_lock_inherited(p)? {
999 raws.push(raw);
1000 }
1001 }
1002 match raws.len() {
1003 0 => return Ok(None),
1004 // A single governing lock (one parent, or only one side has deps):
1005 // commit it verbatim — nothing to union, no reformatting churn.
1006 1 => return Ok(raws.pop()),
1007 _ if raws.iter().all(|r| r == &raws[0]) => return Ok(raws.pop()),
1008 _ => {}
1009 }
1010 let mut locks: Vec<LockFile> = Vec::with_capacity(raws.len());
1011 for raw in &raws {
1012 match LockFile::from_toml(raw) {
1013 Ok(lf) => locks.push(lf),
1014 Err(_) => return Ok(None),
1015 }
1016 }
1017 let mut merged: BTreeMap<String, LockEntry> = BTreeMap::new();
1018 for lf in &locks {
1019 for e in &lf.packages {
1020 match merged.get_mut(&e.name) {
1021 None => {
1022 merged.insert(e.name.clone(), e.clone());
1023 }
1024 Some(kept) if kept.version != e.version => {
1025 return Err(StoreError::DependencyConflict {
1026 package: e.name.clone(),
1027 dst_version: kept.version.clone(),
1028 src_version: e.version.clone(),
1029 });
1030 }
1031 Some(kept) => {
1032 if kept.head_op.is_none() {
1033 kept.head_op = e.head_op.clone();
1034 }
1035 }
1036 }
1037 }
1038 }
1039 let version = locks.iter().map(|l| l.version).max().unwrap_or(LOCK_FORMAT_VERSION);
1040 let out = LockFile { version, packages: merged.into_values().collect() };
1041 out.to_toml().map(Some).map_err(|e| StoreError::Io(std::io::Error::other(e)))
1042 }
1043
1044 /// All `key → sha` bindings in a namespace (e.g. every artifact in a
1045 /// sprint). Empty map if the namespace has no bindings yet.
1046 pub fn list_blob_refs(
1047 &self,
1048 namespace: &str,
1049 ) -> Result<std::collections::BTreeMap<String, String>, StoreError> {
1050 let dir = self.blob_ref_namespace_dir(namespace, "x")?;
1051 let mut out = std::collections::BTreeMap::new();
1052 let entries = match fs::read_dir(&dir) {
1053 Ok(e) => e,
1054 Err(e) if e.kind() == std::io::ErrorKind::NotFound => return Ok(out),
1055 Err(e) => return Err(StoreError::Io(e)),
1056 };
1057 for entry in entries {
1058 let entry = entry?;
1059 if entry.file_type()?.is_file() {
1060 let key = entry.file_name().to_string_lossy().to_string();
1061 let sha = fs::read_to_string(entry.path())?.trim().to_string();
1062 out.insert(key, sha);
1063 }
1064 }
1065 Ok(out)
1066 }
1067
1068 // Resolve the on-disk dir for a (namespace, key), rejecting `..` traversal
1069 // and `/` in the key. `key` is validated but not joined here (callers join
1070 // it themselves so `list_blob_refs` can pass a dummy).
1071 fn blob_ref_namespace_dir(&self, namespace: &str, key: &str) -> Result<PathBuf, StoreError> {
1072 if key.contains('/') || key.contains('\\') || key.split('/').any(|c| c == "..") {
1073 return Err(StoreError::UnknownBlobRef {
1074 namespace: namespace.to_string(),
1075 key: key.to_string(),
1076 });
1077 }
1078 let mut dir = self.blob_refs_dir();
1079 for comp in namespace.split('/') {
1080 if comp == ".." || comp.contains('\\') {
1081 return Err(StoreError::UnknownBlobRef {
1082 namespace: namespace.to_string(),
1083 key: key.to_string(),
1084 });
1085 }
1086 if !comp.is_empty() {
1087 dir.push(comp);
1088 }
1089 }
1090 Ok(dir)
1091 }
1092
1093 fn now() -> u64 {
1094 SystemTime::now()
1095 .duration_since(UNIX_EPOCH)
1096 .map(|d| d.as_secs())
1097 .unwrap_or(0)
1098 }
1099
1100 fn sig_dir(&self, sig: &str) -> PathBuf {
1101 self.root.join("stages").join(sig)
1102 }
1103 fn impl_dir(&self, sig: &str) -> PathBuf {
1104 self.sig_dir(sig).join("implementations")
1105 }
1106 fn tests_dir(&self, sig: &str) -> PathBuf {
1107 self.sig_dir(sig).join("tests")
1108 }
1109 fn specs_dir(&self, sig: &str) -> PathBuf {
1110 self.sig_dir(sig).join("specs")
1111 }
1112 fn lifecycle_path(&self, sig: &str) -> PathBuf {
1113 self.sig_dir(sig).join("lifecycle.json")
1114 }
1115
1116 // ---- publish ----
1117
1118 /// Publish a stage as **Draft**. Returns the StageId.
1119 /// Idempotent: republishing the same canonical AST returns the same
1120 /// StageId without writing duplicates.
1121 pub fn publish(&self, stage: &Stage) -> Result<String, StoreError> {
1122 self.publish_signed(stage, None)
1123 }
1124
1125 /// Like [`Self::publish`] but optionally attaches an Ed25519
1126 /// signature over the StageId (#227). When `signer` is `Some`,
1127 /// the persisted metadata gets a `signature` field that
1128 /// downstream consumers can verify via
1129 /// [`lex_vcs::verify_stage_id`].
1130 ///
1131 /// Idempotency: if a metadata file already exists the signature
1132 /// is *not* re-written. This preserves "republishing is a no-op"
1133 /// even across different signers — promoting a signed stage
1134 /// requires a fresh stage hash anyway, so a metadata overwrite
1135 /// would be the wrong primitive.
1136 pub fn publish_signed(
1137 &self,
1138 stage: &Stage,
1139 signer: Option<&lex_vcs::Keypair>,
1140 ) -> Result<String, StoreError> {
1141 let sig = sig_id(stage).ok_or(StoreError::CannotPublishImport)?;
1142 let stage_id = stage_id(stage).ok_or(StoreError::CannotPublishImport)?;
1143 let name = stage_name(stage).to_string();
1144
1145 fs::create_dir_all(self.impl_dir(&sig))?;
1146 fs::create_dir_all(self.tests_dir(&sig))?;
1147 fs::create_dir_all(self.specs_dir(&sig))?;
1148
1149 let ast_path = self.impl_dir(&sig).join(format!("{}.ast.json", stage_id));
1150 let delta_path = self.impl_dir(&sig).join(format!("{}.delta.json", stage_id));
1151 let meta_path = self
1152 .impl_dir(&sig)
1153 .join(format!("{}.metadata.json", stage_id));
1154
1155 // #261 slice 3: try delta encoding against the most recent
1156 // prior stage in this sig's lifecycle. Falls back to a full
1157 // snapshot when (a) no prior stage exists, (b) the diff
1158 // ratio is over the threshold, or (c) the delta chain is
1159 // already at its cap. The decision is internal — callers
1160 // see the same `Stage` object on `get_ast` regardless.
1161 if !ast_path.exists() && !delta_path.exists() {
1162 self.persist_stage_bytes(&sig, &stage_id, stage, &ast_path, &delta_path)?;
1163 }
1164 let doc = stage_doc(stage);
1165 if !meta_path.exists() {
1166 let signature = signer.map(|kp| kp.sign_stage_id(&stage_id));
1167 let metadata = Metadata {
1168 stage_id: stage_id.clone(),
1169 sig_id: sig.clone(),
1170 name,
1171 published_at: Self::now(),
1172 note: None,
1173 doc,
1174 signature,
1175 };
1176 write_canonical_json(&meta_path, &metadata)?;
1177 } else if !doc.is_empty() {
1178 // Comments are deliberately outside the hash, so rewording one
1179 // yields the same StageId and lands here rather than writing a
1180 // fresh metadata file. Without this branch the first version's
1181 // documentation would be pinned forever and a correction could
1182 // never reach the registry. Metadata is not content-addressed, so
1183 // refreshing it in place is safe.
1184 if let Ok(mut existing) = self.get_metadata_for_sig(&sig, &stage_id) {
1185 if existing.doc != doc {
1186 existing.doc = doc;
1187 write_canonical_json(&meta_path, &existing)?;
1188 }
1189 }
1190 }
1191
1192 // Lifecycle: append a Draft transition for first publish.
1193 let mut life = self.read_lifecycle(&sig).unwrap_or_else(|_| Lifecycle {
1194 sig_id: sig.clone(),
1195 ..Default::default()
1196 });
1197 if !life.transitions.iter().any(|t| t.stage_id == stage_id) {
1198 life.transitions.push(Transition {
1199 stage_id: stage_id.clone(),
1200 from: StageStatus::Draft, // synthesized; "from" of first transition is itself
1201 to: StageStatus::Draft,
1202 at: Self::now(),
1203 reason: None,
1204 });
1205 self.write_lifecycle(&sig, &life)?;
1206 // Register the new stage_id's owning sig up front so a
1207 // later `lookup_lifecycle` (e.g. `get_ast`) never needs
1208 // to fall back to a full tenant-wide scan for it.
1209 self.append_stage_index_entry(&stage_id, &sig);
1210 }
1211 Ok(stage_id)
1212 }
1213
1214 // ---- lifecycle ----
1215
1216 pub fn activate(&self, stage_id: &str) -> Result<(), StoreError> {
1217 let (sig, mut life) = self.lookup_lifecycle(stage_id)?;
1218 // Demote any currently-Active impls for this SigId to Deprecated.
1219 let active = life.current_active().map(|s| s.to_string());
1220 if let Some(prev) = active {
1221 if prev != stage_id {
1222 life.transitions.push(Transition {
1223 stage_id: prev,
1224 from: StageStatus::Active,
1225 to: StageStatus::Deprecated,
1226 at: Self::now(),
1227 reason: Some("superseded".into()),
1228 });
1229 }
1230 }
1231 let cur = life.status_of(stage_id);
1232 if cur == Some(StageStatus::Tombstone) {
1233 return Err(StoreError::InvalidTransition(
1234 "tombstoned cannot be activated".into(),
1235 ));
1236 }
1237 life.transitions.push(Transition {
1238 stage_id: stage_id.into(),
1239 from: cur.unwrap_or(StageStatus::Draft),
1240 to: StageStatus::Active,
1241 at: Self::now(),
1242 reason: None,
1243 });
1244 self.write_lifecycle(&sig, &life)
1245 }
1246
1247 pub fn deprecate(&self, stage_id: &str, reason: impl Into<String>) -> Result<(), StoreError> {
1248 let (sig, mut life) = self.lookup_lifecycle(stage_id)?;
1249 let cur = life
1250 .status_of(stage_id)
1251 .ok_or_else(|| StoreError::UnknownStage(stage_id.into()))?;
1252 if cur != StageStatus::Active {
1253 return Err(StoreError::InvalidTransition(format!(
1254 "{cur:?} ⇒ Deprecated"
1255 )));
1256 }
1257 life.transitions.push(Transition {
1258 stage_id: stage_id.into(),
1259 from: cur,
1260 to: StageStatus::Deprecated,
1261 at: Self::now(),
1262 reason: Some(reason.into()),
1263 });
1264 self.write_lifecycle(&sig, &life)
1265 }
1266
1267 pub fn tombstone(&self, stage_id: &str) -> Result<(), StoreError> {
1268 let (sig, mut life) = self.lookup_lifecycle(stage_id)?;
1269 let cur = life
1270 .status_of(stage_id)
1271 .ok_or_else(|| StoreError::UnknownStage(stage_id.into()))?;
1272 if cur != StageStatus::Deprecated {
1273 return Err(StoreError::InvalidTransition(format!(
1274 "{cur:?} ⇒ Tombstone"
1275 )));
1276 }
1277 life.transitions.push(Transition {
1278 stage_id: stage_id.into(),
1279 from: cur,
1280 to: StageStatus::Tombstone,
1281 at: Self::now(),
1282 reason: None,
1283 });
1284 self.write_lifecycle(&sig, &life)
1285 }
1286
1287 // ---- queries ----
1288
1289 /// The current Active StageId for a signature, or `None`.
1290 pub fn resolve_sig(&self, sig: &str) -> Result<Option<String>, StoreError> {
1291 let life = match self.read_lifecycle(sig) {
1292 Ok(l) => l,
1293 Err(_) => return Ok(None),
1294 };
1295 Ok(life.current_active().map(|s| s.to_string()))
1296 }
1297
1298 /// Per-stage history for a SigId, ordered chronologically by
1299 /// the *last* transition timestamp. Returns one entry per
1300 /// distinct StageId that has ever been published under `sig`.
1301 /// `Ok(vec![])` if the SigId doesn't exist in the store.
1302 ///
1303 /// Used by `lex blame` to render "where does this fn come from".
1304 pub fn sig_history(&self, sig: &str) -> Result<Vec<StageHistoryEntry>, StoreError> {
1305 let life = match self.read_lifecycle(sig) {
1306 Ok(l) => l,
1307 Err(_) => return Ok(Vec::new()),
1308 };
1309 // Collapse transitions: latest status + last_at per stage,
1310 // plus the timestamp of the first Draft transition (≈ when
1311 // the stage was published) when one exists.
1312 let mut by_stage: indexmap::IndexMap<String, StageHistoryEntry> = indexmap::IndexMap::new();
1313 for t in &life.transitions {
1314 let entry = by_stage
1315 .entry(t.stage_id.clone())
1316 .or_insert(StageHistoryEntry {
1317 stage_id: t.stage_id.clone(),
1318 status: t.to,
1319 last_at: t.at,
1320 published_at: None,
1321 });
1322 entry.status = t.to;
1323 entry.last_at = t.at;
1324 if t.from == StageStatus::Draft && entry.published_at.is_none() {
1325 entry.published_at = Some(t.at);
1326 }
1327 if t.to == StageStatus::Draft && entry.published_at.is_none() {
1328 // Initial publication: Draft is the *destination*.
1329 entry.published_at = Some(t.at);
1330 }
1331 }
1332 let mut out: Vec<StageHistoryEntry> = by_stage.into_values().collect();
1333 // Sort newest first so `lex blame` shows recent activity at top.
1334 out.sort_by_key(|e| std::cmp::Reverse(e.last_at));
1335 Ok(out)
1336 }
1337
1338 pub fn get_ast(&self, stage_id: &str) -> Result<Stage, StoreError> {
1339 let (sig, _) = self.lookup_lifecycle(stage_id)?;
1340 let bytes = self.read_stage_canonical_bytes(&sig, stage_id)?;
1341 Ok(serde_json::from_slice(&bytes)?)
1342 }
1343
1344 /// Bulk AST fetch for callers that already know each stage's
1345 /// **signature** — a branch head map, for instance, which is keyed
1346 /// by SigId and whose values are the StageIds it points at.
1347 ///
1348 /// Prefer this over [`Self::get_asts_bulk`] whenever the SigId is in
1349 /// hand, because resolving a StageId back to a SigId is not
1350 /// reliable: a StageId hashes the structural signature plus the
1351 /// implementation, deliberately *not* the name
1352 /// (`docs/INVARIANTS.md`), so two functions that differ only in name
1353 /// share one StageId while having two distinct SigIds — and two
1354 /// separate ASTs, one under each sig directory. `stage_index` maps
1355 /// each StageId to a single sig, so `get_ast`/`get_asts_bulk` return
1356 /// whichever of those ASTs the index happens to name, i.e. the wrong
1357 /// name half the time (#826). Reading straight from the sig the
1358 /// caller already knows removes the ambiguity — and skips loading
1359 /// the index at all.
1360 ///
1361 /// Returns results in the same order as `pairs`, `Err` for anything
1362 /// that fails to resolve (mirroring `get_ast`'s error semantics).
1363 pub fn get_asts_for_sigs_bulk(
1364 &self,
1365 pairs: &[(String, String)],
1366 ) -> Vec<Result<Stage, StoreError>> {
1367 pairs
1368 .iter()
1369 .map(|(sig_id, stage_id)| {
1370 let bytes = self.read_stage_canonical_bytes(sig_id, stage_id)?;
1371 Ok(serde_json::from_slice(&bytes)?)
1372 })
1373 .collect()
1374 }
1375
1376 /// Bulk variant of [`Self::get_ast`] for callers resolving many
1377 /// stage_ids at once (e.g. `pkg_publish_handler`'s `old_head`
1378 /// scan over every live function in a tenant, once per publish
1379 /// request). `get_ast` in a loop calls `lookup_lifecycle` once
1380 /// per stage_id, and `lookup_lifecycle`'s index-hit path reads
1381 /// and re-parses the *entire* `stage_index.jsonl` on every single
1382 /// call — fine for one call, but O(index size × N) for N calls in
1383 /// a row, which dominates once the index itself is large (#825's
1384 /// follow-up: still correct and far better than the pre-index
1385 /// full-tenant-scan-per-call behavior, but the per-call reparse
1386 /// is itself a real, measured cost — 87.6s for 3,664 calls against
1387 /// a ~14k-line index on the alpibrusl tenant).
1388 ///
1389 /// This loads the index once for the whole batch and keeps it in
1390 /// memory across all `stage_ids`, only touching disk again to
1391 /// append genuinely new entries (a positive backfill or a
1392 /// negative "not found anywhere" cache, same as the single-call
1393 /// path) — never to re-read what's already loaded.
1394 ///
1395 /// Returns results in the same order as `stage_ids`, `Err` for
1396 /// anything that fails to resolve (mirroring `get_ast`'s error
1397 /// semantics per call).
1398 pub fn get_asts_bulk(&self, stage_ids: &[String]) -> Vec<Result<Stage, StoreError>> {
1399 let mut index = self.load_stage_index();
1400 let mut sigs_cache: BTreeMap<String, Option<Lifecycle>> = BTreeMap::new();
1401 let mut all_sigs: Option<Vec<String>> = None;
1402
1403 stage_ids
1404 .iter()
1405 .map(|stage_id| {
1406 self.lookup_lifecycle_bulk(stage_id, &mut index, &mut sigs_cache, &mut all_sigs)
1407 .and_then(|sig| {
1408 let bytes = self.read_stage_canonical_bytes(&sig, stage_id)?;
1409 Ok(serde_json::from_slice(&bytes)?)
1410 })
1411 })
1412 .collect()
1413 }
1414
1415 /// Shared implementation behind [`Self::get_asts_bulk`]: identical
1416 /// logic to [`Self::lookup_lifecycle`], but reads and writes the
1417 /// caller-supplied `index` map instead of reloading it from disk
1418 /// on every call, and memoizes `read_lifecycle` per sig and the
1419 /// `list_sigs()` full-scan list across the whole batch. Disk
1420 /// writes for newly-discovered entries (positive or negative)
1421 /// still happen immediately, same as the single-call path — only
1422 /// the repeated *reads* are batched away.
1423 fn lookup_lifecycle_bulk(
1424 &self,
1425 stage_id: &str,
1426 index: &mut BTreeMap<String, String>,
1427 sigs_cache: &mut BTreeMap<String, Option<Lifecycle>>,
1428 all_sigs: &mut Option<Vec<String>>,
1429 ) -> Result<String, StoreError> {
1430 if let Some(sig) = index.get(stage_id) {
1431 if sig == MISSING_STAGE_MARKER {
1432 return Err(StoreError::UnknownStage(stage_id.into()));
1433 }
1434 let life = sigs_cache
1435 .entry(sig.clone())
1436 .or_insert_with(|| self.read_lifecycle(sig).ok());
1437 if let Some(life) = life {
1438 if life.transitions.iter().any(|t| t.stage_id == stage_id) {
1439 return Ok(sig.clone());
1440 }
1441 }
1442 }
1443 if all_sigs.is_none() {
1444 *all_sigs = Some(self.list_sigs()?);
1445 }
1446 for sig in all_sigs.as_ref().unwrap() {
1447 let life = sigs_cache
1448 .entry(sig.clone())
1449 .or_insert_with(|| self.read_lifecycle(sig).ok());
1450 if let Some(life) = life {
1451 if life.transitions.iter().any(|t| t.stage_id == stage_id) {
1452 self.append_stage_index_entry(stage_id, sig);
1453 index.insert(stage_id.to_string(), sig.clone());
1454 return Ok(sig.clone());
1455 }
1456 }
1457 }
1458 self.append_stage_index_entry(stage_id, MISSING_STAGE_MARKER);
1459 index.insert(stage_id.to_string(), MISSING_STAGE_MARKER.to_string());
1460 Err(StoreError::UnknownStage(stage_id.into()))
1461 }
1462
1463 /// Read the canonical bytes of a stage, walking back through
1464 /// any delta chain (#261 slice 3). The recursion ends at a
1465 /// `<stage_id>.ast.json` file (a full snapshot) or, in the
1466 /// degenerate case of a missing chain, with `UnknownStage`.
1467 fn read_stage_canonical_bytes(&self, sig: &str, stage_id: &str) -> Result<Vec<u8>, StoreError> {
1468 let ast_path = self.impl_dir(sig).join(format!("{}.ast.json", stage_id));
1469 if ast_path.exists() {
1470 return Ok(fs::read(&ast_path)?);
1471 }
1472 let delta_path = self.impl_dir(sig).join(format!("{}.delta.json", stage_id));
1473 if !delta_path.exists() {
1474 return Err(StoreError::UnknownStage(stage_id.into()));
1475 }
1476 let delta_bytes = fs::read(&delta_path)?;
1477 let delta: crate::delta::StageDelta = serde_json::from_slice(&delta_bytes)?;
1478 let base_bytes = self.read_stage_canonical_bytes(sig, &delta.base_stage_id)?;
1479 crate::delta::apply(&base_bytes, &delta).map_err(|e| {
1480 StoreError::Io(std::io::Error::new(
1481 std::io::ErrorKind::InvalidData,
1482 format!("applying delta for {stage_id}: {e}"),
1483 ))
1484 })
1485 }
1486
1487 /// Persist a freshly-published stage's canonical bytes (#261
1488 /// slice 3). Tries delta encoding against the most recent
1489 /// prior stage in the sig's lifecycle; falls back to a full
1490 /// snapshot when no base exists, the diff ratio is too high,
1491 /// or the delta chain is already at its cap.
1492 fn persist_stage_bytes(
1493 &self,
1494 sig: &str,
1495 stage_id: &str,
1496 stage: &Stage,
1497 ast_path: &Path,
1498 delta_path: &Path,
1499 ) -> Result<(), StoreError> {
1500 let new_bytes = canonical_bytes(stage)?;
1501 if let Some((base_stage_id, base_chain_length)) = self.pick_delta_base(sig, stage_id)? {
1502 let base_bytes = self.read_stage_canonical_bytes(sig, &base_stage_id)?;
1503 let (prefix, suffix, middle) = crate::delta::splice(&base_bytes, &new_bytes);
1504 let chain_length = base_chain_length + 1;
1505 if crate::delta::is_worth_encoding(middle.len(), new_bytes.len(), chain_length) {
1506 let delta = crate::delta::StageDelta {
1507 base_stage_id,
1508 chain_length,
1509 common_prefix: prefix,
1510 common_suffix: suffix,
1511 middle_hex: hex::encode(&middle),
1512 };
1513 write_canonical_json(delta_path, &delta)?;
1514 return Ok(());
1515 }
1516 }
1517 // Fall through: full snapshot.
1518 if let Some(parent) = ast_path.parent() {
1519 fs::create_dir_all(parent)?;
1520 }
1521 fs::write(ast_path, &new_bytes)?;
1522 Ok(())
1523 }
1524
1525 /// Pick a base stage for delta encoding from the given sig's
1526 /// lifecycle. Returns `(base_stage_id, base_chain_length)` for
1527 /// the most-recent non-tombstoned prior stage, or `None` when
1528 /// there is no candidate. The chain length is read off the
1529 /// base's `.delta.json` (if any) to enforce the cap.
1530 fn pick_delta_base(
1531 &self,
1532 sig: &str,
1533 new_stage_id: &str,
1534 ) -> Result<Option<(String, usize)>, StoreError> {
1535 let life = self.read_lifecycle(sig).ok();
1536 let Some(life) = life else {
1537 return Ok(None);
1538 };
1539 // Walk transitions newest-first; pick the first prior
1540 // stage that isn't this one and isn't tombstoned.
1541 let mut latest_per_stage: indexmap::IndexMap<&str, StageStatus> = indexmap::IndexMap::new();
1542 for t in &life.transitions {
1543 latest_per_stage.insert(&t.stage_id, t.to);
1544 }
1545 let mut candidates: Vec<&str> = latest_per_stage
1546 .iter()
1547 .filter(|(id, status)| **id != new_stage_id && **status != StageStatus::Tombstone)
1548 .map(|(id, _)| *id)
1549 .collect();
1550 // Reverse to get newest-first (transitions are append-only,
1551 // so latest_per_stage's iteration order matches insertion
1552 // order, oldest-first).
1553 candidates.reverse();
1554 let Some(&base) = candidates.first() else {
1555 return Ok(None);
1556 };
1557 let base_chain_length = self.delta_chain_length(sig, base)?;
1558 Ok(Some((base.to_string(), base_chain_length)))
1559 }
1560
1561 /// Length of the delta chain ending at `stage_id`. Zero when
1562 /// the stage is a full snapshot (`.ast.json` present); the
1563 /// stored `chain_length` from `.delta.json` otherwise.
1564 fn delta_chain_length(&self, sig: &str, stage_id: &str) -> Result<usize, StoreError> {
1565 let ast_path = self.impl_dir(sig).join(format!("{}.ast.json", stage_id));
1566 if ast_path.exists() {
1567 return Ok(0);
1568 }
1569 let delta_path = self.impl_dir(sig).join(format!("{}.delta.json", stage_id));
1570 if !delta_path.exists() {
1571 return Ok(0);
1572 }
1573 let bytes = fs::read(&delta_path)?;
1574 let delta: crate::delta::StageDelta = serde_json::from_slice(&bytes)?;
1575 Ok(delta.chain_length)
1576 }
1577
1578 /// The metadata for a stage **under a named sig**.
1579 ///
1580 /// [`Self::get_metadata`] resolves the sig through the stage index, which
1581 /// is a stage-id-only lookup — and a StageId does not encode the name
1582 /// (#826), so two sigs can share one. Resolving by id then hands both
1583 /// declarations whichever record the index happens to point at. Harmless
1584 /// for a hash, not for content: it gave two different declarations the
1585 /// *same* documentation, duplicating a module header into a package that
1586 /// had it once.
1587 ///
1588 /// When the caller knows the sig — a head map is keyed by it — this is the
1589 /// honest lookup.
1590 pub fn get_metadata_for_sig(
1591 &self,
1592 sig_id: &str,
1593 stage_id: &str,
1594 ) -> Result<Metadata, StoreError> {
1595 let path = self.impl_dir(sig_id).join(format!("{stage_id}.metadata.json"));
1596 let bytes = fs::read(&path)?;
1597 Ok(serde_json::from_slice(&bytes)?)
1598 }
1599
1600 pub fn get_metadata(&self, stage_id: &str) -> Result<Metadata, StoreError> {
1601 let (sig, _) = self.lookup_lifecycle(stage_id)?;
1602 let path = self
1603 .impl_dir(&sig)
1604 .join(format!("{}.metadata.json", stage_id));
1605 let bytes = fs::read(&path)?;
1606 Ok(serde_json::from_slice(&bytes)?)
1607 }
1608
1609 pub fn get_status(&self, stage_id: &str) -> Result<StageStatus, StoreError> {
1610 let (_sig, life) = self.lookup_lifecycle(stage_id)?;
1611 life.status_of(stage_id)
1612 .ok_or_else(|| StoreError::UnknownStage(stage_id.into()))
1613 }
1614
1615 pub fn list_stages_by_name(&self, name: &str) -> Result<Vec<String>, StoreError> {
1616 // Walk every SigId → check metadata of any implementation; if its
1617 // name matches, include the SigId.
1618 let mut out = Vec::new();
1619 let stages_dir = self.root.join("stages");
1620 if !stages_dir.exists() {
1621 return Ok(out);
1622 }
1623 for entry in fs::read_dir(&stages_dir)? {
1624 let entry = entry?;
1625 let sig_dir = entry.path();
1626 if !sig_dir.is_dir() {
1627 continue;
1628 }
1629 let sig = entry.file_name().to_string_lossy().to_string();
1630 // Look at any one metadata file under this SigId.
1631 let impls = self.impl_dir(&sig);
1632 if !impls.exists() {
1633 continue;
1634 }
1635 for f in fs::read_dir(impls)? {
1636 let f = f?;
1637 let p = f.path();
1638 if p.extension().is_some_and(|e| e == "json")
1639 && p.file_name()
1640 .is_some_and(|n| n.to_string_lossy().ends_with(".metadata.json"))
1641 {
1642 if let Ok(bytes) = fs::read(&p) {
1643 if let Ok(m) = serde_json::from_slice::<Metadata>(&bytes) {
1644 if m.name == name {
1645 if !out.contains(&sig) {
1646 out.push(sig.clone());
1647 }
1648 break;
1649 }
1650 }
1651 }
1652 }
1653 }
1654 }
1655 out.sort();
1656 Ok(out)
1657 }
1658
1659 pub fn list_sigs(&self) -> Result<Vec<String>, StoreError> {
1660 let stages_dir = self.root.join("stages");
1661 let mut out = Vec::new();
1662 if !stages_dir.exists() {
1663 return Ok(out);
1664 }
1665 for entry in fs::read_dir(stages_dir)? {
1666 let entry = entry?;
1667 if entry.file_type()?.is_dir() {
1668 out.push(entry.file_name().to_string_lossy().to_string());
1669 }
1670 }
1671 out.sort();
1672 Ok(out)
1673 }
1674
1675 // ---- tests/specs as metadata (§4.4) ----
1676
1677 pub fn attach_test(&self, sig: &str, test: &Test) -> Result<String, StoreError> {
1678 if !self.sig_dir(sig).exists() {
1679 return Err(StoreError::UnknownSig(sig.into()));
1680 }
1681 fs::create_dir_all(self.tests_dir(sig))?;
1682 let path = self.tests_dir(sig).join(format!("{}.json", test.id));
1683 write_canonical_json(&path, test)?;
1684 Ok(test.id.clone())
1685 }
1686
1687 pub fn list_tests(&self, sig: &str) -> Result<Vec<Test>, StoreError> {
1688 let dir = self.tests_dir(sig);
1689 if !dir.exists() {
1690 return Ok(Vec::new());
1691 }
1692 let mut out = Vec::new();
1693 for f in fs::read_dir(dir)? {
1694 let f = f?;
1695 if f.path().extension().is_some_and(|e| e == "json") {
1696 let bytes = fs::read(f.path())?;
1697 out.push(serde_json::from_slice(&bytes)?);
1698 }
1699 }
1700 Ok(out)
1701 }
1702
1703 pub fn attach_spec(&self, sig: &str, spec: &Spec) -> Result<String, StoreError> {
1704 if !self.sig_dir(sig).exists() {
1705 return Err(StoreError::UnknownSig(sig.into()));
1706 }
1707 fs::create_dir_all(self.specs_dir(sig))?;
1708 let path = self.specs_dir(sig).join(format!("{}.json", spec.id));
1709 write_canonical_json(&path, spec)?;
1710 Ok(spec.id.clone())
1711 }
1712
1713 pub fn list_specs(&self, sig: &str) -> Result<Vec<Spec>, StoreError> {
1714 let dir = self.specs_dir(sig);
1715 if !dir.exists() {
1716 return Ok(Vec::new());
1717 }
1718 let mut out = Vec::new();
1719 for f in fs::read_dir(dir)? {
1720 let f = f?;
1721 if f.path().extension().is_some_and(|e| e == "json") {
1722 let bytes = fs::read(f.path())?;
1723 out.push(serde_json::from_slice(&bytes)?);
1724 }
1725 }
1726 Ok(out)
1727 }
1728
1729 // ---- traces (§4.2 / M7) ----
1730
1731 // Native run-trace store — gated behind the `trace` feature (depends on
1732 // lex-trace). Off when a lower crate (lex-runtime) depends on lex-store to
1733 // avoid a dependency cycle; the blob/stage store below is unaffected.
1734 #[cfg(feature = "trace")]
1735 fn trace_path(&self, run_id: &str) -> PathBuf {
1736 self.root.join("traces").join(run_id).join("trace.json")
1737 }
1738
1739 #[cfg(feature = "trace")]
1740 pub fn save_trace(&self, tree: &lex_trace::TraceTree) -> Result<String, StoreError> {
1741 let path = self.trace_path(&tree.run_id);
1742 write_canonical_json(&path, tree)?;
1743 Ok(tree.run_id.clone())
1744 }
1745
1746 #[cfg(feature = "trace")]
1747 pub fn load_trace(&self, run_id: &str) -> Result<lex_trace::TraceTree, StoreError> {
1748 let bytes = fs::read(self.trace_path(run_id))?;
1749 Ok(serde_json::from_slice(&bytes)?)
1750 }
1751
1752 pub fn list_traces(&self) -> Result<Vec<String>, StoreError> {
1753 let dir = self.root.join("traces");
1754 if !dir.exists() {
1755 return Ok(Vec::new());
1756 }
1757 let mut out = Vec::new();
1758 for entry in fs::read_dir(dir)? {
1759 let entry = entry?;
1760 if entry.file_type()?.is_dir() {
1761 out.push(entry.file_name().to_string_lossy().to_string());
1762 }
1763 }
1764 out.sort();
1765 Ok(out)
1766 }
1767
1768 // ---- internals ----
1769
1770 /// `<root>/stage_index.jsonl` — an append-only, best-effort
1771 /// reverse index (`StageId` -> owning `SigId`), one JSON object
1772 /// per line. Backs `lookup_lifecycle`'s fast path; see its doc
1773 /// comment. Not a second source of truth: every entry is
1774 /// reconstructible from `stages/<sig>/lifecycle.json`, so a
1775 /// missing, truncated, or entirely absent index file only costs
1776 /// a slower lookup (the pre-existing full scan), never
1777 /// correctness — matching this module's "filesystem is the
1778 /// source of truth" stance (see the module doc comment) rather
1779 /// than introducing an actual second database.
1780 fn stage_index_path(&self) -> PathBuf {
1781 self.root.join("stage_index.jsonl")
1782 }
1783
1784 /// Best-effort load of the whole reverse index into memory.
1785 /// Tolerates a missing file (no index yet) and a corrupt or
1786 /// torn last line (a crash mid-append under the single-writer
1787 /// Tier-1 assumption) by skipping lines that don't parse,
1788 /// rather than failing the lookup that triggered the load.
1789 fn load_stage_index(&self) -> std::collections::BTreeMap<String, String> {
1790 let mut out = std::collections::BTreeMap::new();
1791 let Ok(raw) = fs::read_to_string(self.stage_index_path()) else {
1792 return out;
1793 };
1794 for line in raw.lines() {
1795 if let Ok(entry) = serde_json::from_str::<StageIndexEntry>(line) {
1796 out.insert(entry.stage_id, entry.sig_id);
1797 }
1798 }
1799 out
1800 }
1801
1802 /// Best-effort append of one new `(stage_id, sig_id)` pair.
1803 /// Failure (e.g. a read-only filesystem) only costs a future
1804 /// full scan for this stage_id, never correctness, so it's
1805 /// swallowed rather than propagated.
1806 fn append_stage_index_entry(&self, stage_id: &str, sig: &str) {
1807 use std::io::Write;
1808 let entry = StageIndexEntry { stage_id: stage_id.into(), sig_id: sig.into() };
1809 let Ok(line) = serde_json::to_string(&entry) else { return };
1810 if let Ok(mut f) = fs::OpenOptions::new()
1811 .create(true)
1812 .append(true)
1813 .open(self.stage_index_path())
1814 {
1815 let _ = writeln!(f, "{line}");
1816 }
1817 }
1818
1819 /// Find which SigId owns a StageId, and that sig's lifecycle.
1820 ///
1821 /// Before the reverse index (#822): a full scan over *every*
1822 /// SigId in the tenant (`list_sigs()`, not scoped to the
1823 /// package being looked at), reading and parsing each one's
1824 /// `lifecycle.json` until a match turned up. `get_ast` — called
1825 /// once per pre-existing function when building a publish
1826 /// request's `old_fns_by_name` (`lex-api/src/handlers.rs`) —
1827 /// calls this once per function, so a tenant with a few thousand
1828 /// published functions turned a single publish into millions of
1829 /// individual file reads; measured at roughly an hour on the
1830 /// `alpibrusl` tenant's ~2,400-function store.
1831 ///
1832 /// Now: check the persisted reverse index first (one sequential
1833 /// file read instead of up to N separate ones). A miss — the
1834 /// index doesn't exist yet, or this stage_id predates it — falls
1835 /// back to the full scan and backfills the index so the next
1836 /// lookup for the same stage_id is fast.
1837 fn lookup_lifecycle(&self, stage_id: &str) -> Result<(String, Lifecycle), StoreError> {
1838 let index = self.load_stage_index();
1839 if let Some(sig) = index.get(stage_id) {
1840 if sig == MISSING_STAGE_MARKER {
1841 // A previous full scan already established this
1842 // stage_id exists nowhere in the store. Re-scanning
1843 // would find nothing again -- see #825: a genuinely
1844 // orphaned reference (e.g. from data predating some
1845 // store migration) is looked up on *every* call that
1846 // needs it, forever, so without this negative cache
1847 // it silently costs a full O(total sigs) scan each
1848 // time, indistinguishable from the positive case at
1849 // the call site. Measured directly: on the alpibrusl
1850 // tenant, 988 of 3,664 branch-head entries are
1851 // orphaned this way, turning one `pkg publish`'s
1852 // old_fns_by_name build into ~16M wasted lifecycle
1853 // reads.
1854 return Err(StoreError::UnknownStage(stage_id.into()));
1855 }
1856 if let Ok(life) = self.read_lifecycle(sig) {
1857 if life.transitions.iter().any(|t| t.stage_id == stage_id) {
1858 return Ok((sig.clone(), life));
1859 }
1860 }
1861 // Index entry is stale or wrong (shouldn't happen in
1862 // practice — sig ownership of a stage_id is permanent).
1863 // Fall through to the full scan below rather than trust it.
1864 }
1865 for sig in self.list_sigs()? {
1866 if let Ok(life) = self.read_lifecycle(&sig) {
1867 if life.transitions.iter().any(|t| t.stage_id == stage_id) {
1868 self.append_stage_index_entry(stage_id, &sig);
1869 return Ok((sig, life));
1870 }
1871 }
1872 }
1873 // Genuinely not found anywhere: cache that fact so the next
1874 // lookup for this exact stage_id is an index hit, not another
1875 // full scan. Safe even if this stage_id somehow gets a real
1876 // sig later (content-addressed publish is idempotent, so
1877 // "later" only means "a byte-identical stage republished
1878 // under a real sig") — `append_stage_index_entry`'s later,
1879 // real entry is a later line in the file, and `load_stage_index`
1880 // folds duplicate keys last-write-wins, so the real entry wins.
1881 self.append_stage_index_entry(stage_id, MISSING_STAGE_MARKER);
1882 Err(StoreError::UnknownStage(stage_id.into()))
1883 }
1884
1885 fn read_lifecycle(&self, sig: &str) -> Result<Lifecycle, StoreError> {
1886 let path = self.lifecycle_path(sig);
1887 if !path.exists() {
1888 return Ok(Lifecycle {
1889 sig_id: sig.into(),
1890 transitions: Vec::new(),
1891 });
1892 }
1893 let bytes = fs::read(&path)?;
1894 Ok(serde_json::from_slice(&bytes)?)
1895 }
1896
1897 fn write_lifecycle(&self, sig: &str, life: &Lifecycle) -> Result<(), StoreError> {
1898 write_canonical_json(&self.lifecycle_path(sig), life)
1899 }
1900
1901 /// Apply a published program to a branch as a sequence of typed
1902 /// operations. Returns the ordered list of op_ids + the new
1903 /// head_op. The caller (`lex publish` CLI, `lex serve`'s HTTP
1904 /// handler) is responsible for computing the `DiffReport` against
1905 /// the current branch head — the diff infrastructure lives in
1906 /// `lex-vcs::compute_diff` (previously `lex-cli`) to keep this
1907 /// layer from owning diffing logic.
1908 ///
1909 /// On success: every op in the returned list is durable in the
1910 /// op log and the branch's head_op points at the last one.
1911 /// On a no-op (no diff): returns empty `ops` and the existing
1912 /// `head_op` unchanged.
1913 pub fn publish_program(
1914 &self,
1915 branch: &str,
1916 stages: &[lex_ast::Stage],
1917 diff: &lex_vcs::DiffReport,
1918 new_imports: &lex_vcs::ImportMap,
1919 activate: bool,
1920 ) -> Result<PublishOutcome, StoreError> {
1921 self.publish_program_signed(branch, stages, diff, new_imports, activate, None)
1922 }
1923
1924 /// Signed variant of [`Self::publish_program`] (#227). Every
1925 /// stage written under this batch gets the same signer; per-stage
1926 /// keys aren't supported because the agent identity model treats
1927 /// a publish as a single authorial act.
1928 pub fn publish_program_signed(
1929 &self,
1930 branch: &str,
1931 stages: &[lex_ast::Stage],
1932 diff: &lex_vcs::DiffReport,
1933 new_imports: &lex_vcs::ImportMap,
1934 activate: bool,
1935 signer: Option<&lex_vcs::Keypair>,
1936 ) -> Result<PublishOutcome, StoreError> {
1937 // Single-file / test callers don't publish a mangled package, so
1938 // there are no module prefixes to record (`in_file` stays `None`).
1939 self.publish_program_with_intent(
1940 branch,
1941 stages,
1942 diff,
1943 new_imports,
1944 activate,
1945 signer,
1946 None,
1947 &std::collections::BTreeMap::new(),
1948 )
1949 }
1950
1951 /// [`Self::publish_program_signed`] plus an optional `intent_id`
1952 /// (#131 / #839): when given, every op this publish emits is stamped
1953 /// with it, so the op log records *why* the change happened — the
1954 /// prompt / model / session an agent was acting under — not only
1955 /// what it was. `lex recall --intent <id>` and `lex op replay` read
1956 /// it back. The caller records the [`lex_vcs::Intent`] in the
1957 /// [`lex_vcs::IntentLog`] beforehand; this only links ops to it.
1958 /// `None` is the existing (intent-less) behavior, so op ids for
1959 /// intent-less publishes are unchanged.
1960 // A batch publish legitimately takes the branch, program, diff,
1961 // imports, activate flag, signer, and now the intent — bundling
1962 // them into a struct for one optional field would obscure more
1963 // than it clarifies.
1964 #[allow(clippy::too_many_arguments)]
1965 pub fn publish_program_with_intent(
1966 &self,
1967 branch: &str,
1968 stages: &[lex_ast::Stage],
1969 diff: &lex_vcs::DiffReport,
1970 new_imports: &lex_vcs::ImportMap,
1971 activate: bool,
1972 signer: Option<&lex_vcs::Keypair>,
1973 intent_id: Option<lex_vcs::IntentId>,
1974 // Mangling prefix → package source file, for a multi-module
1975 // package publish; empty for a single file. Recorded as each
1976 // `AddFunction`/`AddType`'s `in_file` so `export-git` can
1977 // de-flatten the package (#894).
1978 module_prefixes: &std::collections::BTreeMap<String, String>,
1979 ) -> Result<PublishOutcome, StoreError> {
1980 use std::collections::{BTreeMap, BTreeSet};
1981
1982 // #130's write-time gate: verify the candidate program
1983 // typechecks (and effects are correctly declared) before
1984 // any disk side-effect. If anything fails, return the
1985 // structured envelope and leave the branch head unchanged
1986 // — the store's "always-valid HEAD" invariant only holds
1987 // because this is the only batch-publish path that
1988 // advances heads. Single-op writes via the lower-level
1989 // `apply_operation` are not gated yet (#130 follow-up).
1990 // #930: resolve any external dependency edges the head keeps
1991 // (empty when no resolver is installed or the head is inlined).
1992 //
1993 // `head_op` is deliberately `None` here and must stay that way (#945):
1994 // this call CREATES the head, so there is no committed lock keyed to it
1995 // yet — the pins live in the caller's working-copy `lex.lock`, which is
1996 // exactly what a `None` head tells the resolver to use. Passing the
1997 // *parent* head would be wrong: a publish may introduce a brand-new
1998 // dependency whose pin exists only in the working copy. `stages` here is
1999 // the full program as loaded, so it already carries its own `import`
2000 // edges — unlike the map-reconstructed heads in the gates below, which
2001 // must add them via `with_head_imports`.
2002 if let Err(errors) = self.check_with_resolved_deps(stages, None) {
2003 return Err(StoreError::TypeError(errors));
2004 }
2005
2006 // Build old-side views from the current branch. There used to be
2007 // an `old_name_to_sig: BTreeMap<String, SigId>` built here too,
2008 // keyed by bare function name — but a bare name is not unique
2009 // across a package's files (#818: two files can legitimately
2010 // both declare a local `validate` helper with different
2011 // signatures), so a name-keyed map silently collapsed distinct
2012 // SigIds onto one. `diff` now carries each entry's own resolved
2013 // `old_sig_id` directly (see `diff_report`'s doc comments), so
2014 // `diff_to_ops` no longer needs this lookup at all.
2015 let old_head = self.branch_head(branch)?;
2016 // Read every live function's effects through the SigId the head
2017 // names, in one batch. Two reasons, both load-bearing:
2018 //
2019 // * Cost. This was a `get_ast` per live function, and
2020 // `get_ast`'s index-hit path re-reads and re-parses the whole
2021 // `stage_index.jsonl` on every call — O(index × live fns) per
2022 // publish, paid again for every `publish_program` call a
2023 // multi-file publish makes (#828; measured 34s for a no-op
2024 // republish of a real 21-file package against only 698 live
2025 // functions, nearly all of it here).
2026 // * Correctness. A StageId is name-independent, so two live
2027 // functions differing only in name share one and the index
2028 // maps it to a single sig — resolving by StageId therefore
2029 // attributed one function's effects to the *other* one's sig,
2030 // the same ambiguity #826 fixed in `pkg_publish_handler`.
2031 let head_pairs: Vec<(String, String)> = old_head
2032 .iter()
2033 .map(|(sig, stage)| (sig.clone(), stage.clone()))
2034 .collect();
2035 let old_effects: BTreeMap<String, BTreeSet<String>> = head_pairs
2036 .iter()
2037 .zip(self.get_asts_for_sigs_bulk(&head_pairs))
2038 .filter_map(|((sig, _), ast)| match ast.ok()? {
2039 lex_ast::Stage::FnDecl(fd) => {
2040 let s: BTreeSet<String> =
2041 fd.effects.iter().map(|e| e.name.clone()).collect();
2042 Some((sig.clone(), s))
2043 }
2044 _ => None,
2045 })
2046 .collect();
2047 let old_imports = self.derive_imports_from_oplog(branch)?;
2048
2049 let op_kinds = lex_vcs::diff_to_ops(lex_vcs::DiffInputs {
2050 old_head: &old_head,
2051 old_effects: &old_effects,
2052 old_imports: &old_imports,
2053 new_stages: stages,
2054 new_imports,
2055 diff,
2056 module_prefixes,
2057 })
2058 .map_err(|e| StoreError::InvalidTransition(format!("diff_to_ops: {e}")))?;
2059
2060 // Retire any stranded head entry first, so the rest of this publish
2061 // applies to a head every part of which can actually be read back.
2062 // Normally this is empty and costs nothing; it fires only on a head
2063 // already carrying pre-#992 damage.
2064 let op_kinds = {
2065 let mut heal = self.stranded_head_entries(&head_pairs);
2066 if heal.is_empty() {
2067 op_kinds
2068 } else {
2069 heal.extend(op_kinds);
2070 heal
2071 }
2072 };
2073
2074 let mut ops_out: Vec<PublishOp> = Vec::new();
2075 let mut last_op_id: Option<lex_vcs::OpId> = None;
2076 for kind in op_kinds {
2077 // Persist the underlying stage AST/metadata if this op
2078 // produces or replaces one.
2079 if let Some(stg) = stage_for_kind(&kind, stages) {
2080 if !matches!(stg, lex_ast::Stage::Import(_)) {
2081 self.publish_signed(stg, signer)?;
2082 if activate {
2083 if let Some(stage_id_str) = stage_id(stg) {
2084 let _ = self.activate(&stage_id_str);
2085 }
2086 }
2087 }
2088 }
2089 let transition = transition_for_kind(&kind);
2090 // #992: a publish holds every stage it binds, so each pair the op
2091 // leaves at the head must now be on disk. If it is not, the op and
2092 // the stage it wrote disagree about the sig — refuse rather than
2093 // advance onto a head no render can resolve. (The general gate in
2094 // `cas_retry_advance` must tolerate a merely absent stage; this
2095 // path has no such excuse.)
2096 for (sig, stage) in bound_pairs(&transition) {
2097 if !self.stage_file_exists(&sig, &stage) {
2098 return Err(match self.unsatisfiable_owner(&sig, &stage) {
2099 Some(owner) => StoreError::UnsatisfiablePair {
2100 sig_id: sig,
2101 stage_id: stage,
2102 filed_under: owner,
2103 },
2104 None => StoreError::UnknownStage(stage),
2105 });
2106 }
2107 }
2108 let attestable = attestable_stage_ids(&transition);
2109 let head_now = self.get_branch(branch)?.and_then(|b| b.head_op);
2110 let op =
2111 lex_vcs::Operation::new(kind.clone(), head_now.into_iter().collect::<Vec<_>>());
2112 // #131 / #839: stamp the caller's intent so the op log records
2113 // why this change happened, not just what it was. The CAS
2114 // retry path preserves `intent_id` when it rebuilds the op.
2115 let op = match &intent_id {
2116 Some(id) => op.with_intent(id.clone()),
2117 None => op,
2118 };
2119 let op_id = self.apply_operation(branch, op, transition)?;
2120 self.record_typecheck_passed(&attestable, &op_id)?;
2121 ops_out.push(PublishOp {
2122 op_id: op_id.clone(),
2123 kind: serde_json::to_value(&kind).map_err(StoreError::Serde)?,
2124 });
2125 last_op_id = Some(op_id);
2126 }
2127
2128 // Record the documentation of *every* declaration in this program, not
2129 // just the ones that produced an op.
2130 //
2131 // Comments live outside the hash, so a declaration whose code is
2132 // unchanged emits no op — `publish_signed` never runs for it, and its
2133 // metadata keeps whatever doc it had. That is exactly the shape of a
2134 // re-publish that exists only to carry documentation: the first attempt
2135 // emitted 2 ops for a 25-declaration package and the other 23 stayed
2136 // undocumented, so the release was as bare as the one it replaced.
2137 //
2138 // Metadata is not content-addressed, so this is a cheap in-place write
2139 // and a no-op when the doc already matches.
2140 for stg in stages {
2141 let doc = stage_doc(stg);
2142 if doc.is_empty() {
2143 continue;
2144 }
2145 let (Some(sig), Some(sid)) = (sig_id(stg), stage_id(stg)) else { continue };
2146 let meta_path = self.impl_dir(&sig).join(format!("{sid}.metadata.json"));
2147 if !meta_path.exists() {
2148 continue;
2149 }
2150 if let Ok(mut existing) = self.get_metadata_for_sig(&sig, &sid) {
2151 if existing.doc != doc {
2152 existing.doc = doc;
2153 write_canonical_json(&meta_path, &existing)?;
2154 }
2155 }
2156 }
2157
2158 let head_op = match last_op_id {
2159 Some(id) => Some(id),
2160 // No ops applied; return whatever the head was already.
2161 None => self.get_branch(branch)?.and_then(|b| b.head_op),
2162 };
2163
2164 Ok(PublishOutcome {
2165 ops: ops_out,
2166 head_op,
2167 })
2168 }
2169
2170 /// The sig `stage` is actually filed under, when binding it to `sig` is
2171 /// **provably** unsatisfiable (#992): `(sig, stage)` cannot be read, yet
2172 /// the store holds that very stage under a different sig. `None` when the
2173 /// pair reads fine, and also when the stage is simply absent — an absent
2174 /// blob (mid-pull, partial clone) proves nothing about satisfiability.
2175 ///
2176 /// Presence is judged by the stage file existing under the sig (a full
2177 /// snapshot or a delta), not by decoding it: this runs over every head
2178 /// entry on a ref advance, and a stat is all the question needs.
2179 pub fn unsatisfiable_owner(&self, sig: &str, stage: &str) -> Option<String> {
2180 if self.stage_file_exists(sig, stage) {
2181 return None;
2182 }
2183 let (owner, _) = self.lookup_lifecycle(stage).ok()?;
2184 if owner == sig || !self.stage_file_exists(&owner, stage) {
2185 return None;
2186 }
2187 Some(owner)
2188 }
2189
2190 fn stage_file_exists(&self, sig: &str, stage: &str) -> bool {
2191 let dir = self.impl_dir(sig);
2192 dir.join(format!("{stage}.ast.json")).exists()
2193 || dir.join(format!("{stage}.delta.json")).exists()
2194 }
2195
2196 /// Refuse the first provably unsatisfiable pair in `pairs` with
2197 /// [`StoreError::UnsatisfiablePair`] — the write-time half of #992's
2198 /// "always-valid HEAD". See [`Self::unsatisfiable_owner`] for what counts.
2199 pub fn check_pairs_satisfiable<'a>(
2200 &self,
2201 pairs: impl IntoIterator<Item = (&'a String, &'a String)>,
2202 ) -> Result<(), StoreError> {
2203 for (sig, stage) in pairs {
2204 if let Some(owner) = self.unsatisfiable_owner(sig, stage) {
2205 return Err(StoreError::UnsatisfiablePair {
2206 sig_id: sig.clone(),
2207 stage_id: stage.clone(),
2208 filed_under: owner,
2209 });
2210 }
2211 }
2212 Ok(())
2213 }
2214
2215 /// Head entries naming a `(sig, stage)` no store can hold, where that same
2216 /// stage **is** readable under a different sig at the head (#992).
2217 ///
2218 /// This is the fingerprint a pre-#992 `ChangeEffectSig` leaves behind. It
2219 /// bound the *old* sig to the *new* stage, while the store files an
2220 /// implementation under the sig its own AST hashes to — and that AST
2221 /// declares the new effects. The head then names one declaration twice and
2222 /// one of the two can never be resolved, so every render of that head
2223 /// fails and any release cut from it is born broken (`lex-web@0.4.0`).
2224 ///
2225 /// Fixing the op stops new ones appearing but cannot repair a head that
2226 /// already has one: the publish diff is keyed by declaration name and
2227 /// reads the old side through the ASTs it *can* resolve, so a stranded
2228 /// entry is invisible to it and no op is ever emitted.
2229 ///
2230 /// **The second condition is the whole reason this is safe to do without
2231 /// asking.** Requiring the stage to be resolvable under another sig proves
2232 /// the content is present and merely filed elsewhere, so retiring the
2233 /// entry discards nothing. A store that is simply missing blobs — mid-pull,
2234 /// a partial clone, a GC'd object — fails that test, because there neither
2235 /// sig resolves, and is left strictly alone.
2236 pub fn stranded_head_entries(
2237 &self,
2238 head_pairs: &[(String, String)],
2239 ) -> Vec<lex_vcs::OperationKind> {
2240 let asts = self.get_asts_for_sigs_bulk(head_pairs);
2241
2242 head_pairs
2243 .iter()
2244 .zip(asts.iter())
2245 .filter(|(_, ast)| ast.is_err())
2246 .filter_map(|((sig, stage), _)| {
2247 // Is this stage readable under *any* sig in the store?
2248 //
2249 // Deliberately store-wide rather than head-scoped. Scoping it
2250 // to the head looked tighter but made the repair miss the case
2251 // it exists for: republishing re-adds the declaration under a
2252 // freshly computed stage, so the sig that owns the orphaned AST
2253 // no longer points at it *from the head*. The content sits
2254 // right there on disk and a head-scoped check cannot see it,
2255 // which is exactly what happened on the first real `lex-web`
2256 // attempt. `get_ast` resolves a stage through the store's
2257 // index, which is the question actually being asked: does some
2258 // sig here hold this?
2259 //
2260 // The guarantee is unchanged, and it is the entire safety
2261 // argument: retiring the entry discards nothing, because the
2262 // AST demonstrably still exists. A store merely missing blobs
2263 // fails this and is left strictly alone.
2264 let twin = self.get_ast(stage).ok()?;
2265 Some(match &twin {
2266 lex_ast::Stage::TypeDecl(_) => lex_vcs::OperationKind::RemoveType {
2267 sig_id: sig.clone(),
2268 last_stage_id: stage.clone(),
2269 },
2270 _ => lex_vcs::OperationKind::RemoveFunction {
2271 sig_id: sig.clone(),
2272 last_stage_id: stage.clone(),
2273 },
2274 })
2275 })
2276 .collect()
2277 }
2278
2279 pub fn derive_imports_from_oplog(
2280 &self,
2281 branch: &str,
2282 ) -> Result<lex_vcs::ImportMap, StoreError> {
2283 use lex_vcs::OperationKind::*;
2284 let log = lex_vcs::OpLog::open(self.root())?;
2285 let head = match self.get_branch(branch)?.and_then(|b| b.head_op) {
2286 Some(h) => h,
2287 None => return Ok(Default::default()),
2288 };
2289 let mut out: lex_vcs::ImportMap = Default::default();
2290 for r in log.walk_forward(&head, None)? {
2291 match r.op.kind {
2292 AddImport { in_file, module, alias } => {
2293 // The op omits the alias when it's the module's
2294 // default (last path segment) to keep its OpId
2295 // stable; rebuild it the same way on the way out.
2296 let alias =
2297 alias.unwrap_or_else(|| lex_vcs::default_import_alias(&module));
2298 out.entry(in_file)
2299 .or_default()
2300 .insert(lex_vcs::ImportRef { reference: module, alias });
2301 }
2302 RemoveImport { in_file, module } => {
2303 // Removal is keyed by reference (the op carries no
2304 // alias), so drop any binding of that module.
2305 if let Some(set) = out.get_mut(&in_file) {
2306 set.retain(|ir| ir.reference != module);
2307 }
2308 }
2309 _ => {}
2310 }
2311 }
2312 Ok(out)
2313 }
2314
2315 /// Apply an operation to a branch and advance its head_op.
2316 ///
2317 /// The single advance path. Validates parents via `lex_vcs::apply`,
2318 /// persists the operation via the op log, then atomically advances
2319 /// the branch file's head_op via `set_branch_head_op`.
2320 ///
2321 /// Errors:
2322 /// - `UnknownBranch`: branch does not exist (no op is persisted).
2323 /// - `Apply(ApplyError::StaleParent)`: the op's parents don't
2324 /// match the branch head — head is unchanged. Callers that
2325 /// want retry-on-stale (e.g. `lex publish` re-running against
2326 /// a moved head) match on this variant explicitly.
2327 /// - `Apply(ApplyError::UnknownMergeParent)`: a merge op's
2328 /// second parent isn't in the log.
2329 /// - `Io`: filesystem error during persist or branch advance.
2330 ///
2331 /// Crash recovery: between op persist and branch advance, a crash
2332 /// can leave an orphan op record in the log with no branch
2333 /// pointing at it. The op is content-addressed and cheap to
2334 /// re-derive from the same source. See
2335 /// Apply a single op against `branch`, gated on the candidate
2336 /// program typechecking. The per-op variant of #130's
2337 /// write-time gate — counterpart to [`Self::publish_program`]'s
2338 /// batch-mode check.
2339 ///
2340 /// `candidate` is the sequence of `Stage`s that *would* exist
2341 /// on this branch after the op is applied. Caller's
2342 /// responsibility: today neither `lex-store` nor `lex-vcs`
2343 /// reconstruct the candidate from the op + branch state on
2344 /// behalf of the caller. The natural callers (HTTP `POST
2345 /// /v1/publish` for a single op; agent harnesses driving
2346 /// merges via the future #134 API) already have the candidate
2347 /// in memory.
2348 ///
2349 /// On rejection: branch head unchanged, no op record persisted.
2350 /// Same atomicity guarantee as the publish path.
2351 ///
2352 /// # Why a separate method, not a flag on `apply_operation`
2353 ///
2354 /// `apply_operation` accepting `Option<&[Stage]>` and silently
2355 /// skipping the gate on `None` is exactly the kind of
2356 /// "secretly opt-out" path #130 is trying to remove. The honest
2357 /// split: `apply_operation` for the one caller that already
2358 /// typechecked its input up front (`publish_program`),
2359 /// `apply_operation_checked` for callers holding the candidate,
2360 /// [`Self::apply_operation_gated`] for single-parent callers
2361 /// that hold only the transition (`/v1/patch`), and
2362 /// [`Self::apply_merge_op_gated`] for merge commits (#833).
2363 pub fn apply_operation_checked(
2364 &self,
2365 branch: &str,
2366 op: lex_vcs::Operation,
2367 transition: lex_vcs::StageTransition,
2368 candidate: &[lex_ast::Stage],
2369 ) -> Result<lex_vcs::OpId, StoreError> {
2370 self.apply_operation_checked_with_intent(branch, op, transition, candidate, None)
2371 }
2372
2373 /// [`Self::apply_operation_checked`] that also records `intent` in the
2374 /// [`lex_vcs::IntentLog`] (#837 piece A). The caller stamps the intent's id on
2375 /// `op` (`Operation::with_intent`); this only makes sure the intent record
2376 /// exists, and does so *after* the type-check and session-budget gates and
2377 /// before the head moves — so a rejected write leaves no op **and no
2378 /// intent** behind, the same "no footprint" guarantee the gate already
2379 /// gives the op record. `None` is exactly `apply_operation_checked`.
2380 pub fn apply_operation_checked_with_intent(
2381 &self,
2382 branch: &str,
2383 op: lex_vcs::Operation,
2384 transition: lex_vcs::StageTransition,
2385 candidate: &[lex_ast::Stage],
2386 intent: Option<&lex_vcs::Intent>,
2387 ) -> Result<lex_vcs::OpId, StoreError> {
2388 // #945: the candidate comes from the head's SigId→stage map, which holds
2389 // only fn/type declarations — the head's `import` edges are absent. Add
2390 // them and resolve against the head the op applies to (whose committed
2391 // lock pins the dependencies), so a non-inlined head type-checks on this
2392 // write path exactly as it does on publish and on the hub's own gate.
2393 let base_head = self.get_branch(branch).ok().flatten().and_then(|b| b.head_op);
2394 let stages = self.with_head_imports(base_head.as_deref(), candidate.to_vec());
2395 if let Err(errors) = self.check_with_resolved_deps(&stages, base_head.as_deref()) {
2396 // #281: emit a `RepairHint` attestation against each
2397 // candidate stage the transition was about to produce.
2398 // The op record itself isn't persisted (the gate is
2399 // pre-persistence), but the candidate stage IS — the
2400 // transform-flow methods publish before this call.
2401 // The attached hint lets `lex repair <op_id>` and
2402 // future LLM-assisted apply paths read the structured
2403 // errors without re-running the typecheck.
2404 let attestable = attestable_stage_ids(&transition);
2405 let failed_op_id = op.op_id();
2406 let _ = self.record_repair_hint(&attestable, &failed_op_id, &errors);
2407 return Err(StoreError::TypeError(errors));
2408 }
2409 // #292 slice 3: per-session budget gate. After typecheck
2410 // passes, refuse the op if it would push its session's
2411 // monotonic spend over the configured cap. Sessions
2412 // without an intent_id, or with an intent whose session
2413 // has no cap configured, sail through.
2414 self.check_session_budget(&op)?;
2415 if let Some(intent) = intent {
2416 lex_vcs::IntentLog::open(self.root())?.put(intent)?;
2417 }
2418 let attestable = attestable_stage_ids(&transition);
2419 let op_effects = op_declared_effects(&op.kind);
2420 // #262: CAS retry loop. Single-parent ops can be safely
2421 // re-persisted under a new parent on contention (the kind
2422 // is invariant; only `parents` changes). Merge ops (already
2423 // 2-parent) come through the merge engine which has its own
2424 // coordination; we don't retry them here — we'll see the
2425 // first attempt's CAS fail and surface Contention.
2426 self.cas_retry_advance(branch, op, transition, |new_head| {
2427 self.record_typecheck_passed(&attestable, &new_head.op_id)?;
2428 self.run_required_attestations_gate(branch, &new_head.op_id, &attestable, &op_effects)
2429 })
2430 }
2431
2432 /// The program that would exist on `branch` after `transition`
2433 /// is applied: the branch head (snapshot-cached) with the
2434 /// transition replayed over it, every resulting `(sig, stage)`
2435 /// bulk-loaded. Exact for a **single-parent** transition — the
2436 /// candidate [`Self::apply_operation_gated`] wants. Not valid for
2437 /// a merge: a `StageTransition::Merge` pins only the sigs the
2438 /// merge decided, while the op-DAG replay that computes a
2439 /// merge's real head walks both parents (#833).
2440 pub fn candidate_program_for(
2441 &self,
2442 branch: &str,
2443 transition: &lex_vcs::StageTransition,
2444 ) -> Result<Vec<Stage>, StoreError> {
2445 let mut head = self.branch_head(branch)?;
2446 crate::branches::apply_transition(&mut head, transition);
2447 let pairs: Vec<(String, String)> = head.into_iter().collect();
2448 self.get_asts_for_sigs_bulk(&pairs).into_iter().collect()
2449 }
2450
2451 /// [`Self::apply_operation_checked`] for a **single-parent** op
2452 /// where the caller holds only the transition: assembles the
2453 /// candidate via [`Self::candidate_program_for`] and runs the
2454 /// gate. Same rejection semantics — `TypeError`, a `RepairHint`
2455 /// attestation, head unchanged, nothing persisted. This is the
2456 /// write path for `/v1/patch` (#833). Merge ops must not use it
2457 /// (see `candidate_program_for`); they go through
2458 /// [`Self::apply_merge_op_gated`].
2459 pub fn apply_operation_gated(
2460 &self,
2461 branch: &str,
2462 op: lex_vcs::Operation,
2463 transition: lex_vcs::StageTransition,
2464 ) -> Result<lex_vcs::OpId, StoreError> {
2465 self.apply_operation_gated_with_intent(branch, op, transition, None)
2466 }
2467
2468 /// [`Self::apply_operation_gated`] that also records `intent` (see
2469 /// [`Self::apply_operation_checked_with_intent`]) — the write path for an
2470 /// intent-carrying `/v1/patch` (#837 piece A).
2471 pub fn apply_operation_gated_with_intent(
2472 &self,
2473 branch: &str,
2474 op: lex_vcs::Operation,
2475 transition: lex_vcs::StageTransition,
2476 intent: Option<&lex_vcs::Intent>,
2477 ) -> Result<lex_vcs::OpId, StoreError> {
2478 debug_assert!(
2479 op.parents.len() <= 1,
2480 "apply_operation_gated is single-parent only; merges use apply_merge_op_gated"
2481 );
2482 let candidate = self.candidate_program_for(branch, &transition)?;
2483 self.apply_operation_checked_with_intent(branch, op, transition, &candidate, intent)
2484 }
2485
2486 /// The gated write path for **merge** commits (`commit_merge`,
2487 /// `POST /v1/merge/<id>/commit`, `lex merge commit`).
2488 ///
2489 /// A `StageTransition::Merge` pins only the sigs the merge decided
2490 /// (#1062); the sig->stage map every consumer reads is recomputed by
2491 /// replaying the op DAG, which for a merge walks *both* parents
2492 /// and can surface sigs the entries never mention. So the only way
2493 /// to know the true post-merge program is to replay it — land the
2494 /// op and read `branch_head`. This lands the merge op,
2495 /// type-checks the resulting head, and on a failure rolls the
2496 /// head back and returns `TypeError`.
2497 ///
2498 /// Before #833 the merge paths landed through the ungated
2499 /// `apply_operation`, so a merge whose result didn't compose
2500 /// (e.g. dst still calls `helper`, an agent-supplied resolution
2501 /// dropped it) advanced the head with nothing to catch it.
2502 ///
2503 /// Rollback leaves the rejected merge op as an unreachable record
2504 /// (reclaimed by `lex op gc`, the same orphan crash-recovery
2505 /// already tolerates). A stage the merge names that was never
2506 /// published surfaces as the underlying `StoreError` from the
2507 /// bulk read — the "never advance onto content that can't be
2508 /// loaded" invariant from the other side.
2509 pub fn apply_merge_op_gated(
2510 &self,
2511 branch: &str,
2512 op: lex_vcs::Operation,
2513 transition: lex_vcs::StageTransition,
2514 ) -> Result<lex_vcs::OpId, StoreError> {
2515 let head_before = self.get_branch(branch)?.and_then(|b| b.head_op);
2516 // Capture the stages this merge introduces before `transition`
2517 // is moved into `apply_operation`; used for the TypeCheck
2518 // attestation below and for the attestation gate. A merge's entries
2519 // pin every sig it decided (#1062), including the ones dst already
2520 // holds; those stages are not introduced by the merge, so they are
2521 // neither gated nor re-attested — exactly what the entries that
2522 // differed from dst were before pinning.
2523 let attestable = match &transition {
2524 lex_vcs::StageTransition::Merge { entries } => {
2525 let dst = self.branch_head(branch)?;
2526 entries
2527 .iter()
2528 .filter(|(sig, stage)| stage.is_some() && dst.get(*sig) != stage.as_ref())
2529 .filter_map(|(_, stage)| stage.clone())
2530 .collect()
2531 }
2532 other => attestable_stage_ids(other),
2533 };
2534 // #977: the merge commits its OWN lock — the union of its parents'
2535 // pins — so a dependency either branch introduced resolves at the
2536 // merged head. Computed before anything is written: a same-package /
2537 // different-version conflict refuses the merge with no side effect
2538 // (always-valid HEAD), rather than landing the op and rolling back.
2539 //
2540 // `Operation::new` sorts `parents` by op id, so their order says
2541 // nothing about which side is `dst`. The destination branch's current
2542 // head does: a merge op can only land when that head is one of its
2543 // parents, so put it first and a conflict names the sides correctly.
2544 // If it is *not* among a multi-parent op's parents the op is stale
2545 // and `apply_operation` refuses it below (`StaleParent`); computing a
2546 // lock with unknown sides would only mislabel an error, so skip it.
2547 let dst_pos = head_before.as_ref().and_then(|d| op.parents.iter().position(|p| p == d));
2548 let merged_lock = match dst_pos {
2549 Some(pos) => {
2550 let mut ordered = op.parents.clone();
2551 let dst_parent = ordered.remove(pos);
2552 ordered.insert(0, dst_parent);
2553 self.merged_lock_for_parents(&ordered)?
2554 }
2555 // No dst head (merging into an empty branch): the lone parent is
2556 // the source, whose governing lock the merge simply carries.
2557 None if head_before.is_none() => self.merged_lock_for_parents(&op.parents)?,
2558 None => None,
2559 };
2560 let op_id = self.apply_operation_attesting(branch, op, transition, attestable.clone())?;
2561
2562 let verdict = (|| -> Result<(), StoreError> {
2563 // Bind the merged lock to the merge op before gating, so the gate
2564 // (and every later re-verification of this head) resolves against
2565 // the merge's own pins instead of inheriting one parent's whole.
2566 // If the gate then rejects the head, the ref stays bound to an
2567 // op no branch points at — harmless, and content-correct for it.
2568 if let Some(lock) = &merged_lock {
2569 self.set_committed_lock(op_id.as_str(), lock)?;
2570 }
2571 let head = self.branch_head(branch)?;
2572 let pairs: Vec<(String, String)> = head.into_iter().collect();
2573 let decls: Vec<Stage> =
2574 self.get_asts_for_sigs_bulk(&pairs).into_iter().collect::<Result<_, _>>()?;
2575 // #945: the merge op is already applied, so the post-merge head IS
2576 // `op_id` — resolve this head's dependencies against it (its
2577 // committed lock) and prepend its `import` edges, which the
2578 // SigId→stage map omits. Without this a non-inlined head failed the
2579 // merge gate as `unknown_identifier <alias>` even with the
2580 // dependency correctly resolvable.
2581 let merged_head = op_id.as_str();
2582 let stages = self.with_head_imports(Some(merged_head), decls);
2583 if let Err(errors) = self.check_with_resolved_deps(&stages, Some(merged_head)) {
2584 return Err(StoreError::TypeError(errors));
2585 }
2586 Ok(())
2587 })();
2588
2589 if let Err(e) = verdict {
2590 // Roll the head back. The empty-dst case never reaches
2591 // here (it fast-forwards without a merge op), so
2592 // `head_before` is always `Some` on this arm.
2593 if let Some(prev) = head_before {
2594 self.set_branch_head_op(branch, prev)?;
2595 }
2596 return Err(e);
2597 }
2598 // #835: the merge's post-merge head type-checked, but until now
2599 // that verdict left no trace in the attestation log — so a
2600 // merged stage looked un-type-checked to `lex blame
2601 // --with-evidence` and the attestation queries, unlike a
2602 // published or patched stage. Emit `TypeCheck::Passed` for the
2603 // stages the merge introduced, mirroring the publish / patch
2604 // paths (`record_typecheck_passed`). Emitted only after the
2605 // check passes and the head is committed, so a rolled-back
2606 // merge records nothing.
2607 self.record_typecheck_passed(&attestable, &op_id)?;
2608 Ok(op_id)
2609 }
2610
2611 /// Like [`Self::apply_merge_op_gated`], but also appends the
2612 /// `SetFiles` a disagreeing-manifest merge needs (#1007 §1 / PR 7).
2613 ///
2614 /// `manifest` is the already-computed, already-stored merged manifest
2615 /// blob id — see [`crate::files::manifest_merge`] plus a merge
2616 /// session's file-conflict resolutions for how the caller builds it.
2617 /// `None` when dst's and src's manifests already agreed (the common
2618 /// case): nothing to record, behaves exactly like
2619 /// `apply_merge_op_gated`.
2620 ///
2621 /// On any failure — the sig-level type-check gate (as before), or the
2622 /// follow-up `SetFiles` (which should essentially never fail here
2623 /// since the caller already validated the manifest before calling
2624 /// this, but a store can be modified concurrently) — the branch head
2625 /// is rolled all the way back to where it was before this call. A
2626 /// merge op is never left standing as a head with an `Ambiguous`
2627 /// manifest; that would violate the always-valid-HEAD invariant
2628 /// `check_head_files` polices for every other path onto a head.
2629 pub fn apply_merge_op_gated_with_manifest(
2630 &self,
2631 branch: &str,
2632 op: lex_vcs::Operation,
2633 transition: lex_vcs::StageTransition,
2634 manifest: Option<&str>,
2635 intent_id: Option<&lex_vcs::IntentId>,
2636 ) -> Result<lex_vcs::OpId, StoreError> {
2637 let head_before = self.get_branch(branch)?.and_then(|b| b.head_op);
2638 let merge_op_id = self.apply_merge_op_gated(branch, op, transition)?;
2639 let Some(manifest) = manifest else {
2640 return Ok(merge_op_id);
2641 };
2642 match self.apply_set_files(branch, manifest, intent_id) {
2643 Ok(set_files_op_id) => Ok(set_files_op_id),
2644 Err(e) => {
2645 if let Some(prev) = head_before {
2646 self.set_branch_head_op(branch, prev)?;
2647 }
2648 Err(e)
2649 }
2650 }
2651 }
2652
2653 /// Type-check the program that would result from overlaying a merge
2654 /// `delta` onto `branch`'s current head — **without moving the
2655 /// head** (#834). `delta` maps `sig_id -> Some(stage)` to set that
2656 /// sig to `stage`, or `sig_id -> None` to remove it, exactly the
2657 /// `entries` a `StageTransition::Merge` records.
2658 ///
2659 /// This is the read-only, resolve-time counterpart of
2660 /// `apply_merge_op_gated`'s commit-time gate: it lets a merge
2661 /// session tell an agent *which resolution broke type-checking* the
2662 /// moment it is submitted, instead of only after a failed commit.
2663 /// `Ok(())` means the projected program composes; a type failure is
2664 /// `Err(StoreError::TypeError(..))`; a read failure is the
2665 /// corresponding `StoreError` I/O variant.
2666 pub fn typecheck_merge_projection(
2667 &self,
2668 branch: &str,
2669 delta: &std::collections::BTreeMap<String, Option<String>>,
2670 ) -> Result<(), StoreError> {
2671 let mut head = self.branch_head(branch)?;
2672 for (sig, stage) in delta {
2673 match stage {
2674 Some(s) => { head.insert(sig.clone(), s.clone()); }
2675 None => { head.remove(sig); }
2676 }
2677 }
2678 let pairs: Vec<(String, String)> = head.into_iter().collect();
2679 let decls: Vec<Stage> =
2680 self.get_asts_for_sigs_bulk(&pairs).into_iter().collect::<Result<_, _>>()?;
2681 // #945: a projection has no committed op of its own, so resolve against
2682 // the head it is projected ONTO — that head's committed lock pins the
2683 // dependencies, and its `import` edges are the projection's imports too
2684 // (a merge delta maps sig→stage, so it never adds or removes an import).
2685 // Without this, resolve-time checking of a non-inlined head reported a
2686 // bogus `unknown_identifier <alias>` and rejected valid resolutions.
2687 let base_head = self.get_branch(branch).ok().flatten().and_then(|b| b.head_op);
2688 let stages = self.with_head_imports(base_head.as_deref(), decls);
2689 if let Err(errors) = self.check_with_resolved_deps(&stages, base_head.as_deref()) {
2690 return Err(StoreError::TypeError(errors));
2691 }
2692 Ok(())
2693 }
2694
2695 /// #838: attempt a typed three-way merge of a single sig's body for
2696 /// a `ModifyModify` conflict — the intra-function, better-than-git
2697 /// case where two agents edited *disjoint* subtrees of the same
2698 /// function (different match arms, different let bindings).
2699 ///
2700 /// `base` / `ours` (the dst side) / `theirs` (the src side) are the
2701 /// three stage ids the merge engine surfaced for `sig_id`. Loads
2702 /// the three `FnDecl`s, structurally merges the bodies
2703 /// ([`lex_vcs::merge_bodies`]), and accepts the result *only if* the
2704 /// merged function also type-checks against `dst_branch`'s head — a
2705 /// body that composes syntactically but not by type is still a
2706 /// conflict (#838). On success the merged stage is published
2707 /// (content-addressed, idempotent; orphaned and GC-reclaimable if
2708 /// the merge is never committed) and its id returned; `None` means
2709 /// "fall back to a whole-function conflict."
2710 ///
2711 /// Deliberately narrow for this slice: only pure body divergence is
2712 /// merged. If the two sides disagree on anything but the body
2713 /// (examples, type params — the signature is identical by
2714 /// construction, since all three share `sig_id`), or either stage
2715 /// isn't a function, it falls back to a conflict.
2716 pub fn try_semantic_body_merge(
2717 &self,
2718 dst_branch: &str,
2719 sig_id: &str,
2720 base: &str,
2721 ours: &str,
2722 theirs: &str,
2723 ) -> Result<Option<String>, StoreError> {
2724 use lex_ast::Stage::FnDecl;
2725 let (base_fd, ours_fd, theirs_fd) =
2726 match (self.get_ast(base), self.get_ast(ours), self.get_ast(theirs)) {
2727 (Ok(FnDecl(b)), Ok(FnDecl(o)), Ok(FnDecl(t))) => (b, o, t),
2728 // A non-function stage (type decl / import) or a stage
2729 // that can't be loaded isn't an intra-body merge.
2730 _ => return Ok(None),
2731 };
2732
2733 // Only the body may diverge between the two sides.
2734 if !fndecl_same_except_body(&ours_fd, &theirs_fd) {
2735 return Ok(None);
2736 }
2737
2738 let merged_body =
2739 match lex_vcs::merge_bodies(&base_fd.body, &ours_fd.body, &theirs_fd.body) {
2740 lex_vcs::BodyMerge::Merged(b) => b,
2741 lex_vcs::BodyMerge::Conflict => return Ok(None),
2742 };
2743
2744 let mut merged_fd = ours_fd.clone();
2745 merged_fd.body = merged_body;
2746 let merged_stage = lex_ast::Stage::FnDecl(merged_fd);
2747 let new_stage_id = match stage_id(&merged_stage) {
2748 Some(id) => id,
2749 None => return Ok(None),
2750 };
2751
2752 // Type-check the merged fn in context: dst's head with this sig
2753 // swapped to the merged stage. Requires the merged stage to be
2754 // loadable, so publish first (idempotent, content-addressed).
2755 self.publish(&merged_stage)?;
2756 let mut delta = std::collections::BTreeMap::new();
2757 delta.insert(sig_id.to_string(), Some(new_stage_id.clone()));
2758 match self.typecheck_merge_projection(dst_branch, &delta) {
2759 Ok(()) => Ok(Some(new_stage_id)),
2760 // Composes syntactically, not by type → still a conflict.
2761 Err(StoreError::TypeError(_)) => Ok(None),
2762 Err(e) => Err(e),
2763 }
2764 }
2765
2766 /// #836 G3: assemble everything a regenerator needs to *replay* an
2767 /// op — re-derive the change from its recorded cause. Returns the
2768 /// op's recorded intent (prompt / model / session), the target sig
2769 /// and the stage id it produced, and the program the change was
2770 /// made against (the parent state, rendered to source). An external
2771 /// harness feeds the prompt + parent program to the recorded model,
2772 /// then hands the regenerated stage back to [`Self::replay_compare`]
2773 /// (lex owns the deterministic comparison; the model call is the
2774 /// harness's, matching the rest of the architecture).
2775 ///
2776 /// Errors with `UnknownOp` if the op_id is unknown, or
2777 /// `InvalidTransition` if the op didn't produce a stage (a removal /
2778 /// import / merge has nothing to regenerate). A `SetFiles` op (#1007)
2779 /// is refused by kind with `NotReplayable::Files` — typed, so replay
2780 /// coverage can leave it out instead of counting it as a miss.
2781 pub fn replay_request(&self, op_id: &str) -> Result<ReplayRequest, StoreError> {
2782 let log = lex_vcs::OpLog::open(self.root())?;
2783 let record = log
2784 .get(&op_id.to_string())?
2785 .ok_or_else(|| StoreError::UnknownOp(op_id.to_string()))?;
2786 crate::files::refuse_non_semantic(&record)?;
2787 let (target_sig, expected_stage_id) = produced_sig_stage(&record.produces)
2788 .ok_or_else(|| StoreError::InvalidTransition(format!("op {op_id} produced no stage to replay")))?;
2789
2790 let (prompt, model, session_id) = match &record.op.intent_id {
2791 Some(id) => {
2792 let intents = lex_vcs::IntentLog::open(self.root())?;
2793 match intents.get(id)? {
2794 Some(i) => (Some(i.prompt), Some(model_label(&i.model)), Some(i.session_id)),
2795 None => (None, None, None),
2796 }
2797 }
2798 None => (None, None, None),
2799 };
2800
2801 // The program the op was applied against: the head state at its
2802 // (first) parent, rendered with the canonical printer. A root
2803 // op has no parent → empty program.
2804 // The op's own recorded stage is the one thing replay cannot do
2805 // without — it's what a regeneration is compared against — so an
2806 // unloadable target fails clearly here (#868), before any skipping.
2807 let target_stage = self.load_replay_target(op_id, &target_sig, &expected_stage_id)?;
2808
2809 // #868: reconstruction skips (and reports) head declarations the
2810 // store can't load instead of failing the whole replay — a long-lived
2811 // history with GC'd or superseded intermediates is exactly where
2812 // replay-as-verification is most useful.
2813 let (parent_program, mut skipped) = match record.op.parents.first() {
2814 Some(parent) => {
2815 let rec = self.demangled_program_at_op_skipping(parent)?;
2816 (lex_ast::print_stages(&rec.stages), rec.skipped)
2817 }
2818 None => (String::new(), Vec::new()),
2819 };
2820 if !skipped.is_empty() {
2821 let referenced = referenced_names(&target_stage);
2822 for sk in &mut skipped {
2823 // The target's own sig in the parent is its *previous*
2824 // version (a modify), not a callee.
2825 sk.called_by_target = sk.sig_id != target_sig
2826 && sk.name.as_deref().is_some_and(|n| referenced.contains(n));
2827 }
2828 }
2829
2830 // The target function's name + signature, from the recorded stage — a
2831 // regenerator needs the interface, not just the hash.
2832 //
2833 // #980: report the *bare* name. A package publish mangles the recorded
2834 // declaration (`lib_<hash>.twice`), and a dotted name is not valid Lex
2835 // source — so the old, mangled `target_signature` asked regenerators
2836 // for something unwritable, and every package-published op replayed as
2837 // a false negative.
2838 let (target_name, target_signature) = match &target_stage {
2839 lex_ast::Stage::FnDecl(fd) => {
2840 let mut bare = fd.clone();
2841 bare.name = demangled_name(&fd.name).to_string();
2842 (Some(bare.name.clone()), Some(lex_vcs::render_signature(&bare)))
2843 }
2844 _ => (None, None),
2845 };
2846
2847 Ok(ReplayRequest {
2848 op_id: op_id.to_string(),
2849 target_sig,
2850 target_name,
2851 target_signature,
2852 expected_stage_id,
2853 prompt,
2854 model,
2855 session_id,
2856 parent_program,
2857 skipped,
2858 })
2859 }
2860
2861 /// Load the stage a replayable op produced, or fail with
2862 /// [`StoreError::ReplayTargetUnloadable`] (#868). Tries the stage-id
2863 /// lookup first (what replay always used), then the `(sig, stage)` read.
2864 fn load_replay_target(&self, op_id: &str, sig: &str, stage_id: &str) -> Result<Stage, StoreError> {
2865 match self.get_ast(stage_id) {
2866 Ok(st) => Ok(st),
2867 Err(first) => {
2868 let pair = [(sig.to_string(), stage_id.to_string())];
2869 match self.get_asts_for_sigs_bulk(&pair).pop() {
2870 Some(Ok(st)) => Ok(st),
2871 _ => Err(StoreError::ReplayTargetUnloadable {
2872 op_id: op_id.to_string(),
2873 stage_id: stage_id.to_string(),
2874 reason: first.to_string(),
2875 }),
2876 }
2877 }
2878 }
2879 }
2880
2881 /// Load the declarations of a reconstructed head, in `pairs` order.
2882 ///
2883 /// Strict (`skip_unloadable = false`): the first unloadable stage is an
2884 /// error — for callers that need a complete head. Lenient: *absent* stages
2885 /// are dropped and reported as [`SkippedStage`]s (#868); a stage that is
2886 /// present but unreadable is still [`StoreError::StageUnreadable`].
2887 pub(crate) fn load_head_decls(
2888 &self,
2889 pairs: &[(String, String)],
2890 skip_unloadable: bool,
2891 ) -> Result<(Vec<Stage>, Vec<SkippedStage>), StoreError> {
2892 let mut decls = Vec::with_capacity(pairs.len());
2893 let mut skipped = Vec::new();
2894 for ((sig_id, stage_id), ast) in pairs.iter().zip(self.get_asts_for_sigs_bulk(pairs)) {
2895 match ast {
2896 Ok(st) => decls.push(st),
2897 // Only a genuinely *absent* stage is skipped. Anything else
2898 // (corrupt JSON, permissions, other I/O) is corruption, and
2899 // must fail loudly rather than read as "context incomplete".
2900 Err(e) if skip_unloadable && !stage_is_absent(&e) => {
2901 return Err(StoreError::StageUnreadable {
2902 sig_id: sig_id.clone(),
2903 stage_id: stage_id.clone(),
2904 reason: e.to_string(),
2905 })
2906 }
2907 Err(e) if skip_unloadable => skipped.push(SkippedStage {
2908 sig_id: sig_id.clone(),
2909 stage_id: stage_id.clone(),
2910 name: self.recover_decl_name(sig_id, stage_id),
2911 reason: e.to_string(),
2912 called_by_target: false,
2913 }),
2914 Err(e) => return Err(e),
2915 }
2916 }
2917 Ok((decls, skipped))
2918 }
2919
2920 /// Best-effort (de-mangled) name of a declaration whose AST can't be
2921 /// loaded: its own metadata, else the metadata of any other stage of the
2922 /// same sig (a SigId fixes the name, so every stage of it shares one).
2923 fn recover_decl_name(&self, sig_id: &str, stage_id: &str) -> Option<String> {
2924 let bare = |m: Metadata| demangled_name(&m.name).to_string();
2925 if let Ok(m) = self.get_metadata_for_sig(sig_id, stage_id) {
2926 return Some(bare(m));
2927 }
2928 self.sig_history(sig_id)
2929 .ok()?
2930 .into_iter()
2931 .find_map(|h| self.get_metadata_for_sig(sig_id, &h.stage_id).ok())
2932 .map(bare)
2933 }
2934
2935 /// #836 G3: compare a regenerated `candidate` against what the op
2936 /// recorded producing, and emit the `Replay` attestation. The
2937 /// reproducibility claim made concrete — a faithful regeneration of
2938 /// the same function from the same cause yields the same
2939 /// content-addressed stage id.
2940 ///
2941 /// `reproduced` is true iff the candidate is the same sig *and* the
2942 /// same stage id the op recorded. A candidate for a different sig
2943 /// counts as "not reproduced" (`produced_stage_id: None`) rather
2944 /// than an error — it's a legitimate, if negative, replay result.
2945 /// The attestation is addressed to the op's recorded stage, so
2946 /// `list_for_stage` surfaces it alongside the TypeCheck/Examples
2947 /// evidence.
2948 pub fn replay_compare(
2949 &self,
2950 op_id: &str,
2951 candidate: &Stage,
2952 ) -> Result<ReplayOutcome, StoreError> {
2953 let (target_sig, expected_stage_id) = self.replay_target(op_id)?;
2954 let cand_sig = lex_ast::sig_id(candidate);
2955 let cand_stage = stage_id(candidate);
2956 let produced_stage_id = match (cand_sig.as_deref(), &cand_stage) {
2957 // Same function regenerated: the produced stage is
2958 // whatever it content-addresses to.
2959 (Some(s), Some(st)) if s == target_sig => Some(st.clone()),
2960 // A different sig (or an unhashable stage) isn't a
2961 // regeneration of this op's change.
2962 _ => None,
2963 };
2964 let reproduced = produced_stage_id.as_deref() == Some(expected_stage_id.as_str());
2965 let detail = if reproduced {
2966 None
2967 } else {
2968 Some("regeneration did not reproduce the recorded stage".to_string())
2969 };
2970 self.emit_replay(op_id, &expected_stage_id, produced_stage_id, reproduced, None, detail)
2971 }
2972
2973 /// Record a *negative* replay result for a regeneration that never
2974 /// yielded a comparable stage — the output didn't parse, or didn't
2975 /// define the target sig (#836 G3). Emits a `Replay { reproduced:
2976 /// false, produced_stage_id: None }` attestation with `reason` in
2977 /// its `Failed` detail, so an automated `lex op replay` run always
2978 /// records a verdict rather than aborting. `reason` is caller-supplied
2979 /// (e.g. "regenerated source did not parse").
2980 pub fn replay_record_miss(&self, op_id: &str, reason: &str) -> Result<ReplayOutcome, StoreError> {
2981 let (_target_sig, expected_stage_id) = self.replay_target(op_id)?;
2982 self.emit_replay(op_id, &expected_stage_id, None, false, None, Some(reason.to_string()))
2983 }
2984
2985 /// `(target_sig, expected_stage_id)` for a replayable op, or an
2986 /// error if the op is unknown or produced no stage.
2987 fn replay_target(&self, op_id: &str) -> Result<(String, String), StoreError> {
2988 let log = lex_vcs::OpLog::open(self.root())?;
2989 let record = log
2990 .get(&op_id.to_string())?
2991 .ok_or_else(|| StoreError::UnknownOp(op_id.to_string()))?;
2992 crate::files::refuse_non_semantic(&record)?;
2993 produced_sig_stage(&record.produces)
2994 .ok_or_else(|| StoreError::InvalidTransition(format!("op {op_id} produced no stage to replay")))
2995 }
2996
2997 /// Record a replay verdict the caller has already decided — used by
2998 /// the CLI's behavioral tier, which does the (VM-backed) equivalence
2999 /// check the store deliberately can't. `expected_stage_id` is looked
3000 /// up from the op. Set `behavioral_samples` to `Some(n)` when the
3001 /// candidate reproduced *behaviorally* over `n` sampled inputs rather
3002 /// than by exact stage-id match; the attestation then records that
3003 /// weaker-but-real claim distinctly.
3004 pub fn replay_record(
3005 &self,
3006 op_id: &str,
3007 produced_stage_id: Option<String>,
3008 reproduced: bool,
3009 behavioral_samples: Option<usize>,
3010 fail_detail: Option<String>,
3011 ) -> Result<ReplayOutcome, StoreError> {
3012 let (_target_sig, expected_stage_id) = self.replay_target(op_id)?;
3013 self.emit_replay(op_id, &expected_stage_id, produced_stage_id, reproduced, behavioral_samples, fail_detail)
3014 }
3015
3016 /// Compute the exact-match verdict for a candidate *without* emitting
3017 /// an attestation — `(expected_stage_id, produced_stage_id, exact)`.
3018 /// Lets a caller (the CLI) fall back to a behavioral check on a valid
3019 /// but non-identical candidate and emit a single verdict, instead of
3020 /// [`Self::replay_compare`]'s emit-immediately shape.
3021 pub fn replay_stage_of(
3022 &self,
3023 op_id: &str,
3024 candidate: &Stage,
3025 ) -> Result<(String, Option<String>, bool), StoreError> {
3026 let (target_sig, expected_stage_id) = self.replay_target(op_id)?;
3027 let cand_sig = lex_ast::sig_id(candidate);
3028 let cand_stage = stage_id(candidate);
3029 let produced_stage_id = match (cand_sig.as_deref(), &cand_stage) {
3030 (Some(s), Some(st)) if s == target_sig => Some(st.clone()),
3031 _ => None,
3032 };
3033 let exact = produced_stage_id.as_deref() == Some(expected_stage_id.as_str());
3034 Ok((expected_stage_id, produced_stage_id, exact))
3035 }
3036
3037 /// Emit the `Replay` attestation and build the outcome. Shared by
3038 /// [`Self::replay_compare`], [`Self::replay_record_miss`], and
3039 /// [`Self::replay_record`].
3040 fn emit_replay(
3041 &self,
3042 op_id: &str,
3043 expected_stage_id: &str,
3044 produced_stage_id: Option<String>,
3045 reproduced: bool,
3046 behavioral_samples: Option<usize>,
3047 fail_detail: Option<String>,
3048 ) -> Result<ReplayOutcome, StoreError> {
3049 let model = {
3050 let log = lex_vcs::OpLog::open(self.root())?;
3051 match log.get(&op_id.to_string())?.and_then(|r| r.op.intent_id) {
3052 Some(id) => lex_vcs::IntentLog::open(self.root())?
3053 .get(&id)?
3054 .map(|i| model_label(&i.model)),
3055 None => None,
3056 }
3057 };
3058 let result = if reproduced {
3059 lex_vcs::AttestationResult::Passed
3060 } else {
3061 lex_vcs::AttestationResult::Failed {
3062 detail: fail_detail.unwrap_or_else(|| "not reproduced".into()),
3063 }
3064 };
3065 let attestation = lex_vcs::Attestation::new(
3066 expected_stage_id.to_string(),
3067 Some(op_id.to_string()),
3068 None,
3069 lex_vcs::AttestationKind::Replay {
3070 expected_stage_id: expected_stage_id.to_string(),
3071 produced_stage_id: produced_stage_id.clone(),
3072 reproduced,
3073 behavioral_samples,
3074 model,
3075 },
3076 result,
3077 replay_producer(),
3078 None,
3079 );
3080 let attestation_id = attestation.attestation_id.clone();
3081 self.attestation_log()?.put(&attestation)?;
3082 Ok(ReplayOutcome {
3083 op_id: op_id.to_string(),
3084 expected_stage_id: expected_stage_id.to_string(),
3085 produced_stage_id,
3086 reproduced,
3087 behavioral_samples,
3088 attestation_id,
3089 context_incomplete: false,
3090 skipped: Vec::new(),
3091 })
3092 }
3093
3094 /// The program at an op (that op and all its ancestors applied), as
3095 /// canonical stages. The behavioral replay tier needs the whole
3096 /// program — a regenerated function may call helpers from its parent
3097 /// state, so it can only be run in context. Exposed for the CLI's
3098 /// equivalence check; `op_id` may be any op in the log.
3099 pub fn program_stages_at_op(&self, op_id: &str) -> Result<Vec<Stage>, StoreError> {
3100 Ok(self.program_stages_at_op_impl(op_id, false)?.stages)
3101 }
3102
3103 /// [`Self::program_stages_at_op`], but head declarations whose stage
3104 /// can't be loaded are skipped and reported instead of failing the
3105 /// reconstruction (#868).
3106 pub fn program_stages_at_op_skipping(&self, op_id: &str) -> Result<ReconstructedProgram, StoreError> {
3107 self.program_stages_at_op_impl(op_id, true)
3108 }
3109
3110 fn program_stages_at_op_impl(
3111 &self,
3112 op_id: &str,
3113 skip_unloadable: bool,
3114 ) -> Result<ReconstructedProgram, StoreError> {
3115 let oid: lex_vcs::OpId = op_id.to_string();
3116 let log = lex_vcs::OpLog::open(self.root())?;
3117 let mut map: std::collections::BTreeMap<String, String> = std::collections::BTreeMap::new();
3118 for rec in log.walk_forward(&oid, None)? {
3119 crate::branches::apply_transition(&mut map, &rec.produces);
3120 }
3121 let pairs: Vec<(String, String)> = map.into_iter().collect();
3122 let (decls, skipped) = self.load_head_decls(&pairs, skip_unloadable)?;
3123 // #946: include the head's `import` edges. The SigId→stage map holds
3124 // only fn/type declarations, so without this a non-inlined head
3125 // reconstructs to a program whose `<alias>.name` calls have no import
3126 // to bind against — and `replay_request` then hands a regenerator
3127 // parent source that silently omits the very dependencies the code
3128 // calls into, which is worse than useless as context.
3129 Ok(ReconstructedProgram { stages: self.with_head_imports(Some(op_id), decls), skipped })
3130 }
3131
3132 /// See [`demangled_name`]; method form for call sites that already hold a
3133 /// `Store`.
3134 pub fn bare_declaration_name(name: &str) -> &str {
3135 demangled_name(name)
3136 }
3137
3138 /// The program at an op **as an author would write it** — declarations
3139 /// under bare names (`twice`, not `lib_<hash>.twice`), with the head's
3140 /// imports (#980).
3141 ///
3142 /// A package publish mangles every declaration with a path-derived prefix.
3143 /// That prefix is a *storage* detail: a dotted name cannot be written as
3144 /// Lex source at all (`fn lib_abc.twice(..)` is a parse error), so anything
3145 /// that shows a program to an author — or asks one to regenerate a
3146 /// function from it — has to de-mangle first.
3147 ///
3148 /// Falls back to the mangled reconstruction when de-mangling isn't
3149 /// available (a multi-module head, pending #942), so callers degrade to the
3150 /// previous behaviour rather than failing outright.
3151 pub fn demangled_program_at_op(&self, op_id: &str) -> Result<Vec<Stage>, StoreError> {
3152 Ok(self.demangled_program_at_op_impl(op_id, false)?.stages)
3153 }
3154
3155 /// [`Self::demangled_program_at_op`], but head declarations whose stage
3156 /// can't be loaded (GC'd, superseded, never persisted) are skipped and
3157 /// reported instead of failing the reconstruction (#868). This is what
3158 /// replay uses: a partial context is still worth replaying against, and
3159 /// the skips travel with the request and verdict.
3160 pub fn demangled_program_at_op_skipping(&self, op_id: &str) -> Result<ReconstructedProgram, StoreError> {
3161 self.demangled_program_at_op_impl(op_id, true)
3162 }
3163
3164 fn demangled_program_at_op_impl(
3165 &self,
3166 op_id: &str,
3167 skip_unloadable: bool,
3168 ) -> Result<ReconstructedProgram, StoreError> {
3169 match crate::render::demangled_head_stages_impl(self, op_id, skip_unloadable) {
3170 Ok((stages, skipped)) => Ok(ReconstructedProgram { stages, skipped }),
3171 Err(StoreError::UnsupportedMultiModuleDependency) => {
3172 self.program_stages_at_op_impl(op_id, skip_unloadable)
3173 }
3174 Err(e) => Err(e),
3175 }
3176 }
3177
3178 /// Open the attestation log rooted at this store. The log lives
3179 /// under `<root>/attestations/`; opening is idempotent and cheap
3180 /// (`fs::create_dir_all`). Exposed publicly so consumers — `lex
3181 /// blame --with-evidence`, `GET /v1/stage/<id>/attestations` —
3182 /// can read what the store gate emitted without round-tripping
3183 /// through this crate's API surface.
3184 /// Recompute a producer's trust score from its recent
3185 /// attestation history and emit a fresh `ProducerTrust`
3186 /// attestation (#293). Score = `passed / (passed + failed
3187 /// + inconclusive)` over the last `window` attestations
3188 /// produced by `tool_id`, expressed in thousandths
3189 /// (`0..=1000`).
3190 ///
3191 /// Refuses to grant trust when the tool has an active
3192 /// `ProducerBlock` — the block wins as a hard veto. Returns
3193 /// `Ok(None)` for "no attestations to score" (a brand-new
3194 /// producer); the caller can choose how to handle it
3195 /// (typically: skip the publish until evidence accrues).
3196 ///
3197 /// `granted_by` is the identity of the actor running the
3198 /// recompute (typically the human admin, or "lex-ci-bot"
3199 /// for an automated nightly).
3200 pub fn recompute_producer_trust(
3201 &self,
3202 tool_id: &str,
3203 window: usize,
3204 granted_by: &str,
3205 ) -> Result<Option<lex_vcs::AttestationId>, StoreError> {
3206 let log = self.attestation_log()?;
3207 let all = log.list_all()?;
3208 // Hard veto: don't grant trust to a blocked tool.
3209 if lex_vcs::active_producer_block(&all, tool_id).is_some() {
3210 return Err(StoreError::InvalidTransition(format!(
3211 "cannot recompute trust for `{tool_id}` — \
3212 producer is currently blocked"
3213 )));
3214 }
3215 // Filter to attestations from this tool, newest-first by
3216 // timestamp, then take the window.
3217 let mut from_tool: Vec<&lex_vcs::Attestation> = all
3218 .iter()
3219 .filter(|a| a.produced_by.tool == tool_id)
3220 // Ignore self-referential trust attestations (we're
3221 // scoring evidence, not previous trust statements).
3222 .filter(|a| {
3223 !matches!(
3224 a.kind,
3225 lex_vcs::AttestationKind::ProducerTrust { .. }
3226 | lex_vcs::AttestationKind::TrustWaived { .. }
3227 )
3228 })
3229 .collect();
3230 from_tool.sort_by_key(|a| std::cmp::Reverse(a.timestamp));
3231 from_tool.truncate(window);
3232 if from_tool.is_empty() {
3233 return Ok(None);
3234 }
3235 let (mut passed, mut total) = (0u64, 0u64);
3236 for a in &from_tool {
3237 total += 1;
3238 if matches!(a.result, lex_vcs::AttestationResult::Passed) {
3239 passed += 1;
3240 }
3241 }
3242 let score = if total == 0 {
3243 0
3244 } else {
3245 let raw = (passed as f64) * 1000.0 / (total as f64);
3246 raw.round().clamp(0.0, 1000.0) as u32
3247 };
3248 let head_op = self
3249 .list_branches()?
3250 .into_iter()
3251 .find_map(|b| self.get_branch(&b).ok().flatten().and_then(|x| x.head_op))
3252 .unwrap_or_else(|| "fresh".into());
3253 let evidence = format!(
3254 "window={window}, sample={}, head_op={head_op:.16}",
3255 from_tool.len()
3256 );
3257 let attestation = lex_vcs::Attestation::new(
3258 tool_id.to_string(),
3259 None,
3260 None,
3261 lex_vcs::AttestationKind::ProducerTrust {
3262 tool_id: tool_id.into(),
3263 score_thousandths: score,
3264 evidence,
3265 granted_by: granted_by.into(),
3266 },
3267 lex_vcs::AttestationResult::Passed,
3268 producer_trust_producer(),
3269 None,
3270 );
3271 let id = attestation.attestation_id.clone();
3272 log.put(&attestation)?;
3273 Ok(Some(id))
3274 }
3275
3276 /// The latest live `ProducerTrust` score (thousandths, `0..=1000`) for
3277 /// every producer that currently has trust: the newest score per tool by
3278 /// timestamp, excluding any tool under an active `ProducerBlock` (a block
3279 /// is a hard veto over trust, matching `recompute_producer_trust`).
3280 ///
3281 /// Used to export a capsule trusted-keys keyring from *earned* trust — the
3282 /// producer id doubles as the publisher's signing key downstream, so this
3283 /// turns track record into the allowlist `capsule install` consumes.
3284 pub fn live_producer_trust_scores(
3285 &self,
3286 ) -> Result<std::collections::BTreeMap<String, u32>, StoreError> {
3287 let log = self.attestation_log()?;
3288 let all = log.list_all()?;
3289 // Newest score per tool.
3290 let mut latest: std::collections::BTreeMap<String, (u64, u32)> =
3291 std::collections::BTreeMap::new();
3292 for a in &all {
3293 if let lex_vcs::AttestationKind::ProducerTrust {
3294 tool_id,
3295 score_thousandths,
3296 ..
3297 } = &a.kind
3298 {
3299 let entry = latest.entry(tool_id.clone()).or_insert((0, 0));
3300 if a.timestamp >= entry.0 {
3301 *entry = (a.timestamp, *score_thousandths);
3302 }
3303 }
3304 }
3305 // Drop blocked producers; a block vetoes trust.
3306 let mut scores = std::collections::BTreeMap::new();
3307 for (tool, (_, score)) in latest {
3308 if lex_vcs::active_producer_block(&all, &tool).is_some() {
3309 continue;
3310 }
3311 scores.insert(tool, score);
3312 }
3313 Ok(scores)
3314 }
3315
3316 pub fn attestation_log(&self) -> Result<lex_vcs::AttestationLog, StoreError> {
3317 Ok(lex_vcs::AttestationLog::open(self.root())?)
3318 }
3319
3320 /// Emit one `TypeCheck::Passed` attestation per stage produced by
3321 /// a successful gated apply. Idempotent on `attestation_id` —
3322 /// re-running the same gate run dedups via content addressing.
3323 ///
3324 /// Failure modes: `io::Error` from the attestation log (disk
3325 /// full, perms). The op has already landed by the time this
3326 /// runs; an error here means the op is durable but the evidence
3327 /// is missing. We propagate so the caller sees the partial
3328 /// state rather than silently swallowing — re-attesting the
3329 /// same op against the same op_id is idempotent (content
3330 /// addressing) so a retry is safe once the underlying issue is
3331 /// fixed.
3332 fn record_typecheck_passed(
3333 &self,
3334 stage_ids: &[String],
3335 op_id: &lex_vcs::OpId,
3336 ) -> Result<(), StoreError> {
3337 if stage_ids.is_empty() {
3338 return Ok(());
3339 }
3340 let log = self.attestation_log()?;
3341 for stage_id in stage_ids {
3342 let attestation = lex_vcs::Attestation::new(
3343 stage_id.clone(),
3344 Some(op_id.clone()),
3345 None,
3346 lex_vcs::AttestationKind::TypeCheck,
3347 lex_vcs::AttestationResult::Passed,
3348 typecheck_producer(),
3349 None,
3350 );
3351 log.put(&attestation)?;
3352 }
3353 Ok(())
3354 }
3355
3356 /// The hosted CI runner (#93): independently re-run the write-time
3357 /// type-check gate on a branch head and record the verdict as a
3358 /// `lex-hub-ci`-produced `TypeCheck` attestation for the stages the
3359 /// advance introduced. Called after an `op push` fast-forwards the
3360 /// head, so `require-attestation type_check` gates are backed by a
3361 /// producer that actually verified the code server-side, not by
3362 /// whatever attestation a client chose to attach. Does NOT move or
3363 /// roll back the head — the client's own always-valid-HEAD gate is
3364 /// what refuses a bad publish; this produces the trusted verdict on
3365 /// top of an already-committed advance (so a client that bypassed
3366 /// its gate is caught by a `TypeCheck::Failed` from `lex-hub-ci`).
3367 ///
3368 /// `from_head` is the branch head *before* the advance; the ops
3369 /// between it and `to_head` are the ones whose stages get attested.
3370 /// Idempotent: attestations are content-addressed, so re-verifying
3371 /// the same head is a no-op.
3372 pub fn verify_head_and_attest(
3373 &self,
3374 branch: &str,
3375 from_head: Option<&str>,
3376 to_head: &str,
3377 ) -> Result<HubCiVerdict, StoreError> {
3378 // Reconstruct the program at the new head and re-check it.
3379 let head = self.branch_head(branch)?;
3380 let pairs: Vec<(String, String)> = head.into_iter().collect();
3381 let decls: Vec<Stage> =
3382 self.get_asts_for_sigs_bulk(&pairs).into_iter().collect::<Result<_, _>>()?;
3383 let checked_stages = decls.len();
3384 // #930: the SigId→stage map holds only fn/type declarations — the
3385 // head's `import` edges are AddImport ops, absent here. Reconstruct
3386 // them so a non-inlined head's `<alias>.name` references bind: the
3387 // resolver scans these imports to resolve each dependency, and the
3388 // checker's Pass 1 binds the alias to the resolved module. (Without
3389 // this the alias is unbound and the head fails as `unknown_identifier`,
3390 // even with the dependency correctly resolved.)
3391 let head_imports = crate::render::package_head_at_op(self, to_head)
3392 .map(|ph| {
3393 ph.flat_imports
3394 .into_iter()
3395 .map(|(reference, alias)| {
3396 Stage::Import(lex_ast::Import { reference, alias })
3397 })
3398 .collect::<Vec<_>>()
3399 })
3400 .unwrap_or_default();
3401 let mut stages = head_imports;
3402 stages.extend(decls);
3403 // #930: the hub gate resolves this head's external dependencies from
3404 // the lock committed with `to_head` (via the installed cross-store
3405 // resolver); empty when none is installed or the head is inlined.
3406 let result = match self.check_with_resolved_deps(&stages, Some(to_head)) {
3407 Ok(_) => lex_vcs::AttestationResult::Passed,
3408 Err(errors) => lex_vcs::AttestationResult::Failed {
3409 detail: serde_json::to_string(&errors).unwrap_or_else(|_| "type errors".into()),
3410 },
3411 };
3412 let passed = matches!(result, lex_vcs::AttestationResult::Passed);
3413
3414 // Stages introduced by THIS advance (from_head exclusive → to_head).
3415 let log = lex_vcs::OpLog::open(self.root())?;
3416 let to = to_head.to_string();
3417 let records = match from_head {
3418 Some(f) => log
3419 .walk_forward_since(&to, &f.to_string())?
3420 .unwrap_or_else(|| log.walk_forward(&to, None).unwrap_or_default()),
3421 None => log.walk_forward(&to, None)?,
3422 };
3423 let mut introduced: Vec<String> = Vec::new();
3424 for rec in &records {
3425 introduced.extend(attestable_stage_ids(&rec.produces));
3426 }
3427
3428 let alog = self.attestation_log()?;
3429 for sid in &introduced {
3430 let att = lex_vcs::Attestation::new(
3431 sid.clone(),
3432 Some(to_head.to_string()),
3433 None,
3434 lex_vcs::AttestationKind::TypeCheck,
3435 result.clone(),
3436 hub_ci_producer(),
3437 None,
3438 );
3439 alog.put(&att)?;
3440 }
3441
3442 let detail = match &result {
3443 lex_vcs::AttestationResult::Failed { detail } => Some(detail.clone()),
3444 _ => None,
3445 };
3446 Ok(HubCiVerdict { passed, checked_stages, attested_stages: introduced.len(), detail })
3447 }
3448
3449 /// Emit an `Examples::Passed` attestation for a published stage
3450 /// whose behavioral `examples {}` block was run and passed (#835,
3451 /// Tier 1). Mirrors [`Self::record_typecheck_passed`]. The
3452 /// behavioral run itself happens one layer up (lex-api / lex-cli)
3453 /// because it needs the bytecode compiler + VM, which this crate
3454 /// deliberately doesn't depend on; the store only records the
3455 /// verdict. `file_hash` uses the stage id — the stage fully
3456 /// determines its own examples.
3457 pub fn record_examples_passed(
3458 &self,
3459 stage_id: &str,
3460 op_id: &lex_vcs::OpId,
3461 count: usize,
3462 ) -> Result<(), StoreError> {
3463 let log = self.attestation_log()?;
3464 let attestation = lex_vcs::Attestation::new(
3465 stage_id.to_string(),
3466 Some(op_id.clone()),
3467 None,
3468 lex_vcs::AttestationKind::Examples { file_hash: stage_id.to_string(), count },
3469 lex_vcs::AttestationResult::Passed,
3470 examples_producer(),
3471 None,
3472 );
3473 log.put(&attestation)?;
3474 Ok(())
3475 }
3476
3477 /// Record a structured `Review` verdict on a stage (#836 G4).
3478 /// The verdict maps onto the attestation `result` so existing
3479 /// result-based tooling reads it: Approve->Passed,
3480 /// Reject->Failed, RequestChanges->Inconclusive.
3481 pub fn record_review(
3482 &self,
3483 stage_id: &str,
3484 op_id: Option<lex_vcs::OpId>,
3485 reviewer: &str,
3486 verdict: lex_vcs::ReviewVerdict,
3487 notes: Option<String>,
3488 ) -> Result<lex_vcs::AttestationId, StoreError> {
3489 let result = match verdict {
3490 lex_vcs::ReviewVerdict::Approve => lex_vcs::AttestationResult::Passed,
3491 lex_vcs::ReviewVerdict::Reject => lex_vcs::AttestationResult::Failed {
3492 detail: notes.clone().unwrap_or_else(|| "rejected".into()),
3493 },
3494 lex_vcs::ReviewVerdict::RequestChanges => lex_vcs::AttestationResult::Inconclusive {
3495 detail: notes.clone().unwrap_or_else(|| "changes requested".into()),
3496 },
3497 };
3498 let att = lex_vcs::Attestation::new(
3499 stage_id.to_string(),
3500 op_id,
3501 None,
3502 lex_vcs::AttestationKind::Review { reviewer: reviewer.to_string(), verdict, notes },
3503 result,
3504 review_producer(reviewer),
3505 None,
3506 );
3507 let id = att.attestation_id.clone();
3508 self.attestation_log()?.put(&att)?;
3509 Ok(id)
3510 }
3511
3512 /// The latest `Review` verdict recorded on a stage, if any
3513 /// (#836 G4). "Latest" is **arrival order** in the attestation
3514 /// log — a server-assigned sequence number — NOT the
3515 /// attestation's `timestamp`, which the writer chooses and which
3516 /// is excluded from the attestation id. So a verdict pushed with a
3517 /// far-future timestamp does not outrank one that actually arrived
3518 /// later. Legacy attestations with no arrival stamp order by
3519 /// timestamp among themselves and below every stamped one; see
3520 /// [`lex_vcs::AttestationLog::sort_by_arrival`]. Used by
3521 /// `promote_candidate` to honor a standing Reject.
3522 ///
3523 /// This fixes *ordering* only. It does not authenticate who wrote a
3524 /// verdict: any writer that can reach the log can still append a
3525 /// later one.
3526 pub fn latest_review_verdict(
3527 &self,
3528 stage_id: &str,
3529 ) -> Result<Option<lex_vcs::ReviewVerdict>, StoreError> {
3530 let log = self.attestation_log()?;
3531 // Oldest first; the last Review in that order is the latest.
3532 let latest = log
3533 .list_for_stage_by_arrival(&stage_id.to_string())?
3534 .into_iter()
3535 .rev()
3536 .find_map(|a| match a.kind {
3537 lex_vcs::AttestationKind::Review { verdict, .. } => Some(verdict),
3538 _ => None,
3539 });
3540 Ok(latest)
3541 }
3542
3543 /// Consult `policy.session_budgets` for the op's session
3544 /// (resolved via `op.intent_id → Intent.session_id`) and
3545 /// refuse if applying would push the session's monotonic spend
3546 /// over the configured cap (#292 slice 3).
3547 ///
3548 /// Ops without an `intent_id`, or whose intent has no
3549 /// configured cap, return Ok without any disk read.
3550 fn check_session_budget(&self, op: &lex_vcs::Operation) -> Result<(), StoreError> {
3551 let Some(intent_id) = op.intent_id.as_deref() else {
3552 return Ok(());
3553 };
3554 let intent_log = lex_vcs::IntentLog::open(self.root())?;
3555 let Some(intent) = intent_log.get(&intent_id.to_string())? else {
3556 // Dangling intent — treat as "no session" and let it
3557 // sail through. Slice 1's ledger already documents
3558 // this as graceful-degradation semantics.
3559 return Ok(());
3560 };
3561 let policy = crate::policy::load(self.root())?.unwrap_or_default();
3562 let Some(cap) = policy.session_budgets.cap_for(&intent.session_id) else {
3563 return Ok(());
3564 };
3565 // Recompute the session's current spend + the contribution
3566 // from this op. Re-running the ledger walk on every gated
3567 // op is O(branch history); see #292 slice 1's note about
3568 // a future on-disk cache.
3569 let current = self.session_budget(&intent.session_id)?;
3570 let increment = crate::budget::monotonic_spend_of(&op.kind);
3571 let spent_after = current.spent.saturating_add(increment);
3572 if spent_after > cap {
3573 return Err(StoreError::BudgetExceeded {
3574 session_id: intent.session_id,
3575 cap,
3576 spent_after,
3577 });
3578 }
3579 Ok(())
3580 }
3581
3582 /// Emit `RepairHint` attestations for a TypeError-rejected op
3583 /// (#281). One per candidate stage in the transition. The hint
3584 /// records the *would-be* op_id (deterministic, content-
3585 /// addressed even though the op record was never persisted)
3586 /// and the structured errors.
3587 ///
3588 /// #306 slice 3: `suggested_transform` is populated from the
3589 /// static (rule_tag → likely_transform) table for the *first*
3590 /// error in the batch. The LLM-driven `lex repair --apply`
3591 /// flow can still overwrite this with a higher-quality
3592 /// suggestion; the static value is the floor, not the ceiling.
3593 ///
3594 /// Best-effort: a write failure here is swallowed by the
3595 /// caller (the original `TypeError` is the load-bearing
3596 /// signal; missing the hint is recoverable on a retry).
3597 fn record_repair_hint(
3598 &self,
3599 stage_ids: &[String],
3600 failed_op_id: &lex_vcs::OpId,
3601 errors: &[lex_types::TypeError],
3602 ) -> Result<(), StoreError> {
3603 if stage_ids.is_empty() {
3604 return Ok(());
3605 }
3606 let errors_json = serde_json::to_value(errors).map_err(StoreError::Serde)?;
3607 // #306 slice 3: look up the static suggested_transform for
3608 // the first error's rule_tag. Multiple errors per op are
3609 // possible — when they fire in lockstep (e.g. one bad let
3610 // binding propagates to several use sites), the first
3611 // error's rule_tag is usually the load-bearing one to fix.
3612 let suggested_transform = errors
3613 .first()
3614 .and_then(|e| lex_types::suggested_transform_for(e.rule_tag()));
3615 let log = self.attestation_log()?;
3616 for stage_id in stage_ids {
3617 let attestation = lex_vcs::Attestation::new(
3618 stage_id.clone(),
3619 None, // the failed op was never persisted; not the
3620 // attestation's op_id (which is for a
3621 // *successful* op).
3622 None,
3623 lex_vcs::AttestationKind::RepairHint {
3624 failed_op_id: failed_op_id.clone(),
3625 errors: errors_json.clone(),
3626 suggested_transform: suggested_transform.clone(),
3627 },
3628 lex_vcs::AttestationResult::Failed {
3629 detail: format!(
3630 "op {} rejected: {} type error(s)",
3631 failed_op_id,
3632 errors.len()
3633 ),
3634 },
3635 repair_hint_producer(),
3636 None,
3637 );
3638 log.put(&attestation)?;
3639 }
3640 Ok(())
3641 }
3642
3643 /// Emit `Trace` attestations linking an already-committed `op`
3644 /// to the run that produced it (#257). One attestation per
3645 /// produced stage (matching the `TypeCheck` emission contract
3646 /// — see [`Self::apply_operation_checked`]) with
3647 /// `op_id: Some(op_id)` set, so `lex trace --op <op_id>`
3648 /// surfaces the run.
3649 ///
3650 /// Returns the number of attestations emitted (zero for ops
3651 /// that produce no attestable stage, e.g. `Remove` /
3652 /// `ImportOnly`).
3653 ///
3654 /// Idempotent: re-emitting for the same
3655 /// `(run_id, root_target, op_id, stage_id, producer, result)`
3656 /// tuple dedups via content addressing.
3657 ///
3658 /// `op_id` must already exist in the op log — an unknown op
3659 /// surfaces as `StoreError::UnknownOp`.
3660 pub fn record_op_trace(
3661 &self,
3662 run_id: &str,
3663 root_target: &str,
3664 op_id: &lex_vcs::OpId,
3665 result: lex_vcs::AttestationResult,
3666 producer: lex_vcs::ProducerDescriptor,
3667 ) -> Result<usize, StoreError> {
3668 let log = lex_vcs::OpLog::open(self.root())?;
3669 let rec = log
3670 .get(op_id)?
3671 .ok_or_else(|| StoreError::UnknownOp(op_id.clone()))?;
3672 let stage_ids = attestable_stage_ids(&rec.produces);
3673 if stage_ids.is_empty() {
3674 return Ok(0);
3675 }
3676 let attlog = self.attestation_log()?;
3677 let mut emitted = 0;
3678 for stage_id in stage_ids {
3679 let attestation = lex_vcs::Attestation::new(
3680 stage_id,
3681 Some(op_id.clone()),
3682 None,
3683 lex_vcs::AttestationKind::Trace {
3684 run_id: run_id.into(),
3685 root_target: root_target.into(),
3686 },
3687 result.clone(),
3688 producer.clone(),
3689 None,
3690 );
3691 attlog.put(&attestation)?;
3692 emitted += 1;
3693 }
3694 Ok(emitted)
3695 }
3696
3697 /// Walk `ops_since(branch_head, base)` and emit per-stage
3698 /// `Trace` attestations for each new op, linking them to the
3699 /// run that produced them (#257). Used by `lex run --trace`
3700 /// after the VM exits: snapshot `base = branch_head` before
3701 /// the run, then call this with the post-run head.
3702 ///
3703 /// `base = None` means "every op currently reachable from the
3704 /// branch head" — generally not what you want for a single
3705 /// run; pass the pre-run head.
3706 ///
3707 /// Returns the total number of attestations emitted across
3708 /// every new op. Zero is the common case (the run committed no
3709 /// ops).
3710 ///
3711 /// Idempotent on the per-op level via [`Self::record_op_trace`].
3712 pub fn record_run_committed_ops_since(
3713 &self,
3714 run_id: &str,
3715 root_target: &str,
3716 branch: &str,
3717 base: Option<&lex_vcs::OpId>,
3718 result: lex_vcs::AttestationResult,
3719 producer: lex_vcs::ProducerDescriptor,
3720 ) -> Result<usize, StoreError> {
3721 let head = match self.get_branch(branch)?.and_then(|b| b.head_op) {
3722 Some(h) => h,
3723 None => return Ok(0),
3724 };
3725 let log = lex_vcs::OpLog::open(self.root())?;
3726 let new_ops = log.ops_since(&head, base)?;
3727 let mut total = 0;
3728 for rec in new_ops {
3729 total += self.record_op_trace(
3730 run_id,
3731 root_target,
3732 &rec.op_id,
3733 result.clone(),
3734 producer.clone(),
3735 )?;
3736 }
3737 Ok(total)
3738 }
3739
3740 /// Apply a typed `ReplaceMatchArm` transform (#280) and emit a
3741 /// `OperationKind::ReplaceMatchArm` op that records the
3742 /// semantic shape of the edit, not just the byte effect.
3743 ///
3744 /// Steps:
3745 /// 1. Load the source stage's canonical bytes (delta-aware).
3746 /// 2. Run [`lex_ast::replace_match_arm`] to produce the new
3747 /// `Stage`. Pure function, no I/O.
3748 /// 3. Publish the new stage. Idempotent on the
3749 /// content-addressed `to_stage_id`.
3750 /// 4. Assemble the candidate program (every active stage on
3751 /// the branch, with the rewritten one swapped in) and call
3752 /// [`Self::apply_operation_checked`] — re-typechecks and
3753 /// runs every existing gate (TypeCheck attestation,
3754 /// required_attestations, producer-block walk-back).
3755 ///
3756 /// Failure modes:
3757 /// * [`StoreError::TransformError`] — transform didn't apply.
3758 /// The branch is unchanged; no stage published.
3759 /// * [`StoreError::TypeError`] — transform produced an
3760 /// ill-typed program. The new stage is on disk (idempotent
3761 /// on its content hash) but the branch is unchanged. Same
3762 /// "publish without advance" semantics as #245.
3763 /// * Everything else from `apply_operation_checked`.
3764 pub fn apply_replace_match_arm(
3765 &self,
3766 branch: &str,
3767 from_stage_id: &str,
3768 match_node: &lex_ast::NodeId,
3769 arm_index: usize,
3770 new_body: lex_ast::CExpr,
3771 ) -> Result<lex_vcs::OpId, StoreError> {
3772 self.apply_replace_match_arm_with_intent(
3773 branch,
3774 from_stage_id,
3775 match_node,
3776 arm_index,
3777 new_body,
3778 None,
3779 )
3780 }
3781
3782 /// [`Self::apply_replace_match_arm`] with the write attributed to `intent`
3783 /// (#837 piece A): the op is stamped with the intent's id and the intent
3784 /// is recorded once the gate has passed. `None` is the intent-less form.
3785 pub fn apply_replace_match_arm_with_intent(
3786 &self,
3787 branch: &str,
3788 from_stage_id: &str,
3789 match_node: &lex_ast::NodeId,
3790 arm_index: usize,
3791 new_body: lex_ast::CExpr,
3792 intent: Option<&lex_vcs::Intent>,
3793 ) -> Result<lex_vcs::OpId, StoreError> {
3794 let from_stage = self.get_ast(from_stage_id)?;
3795 let new_stage = lex_ast::replace_match_arm(&from_stage, match_node, arm_index, new_body)
3796 .map_err(StoreError::TransformError)?;
3797 let sig = lex_ast::sig_id(&from_stage).ok_or(StoreError::CannotPublishImport)?;
3798 let to_stage_id = self.publish(&new_stage)?;
3799 if to_stage_id == from_stage_id {
3800 // No-op transform — the new body was structurally
3801 // identical to the old. Refuse rather than advancing
3802 // the branch with an empty edit.
3803 return Err(StoreError::InvalidTransition(format!(
3804 "replace_match_arm produced the same stage_id `{from_stage_id}`"
3805 )));
3806 }
3807
3808 // Assemble the candidate program: every active stage on
3809 // the branch, with `from_stage_id` swapped for `new_stage`.
3810 let head = self.branch_head(branch)?;
3811 let mut candidate: Vec<lex_ast::Stage> = Vec::with_capacity(head.len());
3812 for (other_sig, other_stage_id) in &head {
3813 if other_sig == &sig {
3814 candidate.push(new_stage.clone());
3815 } else {
3816 candidate.push(self.get_ast(other_stage_id)?);
3817 }
3818 }
3819 // If the source sig isn't on the current branch head, the
3820 // transform is operating on a stage that hasn't been added
3821 // yet — refuse rather than risking a candidate program
3822 // that doesn't reflect the branch's actual state.
3823 if !head.contains_key(&sig) {
3824 return Err(StoreError::InvalidTransition(format!(
3825 "sig `{sig}` not on branch `{branch}`'s head"
3826 )));
3827 }
3828
3829 // #247: budget delta captured for `lex op log --budget-drift`.
3830 let from_budget = budget_of_stage(&from_stage);
3831 let to_budget = budget_of_stage(&new_stage);
3832
3833 let head_now = self.get_branch(branch)?.and_then(|b| b.head_op);
3834 let kind = lex_vcs::OperationKind::ReplaceMatchArm {
3835 sig_id: sig.clone(),
3836 from_stage_id: from_stage_id.to_string(),
3837 to_stage_id: to_stage_id.clone(),
3838 match_node: match_node.as_str().to_string(),
3839 arm_index,
3840 from_budget,
3841 to_budget,
3842 };
3843 let transition = lex_vcs::StageTransition::Replace {
3844 sig_id: sig.clone(),
3845 from: from_stage_id.to_string(),
3846 to: to_stage_id.clone(),
3847 };
3848 let op = lex_vcs::Operation::new(kind, head_now.into_iter().collect::<Vec<_>>());
3849 let op = match intent {
3850 Some(i) => op.with_intent(i.intent_id.clone()),
3851 None => op,
3852 };
3853 self.apply_operation_checked_with_intent(branch, op, transition, &candidate, intent)
3854 }
3855
3856 /// Apply a typed `RenameLocal` transform (#280) — rename a
3857 /// `let`-bound local within a fn body and emit a matching
3858 /// `OperationKind::RenameLocal`. Same end-to-end shape as
3859 /// [`Self::apply_replace_match_arm`]; see that method for the
3860 /// failure-mode taxonomy.
3861 pub fn apply_rename_local(
3862 &self,
3863 branch: &str,
3864 from_stage_id: &str,
3865 let_node: &lex_ast::NodeId,
3866 new_name: &str,
3867 ) -> Result<lex_vcs::OpId, StoreError> {
3868 self.apply_rename_local_with_intent(branch, from_stage_id, let_node, new_name, None)
3869 }
3870
3871 /// [`Self::apply_rename_local`] with the write attributed to `intent`
3872 /// (#837 piece A); see [`Self::apply_replace_match_arm_with_intent`].
3873 pub fn apply_rename_local_with_intent(
3874 &self,
3875 branch: &str,
3876 from_stage_id: &str,
3877 let_node: &lex_ast::NodeId,
3878 new_name: &str,
3879 intent: Option<&lex_vcs::Intent>,
3880 ) -> Result<lex_vcs::OpId, StoreError> {
3881 let from_stage = self.get_ast(from_stage_id)?;
3882 // Read the old name before running the transform, so the
3883 // op log records the rename target rather than just the
3884 // new value.
3885 let old_name = read_let_name(&from_stage, let_node).map_err(StoreError::TransformError)?;
3886 let new_stage = lex_ast::rename_local(&from_stage, let_node, new_name)
3887 .map_err(StoreError::TransformError)?;
3888 let sig = lex_ast::sig_id(&from_stage).ok_or(StoreError::CannotPublishImport)?;
3889 let to_stage_id = self.publish(&new_stage)?;
3890 if to_stage_id == from_stage_id {
3891 return Err(StoreError::InvalidTransition(format!(
3892 "rename_local produced the same stage_id `{from_stage_id}`"
3893 )));
3894 }
3895 let head = self.branch_head(branch)?;
3896 let mut candidate: Vec<lex_ast::Stage> = Vec::with_capacity(head.len());
3897 for (other_sig, other_stage_id) in &head {
3898 if other_sig == &sig {
3899 candidate.push(new_stage.clone());
3900 } else {
3901 candidate.push(self.get_ast(other_stage_id)?);
3902 }
3903 }
3904 if !head.contains_key(&sig) {
3905 return Err(StoreError::InvalidTransition(format!(
3906 "sig `{sig}` not on branch `{branch}`'s head"
3907 )));
3908 }
3909 let from_budget = budget_of_stage(&from_stage);
3910 let to_budget = budget_of_stage(&new_stage);
3911 let head_now = self.get_branch(branch)?.and_then(|b| b.head_op);
3912 let kind = lex_vcs::OperationKind::RenameLocal {
3913 sig_id: sig.clone(),
3914 from_stage_id: from_stage_id.to_string(),
3915 to_stage_id: to_stage_id.clone(),
3916 let_node: let_node.as_str().to_string(),
3917 old_name,
3918 new_name: new_name.to_string(),
3919 from_budget,
3920 to_budget,
3921 };
3922 let transition = lex_vcs::StageTransition::Replace {
3923 sig_id: sig.clone(),
3924 from: from_stage_id.to_string(),
3925 to: to_stage_id.clone(),
3926 };
3927 let op = lex_vcs::Operation::new(kind, head_now.into_iter().collect::<Vec<_>>());
3928 let op = match intent {
3929 Some(i) => op.with_intent(i.intent_id.clone()),
3930 None => op,
3931 };
3932 self.apply_operation_checked_with_intent(branch, op, transition, &candidate, intent)
3933 }
3934
3935 /// Apply a typed `InlineLet` transform (#280) — eliminate a
3936 /// `let x := v; body` by substituting `v` for every unshadowed
3937 /// `x` in `body`, then replacing the `Let` node with the
3938 /// substituted body. Same end-to-end shape as
3939 /// [`Self::apply_replace_match_arm`].
3940 pub fn apply_inline_let(
3941 &self,
3942 branch: &str,
3943 from_stage_id: &str,
3944 let_node: &lex_ast::NodeId,
3945 ) -> Result<lex_vcs::OpId, StoreError> {
3946 self.apply_inline_let_with_intent(branch, from_stage_id, let_node, None)
3947 }
3948
3949 /// [`Self::apply_inline_let`] with the write attributed to `intent`
3950 /// (#837 piece A); see [`Self::apply_replace_match_arm_with_intent`].
3951 pub fn apply_inline_let_with_intent(
3952 &self,
3953 branch: &str,
3954 from_stage_id: &str,
3955 let_node: &lex_ast::NodeId,
3956 intent: Option<&lex_vcs::Intent>,
3957 ) -> Result<lex_vcs::OpId, StoreError> {
3958 let from_stage = self.get_ast(from_stage_id)?;
3959 let binding_name =
3960 read_let_name(&from_stage, let_node).map_err(StoreError::TransformError)?;
3961 let new_stage =
3962 lex_ast::inline_let(&from_stage, let_node).map_err(StoreError::TransformError)?;
3963 let sig = lex_ast::sig_id(&from_stage).ok_or(StoreError::CannotPublishImport)?;
3964 let to_stage_id = self.publish(&new_stage)?;
3965 if to_stage_id == from_stage_id {
3966 return Err(StoreError::InvalidTransition(format!(
3967 "inline_let produced the same stage_id `{from_stage_id}`"
3968 )));
3969 }
3970 let head = self.branch_head(branch)?;
3971 let mut candidate: Vec<lex_ast::Stage> = Vec::with_capacity(head.len());
3972 for (other_sig, other_stage_id) in &head {
3973 if other_sig == &sig {
3974 candidate.push(new_stage.clone());
3975 } else {
3976 candidate.push(self.get_ast(other_stage_id)?);
3977 }
3978 }
3979 if !head.contains_key(&sig) {
3980 return Err(StoreError::InvalidTransition(format!(
3981 "sig `{sig}` not on branch `{branch}`'s head"
3982 )));
3983 }
3984 let from_budget = budget_of_stage(&from_stage);
3985 let to_budget = budget_of_stage(&new_stage);
3986 let head_now = self.get_branch(branch)?.and_then(|b| b.head_op);
3987 let kind = lex_vcs::OperationKind::InlineLet {
3988 sig_id: sig.clone(),
3989 from_stage_id: from_stage_id.to_string(),
3990 to_stage_id: to_stage_id.clone(),
3991 let_node: let_node.as_str().to_string(),
3992 binding_name,
3993 from_budget,
3994 to_budget,
3995 };
3996 let transition = lex_vcs::StageTransition::Replace {
3997 sig_id: sig.clone(),
3998 from: from_stage_id.to_string(),
3999 to: to_stage_id.clone(),
4000 };
4001 let op = lex_vcs::Operation::new(kind, head_now.into_iter().collect::<Vec<_>>());
4002 let op = match intent {
4003 Some(i) => op.with_intent(i.intent_id.clone()),
4004 None => op,
4005 };
4006 self.apply_operation_checked_with_intent(branch, op, transition, &candidate, intent)
4007 }
4008
4009 /// Apply a typed `ExtractFunction` transform (#280 slice 4) —
4010 /// extract a sub-expression of `from_stage_id`'s body into a
4011 /// new top-level fn defined by `spec`, and emit two ops tied
4012 /// together by a shared synthetic Intent so `lex op log
4013 /// --intent <id>` groups them.
4014 ///
4015 /// The two ops:
4016 /// 1. `AddFunction { sig_id: <new_fn_sig>, stage_id: <new_fn_stage> }`
4017 /// 2. `ModifyBody { sig_id: <source_sig>, from_stage_id, to_stage_id: <modified> }`
4018 ///
4019 /// The shared Intent's prompt is structured (`extract_function:
4020 /// <new_fn_name>` plus the source identity) so downstream
4021 /// tooling can recover the typed-transform shape from the
4022 /// op-log + intent-log join.
4023 ///
4024 /// Returns `(add_fn_op_id, modify_body_op_id)`.
4025 pub fn apply_extract_function(
4026 &self,
4027 branch: &str,
4028 from_stage_id: &str,
4029 expr_node: &lex_ast::NodeId,
4030 spec: lex_ast::ExtractFnSpec,
4031 ) -> Result<(lex_vcs::OpId, lex_vcs::OpId), StoreError> {
4032 self.apply_extract_function_with_intent(branch, from_stage_id, expr_node, spec, None)
4033 }
4034
4035 /// [`Self::apply_extract_function`] with the write attributed to a
4036 /// caller-supplied `intent` (#837 piece A). Both ops carry that intent
4037 /// instead of the synthetic `[lex.transform.extract_function]` one — the
4038 /// typed shape is still recoverable from the op pair itself
4039 /// (`AddFunction` + `ModifyBody` under one intent). `None` keeps the
4040 /// synthetic intent.
4041 ///
4042 /// Two ops means two gate passes, so this method makes the pair atomic
4043 /// itself: the *final* program (source rewritten + new fn) is checked
4044 /// before anything lands, and if the second op is refused anyway (a
4045 /// non-type gate) the head is rolled back to where it started. A rejected
4046 /// extraction therefore never leaves the new fn on the head without the
4047 /// call that uses it.
4048 pub fn apply_extract_function_with_intent(
4049 &self,
4050 branch: &str,
4051 from_stage_id: &str,
4052 expr_node: &lex_ast::NodeId,
4053 spec: lex_ast::ExtractFnSpec,
4054 supplied_intent: Option<&lex_vcs::Intent>,
4055 ) -> Result<(lex_vcs::OpId, lex_vcs::OpId), StoreError> {
4056 let from_stage = self.get_ast(from_stage_id)?;
4057 let new_fn_name = spec.name.clone();
4058 let (modified_stage, new_fn_stage) =
4059 lex_ast::extract_function(&from_stage, expr_node, spec)
4060 .map_err(StoreError::TransformError)?;
4061
4062 let source_sig = lex_ast::sig_id(&from_stage).ok_or(StoreError::CannotPublishImport)?;
4063 let new_fn_sig = lex_ast::sig_id(&new_fn_stage).ok_or(StoreError::CannotPublishImport)?;
4064 if source_sig == new_fn_sig {
4065 return Err(StoreError::InvalidTransition(format!(
4066 "extract_function produced a sig matching the source `{source_sig}`"
4067 )));
4068 }
4069 let new_fn_stage_id = self.publish(&new_fn_stage)?;
4070 let modified_stage_id = self.publish(&modified_stage)?;
4071 if modified_stage_id == from_stage_id {
4072 return Err(StoreError::InvalidTransition(format!(
4073 "extract_function produced the same stage_id `{from_stage_id}` for the source"
4074 )));
4075 }
4076
4077 let head = self.branch_head(branch)?;
4078 if !head.contains_key(&source_sig) {
4079 return Err(StoreError::InvalidTransition(format!(
4080 "sig `{source_sig}` not on branch `{branch}`'s head"
4081 )));
4082 }
4083
4084 // Synthesize an Intent linking the two ops. The session_id
4085 // / model fields here are not load-bearing — they exist to
4086 // make the IntentId content-addressed; downstream tooling
4087 // reads `prompt` to reconstruct the typed-transform shape.
4088 let synthetic = lex_vcs::Intent::new(
4089 format!(
4090 "[lex.transform.extract_function]\nnew_fn={new_fn_name}\nsource_sig={source_sig}\nfrom_stage={from_stage_id}\nexpr_node={node}",
4091 node = expr_node.as_str(),
4092 ),
4093 "lex-store::apply_extract_function",
4094 lex_vcs::ModelDescriptor {
4095 provider: "lex-store".into(),
4096 name: env!("CARGO_PKG_VERSION").into(),
4097 version: None,
4098 },
4099 None,
4100 );
4101 // A caller-supplied intent is recorded by the gated apply below
4102 // (after the gate, so a refusal leaves no intent); the synthetic one
4103 // keeps its original eager write (legacy behaviour, unchanged).
4104 let (intent, intent_is_supplied) = match supplied_intent {
4105 Some(i) => (i.clone(), true),
4106 None => (synthetic, false),
4107 };
4108 let intent_id = intent.intent_id.clone();
4109 if !intent_is_supplied {
4110 lex_vcs::IntentLog::open(self.root())?.put(&intent)?;
4111 }
4112 let record = intent_is_supplied.then_some(&intent);
4113
4114 // The program the pair produces: source rewritten, new fn alongside.
4115 let mut candidate_with_modified: Vec<lex_ast::Stage> = Vec::with_capacity(head.len() + 1);
4116 for (other_sig, other_stage_id) in &head {
4117 if other_sig == &source_sig {
4118 candidate_with_modified.push(modified_stage.clone());
4119 } else {
4120 candidate_with_modified.push(self.get_ast(other_stage_id)?);
4121 }
4122 }
4123 candidate_with_modified.push(new_fn_stage.clone());
4124
4125 // Pre-flight the *final* program before either op lands: the pair is
4126 // not atomic op-by-op (step 1 advances the head), so a rewrite that
4127 // only fails at step 2 would otherwise strand the new fn on the head.
4128 let head_before = self.get_branch(branch)?.and_then(|b| b.head_op);
4129 {
4130 let stages =
4131 self.with_head_imports(head_before.as_deref(), candidate_with_modified.clone());
4132 if let Err(errors) = self.check_with_resolved_deps(&stages, head_before.as_deref()) {
4133 // Same repair hint the gated apply leaves for a refused op,
4134 // against the op that would have been the modify step.
4135 let would_be = lex_vcs::Operation::new(
4136 lex_vcs::OperationKind::ModifyBody {
4137 sig_id: source_sig.clone(),
4138 from_stage_id: from_stage_id.to_string(),
4139 to_stage_id: modified_stage_id.clone(),
4140 from_budget: budget_of_stage(&from_stage),
4141 to_budget: budget_of_stage(&modified_stage),
4142 to_sig_id: None,
4143 },
4144 head_before.iter().cloned().collect::<Vec<_>>(),
4145 )
4146 .with_intent(intent_id.clone());
4147 let replace = lex_vcs::StageTransition::Replace {
4148 sig_id: source_sig.clone(),
4149 from: from_stage_id.to_string(),
4150 to: modified_stage_id.clone(),
4151 };
4152 let _ = self.record_repair_hint(
4153 &attestable_stage_ids(&replace),
4154 &would_be.op_id(),
4155 &errors,
4156 );
4157 return Err(StoreError::TypeError(errors));
4158 }
4159 }
4160
4161 // Step 1 — emit the AddFunction op for the new fn. Build
4162 // the candidate program by appending the new fn to every
4163 // stage on the current branch head.
4164 let new_fn_effects: std::collections::BTreeSet<String> = match &new_fn_stage {
4165 lex_ast::Stage::FnDecl(fd) => fd.effects.iter().map(|e| e.name.clone()).collect(),
4166 _ => Default::default(),
4167 };
4168 let new_fn_budget = budget_of_stage(&new_fn_stage);
4169 let mut candidate_with_new_fn: Vec<lex_ast::Stage> = Vec::with_capacity(head.len() + 1);
4170 for stage_id in head.values() {
4171 candidate_with_new_fn.push(self.get_ast(stage_id)?);
4172 }
4173 candidate_with_new_fn.push(new_fn_stage.clone());
4174 let head_now = self.get_branch(branch)?.and_then(|b| b.head_op);
4175 let add_op = lex_vcs::Operation::new(
4176 lex_vcs::OperationKind::AddFunction {
4177 sig_id: new_fn_sig.clone(),
4178 stage_id: new_fn_stage_id.clone(),
4179 effects: new_fn_effects,
4180 budget_cost: new_fn_budget,
4181 // Single-op apply path — no package context here.
4182 in_file: None,
4183 },
4184 head_now.into_iter().collect::<Vec<_>>(),
4185 )
4186 .with_intent(intent_id.clone());
4187 let add_transition = lex_vcs::StageTransition::Create {
4188 sig_id: new_fn_sig.clone(),
4189 stage_id: new_fn_stage_id.clone(),
4190 };
4191 let add_op_id = self.apply_operation_checked_with_intent(
4192 branch,
4193 add_op,
4194 add_transition,
4195 &candidate_with_new_fn,
4196 record,
4197 )?;
4198
4199 // Step 2 — emit the ModifyBody op for the source, gated against
4200 // `candidate_with_modified` (built up front for the pre-flight).
4201 let from_budget = budget_of_stage(&from_stage);
4202 let to_budget = budget_of_stage(&modified_stage);
4203 let head_now = self.get_branch(branch)?.and_then(|b| b.head_op);
4204 let modify_op = lex_vcs::Operation::new(
4205 lex_vcs::OperationKind::ModifyBody {
4206 sig_id: source_sig.clone(),
4207 from_stage_id: from_stage_id.to_string(),
4208 to_stage_id: modified_stage_id.clone(),
4209 from_budget,
4210 to_budget,
4211 to_sig_id: None,
4212 },
4213 head_now.into_iter().collect::<Vec<_>>(),
4214 )
4215 .with_intent(intent_id);
4216 let modify_transition = lex_vcs::StageTransition::Replace {
4217 sig_id: source_sig,
4218 from: from_stage_id.to_string(),
4219 to: modified_stage_id,
4220 };
4221 let modify_op_id = match self.apply_operation_checked_with_intent(
4222 branch,
4223 modify_op,
4224 modify_transition,
4225 &candidate_with_modified,
4226 record,
4227 ) {
4228 Ok(id) => id,
4229 Err(e) => {
4230 // Step 1 already moved the head; put it back so a refused
4231 // extraction is all-or-nothing. The orphaned `AddFunction`
4232 // record is reclaimed by `lex op gc`. (Only reachable when a
4233 // gate other than the pre-flighted type-check refuses.)
4234 if let Some(prev) = head_before {
4235 self.set_branch_head_op(branch, prev)?;
4236 }
4237 return Err(e);
4238 }
4239 };
4240
4241 Ok((add_op_id, modify_op_id))
4242 }
4243
4244 /// Propose a stage for `sig_id` without advancing the branch
4245 /// head (#294). Multiple agents can call this concurrently
4246 /// for the same sig — every call lands a fresh `Candidate`
4247 /// op chained off the current head_op. The branch head stays
4248 /// where it was; a later [`Self::promote_candidate`] picks
4249 /// the winner.
4250 ///
4251 /// The caller is responsible for typechecking `new_stage`
4252 /// against whatever program context they consider valid —
4253 /// `propose_candidate` doesn't run the gate. Type errors
4254 /// surface at promotion time, where the candidate is
4255 /// composed back into a candidate program via the standard
4256 /// `apply_operation_checked` path.
4257 ///
4258 /// The stage is published (idempotent on content hash). The
4259 /// `intent_id` is required so downstream consumers can
4260 /// distinguish proposals by author.
4261 pub fn propose_candidate(
4262 &self,
4263 branch: &str,
4264 new_stage: &lex_ast::Stage,
4265 intent_id: &lex_vcs::IntentId,
4266 ) -> Result<lex_vcs::OpId, StoreError> {
4267 let sig = lex_ast::sig_id(new_stage).ok_or(StoreError::CannotPublishImport)?;
4268 let stage_id = self.publish(new_stage)?;
4269 let head_now = self.get_branch(branch)?.and_then(|b| b.head_op);
4270 let op = lex_vcs::Operation::new(
4271 lex_vcs::OperationKind::Candidate {
4272 sig_id: sig,
4273 stage_id,
4274 },
4275 head_now.into_iter().collect::<Vec<_>>(),
4276 )
4277 .with_intent(intent_id.clone());
4278 let transition = lex_vcs::StageTransition::ImportOnly;
4279 self.apply_operation(branch, op, transition)
4280 }
4281
4282 /// List every live `Candidate` op for `sig_id` — i.e. those
4283 /// not yet referenced by any `Promote` op (either as the
4284 /// winner or in the `supersedes` set). Used by `lex stage
4285 /// candidates`. Results are sorted by op_id for
4286 /// reproducibility.
4287 pub fn list_candidates(&self, sig_id: &str) -> Result<Vec<CandidateInfo>, StoreError> {
4288 let log = lex_vcs::OpLog::open(self.root())?;
4289 let all = log.list_all()?;
4290 // Collect the set of candidate op_ids referenced by any
4291 // Promote for this sig. Those candidates are no longer
4292 // live.
4293 let mut referenced: std::collections::BTreeSet<lex_vcs::OpId> = Default::default();
4294 for rec in &all {
4295 if let lex_vcs::OperationKind::Promote {
4296 sig_id: s,
4297 winner_candidate,
4298 supersedes,
4299 ..
4300 } = &rec.op.kind
4301 {
4302 if s != sig_id {
4303 continue;
4304 }
4305 referenced.insert(winner_candidate.clone());
4306 for sup in supersedes {
4307 referenced.insert(sup.clone());
4308 }
4309 }
4310 }
4311 let mut out: Vec<CandidateInfo> = Vec::new();
4312 for rec in all {
4313 let lex_vcs::OperationKind::Candidate {
4314 sig_id: s,
4315 stage_id,
4316 } = &rec.op.kind
4317 else {
4318 continue;
4319 };
4320 if s != sig_id {
4321 continue;
4322 }
4323 if referenced.contains(&rec.op_id) {
4324 continue;
4325 }
4326 out.push(CandidateInfo {
4327 op_id: rec.op_id.clone(),
4328 stage_id: stage_id.clone(),
4329 intent_id: rec.op.intent_id.clone(),
4330 });
4331 }
4332 out.sort_by(|a, b| a.op_id.cmp(&b.op_id));
4333 Ok(out)
4334 }
4335
4336 /// Promote a previously-landed `Candidate` op as the new
4337 /// branch head for its sig (#294). Emits a `Promote` op
4338 /// listing every other live `Candidate` for the same sig
4339 /// in its `supersedes` field. After this lands,
4340 /// [`Self::list_candidates`] returns an empty set for the
4341 /// sig.
4342 ///
4343 /// Re-typechecks the candidate program (winner stage + the
4344 /// rest of the branch) through `apply_operation_checked`, so
4345 /// a candidate that doesn't compose with the current branch
4346 /// state surfaces as `StoreError::TypeError`.
4347 pub fn promote_candidate(
4348 &self,
4349 branch: &str,
4350 candidate_op_id: &lex_vcs::OpId,
4351 ) -> Result<lex_vcs::OpId, StoreError> {
4352 let log = lex_vcs::OpLog::open(self.root())?;
4353 let candidate_rec = log
4354 .get(candidate_op_id)?
4355 .ok_or_else(|| StoreError::UnknownOp(candidate_op_id.clone()))?;
4356 let (sig, winner_stage_id) = match &candidate_rec.op.kind {
4357 lex_vcs::OperationKind::Candidate { sig_id, stage_id } => {
4358 (sig_id.clone(), stage_id.clone())
4359 }
4360 other => {
4361 return Err(StoreError::InvalidTransition(format!(
4362 "op `{candidate_op_id}` is a `{:?}`, not a Candidate",
4363 other
4364 )))
4365 }
4366 };
4367
4368 // #836 G4: a candidate carrying a standing `Reject` review must
4369 // not be promoted. "Standing" = the latest `Review` on the
4370 // winner's stage is a Reject; a later `Approve` (or
4371 // `RequestChanges`, which is advisory, not a veto) lifts it.
4372 // Safe by default: a candidate with no review, or an approved
4373 // one, promotes exactly as before.
4374 if let Some(lex_vcs::ReviewVerdict::Reject) = self.latest_review_verdict(&winner_stage_id)? {
4375 return Err(StoreError::InvalidTransition(format!(
4376 "candidate `{candidate_op_id}` has a standing Reject review on stage `{winner_stage_id}`; record an Approve review (or promote a different candidate) before promoting"
4377 )));
4378 }
4379
4380 // Gather every OTHER live candidate for this sig — the
4381 // ones this Promote will supersede.
4382 let live = self.list_candidates(&sig)?;
4383 let mut supersedes: Vec<lex_vcs::OpId> = live
4384 .iter()
4385 .filter(|c| &c.op_id != candidate_op_id)
4386 .map(|c| c.op_id.clone())
4387 .collect();
4388 supersedes.sort();
4389
4390 // Assemble candidate program: winner stage in place of
4391 // the sig's current head (if any), plus every other sig
4392 // unchanged.
4393 let head = self.branch_head(branch)?;
4394 let winner_stage = self.get_ast(&winner_stage_id)?;
4395 let mut candidate_program: Vec<lex_ast::Stage> = Vec::with_capacity(head.len() + 1);
4396 let mut found = false;
4397 for (other_sig, other_stage_id) in &head {
4398 if other_sig == &sig {
4399 candidate_program.push(winner_stage.clone());
4400 found = true;
4401 } else {
4402 candidate_program.push(self.get_ast(other_stage_id)?);
4403 }
4404 }
4405 if !found {
4406 // Sig doesn't have a head yet — append the winner
4407 // stage to make it a Create.
4408 candidate_program.push(winner_stage.clone());
4409 }
4410 let from_stage_id = head.get(&sig).cloned();
4411 // Budget delta from old head to winner — same shape as
4412 // ModifyBody.
4413 let from_budget = from_stage_id
4414 .as_deref()
4415 .and_then(|s| self.get_ast(s).ok())
4416 .and_then(|s| budget_of_stage(&s));
4417 let to_budget = budget_of_stage(&winner_stage);
4418
4419 let head_now = self.get_branch(branch)?.and_then(|b| b.head_op);
4420 let op = lex_vcs::Operation::new(
4421 lex_vcs::OperationKind::Promote {
4422 sig_id: sig.clone(),
4423 winner_candidate: candidate_op_id.clone(),
4424 winner_stage_id: winner_stage_id.clone(),
4425 supersedes,
4426 from_stage_id: from_stage_id.clone(),
4427 from_budget,
4428 to_budget,
4429 },
4430 head_now.into_iter().collect::<Vec<_>>(),
4431 );
4432 let transition = match &from_stage_id {
4433 Some(from) => lex_vcs::StageTransition::Replace {
4434 sig_id: sig,
4435 from: from.clone(),
4436 to: winner_stage_id,
4437 },
4438 None => lex_vcs::StageTransition::Create {
4439 sig_id: sig,
4440 stage_id: winner_stage_id,
4441 },
4442 };
4443 self.apply_operation_checked(branch, op, transition, &candidate_program)
4444 }
4445
4446 /// `set_branch_head_op` for the durability story on the branch
4447 /// file itself.
4448 pub fn apply_operation(
4449 &self,
4450 branch: &str,
4451 op: lex_vcs::Operation,
4452 transition: lex_vcs::StageTransition,
4453 ) -> Result<lex_vcs::OpId, StoreError> {
4454 let attestable = attestable_stage_ids(&transition);
4455 self.apply_operation_attesting(branch, op, transition, attestable)
4456 }
4457
4458 /// [`Self::apply_operation`] with the stages the write-time attestation
4459 /// gate checks given explicitly rather than read off the transition.
4460 fn apply_operation_attesting(
4461 &self,
4462 branch: &str,
4463 op: lex_vcs::Operation,
4464 transition: lex_vcs::StageTransition,
4465 attestable: Vec<String>,
4466 ) -> Result<lex_vcs::OpId, StoreError> {
4467 let op_effects = op_declared_effects(&op.kind);
4468 self.cas_retry_advance(branch, op, transition, |new_head| {
4469 self.run_required_attestations_gate(branch, &new_head.op_id, &attestable, &op_effects)
4470 })
4471 }
4472
4473 /// CAS retry loop for #262. Single-parent ops are rebuilt on
4474 /// each iteration with the current branch head as parent;
4475 /// the per-iteration callback runs the gate (and TypeCheck
4476 /// emission, for the checked path) between persist and CAS.
4477 /// Merge ops (with 2 parents already set) skip the rebuild —
4478 /// their parents are caller-supplied and meaningful — and get
4479 /// a single attempt; on CAS failure they surface `Contention`.
4480 fn cas_retry_advance<F>(
4481 &self,
4482 branch: &str,
4483 op: lex_vcs::Operation,
4484 transition: lex_vcs::StageTransition,
4485 mut between_persist_and_cas: F,
4486 ) -> Result<lex_vcs::OpId, StoreError>
4487 where
4488 F: FnMut(&lex_vcs::NewHead) -> Result<(), StoreError>,
4489 {
4490 // 32 retries handles up to ~32 concurrent writers racing on
4491 // the same branch tip. Beyond that, surfacing `Contention`
4492 // is the right signal — clients should back off or batch.
4493 const MAX_ATTEMPTS: u32 = 32;
4494 // #992: every single-op write funnels through here, so this is where
4495 // "always-valid HEAD" is enforced for local writes. Refuse a
4496 // transition that binds a sig to a stage filed under a *different*
4497 // sig — the pair can never be read, so the head it produced could
4498 // never be rendered. Checked before anything is persisted.
4499 self.check_pairs_satisfiable(bound_pairs(&transition).iter().map(|(a, b)| (a, b)))?;
4500 // Single-parent ops can be rebuilt on retry; merge ops
4501 // can't (their two parents are meaningful, supplied by the
4502 // merge engine). For merges, single attempt: if CAS
4503 // fails, surface Contention.
4504 let is_rebuildable = op.parents.len() <= 1;
4505 let kind = op.kind.clone();
4506 let intent_id = op.intent_id.clone();
4507
4508 let mut last_io_err: Option<StoreError> = None;
4509 let mut current_op = op;
4510 let current_transition = transition;
4511 // Only rebuild on retries — attempt 1 honors the caller's
4512 // exact op so a user-supplied bogus parent (parents =
4513 // ["someone-else"]) surfaces as `StaleParent` instead of
4514 // being silently corrected.
4515 //
4516 // Exception (#262 follow-up): an op with `parents = []`
4517 // means "I don't care; chain off whatever the current
4518 // head is." Under concurrent apply, attempt 1 can read
4519 // `head_op = Some(opA)` after a sibling writer landed,
4520 // and the persist's parent check fails StaleParent
4521 // unprompted. Rebuild attempt 1 for the empty-parents
4522 // case so the legitimate-race path retries cleanly.
4523 let mut rebuilt_already = false;
4524 for attempt in 1..=MAX_ATTEMPTS {
4525 // Read the current head BEFORE we persist — this is
4526 // the value we'll compare against in the CAS.
4527 let parent = self.get_branch(branch)?.and_then(|b| b.head_op);
4528
4529 // Rebuild the op against the current head, but only
4530 // on retries (not the caller's first attempt) and
4531 // only for single-parent operations. Multi-parent
4532 // (merge) ops are passed through unchanged.
4533 //
4534 // Empty-parents ops also rebuild on attempt 1 (see
4535 // the exception note above) so concurrent apply
4536 // doesn't false-positive on StaleParent.
4537 let should_rebuild = is_rebuildable
4538 && (rebuilt_already || (current_op.parents.is_empty() && parent.is_some()));
4539 if should_rebuild {
4540 current_op = lex_vcs::Operation {
4541 kind: kind.clone(),
4542 parents: parent.iter().cloned().collect(),
4543 intent_id: intent_id.clone(),
4544 };
4545 }
4546
4547 // Persist (idempotent). On `StaleParent` from a retry
4548 // attempt (where we already rebuilt), the head changed
4549 // between our `get_branch` and this `lex_vcs::apply`
4550 // — race; rebuild and continue. On `StaleParent` from
4551 // attempt 1 (caller's input), propagate.
4552 let new_head = match self.persist_op_only_with_parent(
4553 branch,
4554 parent.as_ref(),
4555 current_op.clone(),
4556 current_transition.clone(),
4557 ) {
4558 Ok(nh) => nh,
4559 Err(StoreError::Apply(lex_vcs::ApplyError::StaleParent { .. }))
4560 if is_rebuildable && rebuilt_already =>
4561 {
4562 rebuilt_already = true;
4563 continue;
4564 }
4565 Err(e) => return Err(e),
4566 };
4567
4568 // Run the caller's between-persist-and-cas hook
4569 // (TypeCheck emission + gate). If this fails, the op
4570 // record is durable but orphaned — same semantics as
4571 // pre-#262.
4572 between_persist_and_cas(&new_head)?;
4573
4574 // CAS the branch head. On success: done. On mismatch:
4575 // someone advanced in parallel; retry.
4576 match self.set_branch_head_op_cas(branch, parent, new_head.op_id.clone()) {
4577 Ok(()) => return Ok(new_head.op_id),
4578 Err(crate::branches::CasFailed::Mismatch { .. }) if is_rebuildable => {
4579 // Try again with the new head as parent.
4580 rebuilt_already = true;
4581 continue;
4582 }
4583 Err(crate::branches::CasFailed::Mismatch { .. }) => {
4584 // Merge op: surface immediately — we can't
4585 // rebuild without rerunning the merge engine.
4586 let _ = attempt;
4587 return Err(StoreError::Contention {
4588 branch: branch.into(),
4589 attempts: 1,
4590 });
4591 }
4592 Err(crate::branches::CasFailed::UnknownBranch(b)) => {
4593 return Err(StoreError::UnknownBranch(b));
4594 }
4595 Err(crate::branches::CasFailed::Io(e)) => {
4596 last_io_err = Some(StoreError::Io(std::io::Error::other(e)));
4597 continue;
4598 }
4599 }
4600 }
4601 // Retries exhausted. Prefer surfacing the most recent IO
4602 // error if we hit one; otherwise it's pure CAS contention.
4603 match last_io_err {
4604 Some(e) => Err(e),
4605 None => Err(StoreError::Contention {
4606 branch: branch.into(),
4607 attempts: MAX_ATTEMPTS,
4608 }),
4609 }
4610 }
4611
4612 /// Persist an op against an explicitly-supplied parent. Used
4613 /// by the CAS retry loop in `cas_retry_advance` so the
4614 /// `lex_vcs::apply` parent check matches what we read at the
4615 /// top of the loop iteration (avoids a TOCTOU race against
4616 /// `persist_op_only`'s second read).
4617 fn persist_op_only_with_parent(
4618 &self,
4619 branch: &str,
4620 parent: Option<&lex_vcs::OpId>,
4621 op: lex_vcs::Operation,
4622 transition: lex_vcs::StageTransition,
4623 ) -> Result<lex_vcs::NewHead, StoreError> {
4624 if branch != DEFAULT_BRANCH && self.get_branch(branch)?.is_none() {
4625 return Err(StoreError::UnknownBranch(branch.into()));
4626 }
4627 let log = lex_vcs::OpLog::open(self.root())?;
4628 lex_vcs::apply(&log, parent, op, transition).map_err(|e| match e {
4629 lex_vcs::ApplyError::Persist(io) => StoreError::Io(io),
4630 other => StoreError::Apply(other),
4631 })
4632 }
4633
4634 /// Run the `required_attestations` gate (#245) and the
4635 /// retroactive producer-block gate (#248) over a single op
4636 /// against the store's `policy.json` and attestation log.
4637 ///
4638 /// Failure modes (in order):
4639 ///
4640 /// 1. Producer-block first: if any attestation on the op's
4641 /// stage is from a quarantined tool, refuse with
4642 /// `ProducerBlocked` (#248). Surfaces *before* the
4643 /// required-attestations gate so a clearly-malicious record
4644 /// isn't masked by a missing-Spec error.
4645 /// 2. Required-attestations next: if any required attestation
4646 /// kind is missing, refuse with `BranchAdvanceBlocked`
4647 /// (#245).
4648 ///
4649 /// Loads the policy / attestation log lazily; with no policy
4650 /// file and no `ProducerBlock` attestations the gate is a no-op
4651 /// (default-permissive — matches pre-#245 stores).
4652 fn run_required_attestations_gate(
4653 &self,
4654 branch: &str,
4655 op_id: &lex_vcs::OpId,
4656 stage_ids: &[String],
4657 op_effects: &std::collections::BTreeSet<String>,
4658 ) -> Result<(), StoreError> {
4659 // Build the candidate slice for the new op. Ops with no
4660 // attestable stage (imports, empty merges) get a single
4661 // `None`-stage tuple; both gates skip those.
4662 let new_op_candidate: Vec<(
4663 lex_vcs::OpId,
4664 Option<String>,
4665 std::collections::BTreeSet<String>,
4666 )> = if stage_ids.is_empty() {
4667 vec![(op_id.clone(), None, op_effects.clone())]
4668 } else {
4669 stage_ids
4670 .iter()
4671 .map(|sid| (op_id.clone(), Some(sid.clone()), op_effects.clone()))
4672 .collect()
4673 };
4674 let attest_log = self.attestation_log()?;
4675
4676 // #248 + #256: producer-block gate, walk-back style.
4677 //
4678 // The naive #248 gate only checked the new op's stage. That
4679 // missed contamination on ancestors — once `lex attest
4680 // retro-block` lands, every previously-gated op stays in
4681 // the chain even though its attestations are now from a
4682 // quarantined producer.
4683 //
4684 // #256 fixes this by walking the chain from `head_op` back
4685 // to `last_gate_checkpoint` (or genesis when the checkpoint
4686 // is invalidated), collecting each ancestor's attestable
4687 // stages, and running `check_producer_block` on the
4688 // combined set. After a successful advance,
4689 // `set_branch_head_op` moves the checkpoint to the new
4690 // head (steady-state O(new ops) per advance).
4691 let walk_back_candidate = self.collect_ancestor_candidates(branch)?;
4692 let mut producer_block_candidate = walk_back_candidate;
4693 producer_block_candidate.extend(new_op_candidate.iter().cloned());
4694 crate::policy::check_producer_block(&attest_log, &producer_block_candidate)
4695 .map_err(StoreError::ProducerBlocked)?;
4696
4697 // #245: required-attestations gate. Forward-going only —
4698 // only the new op is checked. Walking back makes no sense
4699 // here: the policy is "this advance must carry these
4700 // attestations," not "every prior op must have."
4701 let policy = match crate::policy::load(self.root())? {
4702 Some(p) if !p.required_attestations.is_empty() => p,
4703 _ => return Ok(()),
4704 };
4705 let waivers =
4706 crate::policy::check_required_attestations(&attest_log, &new_op_candidate, &policy)
4707 .map_err(StoreError::BranchAdvanceBlocked)?;
4708 // #293: emit one `TrustWaived` attestation per waiver so
4709 // the audit trail records every skip. Idempotent on
4710 // attestation_id (content-addressed dedup) — re-running
4711 // the gate with the same state writes the same files.
4712 for w in waivers {
4713 let att = lex_vcs::Attestation::new(
4714 w.stage_id,
4715 Some(op_id.clone()),
4716 None,
4717 lex_vcs::AttestationKind::TrustWaived {
4718 producer: w.producer,
4719 score_thousandths: w.score_thousandths,
4720 threshold_thousandths: w.threshold_thousandths,
4721 kind_tag: w.kind_tag,
4722 },
4723 lex_vcs::AttestationResult::Passed,
4724 trust_waived_producer(),
4725 None,
4726 );
4727 attest_log.put(&att)?;
4728 }
4729 Ok(())
4730 }
4731
4732 /// Walk the branch from `head_op` back to `last_gate_checkpoint`
4733 /// (exclusive) and return the `(op_id, stage_id, op_effects)`
4734 /// tuples for every attestable stage touched by an ancestor
4735 /// (#256). Empty when the branch is fresh, when the checkpoint
4736 /// equals the head, or when the head is None.
4737 fn collect_ancestor_candidates(&self, branch: &str) -> Result<Vec<GateCandidate>, StoreError> {
4738 let b = match self.get_branch(branch)? {
4739 Some(b) => b,
4740 None => return Ok(Vec::new()),
4741 };
4742 let Some(head) = b.head_op else {
4743 return Ok(Vec::new());
4744 };
4745 if Some(&head) == b.last_gate_checkpoint.as_ref() {
4746 // Steady-state common case: previous advance left the
4747 // checkpoint at head. Nothing to re-walk.
4748 return Ok(Vec::new());
4749 }
4750
4751 let log = lex_vcs::OpLog::open(self.root())?;
4752 let walk = log.walk_back(&head, None)?;
4753 let stop_at = b.last_gate_checkpoint.clone();
4754 let mut out = Vec::new();
4755 for rec in walk {
4756 if Some(&rec.op_id) == stop_at.as_ref() {
4757 break;
4758 }
4759 let stages = attestable_stage_ids(&rec.produces);
4760 let effects = op_declared_effects(&rec.op.kind);
4761 if stages.is_empty() {
4762 out.push((rec.op_id.clone(), None, effects));
4763 } else {
4764 for sid in stages {
4765 out.push((rec.op_id.clone(), Some(sid), effects.clone()));
4766 }
4767 }
4768 }
4769 Ok(out)
4770 }
4771}
4772
4773fn stage_name(stage: &Stage) -> &str {
4774 match stage {
4775 Stage::FnDecl(fd) => &fd.name,
4776 Stage::TypeDecl(td) => &td.name,
4777 Stage::Import(i) => &i.alias,
4778 }
4779}
4780
4781fn stage_for_kind<'a>(
4782 kind: &lex_vcs::OperationKind,
4783 stages: &'a [lex_ast::Stage],
4784) -> Option<&'a lex_ast::Stage> {
4785 use lex_vcs::OperationKind::*;
4786 let target_sig = match kind {
4787 // #992: the stage this op publishes belongs to the sig it moves *to*.
4788 // Looking it up under the old sig found nothing in the new program —
4789 // the declaration is there, but under its new (effect-bearing) sig —
4790 // so the AST was silently never written, leaving the head naming a
4791 // pair no store held. Same reason `RenameSymbol` already uses `to`.
4792 // The same holds for any sig-moving modification: a type or example
4793 // change under the same name is a new sig too.
4794 ChangeEffectSig { sig_id, to_sig_id, .. }
4795 | ModifyBody { sig_id, to_sig_id, .. }
4796 | ModifyType { sig_id, to_sig_id, .. } => {
4797 Some(to_sig_id.clone().unwrap_or_else(|| sig_id.clone()))
4798 }
4799 AddFunction { sig_id, .. } | AddType { sig_id, .. } => Some(sig_id.clone()),
4800 RenameSymbol { to, .. } => Some(to.clone()),
4801 _ => None,
4802 };
4803 let target_sig = target_sig?;
4804 stages
4805 .iter()
4806 .find(|s| sig_id(s).as_deref() == Some(target_sig.as_str()))
4807}
4808
4809/// A declaration's `#` comments, as carried across the syntax -> AST boundary.
4810/// Empty for imports, which have no metadata record of their own.
4811fn stage_doc(stage: &Stage) -> Vec<String> {
4812 match stage {
4813 Stage::FnDecl(fd) => fd.doc.clone(),
4814 Stage::TypeDecl(td) => td.doc.clone(),
4815 Stage::Import(_) => Vec::new(),
4816 }
4817}
4818
4819/// The head transition an op kind produces — the single source of truth for
4820/// how an op moves the SigId→StageId map. Public so write paths that build
4821/// ops by hand (`/v1/patch`) cannot drift from `publish_program` (#992).
4822pub fn transition_for_kind(kind: &lex_vcs::OperationKind) -> lex_vcs::StageTransition {
4823 use lex_vcs::OperationKind::*;
4824 use lex_vcs::StageTransition;
4825 match kind {
4826 AddFunction {
4827 sig_id, stage_id, ..
4828 }
4829 | AddType { sig_id, stage_id, .. } => StageTransition::Create {
4830 sig_id: sig_id.clone(),
4831 stage_id: stage_id.clone(),
4832 },
4833 RemoveFunction {
4834 sig_id,
4835 last_stage_id,
4836 }
4837 | RemoveType {
4838 sig_id,
4839 last_stage_id,
4840 } => StageTransition::Remove {
4841 sig_id: sig_id.clone(),
4842 last: last_stage_id.clone(),
4843 },
4844 // #992: an effect change is a *sig* change, so the head entry has to
4845 // move to the new sig rather than keep the old one pointing at the new
4846 // stage. `Replace` would leave `(old_sig, new_stage)` at the head — a
4847 // pair no store can hold, since the new AST declares the new effects
4848 // and is filed under the sig it hashes to. `Rename` is exactly the
4849 // right shape: drop the old sig, bind the new one to the body.
4850 //
4851 // Ops written before `to_sig_id` existed decode as `None` and keep the
4852 // old `Replace` behaviour, so historical logs replay unchanged.
4853 //
4854 // The same applies to `ModifyBody` / `ModifyType` whose signature
4855 // (types, examples, type params) changed under an unchanged name.
4856 ChangeEffectSig {
4857 sig_id,
4858 to_stage_id,
4859 to_sig_id: Some(to_sig),
4860 ..
4861 }
4862 | ModifyBody {
4863 sig_id,
4864 to_stage_id,
4865 to_sig_id: Some(to_sig),
4866 ..
4867 }
4868 | ModifyType {
4869 sig_id,
4870 to_stage_id,
4871 to_sig_id: Some(to_sig),
4872 ..
4873 } if to_sig != sig_id => StageTransition::Rename {
4874 from: sig_id.clone(),
4875 to: to_sig.clone(),
4876 body_stage_id: to_stage_id.clone(),
4877 },
4878 ModifyBody {
4879 sig_id,
4880 from_stage_id,
4881 to_stage_id,
4882 ..
4883 }
4884 | ChangeEffectSig {
4885 sig_id,
4886 from_stage_id,
4887 to_stage_id,
4888 ..
4889 }
4890 | ModifyType {
4891 sig_id,
4892 from_stage_id,
4893 to_stage_id,
4894 ..
4895 }
4896 | ReplaceMatchArm {
4897 sig_id,
4898 from_stage_id,
4899 to_stage_id,
4900 ..
4901 }
4902 | RenameLocal {
4903 sig_id,
4904 from_stage_id,
4905 to_stage_id,
4906 ..
4907 }
4908 | InlineLet {
4909 sig_id,
4910 from_stage_id,
4911 to_stage_id,
4912 ..
4913 } => StageTransition::Replace {
4914 sig_id: sig_id.clone(),
4915 from: from_stage_id.clone(),
4916 to: to_stage_id.clone(),
4917 },
4918 RenameSymbol {
4919 from,
4920 to,
4921 body_stage_id,
4922 ..
4923 } => StageTransition::Rename {
4924 from: from.clone(),
4925 to: to.clone(),
4926 body_stage_id: body_stage_id.clone(),
4927 },
4928 AddImport { .. } | RemoveImport { .. } => StageTransition::ImportOnly,
4929 Merge { .. } => StageTransition::Merge {
4930 entries: Default::default(),
4931 },
4932 // #294: a Candidate proposes a stage without advancing
4933 // the branch. ImportOnly keeps the branch head untouched
4934 // — the stage IS published on disk (Store::propose_candidate
4935 // calls publish before apply), but no head delta lands.
4936 Candidate { .. } => StageTransition::ImportOnly,
4937 // #1007: files are not stages; the head map is untouched.
4938 SetFiles { .. } => StageTransition::FilesOnly,
4939 // A Promote advances the head exactly like ModifyBody
4940 // (or Create when the sig had no head). The winner
4941 // stage is the new branch state for that sig.
4942 Promote {
4943 sig_id,
4944 winner_stage_id,
4945 from_stage_id,
4946 ..
4947 } => match from_stage_id {
4948 Some(from) => StageTransition::Replace {
4949 sig_id: sig_id.clone(),
4950 from: from.clone(),
4951 to: winner_stage_id.clone(),
4952 },
4953 None => StageTransition::Create {
4954 sig_id: sig_id.clone(),
4955 stage_id: winner_stage_id.clone(),
4956 },
4957 },
4958 }
4959}
4960
4961/// Producer identity for TypeCheck attestations emitted by the
4962/// store-write gate. Pinned to this crate's name + version so an
4963/// attestation produced by a different `lex-store` revision is
4964/// distinguishable (content-hashed `produced_by`).
4965fn typecheck_producer() -> lex_vcs::ProducerDescriptor {
4966 lex_vcs::ProducerDescriptor {
4967 tool: "lex-store".into(),
4968 version: env!("CARGO_PKG_VERSION").into(),
4969 model: None,
4970 }
4971}
4972
4973/// The `produced_by.tool` name of the hub's own server-side CI
4974/// attestations (see [`Store::verify_head_and_attest`]). Exported so an
4975/// embedder can reserve it (`lex_api::State::with_reserved_producers`)
4976/// without duplicating the string.
4977pub const HUB_CI_PRODUCER_TOOL: &str = "lex-hub-ci";
4978
4979/// Producer for attestations the hosted CI runner writes (#93). A
4980/// distinct tool name so a `require-attestation` gate — via the
4981/// producer-trust model — can weight "the hub verified this
4982/// server-side" above a client-attached `TypeCheck`.
4983fn hub_ci_producer() -> lex_vcs::ProducerDescriptor {
4984 lex_vcs::ProducerDescriptor {
4985 tool: HUB_CI_PRODUCER_TOOL.into(),
4986 version: env!("CARGO_PKG_VERSION").into(),
4987 model: None,
4988 }
4989}
4990
4991/// Verdict of a hosted-CI run over a branch head (#93).
4992#[derive(Debug, Clone, serde::Serialize)]
4993pub struct HubCiVerdict {
4994 pub passed: bool,
4995 pub checked_stages: usize,
4996 pub attested_stages: usize,
4997 #[serde(skip_serializing_if = "Option::is_none")]
4998 pub detail: Option<String>,
4999}
5000
5001/// Producer for the replay-comparison attestation (#836 G3). Distinct
5002/// tool name so the comparison lex performed is attributable
5003/// separately from the (external) regeneration.
5004fn replay_producer() -> lex_vcs::ProducerDescriptor {
5005 lex_vcs::ProducerDescriptor {
5006 tool: "lex-store-replay".into(),
5007 version: env!("CARGO_PKG_VERSION").into(),
5008 model: None,
5009 }
5010}
5011
5012/// Human/audit label for a recorded model: `provider/name` (`@version`
5013/// when pinned).
5014fn model_label(m: &lex_vcs::ModelDescriptor) -> String {
5015 match &m.version {
5016 Some(v) => format!("{}/{}@{}", m.provider, m.name, v),
5017 None => format!("{}/{}", m.provider, m.name),
5018 }
5019}
5020
5021/// Every `(sig, stage)` binding `t` leaves in the head map — what
5022/// `apply_transition` inserts, as opposed to what it removes or supersedes.
5023/// These are the pairs a render of the resulting head must be able to read
5024/// (#992), so they are what the write-time gate checks.
5025fn bound_pairs(t: &lex_vcs::StageTransition) -> Vec<(String, String)> {
5026 use lex_vcs::StageTransition::*;
5027 match t {
5028 Merge { entries } => entries
5029 .iter()
5030 .filter_map(|(sig, stage)| stage.as_ref().map(|st| (sig.clone(), st.clone())))
5031 .collect(),
5032 other => produced_sig_stage(other).into_iter().collect(),
5033 }
5034}
5035
5036/// The `(sig_id, stage_id)` an op recorded producing, or `None` for a
5037/// transition that produces no stage (removal / import / merge) — those
5038/// have nothing to regenerate for a replay.
5039/// Every identifier-like string in a stage, de-mangled — a conservative
5040/// "names this stage references" set for flagging skipped callees (#868).
5041/// Walks the stage's serialized form, so it covers calls, constructors and
5042/// type references alike without a bespoke AST visitor; a local that happens
5043/// to share a skipped declaration's name over-reports, which is the safe side.
5044/// Whether a stage-load error means the stage is simply *absent* (GC'd,
5045/// superseded, never persisted) — the only case lenient reconstruction may
5046/// skip (#868). Corrupt or unreadable stages are not absent.
5047fn stage_is_absent(e: &StoreError) -> bool {
5048 match e {
5049 StoreError::UnknownStage(_) => true,
5050 StoreError::Io(io) => io.kind() == std::io::ErrorKind::NotFound,
5051 _ => false,
5052 }
5053}
5054
5055fn referenced_names(stage: &Stage) -> std::collections::BTreeSet<String> {
5056 fn walk(v: &serde_json::Value, out: &mut std::collections::BTreeSet<String>) {
5057 match v {
5058 serde_json::Value::String(s) => {
5059 out.insert(demangled_name(s).to_string());
5060 }
5061 serde_json::Value::Array(xs) => xs.iter().for_each(|x| walk(x, out)),
5062 serde_json::Value::Object(m) => m.values().for_each(|x| walk(x, out)),
5063 _ => {}
5064 }
5065 }
5066 let mut out = std::collections::BTreeSet::new();
5067 if let Ok(v) = serde_json::to_value(stage) {
5068 walk(&v, &mut out);
5069 }
5070 out
5071}
5072
5073fn produced_sig_stage(t: &lex_vcs::StageTransition) -> Option<(String, String)> {
5074 use lex_vcs::StageTransition::*;
5075 match t {
5076 Create { sig_id, stage_id } => Some((sig_id.clone(), stage_id.clone())),
5077 Replace { sig_id, to, .. } => Some((sig_id.clone(), to.clone())),
5078 Rename { to, body_stage_id, .. } => Some((to.clone(), body_stage_id.clone())),
5079 Remove { .. } | ImportOnly | FilesOnly | Merge { .. } => None,
5080 }
5081}
5082
5083/// Producer identity for `Examples::Passed` attestations emitted by
5084/// [`Store::record_examples_passed`] (#835). Distinct tool name so
5085/// the activity feed can tell an auto-emitted publish-time examples
5086/// verdict apart from an `lex agent-tool --examples` one.
5087fn examples_producer() -> lex_vcs::ProducerDescriptor {
5088 lex_vcs::ProducerDescriptor {
5089 tool: "lex-store::examples".into(),
5090 version: env!("CARGO_PKG_VERSION").into(),
5091 model: None,
5092 }
5093}
5094
5095/// Prefix of the `produced_by.tool` of `Review` attestations
5096/// ([`Store::record_review`]); the suffix is the reviewer name, so it is
5097/// variable. An embedder that stamps reviewer identity server-side reserves
5098/// the whole family with [`REVIEW_PRODUCER_RESERVATION`].
5099pub const REVIEW_PRODUCER_PREFIX: &str = "lex-store::review:";
5100
5101/// [`REVIEW_PRODUCER_PREFIX`] as a `lex_api::State::with_reserved_producers`
5102/// entry (a trailing `*` means "prefix"), so a client cannot post `Review`
5103/// attestations of its own through `POST /v1/attestations/batch`.
5104pub const REVIEW_PRODUCER_RESERVATION: &str = "lex-store::review:*";
5105
5106/// Producer identity for `Review` attestations (#836). The reviewer's
5107/// own id lives in the kind; this records which tool minted the record.
5108fn review_producer(reviewer: &str) -> lex_vcs::ProducerDescriptor {
5109 lex_vcs::ProducerDescriptor {
5110 tool: format!("{REVIEW_PRODUCER_PREFIX}{reviewer}"),
5111 version: env!("CARGO_PKG_VERSION").into(),
5112 model: None,
5113 }
5114}
5115
5116/// Producer identity for `RepairHint` attestations emitted by
5117/// `apply_operation_checked` on TypeError (#281). Distinct tool
5118/// name from `typecheck_producer` so consumers can filter the
5119/// activity feed for repair hints without scanning kinds.
5120fn repair_hint_producer() -> lex_vcs::ProducerDescriptor {
5121 lex_vcs::ProducerDescriptor {
5122 tool: "lex-store::repair_hint".into(),
5123 version: env!("CARGO_PKG_VERSION").into(),
5124 model: None,
5125 }
5126}
5127
5128/// Producer identity for `TrustWaived` attestations emitted by
5129/// the `required_attestations` gate on a trust-driven waiver
5130/// (#293). Distinct from `typecheck_producer` and `repair_hint`
5131/// so the audit trail clearly shows "the gate let this advance
5132/// through because trust > threshold."
5133fn trust_waived_producer() -> lex_vcs::ProducerDescriptor {
5134 lex_vcs::ProducerDescriptor {
5135 tool: "lex-store::trust_waived".into(),
5136 version: env!("CARGO_PKG_VERSION").into(),
5137 model: None,
5138 }
5139}
5140
5141/// Producer identity for `ProducerTrust` attestations emitted by
5142/// [`Store::recompute_producer_trust`]. The score-derivation
5143/// recompute is its own machine-emittable kind, distinct from
5144/// the gate-side `TrustWaived` emit (#293).
5145fn producer_trust_producer() -> lex_vcs::ProducerDescriptor {
5146 lex_vcs::ProducerDescriptor {
5147 tool: "lex-store::producer_trust".into(),
5148 version: env!("CARGO_PKG_VERSION").into(),
5149 model: None,
5150 }
5151}
5152
5153/// The set of stage_ids a transition introduces. These are the
5154/// stages a successful TypeCheck pass attests *about* — the new
5155/// head produced by Create/Replace, the renamed body, or the per-
5156/// sig resolution of a Merge. Removes and ImportOnly produce no
5157/// attestable stage; the program typechecks but no specific stage
5158/// is the subject of the claim.
5159/// One row of input to the producer-block / required-attestations
5160/// gates: `(op_id, stage_id, op_effects)`. The `stage_id` is
5161/// `None` for ops that don't touch a stage (imports, empty
5162/// merges) — the gate skips those.
5163type GateCandidate = (
5164 lex_vcs::OpId,
5165 Option<String>,
5166 std::collections::BTreeSet<String>,
5167);
5168
5169/// Effect set declared *by the operation itself* (#245). Used by
5170/// the `required_attestations` gate's `EffectsIntersect` clause.
5171///
5172/// Only `AddFunction` and `ChangeEffectSig` carry an effect set in
5173/// their op payload; for everything else this returns the empty
5174/// set, which means `EffectsIntersect` rules don't fire on those
5175/// ops. `Always` rules continue to fire regardless. A future
5176/// improvement is to extract effects from the candidate `Stage`
5177/// for `ModifyBody` ops, but the typed-effects-on-ops path (#247)
5178/// is the cleaner solution and lands separately.
5179fn op_declared_effects(kind: &lex_vcs::OperationKind) -> std::collections::BTreeSet<String> {
5180 use lex_vcs::OperationKind::*;
5181 match kind {
5182 AddFunction { effects, .. } => effects.clone(),
5183 ChangeEffectSig { to_effects, .. } => to_effects.clone(),
5184 _ => std::collections::BTreeSet::new(),
5185 }
5186}
5187
5188fn attestable_stage_ids(transition: &lex_vcs::StageTransition) -> Vec<String> {
5189 use lex_vcs::StageTransition::*;
5190 match transition {
5191 Create { stage_id, .. } => vec![stage_id.clone()],
5192 Replace { to, .. } => vec![to.clone()],
5193 Rename { body_stage_id, .. } => vec![body_stage_id.clone()],
5194 Merge { entries } => entries.values().filter_map(|opt| opt.clone()).collect(),
5195 Remove { .. } | ImportOnly | FilesOnly => Vec::new(),
5196 }
5197}
5198
5199/// True when two `FnDecl`s are identical except for their body — the
5200/// precondition for a pure intra-body three-way merge (#838). The
5201/// signature fields are equal by construction when both share a
5202/// `sig_id`; this also guards the non-signature fields (`type_params`,
5203/// `examples`) so a side that changed those isn't silently dropped.
5204fn fndecl_same_except_body(a: &lex_ast::FnDecl, b: &lex_ast::FnDecl) -> bool {
5205 a.name == b.name
5206 && a.type_params == b.type_params
5207 && a.params == b.params
5208 && a.effects == b.effects
5209 && a.effect_row_var == b.effect_row_var
5210 && a.return_type == b.return_type
5211 && a.examples == b.examples
5212}
5213
5214fn write_canonical_json<T: Serialize>(path: &Path, value: &T) -> Result<(), StoreError> {
5215 let v = serde_json::to_value(value)?;
5216 let s = lex_ast::canon_json::to_canonical_string(&v);
5217 if let Some(parent) = path.parent() {
5218 fs::create_dir_all(parent)?;
5219 }
5220 fs::write(path, s)?;
5221 Ok(())
5222}
5223
5224/// Read the `name` of the `Let` expression at `let_node` inside
5225/// `stage`'s body. Used by [`Store::apply_rename_local`] to record
5226/// the rename source. Returns the same `TransformError` shapes as
5227/// the transformer itself so callers see a consistent error
5228/// vocabulary.
5229fn read_let_name(
5230 stage: &Stage,
5231 let_node: &lex_ast::NodeId,
5232) -> Result<String, lex_ast::TransformError> {
5233 // The transformer is itself a pure function; ask it to perform
5234 // a rename to a sentinel value and read the resulting let's
5235 // original name from the output. Cheaper than duplicating the
5236 // node-walk here, and stays correct as the transform evolves.
5237 //
5238 // We use a sentinel that's invalid as a Lex identifier so even
5239 // if the rename somehow lands, downstream parsing would
5240 // surface it loudly. (The transform path discards the renamed
5241 // value — we only need the *original* name.)
5242 let probed = lex_ast::rename_local(stage, let_node, "__lex_rename_probe__")?;
5243 let Stage::FnDecl(fd) = probed else {
5244 return Err(lex_ast::TransformError::NonFnTarget {
5245 stage_kind: "non-FnDecl",
5246 });
5247 };
5248 // Walk back to the probed let to read its old name from the
5249 // *original* stage — the probed stage's let has already been
5250 // renamed.
5251 let Stage::FnDecl(orig_fd) = stage else {
5252 return Err(lex_ast::TransformError::NonFnTarget {
5253 stage_kind: "non-FnDecl",
5254 });
5255 };
5256 // Path-based lookup matches the transformer's navigation.
5257 let path = parse_let_node_path(let_node.as_str())?;
5258 if path.is_empty() {
5259 return Err(lex_ast::TransformError::NotALet {
5260 at: let_node.as_str().into(),
5261 found_kind: "stage_root",
5262 });
5263 }
5264 if path[0] != orig_fd.params.len() + 1 {
5265 return Err(lex_ast::TransformError::UnknownNode {
5266 at: let_node.as_str().into(),
5267 });
5268 }
5269 let inner = &path[1..];
5270 let target = navigate_to_let(&orig_fd.body, inner, let_node.as_str())?;
5271 let _ = fd; // probed stage discarded
5272 Ok(target.to_string())
5273}
5274
5275fn parse_let_node_path(id: &str) -> Result<Vec<usize>, lex_ast::TransformError> {
5276 let s = id
5277 .strip_prefix("n_")
5278 .ok_or_else(|| lex_ast::TransformError::BadNodeId(id.into()))?;
5279 let mut parts = s.split('.');
5280 let head = parts
5281 .next()
5282 .ok_or_else(|| lex_ast::TransformError::BadNodeId(id.into()))?;
5283 if head != "0" {
5284 return Err(lex_ast::TransformError::BadNodeId(id.into()));
5285 }
5286 let mut out = Vec::new();
5287 for p in parts {
5288 out.push(
5289 p.parse::<usize>()
5290 .map_err(|_| lex_ast::TransformError::BadNodeId(id.into()))?,
5291 );
5292 }
5293 Ok(out)
5294}
5295
5296fn navigate_to_let<'a>(
5297 root: &'a lex_ast::CExpr,
5298 path: &[usize],
5299 at: &str,
5300) -> Result<&'a str, lex_ast::TransformError> {
5301 use lex_ast::CExpr::*;
5302 let mut current = root;
5303 for &idx in path {
5304 current = match current {
5305 Call { callee, args } => {
5306 if idx == 0 {
5307 callee
5308 } else {
5309 args.get(idx - 1)
5310 .ok_or_else(|| lex_ast::TransformError::UnknownNode { at: at.into() })?
5311 }
5312 }
5313 Let { value, body, .. } => match idx {
5314 0 => value,
5315 1 => body,
5316 _ => return Err(lex_ast::TransformError::UnknownNode { at: at.into() }),
5317 },
5318 Match { scrutinee, arms } => {
5319 if idx == 0 {
5320 scrutinee
5321 } else {
5322 let arm_off = idx - 1;
5323 if arm_off % 2 != 1 {
5324 return Err(lex_ast::TransformError::UnknownNode { at: at.into() });
5325 }
5326 let arm_index = arm_off / 2;
5327 &arms
5328 .get(arm_index)
5329 .ok_or_else(|| lex_ast::TransformError::UnknownNode { at: at.into() })?
5330 .body
5331 }
5332 }
5333 Block { statements, result } => {
5334 if idx < statements.len() {
5335 &statements[idx]
5336 } else if idx == statements.len() {
5337 result
5338 } else {
5339 return Err(lex_ast::TransformError::UnknownNode { at: at.into() });
5340 }
5341 }
5342 Constructor { args, .. }
5343 | TupleLit { items: args, .. }
5344 | ListLit { items: args, .. } => args
5345 .get(idx)
5346 .ok_or_else(|| lex_ast::TransformError::UnknownNode { at: at.into() })?,
5347 RecordLit { fields } => {
5348 &fields
5349 .get(idx)
5350 .ok_or_else(|| lex_ast::TransformError::UnknownNode { at: at.into() })?
5351 .value
5352 }
5353 FieldAccess { value, .. } if idx == 0 => value,
5354 Lambda { body, .. } if idx == 0 => body,
5355 BinOp { lhs, rhs, .. } => match idx {
5356 0 => lhs,
5357 1 => rhs,
5358 _ => return Err(lex_ast::TransformError::UnknownNode { at: at.into() }),
5359 },
5360 UnaryOp { expr, .. } if idx == 0 => expr,
5361 Return { value } if idx == 0 => value,
5362 _ => return Err(lex_ast::TransformError::UnknownNode { at: at.into() }),
5363 };
5364 }
5365 let Let { name, .. } = current else {
5366 return Err(lex_ast::TransformError::NotALet {
5367 at: at.into(),
5368 found_kind: lex_cexpr_kind(current),
5369 });
5370 };
5371 Ok(name)
5372}
5373
5374fn lex_cexpr_kind(e: &lex_ast::CExpr) -> &'static str {
5375 use lex_ast::CExpr::*;
5376 match e {
5377 Literal { .. } => "Literal",
5378 Var { .. } => "Var",
5379 Call { .. } => "Call",
5380 Let { .. } => "Let",
5381 Match { .. } => "Match",
5382 Block { .. } => "Block",
5383 Constructor { .. } => "Constructor",
5384 RecordLit { .. } => "RecordLit",
5385 TupleLit { .. } => "TupleLit",
5386 ListLit { .. } => "ListLit",
5387 FieldAccess { .. } => "FieldAccess",
5388 Lambda { .. } => "Lambda",
5389 BinOp { .. } => "BinOp",
5390 UnaryOp { .. } => "UnaryOp",
5391 Return { .. } => "Return",
5392 }
5393}
5394
5395/// Extract the declared `[budget(N)]` integer from a stage's
5396/// effect set, if any (#280 + #247). Returns `None` for stages
5397/// that aren't `FnDecl` or don't carry a budget effect — same
5398/// shape as `lex_vcs::budget_from_effects`.
5399fn budget_of_stage(stage: &Stage) -> Option<u64> {
5400 let fd = match stage {
5401 Stage::FnDecl(fd) => fd,
5402 _ => return None,
5403 };
5404 let mut min_cost: Option<u64> = None;
5405 for eff in &fd.effects {
5406 if eff.name != "budget" {
5407 continue;
5408 }
5409 if let Some(lex_ast::EffectArg::Int { value }) = &eff.arg {
5410 let n = *value as u64;
5411 min_cost = Some(min_cost.map(|c| c.min(n)).unwrap_or(n));
5412 }
5413 }
5414 min_cost
5415}
5416
5417/// Serialize a stage to its canonical-JSON byte form. Used by
5418/// `publish_signed` for delta encoding (#261 slice 3) — both the
5419/// "compute the diff" path and the "write a full snapshot"
5420/// fallback need exactly the same bytes.
5421fn canonical_bytes(stage: &Stage) -> Result<Vec<u8>, StoreError> {
5422 let v = serde_json::to_value(stage)?;
5423 Ok(lex_ast::canon_json::to_canonical_string(&v).into_bytes())
5424}
5425
5426#[allow(dead_code)]
5427fn read_json<T: DeserializeOwned>(path: &Path) -> Result<T, StoreError> {
5428 let bytes = fs::read(path)?;
5429 Ok(serde_json::from_slice(&bytes)?)
5430}
5431
5432/// A declaration's name with the package loader's mangle prefix removed —
5433/// `lib_62579d1a.twice` → `twice`; a name with no prefix is returned unchanged
5434/// (#980).
5435///
5436/// The split is unambiguous: a Lex declaration name cannot contain a `.` in
5437/// source (`fn a.b(..)` is a parse error), so any dot in a stored name is the
5438/// loader's `<stem>_<hash>.` separator and never part of the author's name.
5439pub fn demangled_name(name: &str) -> &str {
5440 match name.split_once('.') {
5441 Some((_prefix, bare)) => bare,
5442 None => name,
5443 }
5444}