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 {