Carry required publication into fork declarations

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
This commit is contained in:
Scott Werner 2026-10-07 15:53:08 -04:00
parent 59f9088b04
commit 3443a923c0
5 changed files with 51 additions and 3 deletions

View file

@ -42,6 +42,8 @@ historical-view repair requires a separate, bounded process over preserved
source records. No publication is invoked by projection or replay.
A committed resume returns the same overall result without publishing again.
A fork inherits the source run’s required-finalization declaration while seeding
its records; its worker restores the actual publication hooks before resume.
Unfinished recovery must restore the same finalization requirement. Petri may
call the restored finalizer again after interruption, including a crash after
the callback returns but before the terminal record commits. Existing GitHub

View file

@ -1103,12 +1103,13 @@ fn attach_json_errors_without_prompting_for_human_input() {
"recorded_at": "[EPOCH_MS]",
"body": {
"event": "run.started",
"format_version": 8,
"format_version": 9,
"key": "[ULID]",
"root": 0,
"middleware_chain": [
"circuit-breaker"
]
],
"required_finalization": false
}
}
}

View file

@ -263,6 +263,12 @@ fn bare_fabro_with_unbound_inputs_in_template_partial_validates_structurally_wit
// The include error names the partial relative to the run's working
// directory, so the `../` run is as long as that directory is deep.
let mut filters = context.filters();
// A checkout outside the user home can be rendered through a relative
// path (including macOS's /tmp alias), rather than its canonical root.
filters.push((
r#"(?:\.\./)+[^"\s]*?/test/(templated_unbound_partial/)"#.to_string(),
"[UP][FIXTURES]/$1".to_string(),
));
filters.push((
r"(\.\./)*\.\.\[FIXTURES\]".to_string(),
"[UP][FIXTURES]".to_string(),
@ -333,6 +339,10 @@ fn validate_reports_missing_template_dependency() {
let mut cmd = context.validate();
cmd.arg(fixture("templates/missing_dependency/workflow.fabro"));
let mut filters = context.filters();
filters.push((
r"(?:\.\./)+[^`\s]*?/test/(templates/missing_dependency/)".to_string(),
"[FIXTURES]/$1".to_string(),
));
filters.push((
r"(?:\.\./)*\.\.\[FIXTURES\]/".to_string(),
"[FIXTURES]/".to_string(),

View file

@ -43,7 +43,8 @@ use petri_execution::{
StoreError as CoordinatorStoreError,
};
use petri_runtime::RunOptions;
use petri_runtime::ir::FiringId;
use petri_runtime::driver::lifecycle::{ExecutionHooks, HookContext, RunFinished};
use petri_runtime::ir::{FinalizationFailure, FiringId};
use petri_store::StoreError;
use tracing::{debug, info};
@ -169,6 +170,31 @@ pub async fn check(
Ok(())
}
/// Declaration-only hooks used while copying a fork's records. The fork
/// inherits its source's finalization requirement; its worker installs the
/// actual publisher before resuming. Never execute with these hooks.
struct ForkFinalizationRequirement {
required: bool,
}
#[async_trait::async_trait]
impl ExecutionHooks for ForkFinalizationRequirement {
fn requires_run_finalization(&self) -> bool {
self.required
}
async fn finalize_run(
&self,
_context: &HookContext,
_finished: RunFinished,
) -> Result<(), FinalizationFailure> {
Err(FinalizationFailure::new(
"finalization_unavailable",
"the fork must install its worker publication hooks before execution",
))
}
}
/// Seed the fork: Petri's records, the kept checkpoints and the run branch. The
/// new run must not exist in the store yet.
pub async fn fork(request: ForkRequest) -> Result<Forked, ForkError> {
@ -181,12 +207,17 @@ pub async fn fork(request: ForkRequest) -> Result<Forked, ForkError> {
.await
.map_err(ForkError::Open)?;
let required = host::stored_state(&*source_logs)
.await
.map_err(ForkError::Seed)?
.required_finalization;
let mut options = RunOptions::new(&request.fork_run_dir);
options.run_key = Some(fork_key.clone());
// A fork only copies records and acquires no sandbox, so it needs no
// provider configuration.
let runtime = providers::standard_runtime(&SandboxProviderConfig::default())
.options(options)
.hooks(Arc::new(ForkFinalizationRequirement { required }))
.store(Arc::clone(&request.store));
let forked = host::fork_from(&runtime, &*source_logs, request.position, ForkOptions {
rerun_last: request.rerun_last,

View file

@ -1459,6 +1459,10 @@ async fn assert_shallow_fork(checkpoint_index: usize) {
})
.await
.expect("the fork is seeded");
assert!(
forked.inspection().await.required_finalization,
"a publishing fork inherits its required finalization declaration"
);
assert_eq!(
seeded.start.expect("the selected checkpoint").sha,
*checkpoint_sha