Migrate sandbox config to named environments; add InterviewOption metad… (#372)

## Summary

Two related changes land together: the sandbox configuration surface is
replaced with a named-environment model, and `InterviewOption` gains
`description` and `preview` fields needed for the mid-stage agent
interview tools described in the plan.

## What changed

### Named environments (was `[run.sandbox]`)

`[run.sandbox]` and its provider-specific sub-tables
(`[run.sandbox.daytona]`, `[run.sandbox.docker]`) are replaced by a
two-level model:

- **`[environments.<slug>]`** — reusable catalog entries with a unified
shape: `provider`, `image`, `resources`, `network`, `lifecycle`,
`labels`, `volumes`, `env`.
- **`[run.environment] id = "<slug>"`** — selects which environment a
run uses.
- **`[run.environment.<field>]`** — sparse run-level overrides applied
on top of the selected environment.

The OpenAPI schema drops `RunSandboxSettings`, `DaytonaSettings`,
`DaytonaSnapshotSettings`, `DaytonaNetworkLayer`, and `DockerSettings`
in favour of `EnvironmentSettings`, `RunEnvironmentSettings`, and the
new sub-schemas (`EnvironmentImageSettings`,
`EnvironmentResourcesSettings`, `EnvironmentNetworkSettings`,
`EnvironmentLifecycleSettings`, `EnvironmentVolumeSettings`). The
`--sandbox` CLI flag becomes `--environment`.

All docs, example configs, `.fabro/project.toml`, and the
automation-detail / run-settings UI panels are updated to the new shape.
The run-settings page renames "Sandbox" → "Environment" and reads from
the new field paths.

### `InterviewOption` metadata fields

`description` and `preview` are added to the canonical `InterviewOption`
type (OpenAPI, helpers.ts, interview-dock, human-qa renderer). Both are
treated as untrusted model-authored text — stored and displayed as plain
strings, never rendered as HTML. The `interview-dock` test asserts that
raw HTML in `preview` is not rendered. Option `description` is shown as
secondary text under the label in choice and multi-select buttons.

### `StageModelUsage` projection

`provider_used` on `RunStageInfo` and stage projections is promoted from
a freeform object to a typed `StageModelUsage` schema (with `mode`,
`provider`, `model`, `reasoning_effort`, `speed`). The
`extractStageModel` event-scraping helper is replaced by
`formatStageModelUsageLabel` and `stageModelUsageTitle`, which work
directly from the projection field. The `Stage` interface gains
`providerUsed` and the `EventsToolbar` consumes it.

### Other schema additions

`ReasoningEffort` enum, `small_default` on model info,
`SubAgentProjection`/`SkillsProjection`/`McpServerProjection` inline in
stage projections, and `TodoListProjection` moved from the run-state
top-level `todos_by_list` map into per-stage `todos`.

### Plan summary

- Replace `[run.sandbox]` config with `[environments.<slug>]` +
`[run.environment]` selection across config, OpenAPI, UI, and docs.
- Extend `InterviewOption` with `description` and `preview`; render
`description` in choice/multi-select buttons.
- Promote `provider_used` to a typed `StageModelUsage` schema; drop
event-scraping in favour of the projection field.
- Add `ReasoningEffort`, `small_default`, subagent/skills/MCP
stage-projection schemas to OpenAPI.


### Fabro Details

<details>
<summary>Ran 9 stages in 93m 2s for $48.56</summary>

| Stage | Duration | Cost | Retries |
|---|---|---|---|
| start | 0s | – | 0 |
| toolchain | 1s | – | 0 |
| preflight_compile | 2m 4s | – | 0 |
| preflight_lint | 2m 16s | – | 0 |
| implement | 39m 28s | $35.95 | 0 |
| simplify_opus | 22m 6s | $8.66 | 0 |
| simplify_gpt | 7m 18s | $1.66 | 0 |
| verify | 6m 33s | – | 0 |
| fixup | 12m 34s | $2.29 | 0 |
| **Total** | **93m 2s** | **$48.56** | **0** |

</details>

<details>
<summary>Ran <code>ImplementPlan.fabro</code> (11 nodes and 14
edges)</summary>

```dot
digraph ImplementPlan {
    graph [
        goal="Implement and simplify",
        model_stylesheet="
            * { model: claude-opus-4-7; }
        "
    ]
    rankdir=LR

    start [shape=Mdiamond, label="Start"]
    exit  [shape=Msquare, label="Exit"]

    toolchain         [label="Toolchain", shape=parallelogram, script="command -v cargo >/dev/null || { curl --proto '=https' --tlsv1.2 -sSf https://sh.rustup.rs | sh -s -- -y && sudo ln -sf $HOME/.cargo/bin/* /usr/local/bin/; }; cargo --version 2>&1", max_retries=0]
    preflight_compile [label="Preflight Compile", shape=parallelogram, script="cargo check -q --workspace 2>&1", max_retries=0]
    preflight_lint    [label="Preflight Lint", shape=parallelogram, script="cargo +nightly-2026-04-14 clippy -q --workspace --all-targets -- -D warnings 2>&1", max_retries=0]
    fix_lints         [label="Fix Lints", prompt="The preflight lint step failed. Read the build output from context and fix all clippy lint warnings.", max_visits=3]
    implement         [label="Implement", prompt="Read the plan file referenced in the goal and implement every step. Make all the code changes described in the plan. Use red/green TDD.", model="gpt-55", reasoning_effort="xhigh"]
    simplify_opus     [label="Simplify (Opus)", prompt="@prompts/simplify.md"]
    simplify_gpt      [label="Simplify (GPT-55)", prompt="@prompts/simplify.md", model="gpt-55"]
    verify            [label="Verify", shape=parallelogram, script="git fetch origin main 2>&1 && git merge --no-edit --no-stat origin/main 2>&1 && cargo +nightly-2026-04-14 fmt --all 2>&1 && cargo dev docs refresh 2>&1 && cargo +nightly-2026-04-14 fmt --check --all 2>&1 && ! rg -n 'AuthMode::Disabled|RunAuthMethod|RunSubjectProvenance|\bActorRef\b|\bActorKind\b|AuthenticatedSubject|AuthenticatedService|AuthorizeRunScoped|AuthorizeRunBlob|AuthorizeStageArtifact|AuthorizeCommandLog|auth_method\s*==\s*\"disabled\"' lib/crates apps lib/packages docs/public/api-reference/fabro-api.yaml 2>&1 && cargo +nightly-2026-04-14 clippy --workspace --all-targets -- -D warnings 2>&1 && cargo nextest run --workspace --status-level slow --profile ci 2>&1 && cargo dev docs check 2>&1 && bun install --frozen-lockfile 2>&1 && (cd apps/fabro-web && bun run typecheck) 2>&1 && (cd apps/fabro-web && bun run test) 2>&1 && (cd lib/packages/fabro-api-client && bun run typecheck) 2>&1 && cargo dev build -- -p fabro-cli --release 2>&1", goal_gate=true, retry_target="fixup"]
    fixup             [label="Fixup", prompt="The verify step failed. Read the build output from context and fix all format, clippy, Rust test, docs, TypeScript typecheck/test, and build failures.", max_visits=3]

    start -> toolchain
    toolchain -> preflight_compile [condition="outcome=succeeded"]
    toolchain -> exit
    preflight_compile -> preflight_lint [condition="outcome=succeeded"]
    preflight_compile -> exit
    preflight_lint -> implement [condition="outcome=succeeded"]
    preflight_lint -> fix_lints
    fix_lints -> preflight_lint
    implement -> simplify_opus -> simplify_gpt -> verify
    verify -> exit  [condition="outcome=succeeded"]
    verify -> fixup
    fixup -> verify
}

```

</details>

⚒️ Generated with [Fabro](https://fabro.sh)

---------

Co-authored-by: Fabro <noreply@fabro.sh>
Co-authored-by: Fabro <fabro@fabro.sh>
Co-authored-by: Bryan Helmkamp <bryan@brynary.com>
This commit is contained in:
fabro-sh-0530[bot] 2026-05-23 15:47:33 -04:00 • committed by GitHub
parent f73f2a53f3
commit dbe3e3966d
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
41 changed files with 2058 additions and 359 deletions

View file

@ -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: "<b>not rendered specially</b>",
},
],
});
const tree = render(
<InterviewDock runId="run-1" questions={[question]} />,
);
const text = textContent(tree.toJSON());
expect(text).toContain("Approve");
expect(text).toContain("Deploy the current patch");
expect(text).not.toContain("<b>not rendered specially</b>");
});
test("freeform question renders a textarea and disables send when empty", () => {
const question = makeQuestion({
question_type: QuestionType.FREEFORM,

View file

@ -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<void>;
@ -287,7 +287,7 @@ function ChoiceBody({
onClick={() => void onSubmit({ kind: "selected", option_key: option.key })}
className={CHOICE_BUTTON}
>
{displayLabel(option.label)}
<OptionLabel option={option} />
</button>
))}
</div>
@ -314,7 +314,7 @@ function MultiSelectBody({
submitting,
onSubmit,
}: {
options: ApiQuestionOption[];
options: InterviewOption[];
submitting: boolean;
onSubmit: (answer: SubmitInterviewAnswer) => Promise<void>;
}) {
@ -346,7 +346,7 @@ function MultiSelectBody({
className={isSelected ? CHOICE_BUTTON_SELECTED : CHOICE_BUTTON}
>
{isSelected && <CheckIcon className="size-3.5" aria-hidden="true" />}
{displayLabel(option.label)}
<OptionLabel option={option} />
</button>
);
})}
@ -455,6 +455,19 @@ function FreeformBody({
);
}
function OptionLabel({ option }: { option: InterviewOption }) {
return (
<span className="text-left">
<span className="block">{displayLabel(option.label)}</span>
{option.description && (
<span className="mt-0.5 block text-xs/5 font-normal text-fg-muted">
{option.description}
</span>
)}
</span>
);
}
function Spinner() {
return <ArrowPathIcon className="size-4 animate-spin" aria-hidden="true" />;
}

View file

@ -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, {

View file

@ -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;
}

View file

@ -184,7 +184,14 @@ function QuestionBlock({
<span className="inline-flex size-5 shrink-0 items-center justify-center rounded bg-overlay-strong font-mono text-[11px] text-fg-2">
{option.key}
</span>
<span className="text-sm text-fg-3">{option.label}</span>
<span className="min-w-0 text-sm text-fg-3">
<span>{option.label}</span>
{option.description && (
<span className="mt-0.5 block text-xs/5 text-fg-muted">
{option.description}
</span>
)}
</span>
</li>
))}
{question.allowFreeform && (

View file

@ -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.

View file

@ -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,

View file

@ -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<Arc<dyn AgentQuestionRuntime>>,
}
impl AgentToolRuntime {
#[must_use]
pub fn new() -> Self {
Self::default()
}
#[must_use]
pub fn with_question_runtime(runtime: Arc<dyn AgentQuestionRuntime>) -> Self {
Self {
question_runtime: Some(runtime),
}
}
#[must_use]
pub fn question_runtime(&self) -> Option<Arc<dyn AgentQuestionRuntime>> {
self.question_runtime.clone()
}
}
pub async fn scope_agent_tool_runtime<F>(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<String>,
pub original_question: String,
pub header: Option<String>,
pub text: String,
pub question_type: QuestionType,
pub options: Vec<InterviewOption>,
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<String>,
pub original_question: String,
pub answers: Vec<String>,
pub status: AgentQuestionAnswerStatus,
}
#[async_trait]
pub trait AgentQuestionRuntime: Send + Sync {
async fn ask_questions(
&self,
tool_call_id: &str,
questions: Vec<AgentQuestion>,
cancel_token: CancellationToken,
) -> Result<Vec<AgentQuestionAnswer>, String>;
}
#[derive(Debug, Deserialize)]
struct OpenAiQuestionToolArgs {
questions: Vec<OpenAiQuestion>,
}
#[derive(Debug, Deserialize)]
struct OpenAiQuestion {
id: String,
header: String,
question: String,
#[serde(default)]
options: Vec<OpenAiOption>,
}
#[derive(Debug, Deserialize)]
struct OpenAiOption {
label: String,
#[serde(default)]
description: Option<String>,
}
#[derive(Debug, Deserialize)]
struct AnthropicQuestionToolArgs {
questions: Vec<AnthropicQuestion>,
}
#[derive(Debug, Deserialize)]
#[serde(rename_all = "camelCase")]
struct AnthropicQuestion {
question: String,
#[serde(default)]
header: Option<String>,
#[serde(default)]
options: Vec<AnthropicOption>,
#[serde(default)]
multi_select: bool,
}
#[derive(Debug, Deserialize)]
struct AnthropicOption {
label: String,
#[serde(default)]
description: Option<String>,
#[serde(default)]
preview: Option<String>,
}
#[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<T: for<'de> Deserialize<'de>>(args: serde_json::Value) -> Result<T, String> {
serde_json::from_value(args).map_err(|err| format!("invalid question tool arguments: {err}"))
}
async fn execute_question_tool(
ctx: ToolContext,
questions: Vec<AgentQuestion>,
) -> Result<Vec<AgentQuestionAnswer>, 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<Vec<AgentQuestion>, 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<Vec<AgentQuestion>, 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<OpenAiOption>) -> Vec<InterviewOption> {
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<AnthropicOption>) -> Vec<InterviewOption> {
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<String, String> {
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<String, String> {
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<String, String> {
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::<Vec<_>>()
.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::<serde_json::Value>(&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());
}
}

View file

@ -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();

View file

@ -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<dyn ToolEnvProvider>>,
agent_tool_runtime: &AgentToolRuntime,
) -> Vec<ToolResult> {
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<dyn ToolEnvProvider>>,
agent_tool_runtime: &AgentToolRuntime,
) -> Vec<ToolResult> {
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<dyn ToolEnvProvider>>,
agent_tool_runtime: &AgentToolRuntime,
) -> Vec<ToolResult> {
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<dyn Sandbox>,
tool_hooks: Option<&Arc<dyn ToolHookCallback>>,
cancel_token: &CancellationToken,
config: &SessionOptions,
emitter: &Emitter,
session_id: &str,
root_session_id: &str,
tool_env_provider: Option<&Arc<dyn ToolEnvProvider>>,
agent_tool_runtime: &AgentToolRuntime,
) -> Vec<ToolResult> {
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<dyn ToolEnvProvider>>,
) -> 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<dyn Sandbox>,
tool_hooks: Option<&Arc<dyn ToolHookCallback>>,
cancel_token: CancellationToken,
config: &SessionOptions,
emitter: &Emitter,
session_id: &str,
root_session_id: &str,
tool_env_provider: Option<&Arc<dyn ToolEnvProvider>>,
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<dyn ToolEnvProvider>>,
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<dyn ToolEnvProvider>>,
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<AgentQuestion>,
_cancel_token: CancellationToken,
) -> Result<Vec<AgentQuestionAnswer>, 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,
&registry,
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,
&registry,
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<Mutex<Vec<(String, String, String)>>>,

View file

@ -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();

View file

@ -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,

View file

@ -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,

View file

@ -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,

View file

@ -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,
},
];

View file

@ -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,
})
);
}

View file

@ -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);

View file

@ -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);
}

View file

@ -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,

View file

@ -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(),

View file

@ -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,

View file

@ -153,6 +153,32 @@ fn lead_blocks(question: &Question, run_web_url: Option<&str>) -> Vec<Value> {
blocks
}
fn option_descriptions_section(question: &Question) -> Option<Value> {
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::<Vec<_>>();
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<Value> {
vec![text_block(&format!(
"~{}~\n*Answer:* {}",
@ -168,6 +194,9 @@ pub fn question_to_blocks(
run_web_url: Option<&str>,
) -> Vec<Value> {
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 <now>".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 &lt;now&gt;"));
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 {

View file

@ -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!(

View file

@ -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<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub preview: Option<String>,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]

View file

@ -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<Emitter>,
pub sandbox: &'a Arc<dyn Sandbox>,
pub tool_hooks: Option<Arc<dyn fabro_agent::ToolHookCallback>>,
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<Emitter>,
pub sandbox: &'a Arc<dyn Sandbox>,
pub tool_hooks: Option<Arc<dyn fabro_agent::ToolHookCallback>>,
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 {

View file

@ -272,6 +272,7 @@ async fn llm_evaluate(
sandbox,
tool_hooks: None,
cancel_token,
agent_tool_runtime: fabro_agent::AgentToolRuntime::default(),
})
.await
{

View file

@ -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<dyn Interviewer>,
emitter: Option<Arc<Emitter>>,
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"

View file

@ -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 {

View file

@ -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;

View file

@ -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,

View file

@ -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(&registry),
interviewer,
git_state: std::sync::RwLock::new(None),
base_env: base_env.clone(),
github_token: github_token.clone(),

View file

@ -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<Self>, emitter: Arc<Emitter>) -> 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<RunInterviewBlocker>,
emitter: Arc<Emitter>,
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<dyn Interviewer>,
emitter: Arc<Emitter>,
stage_scope: StageScope,
stage_id: String,
blocker: Arc<RunInterviewBlocker>,
}
impl WorkflowAgentQuestionRuntime {
#[must_use]
pub(crate) fn new(
interviewer: Arc<dyn Interviewer>,
emitter: Arc<Emitter>,
stage_scope: StageScope,
stage_id: impl Into<String>,
blocker: Arc<RunInterviewBlocker>,
) -> Self {
Self {
interviewer,
emitter,
stage_scope,
stage_id: stage_id.into(),
blocker,
}
}
}
struct PreparedQuestion {
agent_question: AgentQuestion,
question: Question,
}
struct PendingAgentQuestionBatch {
emitter: Arc<Emitter>,
stage_scope: StageScope,
stage_id: String,
questions: Vec<(String, String)>,
started_at: Instant,
guard: Option<RunInterviewGuard>,
}
impl PendingAgentQuestionBatch {
fn new(
emitter: Arc<Emitter>,
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<AgentQuestion>,
cancel_token: CancellationToken,
) -> Result<Vec<AgentQuestionAnswer>, 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::<Vec<_>>();
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::<Vec<_>>(),
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::<Vec<_>>(),
};
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<Principal>,
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<String> {
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::<String>();
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::<Vec<_>>()
};
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(_)))
);
}
}

View file

@ -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."

View file

@ -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,

View file

@ -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<Emitter>,
pub sandbox: Arc<dyn Sandbox>,
pub hook_runner: Option<Arc<HookRunner>>,
pub locations: RunLocations,
pub(crate) cancel_token: CancellationToken,
pub provider_id: ProviderId,
pub model: String,
pub llm_source: Arc<dyn CredentialSource>,
pub catalog: Arc<Catalog>,
pub(crate) sandbox_git: Arc<SandboxGitRuntime>,
pub(crate) metadata_runtime: Arc<RunMetadataRuntime>,
pub(crate) metadata_writer: Option<RunMetadataWriterHandle>,
pub run_store: RunStoreHandle,
pub emitter: Arc<Emitter>,
pub sandbox: Arc<dyn Sandbox>,
pub hook_runner: Option<Arc<HookRunner>>,
pub locations: RunLocations,
pub(crate) cancel_token: CancellationToken,
pub provider_id: ProviderId,
pub model: String,
pub llm_source: Arc<dyn CredentialSource>,
pub catalog: Arc<Catalog>,
pub(crate) sandbox_git: Arc<SandboxGitRuntime>,
pub(crate) metadata_runtime: Arc<RunMetadataRuntime>,
pub(crate) metadata_writer: Option<RunMetadataWriterHandle>,
pub(crate) interview_blocker: Arc<RunInterviewBlocker>,
}
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<RunServices>,
pub registry: Arc<HandlerRegistry>,
pub interviewer: Arc<dyn Interviewer>,
/// 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<Option<Arc<GitState>>>,
@ -330,6 +335,7 @@ impl EngineServices {
None,
),
registry: Arc::new(HandlerRegistry::new(Box::new(start::StartHandler))),
interviewer: Arc::new(fabro_interview::AutoApproveInterviewer::engine()),
git_state: std::sync::RwLock::new(None),
base_env: HashMap::new(),
github_token: None,

View file

@ -7,6 +7,7 @@ use std::time::Duration;
use fabro_agent::Sandbox;
use fabro_auth::{CredentialSource, EnvCredentialSource};
use fabro_graphviz::graph::Graph as GvGraph;
use fabro_interview::AutoApproveInterviewer;
use fabro_model::Catalog;
use fabro_store::{ArtifactStore, Database, RunProjection};
use object_store::local::LocalFileSystem;
@ -182,6 +183,7 @@ async fn initialized(
None,
),
registry: Arc::new(registry),
interviewer: Arc::new(AutoApproveInterviewer::engine()),
git_state: std::sync::RwLock::new(None),
base_env: options.env,
github_token: None,

View file

@ -28,7 +28,6 @@ models/agent-skill-activation-source.ts
models/agent-skill-summary.ts
models/aggregate-billing-totals.ts
models/aggregate-billing.ts
models/api-question-option.ts
models/api-question.ts
models/append-event-response.ts
models/approval-mode.ts

View file

@ -1,29 +0,0 @@
/* tslint:disable */
/* eslint-disable */
/**
* Fabro Run API
* HTTP API for managing Fabro workflow run executions.
*
* The version of the OpenAPI document: 0.1.0
*
*
* NOTE: This class is auto generated by OpenAPI Generator (https://openapi-generator.tech).
* https://openapi-generator.tech
* Do not edit the class manually.
*/
/**
* A selectable option for a multiple-choice or multi-select question.
*/
export interface ApiQuestionOption {
/**
* Machine-readable option key used when submitting an answer.
*/
'key': string;
/**
* Human-readable label displayed to the user.
*/
'label': string;
}

View file

@ -15,7 +15,7 @@
// May contain unused imports in some cases
// @ts-ignore
import type { ApiQuestionOption } from './api-question-option';
import type { InterviewOption } from './interview-option';
// May contain unused imports in some cases
// @ts-ignore
import type { QuestionType } from './question-type';
@ -40,7 +40,7 @@ export interface ApiQuestion {
/**
* Available options for selection-based questions. Empty for freeform questions.
*/
'options': Array<ApiQuestionOption>;
'options': Array<InterviewOption>;
/**
* Whether the user may provide freeform text in addition to selecting options.
*/

View file

@ -6,7 +6,6 @@ export * from './agent-skill-summary';
export * from './aggregate-billing';
export * from './aggregate-billing-totals';
export * from './api-question';
export * from './api-question-option';
export * from './append-event-response';
export * from './approval-mode';
export * from './artifact-batch-upload-entry';

View file

@ -18,6 +18,20 @@
* Option stored with an interview question in the event log.
*/
export interface InterviewOption {
/**
* Machine-readable option key used when submitting an answer.
*/
'key': string;
/**
* Human-readable label displayed to the user.
*/
'label': string;
/**
* Optional untrusted model-authored option description for display.
*/
'description'?: string | null;
/**
* Optional untrusted model-authored option preview captured for clients.
*/
'preview'?: string | null;
}