checkpoint

⚒️ Generated with [Fabro](https://fabro.sh)
This commit is contained in:
Fabro 2026-05-25 18:52:12 -04:00
parent ab3607cb29
commit 0f0df5ea17
6 changed files with 1696 additions and 41 deletions

554
run.json

File diff suppressed because one or more lines are too long

View file

@ -0,0 +1,964 @@
diff --git a/lib/crates/fabro-agent/src/session.rs b/lib/crates/fabro-agent/src/session.rs
index 24df5149d..0cfe25d98 100644
--- a/lib/crates/fabro-agent/src/session.rs
+++ b/lib/crates/fabro-agent/src/session.rs
@@ -76,6 +76,15 @@ pub struct SessionInputTiming {
pub tool: Duration,
}
+/// Take the value out of `start`, add its elapsed time to `total`. Used by
+/// `run_single_input` to accumulate inference and tool spans at well-defined
+/// boundaries (stream open, retry, error, cancel, end-of-loop).
+fn record_elapsed(start: &mut Option<Instant>, total: &mut Duration) {
+ if let Some(s) = start.take() {
+ *total = total.saturating_add(s.elapsed());
+ }
+}
+
impl SteeringItem {
#[must_use]
pub fn actor(&self) -> Option<&Principal> {
@@ -343,8 +352,6 @@ pub struct Session {
tool_env_provider: Option<Arc<dyn ToolEnvProvider>>,
subagent_manager: Option<Arc<AsyncMutex<SubAgentManager>>>,
completion_coordinator: Option<Arc<dyn CompletionCoordinator>>,
- last_input_inference_duration: Duration,
- last_input_tool_duration: Duration,
}
impl Session {
@@ -382,8 +389,6 @@ impl Session {
tool_env_provider: None,
subagent_manager,
completion_coordinator: None,
- last_input_inference_duration: Duration::ZERO,
- last_input_tool_duration: Duration::ZERO,
}
}
@@ -1198,28 +1203,22 @@ impl Session {
&self.file_tracker
}
- #[must_use]
- pub const fn last_input_timing(&self) -> SessionInputTiming {
- SessionInputTiming {
- inference: self.last_input_inference_duration,
- tool: self.last_input_tool_duration,
- }
- }
-
pub async fn process_input(&mut self, input: &str) -> Result<(), Error> {
self.process_input_with_runtime(input, AgentToolRuntime::default())
.await
+ .1
}
+ /// Process an input. Returns the inference/tool timing accumulated during
+ /// the call alongside the call result; timing is observed even on error.
pub async fn process_input_with_runtime(
&mut self,
input: &str,
agent_tool_runtime: AgentToolRuntime,
- ) -> Result<(), Error> {
- self.last_input_inference_duration = Duration::ZERO;
- self.last_input_tool_duration = Duration::ZERO;
+ ) -> (SessionInputTiming, Result<(), Error>) {
+ let mut timing = SessionInputTiming::default();
if self.state == SessionState::Closed {
- return Err(Error::SessionClosed);
+ return (timing, Err(Error::SessionClosed));
}
// Spawn wall-clock timeout task if configured
@@ -1241,7 +1240,9 @@ impl Session {
});
// Process the initial input, then drain any followups
- let mut result = self.run_single_input(input, &agent_tool_runtime).await;
+ let mut result = self
+ .run_single_input(input, &agent_tool_runtime, &mut timing)
+ .await;
if result.is_ok() {
loop {
@@ -1251,7 +1252,9 @@ impl Session {
.expect("followup queue lock poisoned")
.pop_front();
let Some(followup) = followup else { break };
- result = self.run_single_input(&followup, &agent_tool_runtime).await;
+ result = self
+ .run_single_input(&followup, &agent_tool_runtime, &mut timing)
+ .await;
if result.is_err() {
break;
}
@@ -1268,13 +1271,14 @@ impl Session {
self.transition(SessionState::Idle);
}
- result
+ (timing, result)
}
async fn run_single_input(
&mut self,
input: &str,
agent_tool_runtime: &AgentToolRuntime,
+ timing: &mut SessionInputTiming,
) -> Result<(), Error> {
const STREAM_CONSUME_RETRIES: usize = 3;
@@ -1409,15 +1413,6 @@ impl Session {
let client = self.llm_client.clone();
let cancel_token_for_select = self.cancel_token.clone();
let mut inference_start = Some(Instant::now());
- macro_rules! record_inference_duration {
- () => {
- if let Some(start) = inference_start.take() {
- self.last_input_inference_duration = self
- .last_input_inference_duration
- .saturating_add(start.elapsed());
- }
- };
- }
let stream_outcome: Option<Result<StreamEventStream, Error>> = tokio::select! {
biased;
() = round_token.cancelled() => None,
@@ -1428,12 +1423,12 @@ impl Session {
match stream {
Ok(stream) => stream,
Err(err) => {
- record_inference_duration!();
+ record_elapsed(&mut inference_start, &mut timing.inference);
return Err(err);
}
}
} else {
- record_inference_duration!();
+ record_elapsed(&mut inference_start, &mut timing.inference);
if self.cancel_token.is_cancelled() {
self.close();
return Err(self.interrupted_error());
@@ -1509,7 +1504,7 @@ impl Session {
// If terminal cancel fired, drop the stream and bail out.
if self.cancel_token.is_cancelled() {
drop(event_stream);
- record_inference_duration!();
+ record_elapsed(&mut inference_start, &mut timing.inference);
self.close();
return Err(self.interrupted_error());
}
@@ -1579,7 +1574,7 @@ impl Session {
match stream {
Ok(stream) => stream,
Err(err) => {
- record_inference_duration!();
+ record_elapsed(&mut inference_start, &mut timing.inference);
return Err(err);
}
}
@@ -1600,7 +1595,7 @@ impl Session {
},
);
}
- record_inference_duration!();
+ record_elapsed(&mut inference_start, &mut timing.inference);
return Err(self.emit_llm_error(err));
}
@@ -1632,7 +1627,7 @@ impl Session {
match stream {
Ok(stream) => stream,
Err(err) => {
- record_inference_duration!();
+ record_elapsed(&mut inference_start, &mut timing.inference);
return Err(err);
}
}
@@ -1643,7 +1638,7 @@ impl Session {
};
}
}
- record_inference_duration!();
+ record_elapsed(&mut inference_start, &mut timing.inference);
// Mid-LLM steer interrupt: drop the unrecorded turn, clear any
// partial visible output, and re-iterate. The next turn's
@@ -1770,9 +1765,7 @@ impl Session {
agent_tool_runtime,
)
.await;
- self.last_input_tool_duration = self
- .last_input_tool_duration
- .saturating_add(tool_start.elapsed());
+ timing.tool = timing.tool.saturating_add(tool_start.elapsed());
composite_watcher.abort();
if tool_calls
.iter()
@@ -2292,8 +2285,10 @@ mod tests {
let env = Arc::new(MockSandbox::default());
let mut session = Session::new(client, profile, env, SessionOptions::default(), None);
- session.process_input("use the slow tool").await.unwrap();
- let first = session.last_input_timing();
+ let (first, result) = session
+ .process_input_with_runtime("use the slow tool", AgentToolRuntime::default())
+ .await;
+ result.unwrap();
assert!(
first.inference >= Duration::from_millis(35),
"expected non-zero inference timing for first input, got {first:?}"
@@ -2303,8 +2298,10 @@ mod tests {
"expected non-zero tool timing for first input, got {first:?}"
);
- session.process_input("no tools this time").await.unwrap();
- let second = session.last_input_timing();
+ let (second, result) = session
+ .process_input_with_runtime("no tools this time", AgentToolRuntime::default())
+ .await;
+ result.unwrap();
assert!(
second.inference >= Duration::from_millis(15),
"expected per-input inference timing for second input, got {second:?}"
diff --git a/lib/crates/fabro-types/src/timing.rs b/lib/crates/fabro-types/src/timing.rs
index 6ca201143..49fa57b06 100644
--- a/lib/crates/fabro-types/src/timing.rs
+++ b/lib/crates/fabro-types/src/timing.rs
@@ -53,6 +53,14 @@ impl StageTiming {
Self::new(wall_time_ms, 0, 0)
}
+ /// Active-only timing for stages whose wall time will be supplied
+ /// separately by the executor's own stopwatch (current shape of the
+ /// handler → executor hop).
+ #[must_use]
+ pub fn active_only(inference_time_ms: u64, tool_time_ms: u64) -> Self {
+ Self::new(0, inference_time_ms, tool_time_ms)
+ }
+
/// Sum two timings field-by-field. Used to aggregate visits of one node
/// and to accumulate run-level rollups.
#[must_use]
diff --git a/lib/crates/fabro-workflow/src/billing_rollup.rs b/lib/crates/fabro-workflow/src/billing_rollup.rs
index 0e909ce08..feceb0e12 100644
--- a/lib/crates/fabro-workflow/src/billing_rollup.rs
+++ b/lib/crates/fabro-workflow/src/billing_rollup.rs
@@ -164,32 +164,12 @@ mod tests {
use fabro_model::{Catalog, ModelRef, ProviderId};
use fabro_types::{
- AttrValue, BilledModelUsage, BilledTokenCounts, Graph, Node, RunProjection, RunSpec,
- StageCompletion, StageOutcome, WorkflowSettings, first_event_seq, fixtures,
+ AttrValue, BilledTokenCounts, Graph, Node, RunProjection, RunSpec, StageCompletion,
+ StageOutcome, WorkflowSettings, first_event_seq, fixtures,
};
- use serde_json::json;
use super::billing_rollup_from_projection;
-
- fn test_usage(model_id: &str, input_tokens: i64, output_tokens: i64) -> BilledModelUsage {
- serde_json::from_value(json!({
- "input": {
- "usage": {
- "model": {
- "provider": "openai",
- "model_id": model_id
- },
- "tokens": {
- "input_tokens": input_tokens,
- "output_tokens": output_tokens
- }
- },
- "facts": { "algorithm": "openai" }
- },
- "total_usd_micros": input_tokens + output_tokens
- }))
- .unwrap()
- }
+ use crate::test_support::test_usage;
fn test_projection() -> RunProjection {
RunProjection::new(
diff --git a/lib/crates/fabro-workflow/src/event/convert.rs b/lib/crates/fabro-workflow/src/event/convert.rs
index 2e2014dc1..86a0fef41 100644
--- a/lib/crates/fabro-workflow/src/event/convert.rs
+++ b/lib/crates/fabro-workflow/src/event/convert.rs
@@ -1424,7 +1424,7 @@ mod tests {
use crate::error::Error;
use crate::event::test_support::user_principal;
use crate::event::{Event, StageScope};
- use crate::outcome::{BilledModelUsage, FailureDetail};
+ use crate::outcome::FailureDetail;
#[derive(Debug)]
struct EventTestCause;
@@ -1446,25 +1446,7 @@ mod tests {
}
}
- fn test_usage(model_id: &str, input_tokens: i64, output_tokens: i64) -> BilledModelUsage {
- serde_json::from_value(serde_json::json!({
- "input": {
- "usage": {
- "model": {
- "provider": "openai",
- "model_id": model_id
- },
- "tokens": {
- "input_tokens": input_tokens,
- "output_tokens": output_tokens
- }
- },
- "facts": { "algorithm": "openai" }
- },
- "total_usd_micros": input_tokens + output_tokens
- }))
- .unwrap()
- }
+ use crate::test_support::test_usage;
#[test]
fn run_event_stage_completed_places_node_fields_in_header() {
diff --git a/lib/crates/fabro-workflow/src/handler/agent.rs b/lib/crates/fabro-workflow/src/handler/agent.rs
index f7a9048a5..648f00c4b 100644
--- a/lib/crates/fabro-workflow/src/handler/agent.rs
+++ b/lib/crates/fabro-workflow/src/handler/agent.rs
@@ -20,10 +20,14 @@ use crate::interview_runtime::WorkflowAgentQuestionRuntime;
use crate::outcome::{BilledModelUsage, Outcome, OutcomeExt};
/// Result from a `CodergenBackend` invocation.
+#[allow(
+ clippy::large_enum_variant,
+ reason = "Text payload is the common case; Full(Box<Outcome>) is the rare alternative."
+)]
pub enum CodergenResult {
Text {
text: String,
- usage: Option<Box<BilledModelUsage>>,
+ usage: Option<BilledModelUsage>,
files_touched: Vec<String>,
last_file_touched: Option<String>,
/// Active timing observed by the backend. The wall field is ignored by
@@ -302,13 +306,7 @@ impl Handler for AgentHandler {
files_touched,
last_file_touched,
timing,
- }) => (
- text,
- usage.map(|usage| *usage),
- files_touched,
- last_file_touched,
- timing,
- ),
+ }) => (text, usage, files_touched, last_file_touched, timing),
Err(Error::Cancelled) => return Err(Error::Cancelled),
Err(e) if e.is_retryable() => {
return Err(e);
diff --git a/lib/crates/fabro-workflow/src/handler/command.rs b/lib/crates/fabro-workflow/src/handler/command.rs
index 736731a26..dde68dd5d 100644
--- a/lib/crates/fabro-workflow/src/handler/command.rs
+++ b/lib/crates/fabro-workflow/src/handler/command.rs
@@ -178,7 +178,7 @@ impl Handler for CommandHandler {
serde_json::json!(finalized.output_ref),
);
outcome.notes = Some(format!("Script completed: {script}"));
- outcome.timing = Some(StageTiming::new(0, 0, result.duration_ms));
+ outcome.timing = Some(StageTiming::active_only(0, result.duration_ms));
Ok(outcome)
} else {
let mut reason = format!(
@@ -191,7 +191,7 @@ impl Handler for CommandHandler {
keys::COMMAND_OUTPUT.to_string(),
serde_json::json!(finalized.output_ref),
);
- outcome.timing = Some(StageTiming::new(0, 0, result.duration_ms));
+ outcome.timing = Some(StageTiming::active_only(0, result.duration_ms));
Ok(outcome)
}
}
diff --git a/lib/crates/fabro-workflow/src/handler/llm/acp.rs b/lib/crates/fabro-workflow/src/handler/llm/acp.rs
index 4b867b7b8..e83d50cf9 100644
--- a/lib/crates/fabro-workflow/src/handler/llm/acp.rs
+++ b/lib/crates/fabro-workflow/src/handler/llm/acp.rs
@@ -233,7 +233,7 @@ impl AgentAcpBackend {
usage: None,
files_touched,
last_file_touched,
- timing: StageTiming::new(0, result.duration_ms, 0),
+ timing: StageTiming::active_only(result.duration_ms, 0),
})
}
diff --git a/lib/crates/fabro-workflow/src/handler/llm/api.rs b/lib/crates/fabro-workflow/src/handler/llm/api.rs
index df5af98c2..46e2e5d90 100644
--- a/lib/crates/fabro-workflow/src/handler/llm/api.rs
+++ b/lib/crates/fabro-workflow/src/handler/llm/api.rs
@@ -120,12 +120,6 @@ pub struct EffectiveRequestControls {
pub(crate) speed: Option<Speed>,
}
-fn active_stage_timing(inference: Duration, tool: Duration) -> StageTiming {
- // The executor ignores this wall field and supplies its own stopwatch-based
- // value when converting Outcome.timing into the emitted stage timing.
- StageTiming::new(0, crate::millis_u64(inference), crate::millis_u64(tool))
-}
-
fn classify_agent_error(err: fabro_agent::Error, allow_failover: bool) -> AgentApiErrorDisposition {
match err {
fabro_agent::Error::Interrupted(fabro_agent::InterruptReason::Cancelled) => {
@@ -1119,10 +1113,13 @@ impl CodergenBackend for AgentApiBackend {
return Ok(CodergenResult::Text {
text: response_text,
- usage: Some(Box::new(stage_usage)),
+ usage: Some(stage_usage),
files_touched: Vec::new(),
last_file_touched: None,
- timing: active_stage_timing(inference_duration, Duration::ZERO),
+ timing: StageTiming::active_only(
+ crate::millis_u64(inference_duration),
+ 0,
+ ),
});
}
}
@@ -1257,10 +1254,9 @@ impl CodergenBackend for AgentApiBackend {
if !is_reused {
emit_agent_tools_available(&session, &node.id, &stage_id, emitter);
}
- let process_result = session
+ let (timing, process_result) = session
.process_input_with_runtime(prompt, agent_tool_runtime.clone())
.await;
- let timing = session.last_input_timing();
inference_duration = inference_duration.saturating_add(timing.inference);
tool_duration = tool_duration.saturating_add(timing.tool);
process_result
@@ -1389,10 +1385,9 @@ impl CodergenBackend for AgentApiBackend {
}
}
emit_agent_tools_available(&session, &node.id, &stage_id, emitter);
- let process_result = session
+ let (timing, process_result) = session
.process_input_with_runtime(prompt, agent_tool_runtime.clone())
.await;
- let timing = session.last_input_timing();
inference_duration = inference_duration.saturating_add(timing.inference);
tool_duration = tool_duration.saturating_add(timing.tool);
match process_result {
@@ -1456,8 +1451,12 @@ impl CodergenBackend for AgentApiBackend {
));
}
let repair_message = error.repair_message(schema);
- let repair_result = session.process_input(&repair_message).await;
- let timing = session.last_input_timing();
+ let (timing, repair_result) = session
+ .process_input_with_runtime(
+ &repair_message,
+ fabro_agent::AgentToolRuntime::default(),
+ )
+ .await;
inference_duration = inference_duration.saturating_add(timing.inference);
tool_duration = tool_duration.saturating_add(timing.tool);
match repair_result {
@@ -1532,10 +1531,13 @@ impl CodergenBackend for AgentApiBackend {
Ok(CodergenResult::Text {
text: response,
- usage: Some(Box::new(stage_usage)),
+ usage: Some(stage_usage),
files_touched,
last_file_touched,
- timing: active_stage_timing(inference_duration, tool_duration),
+ timing: StageTiming::active_only(
+ crate::millis_u64(inference_duration),
+ crate::millis_u64(tool_duration),
+ ),
})
}
}
diff --git a/lib/crates/fabro-workflow/src/handler/prompt.rs b/lib/crates/fabro-workflow/src/handler/prompt.rs
index 94674421f..8d5b2fdf6 100644
--- a/lib/crates/fabro-workflow/src/handler/prompt.rs
+++ b/lib/crates/fabro-workflow/src/handler/prompt.rs
@@ -138,7 +138,7 @@ impl Handler for PromptHandler {
files_touched,
timing,
..
- }) => (text, usage.map(|usage| *usage), files_touched, timing),
+ }) => (text, usage, files_touched, timing),
Err(Error::Cancelled) => return Err(Error::Cancelled),
Err(e) if e.is_retryable() => {
return Err(e);
diff --git a/lib/crates/fabro-workflow/src/operations/start.rs b/lib/crates/fabro-workflow/src/operations/start.rs
index 0adc853ac..ec300d581 100644
--- a/lib/crates/fabro-workflow/src/operations/start.rs
+++ b/lib/crates/fabro-workflow/src/operations/start.rs
@@ -209,9 +209,12 @@ pub(super) async fn execute_persisted_run(
return Err(error);
}
- let mut bootstrap_guard =
- DetachedRunBootstrapGuard::arm(run_id, run_dir, event_sink.clone(), cancel_token.clone());
- bootstrap_guard.run_store = Some(run_store.clone());
+ let mut bootstrap_guard = DetachedRunBootstrapGuard::arm(
+ run_id,
+ run_store.clone(),
+ event_sink.clone(),
+ cancel_token.clone(),
+ );
let persisted = match Persisted::load_from_store(&services.run_store, run_dir).await {
Ok(persisted) => persisted,
@@ -280,28 +283,28 @@ pub(super) async fn execute_persisted_run(
}
}
-async fn persist_terminal_engine_failure(
+/// Build a conclusion from the store and emit `run.failed` carrying the
+/// rolled-up timing and billing. Shared by the engine-failure terminal path,
+/// the bootstrap/completion drop guards, and `persist_detached_failure`.
+async fn emit_workflow_run_failed(
run_id: RunId,
run_store: &RunStoreHandle,
event_sink: &RunEventSink,
- _run_dir: &Path,
error: &Error,
- duration: Duration,
+ reason: FailureReason,
+ wall_duration_ms: u64,
) {
- let engine_result: Result<Outcome, Error> = Err(error.clone());
- let (final_status, failure_reason, run_status) = classify_engine_result(&engine_result);
+ let failure = Some(error::run_failure_from_error(error, reason));
let conclusion = build_conclusion_from_store(
run_store,
- final_status,
- failure_reason,
- crate::millis_u64(duration),
+ StageOutcome::Failed {
+ retry_requested: false,
+ },
+ failure,
+ wall_duration_ms,
None,
)
.await;
- let reason = match run_status {
- RunStatus::Failed { reason } => reason,
- _ => FailureReason::WorkflowError,
- };
let failure_event = Event::workflow_run_failed_from_error(
error,
conclusion.timing,
@@ -309,13 +312,38 @@ async fn persist_terminal_engine_failure(
None,
None,
None,
- conclusion.billing.clone(),
+ conclusion.billing,
);
if let Err(err) = append_event_to_sink(event_sink, &run_id, &failure_event).await {
- tracing::warn!(error = %err, "Failed to append terminal engine failure event");
+ tracing::warn!(error = %err, "Failed to append run.failed event");
}
}
+async fn persist_terminal_engine_failure(
+ run_id: RunId,
+ run_store: &RunStoreHandle,
+ event_sink: &RunEventSink,
+ _run_dir: &Path,
+ error: &Error,
+ duration: Duration,
+) {
+ let engine_result: Result<Outcome, Error> = Err(error.clone());
+ let (_, _, run_status) = classify_engine_result(&engine_result);
+ let reason = match run_status {
+ RunStatus::Failed { reason } => reason,
+ _ => FailureReason::WorkflowError,
+ };
+ emit_workflow_run_failed(
+ run_id,
+ run_store,
+ event_sink,
+ error,
+ reason,
+ crate::millis_u64(duration),
+ )
+ .await;
+}
+
impl RunSession {
async fn new(persisted: &Persisted, services: StartServices) -> Result<Self, Error> {
let record = persisted.run_spec();
@@ -902,7 +930,7 @@ impl RunSession {
struct DetachedRunBootstrapGuard {
run_id: RunId,
- run_store: Option<RunStoreHandle>,
+ run_store: RunStoreHandle,
event_sink: RunEventSink,
cancel_token: CancellationToken,
active: bool,
@@ -911,13 +939,13 @@ struct DetachedRunBootstrapGuard {
impl DetachedRunBootstrapGuard {
fn arm(
run_id: RunId,
- _run_dir: &Path,
+ run_store: RunStoreHandle,
event_sink: RunEventSink,
cancel_token: CancellationToken,
) -> Self {
Self {
run_id,
- run_store: None,
+ run_store,
event_sink,
cancel_token,
active: true,
@@ -931,45 +959,29 @@ impl DetachedRunBootstrapGuard {
impl Drop for DetachedRunBootstrapGuard {
fn drop(&mut self) {
- if self.active {
- let cancelled = self.cancel_token.is_cancelled();
- let reason = if cancelled {
- FailureReason::Cancelled
- } else {
- FailureReason::SandboxInitFailed
- };
- let run_id = self.run_id;
- let run_store = self.run_store.clone();
- let event_sink = self.event_sink.clone();
- if let Ok(handle) = Handle::try_current() {
- handle.spawn(async move {
- let (timing, billing) = if let Some(run_store) = run_store {
- let final_status = StageOutcome::Failed {
- retry_requested: false,
- };
- let failure = Some(error::run_failure_from_error(
- &Error::engine(reason.to_string()),
- reason,
- ));
- let conclusion =
- build_conclusion_from_store(&run_store, final_status, failure, 0, None)
- .await;
- (conclusion.timing, conclusion.billing)
- } else {
- (fabro_types::RunTiming::default(), None)
- };
- let failure_event = Event::workflow_run_failed_from_error(
- &Error::engine(reason.to_string()),
- timing,
- reason,
- None,
- None,
- None,
- billing,
- );
- let _ = append_event_to_sink(&event_sink, &run_id, &failure_event).await;
- });
- }
+ if !self.active {
+ return;
+ }
+ let reason = if self.cancel_token.is_cancelled() {
+ FailureReason::Cancelled
+ } else {
+ FailureReason::SandboxInitFailed
+ };
+ let run_id = self.run_id;
+ let run_store = self.run_store.clone();
+ let event_sink = self.event_sink.clone();
+ if let Ok(handle) = Handle::try_current() {
+ handle.spawn(async move {
+ emit_workflow_run_failed(
+ run_id,
+ &run_store,
+ &event_sink,
+ &Error::engine(reason.to_string()),
+ reason,
+ 0,
+ )
+ .await;
+ });
}
}
}
@@ -1033,25 +1045,15 @@ impl Drop for DetachedRunCompletionGuard {
let run_store = self.run_store.clone();
if let Ok(handle) = Handle::try_current() {
handle.spawn(async move {
- let final_status = StageOutcome::Failed {
- retry_requested: false,
- };
- let failure = Some(error::run_failure_from_error(
- &Error::engine(message.to_string()),
- reason,
- ));
- let conclusion =
- build_conclusion_from_store(&run_store, final_status, failure, 0, None).await;
- let failure_event = Event::workflow_run_failed_from_error(
+ emit_workflow_run_failed(
+ run_id,
+ &run_store,
+ &event_sink,
&Error::engine(message.to_string()),
- conclusion.timing,
reason,
- None,
- None,
- None,
- conclusion.billing,
- );
- let _ = append_event_to_sink(&event_sink, &run_id, &failure_event).await;
+ 0,
+ )
+ .await;
let _ = append_event_to_sink(&event_sink, &run_id, &Event::RunNotice {
level: RunNoticeLevel::Error,
code: code.to_string(),
@@ -1073,30 +1075,12 @@ async fn persist_detached_failure(
reason: FailureReason,
error: &Error,
) -> Result<(), Error> {
- let message = error.to_string();
- let final_status = StageOutcome::Failed {
- retry_requested: false,
- };
- let failure = Some(error::run_failure_from_error(error, reason));
- let conclusion = build_conclusion_from_store(run_store, final_status, failure, 0, None).await;
-
- let failure_event = Event::workflow_run_failed_from_error(
- error,
- conclusion.timing,
- reason,
- None,
- None,
- None,
- conclusion.billing,
- );
- if let Err(err) = append_event_to_sink(event_sink, &run_id, &failure_event).await {
- tracing::warn!(error = %err, "Failed to append detached failure event");
- }
+ emit_workflow_run_failed(run_id, run_store, event_sink, error, reason, 0).await;
let event = Event::RunNotice {
level: RunNoticeLevel::Error,
code: format!("{phase}_failed"),
- message: message.clone(),
+ message: error.to_string(),
exec_output_tail: None,
};
if let Err(err) = append_event_to_sink(event_sink, &run_id, &event).await {
@@ -1473,25 +1457,7 @@ reasoning = false
}
}
- fn test_usage(model_id: &str, input_tokens: i64, output_tokens: i64) -> BilledModelUsage {
- serde_json::from_value(serde_json::json!({
- "input": {
- "usage": {
- "model": {
- "provider": "openai",
- "model_id": model_id
- },
- "tokens": {
- "input_tokens": input_tokens,
- "output_tokens": output_tokens
- }
- },
- "facts": { "algorithm": "openai" }
- },
- "total_usd_micros": input_tokens + output_tokens
- }))
- .unwrap()
- }
+ use crate::test_support::{mark_run_running, test_usage};
async fn append_completed_stage(
run_store: &fabro_store::RunDatabase,
@@ -1525,27 +1491,6 @@ reasoning = false
.unwrap();
}
- async fn mark_run_running(run_store: &fabro_store::RunDatabase) {
- crate::event::append_event(run_store, &fixtures::RUN_1, &Event::RunStartRequested {
- resume: false,
- actor: None,
- })
- .await
- .unwrap();
- crate::event::append_event(run_store, &fixtures::RUN_1, &Event::RunRunnable {
- source: RunRunnableSource::StartRequested,
- actor: None,
- })
- .await
- .unwrap();
- crate::event::append_event(run_store, &fixtures::RUN_1, &Event::RunStarting)
- .await
- .unwrap();
- crate::event::append_event(run_store, &fixtures::RUN_1, &Event::RunRunning)
- .await
- .unwrap();
- }
-
async fn wait_for_conclusion(
run_store: &fabro_store::RunDatabase,
) -> crate::records::Conclusion {
@@ -1671,7 +1616,7 @@ reasoning = false
let (storage_root, run_dir) = storage_root_and_run_dir(&temp);
let (_persisted, store) = persisted_workflow(MINIMAL_DOT, &storage_root).await;
let run_store = store.open_run(&fixtures::RUN_1).await.unwrap();
- mark_run_running(&run_store).await;
+ mark_run_running(&run_store, &fixtures::RUN_1).await;
append_completed_stage(
&run_store,
"implement",
@@ -1717,12 +1662,12 @@ reasoning = false
}
#[tokio::test]
- async fn bootstrap_guard_failure_uses_conclusion_timing_and_billing_when_store_exists() {
+ async fn bootstrap_guard_failure_uses_conclusion_timing_and_billing() {
let temp = tempfile::tempdir().unwrap();
- let (storage_root, run_dir) = storage_root_and_run_dir(&temp);
+ let (storage_root, _run_dir) = storage_root_and_run_dir(&temp);
let (_persisted, store) = persisted_workflow(MINIMAL_DOT, &storage_root).await;
let run_store = store.open_run(&fixtures::RUN_1).await.unwrap();
- mark_run_running(&run_store).await;
+ mark_run_running(&run_store, &fixtures::RUN_1).await;
append_completed_stage(
&run_store,
"implement",
@@ -1734,13 +1679,12 @@ reasoning = false
let event_sink = RunEventSink::store(run_store.clone());
{
- let mut guard = DetachedRunBootstrapGuard::arm(
+ let _guard = DetachedRunBootstrapGuard::arm(
fixtures::RUN_1,
- &run_dir,
+ run_store_handle,
event_sink,
CancellationToken::new(),
);
- guard.run_store = Some(run_store_handle);
}
let conclusion = wait_for_conclusion(&run_store).await;
@@ -1762,7 +1706,7 @@ reasoning = false
let (storage_root, _run_dir) = storage_root_and_run_dir(&temp);
let (_persisted, store) = persisted_workflow(MINIMAL_DOT, &storage_root).await;
let run_store = store.open_run(&fixtures::RUN_1).await.unwrap();
- mark_run_running(&run_store).await;
+ mark_run_running(&run_store, &fixtures::RUN_1).await;
append_completed_stage(
&run_store,
"implement",
diff --git a/lib/crates/fabro-workflow/src/pipeline/finalize.rs b/lib/crates/fabro-workflow/src/pipeline/finalize.rs
index 53692daab..05d6db1d0 100644
--- a/lib/crates/fabro-workflow/src/pipeline/finalize.rs
+++ b/lib/crates/fabro-workflow/src/pipeline/finalize.rs
@@ -650,8 +650,8 @@ mod tests {
use fabro_store::{Database, EventEnvelope, RunDatabase, RunProjection};
use fabro_types::run_event::{MetadataSnapshotFailureKind, MetadataSnapshotPhase};
use fabro_types::{
- BilledModelUsage, BilledTokenCounts, EventBody, RunBlobId, RunEvent, RunId, RunSpec,
- StageCompletion, WorkflowSettings, first_event_seq, fixtures,
+ BilledTokenCounts, EventBody, RunBlobId, RunEvent, RunId, RunSpec, StageCompletion,
+ WorkflowSettings, first_event_seq, fixtures,
};
use object_store::memory::InMemory;
@@ -864,25 +864,7 @@ mod tests {
)
}
- fn test_usage(model_id: &str, input_tokens: i64, output_tokens: i64) -> BilledModelUsage {
- serde_json::from_value(serde_json::json!({
- "input": {
- "usage": {
- "model": {
- "provider": "openai",
- "model_id": model_id
- },
- "tokens": {
- "input_tokens": input_tokens,
- "output_tokens": output_tokens
- }
- },
- "facts": { "algorithm": "openai" }
- },
- "total_usd_micros": input_tokens + output_tokens
- }))
- .unwrap()
- }
+ use crate::test_support::test_usage;
#[test]
fn conclusion_stage_order_follows_projection_first_event_order() {
diff --git a/lib/crates/fabro-workflow/src/test_support.rs b/lib/crates/fabro-workflow/src/test_support.rs
index 7db2e1ddd..c587e49fa 100644
--- a/lib/crates/fabro-workflow/src/test_support.rs
+++ b/lib/crates/fabro-workflow/src/test_support.rs
@@ -53,6 +53,56 @@ async fn execute_and_emit_terminal(initialized: InitializedState) -> Executed {
executed
}
+/// Construct a fully-populated `BilledModelUsage` for tests. Centralised so
+/// callers don't keep rebuilding the same JSON skeleton.
+#[must_use]
+pub fn test_usage(
+ model_id: &str,
+ input_tokens: i64,
+ output_tokens: i64,
+) -> fabro_types::BilledModelUsage {
+ serde_json::from_value(serde_json::json!({
+ "input": {
+ "usage": {
+ "model": {
+ "provider": "openai",
+ "model_id": model_id
+ },
+ "tokens": {
+ "input_tokens": input_tokens,
+ "output_tokens": output_tokens
+ }
+ },
+ "facts": { "algorithm": "openai" }
+ },
+ "total_usd_micros": input_tokens + output_tokens
+ }))
+ .expect("test_usage JSON must deserialise")
+}
+
+/// Append the `RunStartRequested → RunRunnable → RunStarting → RunRunning`
+/// sequence so subsequent calls observe the run as live.
+pub async fn mark_run_running(run_store: &fabro_store::RunDatabase, run_id: &fabro_types::RunId) {
+ append_event(run_store, run_id, &Event::RunStartRequested {
+ resume: false,
+ actor: None,
+ })
+ .await
+ .expect("seed run.start_requested");
+ append_event(run_store, run_id, &Event::RunRunnable {
+ source: fabro_types::RunRunnableSource::StartRequested,
+ actor: None,
+ })
+ .await
+ .expect("seed run.runnable");
+ append_event(run_store, run_id, &Event::RunStarting)
+ .await
+ .expect("seed run.starting");
+ append_event(run_store, run_id, &Event::RunRunning)
+ .await
+ .expect("seed run.running");
+}
+
pub fn test_store_dir(run_dir: &std::path::Path) -> PathBuf {
let mut hasher = std::collections::hash_map::DefaultHasher::new();
std::process::id().hash(&mut hasher);

View file

@ -0,0 +1,6 @@
{
"outcome": "succeeded",
"notes": "Stage completed: simplify_opus",
"failure_reason": null,
"timestamp": "2026-05-25T22:48:07.333773Z"
}

View file

@ -0,0 +1,189 @@
Goal: # Plan: Fix stage timing (inference + tool) reporting
## Context
The web UI's Duration popover shows `Active (inference + tools): 0ms` for every run, including agent-heavy runs that obviously did substantial LLM and tool work. Verified on `01KSE2PAVXD56N4TWNK4T5H5VA`: 10 stage.completed events and 1 run.failed event all carry `inference_time_ms: 0, tool_time_ms: 0`, even though stages like `implement@1` (94 min wall) and `simplify_opus@1` (29 min wall) were doing nothing but inference and tool calls.
Two independent bugs:
1. **No production handler ever populates `Outcome.timing`.** The plumbing from `Outcome.timing` → `NodeResult` (`lib/crates/fabro-core/src/executor.rs:30-37`) → `StageTiming` → `stage.completed` props → projection → billing rollup → run.completed/failed → UI is fully wired and shipped as of #343 (2026-05-21), but `AgentHandler::execute`, `PromptHandler::execute`, `CommandHandler::execute`, and `FanInHandler` all build `Outcome::success()` and never touch `.timing`. The executor falls back to zero, and every downstream consumer faithfully aggregates zero.
2. **`persist_terminal_engine_failure` and its sibling Drop-guard failure paths discard timing/billing entirely.** When the engine returns `Err` (e.g. `VisitLimitExceeded`, which is what killed the user's run), `lib/crates/fabro-workflow/src/operations/start.rs:284-308` builds a `Conclusion` via `build_conclusion_from_store`, then throws it away (`let _conclusion = ...`) and emits `WorkflowRunFailed` with `RunTiming::wall_only(...)` and `None` for billing/diff. The three Drop-guard paths (`start.rs:934`, `1001`, `1033`) do similar with `RunTiming::default()` and never even build a conclusion.
Goal: stage and run events carry real per-stage `inference_time_ms` + `tool_time_ms`; engine-failure terminal events preserve the conclusion's rolled-up timing and billing.
## Approach
### Part A — Capture inference + tool time in handlers (Bug 1)
**A1. `fabro-agent` — accumulate per-input timing in `Session`**
`lib/crates/fabro-agent/src/session.rs`
Add two `Duration` accumulators to `Session` (initialised to `Duration::ZERO`):
- `last_input_inference_duration`
- `last_input_tool_duration`
In `process_input_with_runtime` (line 1196), zero them at entry so each call's totals are independent.
In `run_single_input` (line 1254):
- Wrap the inference span: capture `Instant::now()` immediately before opening the stream at line 1391, and add `.elapsed()` to `last_input_inference_duration` once `response = Some(resp)` (line 1487-1490) OR when the loop exits with an error/cancellation. The whole `'streamattempts` loop counts as inference work — retries included.
- Wrap the tool span around `execute_tool_calls` at line 1705-1719: `Instant::now()` before, accumulate `.elapsed()` after `.await`.
Expose a getter:
```rust
pub fn last_input_timing(&self) -> SessionInputTiming { ... }
```
where `SessionInputTiming { pub inference: Duration, pub tool: Duration }` is a new tiny struct in `fabro-agent`.
**A2. `fabro-workflow` — thread timing through the backend boundary**
`lib/crates/fabro-workflow/src/handler/agent.rs`
Extend `CodergenResult::Text` with a `timing: fabro_types::StageTiming` field (wall is irrelevant — see note below). Update the few `CodergenResult::Text { ... }` constructions found by the explore agent to populate it; existing match-bindings only read `text`/`usage`/`files_touched` so they keep compiling with `..` patterns. `CodergenResult::Full(outcome)` keeps current behaviour — the outcome itself already carries any timing.
Note on wall: `lib/crates/fabro-core/src/executor.rs:30-37` reads ONLY `inference_time_ms` and `tool_time_ms` out of `outcome.timing`. The wall comes from the executor's own stopwatch. So we construct `StageTiming::new(0, inference_ms, tool_ms)` and document that the wall field is ignored in this hop.
`lib/crates/fabro-workflow/src/handler/llm/api.rs`
- `AgentApiBackend::run` (line 1103): after `session.process_input_with_runtime(...)` returns, read `session.last_input_timing()` and set the new `timing` on `CodergenResult::Text` at line 1094.
- `AgentApiBackend::one_shot` (line 994): wrap the `complete_one_shot_request` call at line 1053 with `Instant::now()` / `.elapsed()`. Accumulate across repair iterations of the surrounding loop. All of it counts as inference; no tool work happens in `one_shot`. Set `timing` on `CodergenResult::Text` at line 1094.
`lib/crates/fabro-workflow/src/handler/llm/acp.rs`
`AgentAcpBackend::run` (line ~140): already exposes `result.duration_ms`. Set `timing: StageTiming::new(0, duration_ms, 0)` on the returned `CodergenResult::Text` (per user decision: attribute all ACP duration to inference; ACP is opaque about the split).
**A3. Consume timing in stage handlers and set `outcome.timing`**
- `lib/crates/fabro-workflow/src/handler/agent.rs:341` — after building `outcome`, before the final `Ok(outcome)`, set `outcome.timing = Some(timing_from_codergen_result)`.
- `lib/crates/fabro-workflow/src/handler/prompt.rs:180` — same pattern.
- `lib/crates/fabro-workflow/src/handler/fan_in.rs:266` — backend returns timing; pass it onto the outcome built from the fan-in response.
- `lib/crates/fabro-workflow/src/handler/command.rs:175` — `outcome.timing = Some(StageTiming::new(0, 0, result.duration_ms))`. All command wall-time is tool time. `result.duration_ms` is already at line 154 in scope.
Other handlers (`human`, `wait`, `conditional`, `parallel`, `start`, `exit`, `structured_output`, `manager_loop`) do no inference or tool work. Leave `outcome.timing` as `None`; the executor will naturally produce `inference: 0, tool: 0` for those stages, which is correct.
### Part B — Preserve conclusion timing on engine failure (Bug 2)
`lib/crates/fabro-workflow/src/operations/start.rs`
**B1. Main path** (`persist_terminal_engine_failure`, line 274-308):
- Rename `_conclusion` → `conclusion` and use it:
- Pass `conclusion.timing` (already a `RunTiming` with the proper inference/tool/wall rollup from `build_conclusion_from_parts`) instead of `RunTiming::wall_only(...)`.
- Pass `conclusion.billing.clone()` instead of `None` for the billing arg of `workflow_run_failed_from_error`.
- `final_git_commit_sha`, `final_patch`, `diff_summary` stay `None` — those require the finalize-side workspace diff computation that this path deliberately skips.
**B2. Drop-guard paths** (per user decision: fix them too):
- `DetachedRunBootstrapGuard` (line 882-948): add an `Option<RunStoreHandle>` field. The bootstrap function builds the guard before the store exists, then mutates `bootstrap_guard.run_store = Some(store.clone())` once the store is in scope. On Drop, if the store is `Some`, the spawned task calls `build_conclusion_from_store` and uses its timing/billing; otherwise falls back to `RunTiming::default()` (pre-store failure means no stages can possibly exist).
- `DetachedRunCompletionGuard` (line 953-1021): armed after the store exists, so add a non-optional `run_store: RunStoreHandle`. Drop's spawned task builds the conclusion and uses it.
- `persist_detached_failure` (line 1023): add a `run_store: &RunStoreHandle` parameter. Call `build_conclusion_from_store` and forward `timing` + `billing` to the failure event. Update the two callers (postrun-related) to pass the store they already have in scope.
`RunStoreHandle` is already `Clone` (the surrounding code clones it routinely), so move-into-spawned-task is fine.
### Critical existing utilities to reuse (do not duplicate)
- `fabro_types::StageTiming::new(wall, inference, tool)` and `RunTiming::new(...)` — invariant-enforcing constructors at `lib/crates/fabro-types/src/timing.rs:38, 91`.
- `crate::millis_u64(duration)` helper for `Duration → u64` ms in `fabro-workflow` (used widely; see `lifecycle/event.rs:80-86`).
- `build_conclusion_from_store` at `lib/crates/fabro-workflow/src/pipeline/finalize.rs:71` already does the rollup we need on the engine-failure path.
- `billing_rollup_from_projection` (called inside `build_conclusion_from_parts`) sums per-stage timings into `RunTiming` — no need to reimplement.
## Tests
- **`fabro-agent` unit test**: feed `Session` a fake `LlmClient` whose `stream` sleeps a known duration and a fake tool that sleeps another known duration. Drive one `process_input_with_runtime` call. Assert `session.last_input_timing()` reports both non-zero and roughly matching the sleeps. Then call again and assert it's per-call (not cumulative).
- **`fabro-workflow` handler tests**: in `handler/agent.rs`'s test module, wire a `CodergenBackend` that returns `CodergenResult::Text { timing: StageTiming::new(0, 200, 300), .. }` and assert `AgentHandler::execute`'s returned `Outcome.timing` carries those values. Mirror for `prompt.rs` and `fan_in.rs`. Add a `command.rs` test that mocks a `sandbox.exec_command_streaming` returning `duration_ms = 500` and asserts `outcome.timing.tool_time_ms == 500`.
- **Executor integration**: add a test in `fabro-workflow` (or extend an existing one in `pipeline/finalize.rs` tests) that runs a tiny graph with a handler producing `Outcome.timing = Some(StageTiming::new(0, 100, 50))` and asserts the emitted `stage.completed` event carries those values, and that `run.completed` carries the summed rollup.
- **`persist_terminal_engine_failure` test**: seed a `RunStore` with a couple of `stage.completed` events whose timing is non-zero, drive the engine-failure path, and assert the emitted `WorkflowRunFailed` event has `timing.inference_time_ms` and `tool_time_ms` matching the per-stage sum and `billing` populated.
- **Drop guard tests**: trickier because of `Handle::try_current` + spawn. Add focused tests that arm a guard, drop it, and `tokio::task::yield_now().await` enough times to let the spawned task run, then assert the emitted failure event carries non-zero timing.
- Run `cargo nextest run -p fabro-agent -p fabro-workflow -p fabro-store -p fabro-core`.
- Run formatter and lints per CLAUDE.md: `cargo +nightly-2026-04-14 fmt --check --all` and `cargo +nightly-2026-04-14 clippy --workspace --all-targets -- -D warnings`.
## End-to-end verification
1. Build the server: `cargo build -p fabro-server`.
2. Start server: `fabro server start`.
3. Run a small agent-backed workflow (e.g. `fabro run repl` with a short prompt that fires at least one tool call).
4. `fabro events <run_id> --json | jq -s '[.[] | select(.event=="stage.completed")] | .[].properties.timing'` — confirm `inference_time_ms > 0` and `tool_time_ms > 0` for the agent stage.
5. `fabro events <run_id> --json | jq -s '[.[] | select(.event=="run.completed" or .event=="run.failed")] | .[].properties.timing'` — confirm `active_time_ms == inference_time_ms + tool_time_ms` and both are non-zero.
6. Open the run in the web UI (start the SPA dev build per CLAUDE.md or rebuild the embedded SPA with `cargo dev build`), hover the Duration chip, confirm **Active (inference + tools)** is non-zero.
7. For Bug 2: force an engine failure by setting a very low visit limit and rerunning the same workflow; confirm the `run.failed` event timing breakdown is non-zero and matches the per-stage sum.
## Out of scope
- Adding `wall_time_ms` correctness to `Outcome.timing` (executor ignores it; doc tweak only if necessary).
- Surfacing inference vs tool split for ACP backend beyond "all-inference" attribution.
- Backfilling timing for historical runs that have already emitted zero events — past events are immutable.
- Web UI changes beyond what the existing popover already renders.
## Completed stages
- **toolchain**: succeeded
- 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`
- Output:
```
cargo 1.95.0 (f2d3ce0bd 2026-03-21)
```
- **preflight_compile**: succeeded
- Script: `cargo check -q --workspace 2>&1`
- Output: (empty)
- **preflight_lint**: succeeded
- Script: `cargo +nightly-2026-04-14 clippy -q --workspace --all-targets -- -D warnings 2>&1`
- Output: (empty)
- **implement**: succeeded
- Model: gpt-5.5, 2.8m tokens in / 19.0k out
- **simplify_opus**: succeeded
- Model: claude-opus-4-7, 146.9k tokens in / 49.4k out
- Files: /home/daytona/workspace/fabro/lib/crates/fabro-agent/src/session.rs, /home/daytona/workspace/fabro/lib/crates/fabro-types/src/timing.rs, /home/daytona/workspace/fabro/lib/crates/fabro-workflow/src/billing_rollup.rs, /home/daytona/workspace/fabro/lib/crates/fabro-workflow/src/event/convert.rs, /home/daytona/workspace/fabro/lib/crates/fabro-workflow/src/handler/agent.rs, /home/daytona/workspace/fabro/lib/crates/fabro-workflow/src/handler/command.rs, /home/daytona/workspace/fabro/lib/crates/fabro-workflow/src/handler/llm/acp.rs, /home/daytona/workspace/fabro/lib/crates/fabro-workflow/src/handler/llm/api.rs, /home/daytona/workspace/fabro/lib/crates/fabro-workflow/src/handler/prompt.rs, /home/daytona/workspace/fabro/lib/crates/fabro-workflow/src/operations/start.rs, /home/daytona/workspace/fabro/lib/crates/fabro-workflow/src/pipeline/finalize.rs, /home/daytona/workspace/fabro/lib/crates/fabro-workflow/src/test_support.rs
# Simplify: Code Review and Cleanup
Review changes vs. origin for reuse, quality, and efficiency. Fix any issues found.
## Phase 1: Identify Changes
Run git diff (or git diff HEAD if there are staged changes) to see what changed. If there are no git changes, review the most recently modified files that the user mentioned or that you edited earlier in this conversation.
## Phase 2: Launch Three Review Agents in Parallel
Use the Agent tool to launch all three agents concurrently in a single message. Pass each agent the full diff so it has the complete context.
### Agent 1: Code Reuse Review
For each change:
1. Search for existing utilities and helpers that could replace newly written code. Use Grep to find similar patterns elsewhere in the codebase — common locations are utility directories, shared modules, and files adjacent to the changed ones.
2. Flag any new function that duplicates existing functionality. Suggest the existing function to use instead.
3. Flag any inline logic that could use an existing utility — hand-rolled string manipulation, manual path handling, custom environment checks, ad-hoc type guards, and similar patterns are common candidates.
Note: This is a greenfield app, so focus on maximizing simplicity and don't worry about changing things to achieve it.
### Agent 2: Code Quality Review
Review the same changes for hacky patterns:
1. Redundant state: state that duplicates existing state, cached values that could be derived, observers/effects that could be direct calls
2. Parameter sprawl: adding new parameters to a function instead of generalizing or restructuring existing ones
3. Copy-paste with slight variation: near-duplicate code blocks that should be unified with a shared abstraction
4. Leaky abstractions: exposing internal details that should be encapsulated, or breaking existing abstraction boundaries
5. Stringly-typed code: using raw strings where constants, enums (string unions), or branded types already exist in the codebase
Note: This is a greenfield app, so be aggressive in optimizing quality.
### Agent 3: Efficiency Review
Review the same changes for efficiency:
1. Unnecessary work: redundant computations, repeated file reads, duplicate network/API calls, N+1 patterns
2. Missed concurrency: independent operations run sequentially when they could run in parallel
3. Hot-path bloat: new blocking work added to startup or per-request/per-render hot paths
4. Unnecessary existence checks: pre-checking file/resource existence before operating (TOCTOU anti-pattern) — operate directly and handle the error
5. Memory: unbounded data structures, missing cleanup, event listener leaks
6. Overly broad operations: reading entire files when only a portion is needed, loading all items when filtering for one
## Phase 3: Fix Issues
Wait for all three agents to complete. Aggregate their findings and fix each issue directly. If a finding is a false positive or not worth addressing, note it and move on — do not argue with the finding, just skip it.
When done, briefly summarize what was fixed (or confirm the code was already clean).

View file

@ -0,0 +1,5 @@
{
"mode": "agent",
"provider": "openai",
"model": "gpt-5.5"
}

View file

@ -0,0 +1,19 @@
Reviewed the branch diff vs merge base and ran the parallel review pass. I found and fixed one cleanup issue:
- Restored `Session::process_input_with_runtime` to its original public `Result<(), Error>` shape instead of returning a tuple.
- Added the requested `Session::last_input_timing()` getter and `last_input_timing` state on `Session`.
- Updated workflow callers to read timing through the getter after `process_input_with_runtime(...)`.
- Updated the agent timing unit test accordingly.
Validation passed:
- `cargo check -q -p fabro-agent -p fabro-workflow`
- `cargo nextest run -p fabro-agent -p fabro-workflow`
- `cargo nextest run -p fabro-store -p fabro-core`
- `cargo +nightly-2026-04-14 fmt --check --all`
- `cargo +nightly-2026-04-14 clippy -q -p fabro-agent -p fabro-workflow --all-targets -- -D warnings`
Working tree now has the cleanup edits in:
- `lib/crates/fabro-agent/src/session.rs`
- `lib/crates/fabro-workflow/src/handler/llm/api.rs`