From a039274c92027935d8892fffb2582f9d980d9673 Mon Sep 17 00:00:00 2001 From: Bryan Helmkamp Date: Sun, 20 Sep 2026 14:53:46 -0400 Subject: [PATCH] Group the hooks' bookkeeping into three ledgers with one loaded-once shape FabroHooks held twenty-one flat fields. Three of them were "a set read from the store once, then kept current", in two shapes: a Mutex beside a OnceCell<()> that had to be initialised first, and, for the restore plan, a OnceCell> that could not be read uninitialised. The recorded checkpoints and the collected artifacts now use the second shape too. The checkpoint state (committed, recorded, last, branch), the artifact state (globs, collected) and the scope state (environments, inherited workspaces, per-workspace locks) are three structs with their own methods, so each invariant lives in one type instead of in every caller. Co-Authored-By: Claude Fable 5.1 --- lib/components/fabro-petri/src/hooks.rs | 360 ++++++++++++++---------- 1 file changed, 213 insertions(+), 147 deletions(-) diff --git a/lib/components/fabro-petri/src/hooks.rs b/lib/components/fabro-petri/src/hooks.rs index 73dc9d6da..04443fc71 100644 --- a/lib/components/fabro-petri/src/hooks.rs +++ b/lib/components/fabro-petri/src/hooks.rs @@ -194,61 +194,141 @@ type AcquiredEnv = (String, Arc); /// The identity of a collected file: its path and content digest. type ArtifactIdentity = (String, String); -/// The last checkpoint recorded: its workspace and commit. -type LastCheckpoint = (String, String); +/// A checkpoint's workspace and commit. +type WorkspaceCommit = (String, String); -/// Fabro's `ExecutionHooks`, around the hooks the runtime installed. -pub struct FabroHooks { - inner: Arc, - run_id: RunId, - records: Arc, - /// Where an artifact's bytes and a diff's patch go; `None` records - /// summaries alone. - blobs: Option>, - workspaces: RunWorkspaces, - lookup: WorkspaceLookup, - identity: GitIdentity, - artifact_globs: Result, - host_workspaces: bool, - test_gates: Option, - handle: OnceLock, +/// The run's checkpoints as the hooks track them: what this process +/// committed, what has its platform record, which was recorded last, and +/// the run branch. +#[derive(Default)] +struct CheckpointLedger { /// The workspace and commit of every checkpoint this process made. - committed: Mutex>, - /// Which checkpoints have their platform record, loaded from the store - /// once and kept up to date with every append. - recorded: Mutex>, - recorded_loaded: OnceCell<()>, + committed: Mutex>, + /// Which checkpoints have their platform record: read from the store + /// once, then kept current with every append. + recorded: OnceCell>>, /// The workspace and commit of the checkpoint recorded last: the head /// the run's diff is measured to. - last_checkpoint: Mutex>, + last: Mutex>, /// The run branch as recorded, once: read from the store, or written /// by the commit that created the branch. - branch: OnceCell, - /// Every artifact collected so far, by path and digest, loaded from the - /// store once and kept up to date with every append. - collected: Mutex>, - collected_loaded: OnceCell<()>, - /// Inherited workspaces resolved through the run's records. - inherited: Mutex>>, - /// One lock per workspace: the branches of a parallel node and a nested - /// invocation share their caller's workspace, and Git allows one index - /// operation at a time in it. - workspace_locks: Mutex>>>, - /// The checkpoint failure that ended the run, when one did. - failure: Mutex>, + branch: OnceCell, +} + +impl CheckpointLedger { + fn remember_commit(&self, key: CheckpointKey, workspace: &str, sha: &str) { + sync::lock(&self.committed).insert(key, (workspace.to_string(), sha.to_string())); + } + + fn commit_of(&self, key: CheckpointKey) -> Option { + sync::lock(&self.committed).get(&key).cloned() + } + + fn last(&self) -> Option { + sync::lock(&self.last).clone() + } + + fn set_last(&self, last: WorkspaceCommit) { + *sync::lock(&self.last) = Some(last); + } + + /// The last checkpoint the store named, kept only when this process + /// has not recorded one yet. + fn set_last_if_unset(&self, last: Option) { + let mut current = sync::lock(&self.last); + if current.is_none() { + *current = last; + } + } +} + +/// The run's artifacts as the hooks track them: which files are collected, +/// and which were collected already. +struct ArtifactLedger { + /// The `[run.artifacts] include` patterns, or why they do not parse. + globs: Result, + /// Every artifact collected so far, by path and digest: read from the + /// store once, then kept current with every append. + collected: OnceCell>>, +} + +/// Where each acquired scope's workspace is, and the locks that serialize +/// Git work in it. +#[derive(Default)] +struct ScopeEnvs { /// The environment of every acquired scope, by execution and scope, /// with the workspace id the executor named: where `git` runs when the /// workspaces are not on this host, and where artifacts are read from /// on every provider. Dropped at release. - envs: Mutex>, + acquired: Mutex>, + /// Inherited workspaces resolved through the run's records. + inherited: Mutex>>, + /// One lock per workspace: the branches of a parallel node and a nested + /// invocation share their caller's workspace, and Git allows one index + /// operation at a time in it. + locks: Mutex>>>, +} + +impl ScopeEnvs { + fn insert(&self, execution: ExecutionId, scope: ScopeId, env: AcquiredEnv) { + sync::lock(&self.acquired).insert((execution, scope), env); + } + + fn remove(&self, execution: ExecutionId, scope: ScopeId) { + sync::lock(&self.acquired).remove(&(execution, scope)); + } + + fn get(&self, execution: ExecutionId, scope: ScopeId) -> Option { + sync::lock(&self.acquired).get(&(execution, scope)).cloned() + } + + /// The inherited workspace of an invocation, once resolved: `None` + /// when not resolved yet, `Some(None)` when it inherits none. + fn inherited(&self, invocation: InvocationId) -> Option> { + sync::lock(&self.inherited).get(&invocation).cloned() + } + + fn remember_inherited(&self, invocation: InvocationId, inherited: Option) { + sync::lock(&self.inherited).insert(invocation, inherited); + } + + /// The lock that serializes Git work in one workspace. + fn lock_for(&self, workspace: &str) -> Arc> { + Arc::clone( + sync::lock(&self.locks) + .entry(workspace.to_string()) + .or_default(), + ) + } +} + +/// Fabro's `ExecutionHooks`, around the hooks the runtime installed. +pub struct FabroHooks { + inner: Arc, + run_id: RunId, + records: Arc, + /// Where an artifact's bytes and a diff's patch go; `None` records + /// summaries alone. + blobs: Option>, + workspaces: RunWorkspaces, + lookup: WorkspaceLookup, + identity: GitIdentity, + host_workspaces: bool, + test_gates: Option, + handle: OnceLock, + checkpoints: CheckpointLedger, + artifacts: ArtifactLedger, + scopes: ScopeEnvs, + /// The checkpoint failure that ended the run, when one did. + failure: Mutex>, /// Whether the run continues from its records: a sandbox workspace is /// then brought to its snapshot when its scope is first acquired. - resumed: bool, + resumed: bool, /// The snapshot every live sandbox workspace must sit on before work /// resumes in it, read once from the records; an entry leaves when it /// is applied. - restore: OnceCell>>, - store: Arc, + restore: OnceCell>>, + store: Arc, } impl FabroHooks { @@ -284,21 +364,16 @@ impl FabroHooks { workspaces, lookup: WorkspaceLookup::new(Arc::clone(&store), run_key), identity, - artifact_globs: WorkspaceGlobSet::try_new(&spec.artifacts), host_workspaces: spec.host_workspaces, test_gates: spec.test_gates, handle: OnceLock::new(), - committed: Mutex::default(), - recorded: Mutex::default(), - recorded_loaded: OnceCell::new(), - last_checkpoint: Mutex::default(), - branch: OnceCell::new(), - collected: Mutex::default(), - collected_loaded: OnceCell::new(), - inherited: Mutex::default(), - workspace_locks: Mutex::default(), + checkpoints: CheckpointLedger::default(), + artifacts: ArtifactLedger { + globs: WorkspaceGlobSet::try_new(&spec.artifacts), + collected: OnceCell::new(), + }, + scopes: ScopeEnvs::default(), failure: Mutex::default(), - envs: Mutex::default(), resumed, restore: OnceCell::new(), store, @@ -326,15 +401,6 @@ impl FabroHooks { &self.workspaces } - /// The lock that serializes Git work in one workspace. - fn workspace_lock(&self, workspace: &str) -> Arc> { - Arc::clone( - sync::lock(&self.workspace_locks) - .entry(workspace.to_string()) - .or_default(), - ) - } - fn fail_run(&self, message: &str) { let mut failure = sync::lock(&self.failure); if failure.is_none() { @@ -360,10 +426,7 @@ impl FabroHooks { if self.workspaces.workspace_exists(&isolated).await { return Ok(isolated); } - let cached = sync::lock(&self.inherited) - .get(&context.invocation) - .cloned(); - let inherited = if let Some(inherited) = cached { + let inherited = if let Some(inherited) = self.scopes.inherited(context.invocation) { inherited } else { let inherited = self @@ -377,7 +440,8 @@ impl FabroHooks { collect_chain(&error).join(": ") ) })?; - sync::lock(&self.inherited).insert(context.invocation, inherited.clone()); + self.scopes + .remember_inherited(context.invocation, inherited.clone()); inherited }; Ok(inherited.unwrap_or(isolated)) @@ -386,9 +450,7 @@ impl FabroHooks { /// The environment of `scope` in the context's execution, as /// `scope_acquired` kept it, with the workspace id the executor named. fn env_of(&self, context: &HookContext, scope: ScopeId) -> Option { - sync::lock(&self.envs) - .get(&(context.execution, scope)) - .cloned() + self.scopes.get(context.execution, scope) } /// Where `scope`'s workspace is and what to call it: on this host, the @@ -445,7 +507,7 @@ impl FabroHooks { )); }; self.gate("commit", node).await; - let serialized = self.workspace_lock(&workspace); + let serialized = self.scopes.lock_for(&workspace); let _held = serialized.lock().await; match self .workspaces @@ -486,7 +548,8 @@ impl FabroHooks { /// Remember a commit this process made, and record the run branch when /// this commit created it. async fn committed(&self, key: CheckpointKey, workspace: &str, snapshot: &Snapshot) { - sync::lock(&self.committed).insert(key, (workspace.to_string(), snapshot.sha.clone())); + self.checkpoints + .remember_commit(key, workspace, &snapshot.sha); let Some(branched) = &snapshot.branched else { return; }; @@ -520,6 +583,7 @@ impl FabroHooks { firing: key.firing, }; let branch = self + .checkpoints .branch .get_or_try_init(|| async { if let Some(stored) = self.stored_branch().await? { @@ -626,7 +690,7 @@ impl FabroHooks { let Some(target) = target else { return Ok(()); }; - let serialized = self.workspace_lock(workspace); + let serialized = self.scopes.lock_for(workspace); let _held = serialized.lock().await; let action = recovery::bring_to(&self.workspaces, site, workspace, &target) .await @@ -655,14 +719,11 @@ impl FabroHooks { scope: ScopeId, key: CheckpointKey, ) -> Result<(), String> { - self.recorded_loaded - .get_or_try_init(|| self.load_recorded()) - .await?; - if sync::lock(&self.recorded).contains(&key) { + let recorded = self.recorded_checkpoints().await?; + if sync::lock(recorded).contains(&key) { return Ok(()); } - let committed = sync::lock(&self.committed).get(&key).cloned(); - let (workspace, sha) = if let Some(committed) = committed { + let (workspace, sha) = if let Some(committed) = self.checkpoints.commit_of(key) { committed } else { let acquired = self.env_of(context, scope).map(|(workspace, _)| workspace); @@ -670,7 +731,7 @@ impl FabroHooks { Some(workspace) => workspace, None => self.workspace_of(context, scope).await?, }; - let serialized = self.workspace_lock(&workspace); + let serialized = self.scopes.lock_for(&workspace); let held = serialized.lock().await; let found = self.workspaces.find(&workspace, key).await; drop(held); @@ -729,8 +790,8 @@ impl FabroHooks { collect_chain(&error).join(": ") ) })?; - sync::lock(&self.recorded).insert(key); - *sync::lock(&self.last_checkpoint) = Some((workspace, sha)); + sync::lock(recorded).insert(key); + self.checkpoints.set_last((workspace, sha)); Ok(()) } @@ -780,42 +841,43 @@ impl FabroHooks { /// resume's reissued routing decisions must not record again, and /// where the run's diff is measured to when this process made no /// checkpoint yet. - async fn load_recorded(&self) -> Result<(), String> { - let stored = self - .records - .read_kind(&self.run_id, PlatformRecordKind::Checkpoint) + async fn recorded_checkpoints(&self) -> Result<&Mutex>, String> { + self.checkpoints + .recorded + .get_or_try_init(|| async { + let stored = self + .records + .read_kind(&self.run_id, PlatformRecordKind::Checkpoint) + .await + .map_err(|error| { + format!( + "the run's checkpoint records could not be read: {}", + collect_chain(&error).join(": ") + ) + })?; + let mut recorded = HashSet::new(); + let mut last = None; + for record in stored { + let PlatformRecord::Checkpoint(checkpoint) = &record.record else { + continue; + }; + if let Some(key) = checkpoint + .operation + .as_ref() + .and_then(CheckpointKey::from_operation) + { + recorded.insert(key); + } + if let (Some(workspace), Some(sha)) = + (&checkpoint.workspace, &checkpoint.git_commit_sha) + { + last = Some((workspace.clone(), sha.clone())); + } + } + self.checkpoints.set_last_if_unset(last); + Ok(Mutex::new(recorded)) + }) .await - .map_err(|error| { - format!( - "the run's checkpoint records could not be read: {}", - collect_chain(&error).join(": ") - ) - })?; - let mut recorded = sync::lock(&self.recorded); - let mut last = None; - for record in stored { - let PlatformRecord::Checkpoint(checkpoint) = &record.record else { - continue; - }; - if let Some(key) = checkpoint - .operation - .as_ref() - .and_then(CheckpointKey::from_operation) - { - recorded.insert(key); - } - if let (Some(workspace), Some(sha)) = - (&checkpoint.workspace, &checkpoint.git_commit_sha) - { - last = Some((workspace.clone(), sha.clone())); - } - } - drop(recorded); - let mut last_checkpoint = sync::lock(&self.last_checkpoint); - if last_checkpoint.is_none() { - *last_checkpoint = last; - } - Ok(()) } /// The artifacts of a finished attempt: every file of its workspace @@ -828,7 +890,7 @@ impl FabroHooks { scope: ScopeId, key: CheckpointKey, ) -> Result { - let globs = match &self.artifact_globs { + let globs = match &self.artifacts.globs { Ok(globs) => globs, Err(error) => return Err(format!("invalid run.artifacts.include pattern: {error}")), }; @@ -843,9 +905,7 @@ impl FabroHooks { let Some(blobs) = &self.blobs else { return Err("the run has no blob table to collect artifacts into".to_string()); }; - self.collected_loaded - .get_or_try_init(|| self.load_collected()) - .await?; + let already = self.collected_artifacts().await?; let candidates = list_artifacts(env.as_ref(), globs).await?; let limit = usize::try_from(ARTIFACT_MAX_FILE_BYTES).unwrap_or(usize::MAX); let mut collected = 0; @@ -864,7 +924,7 @@ impl FabroHooks { }; let digest = BlobHash::new(&bytes); let identity = (path.clone(), digest.to_string()); - if sync::lock(&self.collected).contains(&identity) { + if sync::lock(already).contains(&identity) { continue; } let blob = blobs @@ -897,7 +957,7 @@ impl FabroHooks { collect_chain(&error).join(": ") ) })?; - sync::lock(&self.collected).insert(identity); + sync::lock(already).insert(identity); total_bytes = total_bytes.saturating_add(size); collected += 1; } @@ -906,24 +966,32 @@ impl FabroHooks { /// The artifacts already collected for the run, read once: a file that /// is unchanged since it was collected is not collected again. - async fn load_collected(&self) -> Result<(), String> { - let stored = self - .records - .read_kind(&self.run_id, PlatformRecordKind::ArtifactCollected) + async fn collected_artifacts(&self) -> Result<&Mutex>, String> { + self.artifacts + .collected + .get_or_try_init(|| async { + let stored = self + .records + .read_kind(&self.run_id, PlatformRecordKind::ArtifactCollected) + .await + .map_err(|error| { + format!( + "the run's artifact records could not be read: {}", + collect_chain(&error).join(": ") + ) + })?; + let collected = stored + .into_iter() + .filter_map(|record| match record.record { + PlatformRecord::ArtifactCollected(artifact) => { + Some((artifact.path, artifact.digest)) + } + _ => None, + }) + .collect(); + Ok(Mutex::new(collected)) + }) .await - .map_err(|error| { - format!( - "the run's artifact records could not be read: {}", - collect_chain(&error).join(": ") - ) - })?; - let mut collected = sync::lock(&self.collected); - for record in stored { - if let PlatformRecord::ArtifactCollected(artifact) = record.record { - collected.insert((artifact.path, artifact.digest)); - } - } - Ok(()) } /// The run's diff: the run branch's last checkpoint against the base @@ -931,10 +999,8 @@ impl FabroHooks { /// Nothing is recorded for a run that never created its branch or /// never checkpointed. async fn record_run_diff(&self) -> Result<(), String> { - self.recorded_loaded - .get_or_try_init(|| self.load_recorded()) - .await?; - let branch = match self.branch.get() { + self.recorded_checkpoints().await?; + let branch = match self.checkpoints.branch.get() { Some(branch) => Some(branch.clone()), None => self.stored_branch().await?, }; @@ -945,8 +1011,7 @@ impl FabroHooks { let Some(base_sha) = branch.base_sha.clone() else { return Ok(()); }; - let last = sync::lock(&self.last_checkpoint).clone(); - let Some((workspace, head_sha)) = last else { + let Some((workspace, head_sha)) = self.checkpoints.last() else { debug!(run_id = %self.run_id, "no checkpoint is recorded; no run diff"); return Ok(()); }; @@ -1216,7 +1281,7 @@ impl ExecutionHooks for FabroHooks { ); let scope = released.scope; let notes = self.inner.scope_released(context, released).await; - sync::lock(&self.envs).remove(&(context.execution, scope)); + self.scopes.remove(context.execution, scope); notes } @@ -1227,8 +1292,9 @@ impl ExecutionHooks for FabroHooks { ) -> Result<(), ScopeAcquiredError> { self.inner.scope_acquired(context, acquired.clone()).await?; let workspace = acquired.workspace.as_str().to_owned(); - sync::lock(&self.envs).insert( - (context.execution, acquired.scope), + self.scopes.insert( + context.execution, + acquired.scope, (workspace.clone(), Arc::clone(&acquired.env)), ); if !self.resumed {