mirror of
https://github.com/fabro-sh/fabro.git
synced 2026-10-03 02:24:33 +00:00
Let the Petri engine start or resume a run over any store
`fabro_petri::engine` is now the one assembly the worker process and the server share: `RunRequest` takes the run's store as `Arc<dyn RunStore>` and an `Execution`, either `Start` with the admitted graphs or `Resume` from the run's records through `host::resume_configured`, with the same interview observer a start installs. A resume whose record has no root invocation is refused with a named error instead of a panic in the host. The outcome is mapped to a `Conclusion` (succeeded, or failed with Fabro's reason and a message) so both callers record the same terminal event. `admission::load_with` loads the admitted graphs through any blob read, so a worker loads them through its client; `admission::load` over the server's `BlobStore` delegates to it. `HttpRunStore::for_worker` takes every lease for the worker's launch id, whatever owner Petri minted for the run runtime, and the open logs the owner. The module docs state the rule. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
This commit is contained in:
parent
832f39f704
commit
a621fb72e1
7 changed files with 345 additions and 59 deletions
|
|
@ -34,6 +34,7 @@ petri_frontend_fabro.workspace = true
|
|||
lithos-llm = { workspace = true, features = ["runtime"] }
|
||||
petri_testkit = { workspace = true, optional = true }
|
||||
anyhow.workspace = true
|
||||
bytes.workspace = true
|
||||
async-trait.workspace = true
|
||||
serde.workspace = true
|
||||
serde_json.workspace = true
|
||||
|
|
|
|||
|
|
@ -29,16 +29,20 @@ Every adapter the integration plan describes lands here.
|
|||
Petri's diagnostics come back in a shape the server maps onto Fabro's.
|
||||
- `admission`: the admitted graphs in Fabro's blob store, named on the run
|
||||
spec as `RunEngine::Petri(PetriAdmission)`, verified by digest on load.
|
||||
- `engine`: a run executed by Petri in the server process over
|
||||
`SqliteRunStore`, with the outcome read from the run's record through
|
||||
`inspect_run`; `interviewer::Unattended` fails any question until the
|
||||
- `engine`: a run executed by Petri, started from its admitted graphs or
|
||||
resumed from its records, with the outcome read from the run's record
|
||||
through `inspect_run` and mapped to the conclusion Fabro's read side
|
||||
records. The run's worker process runs it over `HttpRunStore`; the server
|
||||
runs it in its own process only under its test override, over
|
||||
`SqliteRunStore`. `interviewer::Unattended` fails any question until the
|
||||
interview adapter lands.
|
||||
- `HttpRunStore`: the same store as a run's worker process reaches it, over
|
||||
the server's `/api/v1/runs/{id}/petri/*` endpoints with the worker's token.
|
||||
The server answers from its `SqliteRunStore`, so the lease and the
|
||||
`(log, seq)` rule are the store's; this layer carries requests, resends a
|
||||
request whose reply was lost, and maps the server's error codes back to
|
||||
`StoreError`. The module docs state the rules.
|
||||
request whose reply was lost, maps the server's error codes back to
|
||||
`StoreError`, and, for a worker, takes every lease for the worker's launch
|
||||
id. The module docs state the rules.
|
||||
- `petri`: the Petri store vocabulary re-exported for the server, which
|
||||
answers the worker endpoints from a `SqliteRunStore` without naming a Petri
|
||||
package in its own manifest.
|
||||
|
|
@ -49,7 +53,12 @@ A run goes to Petri when its workflow version's `workflow.toml` names
|
|||
`engine = "petri"` in `[workflow]`, or when the server's
|
||||
`[server.execution] engine` (`FABRO_SERVER_ENGINE`, `fabro server start
|
||||
--engine`) says so for versions that name none. The server side of both
|
||||
halves is `fabro-server`'s `server::petri_runs`.
|
||||
halves is `fabro-server`'s `server::petri_runs`; the worker side is
|
||||
`fabro-cli`'s `commands::run::petri_worker`, which `fabro run __run-worker`
|
||||
takes when the run's stored spec names Petri. After a server restart, a
|
||||
Petri run left in flight goes back to a worker in `--mode resume`: the run
|
||||
continues from its records, as Petri's own resume does, and full recovery
|
||||
of the workspace to a durable snapshot is the plan's F3.5.
|
||||
|
||||
## How it is tested
|
||||
|
||||
|
|
@ -83,6 +92,15 @@ ulimit -n 4096 && cargo nextest run -p fabro-petri
|
|||
|
||||
The server's end-to-end coverage is `lib/apps/fabro-server/tests/it/scenario/petri.rs`:
|
||||
the `hello` bundle on the OpenAI twin and a command-only bundle run to
|
||||
completion through the create handler and the scheduler, under the version
|
||||
flag and under the server setting, and Petri's diagnostics refuse a run at
|
||||
create.
|
||||
completion through the create handler and the scheduler, in the server
|
||||
process under its test override, under the version flag and under the
|
||||
server setting, and Petri's diagnostics refuse a run at create. The
|
||||
server's `petri_runs` unit tests cover the lease ending at worker exit and
|
||||
the restart reconcile that relaunches a worker in resume mode.
|
||||
|
||||
The worker path is covered with the real binary in
|
||||
`lib/apps/fabro-cli/tests/it/scenario/petri.rs`: a command-only Petri run
|
||||
executes in the worker a foreground server launched, its records reach
|
||||
`petri_records` over the HTTP store and its lease ends with the worker; and
|
||||
a run whose server and worker are both killed mid-stage resumes in a new
|
||||
worker after the server restarts, with one `run.completed`.
|
||||
|
|
|
|||
|
|
@ -7,19 +7,38 @@
|
|||
//! which is the key the coordinator registers the graph under and the name a
|
||||
//! nested-workflow step invokes its child by. Loading verifies the digest,
|
||||
//! so a blob that does not decode to the graph it claims is refused.
|
||||
//!
|
||||
//! The server loads through its [`BlobStore`]; a run's worker loads through
|
||||
//! its client's blob read with [`load_with`], since the run's blobs are the
|
||||
//! blob store the server answers `GET /runs/{id}/blobs/{hash}` from.
|
||||
|
||||
use std::future::Future;
|
||||
|
||||
use fabro_store::BlobStore;
|
||||
use fabro_types::{PetriAdmission, PetriGraphRef};
|
||||
use fabro_types::{BlobHash, PetriAdmission, PetriGraphRef};
|
||||
use petri_runtime::frontend::graph_digest;
|
||||
use petri_runtime::ir::Graph;
|
||||
|
||||
use crate::check::Admitted;
|
||||
|
||||
/// The graphs a run starts from: the admitted root and its pre-lowered
|
||||
/// children, loaded and verified.
|
||||
pub struct AdmittedGraphs {
|
||||
pub graph: Graph,
|
||||
pub children: Vec<Graph>,
|
||||
}
|
||||
|
||||
/// Why an admission could not be stored or loaded.
|
||||
#[derive(Debug, thiserror::Error)]
|
||||
pub enum AdmissionError {
|
||||
#[error("the blob store failed")]
|
||||
Store(#[source] fabro_store::Error),
|
||||
#[error("blob `{blob}` could not be read")]
|
||||
Read {
|
||||
blob: String,
|
||||
#[source]
|
||||
source: anyhow::Error,
|
||||
},
|
||||
#[error("graph `{digest}` is not in the blob store")]
|
||||
Missing { digest: String },
|
||||
#[error("graph `{digest}` does not encode as JSON")]
|
||||
|
|
@ -60,13 +79,30 @@ pub async fn persist(
|
|||
pub async fn load(
|
||||
blobs: &BlobStore,
|
||||
admission: &PetriAdmission,
|
||||
) -> Result<(Graph, Vec<Graph>), AdmissionError> {
|
||||
let graph = load_graph(blobs, &admission.graph).await?;
|
||||
) -> Result<AdmittedGraphs, AdmissionError> {
|
||||
load_with(
|
||||
|blob| async move { blobs.read(&blob).await.map_err(anyhow::Error::from) },
|
||||
admission,
|
||||
)
|
||||
.await
|
||||
}
|
||||
|
||||
/// [`load`] over any blob read: `read` answers a hash with the blob's
|
||||
/// bytes, or `None` when the store lacks it.
|
||||
pub async fn load_with<F, Fut>(
|
||||
read: F,
|
||||
admission: &PetriAdmission,
|
||||
) -> Result<AdmittedGraphs, AdmissionError>
|
||||
where
|
||||
F: Fn(BlobHash) -> Fut,
|
||||
Fut: Future<Output = anyhow::Result<Option<bytes::Bytes>>>,
|
||||
{
|
||||
let graph = load_graph(&read, &admission.graph).await?;
|
||||
let mut children = Vec::with_capacity(admission.children.len());
|
||||
for child in &admission.children {
|
||||
children.push(load_graph(blobs, child).await?);
|
||||
children.push(load_graph(&read, child).await?);
|
||||
}
|
||||
Ok((graph, children))
|
||||
Ok(AdmittedGraphs { graph, children })
|
||||
}
|
||||
|
||||
async fn persist_graph(blobs: &BlobStore, graph: &Graph) -> Result<PetriGraphRef, AdmissionError> {
|
||||
|
|
@ -79,11 +115,17 @@ async fn persist_graph(blobs: &BlobStore, graph: &Graph) -> Result<PetriGraphRef
|
|||
Ok(PetriGraphRef { blob, digest })
|
||||
}
|
||||
|
||||
async fn load_graph(blobs: &BlobStore, graph: &PetriGraphRef) -> Result<Graph, AdmissionError> {
|
||||
let bytes = blobs
|
||||
.read(&graph.blob)
|
||||
async fn load_graph<F, Fut>(read: &F, graph: &PetriGraphRef) -> Result<Graph, AdmissionError>
|
||||
where
|
||||
F: Fn(BlobHash) -> Fut,
|
||||
Fut: Future<Output = anyhow::Result<Option<bytes::Bytes>>>,
|
||||
{
|
||||
let bytes = read(graph.blob)
|
||||
.await
|
||||
.map_err(AdmissionError::Store)?
|
||||
.map_err(|source| AdmissionError::Read {
|
||||
blob: graph.blob.to_string(),
|
||||
source,
|
||||
})?
|
||||
.ok_or_else(|| AdmissionError::Missing {
|
||||
digest: graph.digest.clone(),
|
||||
})?;
|
||||
|
|
|
|||
|
|
@ -1,11 +1,17 @@
|
|||
//! A Fabro run executed by Petri, in the server process.
|
||||
//! A Fabro run executed by Petri: the one assembly the run's worker process
|
||||
//! and the server share.
|
||||
//!
|
||||
//! Until the worker's HTTP run store lands, a Petri run executes where the
|
||||
//! server is: the runtime is assembled the same way the create handler
|
||||
//! assembled it for `Runtime::check`, the run's records go to
|
||||
//! [`SqliteRunStore`] under the Fabro run id as the run key, the admitted
|
||||
//! graphs are loaded from the blob store, and `execution::host::run_configured`
|
||||
//! runs the root invocation to its end. The outcome is then derived from
|
||||
//! The worker is where a Petri run executes, as a legacy run does: it
|
||||
//! reaches the run's record through [`HttpRunStore`](crate::HttpRunStore)
|
||||
//! with its token, and everything else here is the same as in the server.
|
||||
//! The server itself executes a run only under its test override, over
|
||||
//! [`SqliteRunStore`](crate::SqliteRunStore) in its own process. Both build
|
||||
//! the runtime the same way the create handler built it for
|
||||
//! `Runtime::check`, name the Fabro run id as the run key, and hand the run
|
||||
//! to `execution::host`: [`Execution::Start`] runs the admitted graphs
|
||||
//! through `run_configured`; [`Execution::Resume`] continues the run from
|
||||
//! its records through `resume_configured`, with the same observers a start
|
||||
//! installs, as the host's docs require. The outcome is then derived from
|
||||
//! `inspect_run` over a read handle of the same store, so what the caller
|
||||
//! reports is what the durable record says.
|
||||
//!
|
||||
|
|
@ -16,30 +22,46 @@
|
|||
//! rides the caller's token: when it fires, the root invocation is cancelled
|
||||
//! politely and Petri records why.
|
||||
//!
|
||||
//! A resume here is Petri's own: the run continues from its records, and
|
||||
//! sandbox leases are reconciled by label. Full recovery, where the
|
||||
//! workspace a resumed stage sees is restored to the snapshot its durable
|
||||
//! state names, is the integration plan's F3.5 and lands after this.
|
||||
//!
|
||||
//! No stage or agent event is projected into Fabro's tables here; the
|
||||
//! caller appends only the run lifecycle events Fabro's read side needs to
|
||||
//! finish the run. The projection over Petri's records is the read-side
|
||||
//! item that follows.
|
||||
//! finish the run, from the [`Conclusion`] this module derives. The
|
||||
//! projection over Petri's records is the read-side item that follows.
|
||||
|
||||
use std::path::PathBuf;
|
||||
use std::sync::Arc;
|
||||
|
||||
use fabro_store::BlobStore;
|
||||
use fabro_types::{PetriAdmission, SandboxProviderKind};
|
||||
use fabro_types::{FailureReason, SandboxProviderKind};
|
||||
use petri_execution::host::{self, HostError, HostRun};
|
||||
use petri_execution::inspect::{self, InspectError, RunInspection};
|
||||
use petri_execution::{Access, CancelReason, InterviewDispatcher, RECEIPT_FILE, RunKey, RunStore};
|
||||
use petri_execution::{
|
||||
Access, CancelReason, InterviewDispatcher, InvocationId, RECEIPT_FILE, RunKey, RunStore,
|
||||
};
|
||||
use petri_runtime::executor::Retention;
|
||||
use petri_runtime::{RunOptions, SandboxBackend};
|
||||
use tokio::fs;
|
||||
use tokio_util::sync::CancellationToken;
|
||||
use tracing::{debug, info, warn};
|
||||
|
||||
use crate::admission::{self, AdmissionError};
|
||||
use crate::admission::AdmittedGraphs;
|
||||
use crate::interviewer::Unattended;
|
||||
use crate::run_store::SqliteRunStore;
|
||||
use crate::runtime::RuntimeSpec;
|
||||
|
||||
/// How the run is entered: fresh, from the admitted graphs, or continued
|
||||
/// from its records.
|
||||
pub enum Execution {
|
||||
/// Run the admitted graphs from the start; the run must not exist in
|
||||
/// the store yet.
|
||||
Start(AdmittedGraphs),
|
||||
/// Continue the run from its records; the run must exist in the store
|
||||
/// with its root invocation declared.
|
||||
Resume,
|
||||
}
|
||||
|
||||
/// One run to execute.
|
||||
pub struct RunRequest {
|
||||
/// The Fabro run id, which becomes Petri's run key: the run's identity
|
||||
|
|
@ -47,12 +69,10 @@ pub struct RunRequest {
|
|||
pub run_id: String,
|
||||
/// Where the run's workspaces, step output and blobs live.
|
||||
pub run_dir: PathBuf,
|
||||
/// What the create handler admitted.
|
||||
pub admission: PetriAdmission,
|
||||
/// The blob store the admitted graphs are read from.
|
||||
pub blobs: Arc<BlobStore>,
|
||||
/// The run's durable record.
|
||||
pub store: Arc<SqliteRunStore>,
|
||||
pub execution: Execution,
|
||||
/// The run's durable record: the worker's HTTP store, or the server's
|
||||
/// SQLite store under the test override.
|
||||
pub store: Arc<dyn RunStore>,
|
||||
pub runtime: RuntimeSpec,
|
||||
/// The sandbox provider Fabro resolved for the run's environment.
|
||||
pub provider: SandboxProviderKind,
|
||||
|
|
@ -86,32 +106,47 @@ pub struct RunOutcome {
|
|||
pub enum RunError {
|
||||
#[error("the run's sandbox provider `{provider}` is not one Petri serves")]
|
||||
UnsupportedProvider { provider: SandboxProviderKind },
|
||||
#[error("the admitted graphs could not be loaded")]
|
||||
Admission(#[from] AdmissionError),
|
||||
#[error("the run's record could not be opened")]
|
||||
Open(#[source] petri_store::StoreError),
|
||||
#[error("the run's record could not be read")]
|
||||
Read(#[source] HostError),
|
||||
#[error("the run's record has no root invocation, so there is nothing to resume")]
|
||||
NothingToResume,
|
||||
#[error("the run's record could not be inspected")]
|
||||
Inspect(#[source] InspectError),
|
||||
#[error("the run ended without recording a status; the record says: {}", .0.join("; "))]
|
||||
Unfinished(Vec<String>),
|
||||
}
|
||||
|
||||
/// How Fabro reports the run: what its read side records as the run's
|
||||
/// terminal event.
|
||||
#[derive(Clone, Debug, PartialEq, Eq)]
|
||||
pub enum Conclusion {
|
||||
/// The record says the run succeeded and is whole.
|
||||
Succeeded,
|
||||
/// Anything else: the record says the run failed or was cancelled, the
|
||||
/// record is incomplete, or the run could not be executed at all.
|
||||
Failed {
|
||||
reason: FailureReason,
|
||||
message: String,
|
||||
},
|
||||
}
|
||||
|
||||
/// Execute the run to its end and report what the record says.
|
||||
pub async fn run(request: RunRequest) -> Result<RunOutcome, RunError> {
|
||||
let backend = backend(&request.provider)?;
|
||||
let (graph, children) = admission::load(&request.blobs, &request.admission).await?;
|
||||
let key = RunKey::new(request.run_id.as_str());
|
||||
let mut options = RunOptions::new(&request.run_dir);
|
||||
options.run_key = Some(key.clone());
|
||||
options.retention = Retention::Always;
|
||||
options.sandbox.backend = backend;
|
||||
let store: Arc<dyn RunStore> = request.store.clone();
|
||||
let runtime = request.runtime.runtime(true).store(store).options(options);
|
||||
let runtime = request
|
||||
.runtime
|
||||
.runtime(true)
|
||||
.store(Arc::clone(&request.store))
|
||||
.options(options);
|
||||
|
||||
let dispatcher = InterviewDispatcher::new(Arc::new(Unattended));
|
||||
let host_run = HostRun::new(graph)
|
||||
.with_children(children)
|
||||
.observe(Arc::new(dispatcher.clone()));
|
||||
let cancel = request.cancel.clone();
|
||||
let mut cancel_task = None;
|
||||
let with_handle = |handle: petri_execution::CoordinatorHandle, secrets| {
|
||||
|
|
@ -122,8 +157,26 @@ pub async fn run(request: RunRequest) -> Result<RunOutcome, RunError> {
|
|||
handle.cancel_root_for(CancelReason::Control);
|
||||
}));
|
||||
};
|
||||
info!(run_id = %request.run_id, backend = %backend, "Starting Petri run");
|
||||
let result = Box::pin(host::run_configured(&runtime, host_run, with_handle)).await;
|
||||
let result = match request.execution {
|
||||
Execution::Start(graphs) => {
|
||||
info!(run_id = %request.run_id, backend = %backend, "Starting Petri run");
|
||||
let host_run = HostRun::new(graphs.graph)
|
||||
.with_children(graphs.children)
|
||||
.observe(Arc::new(dispatcher.clone()));
|
||||
Box::pin(host::run_configured(&runtime, host_run, with_handle)).await
|
||||
}
|
||||
Execution::Resume => {
|
||||
check_resumable(request.store.as_ref(), &key).await?;
|
||||
info!(run_id = %request.run_id, backend = %backend, "Resuming Petri run");
|
||||
Box::pin(host::resume_configured(
|
||||
&runtime,
|
||||
Vec::new(),
|
||||
vec![Arc::new(dispatcher.clone())],
|
||||
with_handle,
|
||||
))
|
||||
.await
|
||||
}
|
||||
};
|
||||
if let Some(task) = cancel_task {
|
||||
task.abort();
|
||||
}
|
||||
|
|
@ -133,18 +186,74 @@ pub async fn run(request: RunRequest) -> Result<RunOutcome, RunError> {
|
|||
Ok(report) => debug!(status = %report.status, "Petri run ended"),
|
||||
Err(error) => warn!(error = %error, "Petri run ended with a host error"),
|
||||
}
|
||||
let inspection = inspect(&request.store, &key).await?;
|
||||
let inspection = inspect(request.store.as_ref(), &key).await?;
|
||||
outcome(inspection, result.err())
|
||||
}
|
||||
|
||||
/// What the run's record says, read through a handle that holds no lease:
|
||||
/// the same derivation [`run`] ends with, for a caller that only holds the
|
||||
/// store, such as a test checking a finished run.
|
||||
pub async fn outcome_of(store: &SqliteRunStore, run_id: &str) -> Result<RunOutcome, RunError> {
|
||||
pub async fn outcome_of(store: &dyn RunStore, run_id: &str) -> Result<RunOutcome, RunError> {
|
||||
let inspection = inspect(store, &RunKey::new(run_id)).await?;
|
||||
outcome(inspection, None)
|
||||
}
|
||||
|
||||
/// How Fabro reports what [`run`] returned. A cancelled run is a failure
|
||||
/// with the cancelled reason, as the legacy executor reports one; every
|
||||
/// other shortfall is a workflow error whose message says what the record,
|
||||
/// or the host, said.
|
||||
#[must_use]
|
||||
pub fn conclusion(result: &Result<RunOutcome, RunError>) -> Conclusion {
|
||||
match result {
|
||||
Ok(RunOutcome {
|
||||
status: RunStatus::Success,
|
||||
complete: true,
|
||||
..
|
||||
}) => Conclusion::Succeeded,
|
||||
Ok(outcome) => {
|
||||
let reason = match outcome.status {
|
||||
RunStatus::Cancelled => FailureReason::Cancelled,
|
||||
RunStatus::Success | RunStatus::Failed => FailureReason::WorkflowError,
|
||||
};
|
||||
Conclusion::Failed {
|
||||
reason,
|
||||
message: failure_message(outcome),
|
||||
}
|
||||
}
|
||||
Err(error) => Conclusion::Failed {
|
||||
reason: FailureReason::WorkflowError,
|
||||
message: error_chain(error),
|
||||
},
|
||||
}
|
||||
}
|
||||
|
||||
/// The failure of a run whose record says it did not succeed.
|
||||
fn failure_message(outcome: &RunOutcome) -> String {
|
||||
let mut message = match (&outcome.status, &outcome.failure) {
|
||||
(RunStatus::Cancelled, _) => "the run was cancelled".to_string(),
|
||||
(_, Some(failure)) => failure.clone(),
|
||||
(RunStatus::Failed, None) => "the run failed".to_string(),
|
||||
(RunStatus::Success, None) => "the run's record is incomplete".to_string(),
|
||||
};
|
||||
if !outcome.complete {
|
||||
message.push_str(" (record incomplete: ");
|
||||
message.push_str(&outcome.incomplete.join("; "));
|
||||
message.push(')');
|
||||
}
|
||||
message
|
||||
}
|
||||
|
||||
/// The error and every cause under it, as one line.
|
||||
fn error_chain(error: &RunError) -> String {
|
||||
let mut parts = vec![error.to_string()];
|
||||
let mut cause = std::error::Error::source(error);
|
||||
while let Some(next) = cause {
|
||||
parts.push(next.to_string());
|
||||
cause = next.source();
|
||||
}
|
||||
parts.join(": ")
|
||||
}
|
||||
|
||||
/// The sandbox backend for Fabro's provider kind.
|
||||
fn backend(provider: &SandboxProviderKind) -> Result<SandboxBackend, RunError> {
|
||||
if *provider == SandboxProviderKind::LOCAL {
|
||||
|
|
@ -160,8 +269,25 @@ fn backend(provider: &SandboxProviderKind) -> Result<SandboxBackend, RunError> {
|
|||
}
|
||||
}
|
||||
|
||||
/// Refuse a resume the host would not survive: `resume_configured` indexes
|
||||
/// the root invocation of the stored state, so a record with none (the run
|
||||
/// was created in the store and nothing more) is refused here with a named
|
||||
/// error instead.
|
||||
async fn check_resumable(store: &dyn RunStore, key: &RunKey) -> Result<(), RunError> {
|
||||
let logs = store
|
||||
.open(key, Access::Read)
|
||||
.await
|
||||
.map_err(RunError::Open)?;
|
||||
let state = host::stored_state(&*logs).await.map_err(RunError::Read)?;
|
||||
if state.invocations.contains_key(&InvocationId::ROOT) {
|
||||
Ok(())
|
||||
} else {
|
||||
Err(RunError::NothingToResume)
|
||||
}
|
||||
}
|
||||
|
||||
/// Read the run back through a handle that holds no lease.
|
||||
async fn inspect(store: &SqliteRunStore, key: &RunKey) -> Result<RunInspection, RunError> {
|
||||
async fn inspect(store: &dyn RunStore, key: &RunKey) -> Result<RunInspection, RunError> {
|
||||
let logs = store
|
||||
.open(key, Access::Read)
|
||||
.await
|
||||
|
|
@ -224,3 +350,66 @@ async fn write_receipt(run_dir: &std::path::Path, receipt: &petri_execution::Int
|
|||
warn!(path = %path.display(), error = %error, "could not write the interview receipt");
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
|
||||
fn outcome_with(status: RunStatus, failure: Option<&str>, complete: bool) -> RunOutcome {
|
||||
RunOutcome {
|
||||
status,
|
||||
failure: failure.map(ToOwned::to_owned),
|
||||
complete,
|
||||
incomplete: if complete {
|
||||
Vec::new()
|
||||
} else {
|
||||
vec!["execution 0 did not finish".to_string()]
|
||||
},
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn a_whole_successful_record_concludes_succeeded() {
|
||||
assert_eq!(
|
||||
conclusion(&Ok(outcome_with(RunStatus::Success, None, true))),
|
||||
Conclusion::Succeeded
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn a_cancelled_record_concludes_cancelled() {
|
||||
assert_eq!(
|
||||
conclusion(&Ok(outcome_with(RunStatus::Cancelled, None, true))),
|
||||
Conclusion::Failed {
|
||||
reason: FailureReason::Cancelled,
|
||||
message: "the run was cancelled".to_string(),
|
||||
}
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn a_failed_record_carries_the_root_failure_and_the_incomplete_reasons() {
|
||||
assert_eq!(
|
||||
conclusion(&Ok(outcome_with(
|
||||
RunStatus::Failed,
|
||||
Some("step `say` failed"),
|
||||
false
|
||||
))),
|
||||
Conclusion::Failed {
|
||||
reason: FailureReason::WorkflowError,
|
||||
message: "step `say` failed (record incomplete: execution 0 did not finish)"
|
||||
.to_string(),
|
||||
}
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn a_host_error_concludes_with_its_chain() {
|
||||
let error = RunError::Unfinished(vec!["no status".to_string()]);
|
||||
assert_eq!(conclusion(&Err(error)), Conclusion::Failed {
|
||||
reason: FailureReason::WorkflowError,
|
||||
message: "the run ended without recording a status; the record says: no status"
|
||||
.to_string(),
|
||||
});
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -26,6 +26,16 @@
|
|||
//! runtime at drop, the server's worker-exit release is the backstop, and
|
||||
//! the drop says so in the log.
|
||||
//!
|
||||
//! # The owner
|
||||
//!
|
||||
//! A store built with [`HttpRunStore::for_worker`] names one owner for the
|
||||
//! whole process: every `Create` and `Write` takes the lease for the
|
||||
//! worker's launch id, whatever owner Petri minted for the run runtime that
|
||||
//! asked. One worker process executes one run, so the lease is the
|
||||
//! launch's, the worker logs it once at start, and the server's lease row
|
||||
//! names the launch that holds it. A store built with [`HttpRunStore::new`]
|
||||
//! passes Petri's owner through unchanged.
|
||||
//!
|
||||
//! # Lost replies
|
||||
//!
|
||||
//! Every call is one request. A reply that never arrives (a transport error,
|
||||
|
|
@ -92,6 +102,9 @@ impl fmt::Debug for HttpRunStore {
|
|||
/// What the store and every handle it opens share.
|
||||
struct Shared {
|
||||
client: Client,
|
||||
/// The owner every writer open takes the lease for, when the store is
|
||||
/// a worker's; `None` passes Petri's owner through.
|
||||
owner: Option<OwnerId>,
|
||||
/// The writer handle alive in this process per run and owner, so a
|
||||
/// same-owner reopen shares it and the lease lasts while any handle
|
||||
/// does.
|
||||
|
|
@ -101,12 +114,25 @@ struct Shared {
|
|||
}
|
||||
|
||||
impl HttpRunStore {
|
||||
/// A store over a client that carries the worker's token.
|
||||
/// A store over a client that carries the worker's token, taking each
|
||||
/// lease for the owner Petri names.
|
||||
#[must_use]
|
||||
pub fn new(client: Client) -> Self {
|
||||
Self::build(client, None)
|
||||
}
|
||||
|
||||
/// A worker's store: every lease is taken for `owner`, the worker's
|
||||
/// launch id, whatever owner Petri names.
|
||||
#[must_use]
|
||||
pub fn for_worker(client: Client, owner: OwnerId) -> Self {
|
||||
Self::build(client, Some(owner))
|
||||
}
|
||||
|
||||
fn build(client: Client, owner: Option<OwnerId>) -> Self {
|
||||
Self {
|
||||
shared: Arc::new(Shared {
|
||||
client,
|
||||
owner,
|
||||
live: Mutex::default(),
|
||||
releases: Mutex::default(),
|
||||
}),
|
||||
|
|
@ -290,6 +316,9 @@ impl RunStore for HttpRunStore {
|
|||
Access::Write { owner } => (PetriAccess::Write, Some(owner)),
|
||||
Access::Read => (PetriAccess::Read, None),
|
||||
};
|
||||
// A worker's store leases for its launch, not for the owner Petri
|
||||
// minted for this run runtime.
|
||||
let owner = owner.map(|named| shared.owner.as_ref().unwrap_or(named));
|
||||
let request = PetriOpenRequest {
|
||||
access: api_access,
|
||||
owner: owner.map(|owner| owner.as_str().to_string()),
|
||||
|
|
@ -326,7 +355,12 @@ impl RunStore for HttpRunStore {
|
|||
};
|
||||
let opened =
|
||||
opened.map_err(|error| shared.store_error(key, "open the run", None, error))?;
|
||||
debug!(run_id = %key, access = ?request.access, "Petri run opened over the API");
|
||||
debug!(
|
||||
run_id = %key,
|
||||
access = ?request.access,
|
||||
owner = owner.map(OwnerId::as_str),
|
||||
"Petri run opened over the API"
|
||||
);
|
||||
match owner {
|
||||
Some(owner) => Ok(self.writer(key, run_id, owner.clone(), opened.locator)),
|
||||
None => Ok(Arc::new(HttpRunLogs {
|
||||
|
|
|
|||
|
|
@ -16,11 +16,13 @@
|
|||
//! its diagnostics come back in a shape Fabro maps onto its own;
|
||||
//! - [`admission`]: the admitted graphs in Fabro's blob store, named on the run
|
||||
//! spec;
|
||||
//! - [`engine`]: a run executed by Petri in the server process, with the
|
||||
//! outcome read from its record;
|
||||
//! - [`engine`]: a run executed by Petri, started or resumed, in the run's
|
||||
//! worker process over the HTTP store (or in the server process under its
|
||||
//! test override), with the outcome read from its record;
|
||||
//! - [`interviewer`]: the interviewer of a run nobody is watching;
|
||||
//! - [`HttpRunStore`]: the same store as a run's worker process reaches it,
|
||||
//! over the server's API with the worker's token;
|
||||
//! over the server's API with the worker's token and its launch id as the
|
||||
//! lease owner;
|
||||
//! - the platform adapters still to come: hooks, interviews over Fabro's API,
|
||||
//! secrets, output storage, the run tools, the event projection.
|
||||
//!
|
||||
|
|
|
|||
|
|
@ -118,11 +118,11 @@ async fn the_hello_bundle_is_admitted_and_round_trips_through_the_blob_store() {
|
|||
.await
|
||||
.expect("the graphs persist");
|
||||
assert!(record.children.is_empty());
|
||||
let (graph, children) = admission::load(&blobs, &record)
|
||||
let graphs = admission::load(&blobs, &record)
|
||||
.await
|
||||
.expect("the graphs load");
|
||||
assert_eq!(graph, admitted.graph);
|
||||
assert!(children.is_empty());
|
||||
assert_eq!(graphs.graph, admitted.graph);
|
||||
assert!(graphs.children.is_empty());
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue