返回 CodeWhale
ledger.rs
根目录 / crates / tui / src / tools / subagent / coord / ledger.rs
1 //! Durable coordination records for delegated Work (#4647).
2 //!
3 //! Split out of `coord.rs` unchanged (#5462): the decision/claim/contention
4 //! ledger and its receipt types are the half of that file with consumers
5 //! outside the tool layer — `tui::coordination_detail`, `tui::work_surface`,
6 //! `tui::ui::tests`, and `core::engine::tests` all name these types — while
7 //! the tool wrappers around them are model-surface code. Keeping both in one
8 //! 3.8k-line file meant every read of either started by scrolling past the
9 //! other.
10 //!
11 //! This is a pure move with re-exports: `coord` re-publishes every public item
12 //! under its original path, so no consumer needed an edited import and no
13 //! behavior changed. New coordination *state* belongs here; new coordination
14 //! *tools* belong in `coord.rs`.
15
16 use std::collections::{BTreeSet, HashMap};
17 use std::path::{Path, PathBuf};
18
19 use serde::{Deserialize, Serialize};
20 use serde_json::Value;
21
22 use super::{
23 COORDINATION_PROJECTION_BYTE_LIMIT, COORDINATION_PROJECTION_DECISION_LIMIT,
24 COORDINATION_RECORD_LIMIT,
25 };
26 use crate::tools::subagent::normalize_claim_path;
27
28 /// Coordination records for delegated Work (#4647).
29 ///
30 /// Decision records, write-scope claims, and contention detection for parallel
31 /// agent work. Parallel work may proceed only when scopes and contracts do not
32 /// collide silently.
33 /// Status of a coordination decision.
34 #[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
35 #[serde(rename_all = "snake_case")]
36 pub enum DecisionStatus {
37 Proposed,
38 Accepted,
39 Superseded,
40 }
41
42 /// Serialized coordination state schema. Increment only with an explicit
43 /// migration; restart/replay must never infer a newer contract from old data.
44 pub const COORDINATION_SCHEMA_VERSION: u32 = 1;
45
46 pub(super) const MAX_RECONCILIATION_RETRIES: u32 = 3;
47
48 const fn coordination_schema_version() -> u32 {
49 COORDINATION_SCHEMA_VERSION
50 }
51
52 /// A bounded coordination decision record (#4647).
53 ///
54 /// Persisted with stable subject, concise constraints, one active owner,
55 /// applicability scope, evidence handles, and sequence/version.
56 #[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
57 pub struct DecisionRecord {
58 pub decision_id: String,
59 pub subject: String,
60 pub status: DecisionStatus,
61 pub owner: String,
62 pub scope: Vec<String>,
63 pub constraints: Vec<String>,
64 pub evidence_handles: Vec<String>,
65 pub version: u32,
66 pub sequence: u64,
67 }
68
69 /// A write-scope claim for a write-capable child (#4647).
70 ///
71 /// Declares expected repo-relative paths/trees and named contracts.
72 /// This is coordination metadata, not another approval system.
73 ///
74 /// An `exact_files` entry beneath a declared root binds that root to the
75 /// files listed under it (#6278): peers may then share the root with
76 /// disjoint file claims, and the claim's authorized surface is what
77 /// [`WriteScopeClaim::contains_path`] and [`WriteScopeClaim::overlaps`]
78 /// both see.
79 #[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
80 pub struct WriteScopeClaim {
81 pub owner: String,
82 pub roots: Vec<String>,
83 pub exact_files: Vec<String>,
84 pub contracts: Vec<String>,
85 }
86
87 impl WriteScopeClaim {
88 /// Roots that authorize their whole tree. A root with declared
89 /// `exact_files` beneath it is bound to those files instead: the files
90 /// are the real exclusion boundary, so peers may share one root while
91 /// claiming disjoint outputs (#6278). A root carrying no declared file
92 /// keeps tree-wide authority, and an exact file outside every root
93 /// stands alone.
94 fn open_roots(&self) -> impl Iterator<Item = &String> {
95 self.roots.iter().filter(|root| {
96 !self
97 .exact_files
98 .iter()
99 .any(|file| paths_overlap_by_containment(root, file))
100 })
101 }
102
103 /// Check whether this claim overlaps with another. A claim overlaps when
104 /// either normalized open tree contains the other or exact files collide.
105 #[must_use]
106 pub fn overlaps(&self, other: &WriteScopeClaim) -> bool {
107 for root_a in self.open_roots() {
108 for root_b in other.open_roots() {
109 if paths_overlap_by_containment(root_a, root_b)
110 || paths_overlap_by_containment(root_b, root_a)
111 {
112 return true;
113 }
114 }
115 }
116 for file_a in &self.exact_files {
117 if other
118 .exact_files
119 .iter()
120 .any(|file| paths_overlap_equal(file, file_a))
121 || other
122 .open_roots()
123 .any(|root| paths_overlap_by_containment(root, file_a))
124 {
125 return true;
126 }
127 }
128 for file_b in &other.exact_files {
129 if self
130 .open_roots()
131 .any(|root| paths_overlap_by_containment(root, file_b))
132 {
133 return true;
134 }
135 }
136 if self
137 .contracts
138 .iter()
139 .any(|contract| other.contracts.iter().any(|other| other == contract))
140 {
141 return true;
142 }
143 false
144 }
145
146 #[must_use]
147 pub fn contains_path(&self, path: &str) -> bool {
148 self.exact_files.iter().any(|file| file == path)
149 || self.open_roots().any(|root| path_contains(root, path))
150 }
151 }
152
153 fn path_contains(root: &str, candidate: &str) -> bool {
154 let root = root.trim_end_matches('/');
155 let candidate = candidate.trim_end_matches('/');
156 root == "."
157 || root == candidate
158 || candidate
159 .strip_prefix(root)
160 .is_some_and(|suffix| suffix.starts_with('/'))
161 }
162
163 fn paths_overlap_equal(left: &str, right: &str) -> bool {
164 left == right || left.to_lowercase() == right.to_lowercase()
165 }
166
167 fn paths_overlap_by_containment(root: &str, candidate: &str) -> bool {
168 path_contains(root, candidate) || path_contains(&root.to_lowercase(), &candidate.to_lowercase())
169 }
170
171 #[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
172 pub struct PersistedWriteClaim {
173 pub claim: WriteScopeClaim,
174 pub sequence: u64,
175 #[serde(default)]
176 pub isolated_worktree: bool,
177 /// Which of this claim's roots/exact files existed on disk when it was
178 /// registered (#5906).
179 ///
180 /// Recorded because "the claimed path is gone" and "the claimed path does
181 /// not exist yet" look identical from a later `exists()` call and mean
182 /// opposite things. A claim on a worktree that has since been removed
183 /// describes work nobody can be doing and must stop refusing claimants; a
184 /// claim on a directory its owner is about to create is exactly the
185 /// forward-looking reservation the ledger exists to honor, and disarming
186 /// that would let two writers into the same new tree.
187 ///
188 /// Empty for claims persisted before this field existed and for claims
189 /// that named nothing on disk at the time — both mean "no evidence of
190 /// disappearance", so neither is ever treated as stale.
191 #[serde(default, skip_serializing_if = "Vec::is_empty")]
192 pub present_at_claim: Vec<String>,
193 }
194
195 impl PersistedWriteClaim {
196 /// Whether every path this claim was recorded as actually covering has
197 /// since disappeared from disk (#5906).
198 ///
199 /// Isolated-worktree claims are excluded because they never contend at
200 /// all — the answer would be unused, and their paths are not resolved
201 /// against the coordination root.
202 pub fn names_only_vanished_paths<F>(&self, mut path_exists: F) -> bool
203 where
204 F: FnMut(&str) -> bool,
205 {
206 !self.isolated_worktree
207 && !self.present_at_claim.is_empty()
208 && !self
209 .present_at_claim
210 .iter()
211 .any(|path| path_exists(path.as_str()))
212 }
213 }
214
215 #[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
216 pub struct ReconciliationReceipt {
217 pub reconciliation_id: String,
218 pub subject: String,
219 pub owner: String,
220 pub input_decisions: Vec<String>,
221 pub outcome: String,
222 pub evidence_handles: Vec<String>,
223 /// Preserved candidate branches, patches, or artifact handles. A fan-in
224 /// receipt is not valid if either conflicting candidate was discarded.
225 #[serde(default)]
226 pub candidate_handles: Vec<String>,
227 #[serde(default)]
228 pub retry_count: u32,
229 #[serde(default)]
230 pub retry_limit: u32,
231 #[serde(default)]
232 pub reviewer_evidence_handles: Vec<String>,
233 #[serde(default)]
234 pub verifier_evidence_handles: Vec<String>,
235 #[serde(default)]
236 pub verification_outcome: String,
237 pub sequence: u64,
238 }
239
240 /// Durable receipt for the minimal accepted-decision context projected into a
241 /// child. It records counts and stable ids, never the child's transcript.
242 #[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
243 pub struct ContextProjectionReceipt {
244 pub child_id: String,
245 pub decision_ids: Vec<String>,
246 pub projected_bytes: usize,
247 /// Repeated constraint facts elided across otherwise distinct decisions.
248 /// Decision records themselves are never collapsed by this count.
249 pub deduplicated: usize,
250 /// Relevant unique decisions omitted solely because the hard count or
251 /// byte bound was reached. This must not be conflated with deduplication.
252 #[serde(default)]
253 pub omitted: usize,
254 pub sequence: u64,
255 }
256
257 /// Admission outcome persisted with a write-contention receipt.
258 #[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
259 #[serde(rename_all = "snake_case")]
260 pub enum WriteContentionDisposition {
261 BlockedPendingIsolationOrSerialization,
262 ResolvedBySuccessfulClaim,
263 }
264
265 impl WriteContentionDisposition {
266 #[must_use]
267 pub const fn as_str(self) -> &'static str {
268 match self {
269 Self::BlockedPendingIsolationOrSerialization => {
270 "blocked_pending_isolation_or_serialization"
271 }
272 Self::ResolvedBySuccessfulClaim => "resolved_by_successful_claim",
273 }
274 }
275
276 #[must_use]
277 pub const fn blocks_admission(self) -> bool {
278 matches!(self, Self::BlockedPendingIsolationOrSerialization)
279 }
280 }
281
282 /// Durable non-secret receipt emitted when two active shared-workspace claims
283 /// collide. Rejected scope expansion remains visible after restart.
284 #[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
285 pub struct WriteContentionReceipt {
286 pub claimant: String,
287 pub conflicting_owner: String,
288 pub roots: Vec<String>,
289 pub exact_files: Vec<String>,
290 pub contracts: Vec<String>,
291 pub disposition: WriteContentionDisposition,
292 /// Sequence of the later successful claim that resolved this receipt.
293 /// It intentionally references that claim's sequence instead of consuming
294 /// another ledger sequence.
295 #[serde(default, skip_serializing_if = "Option::is_none")]
296 pub resolution_sequence: Option<u64>,
297 pub sequence: u64,
298 }
299
300 #[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
301 pub struct CoordinationHotPath {
302 pub path: String,
303 pub active_claims: usize,
304 }
305
306 #[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
307 pub struct CoordinationDetailMetrics {
308 pub hottest_paths: Vec<CoordinationHotPath>,
309 pub package_or_module_growth: Option<Value>,
310 pub route_or_cost: Option<Value>,
311 pub note: String,
312 }
313
314 /// One bounded typed projection shared by headless inspection and the TUI.
315 /// It contains durable coordination facts only, never raw reasoning or a
316 /// delegated transcript.
317 #[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
318 pub struct CoordinationDetailProjection {
319 pub schema_version: u32,
320 pub sequence: u64,
321 pub decisions: Vec<DecisionRecord>,
322 pub write_claims: Vec<PersistedWriteClaim>,
323 pub reconciliations: Vec<ReconciliationReceipt>,
324 pub context_projections: Vec<ContextProjectionReceipt>,
325 pub contentions: Vec<WriteContentionReceipt>,
326 pub metrics: CoordinationDetailMetrics,
327 pub bounded: bool,
328 pub limit: usize,
329 /// Whether this process currently holds the workspace coordination flock.
330 /// When false, durable ledger writes are skipped and the UI must say so —
331 /// a counter must never tick on a turn the engine has already settled.
332 #[serde(default = "default_process_lock_held")]
333 pub process_lock_held: bool,
334 /// Human-readable reason when [`Self::process_lock_held`] is false.
335 #[serde(default, skip_serializing_if = "Option::is_none")]
336 pub process_lock_note: Option<String>,
337 }
338
339 fn default_process_lock_held() -> bool {
340 // Legacy projections (tests, older sessions) assume the lock is held so
341 // they do not spuriously light the unavailable banner.
342 true
343 }
344
345 /// Attribution minted from an owner-admitted directory. Serialized data never
346 /// admits that directory: the manager requires its current held identity.
347 #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
348 pub(crate) struct CoordinationClaimScope {
349 pub canonical_root: PathBuf,
350 pub platform: String,
351 pub volume: u64,
352 pub index: u64,
353 }
354
355 impl CoordinationClaimScope {
356 pub(crate) fn validate_shape(&self) -> Result<(), String> {
357 if !self.canonical_root.is_absolute()
358 || self.canonical_root.as_os_str().len() > 32768
359 || self.canonical_root.components().any(|component| {
360 matches!(
361 component,
362 std::path::Component::ParentDir | std::path::Component::CurDir
363 )
364 })
365 || !matches!(self.platform.as_str(), "unix" | "windows")
366 {
367 return Err("coordination root attribution is not bounded and canonical".into());
368 }
369 Ok(())
370 }
371 }
372
373 fn projected_claim(claim: &WriteScopeClaim, root: &Path) -> WriteScopeClaim {
374 let project = |path: &String| {
375 let joined = if path == "." {
376 root.to_path_buf()
377 } else {
378 root.join(path)
379 };
380 joined.to_string_lossy().replace('\\', "/")
381 };
382 WriteScopeClaim {
383 owner: claim.owner.clone(),
384 roots: claim.roots.iter().map(project).collect(),
385 exact_files: claim.exact_files.iter().map(project).collect(),
386 contracts: claim.contracts.clone(),
387 }
388 }
389
390 /// Durable, bounded coordination state owned by `SubAgentManager`.
391 #[derive(Debug, Clone, Serialize, Deserialize)]
392 pub struct CoordinationLedger {
393 #[serde(default = "coordination_schema_version")]
394 pub schema_version: u32,
395 #[serde(default)]
396 pub sequence: u64,
397 #[serde(default)]
398 pub decisions: Vec<DecisionRecord>,
399 #[serde(default)]
400 pub write_claims: Vec<PersistedWriteClaim>,
401 #[serde(default)]
402 pub reconciliations: Vec<ReconciliationReceipt>,
403 #[serde(default)]
404 pub projections: Vec<ContextProjectionReceipt>,
405 #[serde(default)]
406 pub contentions: Vec<WriteContentionReceipt>,
407 /// Root-session provenance for records whose logical owner (`root`) is
408 /// otherwise shared by every conversation in a workspace. Agent-owned
409 /// records can also be stamped here, but may be resolved through their
410 /// immutable `SubAgent.owner_session_id`. Missing legacy provenance is
411 /// never guessed by active-session views.
412 #[serde(default, skip_serializing_if = "HashMap::is_empty")]
413 pub record_sessions: HashMap<u64, String>,
414 /// One bounded sequence-root map for shared claims and path-qualified
415 /// decisions. The historical wire name stays stable; original schema1
416 /// omits it, and scope-bearing schema2 never guesses an absent map.
417 #[serde(default, skip_serializing_if = "Option::is_none")]
418 pub(crate) claim_scopes: Option<HashMap<u64, CoordinationClaimScope>>,
419 }
420
421 impl Default for CoordinationLedger {
422 fn default() -> Self {
423 Self {
424 schema_version: COORDINATION_SCHEMA_VERSION,
425 sequence: 0,
426 decisions: Vec::new(),
427 write_claims: Vec::new(),
428 reconciliations: Vec::new(),
429 projections: Vec::new(),
430 contentions: Vec::new(),
431 record_sessions: HashMap::new(),
432 claim_scopes: None,
433 }
434 }
435 }
436
437 impl CoordinationLedger {
438 fn next_sequence(&mut self) -> u64 {
439 self.sequence = self.sequence.saturating_add(1);
440 self.sequence
441 }
442
443 #[cfg(test)]
444 pub fn record_decision(&mut self, decision: DecisionRecord) -> Result<DecisionRecord, String> {
445 self.record_decision_in_scope(decision, None)
446 }
447
448 pub(crate) fn record_decision_in_scope(
449 &mut self,
450 mut decision: DecisionRecord,
451 scope: Option<CoordinationClaimScope>,
452 ) -> Result<DecisionRecord, String> {
453 self.validate_schema()?;
454 if let Some(scope) = &scope {
455 scope.validate_shape()?;
456 if self
457 .claim_scopes
458 .as_ref()
459 .is_some_and(|scopes| scopes.len() >= COORDINATION_RECORD_LIMIT)
460 {
461 return Err("coordination root attribution capacity reached".into());
462 }
463 }
464 decision.decision_id = decision.decision_id.trim().to_string();
465 if !decision.decision_id.is_empty() {
466 decision.decision_id = bounded_coordination_atom("decision id", &decision.decision_id)?;
467 }
468 decision.subject = bounded_coordination_atom("decision subject", &decision.subject)?;
469 decision.owner = bounded_coordination_atom("decision owner", &decision.owner)?;
470 decision.scope = normalize_coordination_values("decision scope", &decision.scope, 24)?;
471 decision.constraints =
472 normalize_coordination_values("decision constraints", &decision.constraints, 24)?;
473 decision.evidence_handles = normalize_coordination_values(
474 "decision evidence handles",
475 &decision.evidence_handles,
476 24,
477 )?;
478 reject_sensitive_coordination_values(&decision.constraints)?;
479 reject_sensitive_coordination_values(&decision.evidence_handles)?;
480 if decision.subject.trim().is_empty() || decision.owner.trim().is_empty() {
481 return Err("decision subject and owner are required".to_string());
482 }
483 if !decision.decision_id.trim().is_empty()
484 && self
485 .decisions
486 .iter()
487 .any(|existing| existing.decision_id == decision.decision_id)
488 {
489 return Err(format!(
490 "decision id '{}' already exists",
491 decision.decision_id
492 ));
493 }
494 if decision.status == DecisionStatus::Accepted
495 && let Some(existing) = self.decisions.iter().find(|existing| {
496 existing.subject == decision.subject
497 && existing.status == DecisionStatus::Accepted
498 && existing.decision_id != decision.decision_id
499 })
500 {
501 return Err(format!(
502 "subject '{}' already has accepted decision '{}' owned by '{}'; preserve both candidates and use neutral reconciliation",
503 decision.subject, existing.decision_id, existing.owner
504 ));
505 }
506 let next_version = self
507 .decisions
508 .iter()
509 .filter(|existing| existing.subject == decision.subject)
510 .map(|existing| existing.version)
511 .max()
512 .unwrap_or(0)
513 .saturating_add(1);
514 decision.version = decision.version.max(next_version);
515 decision.sequence = self.next_sequence();
516 if decision.decision_id.trim().is_empty() {
517 decision.decision_id = format!("decision_{}", decision.sequence);
518 }
519 self.decisions.push(decision.clone());
520 if self.decisions.len() > COORDINATION_RECORD_LIMIT {
521 let referenced = self
522 .reconciliations
523 .iter()
524 .flat_map(|receipt| receipt.input_decisions.iter())
525 .cloned()
526 .collect::<BTreeSet<_>>();
527 if let Some(index) = self.decisions.iter().position(|existing| {
528 existing.status != DecisionStatus::Accepted
529 && !referenced.contains(&existing.decision_id)
530 }) {
531 self.decisions.remove(index);
532 } else {
533 self.decisions.pop();
534 return Err(
535 "coordination decision capacity is occupied by accepted or reconciled records"
536 .to_string(),
537 );
538 }
539 }
540 self.prune_claim_scopes();
541 if let Some(scope) = scope {
542 self.schema_version = 2;
543 self.claim_scopes
544 .get_or_insert_with(HashMap::new)
545 .insert(decision.sequence, scope);
546 }
547 Ok(decision)
548 }
549
550 pub fn update_decision_status(
551 &mut self,
552 decision_id: &str,
553 status: DecisionStatus,
554 owner: &str,
555 expected_version: u32,
556 ) -> Result<DecisionRecord, String> {
557 self.validate_schema()?;
558 let Some(index) = self
559 .decisions
560 .iter()
561 .position(|decision| decision.decision_id == decision_id)
562 else {
563 return Err(format!("decision '{decision_id}' not found"));
564 };
565 if self.decisions[index].owner != owner {
566 return Err(format!(
567 "decision '{decision_id}' is owned by '{}'; caller '{owner}' cannot change it",
568 self.decisions[index].owner
569 ));
570 }
571 if self.decisions[index].version != expected_version {
572 return Err(format!(
573 "decision '{decision_id}' version changed: expected {expected_version}, current {}",
574 self.decisions[index].version
575 ));
576 }
577 let subject = self.decisions[index].subject.clone();
578 if status == DecisionStatus::Accepted
579 && let Some(existing) =
580 self.decisions
581 .iter()
582 .enumerate()
583 .find_map(|(other_index, existing)| {
584 (other_index != index
585 && existing.subject == subject
586 && existing.status == DecisionStatus::Accepted)
587 .then_some(existing)
588 })
589 {
590 return Err(format!(
591 "subject '{subject}' already has accepted decision '{}' owned by '{}'; preserve both candidates and use neutral reconciliation",
592 existing.decision_id, existing.owner
593 ));
594 }
595 let previous_sequence = self.decisions[index].sequence;
596 let previous_scope = self.claim_scope(previous_sequence).cloned();
597 let sequence = self.next_sequence();
598 let decision = &mut self.decisions[index];
599 decision.status = status;
600 decision.version = decision.version.saturating_add(1);
601 decision.sequence = sequence;
602 let result = decision.clone();
603 if let Some(scopes) = &mut self.claim_scopes {
604 scopes.remove(&previous_sequence);
605 if let Some(scope) = previous_scope {
606 scopes.insert(sequence, scope);
607 }
608 }
609 Ok(result)
610 }
611
612 /// Register a bounded write claim.
613 ///
614 /// `owner_is_active` is the caller's whole answer to "may this existing
615 /// claim refuse a new claimant" — liveness and, since #5906, whether the
616 /// claim's recorded paths still exist. The ledger has no filesystem of its
617 /// own; [`super::SubAgentManager::admissible_coordination_owners`] resolves
618 /// both before calling, and stamps
619 /// [`PersistedWriteClaim::present_at_claim`] on the record afterwards.
620 #[cfg(test)]
621 pub fn register_claim<F>(
622 &mut self,
623 claim: WriteScopeClaim,
624 isolated_worktree: bool,
625 owner_is_active: F,
626 ) -> Result<PersistedWriteClaim, String>
627 where
628 F: FnMut(&str) -> bool,
629 {
630 self.register_claim_in_scope(claim, isolated_worktree, None, None, owner_is_active)
631 }
632
633 pub(crate) fn claim_scope(&self, sequence: u64) -> Option<&CoordinationClaimScope> {
634 self.claim_scopes
635 .as_ref()
636 .and_then(|scopes| scopes.get(&sequence))
637 }
638
639 pub(in crate::tools::subagent) fn prune_claim_scopes(&mut self) {
640 if let Some(scopes) = &mut self.claim_scopes {
641 scopes.retain(|sequence, _| {
642 self.write_claims
643 .iter()
644 .any(|claim| claim.sequence == *sequence)
645 || self
646 .decisions
647 .iter()
648 .any(|decision| decision.sequence == *sequence)
649 });
650 }
651 }
652
653 pub(crate) fn register_claim_in_scope<F>(
654 &mut self,
655 mut claim: WriteScopeClaim,
656 isolated_worktree: bool,
657 scope: Option<CoordinationClaimScope>,
658 original_root: Option<&Path>,
659 mut owner_is_active: F,
660 ) -> Result<PersistedWriteClaim, String>
661 where
662 F: FnMut(&str) -> bool,
663 {
664 self.validate_schema()?;
665 if let Some(scope) = &scope {
666 scope.validate_shape()?;
667 if isolated_worktree || original_root.is_none() {
668 return Err(
669 "attributed claims require the held original root and shared execution".into(),
670 );
671 }
672 }
673 if self
674 .claim_scopes
675 .as_ref()
676 .is_some_and(|scopes| !scopes.is_empty())
677 && original_root.is_none()
678 {
679 return Err("attributed coordination state requires its held original root".into());
680 }
681 claim.owner = bounded_coordination_atom("write claim owner", &claim.owner)?;
682 claim.roots = normalize_claim_paths(&claim.roots)?;
683 claim.exact_files = normalize_claim_paths(&claim.exact_files)?;
684 claim.contracts = normalize_claim_strings(&claim.contracts, 16, 128, "contracts")?;
685 if claim.roots.is_empty() && claim.exact_files.is_empty() && claim.contracts.is_empty() {
686 return Err(
687 "write claim requires an owner and at least one root, file, or contract"
688 .to_string(),
689 );
690 }
691 let replacing_existing_owner = self
692 .write_claims
693 .iter()
694 .any(|existing| existing.claim.owner == claim.owner);
695 if scope.is_some()
696 && self
697 .claim_scopes
698 .as_ref()
699 .is_some_and(|scopes| scopes.len() >= COORDINATION_RECORD_LIMIT)
700 && !self.write_claims.iter().any(|existing| {
701 existing.claim.owner == claim.owner && self.claim_scope(existing.sequence).is_some()
702 })
703 {
704 return Err("coordination root attribution capacity reached".into());
705 }
706 if !replacing_existing_owner && self.write_claims.len() >= COORDINATION_RECORD_LIMIT {
707 let mut inactive = Vec::new();
708 for existing in &self.write_claims {
709 if !owner_is_active(&existing.claim.owner) {
710 inactive.push((existing.sequence, existing.claim.owner.clone()));
711 }
712 }
713 inactive.sort_by_key(|(sequence, _)| *sequence);
714 for (_, owner) in inactive {
715 if self.write_claims.len() < COORDINATION_RECORD_LIMIT {
716 break;
717 }
718 self.write_claims
719 .retain(|existing| existing.claim.owner != owner);
720 }
721 self.prune_claim_scopes();
722 if self.write_claims.len() >= COORDINATION_RECORD_LIMIT {
723 return Err(format!(
724 "write-claim capacity is {COORDINATION_RECORD_LIMIT} active owners; complete, serialize, or isolate existing work before admitting another writer"
725 ));
726 }
727 }
728 if !isolated_worktree
729 && let Some(existing) = self
730 .write_claims
731 .iter()
732 .find(|existing| {
733 !existing.isolated_worktree
734 && existing.claim.owner != claim.owner
735 && owner_is_active(&existing.claim.owner)
736 && match original_root {
737 Some(original) => {
738 let existing_root = self
739 .claim_scope(existing.sequence)
740 .map_or(original, |scope| scope.canonical_root.as_path());
741 let candidate_root = scope
742 .as_ref()
743 .map_or(original, |scope| scope.canonical_root.as_path());
744 projected_claim(&existing.claim, existing_root)
745 .overlaps(&projected_claim(&claim, candidate_root))
746 }
747 None => existing.claim.overlaps(&claim),
748 }
749 })
750 .cloned()
751 {
752 let receipt = WriteContentionReceipt {
753 claimant: claim.owner.clone(),
754 conflicting_owner: existing.claim.owner.clone(),
755 roots: claim.roots.clone(),
756 exact_files: claim.exact_files.clone(),
757 contracts: claim.contracts.clone(),
758 disposition: WriteContentionDisposition::BlockedPendingIsolationOrSerialization,
759 resolution_sequence: None,
760 sequence: self.next_sequence(),
761 };
762 self.contentions.push(receipt);
763 trim_front(&mut self.contentions, COORDINATION_RECORD_LIMIT);
764 return Err(format!(
765 "write-scope contention with {}: requested roots {:?}, files {:?}, contracts {:?} overlap its writable roots {:?}, files {:?}, contracts {:?}. Claim disjoint sibling write_roots (for example tmp/scan/worker-a and tmp/scan/worker-b) or exact_files for each output. A read-only worker uses write_authority=read_only without a write claim. Otherwise serialize the writers (wait for that owner to settle, or cancel it) or use worktree isolation; a nested path inside an existing writable root still overlaps.",
766 existing.claim.owner,
767 claim.roots,
768 claim.exact_files,
769 claim.contracts,
770 existing.claim.roots,
771 existing.claim.exact_files,
772 existing.claim.contracts
773 ));
774 }
775 self.write_claims
776 .retain(|existing| existing.claim.owner != claim.owner);
777 let record = PersistedWriteClaim {
778 claim,
779 sequence: self.next_sequence(),
780 isolated_worktree,
781 // Stamped by the caller, which owns the filesystem.
782 present_at_claim: Vec::new(),
783 };
784 for contention in &mut self.contentions {
785 if contention.claimant == record.claim.owner
786 && contention.disposition.blocks_admission()
787 {
788 contention.disposition = WriteContentionDisposition::ResolvedBySuccessfulClaim;
789 contention.resolution_sequence = Some(record.sequence);
790 }
791 }
792 self.write_claims.push(record.clone());
793 self.prune_claim_scopes();
794 if let Some(scope) = scope {
795 self.schema_version = 2;
796 self.claim_scopes
797 .get_or_insert_with(HashMap::new)
798 .insert(record.sequence, scope);
799 }
800 Ok(record)
801 }
802
803 /// Release write claims whose owner is not a live claimant (#5562).
804 ///
805 /// Claims are durable records and legitimately outlive the agents that
806 /// registered them; that is fine for history, but a stale claim must never
807 /// keep blocking a later writer. An optional owner restricts the sweep to
808 /// one claim owner; live claimants are never removed. Returns the released
809 /// owners (one entry per claim, in ledger order).
810 pub fn release_stale_claims<F>(
811 &mut self,
812 owner: Option<&str>,
813 mut owner_is_active: F,
814 ) -> Result<Vec<String>, String>
815 where
816 F: FnMut(&str) -> bool,
817 {
818 self.validate_schema()?;
819 let owner_filter = match owner {
820 Some(value) => Some(bounded_coordination_atom("write claim owner", value)?),
821 None => None,
822 };
823 let mut released = Vec::new();
824 self.write_claims.retain(|record| {
825 if owner_filter
826 .as_deref()
827 .is_some_and(|wanted| wanted != record.claim.owner)
828 {
829 return true;
830 }
831 if owner_is_active(&record.claim.owner) {
832 return true;
833 }
834 released.push(record.claim.owner.clone());
835 false
836 });
837 // A release is a mutation: advance the ledger sequence so receipts
838 // stamped by the caller stay coherent (coordinator contract).
839 if !released.is_empty() {
840 self.next_sequence();
841 self.prune_claim_scopes();
842 }
843 Ok(released)
844 }
845
846 #[allow(clippy::too_many_arguments)]
847 pub fn reconcile(
848 &mut self,
849 subject: String,
850 owner: String,
851 input_decisions: Vec<String>,
852 outcome: String,
853 evidence_handles: Vec<String>,
854 candidate_handles: Vec<String>,
855 retry_count: u32,
856 retry_limit: u32,
857 reviewer_evidence_handles: Vec<String>,
858 verifier_evidence_handles: Vec<String>,
859 verification_outcome: String,
860 ) -> Result<ReconciliationReceipt, String> {
861 self.validate_schema()?;
862 let subject = bounded_coordination_atom("reconciliation subject", &subject)?;
863 let owner = bounded_coordination_atom("reconciliation owner", &owner)?;
864 let outcome = bounded_coordination_atom("reconciliation outcome", &outcome)?;
865 let verification_outcome = bounded_coordination_atom(
866 "reconciliation verification outcome",
867 &verification_outcome,
868 )?;
869 if input_decisions.len() < 2 {
870 return Err("neutral fan-in requires at least two input decisions".to_string());
871 }
872 if input_decisions.iter().collect::<BTreeSet<_>>().len() != input_decisions.len() {
873 return Err("neutral fan-in decision ids must be distinct".to_string());
874 }
875 if candidate_handles.len() < 2
876 || candidate_handles
877 .iter()
878 .any(|handle| handle.trim().is_empty())
879 {
880 return Err(
881 "neutral fan-in must preserve at least two candidate branch, patch, or artifact handles"
882 .to_string(),
883 );
884 }
885 if candidate_handles.iter().collect::<BTreeSet<_>>().len() != candidate_handles.len() {
886 return Err("neutral fan-in candidate handles must be distinct".to_string());
887 }
888 let input_decisions =
889 normalize_coordination_values("input decision ids", &input_decisions, 24)?;
890 let evidence_handles = normalize_coordination_values(
891 "reconciliation evidence handles",
892 &evidence_handles,
893 24,
894 )?;
895 let candidate_handles =
896 normalize_coordination_values("candidate handles", &candidate_handles, 24)?;
897 if input_decisions.len() < 2 {
898 return Err(
899 "neutral fan-in requires at least two distinct normalized input decisions"
900 .to_string(),
901 );
902 }
903 if candidate_handles.len() < 2 {
904 return Err(
905 "neutral fan-in must preserve at least two distinct normalized candidate handles"
906 .to_string(),
907 );
908 }
909 let reviewer_evidence_handles = normalize_coordination_values(
910 "Reviewer evidence handles",
911 &reviewer_evidence_handles,
912 24,
913 )?;
914 let verifier_evidence_handles = normalize_coordination_values(
915 "Verifier evidence handles",
916 &verifier_evidence_handles,
917 24,
918 )?;
919 reject_sensitive_coordination_values(&evidence_handles)?;
920 reject_sensitive_coordination_values(&candidate_handles)?;
921 reject_sensitive_coordination_values(&reviewer_evidence_handles)?;
922 reject_sensitive_coordination_values(&verifier_evidence_handles)?;
923 if retry_limit == 0 || retry_limit > MAX_RECONCILIATION_RETRIES {
924 return Err(format!(
925 "reconciliation retry_limit must be between 1 and {MAX_RECONCILIATION_RETRIES}"
926 ));
927 }
928 if retry_count > retry_limit {
929 return Err("reconciliation retry_count exceeds retry_limit".to_string());
930 }
931 if reviewer_evidence_handles.is_empty() || verifier_evidence_handles.is_empty() {
932 return Err(
933 "neutral fan-in requires independent Reviewer and Verifier evidence handles"
934 .to_string(),
935 );
936 }
937 if reviewer_evidence_handles.iter().any(|review| {
938 verifier_evidence_handles
939 .iter()
940 .any(|verify| verify == review)
941 }) {
942 return Err("Reviewer and Verifier evidence handles must be independent".to_string());
943 }
944 if !matches!(
945 verification_outcome.as_str(),
946 "verified" | "failed" | "blocked"
947 ) {
948 return Err(
949 "neutral fan-in verification_outcome must be verified, failed, or blocked"
950 .to_string(),
951 );
952 }
953 if input_decisions.iter().any(|id| {
954 !self
955 .decisions
956 .iter()
957 .any(|decision| &decision.decision_id == id)
958 }) {
959 return Err("reconciliation references an unknown decision".to_string());
960 }
961 let inputs = input_decisions
962 .iter()
963 .filter_map(|id| {
964 self.decisions
965 .iter()
966 .find(|decision| &decision.decision_id == id)
967 })
968 .collect::<Vec<_>>();
969 if inputs.iter().any(|decision| decision.subject != subject) {
970 return Err("reconciliation inputs must share the requested subject".to_string());
971 }
972 if inputs.iter().any(|decision| decision.owner == owner) {
973 return Err(
974 "neutral fan-in owner must differ from every input decision owner".to_string(),
975 );
976 }
977 let sequence = self.next_sequence();
978 let receipt = ReconciliationReceipt {
979 reconciliation_id: format!("reconcile_{sequence}"),
980 subject,
981 owner,
982 input_decisions,
983 outcome,
984 evidence_handles,
985 candidate_handles,
986 retry_count,
987 retry_limit,
988 reviewer_evidence_handles,
989 verifier_evidence_handles,
990 verification_outcome,
991 sequence,
992 };
993 self.reconciliations.push(receipt.clone());
994 trim_front(&mut self.reconciliations, COORDINATION_RECORD_LIMIT);
995 Ok(receipt)
996 }
997
998 #[cfg(test)]
999 pub fn project_relevant_decisions(
1000 &mut self,
1001 child_id: &str,
1002 claim: Option<&WriteScopeClaim>,
1003 capabilities: &[String],
1004 ) -> (String, ContextProjectionReceipt) {
1005 self.project_relevant_decisions_in_scope(child_id, claim, capabilities, None, None)
1006 }
1007
1008 pub(crate) fn project_relevant_decisions_in_scope(
1009 &mut self,
1010 child_id: &str,
1011 claim: Option<&WriteScopeClaim>,
1012 capabilities: &[String],
1013 claim_root: Option<&Path>,
1014 original_root: Option<&Path>,
1015 ) -> (String, ContextProjectionReceipt) {
1016 const HEADER: &str = "Accepted coordination decisions relevant to this child (bounded):\n";
1017 let mut seen_constraint_facts = BTreeSet::new();
1018 let mut decision_ids = Vec::new();
1019 let mut lines = Vec::new();
1020 let mut projected_bytes = 0usize;
1021 let mut deduplicated = 0usize;
1022 let mut omitted = 0usize;
1023 for decision in self
1024 .decisions
1025 .iter()
1026 .rev()
1027 .filter(|decision| decision.status == DecisionStatus::Accepted)
1028 .filter(|decision| {
1029 decision_is_relevant(
1030 decision,
1031 claim,
1032 capabilities,
1033 claim_root,
1034 self.claim_scope(decision.sequence)
1035 .map(|scope| scope.canonical_root.as_path())
1036 .or(original_root),
1037 )
1038 })
1039 {
1040 if decision_ids.len() >= COORDINATION_PROJECTION_DECISION_LIMIT {
1041 omitted = omitted.saturating_add(1);
1042 continue;
1043 }
1044 let constraints = decision
1045 .constraints
1046 .iter()
1047 .filter_map(|value| {
1048 let value = bounded_utf8(value, 192);
1049 if seen_constraint_facts.insert(value.clone()) {
1050 Some(value)
1051 } else {
1052 deduplicated = deduplicated.saturating_add(1);
1053 None
1054 }
1055 })
1056 .take(8)
1057 .collect::<Vec<_>>()
1058 .join("; ");
1059 let mut line = format!(
1060 "- {} v{} [{}] owner={}",
1061 decision.subject, decision.version, decision.decision_id, decision.owner,
1062 );
1063 if !constraints.is_empty() {
1064 line.push_str(": ");
1065 line.push_str(&constraints);
1066 }
1067 let line = bounded_utf8(&line, 512);
1068 let added_bytes = line.len().saturating_add(1);
1069 if HEADER
1070 .len()
1071 .saturating_add(projected_bytes)
1072 .saturating_add(added_bytes)
1073 > COORDINATION_PROJECTION_BYTE_LIMIT
1074 {
1075 omitted = omitted.saturating_add(1);
1076 continue;
1077 }
1078 projected_bytes = projected_bytes.saturating_add(added_bytes);
1079 decision_ids.push(decision.decision_id.clone());
1080 lines.push(line);
1081 }
1082 let projection = if lines.is_empty() {
1083 String::new()
1084 } else {
1085 format!("{HEADER}{}", lines.join("\n"))
1086 };
1087 let receipt = ContextProjectionReceipt {
1088 child_id: child_id.to_string(),
1089 decision_ids,
1090 projected_bytes: projection.len(),
1091 deduplicated,
1092 omitted,
1093 sequence: self.next_sequence(),
1094 };
1095 self.projections.push(receipt.clone());
1096 trim_front(&mut self.projections, COORDINATION_RECORD_LIMIT);
1097 (projection, receipt)
1098 }
1099
1100 pub(in crate::tools::subagent) fn validate_replay(&mut self) -> Result<(), String> {
1101 self.validate_schema()?;
1102 if self.decisions.len() > COORDINATION_RECORD_LIMIT
1103 || self.write_claims.len() > COORDINATION_RECORD_LIMIT
1104 || self.reconciliations.len() > COORDINATION_RECORD_LIMIT
1105 || self.projections.len() > COORDINATION_RECORD_LIMIT
1106 || self.contentions.len() > COORDINATION_RECORD_LIMIT
1107 {
1108 return Err("coordination record count exceeds the durable bound".to_string());
1109 }
1110
1111 let mut sequences = BTreeSet::new();
1112 let mut max_sequence = 0_u64;
1113 let mut decision_ids = BTreeSet::new();
1114 let mut accepted_subjects = BTreeSet::new();
1115 for decision in &self.decisions {
1116 bounded_coordination_atom("decision id", &decision.decision_id)?;
1117 bounded_coordination_atom("decision subject", &decision.subject)?;
1118 bounded_coordination_atom("decision owner", &decision.owner)?;
1119 if decision.version == 0 {
1120 return Err(format!(
1121 "decision '{}' has zero version",
1122 decision.decision_id
1123 ));
1124 }
1125 validate_sequence(
1126 decision.sequence,
1127 "decision",
1128 &mut sequences,
1129 &mut max_sequence,
1130 )?;
1131 if !decision_ids.insert(decision.decision_id.clone()) {
1132 return Err(format!("duplicate decision id '{}'", decision.decision_id));
1133 }
1134 if decision.status == DecisionStatus::Accepted
1135 && !accepted_subjects.insert(decision.subject.clone())
1136 {
1137 return Err(format!(
1138 "multiple accepted decisions own subject '{}'",
1139 decision.subject
1140 ));
1141 }
1142 validate_normalized_coordination_values("decision scope", &decision.scope, 24)?;
1143 validate_normalized_coordination_values(
1144 "decision constraints",
1145 &decision.constraints,
1146 24,
1147 )?;
1148 validate_normalized_coordination_values(
1149 "decision evidence handles",
1150 &decision.evidence_handles,
1151 24,
1152 )?;
1153 reject_sensitive_coordination_values(&decision.constraints)?;
1154 reject_sensitive_coordination_values(&decision.evidence_handles)?;
1155 }
1156
1157 let mut claim_owners = BTreeSet::new();
1158 for claim in &self.write_claims {
1159 validate_sequence(
1160 claim.sequence,
1161 "write claim",
1162 &mut sequences,
1163 &mut max_sequence,
1164 )?;
1165 bounded_coordination_atom("write claim owner", &claim.claim.owner)?;
1166 if !claim_owners.insert(claim.claim.owner.clone()) {
1167 return Err(format!(
1168 "duplicate write claim owner '{}'",
1169 claim.claim.owner
1170 ));
1171 }
1172 let roots = normalize_claim_paths(&claim.claim.roots)?;
1173 let exact_files = normalize_claim_paths(&claim.claim.exact_files)?;
1174 let contracts = normalize_claim_strings(&claim.claim.contracts, 16, 128, "contracts")?;
1175 if roots != claim.claim.roots
1176 || exact_files != claim.claim.exact_files
1177 || contracts != claim.claim.contracts
1178 || (roots.is_empty() && exact_files.is_empty() && contracts.is_empty())
1179 {
1180 return Err(format!(
1181 "write claim for '{}' is not normalized and bounded",
1182 claim.claim.owner
1183 ));
1184 }
1185 }
1186
1187 if let Some(scopes) = &self.claim_scopes {
1188 if scopes.len() > COORDINATION_RECORD_LIMIT {
1189 return Err("coordination root attribution exceeds claim capacity".into());
1190 }
1191 for (sequence, scope) in scopes {
1192 scope.validate_shape()?;
1193 if !self
1194 .write_claims
1195 .iter()
1196 .any(|claim| claim.sequence == *sequence && !claim.isolated_worktree)
1197 && !self
1198 .decisions
1199 .iter()
1200 .any(|decision| decision.sequence == *sequence)
1201 {
1202 return Err(
1203 "coordination root attribution has no matching shared claim or decision"
1204 .into(),
1205 );
1206 }
1207 }
1208 }
1209
1210 for receipt in &self.reconciliations {
1211 validate_sequence(
1212 receipt.sequence,
1213 "reconciliation",
1214 &mut sequences,
1215 &mut max_sequence,
1216 )?;
1217 validate_reconciliation_receipt(receipt, &self.decisions)?;
1218 }
1219 for projection in &self.projections {
1220 validate_sequence(
1221 projection.sequence,
1222 "context projection",
1223 &mut sequences,
1224 &mut max_sequence,
1225 )?;
1226 bounded_coordination_atom("projection child", &projection.child_id)?;
1227 if projection.decision_ids.len() > COORDINATION_PROJECTION_DECISION_LIMIT
1228 || projection.projected_bytes > COORDINATION_PROJECTION_BYTE_LIMIT
1229 || projection
1230 .decision_ids
1231 .iter()
1232 .collect::<BTreeSet<_>>()
1233 .len()
1234 != projection.decision_ids.len()
1235 {
1236 return Err(format!(
1237 "context projection for '{}' exceeds its bounds or duplicates decisions",
1238 projection.child_id
1239 ));
1240 }
1241 }
1242 for contention in &self.contentions {
1243 validate_sequence(
1244 contention.sequence,
1245 "contention",
1246 &mut sequences,
1247 &mut max_sequence,
1248 )?;
1249 bounded_coordination_atom("contention claimant", &contention.claimant)?;
1250 bounded_coordination_atom(
1251 "contention conflicting owner",
1252 &contention.conflicting_owner,
1253 )?;
1254 match (contention.disposition, contention.resolution_sequence) {
1255 (WriteContentionDisposition::BlockedPendingIsolationOrSerialization, None) => {}
1256 (WriteContentionDisposition::ResolvedBySuccessfulClaim, Some(sequence))
1257 if sequence > contention.sequence && sequence <= self.sequence => {}
1258 (WriteContentionDisposition::BlockedPendingIsolationOrSerialization, Some(_)) => {
1259 return Err(
1260 "blocked contention receipt cannot carry a resolution sequence".to_string(),
1261 );
1262 }
1263 (WriteContentionDisposition::ResolvedBySuccessfulClaim, _) => {
1264 return Err(
1265 "resolved contention receipt requires a later valid resolution sequence"
1266 .to_string(),
1267 );
1268 }
1269 }
1270 if normalize_claim_paths(&contention.roots)? != contention.roots
1271 || normalize_claim_paths(&contention.exact_files)? != contention.exact_files
1272 || normalize_claim_strings(&contention.contracts, 16, 128, "contracts")?
1273 != contention.contracts
1274 {
1275 return Err("contention receipt paths/contracts are not normalized".to_string());
1276 }
1277 }
1278 if self.sequence < max_sequence {
1279 return Err(format!(
1280 "coordination sequence {} is behind record sequence {max_sequence}",
1281 self.sequence
1282 ));
1283 }
1284 Ok(())
1285 }
1286
1287 fn validate_schema(&self) -> Result<(), String> {
1288 match (self.schema_version, self.claim_scopes.as_ref()) {
1289 (COORDINATION_SCHEMA_VERSION, None) | (2, Some(_)) => Ok(()),
1290 (2, None) => {
1291 Err("coordination schema 2 requires explicit root attribution state".into())
1292 }
1293 (COORDINATION_SCHEMA_VERSION, Some(_)) => {
1294 Err("legacy coordination schema cannot contain attributed roots".into())
1295 }
1296 (version, _) => Err(format!(
1297 "unsupported coordination schema {version}; expected 1 or 2"
1298 )),
1299 }
1300 }
1301 }
1302
1303 fn normalize_claim_paths(paths: &[String]) -> Result<Vec<String>, String> {
1304 if paths.len() > 32 {
1305 return Err("write claim paths accept at most 32 entries".to_string());
1306 }
1307 let mut normalized = Vec::new();
1308 for path in paths {
1309 let path = normalize_claim_path(path)?;
1310 if !normalized.contains(&path) {
1311 normalized.push(path);
1312 }
1313 }
1314 Ok(normalized)
1315 }
1316
1317 fn normalize_claim_strings(
1318 values: &[String],
1319 count_limit: usize,
1320 char_limit: usize,
1321 field: &str,
1322 ) -> Result<Vec<String>, String> {
1323 if values.len() > count_limit {
1324 return Err(format!(
1325 "write claim {field} accepts at most {count_limit} entries"
1326 ));
1327 }
1328 let mut normalized = Vec::new();
1329 for value in values {
1330 let value = value.trim();
1331 if value.is_empty()
1332 || value.chars().count() > char_limit
1333 || value.chars().any(char::is_control)
1334 {
1335 return Err(format!(
1336 "write claim {field} entries must be 1..={char_limit} characters"
1337 ));
1338 }
1339 if !normalized.iter().any(|existing| existing == value) {
1340 normalized.push(value.to_string());
1341 }
1342 }
1343 Ok(normalized)
1344 }
1345
1346 fn bounded_coordination_atom(field: &str, value: &str) -> Result<String, String> {
1347 let value = value.trim();
1348 if value.is_empty()
1349 || value.chars().count() > 512
1350 || value.chars().any(|ch| matches!(ch, '\r' | '\n'))
1351 {
1352 return Err(format!(
1353 "{field} must be one non-empty line of at most 512 characters"
1354 ));
1355 }
1356 Ok(value.to_string())
1357 }
1358
1359 fn normalize_coordination_values(
1360 field: &str,
1361 values: &[String],
1362 limit: usize,
1363 ) -> Result<Vec<String>, String> {
1364 if values.len() > limit {
1365 return Err(format!("{field} accepts at most {limit} entries"));
1366 }
1367 let mut normalized = Vec::new();
1368 for value in values {
1369 let value = bounded_coordination_atom(field, value)?;
1370 if !normalized.contains(&value) {
1371 normalized.push(value);
1372 }
1373 }
1374 Ok(normalized)
1375 }
1376
1377 fn validate_normalized_coordination_values(
1378 field: &str,
1379 values: &[String],
1380 limit: usize,
1381 ) -> Result<(), String> {
1382 if normalize_coordination_values(field, values, limit)? != values {
1383 return Err(format!("{field} is not trimmed and deduplicated"));
1384 }
1385 Ok(())
1386 }
1387
1388 fn reject_sensitive_coordination_values(values: &[String]) -> Result<(), String> {
1389 const SENSITIVE_MARKERS: &[&str] = &[
1390 "secret",
1391 "password",
1392 "api_key",
1393 "api-key",
1394 "authorization:",
1395 "bearer ",
1396 "token=",
1397 "sk-",
1398 "ghp_",
1399 "xoxb-",
1400 "<thinking",
1401 "chain of thought",
1402 "raw reasoning",
1403 ];
1404 for value in values {
1405 let lower = value.to_ascii_lowercase();
1406 if let Some(marker) = SENSITIVE_MARKERS
1407 .iter()
1408 .find(|marker| lower.contains(**marker))
1409 {
1410 return Err(format!(
1411 "coordination metadata rejected sensitive or raw-reasoning marker '{marker}'"
1412 ));
1413 }
1414 }
1415 Ok(())
1416 }
1417
1418 fn validate_sequence(
1419 sequence: u64,
1420 kind: &str,
1421 sequences: &mut BTreeSet<u64>,
1422 max_sequence: &mut u64,
1423 ) -> Result<(), String> {
1424 if sequence == 0 || !sequences.insert(sequence) {
1425 return Err(format!(
1426 "{kind} has a zero or duplicate sequence {sequence}"
1427 ));
1428 }
1429 *max_sequence = (*max_sequence).max(sequence);
1430 Ok(())
1431 }
1432
1433 fn validate_reconciliation_receipt(
1434 receipt: &ReconciliationReceipt,
1435 decisions: &[DecisionRecord],
1436 ) -> Result<(), String> {
1437 bounded_coordination_atom("reconciliation id", &receipt.reconciliation_id)?;
1438 bounded_coordination_atom("reconciliation subject", &receipt.subject)?;
1439 bounded_coordination_atom("reconciliation owner", &receipt.owner)?;
1440 bounded_coordination_atom("reconciliation outcome", &receipt.outcome)?;
1441 if receipt.input_decisions.len() < 2
1442 || receipt
1443 .input_decisions
1444 .iter()
1445 .collect::<BTreeSet<_>>()
1446 .len()
1447 != receipt.input_decisions.len()
1448 {
1449 return Err("reconciliation requires at least two distinct decision ids".to_string());
1450 }
1451 let inputs = receipt
1452 .input_decisions
1453 .iter()
1454 .map(|id| {
1455 decisions
1456 .iter()
1457 .find(|decision| &decision.decision_id == id)
1458 .ok_or_else(|| format!("reconciliation references unknown decision '{id}'"))
1459 })
1460 .collect::<Result<Vec<_>, _>>()?;
1461 if inputs
1462 .iter()
1463 .any(|decision| decision.subject != receipt.subject)
1464 {
1465 return Err("reconciliation inputs must share the requested subject".to_string());
1466 }
1467 if inputs
1468 .iter()
1469 .any(|decision| decision.owner == receipt.owner)
1470 {
1471 return Err("neutral fan-in owner must differ from every candidate owner".to_string());
1472 }
1473 if receipt.candidate_handles.len() < 2
1474 || receipt
1475 .candidate_handles
1476 .iter()
1477 .collect::<BTreeSet<_>>()
1478 .len()
1479 != receipt.candidate_handles.len()
1480 {
1481 return Err("reconciliation requires at least two distinct candidate handles".to_string());
1482 }
1483 validate_normalized_coordination_values("candidate handles", &receipt.candidate_handles, 24)?;
1484 validate_normalized_coordination_values(
1485 "reconciliation evidence handles",
1486 &receipt.evidence_handles,
1487 24,
1488 )?;
1489 validate_normalized_coordination_values(
1490 "Reviewer evidence handles",
1491 &receipt.reviewer_evidence_handles,
1492 24,
1493 )?;
1494 validate_normalized_coordination_values(
1495 "Verifier evidence handles",
1496 &receipt.verifier_evidence_handles,
1497 24,
1498 )?;
1499 reject_sensitive_coordination_values(&receipt.candidate_handles)?;
1500 reject_sensitive_coordination_values(&receipt.evidence_handles)?;
1501 reject_sensitive_coordination_values(&receipt.reviewer_evidence_handles)?;
1502 reject_sensitive_coordination_values(&receipt.verifier_evidence_handles)?;
1503 if receipt.retry_limit == 0
1504 || receipt.retry_limit > MAX_RECONCILIATION_RETRIES
1505 || receipt.retry_count > receipt.retry_limit
1506 {
1507 return Err("reconciliation retry count/limit is invalid".to_string());
1508 }
1509 if receipt.reviewer_evidence_handles.is_empty()
1510 || receipt.verifier_evidence_handles.is_empty()
1511 || receipt.reviewer_evidence_handles.iter().any(|review| {
1512 receipt
1513 .verifier_evidence_handles
1514 .iter()
1515 .any(|verify| verify == review)
1516 })
1517 {
1518 return Err("Reviewer and Verifier evidence must be present and independent".to_string());
1519 }
1520 if !matches!(
1521 receipt.verification_outcome.as_str(),
1522 "verified" | "failed" | "blocked"
1523 ) {
1524 return Err("reconciliation verification outcome is invalid".to_string());
1525 }
1526 Ok(())
1527 }
1528
1529 fn decision_is_relevant(
1530 decision: &DecisionRecord,
1531 claim: Option<&WriteScopeClaim>,
1532 capabilities: &[String],
1533 claim_root: Option<&Path>,
1534 decision_root: Option<&Path>,
1535 ) -> bool {
1536 let relevant_path = |claim: &WriteScopeClaim, value: &str| match (claim_root, decision_root) {
1537 (Some(claim_root), Some(decision_root)) => {
1538 let Ok(relative) = normalize_claim_path(value) else {
1539 return false;
1540 };
1541 let path = if relative == "." {
1542 decision_root.to_path_buf()
1543 } else {
1544 decision_root.join(relative)
1545 };
1546 let contained = path
1547 .strip_prefix(claim_root)
1548 .ok()
1549 .is_some_and(|relative| claim.contains_path(&relative.to_string_lossy()));
1550 contained
1551 || claim
1552 .roots
1553 .iter()
1554 .chain(&claim.exact_files)
1555 .any(|relative| claim_root.join(relative).starts_with(&path))
1556 }
1557 _ => claim_reaches_path(claim, value),
1558 };
1559 if decision.scope.is_empty() {
1560 return true;
1561 }
1562 decision.scope.iter().any(|raw| {
1563 let value = raw.trim();
1564 let (kind, value) = value
1565 .split_once(':')
1566 .map_or(("", value), |(kind, value)| (kind.trim(), value.trim()));
1567 match kind {
1568 "capability" => capabilities.iter().any(|capability| capability == value),
1569 "contract" => {
1570 claim.is_some_and(|claim| claim.contracts.iter().any(|contract| contract == value))
1571 }
1572 "path" => claim.is_some_and(|claim| relevant_path(claim, value)),
1573 _ => {
1574 capabilities.iter().any(|capability| capability == value)
1575 || claim.is_some_and(|claim| {
1576 claim.contracts.iter().any(|contract| contract == value)
1577 || relevant_path(claim, value)
1578 })
1579 }
1580 }
1581 })
1582 }
1583
1584 fn claim_reaches_path(claim: &WriteScopeClaim, path: &str) -> bool {
1585 let Ok(path) = normalize_claim_path(path) else {
1586 return false;
1587 };
1588 claim.contains_path(&path)
1589 || claim.roots.iter().any(|root| path_contains(&path, root))
1590 || claim
1591 .exact_files
1592 .iter()
1593 .any(|file| path_contains(&path, file))
1594 }
1595
1596 fn bounded_utf8(value: &str, byte_limit: usize) -> String {
1597 if value.len() <= byte_limit {
1598 return value.to_string();
1599 }
1600 let mut end = byte_limit;
1601 while !value.is_char_boundary(end) {
1602 end = end.saturating_sub(1);
1603 }
1604 value[..end].to_string()
1605 }
1606
1607 fn trim_front<T>(records: &mut Vec<T>, limit: usize) {
1608 if records.len() > limit {
1609 records.drain(..records.len() - limit);
1610 }
1611 }
1612
1613 #[cfg(test)]
1614 mod records_tests {
1615 use super::*;
1616 use serde_json::json;
1617
1618 fn scoped_root(path: impl AsRef<Path>) -> CoordinationClaimScope {
1619 CoordinationClaimScope {
1620 canonical_root: path.as_ref().to_path_buf(),
1621 platform: if cfg!(windows) { "windows" } else { "unix" }.into(),
1622 volume: 1,
1623 index: 1,
1624 }
1625 }
1626
1627 fn root_claim(owner: &str, path: &str, contract: Option<&str>) -> WriteScopeClaim {
1628 WriteScopeClaim {
1629 owner: owner.into(),
1630 roots: vec![path.into()],
1631 exact_files: Vec::new(),
1632 contracts: contract.map(String::from).into_iter().collect(),
1633 }
1634 }
1635
1636 #[test]
1637 fn attributed_claims_project_nested_and_alias_roots_and_keep_contracts_global() {
1638 let temp = tempfile::tempdir().unwrap();
1639 let root = temp.path().canonicalize().unwrap();
1640 let original = root.join("original");
1641 let mut ledger = CoordinationLedger::default();
1642 ledger
1643 .register_claim_in_scope(
1644 root_claim("a", "src", None),
1645 false,
1646 Some(scoped_root(root.join("shared"))),
1647 Some(&original),
1648 |_| true,
1649 )
1650 .unwrap();
1651 assert!(
1652 ledger
1653 .register_claim_in_scope(
1654 root_claim("nested", ".", None),
1655 false,
1656 Some(scoped_root(root.join("shared/src"))),
1657 Some(&original),
1658 |_| true
1659 )
1660 .is_err()
1661 );
1662 assert!(
1663 ledger
1664 .register_claim_in_scope(
1665 root_claim("alias", "src", None),
1666 false,
1667 Some(scoped_root(root.join("shared"))),
1668 Some(&original),
1669 |_| true
1670 )
1671 .is_err()
1672 );
1673 ledger
1674 .register_claim_in_scope(
1675 root_claim("disjoint", "src", Some("release")),
1676 false,
1677 Some(scoped_root(root.join("other"))),
1678 Some(&original),
1679 |_| true,
1680 )
1681 .unwrap();
1682 assert!(
1683 ledger
1684 .register_claim_in_scope(
1685 root_claim("global-contract", "unrelated", Some("release")),
1686 false,
1687 None,
1688 Some(&original),
1689 |_| true
1690 )
1691 .is_err()
1692 );
1693 }
1694
1695 #[test]
1696 fn attributed_schema_is_sticky_and_absent_map_or_orphan_receipt_refuses() {
1697 let temp = tempfile::tempdir().unwrap();
1698 let root = temp.path().canonicalize().unwrap();
1699 let original = root.join("original");
1700 let mut legacy = CoordinationLedger::default();
1701 legacy
1702 .register_claim(root_claim("legacy", "src", None), false, |_| false)
1703 .unwrap();
1704 let legacy_bytes = serde_json::to_vec(&legacy).unwrap();
1705 assert!(
1706 !String::from_utf8(legacy_bytes.clone())
1707 .unwrap()
1708 .contains("claim_scopes")
1709 );
1710 let mut decoded: CoordinationLedger = serde_json::from_slice(&legacy_bytes).unwrap();
1711 decoded.validate_replay().unwrap();
1712 assert_eq!(serde_json::to_vec(&decoded).unwrap(), legacy_bytes);
1713 decoded
1714 .register_claim_in_scope(
1715 root_claim("extra", "src", None),
1716 false,
1717 Some(scoped_root(root.join("other"))),
1718 Some(&original),
1719 |_| false,
1720 )
1721 .unwrap();
1722 assert_eq!(decoded.schema_version, 2);
1723 decoded.release_stale_claims(None, |_| false).unwrap();
1724 assert!(decoded.claim_scopes.as_ref().unwrap().is_empty());
1725 assert_eq!(decoded.schema_version, 2);
1726 decoded.validate_replay().unwrap();
1727 decoded.claim_scopes = None;
1728 assert!(
1729 decoded
1730 .validate_replay()
1731 .unwrap_err()
1732 .contains("requires explicit")
1733 );
1734 decoded.claim_scopes = Some(HashMap::from([(100, scoped_root(root.join("other")))]));
1735 assert!(
1736 decoded
1737 .validate_replay()
1738 .unwrap_err()
1739 .contains("matching shared claim")
1740 );
1741 decoded.schema_version = 1;
1742 assert!(
1743 decoded
1744 .validate_replay()
1745 .unwrap_err()
1746 .contains("legacy coordination schema")
1747 );
1748 }
1749
1750 #[test]
1751 fn overlapping_roots_detected() {
1752 let a = WriteScopeClaim {
1753 owner: "agent-a".into(),
1754 roots: vec!["src/tui/".into()],
1755 exact_files: vec![],
1756 contracts: vec![],
1757 };
1758 let b = WriteScopeClaim {
1759 owner: "agent-b".into(),
1760 roots: vec!["src/tui/widgets/".into()],
1761 exact_files: vec![],
1762 contracts: vec![],
1763 };
1764 assert!(a.overlaps(&b));
1765 }
1766
1767 #[test]
1768 fn disjoint_roots_no_overlap() {
1769 let a = WriteScopeClaim {
1770 owner: "agent-a".into(),
1771 roots: vec!["src/tui/".into()],
1772 exact_files: vec![],
1773 contracts: vec![],
1774 };
1775 let b = WriteScopeClaim {
1776 owner: "agent-b".into(),
1777 roots: vec!["src/core/".into()],
1778 exact_files: vec![],
1779 contracts: vec![],
1780 };
1781 assert!(!a.overlaps(&b));
1782 }
1783
1784 #[test]
1785 fn exact_file_collision_detected() {
1786 let a = WriteScopeClaim {
1787 owner: "agent-a".into(),
1788 roots: vec![],
1789 exact_files: vec!["src/main.rs".into()],
1790 contracts: vec![],
1791 };
1792 let b = WriteScopeClaim {
1793 owner: "agent-b".into(),
1794 roots: vec![],
1795 exact_files: vec!["src/main.rs".into()],
1796 contracts: vec![],
1797 };
1798 assert!(a.overlaps(&b));
1799 }
1800
1801 #[test]
1802 fn release_stale_claims_releases_only_inactive_owners_and_honours_an_owner_filter() {
1803 let mut ledger = CoordinationLedger::default();
1804 for (owner, root) in [
1805 ("live-builder", "src/live"),
1806 ("zombie-builder", "src/zombie"),
1807 ("prior-session-builder", "src/prior"),
1808 ] {
1809 ledger
1810 .register_claim(
1811 WriteScopeClaim {
1812 owner: owner.into(),
1813 roots: vec![root.into()],
1814 exact_files: Vec::new(),
1815 contracts: Vec::new(),
1816 },
1817 false,
1818 |candidate| candidate == "live-builder",
1819 )
1820 .expect("non-overlapping claim registers");
1821 }
1822 assert_eq!(ledger.write_claims.len(), 3);
1823
1824 let released = ledger
1825 .release_stale_claims(None, |candidate| candidate == "live-builder")
1826 .expect("release sweeps stale claims");
1827 assert_eq!(released, vec!["zombie-builder", "prior-session-builder"]);
1828 assert_eq!(ledger.write_claims.len(), 1);
1829 assert_eq!(ledger.write_claims[0].claim.owner, "live-builder");
1830
1831 // A second sweep is idempotent and never touches a live claimant.
1832 let released = ledger
1833 .release_stale_claims(None, |candidate| candidate == "live-builder")
1834 .expect("idempotent sweep");
1835 assert!(released.is_empty());
1836 }
1837
1838 #[test]
1839 fn release_stale_claims_refuses_to_remove_an_owner_filter_for_a_live_claimant() {
1840 let mut ledger = CoordinationLedger::default();
1841 ledger
1842 .register_claim(
1843 WriteScopeClaim {
1844 owner: "live-builder".into(),
1845 roots: vec!["src/live".into()],
1846 exact_files: Vec::new(),
1847 contracts: Vec::new(),
1848 },
1849 false,
1850 |candidate| candidate == "live-builder",
1851 )
1852 .expect("claim registers");
1853 let released = ledger
1854 .release_stale_claims(Some("live-builder"), |candidate| {
1855 candidate == "live-builder"
1856 })
1857 .expect("a live claimant is never released");
1858 assert!(released.is_empty());
1859
1860 // A stale owner named in the filter IS released; other stale owners are not touched.
1861 ledger
1862 .register_claim(
1863 WriteScopeClaim {
1864 owner: "zombie-builder".into(),
1865 roots: vec!["src/zombie".into()],
1866 exact_files: Vec::new(),
1867 contracts: Vec::new(),
1868 },
1869 false,
1870 |candidate| candidate == "live-builder",
1871 )
1872 .expect("claim registers");
1873 let released = ledger
1874 .release_stale_claims(Some("zombie-builder"), |candidate| {
1875 candidate == "live-builder"
1876 })
1877 .expect("filtered release");
1878 assert_eq!(released, vec!["zombie-builder"]);
1879 assert_eq!(ledger.write_claims.len(), 1);
1880 }
1881
1882 #[test]
1883 fn path_overlap_respects_component_boundaries_and_root_coverage() {
1884 let root = WriteScopeClaim {
1885 owner: "agent-a".into(),
1886 roots: vec!["src".into()],
1887 exact_files: vec![],
1888 contracts: vec![],
1889 };
1890 let sibling = WriteScopeClaim {
1891 owner: "agent-b".into(),
1892 roots: vec!["src2".into()],
1893 exact_files: vec![],
1894 contracts: vec![],
1895 };
1896 let child_file = WriteScopeClaim {
1897 owner: "agent-c".into(),
1898 roots: vec![],
1899 exact_files: vec!["src/lib.rs".into()],
1900 contracts: vec![],
1901 };
1902 assert!(!root.overlaps(&sibling));
1903 assert!(root.overlaps(&child_file));
1904 }
1905
1906 #[test]
1907 fn disjoint_exact_files_share_a_root() {
1908 // #6278: N workers, one results directory, one disjoint file each.
1909 let a = WriteScopeClaim {
1910 owner: "agent-a".into(),
1911 roots: vec!["tmp/scan".into()],
1912 exact_files: vec!["tmp/scan/a.txt".into()],
1913 contracts: vec![],
1914 };
1915 let b = WriteScopeClaim {
1916 owner: "agent-b".into(),
1917 roots: vec!["tmp/scan".into()],
1918 exact_files: vec!["tmp/scan/b.txt".into()],
1919 contracts: vec![],
1920 };
1921 assert!(!a.overlaps(&b));
1922 assert!(!b.overlaps(&a));
1923 // The bound root is the file boundary: the shared tree and the
1924 // peer's file are outside this claim's authorized surface.
1925 assert!(a.contains_path("tmp/scan/a.txt"));
1926 assert!(!a.contains_path("tmp/scan/b.txt"));
1927 assert!(!a.contains_path("tmp/scan/scratch.txt"));
1928 }
1929
1930 #[test]
1931 fn file_bound_roots_still_refuse_real_overlap() {
1932 let a = WriteScopeClaim {
1933 owner: "agent-a".into(),
1934 roots: vec!["tmp/scan".into()],
1935 exact_files: vec!["tmp/scan/a.txt".into()],
1936 contracts: vec![],
1937 };
1938 // The same file under the shared root still contends.
1939 let same_file = WriteScopeClaim {
1940 owner: "agent-b".into(),
1941 roots: vec!["tmp/scan".into()],
1942 exact_files: vec!["tmp/scan/a.txt".into()],
1943 contracts: vec![],
1944 };
1945 assert!(a.overlaps(&same_file));
1946 // An open root — no files declared beneath it — keeps tree-wide
1947 // authority and still contends with a file beneath it.
1948 let open_root = WriteScopeClaim {
1949 owner: "agent-c".into(),
1950 roots: vec!["tmp/scan".into()],
1951 exact_files: vec![],
1952 contracts: vec![],
1953 };
1954 assert!(a.overlaps(&open_root));
1955 // A root whose files sit elsewhere stays open: union claims are
1956 // unchanged for files outside the claimed trees.
1957 let mixed = WriteScopeClaim {
1958 owner: "agent-d".into(),
1959 roots: vec!["src".into()],
1960 exact_files: vec!["README.md".into()],
1961 contracts: vec![],
1962 };
1963 assert!(mixed.contains_path("src/lib.rs"));
1964 assert!(mixed.contains_path("README.md"));
1965 }
1966
1967 #[test]
1968 fn disjoint_file_claims_under_one_root_both_register() {
1969 let mut ledger = CoordinationLedger::default();
1970 for (owner, file) in [("agent-a", "tmp/scan/a.txt"), ("agent-b", "tmp/scan/b.txt")] {
1971 ledger
1972 .register_claim(
1973 WriteScopeClaim {
1974 owner: owner.into(),
1975 roots: vec!["tmp/scan".into()],
1976 exact_files: vec![file.into()],
1977 contracts: vec![],
1978 },
1979 false,
1980 |_| true,
1981 )
1982 .expect("disjoint files under a shared root register");
1983 }
1984 assert_eq!(ledger.write_claims.len(), 2);
1985 assert!(ledger.contentions.is_empty());
1986 }
1987
1988 #[test]
1989 fn legacy_blocked_contention_wire_defaults_to_unresolved() {
1990 let receipt: WriteContentionReceipt = serde_json::from_value(json!({
1991 "claimant": "agent-b",
1992 "conflicting_owner": "agent-a",
1993 "roots": ["src"],
1994 "exact_files": [],
1995 "contracts": ["public-api"],
1996 "disposition": "blocked_pending_isolation_or_serialization",
1997 "sequence": 3
1998 }))
1999 .expect("existing blocked wire receipt remains readable");
2000
2001 assert_eq!(
2002 receipt.disposition,
2003 WriteContentionDisposition::BlockedPendingIsolationOrSerialization
2004 );
2005 assert_eq!(receipt.resolution_sequence, None);
2006 assert!(receipt.disposition.blocks_admission());
2007 }
2008
2009 #[test]
2010 fn active_shared_claims_contend_but_isolated_claims_do_not() {
2011 let mut ledger = CoordinationLedger::default();
2012 let first = WriteScopeClaim {
2013 owner: "agent-a".into(),
2014 roots: vec!["src".into()],
2015 exact_files: vec![],
2016 contracts: vec!["public-api".into()],
2017 };
2018 ledger.register_claim(first, false, |_| false).unwrap();
2019 let second = WriteScopeClaim {
2020 owner: "agent-b".into(),
2021 roots: vec!["docs".into()],
2022 exact_files: vec![],
2023 contracts: vec!["public-api".into()],
2024 };
2025 let err = ledger
2026 .register_claim(second.clone(), false, |owner| owner == "agent-a")
2027 .unwrap_err();
2028 assert!(
2029 err.contains("contention") && err.contains("agent-a"),
2030 "{err}"
2031 );
2032 assert_eq!(ledger.contentions.len(), 1);
2033 assert_eq!(ledger.contentions[0].claimant, "agent-b");
2034 assert_eq!(ledger.contentions[0].conflicting_owner, "agent-a");
2035 assert_eq!(
2036 ledger.contentions[0].disposition,
2037 WriteContentionDisposition::BlockedPendingIsolationOrSerialization
2038 );
2039 assert_eq!(
2040 serde_json::to_value(&ledger.contentions[0]).unwrap()["disposition"],
2041 json!("blocked_pending_isolation_or_serialization")
2042 );
2043 let resolving_claim = ledger
2044 .register_claim(second, true, |owner| owner == "agent-a")
2045 .expect("isolated claim resolves the blocked admission");
2046 assert_eq!(
2047 ledger.contentions[0].disposition,
2048 WriteContentionDisposition::ResolvedBySuccessfulClaim
2049 );
2050 assert_eq!(
2051 ledger.contentions[0].resolution_sequence,
2052 Some(resolving_claim.sequence)
2053 );
2054 }
2055
2056 #[test]
2057 fn active_write_claims_are_never_evicted_by_receipt_retention() {
2058 let mut ledger = CoordinationLedger::default();
2059 for index in 0..COORDINATION_RECORD_LIMIT {
2060 ledger
2061 .register_claim(
2062 WriteScopeClaim {
2063 owner: format!("agent-{index:03}"),
2064 roots: vec![format!("pkg-{index:03}")],
2065 exact_files: vec![],
2066 contracts: vec![],
2067 },
2068 false,
2069 |_| true,
2070 )
2071 .unwrap();
2072 }
2073 let error = ledger
2074 .register_claim(
2075 WriteScopeClaim {
2076 owner: "agent-over-cap".into(),
2077 roots: vec!["new-package".into()],
2078 exact_files: vec![],
2079 contracts: vec![],
2080 },
2081 false,
2082 |_| true,
2083 )
2084 .expect_err("all-active capacity must fail before evicting ownership");
2085 assert!(error.contains("active owners"), "{error}");
2086 assert_eq!(ledger.write_claims.len(), COORDINATION_RECORD_LIMIT);
2087 assert!(
2088 ledger
2089 .write_claims
2090 .iter()
2091 .any(|record| record.claim.owner == "agent-000")
2092 );
2093 }
2094
2095 #[test]
2096 fn accepted_decisions_require_owner_and_explicit_neutral_reconciliation() {
2097 let mut ledger = CoordinationLedger::default();
2098 let make = |id: &str, owner: &str, status| DecisionRecord {
2099 decision_id: id.into(),
2100 subject: "storage".into(),
2101 status,
2102 owner: owner.into(),
2103 scope: vec!["router".into()],
2104 constraints: vec![],
2105 evidence_handles: vec![format!("receipt:{id}")],
2106 version: 1,
2107 sequence: 0,
2108 };
2109 ledger
2110 .record_decision(make("a", "agent-a", DecisionStatus::Accepted))
2111 .unwrap();
2112 ledger
2113 .record_decision(make("b", "agent-b", DecisionStatus::Proposed))
2114 .unwrap();
2115 let owner_error = ledger
2116 .update_decision_status("b", DecisionStatus::Accepted, "root", 2)
2117 .unwrap_err();
2118 assert!(owner_error.contains("owned by 'agent-b'"), "{owner_error}");
2119 let stale = ledger
2120 .update_decision_status("b", DecisionStatus::Accepted, "agent-b", 1)
2121 .unwrap_err();
2122 assert!(stale.contains("expected 1, current 2"), "{stale}");
2123 let conflict = ledger
2124 .update_decision_status("b", DecisionStatus::Accepted, "agent-b", 2)
2125 .unwrap_err();
2126 assert!(conflict.contains("neutral reconciliation"), "{conflict}");
2127 ledger
2128 .update_decision_status("a", DecisionStatus::Superseded, "agent-a", 1)
2129 .unwrap();
2130 ledger
2131 .update_decision_status("b", DecisionStatus::Accepted, "agent-b", 2)
2132 .unwrap();
2133 let receipt = ledger
2134 .reconcile(
2135 "storage".into(),
2136 "root".into(),
2137 vec!["a".into(), "b".into()],
2138 "use bounded origin-session artifacts".into(),
2139 vec!["test:coord".into()],
2140 vec!["branch:agent-a".into(), "branch:agent-b".into()],
2141 1,
2142 3,
2143 vec!["review:independent".into()],
2144 vec!["verify:locked".into()],
2145 "verified".into(),
2146 )
2147 .unwrap();
2148 assert_eq!(receipt.input_decisions.len(), 2);
2149 assert!(receipt.sequence > ledger.decisions[1].sequence);
2150 }
2151
2152 #[test]
2153 fn relevant_decision_projection_is_deduplicated_bounded_and_receipted() {
2154 let mut ledger = CoordinationLedger::default();
2155 for (id, subject, scope) in [
2156 ("file", "file-contract", "path:src"),
2157 ("docs", "docs-contract", "path:docs"),
2158 ("api", "api-contract", "contract:public-api"),
2159 ] {
2160 ledger
2161 .record_decision(DecisionRecord {
2162 decision_id: id.into(),
2163 subject: subject.into(),
2164 status: DecisionStatus::Accepted,
2165 owner: "planner".into(),
2166 scope: vec![scope.into()],
2167 constraints: vec!["bounded".into(), "bounded".into()],
2168 evidence_handles: vec![format!("receipt:{id}")],
2169 version: 1,
2170 sequence: 0,
2171 })
2172 .unwrap();
2173 }
2174 let claim = WriteScopeClaim {
2175 owner: "worker".into(),
2176 roots: vec!["src/tui".into()],
2177 exact_files: vec![],
2178 contracts: vec!["public-api".into()],
2179 };
2180 let (projection, receipt) =
2181 ledger.project_relevant_decisions("worker", Some(&claim), &["File".into()]);
2182 assert!(projection.contains("file-contract"), "{projection}");
2183 assert!(projection.contains("api-contract"), "{projection}");
2184 assert!(!projection.contains("docs-contract"), "{projection}");
2185 assert!(projection.len() <= COORDINATION_PROJECTION_BYTE_LIMIT);
2186 assert_eq!(receipt.decision_ids, vec!["api", "file"]);
2187 assert_eq!(receipt.deduplicated, 1);
2188 assert_eq!(ledger.projections.last(), Some(&receipt));
2189 }
2190
2191 #[test]
2192 fn projection_receipt_distinguishes_unique_omissions_from_deduplication() {
2193 let mut ledger = CoordinationLedger::default();
2194 for index in 0..(COORDINATION_PROJECTION_DECISION_LIMIT + 2) {
2195 ledger
2196 .record_decision(DecisionRecord {
2197 decision_id: format!("decision-{index}"),
2198 subject: format!("subject-{index}"),
2199 status: DecisionStatus::Accepted,
2200 owner: "planner".into(),
2201 scope: vec!["path:src".into()],
2202 constraints: vec![format!("constraint-{index}")],
2203 evidence_handles: vec![format!("receipt:{index}")],
2204 version: 1,
2205 sequence: 0,
2206 })
2207 .unwrap();
2208 }
2209 let claim = WriteScopeClaim {
2210 owner: "worker".into(),
2211 roots: vec!["src".into()],
2212 exact_files: vec![],
2213 contracts: vec![],
2214 };
2215
2216 let (projection, receipt) = ledger.project_relevant_decisions("worker", Some(&claim), &[]);
2217
2218 assert!(projection.len() <= COORDINATION_PROJECTION_BYTE_LIMIT);
2219 assert_eq!(
2220 receipt.decision_ids.len(),
2221 COORDINATION_PROJECTION_DECISION_LIMIT
2222 );
2223 assert_eq!(receipt.deduplicated, 0);
2224 assert_eq!(receipt.omitted, 2);
2225 }
2226
2227 #[test]
2228 fn coordination_schema_drift_fails_closed_before_mutation() {
2229 let mut ledger = CoordinationLedger {
2230 schema_version: COORDINATION_SCHEMA_VERSION + 2,
2231 ..CoordinationLedger::default()
2232 };
2233 let error = ledger
2234 .register_claim(
2235 WriteScopeClaim {
2236 owner: "worker".into(),
2237 roots: vec!["src".into()],
2238 exact_files: vec![],
2239 contracts: vec![],
2240 },
2241 false,
2242 |_| false,
2243 )
2244 .unwrap_err();
2245 assert!(error.contains("unsupported coordination schema"), "{error}");
2246 assert!(ledger.write_claims.is_empty());
2247 }
2248 }
2249
2249 lines RUST