diff --git a/Cargo.lock b/Cargo.lock index dbc1f6322..cb7206de0 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -2427,12 +2427,9 @@ name = "fabro-interview" version = "0.361.0-nightly.0" dependencies = [ "async-trait", - "dialoguer", "fabro-types", - "fabro-util", "serde", "serde_json", - "tempfile", "tokio", "tracing", ] diff --git a/lib/components/fabro-interview/Cargo.toml b/lib/components/fabro-interview/Cargo.toml index 4c450f54a..c236396fa 100644 --- a/lib/components/fabro-interview/Cargo.toml +++ b/lib/components/fabro-interview/Cargo.toml @@ -18,10 +18,7 @@ serde_json.workspace = true async-trait.workspace = true tokio.workspace = true tracing.workspace = true -dialoguer.workspace = true -fabro-util = { path = "../../foundation/fabro-util" } fabro-types = { path = "../../foundation/fabro-types" } [dev-dependencies] tokio = { workspace = true, features = ["test-util", "macros"] } -tempfile = "3" diff --git a/lib/components/fabro-interview/src/callback.rs b/lib/components/fabro-interview/src/callback.rs deleted file mode 100644 index 2159a7eeb..000000000 --- a/lib/components/fabro-interview/src/callback.rs +++ /dev/null @@ -1,73 +0,0 @@ -use async_trait::async_trait; -use fabro_types::{Principal, SystemActorKind}; - -use crate::{Answer, AnswerSubmission, Interviewer, Question}; - -/// Delegates question answering to a provided callback function. -pub struct CallbackInterviewer { - callback: Box Answer + Send + Sync>, - actor: Principal, -} - -impl CallbackInterviewer { - pub fn new(callback: impl Fn(Question) -> Answer + Send + Sync + 'static) -> Self { - Self::with_actor( - Principal::System { - system_kind: SystemActorKind::Engine, - }, - callback, - ) - } - - pub fn with_actor( - actor: Principal, - callback: impl Fn(Question) -> Answer + Send + Sync + 'static, - ) -> Self { - Self { - callback: Box::new(callback), - actor, - } - } -} - -#[async_trait] -impl Interviewer for CallbackInterviewer { - async fn ask(&self, question: Question) -> AnswerSubmission { - AnswerSubmission::new((self.callback)(question), self.actor.clone()) - } -} - -#[cfg(test)] -mod tests { - use fabro_types::QuestionType; - - use super::*; - use crate::AnswerValue; - - #[tokio::test] - async fn calls_callback_with_question() { - let interviewer = CallbackInterviewer::new(|q| { - if q.question_type == QuestionType::YesNo { - Answer::yes() - } else { - Answer::no() - } - }); - - let yes_q = Question::new("approve?", QuestionType::YesNo); - let answer = interviewer.ask(yes_q).await.answer; - assert_eq!(answer.value, AnswerValue::Yes); - - let no_q = Question::new("choose:", QuestionType::MultipleChoice); - let answer = interviewer.ask(no_q).await.answer; - assert_eq!(answer.value, AnswerValue::No); - } - - #[tokio::test] - async fn callback_receives_question_text() { - let interviewer = CallbackInterviewer::new(|q| Answer::text(q.text)); - let q = Question::new("hello world", QuestionType::Freeform); - let answer = interviewer.ask(q).await.answer; - assert_eq!(answer.text, Some("hello world".to_string())); - } -} diff --git a/lib/components/fabro-interview/src/console.rs b/lib/components/fabro-interview/src/console.rs deleted file mode 100644 index b96114765..000000000 --- a/lib/components/fabro-interview/src/console.rs +++ /dev/null @@ -1,445 +0,0 @@ -use std::io::IsTerminal; - -use async_trait::async_trait; -use dialoguer::console::Term; -use dialoguer::theme::ColorfulTheme; -use fabro_types::{InterviewOption, Principal, QuestionType}; -use fabro_util::terminal::Styles; -use tokio::io::{self, AsyncBufReadExt, BufReader}; -use tokio::task; - -use crate::{Answer, AnswerSubmission, AnswerValue, Interviewer, Question}; - -enum PromptRead { - Line(String), - Eof, - Error, -} - -/// Reads from stdin to collect answers. Displays formatted prompts per spec -/// 6.4. -pub struct ConsoleInterviewer { - styles: &'static Styles, - actor: Principal, -} - -impl ConsoleInterviewer { - #[must_use] - pub fn new(styles: &'static Styles, actor: Principal) -> Self { - Self { styles, actor } - } -} - -fn find_matching_option(response: &str, options: &[InterviewOption]) -> Option { - let trimmed = response.trim(); - // Try matching by key (case-insensitive) - for opt in options { - if opt.key.eq_ignore_ascii_case(trimmed) { - return Some(Answer { - value: AnswerValue::Selected(opt.key.clone()), - selected_option: Some(opt.clone()), - text: None, - }); - } - } - // Try matching by 1-based index - if let Ok(idx) = trimmed.parse::() { - if idx >= 1 && idx <= options.len() { - let opt = &options[idx - 1]; - return Some(Answer { - value: AnswerValue::Selected(opt.key.clone()), - selected_option: Some(opt.clone()), - text: None, - }); - } - } - None -} - -#[allow( - clippy::print_stderr, - reason = "Prompts go to stderr so piped stdout stays machine-readable." -)] -async fn read_line(prompt: &str) -> PromptRead { - // Print the prompt to stderr so it doesn't interfere with piped stdout - eprint!("{prompt}"); - let stdin = io::stdin(); - let mut reader = BufReader::new(stdin); - let mut line = String::new(); - match reader.read_line(&mut line).await { - Ok(0) => PromptRead::Eof, - Ok(_) => PromptRead::Line(line.trim_end().to_string()), - Err(_) => PromptRead::Error, - } -} - -fn parse_non_tty_choice_response(question: &Question, prompt_read: PromptRead) -> Answer { - let PromptRead::Line(response) = prompt_read else { - return Answer::interrupted(); - }; - if response.trim().is_empty() { - return Answer::interrupted(); - } - if let Some(answer) = find_matching_option(&response, &question.options) { - return answer; - } - if question.allow_freeform { - return Answer::text(response); - } - find_matching_option(&response, &question.options).unwrap_or_else(Answer::interrupted) -} - -fn parse_non_tty_confirm_response(prompt_read: PromptRead) -> Answer { - let PromptRead::Line(response) = prompt_read else { - return Answer::interrupted(); - }; - match response.trim().to_lowercase().as_str() { - "y" | "yes" => Answer::yes(), - "n" | "no" => Answer::no(), - _ => Answer::interrupted(), - } -} - -fn parse_non_tty_freeform_response(prompt_read: PromptRead) -> Answer { - let PromptRead::Line(response) = prompt_read else { - return Answer::interrupted(); - }; - if response.trim().is_empty() { - Answer::interrupted() - } else { - Answer::text(response) - } -} - -/// The review target line printed above a question in terminal clients, which -/// cannot render a hyperlink label. The label and resource noun are already in -/// `question.text`, so only the URL is shown. Shared with `fabro-cli`'s attach -/// client. -#[must_use] -pub fn review_target_line(question: &Question) -> Option { - question - .review_target - .as_ref() - .map(|target| format!("Review link: {}", target.url())) -} - -/// Ask a multiple-choice question using dialoguer's `Select` widget on a TTY. -fn ask_select_interactive(question: &Question) -> Answer { - let items: Vec = question - .options - .iter() - .map(|opt| format!("{} - {}", opt.key, opt.label)) - .collect(); - - let has_freeform = question.allow_freeform; - let mut all_items = items; - if has_freeform { - all_items.push("Other (free text)...".to_string()); - } - - let selection = dialoguer::Select::with_theme(&ColorfulTheme::default()) - .with_prompt(&question.text) - .items(&all_items) - .default(0) - .interact_on_opt(&Term::stderr()); - - match selection { - Ok(Some(idx)) if has_freeform && idx == question.options.len() => { - // User chose the free-text option - dialoguer::Input::::with_theme(&ColorfulTheme::default()) - .with_prompt("Enter your response") - .interact_on(&Term::stderr()) - .map_or_else( - |_| Answer::interrupted(), - |response| { - if response.trim().is_empty() { - Answer::interrupted() - } else { - Answer::text(response) - } - }, - ) - } - Ok(Some(idx)) if idx < question.options.len() => { - let opt = &question.options[idx]; - Answer { - value: AnswerValue::Selected(opt.key.clone()), - selected_option: Some(opt.clone()), - text: None, - } - } - _ => Answer::interrupted(), - } -} - -/// Ask a multi-select question using dialoguer's `MultiSelect` widget on a TTY. -fn ask_multi_select_interactive(question: &Question) -> Answer { - let items: Vec = question - .options - .iter() - .map(|opt| format!("{} - {}", opt.key, opt.label)) - .collect(); - - let selection = dialoguer::MultiSelect::with_theme(&ColorfulTheme::default()) - .with_prompt(&question.text) - .items(&items) - .interact_on_opt(&Term::stderr()); - - match selection { - Ok(Some(indices)) if !indices.is_empty() => { - let keys: Vec = indices - .iter() - .map(|&i| question.options[i].key.clone()) - .collect(); - Answer::multi_selected(keys) - } - _ => Answer::interrupted(), - } -} - -/// Ask a yes/no or confirmation question using dialoguer's `Confirm` widget on -/// a TTY. -fn ask_confirm_interactive(question: &Question) -> Answer { - let confirmed = dialoguer::Confirm::with_theme(&ColorfulTheme::default()) - .with_prompt(&question.text) - .default(true) - .interact_on_opt(&Term::stderr()); - - match confirmed { - Ok(Some(true)) => Answer::yes(), - Ok(Some(false)) => Answer::no(), - _ => Answer::interrupted(), - } -} - -/// Ask a freeform question using dialoguer's `Input` widget on a TTY. -fn ask_freeform_interactive(question: &Question) -> Answer { - dialoguer::Input::::with_theme(&ColorfulTheme::default()) - .with_prompt(&question.text) - .interact_on(&Term::stderr()) - .map_or_else( - |_| Answer::interrupted(), - |response| { - if response.trim().is_empty() { - Answer::interrupted() - } else { - Answer::text(response) - } - }, - ) -} - -#[async_trait] -impl Interviewer for ConsoleInterviewer { - #[allow( - clippy::print_stderr, - reason = "Interactive questions and options belong on stderr, not captured stdout." - )] - async fn ask(&self, question: Question) -> AnswerSubmission { - // If stdin is a TTY, use dialoguer for interactive arrow-key navigation. - // Otherwise, fall back to the line-based reader for piped input. - #[expect( - clippy::disallowed_methods, - reason = "is_terminal() on the std stdin handle is a non-blocking fstat check; no \ - actual I/O performed. The real blocking read runs inside spawn_blocking \ - below." - )] - if std::io::stdin().is_terminal() { - if let Some(ref context_text) = question.context_display { - let rendered = self.styles.render_markdown(context_text); - eprint!("{rendered}"); - } - if let Some(line) = review_target_line(&question) { - eprintln!("{line}"); - } - let q = question; - let answer = task::spawn_blocking(move || match q.question_type { - QuestionType::MultipleChoice => ask_select_interactive(&q), - QuestionType::MultiSelect => ask_multi_select_interactive(&q), - QuestionType::YesNo | QuestionType::Confirmation => ask_confirm_interactive(&q), - QuestionType::Freeform => ask_freeform_interactive(&q), - }) - .await - .unwrap_or_else(|_| Answer::interrupted()); - return AnswerSubmission::new(answer, self.actor.clone()); - } - - // Non-TTY fallback: line-based stdin reading - let s = self.styles; - if let Some(line) = review_target_line(&question) { - eprintln!("{line}"); - } - eprintln!("{} {}", s.bold_cyan.apply_to("?"), question.text); - - let answer = match question.question_type { - QuestionType::MultipleChoice | QuestionType::MultiSelect => { - for (i, opt) in question.options.iter().enumerate() { - eprintln!( - " {}{}{} {} - {}", - s.dim.apply_to("["), - s.bold.apply_to(i + 1), - s.dim.apply_to("]"), - opt.key, - opt.label, - ); - } - if question.allow_freeform { - eprintln!(" Or type a free-text response"); - } - parse_non_tty_choice_response(&question, read_line("Select: ").await) - } - QuestionType::YesNo | QuestionType::Confirmation => { - parse_non_tty_confirm_response(read_line("[Y/N]: ").await) - } - QuestionType::Freeform => parse_non_tty_freeform_response(read_line("> ").await), - }; - AnswerSubmission::new(answer, self.actor.clone()) - } - - #[allow( - clippy::print_stderr, - reason = "Stage notices belong on stderr, not captured stdout." - )] - async fn inform(&self, message: &str, stage: &str) { - let s = self.styles; - eprintln!("{} {message}", s.dim.apply_to(format!("[{stage}]"))); - } -} - -#[cfg(test)] -mod tests { - use super::*; - - #[test] - fn find_matching_option_by_key() { - let options = vec![ - InterviewOption { - key: "A".to_string(), - label: "Approve".to_string(), - description: None, - preview: None, - }, - InterviewOption { - key: "R".to_string(), - label: "Reject".to_string(), - description: None, - preview: None, - }, - ]; - let result = find_matching_option("A", &options); - assert!(result.is_some()); - let answer = result.unwrap(); - assert_eq!(answer.value, AnswerValue::Selected("A".to_string())); - } - - #[test] - fn review_target_line_shows_only_the_url() { - let target = fabro_types::ReviewTarget::new( - "Quarry review exercise", - "https://quarry.lithos.computer/tmp/0123456789abcdef0123456789abcdef", - fabro_types::ReviewTargetKind::Document, - ) - .unwrap(); - let mut question = Question::new(target.question_text(), QuestionType::MultipleChoice); - question.review_target = Some(target); - - assert_eq!( - review_target_line(&question).as_deref(), - Some( - "Review link: \ - https://quarry.lithos.computer/tmp/0123456789abcdef0123456789abcdef" - ) - ); - } - - #[test] - fn review_target_line_is_absent_without_a_target() { - let question = Question::new("Approve?", QuestionType::YesNo); - - assert_eq!(review_target_line(&question), None); - } - - #[test] - fn find_matching_option_by_key_case_insensitive() { - let options = vec![InterviewOption { - key: "Y".to_string(), - label: "Yes".to_string(), - description: None, - preview: None, - }]; - let result = find_matching_option("y", &options); - assert!(result.is_some()); - } - - #[test] - fn find_matching_option_by_index() { - let options = vec![ - InterviewOption { - key: "A".to_string(), - label: "Alpha".to_string(), - description: None, - preview: None, - }, - InterviewOption { - key: "B".to_string(), - label: "Beta".to_string(), - description: None, - preview: None, - }, - ]; - let result = find_matching_option("2", &options); - assert!(result.is_some()); - let answer = result.unwrap(); - assert_eq!(answer.value, AnswerValue::Selected("B".to_string())); - } - - #[test] - fn find_matching_option_no_match() { - let options = vec![InterviewOption { - key: "A".to_string(), - label: "Alpha".to_string(), - description: None, - preview: None, - }]; - let result = find_matching_option("zzz", &options); - assert!(result.is_none()); - } - - #[test] - fn find_matching_option_index_out_of_range() { - let options = vec![InterviewOption { - key: "A".to_string(), - label: "Alpha".to_string(), - description: None, - preview: None, - }]; - let result = find_matching_option("5", &options); - assert!(result.is_none()); - } - - #[test] - fn non_tty_multiple_choice_eof_returns_interrupted() { - let mut question = Question::new("Approve?", QuestionType::MultipleChoice); - question.options = vec![InterviewOption { - key: "A".to_string(), - label: "Approve".to_string(), - description: None, - preview: None, - }]; - - let answer = parse_non_tty_choice_response(&question, PromptRead::Eof); - assert_eq!(answer.value, AnswerValue::Interrupted); - } - - #[test] - fn non_tty_confirmation_invalid_response_returns_interrupted() { - let answer = parse_non_tty_confirm_response(PromptRead::Line(String::new())); - assert_eq!(answer.value, AnswerValue::Interrupted); - } - - #[test] - fn non_tty_freeform_blank_response_returns_interrupted() { - let answer = parse_non_tty_freeform_response(PromptRead::Line(" ".to_string())); - assert_eq!(answer.value, AnswerValue::Interrupted); - } -} diff --git a/lib/components/fabro-interview/src/lib.rs b/lib/components/fabro-interview/src/lib.rs index d98ca8e8f..10be3d737 100644 --- a/lib/components/fabro-interview/src/lib.rs +++ b/lib/components/fabro-interview/src/lib.rs @@ -1,18 +1,12 @@ mod auto_approve; -mod callback; -mod console; mod control; mod control_protocol; -mod queue; -mod recording; -mod replay; use std::collections::HashMap; use async_trait::async_trait; use fabro_types::{InterviewOption, Principal, QuestionType, ReviewTarget, SystemActorKind}; use serde::{Deserialize, Serialize}; -use tokio::time; /// A question presented to the user. #[derive(Debug, Clone, Serialize, Deserialize)] @@ -177,28 +171,14 @@ impl AnswerSubmission { } } -/// Apply timeout enforcement to an interviewer ask call. -/// Per spec 6.5: if `timeout_seconds` is set, returns default answer or -/// `Answer::timeout()`. -pub async fn ask_with_timeout( - interviewer: &dyn Interviewer, - question: Question, -) -> AnswerSubmission { - let timeout_secs = question.timeout_seconds; - let default_answer = question.default.clone(); - - if let Some(secs) = timeout_secs { - let duration = std::time::Duration::from_secs_f64(secs); - match time::timeout(duration, interviewer.ask(question)).await { - Ok(answer) => answer, - Err(_elapsed) => AnswerSubmission::system( - default_answer.unwrap_or_else(Answer::timeout), - SystemActorKind::Timeout, - ), - } - } else { - interviewer.ask(question).await - } +/// The line that points a reviewer at the question's review target, when it +/// has one. +#[must_use] +pub fn review_target_line(question: &Question) -> Option { + question + .review_target + .as_ref() + .map(|target| format!("Review link: {}", target.url())) } /// The interviewer trait for human-in-the-loop interactions. @@ -221,8 +201,6 @@ pub trait Interviewer: Send + Sync { // Re-export all implementors at the crate root pub use auto_approve::AutoApproveInterviewer; -pub use callback::CallbackInterviewer; -pub use console::{ConsoleInterviewer, review_target_line}; pub use control::{ControlInterviewer, SubmitError}; pub use control_protocol::{ WORKER_CONTROL_INVALID_CURSOR_REASON, WORKER_CONTROL_PONG_TIMEOUT_REASON, @@ -230,9 +208,6 @@ pub use control_protocol::{ WORKER_CONTROL_WS_PING_INTERVAL, WorkerControlAck, WorkerControlAnswer, WorkerControlDeliveryFrame, WorkerControlEnvelope, WorkerControlMessage, WorkerControlOutcome, }; -pub use queue::QueueInterviewer; -pub use recording::RecordingInterviewer; -pub use replay::ReplayInterviewer; #[cfg(test)] mod tests { @@ -363,47 +338,6 @@ mod tests { assert_eq!(q.question_type, QuestionType::MultiSelect); } - /// A slow interviewer that waits before answering -- for testing timeouts. - struct SlowInterviewer; - - #[async_trait] - impl Interviewer for SlowInterviewer { - async fn ask(&self, _question: Question) -> AnswerSubmission { - time::sleep(std::time::Duration::from_mins(1)).await; - AnswerSubmission::system(Answer::yes(), SystemActorKind::Engine) - } - } - - #[tokio::test] - async fn ask_with_timeout_returns_timeout_when_expired() { - let interviewer = SlowInterviewer; - let mut q = Question::new("approve?", QuestionType::YesNo); - q.timeout_seconds = Some(0.01); - - let answer = ask_with_timeout(&interviewer, q).await.answer; - assert_eq!(answer.value, AnswerValue::Timeout); - } - - #[tokio::test] - async fn ask_with_timeout_returns_default_when_set() { - let interviewer = SlowInterviewer; - let mut q = Question::new("approve?", QuestionType::YesNo); - q.timeout_seconds = Some(0.01); - q.default = Some(Answer::no()); - - let answer = ask_with_timeout(&interviewer, q).await.answer; - assert_eq!(answer.value, AnswerValue::No); - } - - #[tokio::test] - async fn ask_with_timeout_no_timeout_returns_normally() { - let interviewer = AutoApproveInterviewer::engine(); - let q = Question::new("approve?", QuestionType::YesNo); - - let answer = ask_with_timeout(&interviewer, q).await.answer; - assert_eq!(answer.value, AnswerValue::Yes); - } - #[tokio::test] async fn control_interviewer_routes_answers_by_question_id() { let interviewer = Arc::new(ControlInterviewer::new()); diff --git a/lib/components/fabro-interview/src/queue.rs b/lib/components/fabro-interview/src/queue.rs deleted file mode 100644 index 7289fc326..000000000 --- a/lib/components/fabro-interview/src/queue.rs +++ /dev/null @@ -1,81 +0,0 @@ -use std::collections::VecDeque; -use std::sync::Mutex; - -use async_trait::async_trait; -use fabro_types::{Principal, SystemActorKind}; - -use crate::{Answer, AnswerSubmission, Interviewer, Question}; - -/// Reads answers from a pre-filled queue. Returns Interrupted when empty. -pub struct QueueInterviewer { - answers: Mutex>, - actor: Principal, -} - -impl QueueInterviewer { - #[must_use] - pub fn new(answers: VecDeque) -> Self { - Self::with_actor(answers, Principal::System { - system_kind: SystemActorKind::Engine, - }) - } - - #[must_use] - pub fn with_actor(answers: VecDeque, actor: Principal) -> Self { - Self { - answers: Mutex::new(answers), - actor, - } - } -} - -#[async_trait] -impl Interviewer for QueueInterviewer { - async fn ask(&self, _question: Question) -> AnswerSubmission { - let mut queue = self.answers.lock().expect("queue lock poisoned"); - AnswerSubmission::new( - queue.pop_front().unwrap_or_else(Answer::interrupted), - self.actor.clone(), - ) - } -} - -#[cfg(test)] -mod tests { - use fabro_types::QuestionType; - - use super::*; - use crate::AnswerValue; - - #[tokio::test] - async fn returns_queued_answers_in_order() { - let answers = VecDeque::from([Answer::yes(), Answer::no()]); - let interviewer = QueueInterviewer::new(answers); - let q = Question::new("q1", QuestionType::YesNo); - - let a1 = interviewer.ask(q.clone()).await.answer; - assert_eq!(a1.value, AnswerValue::Yes); - - let a2 = interviewer.ask(q).await.answer; - assert_eq!(a2.value, AnswerValue::No); - } - - #[tokio::test] - async fn returns_interrupted_when_empty() { - let interviewer = QueueInterviewer::new(VecDeque::new()); - let q = Question::new("q", QuestionType::YesNo); - let answer = interviewer.ask(q).await.answer; - assert_eq!(answer.value, AnswerValue::Interrupted); - } - - #[tokio::test] - async fn returns_interrupted_after_exhausted() { - let answers = VecDeque::from([Answer::yes()]); - let interviewer = QueueInterviewer::new(answers); - let q = Question::new("q", QuestionType::YesNo); - - let _ = interviewer.ask(q.clone()).await; - let answer = interviewer.ask(q).await.answer; - assert_eq!(answer.value, AnswerValue::Interrupted); - } -} diff --git a/lib/components/fabro-interview/src/recording.rs b/lib/components/fabro-interview/src/recording.rs deleted file mode 100644 index 66a5a19cc..000000000 --- a/lib/components/fabro-interview/src/recording.rs +++ /dev/null @@ -1,215 +0,0 @@ -use std::path::Path; -use std::sync::Mutex; - -use async_trait::async_trait; - -use crate::{AnswerSubmission, Interviewer, Question}; - -/// Wraps another interviewer and records all question-answer pairs. -pub struct RecordingInterviewer { - inner: Box, - submissions: Mutex>, -} - -impl RecordingInterviewer { - #[must_use] - pub fn new(inner: Box) -> Self { - Self { - inner, - submissions: Mutex::new(Vec::new()), - } - } - - /// # Panics - /// Panics if the internal mutex is poisoned. - #[must_use] - pub fn recordings(&self) -> Vec<(Question, AnswerSubmission)> { - self.submissions - .lock() - .expect("recordings lock poisoned") - .clone() - } - - /// Serializes all recordings to a JSON string. - /// - /// # Errors - /// Returns an error if serialization fails. - pub fn to_json(&self) -> std::io::Result { - let recordings = self.recordings(); - serde_json::to_string_pretty(&recordings).map_err(std::io::Error::other) - } - - /// Deserializes recordings from a JSON string. - /// - /// # Errors - /// Returns an error if deserialization fails. - pub fn from_json(json: &str) -> std::io::Result> { - serde_json::from_str(json).map_err(std::io::Error::other) - } - - /// Saves recordings to a file as JSON. - /// - /// # Errors - /// Returns an error if serialization or file writing fails. - #[expect( - clippy::disallowed_methods, - reason = "sync helper for test-mode interview recording storage; not on a Tokio path" - )] - pub fn save_to_file(&self, path: &Path) -> std::io::Result<()> { - let json = self.to_json()?; - std::fs::write(path, json).map_err(|err| { - std::io::Error::new( - err.kind(), - format!("write interview recording {}: {err}", path.display()), - ) - })?; - Ok(()) - } - - /// Loads recordings from a JSON file. - /// - /// # Errors - /// Returns an error if file reading or deserialization fails. - #[expect( - clippy::disallowed_methods, - reason = "sync helper for test-mode interview recording storage; not on a Tokio path" - )] - pub fn load_from_file(path: &Path) -> std::io::Result> { - let json = std::fs::read_to_string(path).map_err(|err| { - std::io::Error::new( - err.kind(), - format!("read interview recording {}: {err}", path.display()), - ) - })?; - Self::from_json(&json) - } -} - -#[async_trait] -impl Interviewer for RecordingInterviewer { - async fn ask(&self, question: Question) -> AnswerSubmission { - let submission = self.inner.ask(question.clone()).await; - self.submissions - .lock() - .expect("recordings lock poisoned") - .push((question, submission.clone())); - submission - } -} - -#[cfg(test)] -mod tests { - use fabro_types::QuestionType; - - use super::*; - use crate::{AnswerValue, AutoApproveInterviewer}; - - #[tokio::test] - async fn records_question_answer_pairs() { - let inner = Box::new(AutoApproveInterviewer::engine()); - let recorder = RecordingInterviewer::new(inner); - - let q1 = Question::new("approve?", QuestionType::YesNo); - let q2 = Question::new("confirm?", QuestionType::Confirmation); - - let a1 = recorder.ask(q1).await.answer; - assert_eq!(a1.value, AnswerValue::Yes); - - let a2 = recorder.ask(q2).await.answer; - assert_eq!(a2.value, AnswerValue::Yes); - - let recs = recorder.recordings(); - assert_eq!(recs.len(), 2); - assert_eq!(recs[0].0.text, "approve?"); - assert_eq!(recs[1].0.text, "confirm?"); - } - - #[tokio::test] - async fn delegates_to_inner() { - let inner = Box::new(AutoApproveInterviewer::engine()); - let recorder = RecordingInterviewer::new(inner); - - let q = Question::new("text input", QuestionType::Freeform); - let answer = recorder.ask(q).await.answer; - assert_eq!(answer.value, AnswerValue::Text("auto-approved".to_string())); - } - - #[tokio::test] - async fn recordings_empty_initially() { - let inner = Box::new(AutoApproveInterviewer::engine()); - let recorder = RecordingInterviewer::new(inner); - assert!(recorder.recordings().is_empty()); - } - - #[tokio::test] - async fn to_json_serializes_recordings() { - let inner = Box::new(AutoApproveInterviewer::engine()); - let recorder = RecordingInterviewer::new(inner); - - let q = Question::new("approve?", QuestionType::YesNo); - recorder.ask(q).await; - - let json = recorder.to_json().unwrap(); - assert!(json.contains("approve?")); - assert!(json.contains("yes_no")); - } - - #[test] - fn from_json_deserializes_recordings() { - let json = r#"[ - [ - {"text":"approve?","question_type":"yes_no","options":[],"allow_freeform":false,"default":null,"timeout_seconds":null,"stage":"","metadata":{}}, - { - "answer":{"value":"Yes","selected_option":null,"text":null}, - "actor":{"kind":"system","system_kind":"engine"} - } - ] - ]"#; - - let recordings = RecordingInterviewer::from_json(json).unwrap(); - assert_eq!(recordings.len(), 1); - assert_eq!(recordings[0].0.text, "approve?"); - assert_eq!(recordings[0].1.answer.value, AnswerValue::Yes); - } - - #[tokio::test] - async fn save_to_file_and_load_from_file() { - let inner = Box::new(AutoApproveInterviewer::engine()); - let recorder = RecordingInterviewer::new(inner); - - let q = Question::new("approve?", QuestionType::YesNo); - recorder.ask(q).await; - - let dir = tempfile::tempdir().unwrap(); - let path = dir.path().join("recordings.json"); - - recorder.save_to_file(&path).unwrap(); - let loaded = RecordingInterviewer::load_from_file(&path).unwrap(); - - assert_eq!(loaded.len(), 1); - assert_eq!(loaded[0].0.text, "approve?"); - assert_eq!(loaded[0].1.answer.value, AnswerValue::Yes); - } - - #[tokio::test] - async fn round_trip_serialize_deserialize() { - let inner = Box::new(AutoApproveInterviewer::engine()); - let recorder = RecordingInterviewer::new(inner); - - let q1 = Question::new("approve?", QuestionType::YesNo); - let q2 = Question::new("confirm?", QuestionType::Confirmation); - recorder.ask(q1).await; - recorder.ask(q2).await; - - let json = recorder.to_json().unwrap(); - let restored = RecordingInterviewer::from_json(&json).unwrap(); - - assert_eq!(restored.len(), 2); - assert_eq!(restored[0].0.text, "approve?"); - assert_eq!(restored[0].0.question_type, QuestionType::YesNo); - assert_eq!(restored[0].1.answer.value, AnswerValue::Yes); - assert_eq!(restored[1].0.text, "confirm?"); - assert_eq!(restored[1].0.question_type, QuestionType::Confirmation); - assert_eq!(restored[1].1.answer.value, AnswerValue::Yes); - } -} diff --git a/lib/components/fabro-interview/src/replay.rs b/lib/components/fabro-interview/src/replay.rs deleted file mode 100644 index 5ba3f5b14..000000000 --- a/lib/components/fabro-interview/src/replay.rs +++ /dev/null @@ -1,107 +0,0 @@ -use std::collections::VecDeque; -use std::sync::Mutex; - -use async_trait::async_trait; -use fabro_types::SystemActorKind; - -use crate::{Answer, AnswerSubmission, Interviewer, Question}; - -/// Replays recorded answers in sequence. When recordings are exhausted, -/// returns `Answer::interrupted()`. -pub struct ReplayInterviewer { - submissions: Mutex>, -} - -impl ReplayInterviewer { - /// Creates a new `ReplayInterviewer` from recorded question-answer - /// submissions. - #[must_use] - pub fn new(recordings: Vec<(Question, AnswerSubmission)>) -> Self { - let submissions = recordings - .into_iter() - .map(|(_, submission)| submission) - .collect(); - Self { - submissions: Mutex::new(submissions), - } - } -} - -#[async_trait] -impl Interviewer for ReplayInterviewer { - async fn ask(&self, _question: Question) -> AnswerSubmission { - let mut submissions = self.submissions.lock().expect("answers lock poisoned"); - submissions.pop_front().unwrap_or_else(|| { - AnswerSubmission::system(Answer::interrupted(), SystemActorKind::Engine) - }) - } -} - -#[cfg(test)] -mod tests { - use fabro_types::{AuthMethod, IdpIdentity, Principal, QuestionType}; - - use super::*; - use crate::AnswerValue; - - #[tokio::test] - async fn replays_recorded_answers() { - let actor = Principal::user( - IdpIdentity::new("https://github.com", "12345").unwrap(), - "octocat".to_string(), - AuthMethod::Github, - ); - let recordings = vec![ - ( - Question::new("approve?", QuestionType::YesNo), - AnswerSubmission::new(Answer::yes(), actor.clone()), - ), - ( - Question::new("name?", QuestionType::Freeform), - AnswerSubmission::new(Answer::text("Alice"), actor.clone()), - ), - ]; - - let replayer = ReplayInterviewer::new(recordings); - - let s1 = replayer - .ask(Question::new("anything", QuestionType::YesNo)) - .await; - assert_eq!(s1.answer.value, AnswerValue::Yes); - assert_eq!(s1.actor, actor); - - let s2 = replayer - .ask(Question::new("anything", QuestionType::Freeform)) - .await; - assert_eq!(s2.answer.value, AnswerValue::Text("Alice".to_string())); - assert_eq!(s2.actor, actor); - } - - #[tokio::test] - async fn returns_interrupted_when_exhausted() { - let recordings = vec![( - Question::new("approve?", QuestionType::YesNo), - AnswerSubmission::system(Answer::yes(), SystemActorKind::Engine), - )]; - - let replayer = ReplayInterviewer::new(recordings); - - let a1 = replayer - .ask(Question::new("first", QuestionType::YesNo)) - .await - .answer; - assert_eq!(a1.value, AnswerValue::Yes); - - let a2 = replayer - .ask(Question::new("second", QuestionType::YesNo)) - .await - .answer; - assert_eq!(a2.value, AnswerValue::Interrupted); - - let a3 = replayer - .ask(Question::new("third", QuestionType::YesNo)) - .await - .answer; - assert_eq!(a3.value, AnswerValue::Interrupted); - } -} diff --git a/lib/components/fabro-workflow/README.md b/lib/components/fabro-workflow/README.md index 007051592..717fe6806 100644 --- a/lib/components/fabro-workflow/README.md +++ b/lib/components/fabro-workflow/README.md @@ -10,7 +10,7 @@ A DOT-based pipeline runner for multi-stage AI workflows. Define workflows as Gr - **Handler** -- An async trait implementation that executes a node and returns an `Outcome`. Built-in handlers include `StartHandler`, `ExitHandler`, `AgentHandler`, `PromptHandler`, `ConditionalHandler`, `HumanHandler`, `ParallelHandler`, `FanInHandler`, `CommandHandler`, and `SubWorkflowHandler`. - **Outcome** -- The result of executing a handler, carrying a `StageOutcome` (Success, Fail, PartialSuccess, Retry, Skipped), optional routing hints (`preferred_label`, `suggested_next_ids`), and context updates. - **Context** -- A thread-safe key-value store shared across pipeline stages, supporting snapshots and isolated cloning for parallel branches. -- **Interviewer** -- A trait for human-in-the-loop interactions. Implementations include `AutoApproveInterviewer`, `QueueInterviewer`, `CallbackInterviewer`, `ConsoleInterviewer`, and `RecordingInterviewer`. +- **Interviewer** -- A trait for human-in-the-loop interactions. Implementations include `AutoApproveInterviewer` and `ControlInterviewer`. - **Checkpoint** -- A serializable snapshot of execution state (completed nodes, context values) for crash recovery and resume. ## Pipeline Definition