mirror of
https://github.com/fabro-sh/fabro.git
synced 2026-10-01 02:04:24 +00:00
Delete the interviewers the Petri cutover left unused
ConsoleInterviewer, RecordingInterviewer, ReplayInterviewer, QueueInterviewer, CallbackInterviewer, and ask_with_timeout had no production caller once every run executes on Petri. review_target_line moves to lib.rs for the CLI's attach prompt. fabro-interview drops dialoguer and fabro-util. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
This commit is contained in:
parent
0e18253d7e
commit
c0fb71a467
9 changed files with 9 additions and 1002 deletions
3
Cargo.lock
generated
3
Cargo.lock
generated
|
|
@ -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",
|
||||
]
|
||||
|
|
|
|||
|
|
@ -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"
|
||||
|
|
|
|||
|
|
@ -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<dyn Fn(Question) -> 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()));
|
||||
}
|
||||
}
|
||||
|
|
@ -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<Answer> {
|
||||
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::<usize>() {
|
||||
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<String> {
|
||||
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<String> = 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::<String>::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<String> = 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<String> = 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::<String>::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);
|
||||
}
|
||||
}
|
||||
|
|
@ -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<String> {
|
||||
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());
|
||||
|
|
|
|||
|
|
@ -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<VecDeque<Answer>>,
|
||||
actor: Principal,
|
||||
}
|
||||
|
||||
impl QueueInterviewer {
|
||||
#[must_use]
|
||||
pub fn new(answers: VecDeque<Answer>) -> Self {
|
||||
Self::with_actor(answers, Principal::System {
|
||||
system_kind: SystemActorKind::Engine,
|
||||
})
|
||||
}
|
||||
|
||||
#[must_use]
|
||||
pub fn with_actor(answers: VecDeque<Answer>, 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);
|
||||
}
|
||||
}
|
||||
|
|
@ -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<dyn Interviewer>,
|
||||
submissions: Mutex<Vec<(Question, AnswerSubmission)>>,
|
||||
}
|
||||
|
||||
impl RecordingInterviewer {
|
||||
#[must_use]
|
||||
pub fn new(inner: Box<dyn Interviewer>) -> 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<String> {
|
||||
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<Vec<(Question, AnswerSubmission)>> {
|
||||
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<Vec<(Question, AnswerSubmission)>> {
|
||||
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);
|
||||
}
|
||||
}
|
||||
|
|
@ -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<VecDeque<AnswerSubmission>>,
|
||||
}
|
||||
|
||||
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);
|
||||
}
|
||||
}
|
||||
|
|
@ -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
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue