From caba862a331567ff9c6a5f737fc6b48ecf5f20f6 Mon Sep 17 00:00:00 2001 From: "fabro-sh-0530[bot]" <281434857+fabro-sh-0530[bot]@users.noreply.github.com> Date: Sat, 23 May 2026 13:31:28 -0400 Subject: [PATCH] Show live timing for in-flight runs via read-time overlay (#361) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit ## Summary `Run.timing` was only populated for terminal runs, so the duration chip and popover were hidden for every queued, running, or blocked run. This change derives a best-effort `RunTiming` at cache read time for started-but-not-terminal runs, making the duration chip appear and tick for in-flight runs everywhere the UI consumes `summary.timing`. ### Plan Summary - Add `RunProjection::live_run_timing(now)` in `fabro-types`: wall time from `start.start_time → now`; active time summed from completed stages' `StageTiming`. - Add `apply_read_overlays(entry, now)` in `projection_cache.rs`: called after mutex release on cloned entries; only fills `timing` when it is `None` (terminal runs are unaffected). - Thread `now: DateTime` through `get_summary` / `list` / `list_cached_runs` / `list_runs` / `list_runs_with_projection` in `slate/mod.rs` and `projection_cache.rs`. - Propagate `Utc::now()` at all HTTP handler and background-task call sites in `fabro-server` and `fabro-workflow`. - Add unit tests covering the three key cases: not-yet-started (returns `None`), in-flight with completed stages, and a terminal run whose live derivation matches `Conclusion.timing`. ### Key design decisions **Overlay happens outside the cache mutex on a cloned copy.** The cached `CachedRunProjection.summary.timing` is never mutated; only the cloned value returned to callers gets the overlay. This means `get_cached_run` (raw cache access, no `now`) still returns `None` for in-flight timing — confirmed by the new integration test `cached_summary_overlays_live_timing_without_mutating_cached_snapshot`. **Known limitation — active time steps, not ticks.** `StageProjection` records inference/tool time only at stage completion, so `active_time_ms` reflects the sum of *completed* stages and jumps forward when a stage finishes. `wall_time_ms` advances continuously. This is intentional and documented in the code comment; live per-stage inference tracking is out of scope. **`build_summary` and `Conclusion.timing` are untouched.** All five internal readers that want "timing as of conclusion" continue to use `Conclusion.timing` directly; no test churn from the existing call sites. ### Fabro Details
Ran 9 stages in 38m 28s for $10.35 | Stage | Duration | Cost | Retries | |---|---|---|---| | start | 0s | – | 0 | | toolchain | 2s | – | 0 | | preflight_compile | 2m 4s | – | 0 | | preflight_lint | 2m 20s | – | 0 | | implement | 11m 27s | $6.21 | 0 | | simplify_opus | 13m 12s | $2.48 | 0 | | simplify_gpt | 4m 20s | $1.66 | 0 | | verify | 4m 9s | – | 0 | | fmt | 3s | – | 0 | | **Total** | **38m 28s** | **$10.35** | **0** |
Ran ImplementPlan.fabro (12 nodes and 15 edges) ```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="cargo +nightly-2026-04-14 clippy -q --workspace --all-targets -- -D warnings 2>&1 && cargo nextest run --cargo-quiet --workspace --status-level fail 2>&1 && cargo dev docs refresh 2>&1 && cargo dev docs check 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 clippy lint warnings, test failures, and generated docs errors.", max_visits=3] fmt [label="Format", shape=parallelogram, script="cargo +nightly-2026-04-14 fmt --all 2>&1", max_retries=0] 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 -> fmt [condition="outcome=succeeded"] verify -> fixup fixup -> verify fmt -> exit } ```
⚒️ Generated with [Fabro](https://fabro.sh) --------- Co-authored-by: Fabro Co-authored-by: Bryan Helmkamp --- lib/crates/fabro-server/src/server.rs | 2 +- .../src/server/handler/lifecycle.rs | 4 +- .../fabro-server/src/server/handler/runs.rs | 45 +++-- .../fabro-server/src/server/handler/system.rs | 6 +- .../src/server/resource_sampler.rs | 2 +- lib/crates/fabro-store/src/run_state.rs | 118 +++++++++++- lib/crates/fabro-store/src/slate/mod.rs | 168 +++++++++++++----- .../fabro-store/src/slate/projection_cache.rs | 41 ++++- lib/crates/fabro-types/src/run_projection.rs | 29 ++- lib/crates/fabro-workflow/src/run_lookup.rs | 2 +- .../tests/it/daytona_integration.rs | 42 +++-- .../fabro-workflow/tests/it/integration.rs | 46 +++-- 12 files changed, 389 insertions(+), 116 deletions(-) diff --git a/lib/crates/fabro-server/src/server.rs b/lib/crates/fabro-server/src/server.rs index 9fd3b990f..b512149e6 100644 --- a/lib/crates/fabro-server/src/server.rs +++ b/lib/crates/fabro-server/src/server.rs @@ -2642,7 +2642,7 @@ pub(crate) async fn reconcile_incomplete_runs_on_startup( ) -> anyhow::Result { let summaries = state .store - .list_runs(&fabro_store::ListRunsQuery::default()) + .list_runs(&fabro_store::ListRunsQuery::default(), chrono::Utc::now()) .await?; let mut reconciled = 0usize; diff --git a/lib/crates/fabro-server/src/server/handler/lifecycle.rs b/lib/crates/fabro-server/src/server/handler/lifecycle.rs index 8360dddf3..569da6d4a 100644 --- a/lib/crates/fabro-server/src/server/handler/lifecycle.rs +++ b/lib/crates/fabro-server/src/server/handler/lifecycle.rs @@ -1,5 +1,7 @@ use std::sync::Arc; +use chrono::Utc; + use super::super::{ ApiError, AppState, FailureReason, ForkRequest, ForkResponse, IntoResponse, Json, Path, Principal, RequireRunScopedOrRunTools, RequiredUser, Response, RewindRequest, RewindResponse, @@ -24,7 +26,7 @@ pub(super) fn routes() -> Router> { } async fn run_response(state: &AppState, id: RunId, status: StatusCode) -> Response { - match state.store.get_cached_summary(&id).await { + match state.store.get_cached_summary(&id, Utc::now()).await { Ok(Some(summary)) => { (status, Json(state.decorate_run_summary(summary).await)).into_response() } diff --git a/lib/crates/fabro-server/src/server/handler/runs.rs b/lib/crates/fabro-server/src/server/handler/runs.rs index 930d3b2cd..0f1057071 100644 --- a/lib/crates/fabro-server/src/server/handler/runs.rs +++ b/lib/crates/fabro-server/src/server/handler/runs.rs @@ -11,6 +11,7 @@ use axum::{Json, Router}; use base64::Engine as _; use base64::engine::general_purpose::STANDARD as BASE64_STANDARD; use bytes::Bytes; +use chrono::Utc; use fabro_api::types::{ BoardColumn, BoardColumnDefinition, RunManifest, SubmitAnswerRequest, UpdateRunParentRequest, UpdateRunRequest, @@ -154,10 +155,13 @@ async fn list_board_runs( ) -> Response { let entries = match state .store - .list_cached_runs(&fabro_store::ListRunsQuery { - parent_id: params.parent_id, - ..fabro_store::ListRunsQuery::default() - }) + .list_cached_runs( + &fabro_store::ListRunsQuery { + parent_id: params.parent_id, + ..fabro_store::ListRunsQuery::default() + }, + Utc::now(), + ) .await { Ok(runs) => runs, @@ -213,7 +217,7 @@ async fn link_run_parent( } }; let _parent_link_guard = state.parent_link_lock.lock().await; - let child = match state.store.get_cached_summary(&child_id).await { + let child = match state.store.get_cached_summary(&child_id, Utc::now()).await { Ok(Some(summary)) => summary, Ok(None) => return ApiError::not_found("Run not found.").into_response(), Err(err) => { @@ -259,7 +263,7 @@ async fn unlink_run_parent( State(state): State>, ) -> Response { let _parent_link_guard = state.parent_link_lock.lock().await; - let child = match state.store.get_cached_summary(&child_id).await { + let child = match state.store.get_cached_summary(&child_id, Utc::now()).await { Ok(Some(summary)) => summary, Ok(None) => return ApiError::not_found("Run not found.").into_response(), Err(err) => { @@ -309,7 +313,7 @@ async fn validate_parent_link( } let summary = state .store - .get_cached_summary(¤t_id) + .get_cached_summary(¤t_id, Utc::now()) .await .map_err(|err| ApiError::new(StatusCode::INTERNAL_SERVER_ERROR, err.to_string()))?; let Some(summary) = summary else { @@ -324,7 +328,7 @@ async fn validate_parent_link( } async fn updated_run_response(state: &AppState, run_id: &RunId) -> Response { - match state.store.get_cached_summary(run_id).await { + match state.store.get_cached_summary(run_id, Utc::now()).await { Ok(Some(summary)) => ( StatusCode::OK, Json(state.decorate_run_summary(summary).await), @@ -344,10 +348,13 @@ async fn list_runs( ) -> Response { match state .store - .list_cached_runs(&fabro_store::ListRunsQuery { - parent_id: params.parent_id, - ..fabro_store::ListRunsQuery::default() - }) + .list_cached_runs( + &fabro_store::ListRunsQuery { + parent_id: params.parent_id, + ..fabro_store::ListRunsQuery::default() + }, + Utc::now(), + ) .await { Ok(entries) => { @@ -415,7 +422,7 @@ async fn resolve_run( ) -> Response { let runs = match state .store - .list_runs(&fabro_store::ListRunsQuery::default()) + .list_runs(&fabro_store::ListRunsQuery::default(), Utc::now()) .await { Ok(runs) => runs, @@ -490,7 +497,7 @@ async fn update_run( Ok(title) => title, Err(err) => return ApiError::bad_request(err.to_string()).into_response(), }; - let current = match state.store.get_cached_summary(&id).await { + let current = match state.store.get_cached_summary(&id, Utc::now()).await { Ok(Some(summary)) => summary, Ok(None) => return ApiError::not_found("Run not found.").into_response(), Err(err) => { @@ -523,7 +530,7 @@ async fn update_run( return ApiError::new(StatusCode::INTERNAL_SERVER_ERROR, err.to_string()).into_response(); } - match state.store.get_cached_summary(&id).await { + match state.store.get_cached_summary(&id, Utc::now()).await { Ok(Some(summary)) => ( StatusCode::OK, Json(state.decorate_run_summary(summary).await), @@ -607,7 +614,11 @@ async fn create_run( } }; let created_at = created.run_id.created_at(); - let summary = match state.store.get_cached_summary(&created.run_id).await { + let summary = match state + .store + .get_cached_summary(&created.run_id, Utc::now()) + .await + { Ok(Some(summary)) => summary, Ok(None) => return ApiError::not_found("Run not found.").into_response(), Err(err) => { @@ -741,7 +752,7 @@ async fn get_run_status( RequireRunScopedOrRunTools(id, _actor): RequireRunScopedOrRunTools, State(state): State>, ) -> Response { - match state.store.get_cached_summary(&id).await { + match state.store.get_cached_summary(&id, Utc::now()).await { Ok(Some(run)) => { (StatusCode::OK, Json(state.decorate_run_summary(run).await)).into_response() } diff --git a/lib/crates/fabro-server/src/server/handler/system.rs b/lib/crates/fabro-server/src/server/handler/system.rs index 3baa77fb5..c32112ab0 100644 --- a/lib/crates/fabro-server/src/server/handler/system.rs +++ b/lib/crates/fabro-server/src/server/handler/system.rs @@ -1,5 +1,7 @@ use std::sync::Arc; +use chrono::Utc; + use super::super::{ AggregateBilling, AggregateBillingTotals, ApiError, AppState, BilledTokenCounts, BillingByModel, DfParams, FABRO_VERSION, GithubIntegrationStrategy, IntoResponse, Json, Path, @@ -97,7 +99,7 @@ async fn get_system_df( let storage_dir = state.server_storage_dir(); let summaries = match state .store - .list_runs(&fabro_store::ListRunsQuery::default()) + .list_runs(&fabro_store::ListRunsQuery::default(), Utc::now()) .await { Ok(summaries) => summaries, @@ -162,7 +164,7 @@ async fn prune_runs( let storage_dir = state.server_storage_dir(); let summaries = match state .store - .list_runs(&fabro_store::ListRunsQuery::default()) + .list_runs(&fabro_store::ListRunsQuery::default(), Utc::now()) .await { Ok(summaries) => summaries, diff --git a/lib/crates/fabro-server/src/server/resource_sampler.rs b/lib/crates/fabro-server/src/server/resource_sampler.rs index bc083c2cb..2001e6de9 100644 --- a/lib/crates/fabro-server/src/server/resource_sampler.rs +++ b/lib/crates/fabro-server/src/server/resource_sampler.rs @@ -138,7 +138,7 @@ impl ResourceSampler { let summaries = state .store - .list_runs(&fabro_store::ListRunsQuery::default()) + .list_runs(&fabro_store::ListRunsQuery::default(), chrono::Utc::now()) .await .context("failed to list runs for resource sampling")?; let storage_path = storage_path.to_path_buf(); diff --git a/lib/crates/fabro-store/src/run_state.rs b/lib/crates/fabro-store/src/run_state.rs index 0fdd4dd3a..eccae3995 100644 --- a/lib/crates/fabro-store/src/run_state.rs +++ b/lib/crates/fabro-store/src/run_state.rs @@ -1023,7 +1023,7 @@ fn merge_agent_process_output(stdout: &str, stderr: &str) -> String { mod tests { use std::collections::{BTreeMap, HashMap}; - use chrono::Utc; + use chrono::{DateTime, Utc}; use fabro_types::run_event::run::RunFailedProps; use fabro_types::run_event::{ AgentAcpCancelledProps, AgentAcpCompletedProps, AgentAcpStartedProps, @@ -1032,9 +1032,9 @@ mod tests { AgentSessionStartedProps, AgentSkillActivatedProps, AgentSkillActivationSource, AgentSkillSummary, AgentSkillsDiscoveredProps, AgentSubClosedProps, AgentSubCompletedProps, AgentSubFailedProps, AgentSubSpawnedProps, CheckpointCompletedProps, - InterviewCompletedProps, InterviewOption, InterviewStartedProps, RunControlEffectProps, - StageCompletedProps, StageFailedProps, StagePromptProps, StageRetryingProps, - StageStartedProps, + InterviewCompletedProps, InterviewOption, InterviewStartedProps, RunCompletedProps, + RunControlEffectProps, StageCompletedProps, StageFailedProps, StagePromptProps, + StageRetryingProps, StageStartedProps, }; use fabro_types::{ AgentBackend, BilledModelUsage, BilledTokenCounts, BlockedReason, Checkpoint, @@ -1124,6 +1124,10 @@ mod tests { RunProjection::new("Test run".to_string(), test_run_spec(), Utc::now()) } + fn test_dt(value: &str) -> DateTime { + value.parse().unwrap() + } + fn running_projection() -> RunProjection { let mut state = initialized_projection(); state @@ -1176,6 +1180,112 @@ mod tests { } } + #[test] + fn live_run_timing_returns_none_before_run_starts() { + let state = initialized_projection(); + + assert_eq!(state.live_run_timing(Utc::now()), None); + } + + #[test] + fn live_run_timing_derives_wall_and_completed_stage_active_for_in_flight_run() { + let mut state = initialized_projection(); + state + .apply_event(&test_raw_event_at( + 1, + "2026-04-07T12:00:00Z", + "run.started", + &json!({ "name": "Test run" }), + None, + )) + .unwrap(); + state.stage_entry("plan", 1, first_event_seq(2)).timing = + Some(fabro_types::StageTiming::new(2_000, 700, 300)); + state.stage_entry("code", 1, first_event_seq(3)).timing = + Some(fabro_types::StageTiming::new(3_000, 20, 80)); + state.stage_entry("running", 1, first_event_seq(4)).timing = None; + + assert_eq!( + state.live_run_timing(test_dt("2026-04-07T12:00:12.345Z")), + Some(fabro_types::RunTiming::new(12_345, 720, 380)) + ); + } + + #[test] + fn live_run_timing_matches_conclusion_timing_at_conclusion_moment() { + let mut state = initialized_projection(); + let started_at = test_dt("2026-04-07T12:00:00Z"); + let completed_at = test_dt("2026-04-07T12:00:10Z"); + state + .apply_event(&test_raw_event_at( + 1, + "2026-04-07T12:00:00Z", + "run.started", + &json!({ "name": "Test run" }), + None, + )) + .unwrap(); + state + .apply_event(&test_raw_event_at( + 2, + "2026-04-07T12:00:00Z", + "run.starting", + &json!({}), + None, + )) + .unwrap(); + state + .apply_event(&test_raw_event_at( + 3, + "2026-04-07T12:00:01Z", + "run.running", + &json!({}), + None, + )) + .unwrap(); + state.stage_entry("plan", 1, first_event_seq(4)).timing = + Some(fabro_types::StageTiming::new(2_000, 700, 300)); + state.stage_entry("code", 1, first_event_seq(5)).timing = + Some(fabro_types::StageTiming::new(3_000, 50, 200)); + + let conclusion_timing = fabro_types::RunTiming::new( + u64::try_from( + completed_at + .signed_duration_since(started_at) + .num_milliseconds(), + ) + .unwrap(), + 750, + 500, + ); + let mut completed = test_event( + 6, + EventBody::RunCompleted(RunCompletedProps { + timing: conclusion_timing, + artifact_count: 0, + status: "succeeded".to_string(), + reason: SuccessReason::Completed, + total_usd_micros: None, + final_git_commit_sha: None, + final_patch: None, + diff_summary: None, + billing: None, + }), + None, + ); + completed.event.ts = completed_at; + state.apply_event(&completed).unwrap(); + + assert_eq!( + state + .conclusion + .as_ref() + .map(|conclusion| conclusion.timing), + Some(conclusion_timing) + ); + assert_eq!(state.live_run_timing(completed_at), Some(conclusion_timing)); + } + #[test] fn last_event_at_tracks_most_recent_event_timestamp() { let mut state = initialized_projection(); diff --git a/lib/crates/fabro-store/src/slate/mod.rs b/lib/crates/fabro-store/src/slate/mod.rs index 2b0376a85..15807a72c 100644 --- a/lib/crates/fabro-store/src/slate/mod.rs +++ b/lib/crates/fabro-store/src/slate/mod.rs @@ -196,9 +196,9 @@ impl Database { RunDatabase::open_reader(*run_id, db, Arc::clone(&self.projection_cache)).await } - pub async fn list_runs(&self, query: &ListRunsQuery) -> Result> { + pub async fn list_runs(&self, query: &ListRunsQuery, now: DateTime) -> Result> { Ok(self - .list_cached_runs(query) + .list_cached_runs(query, now) .await? .into_iter() .map(|entry| entry.summary) @@ -208,9 +208,10 @@ impl Database { pub async fn list_runs_with_projection( &self, query: &ListRunsQuery, + now: DateTime, ) -> Result> { Ok(self - .list_cached_runs(query) + .list_cached_runs(query, now) .await? .into_iter() .map(|entry| (entry.summary, (*entry.projection).clone())) @@ -250,9 +251,10 @@ impl Database { pub async fn list_cached_runs( &self, query: &ListRunsQuery, + now: DateTime, ) -> Result> { self.warm_projection_cache().await?; - Ok(self.projection_cache.list(query).await) + Ok(self.projection_cache.list(query, now).await) } pub async fn list_unreadable_runs(&self) -> Result> { @@ -292,9 +294,13 @@ impl Database { Ok(self.projection_cache.get(run_id).await) } - pub async fn get_cached_summary(&self, run_id: &RunId) -> Result> { + pub async fn get_cached_summary( + &self, + run_id: &RunId, + now: DateTime, + ) -> Result> { self.warm_projection_cache().await?; - Ok(self.projection_cache.get_summary(run_id).await) + Ok(self.projection_cache.get_summary(run_id, now).await) } pub async fn put_session_run_index( @@ -428,11 +434,11 @@ impl Runs { } pub async fn find(&self, run_id: &RunId) -> Result> { - self.db.get_cached_summary(run_id).await + self.db.get_cached_summary(run_id, Utc::now()).await } pub async fn list(&self, query: &ListRunsQuery) -> Result> { - self.db.list_runs(query).await + self.db.list_runs(query, Utc::now()).await } } @@ -669,7 +675,10 @@ mod tests { append_completed(&run_1, "run-1", dt("2026-03-27T12:00:00Z")).await; append_created(&run_2, "run-2", dt("2026-03-27T12:00:10Z")).await; - let summary = store.list_runs(&ListRunsQuery::default()).await.unwrap(); + let summary = store + .list_runs(&ListRunsQuery::default(), Utc::now()) + .await + .unwrap(); assert_eq!(summary.len(), 2); assert_eq!(summary[0].id, test_run_id("run-2")); assert_eq!(summary[1].id, test_run_id("run-1")); @@ -686,7 +695,10 @@ mod tests { store.delete_run(&test_run_id("run-1")).await.unwrap(); assert!(store.open_run(&test_run_id("run-1")).await.is_err()); - let remaining = store.list_runs(&ListRunsQuery::default()).await.unwrap(); + let remaining = store + .list_runs(&ListRunsQuery::default(), Utc::now()) + .await + .unwrap(); assert_eq!(remaining.len(), 1); assert_eq!(remaining[0].id, test_run_id("run-2")); assert!(!list_paths(object_store, "runs/").await.is_empty()); @@ -744,7 +756,10 @@ mod tests { .await .unwrap(); - let summary = store.list_runs(&ListRunsQuery::default()).await.unwrap(); + let summary = store + .list_runs(&ListRunsQuery::default(), Utc::now()) + .await + .unwrap(); assert_eq!(summary.len(), 1); assert_eq!(summary[0].lifecycle.status, RunStatus::Running); assert_eq!( @@ -776,7 +791,7 @@ mod tests { ); assert_eq!( store - .get_cached_summary(&test_run_id("run-3")) + .get_cached_summary(&test_run_id("run-3"), Utc::now()) .await .unwrap() .unwrap() @@ -798,7 +813,7 @@ mod tests { .unwrap(); assert_eq!( store - .get_cached_summary(&test_run_id("run-3")) + .get_cached_summary(&test_run_id("run-3"), Utc::now()) .await .unwrap() .unwrap() @@ -807,20 +822,26 @@ mod tests { ); assert!( store - .list_runs(&ListRunsQuery { - parent_id: Some(test_run_id("run-1")), - ..ListRunsQuery::default() - }) + .list_runs( + &ListRunsQuery { + parent_id: Some(test_run_id("run-1")), + ..ListRunsQuery::default() + }, + Utc::now() + ) .await .unwrap() .is_empty() ); assert_eq!( store - .list_runs(&ListRunsQuery { - parent_id: Some(test_run_id("run-2")), - ..ListRunsQuery::default() - }) + .list_runs( + &ListRunsQuery { + parent_id: Some(test_run_id("run-2")), + ..ListRunsQuery::default() + }, + Utc::now() + ) .await .unwrap() .into_iter() @@ -842,7 +863,7 @@ mod tests { .unwrap(); assert_eq!( store - .get_cached_summary(&test_run_id("run-3")) + .get_cached_summary(&test_run_id("run-3"), Utc::now()) .await .unwrap() .unwrap() @@ -851,10 +872,13 @@ mod tests { ); assert!( store - .list_runs(&ListRunsQuery { - parent_id: Some(test_run_id("run-2")), - ..ListRunsQuery::default() - }) + .list_runs( + &ListRunsQuery { + parent_id: Some(test_run_id("run-2")), + ..ListRunsQuery::default() + }, + Utc::now() + ) .await .unwrap() .is_empty() @@ -878,10 +902,13 @@ mod tests { append_created(&unrelated, "run-3", dt("2026-03-27T12:00:20Z")).await; let summaries = store - .list_runs(&ListRunsQuery { - parent_id: Some(test_run_id("run-1")), - ..ListRunsQuery::default() - }) + .list_runs( + &ListRunsQuery { + parent_id: Some(test_run_id("run-1")), + ..ListRunsQuery::default() + }, + Utc::now(), + ) .await .unwrap(); @@ -914,7 +941,10 @@ mod tests { .await; append_created(&unrelated, "run-4", dt("2026-03-27T12:00:30Z")).await; - let summaries = store.list_runs(&ListRunsQuery::default()).await.unwrap(); + let summaries = store + .list_runs(&ListRunsQuery::default(), Utc::now()) + .await + .unwrap(); let parent_summary = summaries .iter() @@ -935,6 +965,44 @@ mod tests { assert_eq!(unrelated_summary.children_count, 0); } + #[tokio::test] + async fn cached_summary_overlays_live_timing_without_mutating_cached_snapshot() { + let (_object_store, store) = make_store(); + let run = store.create_run(&test_run_id("run-1")).await.unwrap(); + append_created(&run, "run-1", dt("2026-03-27T12:00:00Z")).await; + run.append_event(&event_payload( + "run-1", + "2026-03-27T12:00:01Z", + "run.started", + &serde_json::json!({ "name": "Test run" }), + )) + .await + .unwrap(); + + let now = dt("2026-03-27T12:00:06Z"); + let expected = Some(fabro_types::RunTiming::wall_only(5_000)); + + let summary = store + .get_cached_summary(&test_run_id("run-1"), now) + .await + .unwrap() + .unwrap(); + assert_eq!(summary.timing, expected); + + let listed = store + .list_cached_runs(&ListRunsQuery::default(), now) + .await + .unwrap(); + assert_eq!(listed[0].summary.timing, expected); + + let cached = store + .get_cached_run(&test_run_id("run-1")) + .await + .unwrap() + .unwrap(); + assert_eq!(cached.summary.timing, None); + } + #[tokio::test] async fn control_effect_events_clear_pending_control_and_update_status() { let (_object_store, store) = make_store(); @@ -999,7 +1067,10 @@ mod tests { .await .unwrap(); - let summary = store.list_runs(&ListRunsQuery::default()).await.unwrap(); + let summary = store + .list_runs(&ListRunsQuery::default(), Utc::now()) + .await + .unwrap(); assert_eq!(summary.len(), 1); assert_eq!(summary[0].lifecycle.status, RunStatus::Failed { reason: FailureReason::Cancelled, @@ -1060,7 +1131,10 @@ mod tests { append_completed(&run, "run-1", dt("2026-03-27T12:00:00Z")).await; let reopened = Database::new(object_store, "runs", Duration::from_millis(1), None); - let summary = reopened.list_runs(&ListRunsQuery::default()).await.unwrap(); + let summary = reopened + .list_runs(&ListRunsQuery::default(), Utc::now()) + .await + .unwrap(); assert_eq!(summary.len(), 1); assert_eq!(summary[0].id, test_run_id("run-1")); assert_eq!(summary[0].lifecycle.status, RunStatus::Succeeded { @@ -1080,7 +1154,7 @@ mod tests { reopened.warm_projection_cache().await.unwrap(); let entries = reopened - .list_cached_runs(&ListRunsQuery::default()) + .list_cached_runs(&ListRunsQuery::default(), Utc::now()) .await .unwrap(); assert_eq!( @@ -1092,11 +1166,16 @@ mod tests { assert_eq!(entries[0].last_seq, 3); let filtered = reopened - .list_cached_runs(&ListRunsQuery { - start: Some(test_run_id("run-2").created_at()), - end: Some(test_run_id("run-2").created_at() + chrono::Duration::seconds(1)), - parent_id: None, - }) + .list_cached_runs( + &ListRunsQuery { + start: Some(test_run_id("run-2").created_at()), + end: Some( + test_run_id("run-2").created_at() + chrono::Duration::seconds(1), + ), + parent_id: None, + }, + Utc::now(), + ) .await .unwrap(); assert_eq!(filtered.len(), 1); @@ -1138,7 +1217,7 @@ mod tests { reopened.warm_projection_cache().await.unwrap(); let entries = reopened - .list_cached_runs(&ListRunsQuery::default()) + .list_cached_runs(&ListRunsQuery::default(), Utc::now()) .await .unwrap(); assert_eq!(entries.len(), 1); @@ -1316,9 +1395,12 @@ mod tests { Some("abc123") ); - let summaries = store.list_runs(&ListRunsQuery::default()).await.unwrap(); + let summaries = store + .list_runs(&ListRunsQuery::default(), Utc::now()) + .await + .unwrap(); let projected = store - .list_runs_with_projection(&ListRunsQuery::default()) + .list_runs_with_projection(&ListRunsQuery::default(), Utc::now()) .await .unwrap(); assert_eq!(summaries, vec![cached.summary.clone()]); @@ -1348,7 +1430,7 @@ mod tests { ); assert!( store - .list_cached_runs(&ListRunsQuery::default()) + .list_cached_runs(&ListRunsQuery::default(), Utc::now()) .await .unwrap() .is_empty() diff --git a/lib/crates/fabro-store/src/slate/projection_cache.rs b/lib/crates/fabro-store/src/slate/projection_cache.rs index 32127fcea..aa4974c2f 100644 --- a/lib/crates/fabro-store/src/slate/projection_cache.rs +++ b/lib/crates/fabro-store/src/slate/projection_cache.rs @@ -1,6 +1,7 @@ use std::collections::{BTreeSet, HashMap}; use std::sync::Arc; +use chrono::{DateTime, Utc}; use fabro_types::{Run, RunId, RunProjection}; use tokio::sync::Mutex; @@ -92,6 +93,17 @@ impl RunProjectionCacheState { } } +/// Apply read-time overlays to a cached entry. Pure: does not touch the cache +/// state, so it can run outside the cache mutex. +fn apply_read_overlays(entry: &mut CachedRunProjection, now: DateTime) { + // `Conclusion::timing` is the authoritative terminal snapshot and is + // already present in cached terminal summaries. Only fill missing timing + // with the best-effort live projection. + if entry.summary.timing.is_none() { + entry.summary.timing = entry.projection.live_run_timing(now); + } +} + impl RunProjectionCache { pub(crate) async fn replace_all(&self, entries: Vec) { self.state.lock().await.replace_all(entries); @@ -101,7 +113,11 @@ impl RunProjectionCache { self.state.lock().await.insert(entry); } - pub(crate) async fn list(&self, query: &ListRunsQuery) -> Vec { + pub(crate) async fn list( + &self, + query: &ListRunsQuery, + now: DateTime, + ) -> Vec { let entries = { let state = self.state.lock().await; let raw = match query.parent_id { @@ -131,6 +147,11 @@ impl RunProjectionCache { true }) .collect::>(); + // Apply per-entry live overlays outside the cache mutex, after any + // date filtering so skipped entries do not sum stage timings. + for entry in &mut entries { + apply_read_overlays(entry, now); + } entries.sort_by(|left, right| { right .run_id @@ -150,13 +171,17 @@ impl RunProjectionCache { .map(|entry| state.with_children_count(entry)) } - pub(crate) async fn get_summary(&self, run_id: &RunId) -> Option { - let state = self.state.lock().await; - state.entries.get(run_id).map(|entry| { - let mut summary = entry.summary.clone(); - summary.children_count = state.count_children(run_id); - summary - }) + pub(crate) async fn get_summary(&self, run_id: &RunId, now: DateTime) -> Option { + let mut entry = { + let state = self.state.lock().await; + state + .entries + .get(run_id) + .cloned() + .map(|entry| state.with_children_count(entry))? + }; + apply_read_overlays(&mut entry, now); + Some(entry.summary) } pub(crate) async fn apply_event(&self, run_id: &RunId, event: &EventEnvelope) -> Result<()> { diff --git a/lib/crates/fabro-types/src/run_projection.rs b/lib/crates/fabro-types/src/run_projection.rs index 39d331746..70e1ad0e9 100644 --- a/lib/crates/fabro-types/src/run_projection.rs +++ b/lib/crates/fabro-types/src/run_projection.rs @@ -10,7 +10,7 @@ use crate::{ AgentBackend, AgentMcpToolSummary, AgentSkillActivationSource, AgentSkillSummary, BilledTokenCounts, Checkpoint, Conclusion, InterviewQuestionRecord, InvalidTransition, ModelRef, PullRequestLink, RunControlAction, RunDiff, RunId, RunSandbox, RunSpec, RunStatus, - StageCompletion, StageHandler, StageId, StageState, StageTiming, StartRecord, + RunTiming, StageCompletion, StageHandler, StageId, StageState, StageTiming, StartRecord, TodoListProjection, }; @@ -386,6 +386,33 @@ impl RunProjection { self.archived_at.is_some() } + /// Best-effort run timing for a run that has started but has not reached a + /// terminal conclusion yet. + /// + /// Run-level wall time ticks from `run.started` to `now`. Active time sums + /// inference and tool timing from stages that have already emitted a + /// terminal stage event. Stage projections do not currently track live + /// inference/tool time while a stage is still running, so active time steps + /// forward when each stage completes while wall time advances continuously. + #[must_use] + pub fn live_run_timing(&self, now: DateTime) -> Option { + let start = self.start.as_ref()?; + let wall_time_ms = u64::try_from( + now.signed_duration_since(start.start_time) + .num_milliseconds() + .max(0), + ) + .expect("non-negative milliseconds fit in u64"); + let active = self + .stages + .values() + .filter_map(|stage| stage.timing) + .fold(RunTiming::default(), |acc, timing| { + acc.saturating_add(&RunTiming::from(timing)) + }); + Some(active.with_wall_time(wall_time_ms)) + } + pub fn current_checkpoint(&self) -> Option<&Checkpoint> { self.checkpoints.last().map(|record| &record.checkpoint) } diff --git a/lib/crates/fabro-workflow/src/run_lookup.rs b/lib/crates/fabro-workflow/src/run_lookup.rs index 176bbc144..c84973c7f 100644 --- a/lib/crates/fabro-workflow/src/run_lookup.rs +++ b/lib/crates/fabro-workflow/src/run_lookup.rs @@ -240,7 +240,7 @@ fn scan_orphan_runs(base: &Path) -> Result> { pub async fn scan_runs_combined(store: &Database, base: &Path) -> Result> { let store_runs = store - .list_runs(&fabro_store::ListRunsQuery::default()) + .list_runs(&fabro_store::ListRunsQuery::default(), Utc::now()) .await .unwrap_or_default(); scan_runs_with_summaries(&store_runs, base) diff --git a/lib/crates/fabro-workflow/tests/it/daytona_integration.rs b/lib/crates/fabro-workflow/tests/it/daytona_integration.rs index 8d6e20763..0178a52e9 100644 --- a/lib/crates/fabro-workflow/tests/it/daytona_integration.rs +++ b/lib/crates/fabro-workflow/tests/it/daytona_integration.rs @@ -79,23 +79,27 @@ fn load_run_checkpoint(run_dir: &Path) -> Result Result Result Result