diff --git a/lib/apps/fabro-cli/tests/it/scenario/petri.rs b/lib/apps/fabro-cli/tests/it/scenario/petri.rs index a3ea4903c..d6f14a141 100644 --- a/lib/apps/fabro-cli/tests/it/scenario/petri.rs +++ b/lib/apps/fabro-cli/tests/it/scenario/petri.rs @@ -604,13 +604,21 @@ async fn wait_for_questions( } } -/// Answer a question through the API, as the web app and the CLI do. +/// Answer a question through the API, as the web app and the CLI do. The +/// question id is Petri's (`gate#2`), so it travels as one percent-encoded +/// path segment, as the generated clients send it. async fn answer(server: &RunningServer, run_id: &str, question_id: &str, body: serde_json::Value) { + let mut url = fabro_http::Url::parse(&format!( + "{}/api/v1/runs/{run_id}/questions", + server.api_base_url + )) + .expect("the API base URL parses"); + url.path_segments_mut() + .expect("the API URL has a path") + .push(question_id) + .push("answer"); let response = fabro_test::test_http_client() - .post(format!( - "{}/api/v1/runs/{run_id}/questions/{question_id}/answer", - server.api_base_url - )) + .post(url) .bearer_auth(TEST_DEV_TOKEN) .json(&body) .send() @@ -672,9 +680,13 @@ async fn a_human_gate_in_the_worker_is_answered_through_the_api() { let pending = wait_for_questions(&server, &run_id, 1).await; let question = &pending[0]; - assert_eq!(question["stage"], "gate", "{question}"); + assert_eq!(question["stage"], "gate@1", "{question}"); assert_eq!(question["question_type"], "yes_no", "{question}"); let question_id = question["id"].as_str().expect("an id").to_string(); + assert!( + question_id.starts_with("gate#"), + "Petri's id: {question_id}" + ); answer( &server, &run_id, @@ -732,7 +744,7 @@ async fn two_parallel_gates_in_the_worker_each_bind_their_own_answer() { .unwrap_or_else(|| panic!("`{stage}` is pending: {pending:?}")) .to_string() }; - let (a, b) = (id_of("a"), id_of("b")); + let (a, b) = (id_of("a@1"), id_of("b@1")); assert_ne!(a, b); for question in &pending { assert_eq!(question["question_type"], "yes_no", "{question}"); diff --git a/lib/apps/fabro-server/tests/it/scenario/petri.rs b/lib/apps/fabro-server/tests/it/scenario/petri.rs index f035fb540..3816fbde9 100644 --- a/lib/apps/fabro-server/tests/it/scenario/petri.rs +++ b/lib/apps/fabro-server/tests/it/scenario/petri.rs @@ -34,6 +34,7 @@ use fabro_server::test_support::{ test_register_workflow_version, }; use fabro_static::EnvVars; +use fabro_store::platform_records::{PlatformRecord, PlatformRecordKind, PlatformRecordStore}; use fabro_test::{TwinScenario, TwinScenarios, twin_openai}; use fabro_types::{RunId, WorkflowPath, WorkflowVersion}; use tower::ServiceExt; @@ -591,7 +592,7 @@ async fn a_human_gate_is_answered_through_the_questions_api() { create_and_start_run_from_intent(&app, intent(&version_id, workspace.path())).await; let question = wait_for_question(&app, &run_id).await; - assert_eq!(question["stage"], "gate", "{question}"); + assert_eq!(question["stage"], "gate@1", "{question}"); assert_eq!(question["text"], "Go?", "{question}"); assert_eq!(question["question_type"], "yes_no", "{question}"); let keys: Vec<&str> = question["options"] @@ -602,12 +603,20 @@ async fn a_human_gate_is_answered_through_the_questions_api() { .collect(); assert_eq!(keys, vec!["Y", "N"], "{question}"); let question_id = question["id"].as_str().expect("an id").to_string(); - assert!(question_id.starts_with("gate."), "{question_id}"); + assert!( + question_id.starts_with("gate#"), + "Petri's id: {question_id}" + ); + // Petri's id travels as one percent-encoded path segment, as the + // generated clients send it. + let encoded_id = + percent_encoding::utf8_percent_encode(&question_id, percent_encoding::NON_ALPHANUMERIC) + .to_string(); let req = Request::builder() .method("POST") .uri(api(&format!( - "/runs/{run_id}/questions/{question_id}/answer" + "/runs/{run_id}/questions/{encoded_id}/answer" ))) .header("content-type", "application/json") .body(Body::from(r#"{"kind":"no"}"#)) @@ -656,4 +665,24 @@ async fn a_human_gate_is_answered_through_the_questions_api() { "the answered question is no longer pending: {}", state_body["pending_interviews"] ); + // Who answered is a platform record keyed on Petri's id, derived from + // the adapter's `interview.completed` with the API caller as its actor. + let answered = PlatformRecordStore::new(state.test_petri_view_pool()) + .read_kind( + &run_id.parse().expect("the run id parses"), + PlatformRecordKind::InterviewAnswered, + ) + .await + .expect("the platform records read"); + let [answered] = answered.as_slice() else { + panic!("one question was answered: {answered:?}"); + }; + let PlatformRecord::InterviewAnswered(record) = &answered.record else { + panic!("an answered record: {answered:?}"); + }; + assert_eq!(record.question, question_id); + assert!( + record.principal.is_some(), + "the answering principal: {record:?}" + ); } diff --git a/lib/components/fabro-petri/README.md b/lib/components/fabro-petri/README.md index f257e73a6..451686873 100644 --- a/lib/components/fabro-petri/README.md +++ b/lib/components/fabro-petri/README.md @@ -39,17 +39,20 @@ Every adapter the integration plan describes lands here. `SqliteRunStore`. The caller supplies the interviewer, and the secret provider and blob table when it has them. - `interview`: Petri's `Interviewer` over Fabro's questions API and the - worker's control channel. A human gate's question is posted as the - `interview.started` event a legacy `human` stage emits (through the - worker's run event sink, or the run's database in the server process), so - `GET /runs/{id}/questions`, the web app and Slack list it; the answer - posted to `/questions/{qid}/answer` reaches the worker's control - interviewer over the control bus (or the in-process one directly) under - the same id, and is mapped onto Petri's answer. The question id is - derived from Petri's identity (node, execution, firing, occurrence, ask). + worker's control channel. A question has one id in Fabro, Petri's own + (`gate#2`): the projection lists it pending from the `question` record, + `GET /runs/{id}/questions` serves it, and the answer posted to + `/questions/{qid}/answer` is validated against that pending record and + reaches the worker's control interviewer over the control bus (or the + in-process one directly) under the same id, mapped onto Petri's answer. + The adapter still posts the legacy `interview.*` events (through the + worker's run event sink, or the run's database in the server process) + with that id and the projection's stage label, for the readers that + follow the event stream rather than the projection: Slack, `run attach` + and the web app's Q&A renderer. The store derives the `interview.answered` + platform record, with the answering principal, from `interview.completed`. An expired or cancelled question is completed as `interview.timeout` or - `interview.interrupted`; an auto-approved run answers itself. The module - docs mark the hook points the read side takes over. + `interview.interrupted`; an auto-approved run answers itself. - `secrets`: Petri's `SecretProvider` over the vault's token entries, so a `{{ secrets.NAME }}` reference resolves at spawn into a command's environment and is masked in every record; a sensitive answer registers diff --git a/lib/components/fabro-petri/src/interview.rs b/lib/components/fabro-petri/src/interview.rs index fb62019d6..dd74d3e2a 100644 --- a/lib/components/fabro-petri/src/interview.rs +++ b/lib/components/fabro-petri/src/interview.rs @@ -8,20 +8,45 @@ //! question to Fabro the way a legacy `human` stage does, waits for the //! answer the way the legacy worker does, and hands Petri the reply. //! +//! # One identity +//! +//! A question has one id in Fabro: Petri's [`Question::id`] (`gate#2`), +//! as the `question` record names it. The projection over Petri's records +//! keys `pending_interviews` by it, `GET /runs/{id}/questions` lists it, +//! `POST /runs/{id}/questions/{qid}/answer` validates the answer against +//! the projection's pending question under it, and this adapter waits on +//! the worker's [`ControlInterviewer`] under it. The rest of Petri's +//! identity rides on [`AskedQuestion::identity`] for the record of who +//! answered. The stage a question names is the label the projection gives +//! the asking firing (`gate@1`, or `gate/e3@1` when another execution took +//! that label): [`Observed`] labels every firing through the projection's +//! own rule ([`projection::stage_label`]) as the run's records go by. +//! //! # How a question reaches a person //! -//! The legacy stage emits `interview.started` on the run's event stream; -//! the read side keeps it in the projection's `pending_interviews`, keyed -//! by question id, and that is what `GET /runs/{id}/questions`, the web -//! app's interview dock and the Slack integration read pending questions -//! from. This adapter posts the same event through a [`QuestionSink`]: the -//! worker's [`EventSinkQuestions`] appends it over the run event sink the -//! worker already carries lifecycle events on, and the server's in-process -//! path appends it through [`DatabaseQuestions`]. The question id is -//! Fabro's key for the question and is derived from Petri's identity -//! ([`question_id`]); the node name is the event's `stage`, and the Fabro -//! question type, options, freeform flag, deadline and review target are -//! mapped from Petri's [`Question`]. +//! The projection derives the pending question, its answer and its expiry +//! from Petri's records alone; the run's own record is the source of truth +//! and nothing the adapter posts is folded into it. The adapter still +//! posts the legacy `interview.*` events through a [`QuestionSink`], with +//! Petri's id and the projection's stage label, for the readers that +//! follow the run's event stream rather than its projection: +//! +//! - the server's Slack service posts a question to the channel on +//! `interview.started` and finishes it on `interview.completed`, +//! `interview.timeout` or `interview.interrupted`; +//! - `fabro run attach` polls the questions API when `interview.started` +//! arrives and stops waiting on the question's closing event; +//! - the web app's human Q&A renderer pairs `interview.started` with its +//! closing event by question id in the stage's event list; +//! - the server clears its record of an accepted answer on the closing event, +//! so the transport can be claimed again. +//! +//! The worker's [`EventSinkQuestions`] appends them over the run event sink +//! the worker already carries lifecycle events on, and the server's +//! in-process path appends them through [`DatabaseQuestions`]. For a Petri +//! run the store derives the `interview.answered` platform record from +//! `interview.completed`: the question's Petri id and the principal that +//! answered, which is the actor the adapter stamps on the event. //! //! # How the answer comes back //! @@ -45,16 +70,17 @@ //! adapter's cancel token, as it does when the firing ends without an //! answer or the run is cancelled. The adapter returns promptly with //! [`InterviewReply::Cancelled`] and posts `interview.timeout` when the -//! gate reported the expiry, else `interview.interrupted`, so the pending -//! question clears from Fabro's view. The expiry report is seen by the -//! adapter's own observer ([`FabroInterviewer::observer`]), which the run -//! registers ahead of the dispatcher so the report is noted before the -//! token fires. The dispatcher races the reply against the same token and -//! may drop the reply future the moment the token fires, so the notice is -//! posted from a guard that runs whether the future completes or is -//! dropped, on a task of its own. The dispatcher's own record of the -//! outcome (`TimedOut` with the default taken, `Cancelled`, `Late`) is the -//! authoritative one and reaches the receipt. +//! gate reported the expiry, else `interview.interrupted`, so the readers +//! above see the question end. The expiry report is seen by the adapter's +//! own observer ([`FabroInterviewer::observer`]), which the run registers +//! ahead of the dispatcher so the report is noted before the token fires. +//! The dispatcher races the reply against the same token and may drop the +//! reply future the moment the token fires, so the notice is posted from a +//! guard that runs whether the future completes or is dropped, on a task +//! of its own. The dispatcher's own record of the outcome (`TimedOut` with +//! the default taken, `Cancelled`, `Late`) is the authoritative one and +//! reaches the receipt, and the projection closes the question on Petri's +//! `question_expired` record or the cancelled attempt. //! //! # Auto-approval //! @@ -62,21 +88,9 @@ //! at once as the legacy runner's auto-approve interviewer does (`yes`, //! the first option, or `auto-approved` text), attributed to the engine. //! The question is still posted and completed, so the run's stream shows -//! what was decided. -//! -//! # Hook points for the read side -//! -//! The events posted here are the interim bridge to Fabro's read side. -//! Once the projection over Petri's records derives pending questions from -//! the `question` and `question_expired` records and the delivered answer, -//! the sink can become a no-op: the [`QuestionSink`] is the one seam to -//! replace. Two Fabro facts a Petri record does not carry are marked in -//! [`FabroInterviewer::reply`]: who answered (`AnswerSubmission::actor`, -//! carried on `interview.completed` for now) and the Fabro question id -//! that Petri's identity was mapped to. Both belong in a platform record -//! keyed on the same identity when that record kind exists. +//! what was decided, and the projection closes it on the delivered answer. -use std::collections::HashSet; +use std::collections::{BTreeSet, HashMap, HashSet}; use std::sync::{Arc, Mutex, PoisonError}; use std::time::{Duration, Instant}; @@ -86,20 +100,24 @@ use fabro_interview::{ }; use fabro_store::RunDatabase; use fabro_types::{ - InterviewOption, Principal, QuestionType, ReviewTarget, ReviewTargetKind, RunId, + InterviewOption, Principal, QuestionType, ReviewTarget, ReviewTargetKind, RunId, StageId, SystemActorKind, }; use fabro_workflow::event::{self as workflow_event, Event, RunEventSink}; +use petri_execution::events::{Parsed, Projection, ViewEvent}; use petri_execution::{ CoordinatorRecord, ExecutionId, ExecutionObserver, InterviewError, InterviewReply, InterviewRequest, Interviewer, }; -use petri_runtime::engine::{EngineState, Event as EngineEvent, EventRecord}; -use petri_runtime::steps::{Answer, Question, QuestionExpired, QuestionOption}; +use petri_runtime::engine::{EngineState, EventRecord}; +use petri_runtime::ir::FiringId; +use petri_runtime::steps::{Answer, Question, QuestionOption}; use tokio::runtime::Handle; use tokio_util::sync::CancellationToken; use tracing::{debug, warn}; +use crate::projection; + /// Whether a run answers its own questions. #[derive(Clone, Copy, Debug, PartialEq, Eq)] pub enum Approval { @@ -109,7 +127,8 @@ pub enum Approval { Auto, } -/// Petri's identity for one question, as the read side keys it. +/// Petri's full identity for one question, beyond its id: where in the run +/// it was asked, for the record of who answered it. #[derive(Clone, Debug, PartialEq, Eq)] pub struct QuestionIdentity { pub invocation_path: String, @@ -138,9 +157,11 @@ impl QuestionIdentity { /// A question as Fabro shows it: the fields of `interview.started`. #[derive(Clone, Debug, PartialEq)] pub struct AskedQuestion { + /// Petri's id for the question, the one id Fabro knows it by. pub question_id: String, pub identity: QuestionIdentity, pub text: String, + /// The label the projection gives the asking firing. pub stage: String, pub question_type: QuestionType, pub options: Vec, @@ -294,42 +315,95 @@ impl QuestionSink for DatabaseQuestions { } } -/// The questions whose expiry the gate reported, by execution and Petri -/// question id: an observer the run registers ahead of the dispatcher. +/// What the adapter learns from the run's records ahead of the dispatcher: +/// the label the projection gives each firing, and the questions whose +/// expiry the gate reported. An observer the run registers ahead of the +/// dispatcher, fed the same records the projection folds, derived through +/// Petri's own [`Projection`] so a firing's visit and label come out as the +/// read side computes them. #[derive(Default)] -pub struct Expiries { - expired: Mutex>, +pub struct Observed { + state: Mutex, } -impl Expiries { - fn contains(&self, execution: ExecutionId, question: &str) -> bool { - self.expired +#[derive(Default)] +struct ObservedState { + projection: Projection, + /// Every label given so far, for the projection's collision rule. + labels: BTreeSet, + /// The label of each shown firing. + stages: HashMap<(ExecutionId, FiringId), StageId>, + expired: HashSet<(ExecutionId, String)>, +} + +impl Observed { + /// The label the projection gives `firing`, once its `visit.started` + /// was seen. + fn label(&self, execution: ExecutionId, firing: FiringId) -> Option { + self.state .lock() .unwrap_or_else(PoisonError::into_inner) + .stages + .get(&(execution, firing)) + .cloned() + } + + fn expired(&self, execution: ExecutionId, question: &str) -> bool { + self.state + .lock() + .unwrap_or_else(PoisonError::into_inner) + .expired .contains(&(execution, question.to_string())) } } -impl ExecutionObserver for Expiries { +impl ExecutionObserver for Observed { fn on_engine_record( &self, execution: ExecutionId, record: &EventRecord, - _recorded_at: u64, - _state: &EngineState, + recorded_at: u64, + state: &EngineState, ) { - let EngineEvent::StepProgressRecorded { ev, .. } = &record.event else { - return; - }; - if let Some(expired) = QuestionExpired::from_event(ev) { - self.expired - .lock() - .unwrap_or_else(PoisonError::into_inner) - .insert((execution, expired.question)); + let mut observed = self.state.lock().unwrap_or_else(PoisonError::into_inner); + let events = observed + .projection + .engine(execution, record, recorded_at, state); + for event in &events { + if let Some(ViewEvent::VisitStarted { .. }) = event.view() { + let Some(subject) = event.subject.as_ref() else { + continue; + }; + let Some(firing) = subject.firing else { + continue; + }; + if !projection::is_shown(&subject.node) { + continue; + } + let label = projection::stage_label( + &subject.node.name, + projection::visit_of(subject), + execution, + &observed.labels, + ); + observed.labels.insert(label.to_string()); + observed.stages.insert((execution, firing), label); + } + if let Some(Parsed::QuestionExpired { expired }) = event.parsed() { + observed + .expired + .insert((execution, expired.question.clone())); + } } } - fn on_lifecycle(&self, _record: &CoordinatorRecord) {} + fn on_lifecycle(&self, record: &CoordinatorRecord) { + self.state + .lock() + .unwrap_or_else(PoisonError::into_inner) + .projection + .lifecycle(record); + } } /// The interviewer a Fabro run installs. @@ -337,7 +411,7 @@ pub struct FabroInterviewer { answers: Arc, sink: Arc, approval: Approval, - expiries: Arc, + observed: Arc, } impl FabroInterviewer { @@ -353,17 +427,18 @@ impl FabroInterviewer { answers, sink, approval, - expiries: Arc::new(Expiries::default()), + observed: Arc::new(Observed::default()), } } - /// The observer that sees a gate report a question's expiry. A run - /// registers it ahead of the interview dispatcher, so the adapter - /// tells an expiry from an interruption when the dispatcher ends its - /// wait. + /// The observer that labels each firing as the projection does and + /// sees a gate report a question's expiry. A run registers it ahead of + /// the interview dispatcher, so a question's stage is known when it is + /// asked and the adapter tells an expiry from an interruption when the + /// dispatcher ends its wait. #[must_use] pub fn observer(&self) -> Arc { - self.expiries.clone() + self.observed.clone() } /// Post a notice; a failure after the question was asked is logged, @@ -383,15 +458,26 @@ impl FabroInterviewer { #[async_trait::async_trait] impl Interviewer for FabroInterviewer { async fn reply(&self, request: InterviewRequest, cancel: CancellationToken) -> InterviewReply { - let asked = asked_question(&request); + let stage = self + .observed + .label(request.execution, request.firing) + .map_or_else( + || { + debug!( + node = %request.node, + execution = request.execution.raw(), + firing = request.firing.raw(), + "the asking firing has no label yet; the question names the node" + ); + request.node.to_string() + }, + |label| label.to_string(), + ); + let asked = asked_question(&request, stage); let question_id = asked.question_id.clone(); let text = asked.text.clone(); let stage = asked.stage.clone(); let legacy = legacy_question(&asked); - // HOOK POINT (read side): the mapping from Petri's identity - // (`asked.identity`) to Fabro's question id is a platform fact - // worth a record keyed on that identity; today it lives only in - // the `interview.started` event posted here. if let Err(error) = self.sink.post(QuestionNotice::Asked(asked)).await { return InterviewReply::Failed(InterviewError::with_source( format!("question `{question_id}` could not be published to Fabro"), @@ -400,9 +486,8 @@ impl Interviewer for FabroInterviewer { } let mut outstanding = Outstanding { sink: Arc::clone(&self.sink), - expiries: Arc::clone(&self.expiries), + observed: Arc::clone(&self.observed), execution: request.execution, - question: request.question.id.clone(), question_id: question_id.clone(), text: text.clone(), stage: stage.clone(), @@ -423,9 +508,9 @@ impl Interviewer for FabroInterviewer { outstanding.close_unanswered("cancelled"); return InterviewReply::Cancelled; }; - // HOOK POINT (read side): `submission.actor` is who answered, a - // Fabro fact Petri's answer record does not carry; it rides on - // `interview.completed` until a platform record holds it. + // `submission.actor` is who answered, a Fabro fact Petri's answer + // record does not carry: it goes out on `interview.completed`, from + // which the store derives the `interview.answered` platform record. let Some(answer) = petri_answer(&submission.answer, &request.question) else { outstanding.close_unanswered(&reason_of(&submission.answer.value)); return InterviewReply::Cancelled; @@ -451,10 +536,9 @@ impl Interviewer for FabroInterviewer { /// expiry, else `interview.interrupted`. struct Outstanding { sink: Arc, - expiries: Arc, + observed: Arc, execution: ExecutionId, - /// Petri's question id, as the expiry report names it. - question: String, + /// Petri's question id, as the expiry report names it too. question_id: String, text: String, stage: String, @@ -471,7 +555,7 @@ impl Outstanding { } self.open = false; let duration_ms = millis(self.started.elapsed()); - let expired = self.expiries.contains(self.execution, &self.question); + let expired = self.observed.expired(self.execution, &self.question_id); let notice = if expired { QuestionNotice::Expired { question_id: self.question_id.clone(), @@ -528,38 +612,15 @@ impl std::fmt::Display for AnyhowError { impl std::error::Error for AnyhowError {} -/// Fabro's id for a question, from Petri's identity: the node, then the -/// execution and firing (unique in the run), the occurrence and the ask -/// (a re-asked question is a new one). Only URL-safe characters, so the -/// id travels in the answer endpoint's path as it is. -#[must_use] -pub fn question_id(identity: &QuestionIdentity) -> String { - let node: String = identity - .node - .chars() - .map(|c| { - if c.is_ascii_alphanumeric() || c == '_' || c == '-' { - c - } else { - '_' - } - }) - .collect(); - format!( - "{node}.x{}.f{}.q{}.a{}", - identity.execution, identity.firing, identity.occurrence, identity.ask - ) -} - -/// The question as Fabro shows it. -fn asked_question(request: &InterviewRequest) -> AskedQuestion { - let identity = QuestionIdentity::of(request); +/// The question as Fabro shows it, under Petri's id and the stage label +/// the projection gives the asking firing. +fn asked_question(request: &InterviewRequest, stage: String) -> AskedQuestion { let question = &request.question; AskedQuestion { - question_id: question_id(&identity), - identity, + question_id: question.id.clone(), + identity: QuestionIdentity::of(request), text: question.text.clone(), - stage: request.node.to_string(), + stage, question_type: question_type(question), options: question .options @@ -756,20 +817,6 @@ mod tests { question } - #[test] - fn a_question_id_is_url_safe_and_names_the_identity() { - let identity = QuestionIdentity { - invocation_path: "/branch:fan@2:0:a".into(), - execution: 2, - firing: 3, - attempt: 1, - node: "approve plan".into(), - occurrence: 1, - ask: 2, - }; - assert_eq!(question_id(&identity), "approve_plan.x2.f3.q1.a2"); - } - #[test] fn yes_and_no_name_the_gates_choices_by_key() { let question = yes_no(); diff --git a/lib/components/fabro-petri/src/projection.rs b/lib/components/fabro-petri/src/projection.rs index dc0f03262..bb82449d4 100644 --- a/lib/components/fabro-petri/src/projection.rs +++ b/lib/components/fabro-petri/src/projection.rs @@ -44,7 +44,7 @@ use fabro_types::{ }; use lithos_llm::catalog::{ModelId, ProviderId}; use lithos_llm::types::Usage; -use petri_execution::events::{Derived, Parsed, RunEvent, Subject, ViewEvent, WaitState}; +use petri_execution::events::{Derived, NodeRef, Parsed, RunEvent, Subject, ViewEvent, WaitState}; use petri_execution::{CoordinatorEvent, ExecutionId}; use petri_runtime::engine::{Admission, Event}; use petri_runtime::ir::{Metrics, Status, StepEvent}; @@ -978,28 +978,15 @@ impl RunView { return; } let node_name = subject.node.name.to_string(); - let visit = subject.visit.unwrap_or(1).max(1); - let meta_kind = subject - .node - .meta - .get("kind") - .and_then(Value::as_str) - .unwrap_or(""); - let synthetic = subject - .node - .meta - .get("synthetic") - .and_then(Value::as_bool) - .unwrap_or(false); - let shown = !synthetic && meta_kind != "parallel.branch"; + let visit = visit_of(subject); + let meta_kind = node_meta_kind(&subject.node); + let shown = is_shown(&subject.node); // Only a shown stage takes a label: a lowering node (a branch's // parent-side delegate shares its target's name) never competes with // the stage it stands for. let mut stage_id = StageId::new(node_name.clone(), visit); if shown { - if self.state.labels.contains(&stage_id.to_string()) { - stage_id = StageId::new(format!("{node_name}/e{}", execution.raw()), visit); - } + stage_id = stage_label(&node_name, visit, execution, &self.state.labels); self.state.labels.insert(stage_id.to_string()); } self.state.stages.insert(key, StageRef { @@ -1181,6 +1168,50 @@ pub fn stage_key(execution: u64, firing: u64) -> String { format!("{execution}:{firing}") } +/// Which firing of its node a subject is, 1-based. +#[must_use] +pub fn visit_of(subject: &Subject) -> u32 { + subject.visit.unwrap_or(1).max(1) +} + +/// The role a frontend gave a node under `meta.kind`, or the empty string. +fn node_meta_kind(node: &NodeRef) -> &str { + node.meta.get("kind").and_then(Value::as_str).unwrap_or("") +} + +/// Whether a node is a logical stage the projection shows, or a lowering +/// node it keeps off the list: one a frontend marked synthetic, or a +/// parallel branch's delegate. +#[must_use] +pub fn is_shown(node: &NodeRef) -> bool { + let synthetic = node + .meta + .get("synthetic") + .and_then(Value::as_bool) + .unwrap_or(false); + !synthetic && node_meta_kind(node) != "parallel.branch" +} + +/// The label a shown firing takes, which is the stage id the projection +/// keys it by: `node@visit`, or `node/e@visit` when another +/// execution's firing already took that label. `taken` is every label given +/// so far; the caller adds the one returned. The interview adapter labels a +/// question's stage through this same rule, so the stage a question names +/// is the stage the projection shows. +#[must_use] +pub fn stage_label( + node_name: &str, + visit: u32, + execution: ExecutionId, + taken: &BTreeSet, +) -> StageId { + let stage_id = StageId::new(node_name.to_string(), visit); + if taken.contains(&stage_id.to_string()) { + return StageId::new(format!("{node_name}/e{}", execution.raw()), visit); + } + stage_id +} + /// The fork firing and branch index a branch child's call slot names: /// `branch:@::`. fn branch_slot(slot: &str) -> Option<(u64, u32)> { @@ -1331,7 +1362,6 @@ pub fn run_id_of(key: &str) -> Option { mod tests { use fabro_store::platform_records::RunCreatedRecord; use fabro_types::test_support as types_support; - use petri_execution::events::NodeRef; use petri_runtime::driver::BranchRole; use petri_runtime::ir::{FiringId, NodeId}; diff --git a/lib/components/fabro-petri/tests/interview.rs b/lib/components/fabro-petri/tests/interview.rs index 7415a9e19..9a5f5e29d 100644 --- a/lib/components/fabro-petri/tests/interview.rs +++ b/lib/components/fabro-petri/tests/interview.rs @@ -1,7 +1,8 @@ //! Petri's human gates through Fabro's interview adapter: a question is -//! posted as Fabro's `interview.started`, the answer submitted to the -//! control interviewer under the posted id reaches the gate, two parallel -//! gates each get their own answer, an expired question is completed as a +//! posted as Fabro's `interview.started` under Petri's own id and the +//! projection's stage label, the answer submitted to the control +//! interviewer under that id reaches the gate, two parallel gates each get +//! their own answer, an expired question is completed as a //! timeout with the gate's default, an auto-approved run answers itself, //! and a cancelled run interrupts its question. //! @@ -197,7 +198,7 @@ async fn a_gate_answered_under_the_posted_id_routes_on_the_answer() { let board = gate.board.clone(); let control = gate.control.clone(); tokio::spawn(async move { - let asked = board.wait_asked("gate").await; + let asked = board.wait_asked("gate@1").await; control .submit(&asked.question_id, engine_submission(LegacyAnswer::no())) .await @@ -214,7 +215,10 @@ async fn a_gate_answered_under_the_posted_id_routes_on_the_answer() { gate.marker("no") && !gate.marker("yes"), "the no branch ran" ); - assert_eq!(asked.stage, "gate"); + assert_eq!( + asked.stage, "gate@1", + "the projection's label for the firing" + ); assert_eq!(asked.text, "Go?"); assert_eq!(asked.question_type, QuestionType::YesNo); assert_eq!( @@ -226,11 +230,14 @@ async fn a_gate_answered_under_the_posted_id_routes_on_the_answer() { vec![("Y", "[Y] Yes"), ("N", "[N] No")] ); assert!( - asked.question_id.starts_with("gate.x0.f"), - "{}", + asked.question_id.starts_with("gate#"), + "Petri's id, as the projection serves it: {}", asked.question_id ); assert_eq!(asked.identity.node, "gate"); + assert_eq!(asked.identity.execution, 0); + assert_eq!(asked.identity.occurrence, 1); + assert_eq!(asked.identity.ask, 1); assert_eq!(asked.identity.invocation_path, "/"); let notices = gate.board.notices(); assert!( @@ -280,8 +287,8 @@ async fn two_parallel_gates_each_bind_their_own_answer() { let control = gate.control.clone(); tokio::spawn(async move { // Both are pending before either is answered, and `b` first. - let a = board.wait_asked("a").await; - let b = board.wait_asked("b").await; + let a = board.wait_asked("a@1").await; + let b = board.wait_asked("b@1").await; assert_ne!(a.question_id, b.question_id); control .submit(&b.question_id, engine_submission(LegacyAnswer::yes())) @@ -354,14 +361,14 @@ async fn an_unanswered_question_expires_with_the_gates_default() { assert_eq!(outcome.status, RunStatus::Success, "{outcome:?}"); assert!(gate.marker("no") && !gate.marker("yes"), "the default ran"); - let asked = gate.board.asked("gate").expect("asked"); + let asked = gate.board.asked("gate@1").expect("asked"); let notices = gate.board.wait_ended(&asked.question_id).await; assert_eq!(asked.timeout_seconds, Some(0.3)); assert!( matches!( ¬ices[1], QuestionNotice::Expired { question_id, stage, .. } - if *question_id == asked.question_id && stage == "gate" + if *question_id == asked.question_id && stage == "gate@1" ), "{notices:?}" ); @@ -443,7 +450,7 @@ async fn a_cancelled_run_interrupts_its_pending_question() { let canceller = { let board = gate.board.clone(); tokio::spawn(async move { - board.wait_asked("gate").await; + board.wait_asked("gate@1").await; cancel.cancel(); }) }; @@ -453,7 +460,7 @@ async fn a_cancelled_run_interrupts_its_pending_question() { assert_eq!(outcome.status, RunStatus::Cancelled, "{outcome:?}"); assert!(!gate.marker("yes") && !gate.marker("no"), "no branch ran"); - let asked = gate.board.asked("gate").expect("asked"); + let asked = gate.board.asked("gate@1").expect("asked"); let notices = gate.board.wait_ended(&asked.question_id).await; assert!( matches!( diff --git a/lib/components/fabro-petri/tests/projection.rs b/lib/components/fabro-petri/tests/projection.rs index 69057e9c5..42cf6876d 100644 --- a/lib/components/fabro-petri/tests/projection.rs +++ b/lib/components/fabro-petri/tests/projection.rs @@ -15,15 +15,22 @@ )] #![expect(clippy::print_stderr, reason = "a skipped test says why on its stderr")] +mod support; + use std::collections::BTreeSet; use std::env; use std::path::{Path, PathBuf}; use std::sync::Arc; -use std::time::Duration; +use std::time::{Duration, Instant}; use fabro_db::DbPool; +use fabro_interview::ControlInterviewer; use fabro_petri::SqliteRunStore; +use fabro_petri::check::Launch; +use fabro_petri::engine::{self, RunStatus as EngineRunStatus}; +use fabro_petri::interview::{Approval, FabroInterviewer}; use fabro_petri::projector::{self, Projector}; +use fabro_petri::runtime::RuntimeSpec; use fabro_store::platform_records::{ PlatformRecord, PlatformRecordStore, RunCreatedRecord, RunLifecycleKind, RunLifecycleRecord, }; @@ -40,6 +47,7 @@ use petri_runtime::ir::RunStatus as PetriRunStatus; use petri_runtime::{RunOptions, Runtime}; use petri_store::{RunKey, RunStore}; use tokio::fs; +use tokio::time::sleep; const HOST_PLUGIN: &str = "sandbox-driver-host"; const HOST_PLUGIN_OVERRIDE: &str = "PETRI_SANDBOX_HOST_PLUGIN"; @@ -73,6 +81,26 @@ const PARALLEL_WORKFLOW: &str = r#"digraph Parallel { const SETTINGS: &str = "_version = 1\n\n[workflow]\ngraph = \"workflow.fabro\"\n"; +/// One yes/no gate whose branches leave a marker file each. +fn gate_workflow(markers: &Path, gate_attrs: &str) -> String { + format!( + r#"digraph Gate {{ + graph [goal="Ask once"] + start [shape=Mdiamond] + exit [shape=Msquare] + gate [shape=hexagon, label="Go?", question_type="yes_no"{gate_attrs}] + yes [shape=parallelogram, script="touch {dir}/yes"] + no [shape=parallelogram, script="touch {dir}/no"] + start -> gate + gate -> yes [label="[Y] Yes"] + gate -> no [label="[N] No"] + yes -> exit + no -> exit +}}"#, + dir = markers.display() + ) +} + fn host_plugin() -> Option { let found = env::var_os(HOST_PLUGIN_OVERRIDE) .map(PathBuf::from) @@ -884,3 +912,201 @@ async fn a_torn_tail_holds_the_view_and_reports_the_run_incomplete() { "the projection stands where it was" ); } + +/// A gate scenario runs through the engine assembly with the interview +/// adapter, as a Fabro run does, over a store that signals the projector. +struct GateRun { + scenario: Scenario, + markers: PathBuf, + projector: Arc, + workflow: String, + _root: tempfile::TempDir, +} + +async fn gate_run(gate_attrs: &str) -> GateRun { + let root = tempfile::tempdir().expect("a marker dir"); + let markers = root.path().join("markers"); + fs::create_dir_all(&markers) + .await + .expect("the marker dir creates"); + let workflow = gate_workflow(&markers, gate_attrs); + let scenario = scenario( + "gate", + &[("workflow.fabro", &workflow), ("workflow.toml", SETTINGS)], + false, + ) + .await; + let projector = Projector::new(scenario.pool.clone(), scenario.pool.clone()); + projector.signal(scenario.run_id); + GateRun { + scenario, + markers, + projector, + workflow, + _root: root, + } +} + +impl GateRun { + /// Run the gate to completion through the adapter, under `approval`, + /// with nobody answering. + async fn run(&self, approval: Approval) { + let runtime = RuntimeSpec::default(); + let graphs = support::admit( + &[ + ("workflow.fabro", &self.workflow), + ("workflow.toml", SETTINGS), + ], + Launch::default(), + &runtime, + ); + let interviewer = FabroInterviewer::new( + Arc::new(ControlInterviewer::new()), + Arc::new(support::Silent), + approval, + ); + let store = self + .projector + .observe_store(Arc::new(SqliteRunStore::new(self.scenario.pool.clone()))); + let request = support::run_request( + &self.scenario.run_id.to_string(), + &self.scenario.run_dir, + graphs, + store, + runtime, + interviewer, + ); + let outcome = engine::run(request).await.expect("the run ends"); + assert_eq!(outcome.status, EngineRunStatus::Success, "{outcome:?}"); + self.projector.settle(self.scenario.run_id).await; + } + + async fn stored(&self) -> fabro_types::RunProjection { + projector::stored_projection(&self.scenario.pool, self.scenario.run_id) + .await + .expect("the stored projection reads") + .expect("the run has a stored projection") + } +} + +/// A question Petri expires: while the gate waits, the projection shows +/// the question pending under Petri's id and the firing's label with the +/// run blocked; once `question_expired` lands and the gate takes its +/// default, the question is gone, the run runs on to success, and the view +/// rebuilds the same. +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn an_expired_question_is_pending_while_the_gate_waits_and_closes_on_the_expiry() { + if host_plugin().is_none() { + return; + } + let gate = Arc::new(gate_run(r#", timeout="1500ms", human.default_choice="no""#).await); + let running = { + let gate = Arc::clone(&gate); + tokio::spawn(async move { gate.run(Approval::Prompt).await }) + }; + let pending = { + let pool = gate.scenario.pool.clone(); + let run_id = gate.scenario.run_id; + let deadline = Instant::now() + Duration::from_secs(30); + loop { + let stored = projector::stored_projection(&pool, run_id) + .await + .expect("the stored projection reads"); + if let Some(stored) = stored.filter(|stored| !stored.pending_interviews.is_empty()) { + break stored; + } + assert!( + Instant::now() < deadline, + "the question never showed as pending" + ); + sleep(Duration::from_millis(10)).await; + } + }; + let (id, record) = pending + .pending_interviews + .iter() + .next() + .expect("one pending question"); + assert!(id.starts_with("gate#"), "Petri's id: {id}"); + assert_eq!(&record.question.id, id); + assert_eq!(record.question.stage, "gate@1"); + assert_eq!(record.question.text, "Go?"); + assert_eq!( + record + .question + .options + .iter() + .map(|option| option.key.as_str()) + .collect::>(), + vec!["Y", "N"] + ); + assert_eq!(record.question.timeout_seconds, Some(1.5)); + assert!( + matches!(pending.status, RunStatus::Blocked { .. }), + "{:?}", + pending.status + ); + + running.await.expect("the run task ends"); + + assert!( + gate.markers.join("no").exists() && !gate.markers.join("yes").exists(), + "the default ran" + ); + let stored = gate.stored().await; + assert!( + stored.pending_interviews.is_empty(), + "the expired question is no longer pending: {:?}", + stored.pending_interviews + ); + assert!( + matches!(stored.status, RunStatus::Succeeded { .. }), + "{:?}", + stored.status + ); + let gate_stage = stored + .stage(&StageId::new("gate", 1)) + .expect("the gate is a stage"); + assert_eq!(gate_stage.state, StageState::Succeeded); + assert_view_equals_rebuild(&gate.scenario.pool, gate.scenario.run_id).await; +} + +/// An auto-approved run answers its gate at once: the delivered answer +/// closes the question in the projection, the affirmative branch runs, and +/// the view rebuilds the same. +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn an_auto_approved_answer_closes_the_question_in_the_projection() { + if host_plugin().is_none() { + return; + } + let gate = gate_run("").await; + gate.run(Approval::Auto).await; + + assert!( + gate.markers.join("yes").exists() && !gate.markers.join("no").exists(), + "the yes branch ran" + ); + let stored = gate.stored().await; + assert!( + stored.pending_interviews.is_empty(), + "the answered question is no longer pending: {:?}", + stored.pending_interviews + ); + assert!( + matches!(stored.status, RunStatus::Succeeded { .. }), + "{:?}", + stored.status + ); + let gate_stage = stored + .stage(&StageId::new("gate", 1)) + .expect("the gate is a stage"); + assert_eq!(gate_stage.state, StageState::Succeeded); + let states = stage_states(&gate.scenario.pool, gate.scenario.run_id).await; + assert!( + states + .iter() + .any(|(label, state)| label == "yes@1" && *state == StageState::Succeeded), + "{states:?}" + ); + assert_view_equals_rebuild(&gate.scenario.pool, gate.scenario.run_id).await; +} diff --git a/lib/components/fabro-store/src/platform_records.rs b/lib/components/fabro-store/src/platform_records.rs index 83a436180..13d8d86b7 100644 --- a/lib/components/fabro-store/src/platform_records.rs +++ b/lib/components/fabro-store/src/platform_records.rs @@ -784,7 +784,9 @@ fn interview_answered_record( #[cfg(test)] mod tests { - use fabro_types::{FailureReason, RunStatus, fixtures, test_support as types_support}; + use fabro_types::{ + FailureReason, RunStatus, SystemActorKind, fixtures, test_support as types_support, + }; use serde_json::json; use super::*; @@ -978,6 +980,42 @@ mod tests { ); } + /// The interview adapter completes a Petri question under Petri's own + /// id with the answering principal as the event's actor; the record + /// keeps both, so who answered is a platform fact keyed on that id. + #[test] + fn a_completed_interview_becomes_an_answered_record_under_petris_id_with_its_actor() { + let actor = Principal::System { + system_kind: SystemActorKind::Engine, + }; + let event = fabro_types::RunEvent { + id: "evt".to_string(), + ts: chrono::Utc::now(), + run_id: fixtures::RUN_1, + node_id: None, + node_label: None, + stage_id: None, + parallel_group_id: None, + parallel_branch_id: None, + session_id: None, + parent_session_id: None, + tool_call_id: None, + actor: Some(actor.clone()), + body: EventBody::InterviewCompleted(InterviewCompletedProps { + question_id: "gate#2".to_string(), + question: "Go?".to_string(), + answer: "N".to_string(), + duration_ms: 1_200, + }), + }; + let Some(PlatformRecord::InterviewAnswered(record)) = platform_record_for(&event) else { + panic!("a completed interview maps to an answered record"); + }; + assert_eq!(record.question, "gate#2"); + assert_eq!(record.principal, Some(actor)); + assert_eq!(record.channel, None); + } + #[test] fn a_failed_legacy_event_becomes_a_failed_lifecycle_record_with_its_message() { let event = fabro_types::RunEvent {