mirror of
https://github.com/fabro-sh/fabro.git
synced 2026-10-10 03:30:59 +00:00
Show live timing for in-flight runs via read-time overlay (#361)
## 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<Utc>` 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
<details>
<summary>Ran 9 stages in 38m 28s for $10.35</summary>
| 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** |
</details>
<details>
<summary>Ran <code>ImplementPlan.fabro</code> (12 nodes and 15
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="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
}
```
</details>
⚒️ Generated with [Fabro](https://fabro.sh)
---------
Co-authored-by: Fabro <noreply@fabro.sh>
Co-authored-by: Bryan Helmkamp <bryan@brynary.com>
This commit is contained in:
parent
f81c5b96b1
commit
caba862a33
12 changed files with 389 additions and 116 deletions
|
|
@ -2642,7 +2642,7 @@ pub(crate) async fn reconcile_incomplete_runs_on_startup(
|
|||
) -> anyhow::Result<usize> {
|
||||
let summaries = state
|
||||
.store
|
||||
.list_runs(&fabro_store::ListRunsQuery::default())
|
||||
.list_runs(&fabro_store::ListRunsQuery::default(), chrono::Utc::now())
|
||||
.await?;
|
||||
let mut reconciled = 0usize;
|
||||
|
||||
|
|
|
|||
|
|
@ -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<Arc<AppState>> {
|
|||
}
|
||||
|
||||
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()
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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<Arc<AppState>>,
|
||||
) -> 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<Arc<AppState>>,
|
||||
) -> 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()
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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,
|
||||
|
|
|
|||
|
|
@ -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();
|
||||
|
|
|
|||
|
|
@ -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<Utc> {
|
||||
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();
|
||||
|
|
|
|||
|
|
@ -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<Vec<Run>> {
|
||||
pub async fn list_runs(&self, query: &ListRunsQuery, now: DateTime<Utc>) -> Result<Vec<Run>> {
|
||||
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<Utc>,
|
||||
) -> Result<Vec<(Run, RunProjection)>> {
|
||||
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<Utc>,
|
||||
) -> Result<Vec<CachedRunProjection>> {
|
||||
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<Vec<UnreadableRun>> {
|
||||
|
|
@ -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<Option<Run>> {
|
||||
pub async fn get_cached_summary(
|
||||
&self,
|
||||
run_id: &RunId,
|
||||
now: DateTime<Utc>,
|
||||
) -> Result<Option<Run>> {
|
||||
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<Option<Run>> {
|
||||
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<Vec<Run>> {
|
||||
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()
|
||||
|
|
|
|||
|
|
@ -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<Utc>) {
|
||||
// `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<CachedRunProjection>) {
|
||||
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<CachedRunProjection> {
|
||||
pub(crate) async fn list(
|
||||
&self,
|
||||
query: &ListRunsQuery,
|
||||
now: DateTime<Utc>,
|
||||
) -> Vec<CachedRunProjection> {
|
||||
let entries = {
|
||||
let state = self.state.lock().await;
|
||||
let raw = match query.parent_id {
|
||||
|
|
@ -131,6 +147,11 @@ impl RunProjectionCache {
|
|||
true
|
||||
})
|
||||
.collect::<Vec<_>>();
|
||||
// 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<Run> {
|
||||
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<Utc>) -> Option<Run> {
|
||||
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<()> {
|
||||
|
|
|
|||
|
|
@ -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<Utc>) -> Option<RunTiming> {
|
||||
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)
|
||||
}
|
||||
|
|
|
|||
|
|
@ -240,7 +240,7 @@ fn scan_orphan_runs(base: &Path) -> Result<Vec<RunInfo>> {
|
|||
|
||||
pub async fn scan_runs_combined(store: &Database, base: &Path) -> Result<Vec<RunInfo>> {
|
||||
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)
|
||||
|
|
|
|||
|
|
@ -79,23 +79,27 @@ fn load_run_checkpoint(run_dir: &Path) -> Result<Checkpoint, Box<dyn std::error:
|
|||
let runtime = tokio::runtime::Builder::new_current_thread()
|
||||
.enable_all()
|
||||
.build()?;
|
||||
let run_id = if uses_shared_store {
|
||||
run_dir
|
||||
.file_name()
|
||||
.ok_or("run dir should have file name")?
|
||||
.to_string_lossy()
|
||||
.rsplit('-')
|
||||
.next()
|
||||
.ok_or("run dir should contain run id suffix")?
|
||||
.parse()?
|
||||
} else {
|
||||
runtime
|
||||
.block_on(store.list_runs(&fabro_store::ListRunsQuery::default()))?
|
||||
.into_iter()
|
||||
.next()
|
||||
.ok_or("test store should contain one run")?
|
||||
.id
|
||||
};
|
||||
let run_id =
|
||||
if uses_shared_store {
|
||||
run_dir
|
||||
.file_name()
|
||||
.ok_or("run dir should have file name")?
|
||||
.to_string_lossy()
|
||||
.rsplit('-')
|
||||
.next()
|
||||
.ok_or("run dir should contain run id suffix")?
|
||||
.parse()?
|
||||
} else {
|
||||
runtime
|
||||
.block_on(store.list_runs(
|
||||
&fabro_store::ListRunsQuery::default(),
|
||||
chrono::Utc::now(),
|
||||
))?
|
||||
.into_iter()
|
||||
.next()
|
||||
.ok_or("test store should contain one run")?
|
||||
.id
|
||||
};
|
||||
let run = runtime.block_on(store.open_run_reader(&run_id))?;
|
||||
let state = runtime.block_on(async {
|
||||
for attempt in 0..20 {
|
||||
|
|
@ -128,7 +132,9 @@ fn load_run_checkpoint(run_dir: &Path) -> Result<Checkpoint, Box<dyn std::error:
|
|||
.parse()?
|
||||
} else {
|
||||
runtime
|
||||
.block_on(store.list_runs(&fabro_store::ListRunsQuery::default()))?
|
||||
.block_on(
|
||||
store.list_runs(&fabro_store::ListRunsQuery::default(), chrono::Utc::now()),
|
||||
)?
|
||||
.into_iter()
|
||||
.next()
|
||||
.ok_or("test store should contain one run")?
|
||||
|
|
|
|||
|
|
@ -124,23 +124,27 @@ fn load_run_checkpoint(run_dir: &Path) -> Result<Checkpoint, Box<dyn std::error:
|
|||
let runtime = tokio::runtime::Builder::new_current_thread()
|
||||
.enable_all()
|
||||
.build()?;
|
||||
let run_id = if uses_shared_store {
|
||||
run_dir
|
||||
.file_name()
|
||||
.ok_or("run dir should have file name")?
|
||||
.to_string_lossy()
|
||||
.rsplit('-')
|
||||
.next()
|
||||
.ok_or("run dir should contain run id suffix")?
|
||||
.parse()?
|
||||
} else {
|
||||
runtime
|
||||
.block_on(store.list_runs(&fabro_store::ListRunsQuery::default()))?
|
||||
.into_iter()
|
||||
.next()
|
||||
.ok_or("test store should contain one run")?
|
||||
.id
|
||||
};
|
||||
let run_id =
|
||||
if uses_shared_store {
|
||||
run_dir
|
||||
.file_name()
|
||||
.ok_or("run dir should have file name")?
|
||||
.to_string_lossy()
|
||||
.rsplit('-')
|
||||
.next()
|
||||
.ok_or("run dir should contain run id suffix")?
|
||||
.parse()?
|
||||
} else {
|
||||
runtime
|
||||
.block_on(store.list_runs(
|
||||
&fabro_store::ListRunsQuery::default(),
|
||||
chrono::Utc::now(),
|
||||
))?
|
||||
.into_iter()
|
||||
.next()
|
||||
.ok_or("test store should contain one run")?
|
||||
.id
|
||||
};
|
||||
let run = runtime.block_on(store.open_run_reader(&run_id))?;
|
||||
let state = runtime.block_on(async {
|
||||
for attempt in 0..20 {
|
||||
|
|
@ -173,7 +177,9 @@ fn load_run_checkpoint(run_dir: &Path) -> Result<Checkpoint, Box<dyn std::error:
|
|||
.parse()?
|
||||
} else {
|
||||
runtime
|
||||
.block_on(store.list_runs(&fabro_store::ListRunsQuery::default()))?
|
||||
.block_on(
|
||||
store.list_runs(&fabro_store::ListRunsQuery::default(), chrono::Utc::now()),
|
||||
)?
|
||||
.into_iter()
|
||||
.next()
|
||||
.ok_or("test store should contain one run")?
|
||||
|
|
@ -252,7 +258,9 @@ fn resolve_checkpoint_text(
|
|||
.parse()?
|
||||
} else {
|
||||
runtime
|
||||
.block_on(store.list_runs(&fabro_store::ListRunsQuery::default()))?
|
||||
.block_on(
|
||||
store.list_runs(&fabro_store::ListRunsQuery::default(), chrono::Utc::now()),
|
||||
)?
|
||||
.into_iter()
|
||||
.next()
|
||||
.ok_or("test store should contain one run")?
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue