mirror of
https://github.com/fabro-sh/fabro.git
synced 2026-10-02 02:13:49 +00:00
Prove the checkpoint hooks and the recovery protocol
In-process tests over the memory store: every finish is committed on the run branch with its identity trailers and recorded with its commit, the run-end hooks reach Petri's local service through Fabro's wrapper, a stage that fails on its own terms is committed and its failure route runs on the committed files, a failed checkpoint records `checkpoint_failed` with no route taken and a restart reports the run failed, and a `[[run.hooks]]` hook blocks an agent's tool call through the forwarded service, with the model told why. Real-binary scenarios crash the server and its worker with SIGKILL: after a durable finish the stage's commit is not repeated and the interrupted stage reruns on its snapshot; a crash held before the commit reruns the stage once; a crash held after the commit but before its record reconciles the record from the snapshot repository; a deleted workspace is restored; a failure route sees the same committed files after a crash; a failed checkpoint fails the run and a restart leaves it failed. Recovery selects the executions `inspect_run` reports incomplete, and a run whose coordinator log is still empty is left to the worker's resume. The worker's platform record endpoints get an API test and the generated TypeScript client. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
This commit is contained in:
parent
01beea0a6c
commit
34139b1936
12 changed files with 1032 additions and 14 deletions
1
Cargo.lock
generated
1
Cargo.lock
generated
|
|
@ -2900,6 +2900,7 @@ dependencies = [
|
|||
"fabro-llm",
|
||||
"fabro-petri",
|
||||
"fabro-store",
|
||||
"fabro-test",
|
||||
"fabro-types",
|
||||
"fabro-util",
|
||||
"lithos-llm",
|
||||
|
|
|
|||
|
|
@ -29,12 +29,13 @@ use std::time::{Duration, Instant};
|
|||
use fabro_client::ServerTarget;
|
||||
use fabro_config::{Storage, envfile};
|
||||
use fabro_petri::SqliteRunStore;
|
||||
use fabro_petri::checkpoint::CheckpointKey;
|
||||
use fabro_petri::engine::{self, RunStatus};
|
||||
use fabro_petri::petri::RunKey;
|
||||
use fabro_static::EnvVars;
|
||||
use fabro_store::EventEnvelope;
|
||||
use fabro_store::{EventEnvelope, PlatformRecord, PlatformRecordKind, PlatformRecordStore};
|
||||
use fabro_test::{apply_test_isolation, expect_reqwest_json, isolated_storage_dir, test_context};
|
||||
use fabro_types::EventBody;
|
||||
use fabro_types::{EventBody, RunId};
|
||||
|
||||
use crate::cmd::support::created_run_id;
|
||||
use crate::support::{
|
||||
|
|
@ -81,6 +82,8 @@ struct RunningServer {
|
|||
config_path: PathBuf,
|
||||
port: u16,
|
||||
api_base_url: String,
|
||||
/// The checkpoint gate directory the server forwards to its workers.
|
||||
gates_dir: PathBuf,
|
||||
}
|
||||
|
||||
impl RunningServer {
|
||||
|
|
@ -103,6 +106,8 @@ impl RunningServer {
|
|||
.expect("the server env writes");
|
||||
fabro_util::dev_token::write_dev_token(&runtime_directory.dev_token_path(), TEST_DEV_TOKEN)
|
||||
.expect("the dev token writes");
|
||||
let gates_dir = home_root.path().join("checkpoint-gates");
|
||||
std::fs::create_dir_all(&gates_dir).expect("the gates dir creates");
|
||||
let mut server = Self {
|
||||
child: None,
|
||||
home_root,
|
||||
|
|
@ -111,6 +116,7 @@ impl RunningServer {
|
|||
config_path,
|
||||
port,
|
||||
api_base_url: format!("http://127.0.0.1:{port}"),
|
||||
gates_dir,
|
||||
};
|
||||
server.launch().await;
|
||||
server
|
||||
|
|
@ -129,6 +135,7 @@ impl RunningServer {
|
|||
EnvVars::FABRO_HOME,
|
||||
self.home_root.path().join("fabro-home"),
|
||||
);
|
||||
cmd.env(EnvVars::FABRO_TEST_CHECKPOINT_GATES, &self.gates_dir);
|
||||
cmd.args(["server", "start", "--foreground"])
|
||||
.arg("--storage-dir")
|
||||
.arg(&self.storage_dir)
|
||||
|
|
@ -137,16 +144,16 @@ impl RunningServer {
|
|||
.arg("--config")
|
||||
.arg(&self.config_path)
|
||||
.stdin(Stdio::null())
|
||||
.stdout(Stdio::null())
|
||||
.stdout(self.stderr_log())
|
||||
.stderr(self.stderr_log());
|
||||
let mut child = cmd.spawn().expect("the server spawns");
|
||||
wait_for_http_ready(&self.api_base_url, &mut child).await;
|
||||
self.child = Some(child);
|
||||
}
|
||||
|
||||
/// Where the server's stderr goes: a file beside its storage, so a
|
||||
/// chatty server never blocks on a pipe nobody reads, and a failing
|
||||
/// test can show it.
|
||||
/// Where the server's stdout and stderr go: a file beside its storage,
|
||||
/// so a chatty server never blocks on a pipe nobody reads, and a
|
||||
/// failing test can show its log.
|
||||
fn stderr_log(&self) -> Stdio {
|
||||
let path = self.storage_dir.with_file_name("server.stderr.log");
|
||||
let file = std::fs::OpenOptions::new()
|
||||
|
|
@ -209,6 +216,117 @@ impl RunningServer {
|
|||
.expect("the server database opens");
|
||||
SqliteRunStore::new(database.clone_pool())
|
||||
}
|
||||
|
||||
/// The run's platform records in the server's database.
|
||||
async fn platform_records(&self) -> PlatformRecordStore {
|
||||
let database = fabro_db::Database::connect(Storage::new(&self.storage_dir).sqlite_path())
|
||||
.await
|
||||
.expect("the server database opens");
|
||||
PlatformRecordStore::new(database.clone_pool())
|
||||
}
|
||||
|
||||
/// Where the run's worker ran Petri: the run's scratch under the
|
||||
/// server's storage.
|
||||
fn petri_run_dir(&self, run_id: &str) -> PathBuf {
|
||||
let run_id: RunId = run_id.parse().expect("the run id parses");
|
||||
Storage::new(&self.storage_dir)
|
||||
.run_scratch(&run_id)
|
||||
.root()
|
||||
.join("petri")
|
||||
}
|
||||
|
||||
/// The worker's own log for the run.
|
||||
fn worker_log(&self, run_id: &str) -> PathBuf {
|
||||
let run_id: RunId = run_id.parse().expect("the run id parses");
|
||||
Storage::new(&self.storage_dir)
|
||||
.run_scratch(&run_id)
|
||||
.root()
|
||||
.join("runtime")
|
||||
.join("server.log")
|
||||
}
|
||||
|
||||
/// Hold the worker's checkpoint at `point` (`commit` or `record`) for
|
||||
/// `node` until [`release`](Self::release).
|
||||
fn hold(&self, point: &str, node: &str) {
|
||||
std::fs::write(self.gates_dir.join(format!("{point}.{node}.hold")), "")
|
||||
.expect("the hold file writes");
|
||||
}
|
||||
|
||||
fn release(&self, point: &str, node: &str) {
|
||||
std::fs::write(self.gates_dir.join(format!("{point}.{node}.release")), "")
|
||||
.expect("the release file writes");
|
||||
}
|
||||
|
||||
/// Wait until the worker's log says its checkpoint is held at a gate.
|
||||
fn wait_until_held(&self, run_id: &str, point: &str, node: &str) {
|
||||
let log = self.worker_log(run_id);
|
||||
let needle = format!("checkpoint held at a test gate point=\"{point}\" node=\"{node}\"");
|
||||
let deadline = Instant::now() + RUN_TIMEOUT;
|
||||
loop {
|
||||
let text = std::fs::read_to_string(&log).unwrap_or_default();
|
||||
if text.contains(&needle) {
|
||||
return;
|
||||
}
|
||||
assert!(
|
||||
Instant::now() < deadline,
|
||||
"the worker never held at {point}.{node}; log:\n{text}"
|
||||
);
|
||||
std::thread::sleep(POLL);
|
||||
}
|
||||
}
|
||||
|
||||
/// The one host workspace of the run, and the commits on its run
|
||||
/// branch, oldest first, as `(subject, key)`.
|
||||
fn workspace_commits(&self, run_id: &str) -> (PathBuf, Vec<(String, Option<CheckpointKey>)>) {
|
||||
let scopes = self.petri_run_dir(run_id).join("scopes");
|
||||
let mut workspaces: Vec<PathBuf> = std::fs::read_dir(&scopes)
|
||||
.expect("the scopes directory lists")
|
||||
.map(|entry| entry.expect("an entry reads").path().join("work"))
|
||||
.collect();
|
||||
assert_eq!(workspaces.len(), 1, "one workspace: {workspaces:?}");
|
||||
let workspace = workspaces.remove(0);
|
||||
let output = Command::new("git")
|
||||
.args(["log", "--reverse", "--format=%s%x00%B%x1e"])
|
||||
.current_dir(&workspace)
|
||||
.output()
|
||||
.expect("git runs");
|
||||
assert!(
|
||||
output.status.success(),
|
||||
"git log failed: {}",
|
||||
String::from_utf8_lossy(&output.stderr)
|
||||
);
|
||||
let log = String::from_utf8_lossy(&output.stdout).into_owned();
|
||||
let commits = log
|
||||
.split('\u{1e}')
|
||||
.filter(|entry| !entry.trim().is_empty())
|
||||
.map(|entry| {
|
||||
let mut parts = entry.trim_start().splitn(2, '\0');
|
||||
let subject = parts.next().unwrap_or_default().to_string();
|
||||
let body = parts.next().unwrap_or_default();
|
||||
(subject, CheckpointKey::from_message(body))
|
||||
})
|
||||
.collect();
|
||||
(workspace, commits)
|
||||
}
|
||||
|
||||
/// The run's checkpoint records, in seq order, as `(node position, sha)`.
|
||||
async fn checkpoints(&self, run_id: &str) -> Vec<(CheckpointKey, String)> {
|
||||
let run_id: RunId = run_id.parse().expect("the run id parses");
|
||||
self.platform_records()
|
||||
.await
|
||||
.read_kind(&run_id, PlatformRecordKind::Checkpoint)
|
||||
.await
|
||||
.expect("the checkpoint records read")
|
||||
.into_iter()
|
||||
.filter_map(|stored| match stored.record {
|
||||
PlatformRecord::Checkpoint(record) => Some((
|
||||
CheckpointKey::from_operation(record.operation.as_ref()?)?,
|
||||
record.git_commit_sha?,
|
||||
)),
|
||||
_ => None,
|
||||
})
|
||||
.collect()
|
||||
}
|
||||
}
|
||||
|
||||
impl Drop for RunningServer {
|
||||
|
|
@ -251,17 +369,22 @@ async fn wait_for_http_ready(base_url: &str, child: &mut Child) {
|
|||
/// A workspace holding a command-only bundle whose `workflow.toml` names
|
||||
/// Petri, with the given stage script.
|
||||
fn write_petri_workspace(context: &fabro_test::TestContext, script: &str) -> PathBuf {
|
||||
let workspace = context.temp_dir.join("petri-workspace");
|
||||
std::fs::create_dir_all(&workspace).expect("the workspace creates");
|
||||
std::fs::write(
|
||||
workspace.join("workflow.fabro"),
|
||||
format!(
|
||||
write_petri_bundle(
|
||||
context,
|
||||
&format!(
|
||||
"digraph Command {{\n graph [goal=\"Run one command\", default_max_retries=0]\n start \
|
||||
[shape=Mdiamond]\n exit [shape=Msquare]\n say [shape=parallelogram, \
|
||||
script=\"{script}\", max_retries=0]\n start -> say -> exit\n}}\n"
|
||||
),
|
||||
)
|
||||
.expect("the workflow writes");
|
||||
}
|
||||
|
||||
/// A workspace holding a command-only bundle of the given graph, whose
|
||||
/// `workflow.toml` names Petri.
|
||||
fn write_petri_bundle(context: &fabro_test::TestContext, workflow: &str) -> PathBuf {
|
||||
let workspace = context.temp_dir.join("petri-workspace");
|
||||
std::fs::create_dir_all(&workspace).expect("the workspace creates");
|
||||
std::fs::write(workspace.join("workflow.fabro"), workflow).expect("the workflow writes");
|
||||
std::fs::write(
|
||||
workspace.join("workflow.toml"),
|
||||
"_version = 1\n\n[workflow]\ngraph = \"workflow.fabro\"\nengine = \"petri\"\n\n[run]\ngoal \
|
||||
|
|
@ -566,3 +689,375 @@ async fn run_status_offline(server: &RunningServer) -> Option<String> {
|
|||
.ok()
|
||||
.map(|response| response.status().to_string())
|
||||
}
|
||||
|
||||
/// A shell loop that waits for `gate` to exist.
|
||||
fn wait_for(gate: &Path) -> String {
|
||||
format!("while [ ! -f {} ]; do sleep 0.05; done", gate.display())
|
||||
}
|
||||
|
||||
/// Three command stages: `one` writes a file, `two` writes another after
|
||||
/// the gate opens (and logs each run), `three` checks both files.
|
||||
fn three_stage_bundle(context: &fabro_test::TestContext, gate: &Path) -> PathBuf {
|
||||
write_petri_bundle(
|
||||
context,
|
||||
&format!(
|
||||
"digraph Stages {{\n graph [goal=\"Three stages\", default_max_retries=0]\n start \
|
||||
[shape=Mdiamond]\n exit [shape=Msquare]\n one [shape=parallelogram, script=\"echo \
|
||||
one > one.txt\"]\n two [shape=parallelogram, script=\"echo run >> two.log; {}; \
|
||||
echo two > two.txt\"]\n three [shape=parallelogram, script=\"test \\\"$(cat \
|
||||
one.txt)\\\" = one && test \\\"$(cat two.txt)\\\" = two && cp two.log \
|
||||
three.log\"]\n start -> one -> two -> three -> exit\n}}\n",
|
||||
wait_for(gate)
|
||||
),
|
||||
)
|
||||
}
|
||||
|
||||
/// Kill the server first, so it never observes the worker exit, then the
|
||||
/// worker's whole process group, then the stage's own process group when a
|
||||
/// stage was waiting on `gate`: a stage process runs in a group of its own
|
||||
/// under the host plugin, and a machine crash takes it with everything
|
||||
/// else, where a killed worker alone would leave it writing into the
|
||||
/// workspace.
|
||||
fn crash(server: &mut RunningServer, worker: u32, gate: Option<&Path>) {
|
||||
server.kill();
|
||||
fabro_proc::sigkill_process_group(worker);
|
||||
let deadline = Instant::now() + Duration::from_secs(10);
|
||||
while fabro_proc::process_running(worker) {
|
||||
assert!(Instant::now() < deadline, "the worker did not die");
|
||||
std::thread::sleep(POLL);
|
||||
}
|
||||
let Some(gate) = gate else {
|
||||
return;
|
||||
};
|
||||
let output = Command::new("pgrep")
|
||||
.args(["-f", &gate.display().to_string()])
|
||||
.output()
|
||||
.expect("pgrep runs");
|
||||
for pid in String::from_utf8_lossy(&output.stdout)
|
||||
.lines()
|
||||
.filter_map(|line| line.trim().parse::<u32>().ok())
|
||||
{
|
||||
fabro_proc::sigkill_process_group(pid);
|
||||
}
|
||||
let deadline = Instant::now() + Duration::from_secs(10);
|
||||
while gate_is_polled(gate) {
|
||||
assert!(Instant::now() < deadline, "the stage did not die");
|
||||
std::thread::sleep(POLL);
|
||||
}
|
||||
}
|
||||
|
||||
/// Wait for the run to succeed after a restart, with the server's stderr
|
||||
/// on failure.
|
||||
async fn wait_for_success(server: &RunningServer, run_id: &str) {
|
||||
let status = wait_for_status(server, run_id, &["succeeded", "failed"]).await;
|
||||
assert_eq!(
|
||||
status,
|
||||
"succeeded",
|
||||
"run: {}\nserver stderr:\n{}",
|
||||
run_json(server, &format!("runs/{run_id}")).await,
|
||||
server.stderr_text()
|
||||
);
|
||||
}
|
||||
|
||||
/// The subjects of the commits on the run branch.
|
||||
fn subjects(commits: &[(String, Option<CheckpointKey>)]) -> Vec<&str> {
|
||||
commits
|
||||
.iter()
|
||||
.map(|(subject, _)| subject.as_str())
|
||||
.collect()
|
||||
}
|
||||
|
||||
/// The commit subjects one run of the three-stage bundle produces.
|
||||
fn three_stage_subjects(run_id: &str) -> Vec<String> {
|
||||
["start", "one", "two", "three", "exit"]
|
||||
.iter()
|
||||
.map(|node| format!("fabro({run_id}): {node} (success)"))
|
||||
.collect()
|
||||
}
|
||||
|
||||
/// A worker killed after a stage's finish is durable: on the restart the
|
||||
/// stage's commit is not repeated, the stage in flight reruns on the
|
||||
/// snapshot (its partial output gone), and the next stage sees both.
|
||||
#[tokio::test(flavor = "multi_thread")]
|
||||
async fn a_crash_after_a_durable_finish_keeps_its_one_commit() {
|
||||
if host_plugin().is_none() {
|
||||
return;
|
||||
}
|
||||
let context = test_context!();
|
||||
let mut server = RunningServer::start().await;
|
||||
let gate = context.temp_dir.join("two.gate");
|
||||
let workspace = three_stage_bundle(&context, &gate);
|
||||
let run_id = run_detached(&context, &server, &workspace);
|
||||
|
||||
wait_for_status(&server, &run_id, &["running"]).await;
|
||||
let worker = wait_for_worker(&run_id);
|
||||
wait_until_gate_is_polled(&gate);
|
||||
crash(&mut server, worker, Some(&gate));
|
||||
|
||||
server.launch().await;
|
||||
let resumed = wait_for_worker(&run_id);
|
||||
assert_ne!(resumed, worker);
|
||||
wait_until_gate_is_polled(&gate);
|
||||
std::fs::write(&gate, "go").expect("the gate opens");
|
||||
wait_for_success(&server, &run_id).await;
|
||||
|
||||
let (path, commits) = server.workspace_commits(&run_id);
|
||||
assert_eq!(subjects(&commits), three_stage_subjects(&run_id));
|
||||
assert_eq!(
|
||||
std::fs::read_to_string(path.join("three.log")).expect("three copied the log"),
|
||||
"run\n",
|
||||
"the crashed attempt's partial output was reset before the rerun; server log:\n{}",
|
||||
server.stderr_text()
|
||||
);
|
||||
let checkpoints = server.checkpoints(&run_id).await;
|
||||
assert_eq!(checkpoints.len(), 5, "{checkpoints:?}");
|
||||
let keys: Vec<Option<CheckpointKey>> = checkpoints.iter().map(|(key, _)| Some(*key)).collect();
|
||||
let committed: Vec<Option<CheckpointKey>> = commits.iter().map(|(_, key)| *key).collect();
|
||||
assert_eq!(keys, committed);
|
||||
server.shutdown();
|
||||
}
|
||||
|
||||
/// A worker killed in `prepare_result` before the commit lands: the finish
|
||||
/// is not durable, the stage reruns once, and one commit exists for it.
|
||||
#[tokio::test(flavor = "multi_thread")]
|
||||
async fn a_crash_before_the_commit_lands_reruns_the_stage_once() {
|
||||
if host_plugin().is_none() {
|
||||
return;
|
||||
}
|
||||
let context = test_context!();
|
||||
let mut server = RunningServer::start().await;
|
||||
let gate = context.temp_dir.join("two.gate");
|
||||
std::fs::write(&gate, "open").expect("the script gate is open from the start");
|
||||
let workspace = three_stage_bundle(&context, &gate);
|
||||
server.hold("commit", "two");
|
||||
let run_id = run_detached(&context, &server, &workspace);
|
||||
|
||||
wait_for_status(&server, &run_id, &["running"]).await;
|
||||
let worker = wait_for_worker(&run_id);
|
||||
server.wait_until_held(&run_id, "commit", "two");
|
||||
crash(&mut server, worker, None);
|
||||
|
||||
server.release("commit", "two");
|
||||
server.launch().await;
|
||||
wait_for_success(&server, &run_id).await;
|
||||
|
||||
let (path, commits) = server.workspace_commits(&run_id);
|
||||
assert_eq!(subjects(&commits), three_stage_subjects(&run_id));
|
||||
assert_eq!(
|
||||
std::fs::read_to_string(path.join("three.log")).expect("three copied the log"),
|
||||
"run\n",
|
||||
"the stage reran once, on the snapshot before it"
|
||||
);
|
||||
let store = server.petri_store().await;
|
||||
let outcome = engine::outcome_of(&store, &run_id)
|
||||
.await
|
||||
.expect("the run's Petri record inspects");
|
||||
assert_eq!(outcome.status, RunStatus::Success, "{outcome:?}");
|
||||
assert!(outcome.complete, "{:?}", outcome.incomplete);
|
||||
server.shutdown();
|
||||
}
|
||||
|
||||
/// A worker killed after the commit and its durable finish but before the
|
||||
/// platform record: the restart reconciles the record from the snapshot
|
||||
/// repository, the stage does not rerun, and one commit exists for it.
|
||||
#[tokio::test(flavor = "multi_thread")]
|
||||
async fn a_crash_before_the_record_reconciles_it_from_the_run_branch() {
|
||||
if host_plugin().is_none() {
|
||||
return;
|
||||
}
|
||||
let context = test_context!();
|
||||
let mut server = RunningServer::start().await;
|
||||
let gate = context.temp_dir.join("two.gate");
|
||||
std::fs::write(&gate, "open").expect("the script gate is open from the start");
|
||||
let workspace = three_stage_bundle(&context, &gate);
|
||||
server.hold("record", "two");
|
||||
let run_id = run_detached(&context, &server, &workspace);
|
||||
|
||||
wait_for_status(&server, &run_id, &["running"]).await;
|
||||
let worker = wait_for_worker(&run_id);
|
||||
server.wait_until_held(&run_id, "record", "two");
|
||||
let before = server.checkpoints(&run_id).await;
|
||||
assert_eq!(before.len(), 2, "start and one are recorded: {before:?}");
|
||||
crash(&mut server, worker, None);
|
||||
|
||||
server.release("record", "two");
|
||||
server.launch().await;
|
||||
wait_for_success(&server, &run_id).await;
|
||||
|
||||
let (path, commits) = server.workspace_commits(&run_id);
|
||||
assert_eq!(subjects(&commits), three_stage_subjects(&run_id));
|
||||
assert_eq!(
|
||||
std::fs::read_to_string(path.join("three.log")).expect("three copied the log"),
|
||||
"run\n",
|
||||
"the stage with a durable finish did not rerun"
|
||||
);
|
||||
let checkpoints = server.checkpoints(&run_id).await;
|
||||
assert_eq!(checkpoints.len(), 5, "{checkpoints:?}");
|
||||
let keys: Vec<Option<CheckpointKey>> = checkpoints.iter().map(|(key, _)| Some(*key)).collect();
|
||||
let committed: Vec<Option<CheckpointKey>> = commits.iter().map(|(_, key)| *key).collect();
|
||||
assert_eq!(
|
||||
keys, committed,
|
||||
"the reconciled record names the one commit"
|
||||
);
|
||||
server.shutdown();
|
||||
}
|
||||
|
||||
/// A workspace deleted while the run is down is restored from the snapshot
|
||||
/// repository, and the next stage sees the checkpoint's files.
|
||||
#[tokio::test(flavor = "multi_thread")]
|
||||
async fn a_deleted_workspace_is_restored_from_its_snapshot() {
|
||||
if host_plugin().is_none() {
|
||||
return;
|
||||
}
|
||||
let context = test_context!();
|
||||
let mut server = RunningServer::start().await;
|
||||
let gate = context.temp_dir.join("two.gate");
|
||||
let workspace = three_stage_bundle(&context, &gate);
|
||||
let run_id = run_detached(&context, &server, &workspace);
|
||||
|
||||
wait_for_status(&server, &run_id, &["running"]).await;
|
||||
let worker = wait_for_worker(&run_id);
|
||||
wait_until_gate_is_polled(&gate);
|
||||
crash(&mut server, worker, Some(&gate));
|
||||
let (path, commits) = server.workspace_commits(&run_id);
|
||||
assert_eq!(subjects(&commits), three_stage_subjects(&run_id)[..2]);
|
||||
std::fs::remove_dir_all(&path).expect("the workspace is deleted");
|
||||
|
||||
server.launch().await;
|
||||
wait_until_gate_is_polled(&gate);
|
||||
std::fs::write(&gate, "go").expect("the gate opens");
|
||||
wait_for_success(&server, &run_id).await;
|
||||
|
||||
let (restored, commits) = server.workspace_commits(&run_id);
|
||||
assert_eq!(restored, path);
|
||||
assert_eq!(subjects(&commits), three_stage_subjects(&run_id));
|
||||
assert_eq!(
|
||||
std::fs::read_to_string(restored.join("one.txt")).expect("one.txt was restored"),
|
||||
"one\n"
|
||||
);
|
||||
server.shutdown();
|
||||
}
|
||||
|
||||
/// A stage that fails on its own terms routes to its failure edge on the
|
||||
/// committed files, and after a crash once the failure is durable the
|
||||
/// route reruns on the same files.
|
||||
#[tokio::test(flavor = "multi_thread")]
|
||||
async fn a_failure_route_sees_the_same_committed_files_after_a_crash() {
|
||||
if host_plugin().is_none() {
|
||||
return;
|
||||
}
|
||||
let context = test_context!();
|
||||
let mut server = RunningServer::start().await;
|
||||
let gate = context.temp_dir.join("fix.gate");
|
||||
let workspace = write_petri_bundle(
|
||||
&context,
|
||||
&format!(
|
||||
"digraph Failure {{\n graph [goal=\"Route on failure\", default_max_retries=0]\n \
|
||||
start [shape=Mdiamond]\n exit [shape=Msquare]\n work [shape=parallelogram, \
|
||||
script=\"echo partial > out.txt; exit 1\"]\n fix [shape=parallelogram, script=\"test \
|
||||
\\\"$(cat out.txt)\\\" = partial && echo run >> fix.log && {} && echo fixed >> \
|
||||
out.txt\"]\n start -> work -> exit\n work -> fix [condition=\"outcome=failed\"]\n \
|
||||
fix -> exit\n}}\n",
|
||||
wait_for(&gate)
|
||||
),
|
||||
);
|
||||
let run_id = run_detached(&context, &server, &workspace);
|
||||
|
||||
wait_for_status(&server, &run_id, &["running"]).await;
|
||||
let worker = wait_for_worker(&run_id);
|
||||
wait_until_gate_is_polled(&gate);
|
||||
// The route is running on the committed failure: the crash lands here.
|
||||
crash(&mut server, worker, Some(&gate));
|
||||
|
||||
server.launch().await;
|
||||
wait_until_gate_is_polled(&gate);
|
||||
std::fs::write(&gate, "go").expect("the gate opens");
|
||||
wait_for_success(&server, &run_id).await;
|
||||
|
||||
let (path, commits) = server.workspace_commits(&run_id);
|
||||
assert_eq!(subjects(&commits), vec![
|
||||
format!("fabro({run_id}): start (success)"),
|
||||
format!("fabro({run_id}): work (failure)"),
|
||||
format!("fabro({run_id}): fix (success)"),
|
||||
format!("fabro({run_id}): exit (success)"),
|
||||
]);
|
||||
assert_eq!(
|
||||
std::fs::read_to_string(path.join("out.txt")).expect("out.txt"),
|
||||
"partial\nfixed\n"
|
||||
);
|
||||
assert_eq!(
|
||||
std::fs::read_to_string(path.join("fix.log")).expect("fix.log"),
|
||||
"run\n",
|
||||
"the route saw the failed stage's files, not its own interrupted attempt's"
|
||||
);
|
||||
server.shutdown();
|
||||
}
|
||||
|
||||
/// A checkpoint commit that fails ends the run: `checkpoint_failed` is
|
||||
/// recorded, no route runs, the run is reported failed, and a restart
|
||||
/// leaves it failed without launching a worker.
|
||||
#[tokio::test(flavor = "multi_thread")]
|
||||
async fn a_failed_checkpoint_fails_the_run_and_a_restart_leaves_it_failed() {
|
||||
if host_plugin().is_none() {
|
||||
return;
|
||||
}
|
||||
let context = test_context!();
|
||||
let mut server = RunningServer::start().await;
|
||||
let workspace = write_petri_bundle(
|
||||
&context,
|
||||
"digraph Wreck {\n graph [goal=\"Wreck the repository\", default_max_retries=0]\n start \
|
||||
[shape=Mdiamond]\n exit [shape=Msquare]\n wreck [shape=parallelogram, script=\"rm -rf \
|
||||
.git && echo garbage > .git\"]\n next [shape=parallelogram, script=\"echo next > \
|
||||
next.txt\"]\n fix [shape=parallelogram, script=\"echo fix > fix.txt\"]\n start -> wreck \
|
||||
-> next -> exit\n wreck -> fix [condition=\"outcome=failed\"]\n fix -> exit\n}\n",
|
||||
);
|
||||
let run_id = run_detached(&context, &server, &workspace);
|
||||
let status = wait_for_status(&server, &run_id, &["succeeded", "failed"]).await;
|
||||
let run = run_json(&server, &format!("runs/{run_id}")).await;
|
||||
assert_eq!(status, "failed", "run: {run}");
|
||||
let events = run_events(&server, &run_id).await;
|
||||
let failures: Vec<String> = events
|
||||
.iter()
|
||||
.filter_map(|envelope| match &envelope.event.body {
|
||||
EventBody::RunFailed(props) => Some(props.failure.detail.message.clone()),
|
||||
_ => None,
|
||||
})
|
||||
.collect();
|
||||
assert_eq!(failures.len(), 1, "{failures:?}");
|
||||
assert!(
|
||||
failures[0].contains("checkpoint commit of `wreck` failed"),
|
||||
"{failures:?}"
|
||||
);
|
||||
let scopes = server.petri_run_dir(&run_id).join("scopes");
|
||||
let work = std::fs::read_dir(&scopes)
|
||||
.expect("the scopes directory lists")
|
||||
.map(|entry| entry.expect("an entry reads").path().join("work"))
|
||||
.next()
|
||||
.expect("one workspace");
|
||||
assert!(!work.join("next.txt").exists(), "no route ran");
|
||||
assert!(!work.join("fix.txt").exists(), "no route ran");
|
||||
|
||||
let store = server.petri_store().await;
|
||||
let outcome = engine::outcome_of(&store, &run_id)
|
||||
.await
|
||||
.expect("the run's Petri record inspects");
|
||||
assert_ne!(outcome.status, RunStatus::Success, "{outcome:?}");
|
||||
let checkpoints = server.checkpoints(&run_id).await;
|
||||
assert_eq!(
|
||||
checkpoints.len(),
|
||||
1,
|
||||
"only start was recorded: {checkpoints:?}"
|
||||
);
|
||||
|
||||
// The restart finds the run terminal and launches nothing for it.
|
||||
server.kill();
|
||||
server.launch().await;
|
||||
assert_eq!(run_status(&server, &run_id).await, "failed");
|
||||
std::thread::sleep(Duration::from_secs(1));
|
||||
assert_eq!(
|
||||
worker_pid(&run_id),
|
||||
None,
|
||||
"no worker was launched for the failed run"
|
||||
);
|
||||
server.shutdown();
|
||||
}
|
||||
|
|
|
|||
|
|
@ -346,3 +346,91 @@ async fn a_worker_store_leases_for_its_launch_not_for_petris_owner() {
|
|||
drop(created);
|
||||
wait_until_released(server_store, &key).await;
|
||||
}
|
||||
|
||||
/// A worker stores Fabro's platform records for its run over the API and
|
||||
/// reads them back by kind: what the checkpoint hooks do from the worker
|
||||
/// process.
|
||||
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
|
||||
async fn a_worker_appends_and_reads_platform_records_over_the_api() {
|
||||
use fabro_petri::platform_records::{HttpPlatformRecords, PlatformRecords};
|
||||
use fabro_store::platform_records::{CheckpointRecord, DecisionRef, OperationKey};
|
||||
use fabro_store::{PlatformRecord, PlatformRecordKind, StagePosition};
|
||||
|
||||
let state = test_app_state();
|
||||
let base_url = serve(Arc::clone(&state), |router| router).await;
|
||||
let run_id = RunId::new();
|
||||
let token = state.test_issue_worker_token(&run_id);
|
||||
let records =
|
||||
HttpPlatformRecords::new(worker_client(&base_url, &token, Duration::from_secs(5)).await);
|
||||
|
||||
let checkpoint = PlatformRecord::Checkpoint(CheckpointRecord {
|
||||
execution: 0,
|
||||
firing: 3,
|
||||
attempt: Some(1),
|
||||
workspace: Some("invocation-0-scope-0".to_string()),
|
||||
git_commit_sha: Some("abc123".to_string()),
|
||||
diff_summary: None,
|
||||
patch_blob: None,
|
||||
operation: Some(OperationKey {
|
||||
execution: 0,
|
||||
decision: DecisionRef::AttemptStart {
|
||||
firing: 3,
|
||||
attempt: 1,
|
||||
},
|
||||
effect: "checkpoint".to_string(),
|
||||
}),
|
||||
});
|
||||
let position = Some(StagePosition {
|
||||
execution: 0,
|
||||
firing: 3,
|
||||
});
|
||||
let stored = records
|
||||
.append(&run_id, &checkpoint, position)
|
||||
.await
|
||||
.expect("the record appends over the API");
|
||||
assert_eq!(stored.seq, 1);
|
||||
assert_eq!(stored.position, position);
|
||||
let notice = PlatformRecord::RunArchived;
|
||||
records
|
||||
.append(&run_id, ¬ice, None)
|
||||
.await
|
||||
.expect("a second record appends");
|
||||
|
||||
let checkpoints = records
|
||||
.read_kind(&run_id, PlatformRecordKind::Checkpoint)
|
||||
.await
|
||||
.expect("the checkpoints read back");
|
||||
assert_eq!(checkpoints.len(), 1);
|
||||
assert_eq!(checkpoints[0].seq, 1);
|
||||
assert_eq!(checkpoints[0].position, position);
|
||||
assert!(matches!(
|
||||
&checkpoints[0].record,
|
||||
PlatformRecord::Checkpoint(record) if record.git_commit_sha.as_deref() == Some("abc123")
|
||||
&& record.operation == checkpoint_operation(&checkpoint)
|
||||
));
|
||||
let archived = records
|
||||
.read_kind(&run_id, PlatformRecordKind::RunArchived)
|
||||
.await
|
||||
.expect("the archive record reads back");
|
||||
assert_eq!(archived.len(), 1);
|
||||
assert_eq!(archived[0].seq, 2);
|
||||
assert_eq!(archived[0].position, None);
|
||||
|
||||
// Another run's token cannot read this run's records.
|
||||
let other = state.test_issue_worker_token(&RunId::new());
|
||||
let foreign =
|
||||
HttpPlatformRecords::new(worker_client(&base_url, &other, Duration::from_secs(5)).await);
|
||||
assert!(
|
||||
foreign
|
||||
.read_kind(&run_id, PlatformRecordKind::Checkpoint)
|
||||
.await
|
||||
.is_err(),
|
||||
"a worker token names one run"
|
||||
);
|
||||
}
|
||||
|
||||
fn checkpoint_operation(
|
||||
record: &fabro_store::PlatformRecord,
|
||||
) -> Option<fabro_store::platform_records::OperationKey> {
|
||||
record.operation().cloned()
|
||||
}
|
||||
|
|
|
|||
|
|
@ -53,4 +53,6 @@ fabro-auth = { path = "../../foundation/fabro-auth", features = ["test-support"]
|
|||
fabro-llm = { path = "../fabro-llm", features = ["test-support"] }
|
||||
fabro-store = { path = "../fabro-store", features = ["test-support"] }
|
||||
petri_testkit.workspace = true
|
||||
fabro-test = { workspace = true }
|
||||
lithos-llm = { workspace = true, features = ["runtime"] }
|
||||
tokio = { workspace = true, features = ["macros", "rt-multi-thread"] }
|
||||
|
|
|
|||
|
|
@ -125,6 +125,8 @@ pub enum Recovery {
|
|||
pub enum RecoveryError {
|
||||
#[error("the run's record could not be opened")]
|
||||
Open(#[source] StoreError),
|
||||
#[error("the run's coordinator log could not be read")]
|
||||
Log(#[source] petri_execution::StoreError),
|
||||
#[error("the run's coordinator state could not be read")]
|
||||
State(#[source] HostError),
|
||||
#[error("the run's record could not be inspected")]
|
||||
|
|
@ -159,6 +161,14 @@ pub async fn recover(request: RecoveryRequest) -> Result<Recovery, RecoveryError
|
|||
// A record with no root invocation (the worker died between creating
|
||||
// the run and declaring it) has nothing to reconcile; the worker's
|
||||
// resume reports it as such.
|
||||
let records = petri_execution::read_coordinator_log(&*logs)
|
||||
.await
|
||||
.map_err(RecoveryError::Log)?;
|
||||
if records.is_empty() {
|
||||
return Ok(Recovery::Resume {
|
||||
workspaces: Vec::new(),
|
||||
});
|
||||
}
|
||||
let state = host::stored_state(&*logs)
|
||||
.await
|
||||
.map_err(RecoveryError::State)?;
|
||||
|
|
@ -194,12 +204,14 @@ pub async fn recover(request: RecoveryRequest) -> Result<Recovery, RecoveryError
|
|||
let lookup = WorkspaceLookup::new(Arc::clone(&request.store), key);
|
||||
let recorded = recorded_checkpoints(&*request.records, &request.run_id).await?;
|
||||
|
||||
// The snapshot each live execution's workspace must sit on.
|
||||
// The snapshot each live execution's workspace must sit on. A live
|
||||
// execution is one whose log records no exit: `inspect_run` reports it
|
||||
// as incomplete.
|
||||
let mut candidates: BTreeMap<String, Vec<(Target, String)>> = BTreeMap::new();
|
||||
for execution in inspection
|
||||
.executions
|
||||
.iter()
|
||||
.filter(|execution| execution.status == "running")
|
||||
.filter(|execution| execution.status == "incomplete")
|
||||
{
|
||||
let Some(target) = last_finish(execution) else {
|
||||
continue;
|
||||
|
|
|
|||
|
|
@ -466,3 +466,148 @@ async fn recovery_starts_an_unknown_run_and_resumes_a_finished_one() {
|
|||
3
|
||||
);
|
||||
}
|
||||
|
||||
/// A `[[run.hooks]]` hook that blocks a tool effect keeps working through
|
||||
/// the forwarded local service: the agent's `rm` is refused by the
|
||||
/// `pre_tool_use` hook, the model is told why, and the file it aimed at is
|
||||
/// still in the stage's snapshot.
|
||||
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
|
||||
async fn a_run_hook_blocks_a_tool_effect_through_the_forwarded_service() {
|
||||
use fabro_auth::test_support::env_credential_source;
|
||||
use fabro_llm::test_support::test_catalog_with_provider_base_url;
|
||||
use fabro_petri::runtime;
|
||||
use fabro_test::{TwinScenario, TwinScenarios, TwinToolCall};
|
||||
use lithos_llm::catalog::ProviderId;
|
||||
use serde_json::json;
|
||||
|
||||
const MODEL: &str = "gpt-5.6-sol";
|
||||
if host_plugin().is_none() {
|
||||
return;
|
||||
}
|
||||
let twin = fabro_test::twin_openai().await;
|
||||
let namespace = format!("{}::{}", module_path!(), line!());
|
||||
TwinScenarios::new(namespace.clone())
|
||||
.scenario(
|
||||
TwinScenario::responses(MODEL)
|
||||
.input_contains("Remove the scratch file")
|
||||
.tool_call(TwinToolCall::new(
|
||||
"shell_command",
|
||||
json!({ "command": "rm -f scratch.txt && echo removed" }),
|
||||
)),
|
||||
)
|
||||
.scenario(
|
||||
TwinScenario::responses(MODEL)
|
||||
.input_contains("destructive commands are not allowed")
|
||||
.text("Understood, the file stays."),
|
||||
)
|
||||
.load(twin)
|
||||
.await;
|
||||
let api_key = namespace.clone();
|
||||
let credentials =
|
||||
env_credential_source(move |name| (name == "OPENAI_API_KEY").then(|| api_key.clone()));
|
||||
let client = runtime::model_client(
|
||||
test_catalog_with_provider_base_url("openai", &twin.base_url),
|
||||
credentials,
|
||||
None,
|
||||
&[ProviderId::new("openai")],
|
||||
)
|
||||
.expect("the model client builds")
|
||||
.expect("openai is eligible");
|
||||
|
||||
let harness = Harness::new();
|
||||
let workflow = format!(
|
||||
"digraph Hooks {{\n graph [backend=\"api\", goal=\"Check the tool hooks\", \
|
||||
default_max_retries=0]\n start [shape=Mdiamond]\n exit [shape=Msquare]\n seed \
|
||||
[shape=parallelogram, script=\"echo keep > scratch.txt\"]\n agent [prompt=\"Remove the \
|
||||
scratch file with the shell tool.\", model=\"{MODEL}\", provider=\"openai\", \
|
||||
fidelity=\"full\"]\n check [shape=parallelogram, script=\"test \\\"$(cat scratch.txt)\\\" \
|
||||
= keep\"]\n start -> seed -> agent -> check -> exit\n}}\n"
|
||||
);
|
||||
let settings = format!(
|
||||
"{SETTINGS}\n[[run.hooks]]\nname = \"no-destruction\"\nevent = \"pre_tool_use\"\nmatcher \
|
||||
= \"shell\"\nscript = \"if grep -q 'rm ' \\\"$FABRO_HOOK_CONTEXT\\\"; then echo \
|
||||
'{{\\\"decision\\\":\\\"block\\\",\\\"reason\\\":\\\"destructive commands are not \
|
||||
allowed\\\"}}'; exit 2; fi\"\n"
|
||||
);
|
||||
let request = RunRequest {
|
||||
run_id: harness.run_id.to_string(),
|
||||
run_dir: harness.run_dir.clone(),
|
||||
execution: Execution::Start(admit(&workflow, &settings)),
|
||||
store: Arc::clone(&harness.store) as Arc<dyn petri_store::RunStore>,
|
||||
runtime: RuntimeSpec {
|
||||
model_client: Some(client),
|
||||
..RuntimeSpec::default()
|
||||
},
|
||||
provider: SandboxProviderKind::LOCAL,
|
||||
cancel: CancellationToken::new(),
|
||||
hooks: Some(harness.hooks()),
|
||||
};
|
||||
let outcome = engine::run(request).await.expect("the run executes");
|
||||
assert_eq!(outcome.status, RunStatus::Success, "{outcome:?}");
|
||||
let inspection = harness.inspection().await;
|
||||
assert_eq!(stages(&inspection), vec![
|
||||
("start".to_string(), "success".to_string()),
|
||||
("seed".to_string(), "success".to_string()),
|
||||
("agent".to_string(), "success".to_string()),
|
||||
("check".to_string(), "success".to_string()),
|
||||
("exit".to_string(), "success".to_string()),
|
||||
]);
|
||||
let requests = twin.request_logs(&namespace).await;
|
||||
let inputs: Vec<&str> = requests["requests"]
|
||||
.as_array()
|
||||
.map(|requests| {
|
||||
requests
|
||||
.iter()
|
||||
.filter_map(|request| request["input_text"].as_str())
|
||||
.collect()
|
||||
})
|
||||
.unwrap_or_default();
|
||||
assert_eq!(inputs.len(), 2, "{inputs:?}");
|
||||
assert!(
|
||||
inputs[1].contains("destructive commands are not allowed"),
|
||||
"the model was told why the tool was blocked: {inputs:?}"
|
||||
);
|
||||
let workspace = harness.workspace().await;
|
||||
let path = harness.workspace_path(&workspace);
|
||||
assert_eq!(
|
||||
git(&path, &["show", "HEAD:scratch.txt"]).await,
|
||||
"keep",
|
||||
"the blocked removal never happened"
|
||||
);
|
||||
assert_eq!(harness.checkpoints().len(), 5);
|
||||
}
|
||||
|
||||
/// The run's records name the workspace its root invocation ran in: what
|
||||
/// recovery reads to find the workspace of a live execution.
|
||||
#[tokio::test]
|
||||
async fn the_records_name_the_root_invocations_workspace() {
|
||||
use fabro_petri::workspace::WorkspaceLookup;
|
||||
use petri_execution::InvocationId;
|
||||
|
||||
if host_plugin().is_none() {
|
||||
return;
|
||||
}
|
||||
let harness = Harness::new();
|
||||
let workflow = workflow(
|
||||
" write [shape=parallelogram, script=\"echo one > out.txt\"]",
|
||||
" start -> write -> exit",
|
||||
);
|
||||
let outcome = harness.run(&workflow, SETTINGS).await;
|
||||
assert_eq!(outcome.status, RunStatus::Success, "{outcome:?}");
|
||||
let lookup = WorkspaceLookup::new(
|
||||
Arc::clone(&harness.store) as Arc<dyn petri_store::RunStore>,
|
||||
RunKey::new(harness.run_id.to_string()),
|
||||
);
|
||||
let named = lookup
|
||||
.of_invocation(InvocationId::ROOT)
|
||||
.await
|
||||
.expect("the lookup reads the records");
|
||||
assert_eq!(named, vec![harness.workspace().await]);
|
||||
let inspection = harness.inspection().await;
|
||||
let statuses: Vec<&str> = inspection
|
||||
.executions
|
||||
.iter()
|
||||
.map(|execution| execution.status)
|
||||
.collect();
|
||||
assert_eq!(statuses, vec!["finished"]);
|
||||
}
|
||||
|
|
|
|||
|
|
@ -299,6 +299,9 @@ models/petri-access.ts
|
|||
models/petri-append-request.ts
|
||||
models/petri-open-request.ts
|
||||
models/petri-open-response.ts
|
||||
models/petri-platform-record-append-request.ts
|
||||
models/petri-platform-record-list.ts
|
||||
models/petri-platform-record.ts
|
||||
models/petri-record-list.ts
|
||||
models/petri-record.ts
|
||||
models/petri-release-request.ts
|
||||
|
|
|
|||
|
|
@ -40,6 +40,12 @@ import type { PetriOpenRequest } from '../models';
|
|||
// @ts-ignore
|
||||
import type { PetriOpenResponse } from '../models';
|
||||
// @ts-ignore
|
||||
import type { PetriPlatformRecord } from '../models';
|
||||
// @ts-ignore
|
||||
import type { PetriPlatformRecordAppendRequest } from '../models';
|
||||
// @ts-ignore
|
||||
import type { PetriPlatformRecordList } from '../models';
|
||||
// @ts-ignore
|
||||
import type { PetriRecordList } from '../models';
|
||||
// @ts-ignore
|
||||
import type { PetriReleaseRequest } from '../models';
|
||||
|
|
@ -66,6 +72,51 @@ import type { WriteRunBlobRequest } from '../models';
|
|||
*/
|
||||
export const RunInternalsApiAxiosParamCreator = function (configuration?: Configuration) {
|
||||
return {
|
||||
/**
|
||||
* Stores one platform record at the run\'s next `seq`, tied to the Petri stage named by `execution` and `firing` when it belongs to one. The record is the JSON of a Fabro platform record, tagged by `kind`.
|
||||
* @summary Append Petri Platform Record
|
||||
* @param {string} id Unique run identifier (ULID).
|
||||
* @param {PetriPlatformRecordAppendRequest} petriPlatformRecordAppendRequest
|
||||
* @param {*} [options] Override http request option.
|
||||
* @throws {RequiredError}
|
||||
*/
|
||||
appendPetriPlatformRecord: async (id: string, petriPlatformRecordAppendRequest: PetriPlatformRecordAppendRequest, options: RawAxiosRequestConfig = {}): Promise<RequestArgs> => {
|
||||
// verify required parameter 'id' is not null or undefined
|
||||
assertParamExists('appendPetriPlatformRecord', 'id', id)
|
||||
// verify required parameter 'petriPlatformRecordAppendRequest' is not null or undefined
|
||||
assertParamExists('appendPetriPlatformRecord', 'petriPlatformRecordAppendRequest', petriPlatformRecordAppendRequest)
|
||||
const localVarPath = `/api/v1/runs/{id}/petri/platform-records`
|
||||
.replace(`{${"id"}}`, encodeURIComponent(String(id)));
|
||||
// use dummy base URL string because the URL constructor only accepts absolute URLs.
|
||||
const localVarUrlObj = new URL(localVarPath, DUMMY_BASE_URL);
|
||||
let baseOptions;
|
||||
if (configuration) {
|
||||
baseOptions = configuration.baseOptions;
|
||||
}
|
||||
|
||||
const localVarRequestOptions = { method: 'POST', ...baseOptions, ...options};
|
||||
const localVarHeaderParameter = {} as any;
|
||||
const localVarQueryParameter = {} as any;
|
||||
|
||||
// authentication SessionCookie required
|
||||
|
||||
// authentication BearerAuth required
|
||||
// http bearer authentication required
|
||||
await setBearerAuthToObject(localVarHeaderParameter, configuration)
|
||||
|
||||
localVarHeaderParameter['Content-Type'] = 'application/json';
|
||||
localVarHeaderParameter['Accept'] = 'application/json';
|
||||
|
||||
setSearchParams(localVarUrlObj, localVarQueryParameter);
|
||||
let headersFromBaseOptions = baseOptions && baseOptions.headers ? baseOptions.headers : {};
|
||||
localVarRequestOptions.headers = {...localVarHeaderParameter, ...headersFromBaseOptions, ...options.headers};
|
||||
localVarRequestOptions.data = serializeDataIfNeeded(petriPlatformRecordAppendRequest, localVarRequestOptions, configuration)
|
||||
|
||||
return {
|
||||
url: toPathString(localVarUrlObj),
|
||||
options: localVarRequestOptions,
|
||||
};
|
||||
},
|
||||
/**
|
||||
* Appends one batch of records to one log at the sequences they carry, durably, in one transaction. A record equal to the one already stored at its `seq` is accepted without a second append, so a batch whose reply was lost is safe to resend. A different record at a taken `seq`, or a `seq` past the log\'s end, is refused with `petri_record_conflict` and the batch stores nothing.
|
||||
* @summary Append Petri Records
|
||||
|
|
@ -530,6 +581,51 @@ export const RunInternalsApiAxiosParamCreator = function (configuration?: Config
|
|||
options: localVarRequestOptions,
|
||||
};
|
||||
},
|
||||
/**
|
||||
* The run\'s platform records (Fabro\'s own facts about a Petri run: a checkpoint commit, a pull request, a notification), in `seq` order, optionally of one kind. What a run\'s worker reads to find an effect it already performed before performing it again.
|
||||
* @summary List Petri Platform Records
|
||||
* @param {string} id Unique run identifier (ULID).
|
||||
* @param {string} [kind] Only the platform records of this kind, as its `kind` tag spells it (`checkpoint`, `pull_request.created`, ...).
|
||||
* @param {*} [options] Override http request option.
|
||||
* @throws {RequiredError}
|
||||
*/
|
||||
listPetriPlatformRecords: async (id: string, kind?: string, options: RawAxiosRequestConfig = {}): Promise<RequestArgs> => {
|
||||
// verify required parameter 'id' is not null or undefined
|
||||
assertParamExists('listPetriPlatformRecords', 'id', id)
|
||||
const localVarPath = `/api/v1/runs/{id}/petri/platform-records`
|
||||
.replace(`{${"id"}}`, encodeURIComponent(String(id)));
|
||||
// use dummy base URL string because the URL constructor only accepts absolute URLs.
|
||||
const localVarUrlObj = new URL(localVarPath, DUMMY_BASE_URL);
|
||||
let baseOptions;
|
||||
if (configuration) {
|
||||
baseOptions = configuration.baseOptions;
|
||||
}
|
||||
|
||||
const localVarRequestOptions = { method: 'GET', ...baseOptions, ...options};
|
||||
const localVarHeaderParameter = {} as any;
|
||||
const localVarQueryParameter = {} as any;
|
||||
|
||||
// authentication SessionCookie required
|
||||
|
||||
// authentication BearerAuth required
|
||||
// http bearer authentication required
|
||||
await setBearerAuthToObject(localVarHeaderParameter, configuration)
|
||||
|
||||
if (kind !== undefined) {
|
||||
localVarQueryParameter['kind'] = kind;
|
||||
}
|
||||
|
||||
localVarHeaderParameter['Accept'] = 'application/json';
|
||||
|
||||
setSearchParams(localVarUrlObj, localVarQueryParameter);
|
||||
let headersFromBaseOptions = baseOptions && baseOptions.headers ? baseOptions.headers : {};
|
||||
localVarRequestOptions.headers = {...localVarHeaderParameter, ...headersFromBaseOptions, ...options.headers};
|
||||
|
||||
return {
|
||||
url: toPathString(localVarUrlObj),
|
||||
options: localVarRequestOptions,
|
||||
};
|
||||
},
|
||||
/**
|
||||
* Every record of one log of the run, in `seq` order, unchanged.
|
||||
* @summary List Petri Records
|
||||
|
|
@ -1246,6 +1342,20 @@ export const RunInternalsApiAxiosParamCreator = function (configuration?: Config
|
|||
export const RunInternalsApiFp = function(configuration?: Configuration) {
|
||||
const localVarAxiosParamCreator = RunInternalsApiAxiosParamCreator(configuration)
|
||||
return {
|
||||
/**
|
||||
* Stores one platform record at the run\'s next `seq`, tied to the Petri stage named by `execution` and `firing` when it belongs to one. The record is the JSON of a Fabro platform record, tagged by `kind`.
|
||||
* @summary Append Petri Platform Record
|
||||
* @param {string} id Unique run identifier (ULID).
|
||||
* @param {PetriPlatformRecordAppendRequest} petriPlatformRecordAppendRequest
|
||||
* @param {*} [options] Override http request option.
|
||||
* @throws {RequiredError}
|
||||
*/
|
||||
async appendPetriPlatformRecord(id: string, petriPlatformRecordAppendRequest: PetriPlatformRecordAppendRequest, options?: RawAxiosRequestConfig): Promise<(axios?: AxiosInstance, basePath?: string) => AxiosPromise<PetriPlatformRecord>> {
|
||||
const localVarAxiosArgs = await localVarAxiosParamCreator.appendPetriPlatformRecord(id, petriPlatformRecordAppendRequest, options);
|
||||
const localVarOperationServerIndex = configuration?.serverIndex ?? 0;
|
||||
const localVarOperationServerBasePath = operationServerMap['RunInternalsApi.appendPetriPlatformRecord']?.[localVarOperationServerIndex]?.url;
|
||||
return (axios, basePath) => createRequestFunction(localVarAxiosArgs, globalAxios, BASE_PATH, configuration)(axios, localVarOperationServerBasePath || basePath);
|
||||
},
|
||||
/**
|
||||
* Appends one batch of records to one log at the sequences they carry, durably, in one transaction. A record equal to the one already stored at its `seq` is accepted without a second append, so a batch whose reply was lost is safe to resend. A different record at a taken `seq`, or a `seq` past the log\'s end, is refused with `petri_record_conflict` and the batch stores nothing.
|
||||
* @summary Append Petri Records
|
||||
|
|
@ -1389,6 +1499,20 @@ export const RunInternalsApiFp = function(configuration?: Configuration) {
|
|||
const localVarOperationServerBasePath = operationServerMap['RunInternalsApi.getStageArtifact']?.[localVarOperationServerIndex]?.url;
|
||||
return (axios, basePath) => createRequestFunction(localVarAxiosArgs, globalAxios, BASE_PATH, configuration)(axios, localVarOperationServerBasePath || basePath);
|
||||
},
|
||||
/**
|
||||
* The run\'s platform records (Fabro\'s own facts about a Petri run: a checkpoint commit, a pull request, a notification), in `seq` order, optionally of one kind. What a run\'s worker reads to find an effect it already performed before performing it again.
|
||||
* @summary List Petri Platform Records
|
||||
* @param {string} id Unique run identifier (ULID).
|
||||
* @param {string} [kind] Only the platform records of this kind, as its `kind` tag spells it (`checkpoint`, `pull_request.created`, ...).
|
||||
* @param {*} [options] Override http request option.
|
||||
* @throws {RequiredError}
|
||||
*/
|
||||
async listPetriPlatformRecords(id: string, kind?: string, options?: RawAxiosRequestConfig): Promise<(axios?: AxiosInstance, basePath?: string) => AxiosPromise<PetriPlatformRecordList>> {
|
||||
const localVarAxiosArgs = await localVarAxiosParamCreator.listPetriPlatformRecords(id, kind, options);
|
||||
const localVarOperationServerIndex = configuration?.serverIndex ?? 0;
|
||||
const localVarOperationServerBasePath = operationServerMap['RunInternalsApi.listPetriPlatformRecords']?.[localVarOperationServerIndex]?.url;
|
||||
return (axios, basePath) => createRequestFunction(localVarAxiosArgs, globalAxios, BASE_PATH, configuration)(axios, localVarOperationServerBasePath || basePath);
|
||||
},
|
||||
/**
|
||||
* Every record of one log of the run, in `seq` order, unchanged.
|
||||
* @summary List Petri Records
|
||||
|
|
@ -1615,6 +1739,17 @@ export const RunInternalsApiFp = function(configuration?: Configuration) {
|
|||
export const RunInternalsApiFactory = function (configuration?: Configuration, basePath?: string, axios?: AxiosInstance) {
|
||||
const localVarFp = RunInternalsApiFp(configuration)
|
||||
return {
|
||||
/**
|
||||
* Stores one platform record at the run\'s next `seq`, tied to the Petri stage named by `execution` and `firing` when it belongs to one. The record is the JSON of a Fabro platform record, tagged by `kind`.
|
||||
* @summary Append Petri Platform Record
|
||||
* @param {string} id Unique run identifier (ULID).
|
||||
* @param {PetriPlatformRecordAppendRequest} petriPlatformRecordAppendRequest
|
||||
* @param {*} [options] Override http request option.
|
||||
* @throws {RequiredError}
|
||||
*/
|
||||
appendPetriPlatformRecord(id: string, petriPlatformRecordAppendRequest: PetriPlatformRecordAppendRequest, options?: RawAxiosRequestConfig): AxiosPromise<PetriPlatformRecord> {
|
||||
return localVarFp.appendPetriPlatformRecord(id, petriPlatformRecordAppendRequest, options).then((request) => request(axios, basePath));
|
||||
},
|
||||
/**
|
||||
* Appends one batch of records to one log at the sequences they carry, durably, in one transaction. A record equal to the one already stored at its `seq` is accepted without a second append, so a batch whose reply was lost is safe to resend. A different record at a taken `seq`, or a `seq` past the log\'s end, is refused with `petri_record_conflict` and the batch stores nothing.
|
||||
* @summary Append Petri Records
|
||||
|
|
@ -1728,6 +1863,17 @@ export const RunInternalsApiFactory = function (configuration?: Configuration, b
|
|||
getStageArtifact(id: string, stageId: string, filename: string, retry: number, options?: RawAxiosRequestConfig): AxiosPromise<File> {
|
||||
return localVarFp.getStageArtifact(id, stageId, filename, retry, options).then((request) => request(axios, basePath));
|
||||
},
|
||||
/**
|
||||
* The run\'s platform records (Fabro\'s own facts about a Petri run: a checkpoint commit, a pull request, a notification), in `seq` order, optionally of one kind. What a run\'s worker reads to find an effect it already performed before performing it again.
|
||||
* @summary List Petri Platform Records
|
||||
* @param {string} id Unique run identifier (ULID).
|
||||
* @param {string} [kind] Only the platform records of this kind, as its `kind` tag spells it (`checkpoint`, `pull_request.created`, ...).
|
||||
* @param {*} [options] Override http request option.
|
||||
* @throws {RequiredError}
|
||||
*/
|
||||
listPetriPlatformRecords(id: string, kind?: string, options?: RawAxiosRequestConfig): AxiosPromise<PetriPlatformRecordList> {
|
||||
return localVarFp.listPetriPlatformRecords(id, kind, options).then((request) => request(axios, basePath));
|
||||
},
|
||||
/**
|
||||
* Every record of one log of the run, in `seq` order, unchanged.
|
||||
* @summary List Petri Records
|
||||
|
|
@ -1907,6 +2053,18 @@ export const RunInternalsApiFactory = function (configuration?: Configuration, b
|
|||
* RunInternalsApi - object-oriented interface
|
||||
*/
|
||||
export class RunInternalsApi extends BaseAPI {
|
||||
/**
|
||||
* Stores one platform record at the run\'s next `seq`, tied to the Petri stage named by `execution` and `firing` when it belongs to one. The record is the JSON of a Fabro platform record, tagged by `kind`.
|
||||
* @summary Append Petri Platform Record
|
||||
* @param {string} id Unique run identifier (ULID).
|
||||
* @param {PetriPlatformRecordAppendRequest} petriPlatformRecordAppendRequest
|
||||
* @param {*} [options] Override http request option.
|
||||
* @throws {RequiredError}
|
||||
*/
|
||||
public appendPetriPlatformRecord(id: string, petriPlatformRecordAppendRequest: PetriPlatformRecordAppendRequest, options?: RawAxiosRequestConfig) {
|
||||
return RunInternalsApiFp(this.configuration).appendPetriPlatformRecord(id, petriPlatformRecordAppendRequest, options).then((request) => request(this.axios, this.basePath));
|
||||
}
|
||||
|
||||
/**
|
||||
* Appends one batch of records to one log at the sequences they carry, durably, in one transaction. A record equal to the one already stored at its `seq` is accepted without a second append, so a batch whose reply was lost is safe to resend. A different record at a taken `seq`, or a `seq` past the log\'s end, is refused with `petri_record_conflict` and the batch stores nothing.
|
||||
* @summary Append Petri Records
|
||||
|
|
@ -2030,6 +2188,18 @@ export class RunInternalsApi extends BaseAPI {
|
|||
return RunInternalsApiFp(this.configuration).getStageArtifact(id, stageId, filename, retry, options).then((request) => request(this.axios, this.basePath));
|
||||
}
|
||||
|
||||
/**
|
||||
* The run\'s platform records (Fabro\'s own facts about a Petri run: a checkpoint commit, a pull request, a notification), in `seq` order, optionally of one kind. What a run\'s worker reads to find an effect it already performed before performing it again.
|
||||
* @summary List Petri Platform Records
|
||||
* @param {string} id Unique run identifier (ULID).
|
||||
* @param {string} [kind] Only the platform records of this kind, as its `kind` tag spells it (`checkpoint`, `pull_request.created`, ...).
|
||||
* @param {*} [options] Override http request option.
|
||||
* @throws {RequiredError}
|
||||
*/
|
||||
public listPetriPlatformRecords(id: string, kind?: string, options?: RawAxiosRequestConfig) {
|
||||
return RunInternalsApiFp(this.configuration).listPetriPlatformRecords(id, kind, options).then((request) => request(this.axios, this.basePath));
|
||||
}
|
||||
|
||||
/**
|
||||
* Every record of one log of the run, in `seq` order, unchanged.
|
||||
* @summary List Petri Records
|
||||
|
|
|
|||
|
|
@ -269,6 +269,9 @@ export * from './petri-access';
|
|||
export * from './petri-append-request';
|
||||
export * from './petri-open-request';
|
||||
export * from './petri-open-response';
|
||||
export * from './petri-platform-record';
|
||||
export * from './petri-platform-record-append-request';
|
||||
export * from './petri-platform-record-list';
|
||||
export * from './petri-record';
|
||||
export * from './petri-record-list';
|
||||
export * from './petri-release-request';
|
||||
|
|
|
|||
33
lib/packages/fabro-api-client/src/models/petri-platform-record-append-request.ts
generated
Normal file
33
lib/packages/fabro-api-client/src/models/petri-platform-record-append-request.ts
generated
Normal file
|
|
@ -0,0 +1,33 @@
|
|||
/* tslint:disable */
|
||||
/* eslint-disable */
|
||||
/**
|
||||
* Fabro Run API
|
||||
* HTTP API for managing Fabro workflow run executions.
|
||||
*
|
||||
* The version of the OpenAPI document: 0.2.0
|
||||
*
|
||||
*
|
||||
* NOTE: This class is auto generated by OpenAPI Generator (https://openapi-generator.tech).
|
||||
* https://openapi-generator.tech
|
||||
* Do not edit the class manually.
|
||||
*/
|
||||
|
||||
|
||||
|
||||
/**
|
||||
* One platform record to store for the run.
|
||||
*/
|
||||
export interface PetriPlatformRecordAppendRequest {
|
||||
/**
|
||||
* The platform record, tagged by `kind`.
|
||||
*/
|
||||
'record': { [key: string]: any; };
|
||||
/**
|
||||
* The Petri execution the record belongs to, with `firing`.
|
||||
*/
|
||||
'execution'?: number;
|
||||
/**
|
||||
* The Petri firing the record belongs to, with `execution`.
|
||||
*/
|
||||
'firing'?: number;
|
||||
}
|
||||
25
lib/packages/fabro-api-client/src/models/petri-platform-record-list.ts
generated
Normal file
25
lib/packages/fabro-api-client/src/models/petri-platform-record-list.ts
generated
Normal file
|
|
@ -0,0 +1,25 @@
|
|||
/* tslint:disable */
|
||||
/* eslint-disable */
|
||||
/**
|
||||
* Fabro Run API
|
||||
* HTTP API for managing Fabro workflow run executions.
|
||||
*
|
||||
* The version of the OpenAPI document: 0.2.0
|
||||
*
|
||||
*
|
||||
* NOTE: This class is auto generated by OpenAPI Generator (https://openapi-generator.tech).
|
||||
* https://openapi-generator.tech
|
||||
* Do not edit the class manually.
|
||||
*/
|
||||
|
||||
|
||||
// May contain unused imports in some cases
|
||||
// @ts-ignore
|
||||
import type { PetriPlatformRecord } from './petri-platform-record';
|
||||
|
||||
/**
|
||||
* The run\'s platform records, in `seq` order.
|
||||
*/
|
||||
export interface PetriPlatformRecordList {
|
||||
'records': Array<PetriPlatformRecord>;
|
||||
}
|
||||
41
lib/packages/fabro-api-client/src/models/petri-platform-record.ts
generated
Normal file
41
lib/packages/fabro-api-client/src/models/petri-platform-record.ts
generated
Normal file
|
|
@ -0,0 +1,41 @@
|
|||
/* tslint:disable */
|
||||
/* eslint-disable */
|
||||
/**
|
||||
* Fabro Run API
|
||||
* HTTP API for managing Fabro workflow run executions.
|
||||
*
|
||||
* The version of the OpenAPI document: 0.2.0
|
||||
*
|
||||
*
|
||||
* NOTE: This class is auto generated by OpenAPI Generator (https://openapi-generator.tech).
|
||||
* https://openapi-generator.tech
|
||||
* Do not edit the class manually.
|
||||
*/
|
||||
|
||||
|
||||
|
||||
/**
|
||||
* One of Fabro\'s platform records of a Petri run, as stored: the record\'s JSON tagged by `kind`, its position in the run\'s platform record sequence, and the Petri stage it belongs to when it belongs to one.
|
||||
*/
|
||||
export interface PetriPlatformRecord {
|
||||
/**
|
||||
* The record\'s position in the run\'s platform records, from 1.
|
||||
*/
|
||||
'seq': number;
|
||||
/**
|
||||
* Milliseconds since the Unix epoch when the record was stored.
|
||||
*/
|
||||
'recorded_at': number;
|
||||
/**
|
||||
* The platform record itself, tagged by `kind`.
|
||||
*/
|
||||
'record': { [key: string]: any; };
|
||||
/**
|
||||
* The Petri execution the record belongs to, with `firing`.
|
||||
*/
|
||||
'execution'?: number;
|
||||
/**
|
||||
* The Petri firing the record belongs to, with `execution`.
|
||||
*/
|
||||
'firing'?: number;
|
||||
}
|
||||
Loading…
Add table
Reference in a new issue