mirror of
https://github.com/fabro-sh/fabro.git
synced 2026-10-01 02:04:24 +00:00
Give a Petri question one identity across the adapter and the projection
The interview adapter derived its own question id from Petri's identity and posted it on `interview.started`, while the projection over Petri's records serves the pending question under Petri's `Question.id` with the firing's stage label. The answer endpoint validates against the projection, so an answer under the projection's id never reached the adapter's wait. The adapter now waits under Petri's id and labels the question's stage through the projection's own rule: `stage_label`, `is_shown` and `visit_of` move out of `start_visit` into shared functions, and the adapter's observer derives each firing's `visit.started` through Petri's `Projection`, as the projector does, so the label matches by construction. The full Petri identity stays on `AskedQuestion`. The legacy `interview.*` events are still posted, under Petri's id, for the readers that follow the event stream rather than the projection: the Slack service, `run attach`, the web app's Q&A renderer and the server's answer claim. The store already derives the `interview.answered` platform record from `interview.completed` for a Petri run, so who answered is recorded under Petri's id with the answering principal. The gate scenarios assert the new identity and encode the id as one path segment, as the generated clients do. Projection tests cover an expired question and an auto-approved answer. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
This commit is contained in:
parent
5dae891a98
commit
a6ac3f120e
8 changed files with 569 additions and 177 deletions
|
|
@ -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}");
|
||||
|
|
|
|||
|
|
@ -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:?}"
|
||||
);
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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<InterviewOption>,
|
||||
|
|
@ -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<HashSet<(ExecutionId, String)>>,
|
||||
pub struct Observed {
|
||||
state: Mutex<ObservedState>,
|
||||
}
|
||||
|
||||
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<String>,
|
||||
/// 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<StageId> {
|
||||
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<ControlInterviewer>,
|
||||
sink: Arc<dyn QuestionSink>,
|
||||
approval: Approval,
|
||||
expiries: Arc<Expiries>,
|
||||
observed: Arc<Observed>,
|
||||
}
|
||||
|
||||
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<dyn ExecutionObserver> {
|
||||
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<dyn QuestionSink>,
|
||||
expiries: Arc<Expiries>,
|
||||
observed: Arc<Observed>,
|
||||
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();
|
||||
|
|
|
|||
|
|
@ -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<execution>@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<String>,
|
||||
) -> 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:<fork>@<firing>:<index>:<target>`.
|
||||
fn branch_slot(slot: &str) -> Option<(u64, u32)> {
|
||||
|
|
@ -1331,7 +1362,6 @@ pub fn run_id_of(key: &str) -> Option<RunId> {
|
|||
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};
|
||||
|
||||
|
|
|
|||
|
|
@ -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!(
|
||||
|
|
|
|||
|
|
@ -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<PathBuf> {
|
||||
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<Projector>,
|
||||
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<_>>(),
|
||||
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;
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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 {
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue