diff --git a/apps/fabro-web/app/components/interview-dock.test.tsx b/apps/fabro-web/app/components/interview-dock.test.tsx
index 34a58269f..6fb48bfb9 100644
--- a/apps/fabro-web/app/components/interview-dock.test.tsx
+++ b/apps/fabro-web/app/components/interview-dock.test.tsx
@@ -133,6 +133,27 @@ describe("InterviewDock", () => {
expect(buttons.Revise).toBeDefined();
});
+ test("multiple choice renders option descriptions as display text", () => {
+ const question = makeQuestion({
+ question_type: QuestionType.MULTIPLE_CHOICE,
+ options: [
+ {
+ key: "A",
+ label: "[A] Approve",
+ description: "Deploy the current patch",
+ preview: "not rendered specially",
+ },
+ ],
+ });
+ const tree = render(
+ ,
+ );
+ const text = textContent(tree.toJSON());
+ expect(text).toContain("Approve");
+ expect(text).toContain("Deploy the current patch");
+ expect(text).not.toContain("not rendered specially");
+ });
+
test("freeform question renders a textarea and disables send when empty", () => {
const question = makeQuestion({
question_type: QuestionType.FREEFORM,
diff --git a/apps/fabro-web/app/components/interview-dock.tsx b/apps/fabro-web/app/components/interview-dock.tsx
index 35eb2fcf5..967044907 100644
--- a/apps/fabro-web/app/components/interview-dock.tsx
+++ b/apps/fabro-web/app/components/interview-dock.tsx
@@ -15,7 +15,7 @@ import {
import { QuestionType } from "@qltysh/fabro-api-client";
import type {
ApiQuestion,
- ApiQuestionOption,
+ InterviewOption,
} from "@qltysh/fabro-api-client";
import {
@@ -270,7 +270,7 @@ function ChoiceBody({
submitting,
onSubmit,
}: {
- options: ApiQuestionOption[];
+ options: InterviewOption[];
allowFreeform: boolean;
submitting: boolean;
onSubmit: (answer: SubmitInterviewAnswer) => Promise;
@@ -287,7 +287,7 @@ function ChoiceBody({
onClick={() => void onSubmit({ kind: "selected", option_key: option.key })}
className={CHOICE_BUTTON}
>
- {displayLabel(option.label)}
+
))}
@@ -314,7 +314,7 @@ function MultiSelectBody({
submitting,
onSubmit,
}: {
- options: ApiQuestionOption[];
+ options: InterviewOption[];
submitting: boolean;
onSubmit: (answer: SubmitInterviewAnswer) => Promise;
}) {
@@ -346,7 +346,7 @@ function MultiSelectBody({
className={isSelected ? CHOICE_BUTTON_SELECTED : CHOICE_BUTTON}
>
{isSelected && }
- {displayLabel(option.label)}
+
);
})}
@@ -455,6 +455,19 @@ function FreeformBody({
);
}
+function OptionLabel({ option }: { option: InterviewOption }) {
+ return (
+
+ {displayLabel(option.label)}
+ {option.description && (
+
+ {option.description}
+
+ )}
+
+ );
+}
+
function Spinner() {
return ;
}
diff --git a/apps/fabro-web/app/components/stage-renderers/helpers.test.ts b/apps/fabro-web/app/components/stage-renderers/helpers.test.ts
index 4178c78d3..d5603d8bf 100644
--- a/apps/fabro-web/app/components/stage-renderers/helpers.test.ts
+++ b/apps/fabro-web/app/components/stage-renderers/helpers.test.ts
@@ -79,6 +79,36 @@ describe("parseHumanInterviewPairs", () => {
expect(pairs[0].resolution).toBeNull();
});
+ test("preserves option description and preview metadata from started events", () => {
+ const events: EventEnvelope[] = [
+ envelope(1, {
+ event: "interview.started",
+ properties: {
+ question_id: "q-1",
+ question: "Pick a path",
+ question_type: "multiple_choice",
+ options: [
+ {
+ key: "ship",
+ label: "Ship",
+ description: "Deploy the current patch",
+ preview: "diff preview",
+ },
+ ],
+ },
+ }),
+ ];
+
+ const pairs = parseHumanInterviewPairs(events);
+
+ expect(pairs[0].question.options[0]).toEqual({
+ key: "ship",
+ label: "Ship",
+ description: "Deploy the current patch",
+ preview: "diff preview",
+ });
+ });
+
test("captures timeout and interrupted resolutions", () => {
const events: EventEnvelope[] = [
envelope(1, {
diff --git a/apps/fabro-web/app/components/stage-renderers/helpers.ts b/apps/fabro-web/app/components/stage-renderers/helpers.ts
index 35af48b31..f46f0edcf 100644
--- a/apps/fabro-web/app/components/stage-renderers/helpers.ts
+++ b/apps/fabro-web/app/components/stage-renderers/helpers.ts
@@ -5,6 +5,8 @@ import { getArray, getNumber, getObject, getString, type UnknownRecord } from ".
export interface InterviewOption {
key: string;
label: string;
+ description?: string | null;
+ preview?: string | null;
}
export interface HumanQuestion {
@@ -54,7 +56,14 @@ function parseInterviewOptions(value: unknown): InterviewOption[] {
const record = item as UnknownRecord;
const key = getString(record, "key");
const label = getString(record, "label");
- if (key && label) out.push({ key, label });
+ if (key && label) {
+ const option: InterviewOption = { key, label };
+ const description = getString(record, "description");
+ const preview = getString(record, "preview");
+ if (description !== null) option.description = description;
+ if (preview !== null) option.preview = preview;
+ out.push(option);
+ }
}
return out;
}
diff --git a/apps/fabro-web/app/components/stage-renderers/human-qa.tsx b/apps/fabro-web/app/components/stage-renderers/human-qa.tsx
index 50dc48cdd..37464ce12 100644
--- a/apps/fabro-web/app/components/stage-renderers/human-qa.tsx
+++ b/apps/fabro-web/app/components/stage-renderers/human-qa.tsx
@@ -184,7 +184,14 @@ function QuestionBlock({
{option.key}
- {option.label}
+
+ {option.label}
+ {option.description && (
+
+ {option.description}
+
+ )}
+
))}
{question.allowFreeform && (
diff --git a/docs/public/api-reference/fabro-api.yaml b/docs/public/api-reference/fabro-api.yaml
index c47b053c3..3ed205878 100644
--- a/docs/public/api-reference/fabro-api.yaml
+++ b/docs/public/api-reference/fabro-api.yaml
@@ -7126,22 +7126,6 @@ components:
id:
type: string
- ApiQuestionOption:
- description: A selectable option for a multiple-choice or multi-select question.
- type: object
- required:
- - key
- - label
- properties:
- key:
- type: string
- description: Machine-readable option key used when submitting an answer.
- example: option_a
- label:
- type: string
- description: Human-readable label displayed to the user.
- example: Accept changes
-
ApiQuestion:
description: A pending human-in-the-loop question generated by a workflow stage.
type: object
@@ -7171,7 +7155,7 @@ components:
type: array
description: Available options for selection-based questions. Empty for freeform questions.
items:
- $ref: "#/components/schemas/ApiQuestionOption"
+ $ref: "#/components/schemas/InterviewOption"
allow_freeform:
type: boolean
description: Whether the user may provide freeform text in addition to selecting options.
@@ -8178,8 +8162,16 @@ components:
properties:
key:
type: string
+ description: Machine-readable option key used when submitting an answer.
label:
type: string
+ description: Human-readable label displayed to the user.
+ description:
+ type: ["string", "null"]
+ description: Optional untrusted model-authored option description for display.
+ preview:
+ type: ["string", "null"]
+ description: Optional untrusted model-authored option preview captured for clients.
InterviewQuestionRecord:
description: Storage shape of an interview question recorded in the event log.
diff --git a/lib/crates/fabro-agent/src/lib.rs b/lib/crates/fabro-agent/src/lib.rs
index 04a1c6fdf..010004694 100644
--- a/lib/crates/fabro-agent/src/lib.rs
+++ b/lib/crates/fabro-agent/src/lib.rs
@@ -15,6 +15,7 @@ pub mod loop_detection;
pub mod mcp_integration;
pub mod memory;
pub mod profiles;
+pub mod question_tools;
pub mod read_before_write_sandbox;
pub mod sandbox;
pub mod session;
@@ -46,6 +47,11 @@ pub use local_sandbox::LocalSandbox;
pub use loop_detection::detect_loop;
pub use memory::{MemoryDocument, discover_memory};
pub use profiles::{AnthropicProfile, EnvContext, GeminiProfile, OpenAiProfile};
+pub use question_tools::{
+ ANTHROPIC_ASK_USER_QUESTION_TOOL, AgentQuestion, AgentQuestionAnswer,
+ AgentQuestionAnswerStatus, AgentQuestionRuntime, AgentToolRuntime,
+ OPENAI_REQUEST_USER_INPUT_TOOL, register_question_tools,
+};
pub use read_before_write_sandbox::ReadBeforeWriteSandbox;
pub use sandbox::{
CommandOutputCallback, DirEntry, ExecResult, ExecStreamingResult, GrepOptions, Sandbox,
diff --git a/lib/crates/fabro-agent/src/question_tools.rs b/lib/crates/fabro-agent/src/question_tools.rs
new file mode 100644
index 000000000..d8351d67d
--- /dev/null
+++ b/lib/crates/fabro-agent/src/question_tools.rs
@@ -0,0 +1,588 @@
+//! Model-native tools that let a root workflow agent ask the human for input.
+
+use std::collections::BTreeMap;
+use std::future::Future;
+use std::sync::Arc;
+
+use async_trait::async_trait;
+use fabro_llm::types::ToolDefinition;
+use fabro_model::AgentProfileKind;
+use fabro_types::{InterviewOption, QuestionType};
+use serde::Deserialize;
+use serde_json::json;
+use tokio_util::sync::CancellationToken;
+
+use crate::tool_registry::{RegisteredTool, ToolContext, ToolRegistry};
+
+tokio::task_local! {
+ static CURRENT_AGENT_TOOL_RUNTIME: AgentToolRuntime;
+}
+
+pub const OPENAI_REQUEST_USER_INPUT_TOOL: &str = "request_user_input";
+pub const ANTHROPIC_ASK_USER_QUESTION_TOOL: &str = "AskUserQuestion";
+
+pub const OPTION_DESCRIPTION_MAX_CHARS: usize = 2_000;
+pub const OPTION_PREVIEW_MAX_CHARS: usize = 4_000;
+
+const ROOT_SESSION_REQUIRED_ERROR: &str =
+ "human-question tools are available only during a root workflow agent session";
+
+#[derive(Clone, Default)]
+pub struct AgentToolRuntime {
+ question_runtime: Option>,
+}
+
+impl AgentToolRuntime {
+ #[must_use]
+ pub fn new() -> Self {
+ Self::default()
+ }
+
+ #[must_use]
+ pub fn with_question_runtime(runtime: Arc) -> Self {
+ Self {
+ question_runtime: Some(runtime),
+ }
+ }
+
+ #[must_use]
+ pub fn question_runtime(&self) -> Option> {
+ self.question_runtime.clone()
+ }
+}
+
+pub async fn scope_agent_tool_runtime(runtime: AgentToolRuntime, future: F) -> F::Output
+where
+ F: Future,
+{
+ CURRENT_AGENT_TOOL_RUNTIME.scope(runtime, future).await
+}
+
+fn current_agent_tool_runtime() -> AgentToolRuntime {
+ CURRENT_AGENT_TOOL_RUNTIME
+ .try_with(Clone::clone)
+ .unwrap_or_default()
+}
+
+#[derive(Debug, Clone, PartialEq, Eq)]
+pub struct AgentQuestion {
+ pub original_id: Option,
+ pub original_question: String,
+ pub header: Option,
+ pub text: String,
+ pub question_type: QuestionType,
+ pub options: Vec,
+ pub allow_freeform: bool,
+}
+
+#[derive(Debug, Clone, Copy, PartialEq, Eq)]
+pub enum AgentQuestionAnswerStatus {
+ Answered,
+ Cancelled,
+ Interrupted,
+ Skipped,
+ Timeout,
+}
+
+#[derive(Debug, Clone, PartialEq, Eq)]
+pub struct AgentQuestionAnswer {
+ pub original_id: Option,
+ pub original_question: String,
+ pub answers: Vec,
+ pub status: AgentQuestionAnswerStatus,
+}
+
+#[async_trait]
+pub trait AgentQuestionRuntime: Send + Sync {
+ async fn ask_questions(
+ &self,
+ tool_call_id: &str,
+ questions: Vec,
+ cancel_token: CancellationToken,
+ ) -> Result, String>;
+}
+
+#[derive(Debug, Deserialize)]
+struct OpenAiQuestionToolArgs {
+ questions: Vec,
+}
+
+#[derive(Debug, Deserialize)]
+struct OpenAiQuestion {
+ id: String,
+ header: String,
+ question: String,
+ #[serde(default)]
+ options: Vec,
+}
+
+#[derive(Debug, Deserialize)]
+struct OpenAiOption {
+ label: String,
+ #[serde(default)]
+ description: Option,
+}
+
+#[derive(Debug, Deserialize)]
+struct AnthropicQuestionToolArgs {
+ questions: Vec,
+}
+
+#[derive(Debug, Deserialize)]
+#[serde(rename_all = "camelCase")]
+struct AnthropicQuestion {
+ question: String,
+ #[serde(default)]
+ header: Option,
+ #[serde(default)]
+ options: Vec,
+ #[serde(default)]
+ multi_select: bool,
+}
+
+#[derive(Debug, Deserialize)]
+struct AnthropicOption {
+ label: String,
+ #[serde(default)]
+ description: Option,
+ #[serde(default)]
+ preview: Option,
+}
+
+#[must_use]
+pub fn is_question_tool(name: &str) -> bool {
+ matches!(
+ name,
+ OPENAI_REQUEST_USER_INPUT_TOOL | ANTHROPIC_ASK_USER_QUESTION_TOOL
+ )
+}
+
+pub fn register_question_tools(profile_kind: AgentProfileKind, registry: &mut ToolRegistry) {
+ match profile_kind {
+ AgentProfileKind::OpenAi => registry.register(make_openai_question_tool()),
+ AgentProfileKind::Anthropic => registry.register(make_anthropic_question_tool()),
+ AgentProfileKind::Gemini => {}
+ }
+}
+
+fn make_openai_question_tool() -> RegisteredTool {
+ RegisteredTool {
+ definition: ToolDefinition {
+ name: OPENAI_REQUEST_USER_INPUT_TOOL.to_string(),
+ description: "Ask the human one or more questions and wait for their answers before continuing this stage.".to_string(),
+ parameters: json!({
+ "type": "object",
+ "required": ["questions"],
+ "properties": {
+ "questions": {
+ "type": "array",
+ "minItems": 1,
+ "items": {
+ "type": "object",
+ "required": ["id", "header", "question", "options"],
+ "properties": {
+ "id": { "type": "string" },
+ "header": { "type": "string" },
+ "question": { "type": "string" },
+ "options": {
+ "type": "array",
+ "items": {
+ "type": "object",
+ "required": ["label"],
+ "properties": {
+ "label": { "type": "string" },
+ "description": { "type": "string" }
+ }
+ }
+ }
+ }
+ }
+ }
+ }
+ }),
+ },
+ executor: Arc::new(|args, ctx| {
+ Box::pin(async move {
+ let parsed: OpenAiQuestionToolArgs = parse_tool_args(args)?;
+ let questions = normalize_openai_questions(parsed)?;
+ let answers = execute_question_tool(ctx, questions).await?;
+ format_openai_answers(&answers)
+ })
+ }),
+ }
+}
+
+fn make_anthropic_question_tool() -> RegisteredTool {
+ RegisteredTool {
+ definition: ToolDefinition {
+ name: ANTHROPIC_ASK_USER_QUESTION_TOOL.to_string(),
+ description: "Ask the human one or more questions and wait for their answers before continuing this stage.".to_string(),
+ parameters: json!({
+ "type": "object",
+ "required": ["questions"],
+ "properties": {
+ "questions": {
+ "type": "array",
+ "minItems": 1,
+ "items": {
+ "type": "object",
+ "required": ["question", "options", "multiSelect"],
+ "properties": {
+ "question": { "type": "string" },
+ "header": { "type": "string" },
+ "options": {
+ "type": "array",
+ "items": {
+ "type": "object",
+ "required": ["label"],
+ "properties": {
+ "label": { "type": "string" },
+ "description": { "type": "string" },
+ "preview": { "type": "string" }
+ }
+ }
+ },
+ "multiSelect": { "type": "boolean" }
+ }
+ }
+ }
+ }
+ }),
+ },
+ executor: Arc::new(|args, ctx| {
+ Box::pin(async move {
+ let parsed: AnthropicQuestionToolArgs = parse_tool_args(args)?;
+ let questions = normalize_anthropic_questions(parsed)?;
+ let answers = execute_question_tool(ctx, questions).await?;
+ format_anthropic_answers(&answers)
+ })
+ }),
+ }
+}
+
+fn parse_tool_args Deserialize<'de>>(args: serde_json::Value) -> Result {
+ serde_json::from_value(args).map_err(|err| format!("invalid question tool arguments: {err}"))
+}
+
+async fn execute_question_tool(
+ ctx: ToolContext,
+ questions: Vec,
+) -> Result, String> {
+ let session_id = ctx
+ .session_id
+ .as_deref()
+ .ok_or_else(|| ROOT_SESSION_REQUIRED_ERROR.to_string())?;
+ let root_session_id = ctx
+ .root_session_id
+ .as_deref()
+ .ok_or_else(|| ROOT_SESSION_REQUIRED_ERROR.to_string())?;
+ if session_id != root_session_id {
+ return Err(
+ "human-question tools are only available to the root agent; subagents must report back to their parent".to_string(),
+ );
+ }
+ let tool_call_id = ctx
+ .tool_call_id
+ .as_deref()
+ .ok_or_else(|| "human-question tool call is missing a provider tool_call_id".to_string())?;
+ let runtime = current_agent_tool_runtime().question_runtime().ok_or_else(|| {
+ "human-question tools are available only inside a workflow run with an active interviewer".to_string()
+ })?;
+ runtime
+ .ask_questions(tool_call_id, questions, ctx.cancel.clone())
+ .await
+}
+
+fn normalize_openai_questions(args: OpenAiQuestionToolArgs) -> Result, String> {
+ if args.questions.is_empty() {
+ return Err("questions must contain at least one question".to_string());
+ }
+ args.questions
+ .into_iter()
+ .map(|question| {
+ let original_question = question.question.trim().to_string();
+ Ok(AgentQuestion {
+ original_id: Some(non_empty(&question.id, "question id")?),
+ text: display_text(Some(question.header.as_str()), &question.question),
+ header: Some(question.header),
+ original_question,
+ question_type: QuestionType::MultipleChoice,
+ options: options_from_openai(question.options),
+ allow_freeform: true,
+ })
+ })
+ .collect()
+}
+
+fn normalize_anthropic_questions(
+ args: AnthropicQuestionToolArgs,
+) -> Result, String> {
+ if args.questions.is_empty() {
+ return Err("questions must contain at least one question".to_string());
+ }
+ args.questions
+ .into_iter()
+ .map(|question| {
+ let original_question = non_empty(&question.question, "question")?;
+ Ok(AgentQuestion {
+ original_id: None,
+ text: display_text(question.header.as_deref(), &question.question),
+ header: question.header,
+ original_question,
+ question_type: if question.multi_select {
+ QuestionType::MultiSelect
+ } else {
+ QuestionType::MultipleChoice
+ },
+ options: options_from_anthropic(question.options),
+ allow_freeform: true,
+ })
+ })
+ .collect()
+}
+
+fn options_from_openai(options: Vec) -> Vec {
+ options
+ .into_iter()
+ .enumerate()
+ .map(|(idx, option)| InterviewOption {
+ key: option_key(idx),
+ label: option.label,
+ description: option
+ .description
+ .map(|value| bounded_display_field(&value, OPTION_DESCRIPTION_MAX_CHARS)),
+ preview: None,
+ })
+ .collect()
+}
+
+fn options_from_anthropic(options: Vec) -> Vec {
+ options
+ .into_iter()
+ .enumerate()
+ .map(|(idx, option)| InterviewOption {
+ key: option_key(idx),
+ label: option.label,
+ description: option
+ .description
+ .map(|value| bounded_display_field(&value, OPTION_DESCRIPTION_MAX_CHARS)),
+ preview: option
+ .preview
+ .map(|value| bounded_display_field(&value, OPTION_PREVIEW_MAX_CHARS)),
+ })
+ .collect()
+}
+
+fn option_key(idx: usize) -> String {
+ format!("option_{}", idx + 1)
+}
+
+fn non_empty(value: &str, field: &str) -> Result {
+ let trimmed = value.trim();
+ if trimmed.is_empty() {
+ Err(format!("{field} must not be empty"))
+ } else {
+ Ok(trimmed.to_string())
+ }
+}
+
+fn display_text(header: Option<&str>, question: &str) -> String {
+ let header = header.map(str::trim).filter(|value| !value.is_empty());
+ let question = question.trim();
+ match (header, question.is_empty()) {
+ (Some(header), false) => format!("{header}\n\n{question}"),
+ (Some(header), true) => header.to_string(),
+ (None, false) => question.to_string(),
+ (None, true) => String::new(),
+ }
+}
+
+fn bounded_display_field(value: &str, max_chars: usize) -> String {
+ match value.char_indices().nth(max_chars) {
+ Some((byte_idx, _)) => value[..byte_idx].to_string(),
+ None => value.to_string(),
+ }
+}
+
+fn ensure_all_answered(answers: &[AgentQuestionAnswer]) -> Result<(), String> {
+ if let Some(answer) = answers
+ .iter()
+ .find(|answer| answer.status != AgentQuestionAnswerStatus::Answered)
+ {
+ return Err(format!(
+ "human-question request ended before the user answered `{}`: {}",
+ answer.original_question,
+ answer_status_label(answer.status)
+ ));
+ }
+ Ok(())
+}
+
+fn answer_status_label(status: AgentQuestionAnswerStatus) -> &'static str {
+ match status {
+ AgentQuestionAnswerStatus::Answered => "answered",
+ AgentQuestionAnswerStatus::Cancelled => "cancelled",
+ AgentQuestionAnswerStatus::Interrupted => "interrupted",
+ AgentQuestionAnswerStatus::Skipped => "skipped",
+ AgentQuestionAnswerStatus::Timeout => "timed out",
+ }
+}
+
+fn format_openai_answers(answers: &[AgentQuestionAnswer]) -> Result {
+ ensure_all_answered(answers)?;
+ let mut answer_map = BTreeMap::new();
+ for answer in answers {
+ let Some(original_id) = answer.original_id.as_ref() else {
+ return Err(
+ "OpenAI question answer is missing the original model question id".to_string(),
+ );
+ };
+ answer_map.insert(original_id.clone(), json!({ "answers": answer.answers }));
+ }
+ serde_json::to_string(&json!({ "answers": answer_map }))
+ .map_err(|err| format!("failed to serialize answers: {err}"))
+}
+
+fn format_anthropic_answers(answers: &[AgentQuestionAnswer]) -> Result {
+ ensure_all_answered(answers)?;
+ let pairs = answers
+ .iter()
+ .map(|answer| {
+ let question = json!(answer.original_question);
+ let answer_text = json!(answer.answers.join(", "));
+ format!("{question}={answer_text}")
+ })
+ .collect::>()
+ .join(", ");
+ Ok(format!(
+ "User has answered your questions: {pairs}. You can now continue with the task."
+ ))
+}
+
+#[cfg(test)]
+mod tests {
+ use super::*;
+
+ fn answered(
+ original_id: Option<&str>,
+ question: &str,
+ answers: &[&str],
+ ) -> AgentQuestionAnswer {
+ AgentQuestionAnswer {
+ original_id: original_id.map(str::to_string),
+ original_question: question.to_string(),
+ answers: answers.iter().map(|value| (*value).to_string()).collect(),
+ status: AgentQuestionAnswerStatus::Answered,
+ }
+ }
+
+ #[test]
+ fn openai_request_with_descriptions_normalizes_to_multiple_choice() {
+ let args: OpenAiQuestionToolArgs = serde_json::from_value(json!({
+ "questions": [{
+ "id": "q1",
+ "header": "Decision",
+ "question": "Which path?",
+ "options": [{ "label": "Ship", "description": "Deploy now" }]
+ }]
+ }))
+ .unwrap();
+
+ let questions = normalize_openai_questions(args).unwrap();
+
+ assert_eq!(questions.len(), 1);
+ assert_eq!(questions[0].original_id.as_deref(), Some("q1"));
+ assert_eq!(questions[0].question_type, QuestionType::MultipleChoice);
+ assert!(questions[0].allow_freeform);
+ assert_eq!(questions[0].text, "Decision\n\nWhich path?");
+ assert_eq!(questions[0].options[0].key, "option_1");
+ assert_eq!(questions[0].options[0].label, "Ship");
+ assert_eq!(
+ questions[0].options[0].description.as_deref(),
+ Some("Deploy now")
+ );
+ }
+
+ #[test]
+ fn anthropic_multiselect_preserves_preview_and_formats_comma_joined_answers() {
+ let args: AnthropicQuestionToolArgs = serde_json::from_value(json!({
+ "questions": [{
+ "header": "Pick features",
+ "question": "Which features?",
+ "multiSelect": true,
+ "options": [{
+ "label": "Auth",
+ "description": "Login support",
+ "preview": "auth diff"
+ }]
+ }]
+ }))
+ .unwrap();
+
+ let questions = normalize_anthropic_questions(args).unwrap();
+
+ assert_eq!(questions[0].question_type, QuestionType::MultiSelect);
+ assert_eq!(
+ questions[0].options[0].preview.as_deref(),
+ Some("auth diff")
+ );
+ let text =
+ format_anthropic_answers(&[answered(None, "Which features?", &["Auth", "Billing"])])
+ .unwrap();
+ assert!(text.contains("\"Which features?\"=\"Auth, Billing\""));
+ }
+
+ #[test]
+ fn openai_answers_are_keyed_by_original_model_question_id() {
+ let text = format_openai_answers(&[
+ answered(Some("first"), "First?", &["Yes"]),
+ answered(Some("second"), "Second?", &["No"]),
+ ])
+ .unwrap();
+
+ assert_eq!(
+ serde_json::from_str::(&text).unwrap(),
+ json!({
+ "answers": {
+ "first": { "answers": ["Yes"] },
+ "second": { "answers": ["No"] }
+ }
+ })
+ );
+ }
+
+ #[test]
+ fn option_description_and_preview_are_bounded() {
+ let long = "x".repeat(OPTION_PREVIEW_MAX_CHARS + 10);
+
+ assert_eq!(
+ bounded_display_field(&long, OPTION_DESCRIPTION_MAX_CHARS)
+ .chars()
+ .count(),
+ OPTION_DESCRIPTION_MAX_CHARS
+ );
+ assert_eq!(
+ bounded_display_field(&long, OPTION_PREVIEW_MAX_CHARS)
+ .chars()
+ .count(),
+ OPTION_PREVIEW_MAX_CHARS
+ );
+ }
+
+ #[test]
+ fn question_tool_registration_is_profile_specific() {
+ let mut openai = ToolRegistry::new();
+ register_question_tools(AgentProfileKind::OpenAi, &mut openai);
+ assert!(openai.get(OPENAI_REQUEST_USER_INPUT_TOOL).is_some());
+ assert!(openai.get(ANTHROPIC_ASK_USER_QUESTION_TOOL).is_none());
+
+ let mut anthropic = ToolRegistry::new();
+ register_question_tools(AgentProfileKind::Anthropic, &mut anthropic);
+ assert!(anthropic.get(ANTHROPIC_ASK_USER_QUESTION_TOOL).is_some());
+ assert!(anthropic.get(OPENAI_REQUEST_USER_INPUT_TOOL).is_none());
+
+ let mut gemini = ToolRegistry::new();
+ register_question_tools(AgentProfileKind::Gemini, &mut gemini);
+ assert!(gemini.names().is_empty());
+ }
+}
diff --git a/lib/crates/fabro-agent/src/session.rs b/lib/crates/fabro-agent/src/session.rs
index 71977df82..de822bbbd 100644
--- a/lib/crates/fabro-agent/src/session.rs
+++ b/lib/crates/fabro-agent/src/session.rs
@@ -32,6 +32,7 @@ use crate::history::History;
use crate::loop_detection::detect_loop;
use crate::memory::{BUDGET_BYTES, MemoryDocument, discover_memory};
use crate::profiles::EnvContext;
+use crate::question_tools::AgentToolRuntime;
use crate::sandbox::Sandbox;
use crate::skills::{
ExpandedInput, Skill, default_skill_dirs, discover_skills, expand_skill, make_use_skill_tool,
@@ -1129,6 +1130,15 @@ impl Session {
}
pub async fn process_input(&mut self, input: &str) -> Result<(), Error> {
+ self.process_input_with_runtime(input, AgentToolRuntime::default())
+ .await
+ }
+
+ pub async fn process_input_with_runtime(
+ &mut self,
+ input: &str,
+ agent_tool_runtime: AgentToolRuntime,
+ ) -> Result<(), Error> {
if self.state == SessionState::Closed {
return Err(Error::SessionClosed);
}
@@ -1152,7 +1162,7 @@ impl Session {
});
// Process the initial input, then drain any followups
- let mut result = self.run_single_input(input).await;
+ let mut result = self.run_single_input(input, &agent_tool_runtime).await;
if result.is_ok() {
loop {
@@ -1162,7 +1172,7 @@ impl Session {
.expect("followup queue lock poisoned")
.pop_front();
let Some(followup) = followup else { break };
- result = self.run_single_input(&followup).await;
+ result = self.run_single_input(&followup, &agent_tool_runtime).await;
if result.is_err() {
break;
}
@@ -1182,7 +1192,11 @@ impl Session {
result
}
- async fn run_single_input(&mut self, input: &str) -> Result<(), Error> {
+ async fn run_single_input(
+ &mut self,
+ input: &str,
+ agent_tool_runtime: &AgentToolRuntime,
+ ) -> Result<(), Error> {
const STREAM_CONSUME_RETRIES: usize = 3;
if self.state == SessionState::Closed {
@@ -1633,6 +1647,7 @@ impl Session {
&self.id,
&self.root_session_id,
self.tool_env_provider.as_ref(),
+ agent_tool_runtime,
)
.await;
composite_watcher.abort();
diff --git a/lib/crates/fabro-agent/src/tool_execution.rs b/lib/crates/fabro-agent/src/tool_execution.rs
index 939e6a4e9..c55b9ce3d 100644
--- a/lib/crates/fabro-agent/src/tool_execution.rs
+++ b/lib/crates/fabro-agent/src/tool_execution.rs
@@ -7,6 +7,7 @@ use tracing::debug;
use crate::config::{SessionOptions, ToolHookCallback, ToolHookDecision};
use crate::event::{Emitter, SessionBoundEmitter};
+use crate::question_tools::{self, AgentToolRuntime, is_question_tool};
use crate::sandbox::Sandbox;
use crate::session::ToolEnvProvider;
use crate::tool_registry::{AgentEventEmitter, RegisteredTool, ToolContext, ToolRegistry};
@@ -31,7 +32,25 @@ pub async fn execute_tool_calls(
session_id: &str,
root_session_id: &str,
tool_env_provider: Option<&Arc>,
+ agent_tool_runtime: &AgentToolRuntime,
) -> Vec {
+ if tool_calls.iter().any(|tc| is_question_tool(&tc.name)) {
+ return execute_question_tool_round(
+ tool_calls,
+ registry,
+ env,
+ tool_hooks,
+ cancel_token,
+ config,
+ emitter,
+ session_id,
+ root_session_id,
+ tool_env_provider,
+ agent_tool_runtime,
+ )
+ .await;
+ }
+
if parallel && tool_calls.len() > 1 {
execute_tool_calls_parallel(
tool_calls,
@@ -44,6 +63,7 @@ pub async fn execute_tool_calls(
session_id,
root_session_id,
tool_env_provider,
+ agent_tool_runtime,
)
.await
} else {
@@ -58,6 +78,7 @@ pub async fn execute_tool_calls(
session_id,
root_session_id,
tool_env_provider,
+ agent_tool_runtime,
)
.await
}
@@ -78,6 +99,7 @@ async fn execute_tool_calls_sequential(
session_id: &str,
root_session_id: &str,
tool_env_provider: Option<&Arc>,
+ agent_tool_runtime: &AgentToolRuntime,
) -> Vec {
let mut results = Vec::new();
for tc in tool_calls {
@@ -86,7 +108,7 @@ async fn execute_tool_calls_sequential(
continue;
}
- let result = execute_and_emit_one_tool(
+ let result = execute_and_emit_one_tool_with_runtime(
tc,
registry,
env.clone(),
@@ -97,6 +119,7 @@ async fn execute_tool_calls_sequential(
session_id,
root_session_id,
tool_env_provider,
+ agent_tool_runtime,
)
.await;
results.push(result);
@@ -119,8 +142,10 @@ async fn execute_tool_calls_parallel(
session_id: &str,
root_session_id: &str,
tool_env_provider: Option<&Arc>,
+ agent_tool_runtime: &AgentToolRuntime,
) -> Vec {
let tool_env_provider = tool_env_provider.cloned();
+ let agent_tool_runtime = agent_tool_runtime.clone();
let futures: Vec<_> = tool_calls
.iter()
.map(|tc| {
@@ -133,6 +158,7 @@ async fn execute_tool_calls_parallel(
let root_session_id = root_session_id.to_owned();
let tool_hooks = tool_hooks.cloned();
let tool_env_provider = tool_env_provider.clone();
+ let agent_tool_runtime = agent_tool_runtime.clone();
let access_denial = config.tool_access_denial_reason(&tc.name);
// Look up the tool before spawning since ToolRegistry is not Send.
let registered_tool = if access_denial.is_none() {
@@ -153,6 +179,7 @@ async fn execute_tool_calls_parallel(
&session_id,
&root_session_id,
tool_env_provider.as_ref(),
+ &agent_tool_runtime,
)
.await
}
@@ -162,6 +189,107 @@ async fn execute_tool_calls_parallel(
future::join_all(futures).await
}
+#[allow(
+ clippy::too_many_arguments,
+ reason = "Question-tool round handling needs the same execution context as normal tool dispatch."
+)]
+async fn execute_question_tool_round(
+ tool_calls: &[ToolCall],
+ registry: &ToolRegistry,
+ env: Arc,
+ tool_hooks: Option<&Arc>,
+ cancel_token: &CancellationToken,
+ config: &SessionOptions,
+ emitter: &Emitter,
+ session_id: &str,
+ root_session_id: &str,
+ tool_env_provider: Option<&Arc>,
+ agent_tool_runtime: &AgentToolRuntime,
+) -> Vec {
+ let first_question_index = tool_calls
+ .iter()
+ .position(|tc| is_question_tool(&tc.name))
+ .expect("question-tool round should contain a question tool");
+ let mut results = Vec::with_capacity(tool_calls.len());
+
+ for (index, tc) in tool_calls.iter().enumerate() {
+ if cancel_token.is_cancelled() {
+ results.push(ToolResult::error(tc.id.clone(), "Cancelled"));
+ continue;
+ }
+
+ if index == first_question_index {
+ results.push(
+ execute_and_emit_one_tool_with_runtime(
+ tc,
+ registry,
+ env.clone(),
+ tool_hooks,
+ cancel_token.child_token(),
+ config,
+ emitter,
+ session_id,
+ root_session_id,
+ tool_env_provider,
+ agent_tool_runtime,
+ )
+ .await,
+ );
+ } else if is_question_tool(&tc.name) {
+ results.push(error_tool_result_with_events(
+ tc,
+ emitter,
+ session_id,
+ config,
+ "Only one human-question tool call may be used in a tool round. Combine all questions into a single questions[] batch and call the question tool once.",
+ ));
+ } else {
+ results.push(error_tool_result_with_events(
+ tc,
+ emitter,
+ session_id,
+ config,
+ "This tool call was not executed because human-question tools must run alone in a tool round. Retry non-question tools in a later round after the user answers.",
+ ));
+ }
+ }
+
+ results
+}
+
+fn error_tool_result_with_events(
+ tc: &ToolCall,
+ emitter: &Emitter,
+ session_id: &str,
+ config: &SessionOptions,
+ message: &str,
+) -> ToolResult {
+ emit_tool_call_started(emitter, session_id, tc);
+ let result = ToolResult::error(&tc.id, message);
+ emit_tool_call_result(emitter, session_id, tc, &result);
+ truncate_tool_result(&result, &tc.name, config)
+}
+
+fn emit_tool_call_started(emitter: &Emitter, session_id: &str, tc: &ToolCall) {
+ emitter.emit(session_id.to_owned(), AgentEvent::ToolCallStarted {
+ tool_name: tc.name.clone(),
+ tool_call_id: tc.id.clone(),
+ arguments: tc.arguments.clone(),
+ });
+}
+
+fn emit_tool_call_result(emitter: &Emitter, session_id: &str, tc: &ToolCall, result: &ToolResult) {
+ emitter.emit(session_id.to_owned(), AgentEvent::ToolCallOutputDelta {
+ delta: result.content.to_string(),
+ });
+ emitter.emit(session_id.to_owned(), AgentEvent::ToolCallCompleted {
+ tool_name: tc.name.clone(),
+ tool_call_id: tc.id.clone(),
+ output: result.content.clone(),
+ is_error: result.is_error,
+ });
+}
+
/// Execute a single tool call with event emission and output truncation.
#[allow(
clippy::too_many_arguments,
@@ -178,6 +306,39 @@ pub async fn execute_and_emit_one_tool(
session_id: &str,
root_session_id: &str,
tool_env_provider: Option<&Arc>,
+) -> ToolResult {
+ execute_and_emit_one_tool_with_runtime(
+ tc,
+ registry,
+ env,
+ tool_hooks,
+ cancel_token,
+ config,
+ emitter,
+ session_id,
+ root_session_id,
+ tool_env_provider,
+ &AgentToolRuntime::default(),
+ )
+ .await
+}
+
+#[allow(
+ clippy::too_many_arguments,
+ reason = "Single-tool execution needs the tool, runtime handles, and emission context."
+)]
+async fn execute_and_emit_one_tool_with_runtime(
+ tc: &ToolCall,
+ registry: &ToolRegistry,
+ env: Arc,
+ tool_hooks: Option<&Arc>,
+ cancel_token: CancellationToken,
+ config: &SessionOptions,
+ emitter: &Emitter,
+ session_id: &str,
+ root_session_id: &str,
+ tool_env_provider: Option<&Arc>,
+ agent_tool_runtime: &AgentToolRuntime,
) -> ToolResult {
let access_denial = config.tool_access_denial_reason(&tc.name);
let registered_tool = if access_denial.is_none() {
@@ -197,6 +358,7 @@ pub async fn execute_and_emit_one_tool(
session_id,
root_session_id,
tool_env_provider,
+ agent_tool_runtime,
)
.await
}
@@ -219,26 +381,13 @@ async fn execute_and_emit_one_tool_with_lookup(
session_id: &str,
root_session_id: &str,
tool_env_provider: Option<&Arc>,
+ agent_tool_runtime: &AgentToolRuntime,
) -> ToolResult {
- emitter.emit(session_id.to_owned(), AgentEvent::ToolCallStarted {
- tool_name: tc.name.clone(),
- tool_call_id: tc.id.clone(),
- arguments: tc.arguments.clone(),
- });
+ emit_tool_call_started(emitter, session_id, tc);
if let Some(reason) = access_denial {
let result = ToolResult::error(&tc.id, &reason);
-
- emitter.emit(session_id.to_owned(), AgentEvent::ToolCallOutputDelta {
- delta: result.content.to_string(),
- });
- emitter.emit(session_id.to_owned(), AgentEvent::ToolCallCompleted {
- tool_name: tc.name.clone(),
- tool_call_id: tc.id.clone(),
- output: result.content.clone(),
- is_error: true,
- });
-
+ emit_tool_call_result(emitter, session_id, tc, &result);
return truncate_tool_result(&result, &tc.name, config);
}
@@ -252,17 +401,7 @@ async fn execute_and_emit_one_tool_with_lookup(
if let ToolHookDecision::Block { reason } = decision {
let result = ToolResult::error(&tc.id, &reason);
-
- emitter.emit(session_id.to_owned(), AgentEvent::ToolCallOutputDelta {
- delta: result.content.to_string(),
- });
- emitter.emit(session_id.to_owned(), AgentEvent::ToolCallCompleted {
- tool_name: tc.name.clone(),
- tool_call_id: tc.id.clone(),
- output: result.content.clone(),
- is_error: true,
- });
-
+ emit_tool_call_result(emitter, session_id, tc, &result);
return truncate_tool_result(&result, &tc.name, config);
}
}
@@ -276,19 +415,11 @@ async fn execute_and_emit_one_tool_with_lookup(
session_id,
root_session_id,
tool_env_provider,
+ agent_tool_runtime,
)
.await;
- emitter.emit(session_id.to_owned(), AgentEvent::ToolCallOutputDelta {
- delta: result.content.to_string(),
- });
-
- emitter.emit(session_id.to_owned(), AgentEvent::ToolCallCompleted {
- tool_name: tc.name.clone(),
- tool_call_id: tc.id.clone(),
- output: result.content.clone(),
- is_error: result.is_error,
- });
+ emit_tool_call_result(emitter, session_id, tc, &result);
// Post-tool-use hooks
if let Some(hooks) = tool_hooks {
@@ -329,6 +460,7 @@ async fn execute_one_tool(
session_id: &str,
root_session_id: &str,
tool_env_provider: Option<&Arc>,
+ agent_tool_runtime: &AgentToolRuntime,
) -> ToolResult {
match registered_tool {
Some(tool) => {
@@ -355,7 +487,10 @@ async fn execute_one_tool(
tool_call_id: Some(tc.id.clone()),
agent_event_emitter,
};
- match (tool.executor)(tc.arguments.clone(), ctx).await {
+ let execution = (tool.executor)(tc.arguments.clone(), ctx);
+ match question_tools::scope_agent_tool_runtime(agent_tool_runtime.clone(), execution)
+ .await
+ {
Ok(output) => ToolResult::success(&tc.id, serde_json::json!(output)),
Err(err) => ToolResult::error(&tc.id, err),
}
@@ -418,9 +553,11 @@ pub fn validate_tool_args(
#[cfg(test)]
mod tests {
use std::collections::HashMap;
- use std::sync::Mutex;
+ use std::sync::{Arc, Mutex};
+ use async_trait::async_trait;
use fabro_llm::types::{ToolCall, ToolDefinition};
+ use fabro_model::AgentProfileKind;
use super::*;
use crate::config::{
@@ -428,6 +565,10 @@ mod tests {
};
use crate::event::Emitter;
use crate::local_sandbox::LocalSandbox;
+ use crate::question_tools::{
+ AgentQuestion, AgentQuestionAnswer, AgentQuestionAnswerStatus, AgentQuestionRuntime,
+ AgentToolRuntime, register_question_tools,
+ };
use crate::read_before_write_sandbox::ReadBeforeWriteSandbox;
use crate::test_support::MutableMockSandbox;
use crate::tool_registry::{RegisteredTool, ToolContext, ToolRegistry};
@@ -505,6 +646,125 @@ mod tests {
}
}
+ struct StubQuestionRuntime;
+
+ #[async_trait]
+ impl AgentQuestionRuntime for StubQuestionRuntime {
+ async fn ask_questions(
+ &self,
+ _tool_call_id: &str,
+ questions: Vec,
+ _cancel_token: CancellationToken,
+ ) -> Result, String> {
+ Ok(questions
+ .into_iter()
+ .map(|question| AgentQuestionAnswer {
+ original_id: question.original_id,
+ original_question: question.original_question,
+ answers: vec!["Ship".to_string()],
+ status: AgentQuestionAnswerStatus::Answered,
+ })
+ .collect())
+ }
+ }
+
+ #[tokio::test]
+ async fn question_tool_round_rejects_non_question_peers_and_preserves_order() {
+ let mut registry = ToolRegistry::new();
+ register_question_tools(AgentProfileKind::OpenAi, &mut registry);
+ registry.register(make_echo_tool());
+ let tool_calls = vec![
+ make_tool_call(
+ "request_user_input",
+ "call_question",
+ serde_json::json!({
+ "questions": [{
+ "id": "q1",
+ "header": "Decision",
+ "question": "Ship it?",
+ "options": [{ "label": "Ship" }]
+ }]
+ }),
+ ),
+ make_tool_call("echo", "call_echo", serde_json::json!({"text": "hello"})),
+ ];
+ let runtime = AgentToolRuntime::with_question_runtime(Arc::new(StubQuestionRuntime));
+
+ let results = execute_tool_calls(
+ &tool_calls,
+ true,
+ ®istry,
+ Arc::new(LocalSandbox::new(std::env::current_dir().unwrap())),
+ None,
+ &CancellationToken::new(),
+ &SessionOptions::default(),
+ &Emitter::new(),
+ "root",
+ "root",
+ None,
+ &runtime,
+ )
+ .await;
+
+ assert_eq!(results.len(), 2);
+ assert_eq!(results[0].tool_call_id, "call_question");
+ assert!(!results[0].is_error);
+ assert_eq!(results[1].tool_call_id, "call_echo");
+ assert!(results[1].is_error);
+ assert!(
+ results[1]
+ .content
+ .as_str()
+ .unwrap()
+ .contains("human-question tools must run alone")
+ );
+ }
+
+ #[tokio::test]
+ async fn multiple_question_tool_calls_execute_only_first() {
+ let mut registry = ToolRegistry::new();
+ register_question_tools(AgentProfileKind::OpenAi, &mut registry);
+ let question_args = serde_json::json!({
+ "questions": [{
+ "id": "q1",
+ "header": "Decision",
+ "question": "Ship it?",
+ "options": [{ "label": "Ship" }]
+ }]
+ });
+ let tool_calls = vec![
+ make_tool_call("request_user_input", "call_first", question_args.clone()),
+ make_tool_call("request_user_input", "call_second", question_args),
+ ];
+ let runtime = AgentToolRuntime::with_question_runtime(Arc::new(StubQuestionRuntime));
+
+ let results = execute_tool_calls(
+ &tool_calls,
+ true,
+ ®istry,
+ Arc::new(LocalSandbox::new(std::env::current_dir().unwrap())),
+ None,
+ &CancellationToken::new(),
+ &SessionOptions::default(),
+ &Emitter::new(),
+ "root",
+ "root",
+ None,
+ &runtime,
+ )
+ .await;
+
+ assert!(!results[0].is_error);
+ assert!(results[1].is_error);
+ assert!(
+ results[1]
+ .content
+ .as_str()
+ .unwrap()
+ .contains("Combine all questions into a single questions[] batch")
+ );
+ }
+
struct MockHookCallback {
pre_decision: ToolHookDecision,
post_calls: Arc>>,
diff --git a/lib/crates/fabro-api/tests/interview_option_round_trip.rs b/lib/crates/fabro-api/tests/interview_option_round_trip.rs
index 40fb51b23..06f4cff82 100644
--- a/lib/crates/fabro-api/tests/interview_option_round_trip.rs
+++ b/lib/crates/fabro-api/tests/interview_option_round_trip.rs
@@ -13,7 +13,9 @@ fn interview_option_reuses_canonical_type() {
fn interview_option_round_trips_representative_json() {
let value = json!({
"key": "approve",
- "label": "Approve"
+ "label": "Approve",
+ "description": "Approve the proposed changes.",
+ "preview": "diff --stat output"
});
let option: InterviewOption = serde_json::from_value(value.clone()).unwrap();
diff --git a/lib/crates/fabro-api/tests/interview_question_record_round_trip.rs b/lib/crates/fabro-api/tests/interview_question_record_round_trip.rs
index 8e8f743b6..14338c4b4 100644
--- a/lib/crates/fabro-api/tests/interview_question_record_round_trip.rs
+++ b/lib/crates/fabro-api/tests/interview_question_record_round_trip.rs
@@ -17,7 +17,12 @@ fn interview_question_record_round_trips_representative_json() {
"stage": "gate",
"question_type": "multiple_choice",
"options": [
- { "key": "approve", "label": "Approve" },
+ {
+ "key": "approve",
+ "label": "Approve",
+ "description": "Deploy now",
+ "preview": "deploy --prod"
+ },
{ "key": "reject", "label": "Reject" }
],
"allow_freeform": true,
diff --git a/lib/crates/fabro-api/tests/pending_interview_record_round_trip.rs b/lib/crates/fabro-api/tests/pending_interview_record_round_trip.rs
index 9ce7c4c9c..dae81fdc3 100644
--- a/lib/crates/fabro-api/tests/pending_interview_record_round_trip.rs
+++ b/lib/crates/fabro-api/tests/pending_interview_record_round_trip.rs
@@ -18,7 +18,12 @@ fn pending_interview_record_round_trips_populated_question() {
"stage": "gate",
"question_type": "multiple_choice",
"options": [
- { "key": "approve", "label": "Approve" },
+ {
+ "key": "approve",
+ "label": "Approve",
+ "description": "Deploy now",
+ "preview": "deploy --prod"
+ },
{ "key": "reject", "label": "Reject" }
],
"allow_freeform": true,
diff --git a/lib/crates/fabro-api/tests/run_projection_round_trip.rs b/lib/crates/fabro-api/tests/run_projection_round_trip.rs
index 00aeb7ae2..64a00df91 100644
--- a/lib/crates/fabro-api/tests/run_projection_round_trip.rs
+++ b/lib/crates/fabro-api/tests/run_projection_round_trip.rs
@@ -54,7 +54,12 @@ fn run_projection_round_trips_populated_projection() {
"stage": "gate",
"question_type": "multiple_choice",
"options": [
- { "key": "approve", "label": "Approve" },
+ {
+ "key": "approve",
+ "label": "Approve",
+ "description": "Deploy now",
+ "preview": "deploy --prod"
+ },
{ "key": "reject", "label": "Reject" }
],
"allow_freeform": true,
diff --git a/lib/crates/fabro-cli/src/commands/run/attach.rs b/lib/crates/fabro-cli/src/commands/run/attach.rs
index d66249aad..c8ad2fff2 100644
--- a/lib/crates/fabro-cli/src/commands/run/attach.rs
+++ b/lib/crates/fabro-cli/src/commands/run/attach.rs
@@ -426,8 +426,10 @@ fn api_question_to_question(question: &types::ApiQuestion) -> Question {
.options
.iter()
.map(|option| InterviewOption {
- key: option.key.clone(),
- label: option.label.clone(),
+ key: option.key.clone(),
+ label: option.label.clone(),
+ description: option.description.clone(),
+ preview: option.preview.clone(),
})
.collect();
converted.allow_freeform = question.allow_freeform;
@@ -966,8 +968,10 @@ mod tests {
fn invalid_multiple_choice_without_freeform_is_user_correctable() {
let mut question = Question::new("Pick one.", QuestionType::MultipleChoice);
question.options = vec![InterviewOption {
- key: "A".to_string(),
- label: "Approve".to_string(),
+ key: "A".to_string(),
+ label: "Approve".to_string(),
+ description: None,
+ preview: None,
}];
let response = parse_choice_response(&question, PromptRead::Line("bogus".to_string()));
@@ -979,8 +983,10 @@ mod tests {
fn unmatched_multiple_choice_with_freeform_remains_text() {
let mut question = Question::new("Pick one.", QuestionType::MultipleChoice);
question.options = vec![InterviewOption {
- key: "A".to_string(),
- label: "Approve".to_string(),
+ key: "A".to_string(),
+ label: "Approve".to_string(),
+ description: None,
+ preview: None,
}];
question.allow_freeform = true;
@@ -1000,12 +1006,16 @@ mod tests {
let mut question = Question::new("Pick many.", QuestionType::MultiSelect);
question.options = vec![
InterviewOption {
- key: "A".to_string(),
- label: "Approve".to_string(),
+ key: "A".to_string(),
+ label: "Approve".to_string(),
+ description: None,
+ preview: None,
},
InterviewOption {
- key: "N".to_string(),
- label: "Notify".to_string(),
+ key: "N".to_string(),
+ label: "Notify".to_string(),
+ description: None,
+ preview: None,
},
];
diff --git a/lib/crates/fabro-interview/src/auto_approve.rs b/lib/crates/fabro-interview/src/auto_approve.rs
index faeba9a7b..66a323548 100644
--- a/lib/crates/fabro-interview/src/auto_approve.rs
+++ b/lib/crates/fabro-interview/src/auto_approve.rs
@@ -72,12 +72,16 @@ mod tests {
let mut q = Question::new("Choose:", QuestionType::MultipleChoice);
q.options = vec![
InterviewOption {
- key: "A".to_string(),
- label: "Alpha".to_string(),
+ key: "A".to_string(),
+ label: "Alpha".to_string(),
+ description: None,
+ preview: None,
},
InterviewOption {
- key: "B".to_string(),
- label: "Beta".to_string(),
+ key: "B".to_string(),
+ label: "Beta".to_string(),
+ description: None,
+ preview: None,
},
];
let answer = interviewer.ask(q).await.answer;
@@ -85,8 +89,10 @@ mod tests {
assert_eq!(
answer.selected_option,
Some(InterviewOption {
- key: "A".to_string(),
- label: "Alpha".to_string(),
+ key: "A".to_string(),
+ label: "Alpha".to_string(),
+ description: None,
+ preview: None,
})
);
}
diff --git a/lib/crates/fabro-interview/src/console.rs b/lib/crates/fabro-interview/src/console.rs
index 437025041..0005b1aca 100644
--- a/lib/crates/fabro-interview/src/console.rs
+++ b/lib/crates/fabro-interview/src/console.rs
@@ -296,12 +296,16 @@ mod tests {
fn find_matching_option_by_key() {
let options = vec![
InterviewOption {
- key: "A".to_string(),
- label: "Approve".to_string(),
+ key: "A".to_string(),
+ label: "Approve".to_string(),
+ description: None,
+ preview: None,
},
InterviewOption {
- key: "R".to_string(),
- label: "Reject".to_string(),
+ key: "R".to_string(),
+ label: "Reject".to_string(),
+ description: None,
+ preview: None,
},
];
let result = find_matching_option("A", &options);
@@ -313,8 +317,10 @@ mod tests {
#[test]
fn find_matching_option_by_key_case_insensitive() {
let options = vec![InterviewOption {
- key: "Y".to_string(),
- label: "Yes".to_string(),
+ key: "Y".to_string(),
+ label: "Yes".to_string(),
+ description: None,
+ preview: None,
}];
let result = find_matching_option("y", &options);
assert!(result.is_some());
@@ -324,12 +330,16 @@ mod tests {
fn find_matching_option_by_index() {
let options = vec![
InterviewOption {
- key: "A".to_string(),
- label: "Alpha".to_string(),
+ key: "A".to_string(),
+ label: "Alpha".to_string(),
+ description: None,
+ preview: None,
},
InterviewOption {
- key: "B".to_string(),
- label: "Beta".to_string(),
+ key: "B".to_string(),
+ label: "Beta".to_string(),
+ description: None,
+ preview: None,
},
];
let result = find_matching_option("2", &options);
@@ -341,8 +351,10 @@ mod tests {
#[test]
fn find_matching_option_no_match() {
let options = vec![InterviewOption {
- key: "A".to_string(),
- label: "Alpha".to_string(),
+ key: "A".to_string(),
+ label: "Alpha".to_string(),
+ description: None,
+ preview: None,
}];
let result = find_matching_option("zzz", &options);
assert!(result.is_none());
@@ -351,8 +363,10 @@ mod tests {
#[test]
fn find_matching_option_index_out_of_range() {
let options = vec![InterviewOption {
- key: "A".to_string(),
- label: "Alpha".to_string(),
+ key: "A".to_string(),
+ label: "Alpha".to_string(),
+ description: None,
+ preview: None,
}];
let result = find_matching_option("5", &options);
assert!(result.is_none());
@@ -362,8 +376,10 @@ mod tests {
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(),
+ key: "A".to_string(),
+ label: "Approve".to_string(),
+ description: None,
+ preview: None,
}];
let answer = parse_non_tty_choice_response(&question, PromptRead::Eof);
diff --git a/lib/crates/fabro-interview/src/lib.rs b/lib/crates/fabro-interview/src/lib.rs
index c242cf14a..4c050c4f7 100644
--- a/lib/crates/fabro-interview/src/lib.rs
+++ b/lib/crates/fabro-interview/src/lib.rs
@@ -301,8 +301,10 @@ mod tests {
#[test]
fn answer_selected() {
let opt = InterviewOption {
- key: "A".to_string(),
- label: "Approve".to_string(),
+ key: "A".to_string(),
+ label: "Approve".to_string(),
+ description: None,
+ preview: None,
};
let a = Answer::selected("A", opt.clone());
assert_eq!(a.value, AnswerValue::Selected("A".to_string()));
@@ -319,12 +321,16 @@ mod tests {
#[test]
fn question_option_eq() {
let a = InterviewOption {
- key: "Y".to_string(),
- label: "Yes".to_string(),
+ key: "Y".to_string(),
+ label: "Yes".to_string(),
+ description: None,
+ preview: None,
};
let b = InterviewOption {
- key: "Y".to_string(),
- label: "Yes".to_string(),
+ key: "Y".to_string(),
+ label: "Yes".to_string(),
+ description: None,
+ preview: None,
};
assert_eq!(a, b);
}
diff --git a/lib/crates/fabro-server/src/demo/mod.rs b/lib/crates/fabro-server/src/demo/mod.rs
index 4bbf09241..f13adafd5 100644
--- a/lib/crates/fabro-server/src/demo/mod.rs
+++ b/lib/crates/fabro-server/src/demo/mod.rs
@@ -1660,13 +1660,17 @@ mod runs {
stage: "review".into(),
question_type: QuestionType::YesNo,
options: vec![
- ApiQuestionOption {
- key: "yes".into(),
- label: "Yes".into(),
+ InterviewOption {
+ key: "yes".into(),
+ label: "Yes".into(),
+ description: None,
+ preview: None,
},
- ApiQuestionOption {
- key: "no".into(),
- label: "No".into(),
+ InterviewOption {
+ key: "no".into(),
+ label: "No".into(),
+ description: None,
+ preview: None,
},
],
allow_freeform: false,
@@ -1679,13 +1683,17 @@ mod runs {
stage: "migration".into(),
question_type: QuestionType::MultipleChoice,
options: vec![
- ApiQuestionOption {
- key: "incremental".into(),
- label: "Incremental migration".into(),
+ InterviewOption {
+ key: "incremental".into(),
+ label: "Incremental migration".into(),
+ description: None,
+ preview: None,
},
- ApiQuestionOption {
- key: "big_bang".into(),
- label: "Big-bang rewrite".into(),
+ InterviewOption {
+ key: "big_bang".into(),
+ label: "Big-bang rewrite".into(),
+ description: None,
+ preview: None,
},
],
allow_freeform: true,
diff --git a/lib/crates/fabro-server/src/server.rs b/lib/crates/fabro-server/src/server.rs
index 0063fadbf..cd67951ea 100644
--- a/lib/crates/fabro-server/src/server.rs
+++ b/lib/crates/fabro-server/src/server.rs
@@ -22,10 +22,10 @@ use base64::Engine as _;
use base64::engine::general_purpose::STANDARD as BASE64_STANDARD;
use bytes::Bytes;
pub use fabro_api::types::{
- AggregateBilling, AggregateBillingTotals, ApiQuestion, ApiQuestionOption, AppendEventResponse,
- ArtifactEntry, ArtifactListResponse, BillingByModel, BillingStageRef,
- CloseRunPullRequestResponse, CompletionContentPart, CompletionMessage, CompletionMessageRole,
- CompletionResponse, CompletionToolChoiceMode, CompletionUsage, CreateCompletionRequest,
+ AggregateBilling, AggregateBillingTotals, ApiQuestion, AppendEventResponse, ArtifactEntry,
+ ArtifactListResponse, BillingByModel, BillingStageRef, CloseRunPullRequestResponse,
+ CompletionContentPart, CompletionMessage, CompletionMessageRole, CompletionResponse,
+ CompletionToolChoiceMode, CompletionUsage, CreateCompletionRequest,
CreateRunPullRequestRequest, CreateSecretRequest, DeleteRunResponse, DeleteRunSandbox,
DeleteSecretRequest, DenyRunRequest, DiskUsageResponse, DiskUsageRunRow, DiskUsageSummaryRow,
ForkRequest, ForkResponse, LinkRunPullRequestRequest, MergeRunPullRequestRequest,
@@ -3242,14 +3242,7 @@ fn api_question_from_interview_record(question: &InterviewQuestionRecord) -> Api
text: question.text.clone(),
stage: question.stage.clone(),
question_type: question.question_type,
- options: question
- .options
- .iter()
- .map(|option| ApiQuestionOption {
- key: option.key.clone(),
- label: option.label.clone(),
- })
- .collect(),
+ options: question.options.clone(),
allow_freeform: question.allow_freeform,
timeout_seconds: question.timeout_seconds,
context_display: question.context_display.clone(),
diff --git a/lib/crates/fabro-server/src/server/tests.rs b/lib/crates/fabro-server/src/server/tests.rs
index 41349129e..2861a6bec 100644
--- a/lib/crates/fabro-server/src/server/tests.rs
+++ b/lib/crates/fabro-server/src/server/tests.rs
@@ -5787,8 +5787,10 @@ async fn submit_pending_interview_answer_rejects_invalid_answer_shape() {
stage: "gate".to_string(),
question_type: QuestionType::MultipleChoice,
options: vec![fabro_types::run_event::InterviewOption {
- key: "approve".to_string(),
- label: "Approve".to_string(),
+ key: "approve".to_string(),
+ label: "Approve".to_string(),
+ description: None,
+ preview: None,
}],
allow_freeform: false,
timeout_seconds: None,
@@ -5874,8 +5876,10 @@ fn answer_from_typed_selected_request_validates_and_attaches_option() {
stage: "gate".to_string(),
question_type: QuestionType::MultipleChoice,
options: vec![fabro_types::run_event::InterviewOption {
- key: "approve".to_string(),
- label: "Approve".to_string(),
+ key: "approve".to_string(),
+ label: "Approve".to_string(),
+ description: None,
+ preview: None,
}],
allow_freeform: false,
timeout_seconds: None,
@@ -5905,12 +5909,16 @@ fn answer_from_typed_multi_selected_request_validates_option_keys() {
question_type: QuestionType::MultiSelect,
options: vec![
fabro_types::run_event::InterviewOption {
- key: "approve".to_string(),
- label: "Approve".to_string(),
+ key: "approve".to_string(),
+ label: "Approve".to_string(),
+ description: None,
+ preview: None,
},
fabro_types::run_event::InterviewOption {
- key: "notify".to_string(),
- label: "Notify".to_string(),
+ key: "notify".to_string(),
+ label: "Notify".to_string(),
+ description: None,
+ preview: None,
},
],
allow_freeform: false,
diff --git a/lib/crates/fabro-slack/src/blocks.rs b/lib/crates/fabro-slack/src/blocks.rs
index bea8188c3..0062f2df0 100644
--- a/lib/crates/fabro-slack/src/blocks.rs
+++ b/lib/crates/fabro-slack/src/blocks.rs
@@ -153,6 +153,32 @@ fn lead_blocks(question: &Question, run_web_url: Option<&str>) -> Vec {
blocks
}
+fn option_descriptions_section(question: &Question) -> Option {
+ let rows = question
+ .options
+ .iter()
+ .filter_map(|option| {
+ let description = option.description.as_deref()?.trim();
+ if description.is_empty() {
+ return None;
+ }
+ Some(format!(
+ "• *{}* — {}",
+ escape_slack_controls(&option.label),
+ escape_slack_controls(description)
+ ))
+ })
+ .collect::>();
+ if rows.is_empty() {
+ return None;
+ }
+ Some(text_block(&truncate_to_limit(
+ &rows.join("\n"),
+ SLACK_SECTION_TEXT_LIMIT,
+ HEADER_TRUNCATION_SUFFIX,
+ )))
+}
+
pub fn answered_blocks(question_text: &str, answer_text: &str) -> Vec {
vec![text_block(&format!(
"~{}~\n*Answer:* {}",
@@ -168,6 +194,9 @@ pub fn question_to_blocks(
run_web_url: Option<&str>,
) -> Vec {
let mut blocks = lead_blocks(question, run_web_url);
+ if let Some(descriptions) = option_descriptions_section(question) {
+ blocks.push(descriptions);
+ }
match question.question_type {
QuestionType::YesNo | QuestionType::Confirmation => {
@@ -212,10 +241,26 @@ pub fn question_to_blocks(
.options
.iter()
.map(|opt| {
- json!({
+ let mut option = json!({
"text": { "type": "plain_text", "text": opt.label },
"value": opt.key
- })
+ });
+ if let Some(description) = opt
+ .description
+ .as_deref()
+ .map(str::trim)
+ .filter(|value| !value.is_empty())
+ {
+ option["description"] = json!({
+ "type": "plain_text",
+ "text": truncate_to_limit(
+ description,
+ 75,
+ HEADER_TRUNCATION_SUFFIX
+ )
+ });
+ }
+ option
})
.collect();
blocks.push(json!({
@@ -436,16 +481,22 @@ mod tests {
let mut q = Question::new("Pick a language:", QuestionType::MultipleChoice);
q.options = vec![
InterviewOption {
- key: "rs".to_string(),
- label: "Rust".to_string(),
+ key: "rs".to_string(),
+ label: "Rust".to_string(),
+ description: None,
+ preview: None,
},
InterviewOption {
- key: "ts".to_string(),
- label: "TypeScript".to_string(),
+ key: "ts".to_string(),
+ label: "TypeScript".to_string(),
+ description: None,
+ preview: None,
},
InterviewOption {
- key: "py".to_string(),
- label: "Python".to_string(),
+ key: "py".to_string(),
+ label: "Python".to_string(),
+ description: None,
+ preview: None,
},
];
let blocks = question_to_blocks("run-1", "q-3", &q, None);
@@ -709,12 +760,16 @@ mod tests {
let mut q = Question::new("Select features:", QuestionType::MultiSelect);
q.options = vec![
InterviewOption {
- key: "a".to_string(),
- label: "Auth".to_string(),
+ key: "a".to_string(),
+ label: "Auth".to_string(),
+ description: None,
+ preview: None,
},
InterviewOption {
- key: "b".to_string(),
- label: "Billing".to_string(),
+ key: "b".to_string(),
+ label: "Billing".to_string(),
+ description: None,
+ preview: None,
},
];
let blocks = question_to_blocks("run-1", "q-5", &q, None);
@@ -746,6 +801,24 @@ mod tests {
);
}
+ #[test]
+ fn option_descriptions_are_rendered_and_preview_is_not_special_cased() {
+ let mut q = Question::new("Pick one:", QuestionType::MultipleChoice);
+ q.options = vec![InterviewOption {
+ key: "ship".to_string(),
+ label: "Ship".to_string(),
+ description: Some("Deploy ".to_string()),
+ preview: Some("preview should not render".to_string()),
+ }];
+
+ let blocks_value: Value =
+ serde_json::to_value(question_to_blocks("run-1", "q-6", &q, None)).unwrap();
+ let text = blocks_value.to_string();
+
+ assert!(text.contains("Deploy <now>"));
+ assert!(!text.contains("preview should not render"));
+ }
+
#[test]
fn run_started_blocks_include_run_link_workflow_and_run_id_without_actions() {
let blocks = run_lifecycle_blocks(RunLifecycleKind::Started, &RunLifecycleBlocks {
diff --git a/lib/crates/fabro-store/src/run_state.rs b/lib/crates/fabro-store/src/run_state.rs
index 421885feb..f30e10598 100644
--- a/lib/crates/fabro-store/src/run_state.rs
+++ b/lib/crates/fabro-store/src/run_state.rs
@@ -2059,12 +2059,16 @@ mod tests {
question_type: "multiple_choice".to_string(),
options: vec![
InterviewOption {
- key: "approve".to_string(),
- label: "Approve".to_string(),
+ key: "approve".to_string(),
+ label: "Approve".to_string(),
+ description: Some("Ship it".to_string()),
+ preview: Some("deploy --prod".to_string()),
},
InterviewOption {
- key: "revise".to_string(),
- label: "Revise".to_string(),
+ key: "revise".to_string(),
+ label: "Revise".to_string(),
+ description: None,
+ preview: None,
},
],
allow_freeform: true,
@@ -2083,6 +2087,14 @@ mod tests {
assert_eq!(pending.question.stage, "gate");
assert_eq!(pending.question.question_type, QuestionType::MultipleChoice);
assert_eq!(pending.question.options.len(), 2);
+ assert_eq!(
+ pending.question.options[0].description.as_deref(),
+ Some("Ship it")
+ );
+ assert_eq!(
+ pending.question.options[0].preview.as_deref(),
+ Some("deploy --prod")
+ );
assert!(pending.question.allow_freeform);
assert_eq!(pending.question.timeout_seconds, Some(30.0));
assert_eq!(
diff --git a/lib/crates/fabro-types/src/run_event/misc.rs b/lib/crates/fabro-types/src/run_event/misc.rs
index b2280c243..3f88cdbd3 100644
--- a/lib/crates/fabro-types/src/run_event/misc.rs
+++ b/lib/crates/fabro-types/src/run_event/misc.rs
@@ -4,10 +4,14 @@ use serde_json::Value;
use super::ExecOutputTail;
use crate::{CommandTermination, PullRequestLink};
-#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
+#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, Default)]
pub struct InterviewOption {
- pub key: String,
- pub label: String,
+ pub key: String,
+ pub label: String,
+ #[serde(default, skip_serializing_if = "Option::is_none")]
+ pub description: Option,
+ #[serde(default, skip_serializing_if = "Option::is_none")]
+ pub preview: Option,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
diff --git a/lib/crates/fabro-workflow/src/handler/agent.rs b/lib/crates/fabro-workflow/src/handler/agent.rs
index 23c18477b..0788ee1e5 100644
--- a/lib/crates/fabro-workflow/src/handler/agent.rs
+++ b/lib/crates/fabro-workflow/src/handler/agent.rs
@@ -12,6 +12,7 @@ use super::{EngineServices, Handler, NodeTimeoutPolicy};
use crate::context::{Context, WorkflowContext, keys};
use crate::error::Error;
use crate::event::{Emitter, Event, StageScope};
+use crate::interview_runtime::WorkflowAgentQuestionRuntime;
use crate::outcome::{
BilledModelUsage, FailureCategory, FailureDetail, Outcome, OutcomeExt, StageOutcome,
};
@@ -28,14 +29,15 @@ pub enum CodergenResult {
}
pub struct CodergenRunRequest<'a> {
- pub node: &'a Node,
- pub prompt: &'a str,
- pub context: &'a Context,
- pub thread_id: Option<&'a str>,
- pub emitter: &'a Arc,
- pub sandbox: &'a Arc,
- pub tool_hooks: Option>,
- pub cancel_token: CancellationToken,
+ pub node: &'a Node,
+ pub prompt: &'a str,
+ pub context: &'a Context,
+ pub thread_id: Option<&'a str>,
+ pub emitter: &'a Arc,
+ pub sandbox: &'a Arc,
+ pub tool_hooks: Option>,
+ pub cancel_token: CancellationToken,
+ pub agent_tool_runtime: fabro_agent::AgentToolRuntime,
}
pub struct OneShotRequest<'a> {
@@ -311,6 +313,15 @@ impl Handler for AgentHandler {
StageModelUsage::MODE_AGENT,
self.backend.as_deref(),
)?;
+ let agent_tool_runtime = fabro_agent::AgentToolRuntime::with_question_runtime(Arc::new(
+ WorkflowAgentQuestionRuntime::new(
+ Arc::clone(&services.interviewer),
+ Arc::clone(&services.run.emitter),
+ stage_scope.clone(),
+ node.id.clone(),
+ Arc::clone(&services.run.interview_blocker),
+ ),
+ ));
// 3. Call LLM backend (agent loop)
let thread_id = context.thread_id();
@@ -341,6 +352,7 @@ impl Handler for AgentHandler {
sandbox: &services.run.sandbox,
tool_hooks,
cancel_token: services.run.cancel_token(),
+ agent_tool_runtime: agent_tool_runtime.clone(),
})
.await;
match result {
diff --git a/lib/crates/fabro-workflow/src/handler/fan_in.rs b/lib/crates/fabro-workflow/src/handler/fan_in.rs
index 690479516..8ba1e251f 100644
--- a/lib/crates/fabro-workflow/src/handler/fan_in.rs
+++ b/lib/crates/fabro-workflow/src/handler/fan_in.rs
@@ -272,6 +272,7 @@ async fn llm_evaluate(
sandbox,
tool_hooks: None,
cancel_token,
+ agent_tool_runtime: fabro_agent::AgentToolRuntime::default(),
})
.await
{
diff --git a/lib/crates/fabro-workflow/src/handler/human.rs b/lib/crates/fabro-workflow/src/handler/human.rs
index a46a6e525..77b83454a 100644
--- a/lib/crates/fabro-workflow/src/handler/human.rs
+++ b/lib/crates/fabro-workflow/src/handler/human.rs
@@ -1,13 +1,12 @@
use std::path::Path;
use std::str::FromStr;
use std::sync::Arc;
-use std::sync::atomic::{AtomicUsize, Ordering};
use std::time::Instant;
use async_trait::async_trait;
use fabro_graphviz::graph::{Graph, Node};
use fabro_interview::{Answer, AnswerValue, Interviewer, Question, ask_with_timeout};
-use fabro_types::{BlockedReason, InterviewOption, Principal, QuestionType, SystemActorKind};
+use fabro_types::{InterviewOption, Principal, QuestionType, SystemActorKind};
use ulid::Ulid;
use super::{EngineServices, Handler, NodeTimeoutPolicy};
@@ -112,8 +111,10 @@ fn build_human_gate_question(
question.options = choices
.iter()
.map(|choice| InterviewOption {
- key: choice.key.clone(),
- label: choice.label.clone(),
+ key: choice.key.clone(),
+ label: choice.label.clone(),
+ description: None,
+ preview: None,
})
.collect();
question.allow_freeform = freeform_target.is_some();
@@ -138,61 +139,10 @@ fn build_human_gate_question(
})
}
-/// Refcount of open interviews for this handler's run. Emits `run.blocked`
-/// exactly once when the count transitions 0→1, and `run.unblocked` exactly
-/// once when it transitions back to 0. Internal to `HumanHandler`; shared
-/// across concurrent `execute` calls fanned out by `ParallelHandler`.
-struct BlockedStateTracker {
- unresolved_interviews: AtomicUsize,
-}
-
-impl BlockedStateTracker {
- fn new() -> Self {
- Self {
- unresolved_interviews: AtomicUsize::new(0),
- }
- }
-
- fn interview_started(&self, emitter: &Emitter) {
- if self.unresolved_interviews.fetch_add(1, Ordering::AcqRel) == 0 {
- emitter.emit(&Event::RunBlocked {
- blocked_reason: BlockedReason::HumanInputRequired,
- });
- }
- }
-
- fn interview_resolved(&self, emitter: &Emitter) {
- // Guard against unmatched resolves (e.g., tests that over-resolve) so
- // the counter cannot underflow. `compare_exchange_weak` loops until we
- // either observe zero (and bail) or successfully decrement.
- let mut current = self.unresolved_interviews.load(Ordering::Acquire);
- loop {
- if current == 0 {
- return;
- }
- match self.unresolved_interviews.compare_exchange_weak(
- current,
- current - 1,
- Ordering::AcqRel,
- Ordering::Acquire,
- ) {
- Ok(_) => {
- if current == 1 {
- emitter.emit(&Event::RunUnblocked);
- }
- return;
- }
- Err(observed) => current = observed,
- }
- }
- }
-}
-
/// Blocks until a human selects an option derived from outgoing edges.
pub struct HumanHandler {
interviewer: Arc,
emitter: Option>,
- tracker: BlockedStateTracker,
}
impl HumanHandler {
@@ -200,7 +150,6 @@ impl HumanHandler {
Self {
interviewer,
emitter: None,
- tracker: BlockedStateTracker::new(),
}
}
@@ -295,8 +244,10 @@ impl Handler for HumanHandler {
.options
.iter()
.map(|option| InterviewOption {
- key: option.key.clone(),
- label: option.label.clone(),
+ key: option.key.clone(),
+ label: option.label.clone(),
+ description: option.description.clone(),
+ preview: option.preview.clone(),
})
.collect(),
allow_freeform: question.allow_freeform,
@@ -305,8 +256,10 @@ impl Handler for HumanHandler {
},
&stage_scope,
);
- self.tracker
- .interview_started(services.run.emitter.as_ref());
+ let interview_guard = services
+ .run
+ .interview_blocker
+ .block(Arc::clone(&services.run.emitter));
let interview_start = Instant::now();
let answer_submission = ask_with_timeout(self.interviewer.as_ref(), question).await;
let answer_actor = answer_submission.actor.clone();
@@ -327,8 +280,7 @@ impl Handler for HumanHandler {
},
&stage_scope,
);
- self.tracker
- .interview_resolved(services.run.emitter.as_ref());
+ interview_guard.resolve();
let default_choice = node
.attrs
.get("human.default_choice")
@@ -371,8 +323,7 @@ impl Handler for HumanHandler {
},
&stage_scope,
);
- self.tracker
- .interview_resolved(services.run.emitter.as_ref());
+ interview_guard.resolve();
return Ok(unanswered_human_gate(
"human interaction interrupted before an answer was provided",
));
@@ -389,8 +340,7 @@ impl Handler for HumanHandler {
},
&stage_scope,
);
- self.tracker
- .interview_resolved(services.run.emitter.as_ref());
+ interview_guard.resolve();
return Ok(unanswered_human_gate("human skipped interaction"));
}
@@ -406,8 +356,7 @@ impl Handler for HumanHandler {
},
&stage_scope,
);
- self.tracker
- .interview_resolved(services.run.emitter.as_ref());
+ interview_guard.resolve();
// Try fixed-choice match
if let Some(selected) = find_choice_match(&answer, &choices) {
@@ -894,8 +843,10 @@ mod tests {
async fn wait_human_emits_blocked_then_unblocked_around_interview() {
let interviewer = Arc::new(CallbackInterviewer::new(|_| {
Answer::selected("A", InterviewOption {
- key: "A".to_string(),
- label: "Approve".to_string(),
+ key: "A".to_string(),
+ label: "Approve".to_string(),
+ description: None,
+ preview: None,
})
}));
let handler = HumanHandler::new(interviewer);
@@ -1128,9 +1079,10 @@ mod tests {
#[test]
fn blocked_state_tracker_emits_once_across_parallel_interview_races() {
- let tracker = BlockedStateTracker::new();
+ let blocker = Arc::new(crate::interview_runtime::RunInterviewBlocker::new());
let emitter = Arc::new(Emitter::new(fabro_types::fixtures::RUN_1));
let event_names = Arc::new(Mutex::new(Vec::new()));
+ let guards = Arc::new(Mutex::new(Vec::new()));
emitter.on_event({
let event_names = Arc::clone(&event_names);
@@ -1148,22 +1100,25 @@ mod tests {
std::thread::scope(|scope| {
for _ in 0..8 {
- let tracker = &tracker;
+ let blocker = Arc::clone(&blocker);
let emitter = Arc::clone(&emitter);
- scope.spawn(move || tracker.interview_started(emitter.as_ref()));
+ let guards = Arc::clone(&guards);
+ scope.spawn(move || {
+ guards.lock().unwrap().push(blocker.block(emitter));
+ });
}
});
std::thread::scope(|scope| {
for _ in 0..8 {
- let tracker = &tracker;
- let emitter = Arc::clone(&emitter);
- scope.spawn(move || tracker.interview_resolved(emitter.as_ref()));
+ let guards = Arc::clone(&guards);
+ scope.spawn(move || {
+ let guard = guards.lock().unwrap().pop().unwrap();
+ guard.resolve();
+ });
}
});
- tracker.interview_resolved(emitter.as_ref());
-
assert_eq!(event_names.lock().unwrap().as_slice(), [
"run.blocked",
"run.unblocked"
diff --git a/lib/crates/fabro-workflow/src/handler/llm/acp.rs b/lib/crates/fabro-workflow/src/handler/llm/acp.rs
index 6dfbaaff4..827f09e8a 100644
--- a/lib/crates/fabro-workflow/src/handler/llm/acp.rs
+++ b/lib/crates/fabro-workflow/src/handler/llm/acp.rs
@@ -451,14 +451,15 @@ mod tests {
let context = Context::new();
let result = backend
.run(CodergenRunRequest {
- node: &node,
- prompt: "write hello",
- context: &context,
- thread_id: None,
- emitter: &emitter,
- sandbox: &sandbox,
- tool_hooks: None,
- cancel_token: CancellationToken::new(),
+ node: &node,
+ prompt: "write hello",
+ context: &context,
+ thread_id: None,
+ emitter: &emitter,
+ sandbox: &sandbox,
+ tool_hooks: None,
+ cancel_token: CancellationToken::new(),
+ agent_tool_runtime: fabro_agent::AgentToolRuntime::default(),
})
.await
.unwrap();
@@ -518,14 +519,15 @@ mod tests {
let context = Context::new();
let result = backend
.run(CodergenRunRequest {
- node: &node,
- prompt: "write hello",
- context: &context,
- thread_id: None,
- emitter: &emitter,
- sandbox: &sandbox,
- tool_hooks: None,
- cancel_token: CancellationToken::new(),
+ node: &node,
+ prompt: "write hello",
+ context: &context,
+ thread_id: None,
+ emitter: &emitter,
+ sandbox: &sandbox,
+ tool_hooks: None,
+ cancel_token: CancellationToken::new(),
+ agent_tool_runtime: fabro_agent::AgentToolRuntime::default(),
})
.await
.unwrap();
@@ -565,14 +567,15 @@ mod tests {
let context = Context::new();
let result = backend
.run(CodergenRunRequest {
- node: &node,
- prompt: "write hello",
- context: &context,
- thread_id: None,
- emitter: &emitter,
- sandbox: &sandbox,
- tool_hooks: None,
- cancel_token: CancellationToken::new(),
+ node: &node,
+ prompt: "write hello",
+ context: &context,
+ thread_id: None,
+ emitter: &emitter,
+ sandbox: &sandbox,
+ tool_hooks: None,
+ cancel_token: CancellationToken::new(),
+ agent_tool_runtime: fabro_agent::AgentToolRuntime::default(),
})
.await
.unwrap();
@@ -603,14 +606,15 @@ mod tests {
let context = Context::new();
let result = backend
.run(CodergenRunRequest {
- node: &node,
- prompt: "write hello",
- context: &context,
- thread_id: None,
- emitter: &emitter,
- sandbox: &sandbox_dyn,
- tool_hooks: None,
- cancel_token: CancellationToken::new(),
+ node: &node,
+ prompt: "write hello",
+ context: &context,
+ thread_id: None,
+ emitter: &emitter,
+ sandbox: &sandbox_dyn,
+ tool_hooks: None,
+ cancel_token: CancellationToken::new(),
+ agent_tool_runtime: fabro_agent::AgentToolRuntime::default(),
})
.await;
assert!(result.is_err());
@@ -652,14 +656,15 @@ mod tests {
let context = Context::new();
let result = backend
.run(CodergenRunRequest {
- node: &node,
- prompt: "cancel",
- context: &context,
- thread_id: None,
- emitter: &emitter,
- sandbox: &sandbox,
- tool_hooks: None,
- cancel_token: CancellationToken::new(),
+ node: &node,
+ prompt: "cancel",
+ context: &context,
+ thread_id: None,
+ emitter: &emitter,
+ sandbox: &sandbox,
+ tool_hooks: None,
+ cancel_token: CancellationToken::new(),
+ agent_tool_runtime: fabro_agent::AgentToolRuntime::default(),
})
.await;
let Err(err) = result else {
@@ -705,14 +710,15 @@ mod tests {
let context = Context::new();
backend
.run(CodergenRunRequest {
- node: &node,
- prompt: "write hello",
- context: &context,
- thread_id: None,
- emitter: &emitter,
- sandbox: &sandbox,
- tool_hooks: None,
- cancel_token: CancellationToken::new(),
+ node: &node,
+ prompt: "write hello",
+ context: &context,
+ thread_id: None,
+ emitter: &emitter,
+ sandbox: &sandbox,
+ tool_hooks: None,
+ cancel_token: CancellationToken::new(),
+ agent_tool_runtime: fabro_agent::AgentToolRuntime::default(),
})
.await
.unwrap();
@@ -746,14 +752,15 @@ mod tests {
let context = Context::new();
let result = backend
.run(CodergenRunRequest {
- node: &node,
- prompt: "write hello",
- context: &context,
- thread_id: None,
- emitter: &emitter,
- sandbox: &sandbox_dyn,
- tool_hooks: None,
- cancel_token: CancellationToken::new(),
+ node: &node,
+ prompt: "write hello",
+ context: &context,
+ thread_id: None,
+ emitter: &emitter,
+ sandbox: &sandbox_dyn,
+ tool_hooks: None,
+ cancel_token: CancellationToken::new(),
+ agent_tool_runtime: fabro_agent::AgentToolRuntime::default(),
})
.await;
let Err(err) = result else {
@@ -798,14 +805,15 @@ mod tests {
let context = Context::new();
let result = backend
.run(CodergenRunRequest {
- node: &node,
- prompt: "write hello",
- context: &context,
- thread_id: None,
- emitter: &emitter,
- sandbox: &sandbox_dyn,
- tool_hooks: None,
- cancel_token: CancellationToken::new(),
+ node: &node,
+ prompt: "write hello",
+ context: &context,
+ thread_id: None,
+ emitter: &emitter,
+ sandbox: &sandbox_dyn,
+ tool_hooks: None,
+ cancel_token: CancellationToken::new(),
+ agent_tool_runtime: fabro_agent::AgentToolRuntime::default(),
})
.await;
let Err(err) = result else {
diff --git a/lib/crates/fabro-workflow/src/handler/llm/api.rs b/lib/crates/fabro-workflow/src/handler/llm/api.rs
index 709bd278e..8da02640d 100644
--- a/lib/crates/fabro-workflow/src/handler/llm/api.rs
+++ b/lib/crates/fabro-workflow/src/handler/llm/api.rs
@@ -7,7 +7,7 @@ use fabro_agent::tool_registry::{RegisteredTool, ToolContext, ToolRegistry};
use fabro_agent::{
AgentEvent, AgentProfile, AnthropicProfile, CompletionCoordinator, GeminiProfile,
Message as AgentMessage, OpenAiProfile, Sandbox, Session, SessionOptions, StaticEnvProvider,
- ToolEnvProvider,
+ ToolEnvProvider, register_question_tools,
};
use fabro_auth::{CredentialSource, EnvCredentialSource};
use fabro_graphviz::graph::{AttrValue, Node};
@@ -736,6 +736,7 @@ impl AgentApiBackend {
});
profile.register_subagent_tools(manager, factory, 0);
+ register_question_tools(provider.profile_kind, profile.tool_registry_mut());
if let Some(services) = fabro_run_tools {
register_fabro_run_tools(profile.tool_registry_mut(), &services);
}
@@ -979,6 +980,7 @@ impl CodergenBackend for AgentApiBackend {
let sandbox = request.sandbox;
let tool_hooks = request.tool_hooks;
let cancel_token = request.cancel_token;
+ let agent_tool_runtime = request.agent_tool_runtime;
let fidelity = context.fidelity();
let reuse_key = if fidelity == Fidelity::Full {
@@ -1089,7 +1091,9 @@ impl CodergenBackend for AgentApiBackend {
return Err(err);
}
}
- session.process_input(prompt).await
+ session
+ .process_input_with_runtime(prompt, agent_tool_runtime.clone())
+ .await
}
Err(err) => Err(err),
};
@@ -1213,7 +1217,10 @@ impl CodergenBackend for AgentApiBackend {
return Err(err);
}
}
- match session.process_input(prompt).await {
+ match session
+ .process_input_with_runtime(prompt, agent_tool_runtime.clone())
+ .await
+ {
Ok(()) => {
succeeded = true;
break;
diff --git a/lib/crates/fabro-workflow/src/handler/manager_loop.rs b/lib/crates/fabro-workflow/src/handler/manager_loop.rs
index cf6fcaf8e..d7e1dff7b 100644
--- a/lib/crates/fabro-workflow/src/handler/manager_loop.rs
+++ b/lib/crates/fabro-workflow/src/handler/manager_loop.rs
@@ -231,6 +231,7 @@ impl Handler for SubWorkflowHandler {
let parent_run = Arc::clone(&services.run);
let registry = Arc::clone(&services.registry);
+ let interviewer = Arc::clone(&services.interviewer);
let base_env = services.base_env.clone();
let github_token = services.github_token.clone();
let inputs = services.inputs.clone();
@@ -269,6 +270,7 @@ impl Handler for SubWorkflowHandler {
engine: Arc::new(EngineServices {
run: child_run,
registry,
+ interviewer,
git_state: std::sync::RwLock::new(None),
base_env,
github_token,
diff --git a/lib/crates/fabro-workflow/src/handler/parallel.rs b/lib/crates/fabro-workflow/src/handler/parallel.rs
index eeee23acc..7c85042cc 100644
--- a/lib/crates/fabro-workflow/src/handler/parallel.rs
+++ b/lib/crates/fabro-workflow/src/handler/parallel.rs
@@ -304,6 +304,7 @@ impl Handler for ParallelHandler {
for setup in branch_setups {
let parent_run = Arc::clone(&services.run);
let registry = Arc::clone(&services.registry);
+ let interviewer = Arc::clone(&services.interviewer);
let base_env = services.base_env.clone();
let github_token = services.github_token.clone();
let inputs = services.inputs.clone();
@@ -375,6 +376,7 @@ impl Handler for ParallelHandler {
let branch_services = EngineServices {
run: parent_run.with_sandbox(Arc::clone(&setup.sandbox)),
registry: Arc::clone(®istry),
+ interviewer,
git_state: std::sync::RwLock::new(None),
base_env: base_env.clone(),
github_token: github_token.clone(),
diff --git a/lib/crates/fabro-workflow/src/interview_runtime.rs b/lib/crates/fabro-workflow/src/interview_runtime.rs
new file mode 100644
index 000000000..f0370ae0e
--- /dev/null
+++ b/lib/crates/fabro-workflow/src/interview_runtime.rs
@@ -0,0 +1,625 @@
+use std::sync::Arc;
+use std::sync::atomic::{AtomicUsize, Ordering};
+use std::time::Instant;
+
+use async_trait::async_trait;
+use fabro_agent::{
+ AgentQuestion, AgentQuestionAnswer, AgentQuestionAnswerStatus, AgentQuestionRuntime,
+};
+use fabro_interview::{Answer, AnswerSubmission, AnswerValue, Interviewer, Question};
+use fabro_types::{BlockedReason, InterviewOption, Principal, SystemActorKind};
+use futures::future;
+use tokio_util::sync::CancellationToken;
+use ulid::Ulid;
+
+use crate::event::{Emitter, Event, StageScope};
+use crate::millis_u64;
+
+/// Run-scoped refcount for unresolved human input. Emits `run.blocked` on the
+/// first unresolved human/agent interview and `run.unblocked` after the last
+/// one resolves.
+pub(crate) struct RunInterviewBlocker {
+ unresolved_interviews: AtomicUsize,
+}
+
+impl RunInterviewBlocker {
+ #[must_use]
+ pub(crate) fn new() -> Self {
+ Self {
+ unresolved_interviews: AtomicUsize::new(0),
+ }
+ }
+
+ pub(crate) fn block(self: &Arc, emitter: Arc) -> RunInterviewGuard {
+ if self.unresolved_interviews.fetch_add(1, Ordering::AcqRel) == 0 {
+ emitter.emit(&Event::RunBlocked {
+ blocked_reason: BlockedReason::HumanInputRequired,
+ });
+ }
+ RunInterviewGuard {
+ blocker: Arc::clone(self),
+ emitter,
+ resolved: false,
+ }
+ }
+
+ fn resolved(&self, emitter: &Emitter) {
+ let mut current = self.unresolved_interviews.load(Ordering::Acquire);
+ loop {
+ if current == 0 {
+ return;
+ }
+ match self.unresolved_interviews.compare_exchange_weak(
+ current,
+ current - 1,
+ Ordering::AcqRel,
+ Ordering::Acquire,
+ ) {
+ Ok(_) => {
+ if current == 1 {
+ emitter.emit(&Event::RunUnblocked);
+ }
+ return;
+ }
+ Err(observed) => current = observed,
+ }
+ }
+ }
+}
+
+pub(crate) struct RunInterviewGuard {
+ blocker: Arc,
+ emitter: Arc,
+ resolved: bool,
+}
+
+impl RunInterviewGuard {
+ pub(crate) fn resolve(mut self) {
+ self.resolve_in_place();
+ }
+
+ fn resolve_in_place(&mut self) {
+ if !self.resolved {
+ self.blocker.resolved(self.emitter.as_ref());
+ self.resolved = true;
+ }
+ }
+}
+
+impl Drop for RunInterviewGuard {
+ fn drop(&mut self) {
+ self.resolve_in_place();
+ }
+}
+
+pub(crate) struct WorkflowAgentQuestionRuntime {
+ interviewer: Arc,
+ emitter: Arc,
+ stage_scope: StageScope,
+ stage_id: String,
+ blocker: Arc,
+}
+
+impl WorkflowAgentQuestionRuntime {
+ #[must_use]
+ pub(crate) fn new(
+ interviewer: Arc,
+ emitter: Arc,
+ stage_scope: StageScope,
+ stage_id: impl Into,
+ blocker: Arc,
+ ) -> Self {
+ Self {
+ interviewer,
+ emitter,
+ stage_scope,
+ stage_id: stage_id.into(),
+ blocker,
+ }
+ }
+}
+
+struct PreparedQuestion {
+ agent_question: AgentQuestion,
+ question: Question,
+}
+
+struct PendingAgentQuestionBatch {
+ emitter: Arc,
+ stage_scope: StageScope,
+ stage_id: String,
+ questions: Vec<(String, String)>,
+ started_at: Instant,
+ guard: Option,
+}
+
+impl PendingAgentQuestionBatch {
+ fn new(
+ emitter: Arc,
+ stage_scope: StageScope,
+ stage_id: String,
+ prepared: &[PreparedQuestion],
+ guard: RunInterviewGuard,
+ started_at: Instant,
+ ) -> Self {
+ Self {
+ emitter,
+ stage_scope,
+ stage_id,
+ questions: prepared
+ .iter()
+ .map(|prepared_question| {
+ (
+ prepared_question.question.id.clone(),
+ prepared_question.question.text.clone(),
+ )
+ })
+ .collect(),
+ started_at,
+ guard: Some(guard),
+ }
+ }
+
+ fn resolve(mut self) {
+ if let Some(guard) = self.guard.take() {
+ guard.resolve();
+ }
+ }
+}
+
+impl Drop for PendingAgentQuestionBatch {
+ fn drop(&mut self) {
+ if self.guard.is_none() {
+ return;
+ }
+ let duration_ms = millis_u64(self.started_at.elapsed());
+ for (question_id, question) in &self.questions {
+ self.emitter.emit_scoped(
+ &Event::InterviewInterrupted {
+ actor: Some(Principal::System {
+ system_kind: SystemActorKind::Engine,
+ }),
+ question_id: question_id.clone(),
+ question: question.clone(),
+ stage: self.stage_id.clone(),
+ reason: "interrupted".to_string(),
+ duration_ms,
+ },
+ &self.stage_scope,
+ );
+ }
+ if let Some(guard) = self.guard.take() {
+ guard.resolve();
+ }
+ }
+}
+
+#[async_trait]
+impl AgentQuestionRuntime for WorkflowAgentQuestionRuntime {
+ async fn ask_questions(
+ &self,
+ tool_call_id: &str,
+ questions: Vec,
+ cancel_token: CancellationToken,
+ ) -> Result, String> {
+ if questions.is_empty() {
+ return Ok(Vec::new());
+ }
+
+ let prepared = questions
+ .into_iter()
+ .enumerate()
+ .map(|(index, question)| self.prepare_question(tool_call_id, index, question))
+ .collect::>();
+
+ for prepared_question in &prepared {
+ let question = &prepared_question.question;
+ self.emitter.emit_scoped(
+ &Event::InterviewStarted {
+ question_id: question.id.clone(),
+ question: question.text.clone(),
+ stage: self.stage_id.clone(),
+ question_type: question.question_type.to_string(),
+ options: question.options.clone(),
+ allow_freeform: question.allow_freeform,
+ timeout_seconds: None,
+ context_display: question.context_display.clone(),
+ },
+ &self.stage_scope,
+ );
+ }
+
+ let interview_start = Instant::now();
+ let cleanup = PendingAgentQuestionBatch::new(
+ Arc::clone(&self.emitter),
+ self.stage_scope.clone(),
+ self.stage_id.clone(),
+ &prepared,
+ self.blocker.block(Arc::clone(&self.emitter)),
+ interview_start,
+ );
+ let ask_all = future::join_all(
+ prepared
+ .iter()
+ .map(|prepared_question| self.interviewer.ask(prepared_question.question.clone())),
+ );
+ tokio::pin!(ask_all);
+
+ let answers = tokio::select! {
+ submissions = &mut ask_all => Some(submissions),
+ () = cancel_token.cancelled() => None,
+ };
+
+ let results = match answers {
+ Some(submissions) => prepared
+ .iter()
+ .zip(submissions)
+ .map(|(prepared_question, submission)| {
+ self.emit_submission_event(
+ prepared_question,
+ &submission,
+ millis_u64(interview_start.elapsed()),
+ );
+ answer_from_submission(&prepared_question.agent_question, &submission)
+ })
+ .collect::>(),
+ None => prepared
+ .iter()
+ .map(|prepared_question| {
+ self.emit_interrupted(
+ prepared_question,
+ Some(Principal::System {
+ system_kind: SystemActorKind::Engine,
+ }),
+ "interrupted",
+ millis_u64(interview_start.elapsed()),
+ );
+ AgentQuestionAnswer {
+ original_id: prepared_question.agent_question.original_id.clone(),
+ original_question: prepared_question
+ .agent_question
+ .original_question
+ .clone(),
+ answers: Vec::new(),
+ status: AgentQuestionAnswerStatus::Interrupted,
+ }
+ })
+ .collect::>(),
+ };
+
+ cleanup.resolve();
+ Ok(results)
+ }
+}
+
+impl WorkflowAgentQuestionRuntime {
+ fn prepare_question(
+ &self,
+ tool_call_id: &str,
+ index: usize,
+ agent_question: AgentQuestion,
+ ) -> PreparedQuestion {
+ let mut question = Question::new(agent_question.text.clone(), agent_question.question_type);
+ question.id = internal_question_id(&self.stage_scope, tool_call_id, index);
+ question.options.clone_from(&agent_question.options);
+ question.allow_freeform = agent_question.allow_freeform;
+ question.stage.clone_from(&self.stage_id);
+ question.metadata.insert(
+ "agent.tool_call_id".to_string(),
+ serde_json::json!(tool_call_id),
+ );
+ question.metadata.insert(
+ "agent.original_question".to_string(),
+ serde_json::json!(agent_question.original_question),
+ );
+ if let Some(original_id) = &agent_question.original_id {
+ question.metadata.insert(
+ "agent.original_id".to_string(),
+ serde_json::json!(original_id),
+ );
+ }
+ if let Some(header) = &agent_question.header {
+ question
+ .metadata
+ .insert("agent.header".to_string(), serde_json::json!(header));
+ }
+ PreparedQuestion {
+ agent_question,
+ question,
+ }
+ }
+
+ fn emit_submission_event(
+ &self,
+ prepared: &PreparedQuestion,
+ submission: &AnswerSubmission,
+ duration_ms: u64,
+ ) {
+ match submission.answer.value {
+ AnswerValue::Timeout => self.emitter.emit_scoped(
+ &Event::InterviewTimeout {
+ actor: Some(Principal::System {
+ system_kind: SystemActorKind::Timeout,
+ }),
+ question_id: prepared.question.id.clone(),
+ question: prepared.question.text.clone(),
+ stage: self.stage_id.clone(),
+ duration_ms,
+ },
+ &self.stage_scope,
+ ),
+ AnswerValue::Interrupted => self.emit_interrupted(
+ prepared,
+ Some(submission.actor.clone()),
+ "interrupted",
+ duration_ms,
+ ),
+ AnswerValue::Cancelled => self.emit_interrupted(
+ prepared,
+ Some(submission.actor.clone()),
+ "cancelled",
+ duration_ms,
+ ),
+ _ => self.emitter.emit_scoped(
+ &Event::InterviewCompleted {
+ actor: Some(submission.actor.clone()),
+ question_id: prepared.question.id.clone(),
+ question: prepared.question.text.clone(),
+ answer: answer_labels(&prepared.question.options, &submission.answer)
+ .join(", "),
+ duration_ms,
+ },
+ &self.stage_scope,
+ ),
+ }
+ }
+
+ fn emit_interrupted(
+ &self,
+ prepared: &PreparedQuestion,
+ actor: Option,
+ reason: &str,
+ duration_ms: u64,
+ ) {
+ self.emitter.emit_scoped(
+ &Event::InterviewInterrupted {
+ actor,
+ question_id: prepared.question.id.clone(),
+ question: prepared.question.text.clone(),
+ stage: self.stage_id.clone(),
+ reason: reason.to_string(),
+ duration_ms,
+ },
+ &self.stage_scope,
+ );
+ }
+}
+
+fn answer_from_submission(
+ agent_question: &AgentQuestion,
+ submission: &AnswerSubmission,
+) -> AgentQuestionAnswer {
+ let status = match &submission.answer.value {
+ AnswerValue::Cancelled => AgentQuestionAnswerStatus::Cancelled,
+ AnswerValue::Interrupted => AgentQuestionAnswerStatus::Interrupted,
+ AnswerValue::Skipped => AgentQuestionAnswerStatus::Skipped,
+ AnswerValue::Timeout => AgentQuestionAnswerStatus::Timeout,
+ _ => AgentQuestionAnswerStatus::Answered,
+ };
+ let answers = if status == AgentQuestionAnswerStatus::Answered {
+ answer_labels(&agent_question.options, &submission.answer)
+ } else {
+ Vec::new()
+ };
+ AgentQuestionAnswer {
+ original_id: agent_question.original_id.clone(),
+ original_question: agent_question.original_question.clone(),
+ answers,
+ status,
+ }
+}
+
+fn answer_labels(options: &[InterviewOption], answer: &Answer) -> Vec {
+ match &answer.value {
+ AnswerValue::Selected(key) => vec![label_for_key(options, key)],
+ AnswerValue::MultiSelected(keys) => {
+ keys.iter().map(|key| label_for_key(options, key)).collect()
+ }
+ AnswerValue::Text(text) => vec![text.clone()],
+ AnswerValue::Yes => vec!["yes".to_string()],
+ AnswerValue::No => vec!["no".to_string()],
+ AnswerValue::Cancelled => vec!["cancelled".to_string()],
+ AnswerValue::Interrupted => vec!["interrupted".to_string()],
+ AnswerValue::Skipped => vec!["skipped".to_string()],
+ AnswerValue::Timeout => vec!["timeout".to_string()],
+ }
+}
+
+fn label_for_key(options: &[InterviewOption], key: &str) -> String {
+ options
+ .iter()
+ .find(|option| option.key == key)
+ .map_or_else(|| key.to_string(), |option| option.label.clone())
+}
+
+fn internal_question_id(scope: &StageScope, tool_call_id: &str, index: usize) -> String {
+ format!(
+ "agentq-{}-v{}-{}-{}-{}",
+ slug(&scope.node_id),
+ scope.visit,
+ slug(tool_call_id),
+ index + 1,
+ Ulid::new(),
+ )
+}
+
+fn slug(value: &str) -> String {
+ let mut out = value
+ .chars()
+ .filter_map(|ch| {
+ if ch.is_ascii_alphanumeric() {
+ Some(ch.to_ascii_lowercase())
+ } else if matches!(ch, '-' | '_') {
+ Some(ch)
+ } else {
+ None
+ }
+ })
+ .take(48)
+ .collect::();
+ if out.is_empty() {
+ out.push('x');
+ }
+ out
+}
+
+#[cfg(test)]
+mod tests {
+ use fabro_interview::ControlInterviewer;
+ use fabro_types::{EventBody, RunId};
+
+ use super::*;
+
+ #[test]
+ fn answer_labels_return_user_facing_labels_in_submission_order() {
+ let options = vec![
+ InterviewOption {
+ key: "a".to_string(),
+ label: "Alpha".to_string(),
+ ..InterviewOption::default()
+ },
+ InterviewOption {
+ key: "b".to_string(),
+ label: "Beta".to_string(),
+ ..InterviewOption::default()
+ },
+ ];
+ let answer = Answer::multi_selected(vec!["b".to_string(), "a".to_string()]);
+
+ assert_eq!(answer_labels(&options, &answer), vec!["Beta", "Alpha"]);
+ }
+
+ #[test]
+ fn internal_question_id_includes_stage_visit_and_tool_call_context() {
+ let scope = StageScope {
+ node_id: "Review Changes".to_string(),
+ visit: 3,
+ parallel_group_id: None,
+ parallel_branch_id: None,
+ };
+
+ let id = internal_question_id(&scope, "call_123", 1);
+
+ assert!(id.starts_with("agentq-reviewchanges-v3-call_123-2-"));
+ let ulid = id
+ .rsplit('-')
+ .next()
+ .expect("question id should include a ULID suffix");
+ assert_eq!(ulid.len(), 26);
+ }
+
+ #[tokio::test]
+ async fn batch_questions_are_all_started_before_run_is_blocked_and_return_labels() {
+ let interviewer = Arc::new(ControlInterviewer::new());
+ let emitter = Arc::new(Emitter::new(RunId::new()));
+ let events = Arc::new(std::sync::Mutex::new(Vec::new()));
+ emitter.on_event({
+ let events = Arc::clone(&events);
+ move |event| events.lock().unwrap().push(event.clone())
+ });
+ let runtime = WorkflowAgentQuestionRuntime::new(
+ interviewer.clone(),
+ Arc::clone(&emitter),
+ StageScope {
+ node_id: "ask".to_string(),
+ visit: 1,
+ parallel_group_id: None,
+ parallel_branch_id: None,
+ },
+ "ask",
+ Arc::new(RunInterviewBlocker::new()),
+ );
+ let option = InterviewOption {
+ key: "ship".to_string(),
+ label: "Ship it".to_string(),
+ description: Some("Deploy".to_string()),
+ preview: Some("preview".to_string()),
+ };
+
+ let ask = tokio::spawn(async move {
+ runtime
+ .ask_questions(
+ "call_1",
+ vec![
+ AgentQuestion {
+ original_id: Some("q1".to_string()),
+ original_question: "First?".to_string(),
+ header: None,
+ text: "First?".to_string(),
+ question_type: fabro_types::QuestionType::MultipleChoice,
+ options: vec![option.clone()],
+ allow_freeform: true,
+ },
+ AgentQuestion {
+ original_id: Some("q2".to_string()),
+ original_question: "Second?".to_string(),
+ header: None,
+ text: "Second?".to_string(),
+ question_type: fabro_types::QuestionType::MultipleChoice,
+ options: vec![option.clone()],
+ allow_freeform: true,
+ },
+ ],
+ CancellationToken::new(),
+ )
+ .await
+ .unwrap()
+ });
+
+ tokio::task::yield_now().await;
+ let question_ids = {
+ let events = events.lock().unwrap();
+ assert!(matches!(events[0].body, EventBody::InterviewStarted(_)));
+ assert!(matches!(events[1].body, EventBody::InterviewStarted(_)));
+ assert!(matches!(events[2].body, EventBody::RunBlocked(_)));
+ events
+ .iter()
+ .filter_map(|event| match &event.body {
+ EventBody::InterviewStarted(props) => Some(props.question_id.clone()),
+ _ => None,
+ })
+ .collect::>()
+ };
+
+ for question_id in question_ids {
+ let option = InterviewOption {
+ key: "ship".to_string(),
+ label: "Ship it".to_string(),
+ ..InterviewOption::default()
+ };
+ interviewer
+ .submit(
+ &question_id,
+ AnswerSubmission::system(
+ Answer::selected("ship", option),
+ SystemActorKind::Engine,
+ ),
+ )
+ .await
+ .unwrap();
+ }
+
+ let answers = ask.await.unwrap();
+
+ assert_eq!(answers.len(), 2);
+ assert_eq!(answers[0].answers, vec!["Ship it"]);
+ assert_eq!(answers[1].answers, vec!["Ship it"]);
+ assert!(
+ events
+ .lock()
+ .unwrap()
+ .iter()
+ .any(|event| matches!(event.body, EventBody::RunUnblocked(_)))
+ );
+ }
+}
diff --git a/lib/crates/fabro-workflow/src/lib.rs b/lib/crates/fabro-workflow/src/lib.rs
index 194a0ca2c..42f5c7f43 100644
--- a/lib/crates/fabro-workflow/src/lib.rs
+++ b/lib/crates/fabro-workflow/src/lib.rs
@@ -298,6 +298,7 @@ pub mod github_token_source;
pub(crate) mod graph;
pub mod handler;
mod hook_context;
+mod interview_runtime;
#[allow(
dead_code,
reason = "The lifecycle module remains crate-visible for tests and pending integrations."
diff --git a/lib/crates/fabro-workflow/src/pipeline/initialize.rs b/lib/crates/fabro-workflow/src/pipeline/initialize.rs
index 953f39d2d..3165f7855 100644
--- a/lib/crates/fabro-workflow/src/pipeline/initialize.rs
+++ b/lib/crates/fabro-workflow/src/pipeline/initialize.rs
@@ -675,6 +675,7 @@ pub async fn initialize(
let engine = Arc::new(EngineServices {
run: Arc::clone(&run_services),
registry,
+ interviewer: Arc::clone(&options.interviewer),
git_state: std::sync::RwLock::new(None),
base_env,
github_token,
diff --git a/lib/crates/fabro-workflow/src/services.rs b/lib/crates/fabro-workflow/src/services.rs
index 72c5499f9..9210998d3 100644
--- a/lib/crates/fabro-workflow/src/services.rs
+++ b/lib/crates/fabro-workflow/src/services.rs
@@ -9,6 +9,7 @@ use fabro_auth::CredentialSource;
#[cfg(test)]
use fabro_auth::ResolvedCredentials;
use fabro_hooks::{HookContext, HookDecision, HookExecutionContext, HookRunner};
+use fabro_interview::Interviewer;
use fabro_model::{Catalog, ProviderId};
use fabro_types::{ManifestPath, RunId};
use tokio_util::sync::CancellationToken;
@@ -16,6 +17,7 @@ use tokio_util::sync::CancellationToken;
use crate::event::Emitter;
use crate::github_token_source::GitHubTokenSource;
use crate::handler::HandlerRegistry;
+use crate::interview_runtime::RunInterviewBlocker;
use crate::run_metadata::{RunMetadataRuntime, RunMetadataWriterHandle};
use crate::runtime_store::RunStoreHandle;
use crate::sandbox_git::GitState;
@@ -91,19 +93,20 @@ pub struct FabroRunToolServices {
/// does NOT count as cancellation.
#[derive(Clone)]
pub struct RunServices {
- pub run_store: RunStoreHandle,
- pub emitter: Arc,
- pub sandbox: Arc,
- pub hook_runner: Option>,
- pub locations: RunLocations,
- pub(crate) cancel_token: CancellationToken,
- pub provider_id: ProviderId,
- pub model: String,
- pub llm_source: Arc,
- pub catalog: Arc,
- pub(crate) sandbox_git: Arc,
- pub(crate) metadata_runtime: Arc,
- pub(crate) metadata_writer: Option,
+ pub run_store: RunStoreHandle,
+ pub emitter: Arc,
+ pub sandbox: Arc,
+ pub hook_runner: Option>,
+ pub locations: RunLocations,
+ pub(crate) cancel_token: CancellationToken,
+ pub provider_id: ProviderId,
+ pub model: String,
+ pub llm_source: Arc,
+ pub catalog: Arc,
+ pub(crate) sandbox_git: Arc,
+ pub(crate) metadata_runtime: Arc,
+ pub(crate) metadata_writer: Option,
+ pub(crate) interview_blocker: Arc,
}
impl RunServices {
@@ -137,6 +140,7 @@ impl RunServices {
sandbox_git,
metadata_runtime,
metadata_writer,
+ interview_blocker: Arc::new(RunInterviewBlocker::new()),
})
}
@@ -224,6 +228,7 @@ impl RunServices {
pub struct EngineServices {
pub run: Arc,
pub registry: Arc,
+ pub interviewer: Arc,
/// Git state for the current run. Set via `set_git_state` at the start of
/// `execute` and read by parallel/fan-in handlers.
pub(crate) git_state: std::sync::RwLock