diff --git a/Cargo.lock b/Cargo.lock index 82600c564..a845a3035 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -2900,6 +2900,7 @@ dependencies = [ "fabro-llm", "fabro-petri", "fabro-store", + "fabro-test", "fabro-types", "fabro-util", "lithos-llm", diff --git a/lib/apps/fabro-cli/tests/it/scenario/petri.rs b/lib/apps/fabro-cli/tests/it/scenario/petri.rs index ec95672b1..75eae6873 100644 --- a/lib/apps/fabro-cli/tests/it/scenario/petri.rs +++ b/lib/apps/fabro-cli/tests/it/scenario/petri.rs @@ -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)>) { + let scopes = self.petri_run_dir(run_id).join("scopes"); + let mut workspaces: Vec = 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 { .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::().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)]) -> 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 { + ["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> = checkpoints.iter().map(|(key, _)| Some(*key)).collect(); + let committed: Vec> = 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> = checkpoints.iter().map(|(key, _)| Some(*key)).collect(); + let committed: Vec> = 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 = 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(); +} diff --git a/lib/apps/fabro-server/tests/it/api/petri_store.rs b/lib/apps/fabro-server/tests/it/api/petri_store.rs index 7701dd156..e87ff62a1 100644 --- a/lib/apps/fabro-server/tests/it/api/petri_store.rs +++ b/lib/apps/fabro-server/tests/it/api/petri_store.rs @@ -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 { + record.operation().cloned() +} diff --git a/lib/components/fabro-petri/Cargo.toml b/lib/components/fabro-petri/Cargo.toml index 82967dca2..0d967f79f 100644 --- a/lib/components/fabro-petri/Cargo.toml +++ b/lib/components/fabro-petri/Cargo.toml @@ -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"] } diff --git a/lib/components/fabro-petri/src/recovery.rs b/lib/components/fabro-petri/src/recovery.rs index b80f4b26c..c5da156df 100644 --- a/lib/components/fabro-petri/src/recovery.rs +++ b/lib/components/fabro-petri/src/recovery.rs @@ -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 Result> = 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; diff --git a/lib/components/fabro-petri/tests/hooks.rs b/lib/components/fabro-petri/tests/hooks.rs index 6493ca9c6..3a7154920 100644 --- a/lib/components/fabro-petri/tests/hooks.rs +++ b/lib/components/fabro-petri/tests/hooks.rs @@ -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, + 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, + 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"]); +} diff --git a/lib/packages/fabro-api-client/src/.openapi-generator/FILES b/lib/packages/fabro-api-client/src/.openapi-generator/FILES index 0d0b0f64d..da688b54b 100644 --- a/lib/packages/fabro-api-client/src/.openapi-generator/FILES +++ b/lib/packages/fabro-api-client/src/.openapi-generator/FILES @@ -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 diff --git a/lib/packages/fabro-api-client/src/api/run-internals-api.ts b/lib/packages/fabro-api-client/src/api/run-internals-api.ts index ba15825ed..46f924e9e 100644 --- a/lib/packages/fabro-api-client/src/api/run-internals-api.ts +++ b/lib/packages/fabro-api-client/src/api/run-internals-api.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 => { + // 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 => { + // 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> { + 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> { + 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 { + 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 { 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 { + 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 diff --git a/lib/packages/fabro-api-client/src/models/index.ts b/lib/packages/fabro-api-client/src/models/index.ts index 897027f9d..bb9d7664e 100644 --- a/lib/packages/fabro-api-client/src/models/index.ts +++ b/lib/packages/fabro-api-client/src/models/index.ts @@ -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'; diff --git a/lib/packages/fabro-api-client/src/models/petri-platform-record-append-request.ts b/lib/packages/fabro-api-client/src/models/petri-platform-record-append-request.ts new file mode 100644 index 000000000..72443ade9 --- /dev/null +++ b/lib/packages/fabro-api-client/src/models/petri-platform-record-append-request.ts @@ -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; +} diff --git a/lib/packages/fabro-api-client/src/models/petri-platform-record-list.ts b/lib/packages/fabro-api-client/src/models/petri-platform-record-list.ts new file mode 100644 index 000000000..af6e45cc8 --- /dev/null +++ b/lib/packages/fabro-api-client/src/models/petri-platform-record-list.ts @@ -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; +} diff --git a/lib/packages/fabro-api-client/src/models/petri-platform-record.ts b/lib/packages/fabro-api-client/src/models/petri-platform-record.ts new file mode 100644 index 000000000..5b259840e --- /dev/null +++ b/lib/packages/fabro-api-client/src/models/petri-platform-record.ts @@ -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; +}