From f2a2630ad069bacd71593d87d037e85f9fb015c5 Mon Sep 17 00:00:00 2001 From: Scott Werner Date: Mon, 31 Aug 2026 14:58:23 -0400 Subject: [PATCH] Simplify pull request recovery and summary store queries - Load recovery candidates from the warm projection cache instead of replaying each run's full event history, and check dispatch eligibility before any I/O - Extract the shared can_dispatch predicate used by both the recovery scan and the worker dispatch loop - Drop load_durable_run_status, now identical to durable_run_status - Share parse_stored_run_id across the three stored-id decode sites - Reuse push_order for the canonical run ordering in list_all and list_by_statuses - Replace the test-only queue clear accessor with the existing drain Co-Authored-By: Claude Fable 5 --- lib/apps/fabro-server/src/server.rs | 12 +---- .../src/server/pull_request_supervisor.rs | 52 ++++++++----------- lib/apps/fabro-server/src/server/tests.rs | 2 +- .../fabro-store/src/run_summary_store.rs | 49 ++++++++--------- 4 files changed, 49 insertions(+), 66 deletions(-) diff --git a/lib/apps/fabro-server/src/server.rs b/lib/apps/fabro-server/src/server.rs index eee8d3ec0..cfc671909 100644 --- a/lib/apps/fabro-server/src/server.rs +++ b/lib/apps/fabro-server/src/server.rs @@ -2638,7 +2638,7 @@ async fn delete_run_internal( }; let had_managed_run = managed_run.is_some(); let durable_status = if managed_run.is_some() { - load_durable_run_status(state, &id).await + durable_run_status(state, id).await.ok().flatten() } else { None }; @@ -2705,16 +2705,6 @@ async fn delete_run_internal( } } -async fn load_durable_run_status(state: &AppState, id: &RunId) -> Option { - state - .stores - .run_summaries - .get(id, Utc::now()) - .await - .ok()? - .map(|summary| summary.lifecycle.status) -} - async fn delete_run_sandbox_resource( state: &AppState, id: RunId, diff --git a/lib/apps/fabro-server/src/server/pull_request_supervisor.rs b/lib/apps/fabro-server/src/server/pull_request_supervisor.rs index d7ac3ec31..60fec9a6e 100644 --- a/lib/apps/fabro-server/src/server/pull_request_supervisor.rs +++ b/lib/apps/fabro-server/src/server/pull_request_supervisor.rs @@ -88,15 +88,6 @@ impl AppState { .pop() } - #[cfg(test)] - pub(super) fn clear_pull_request_creation_queue(&self) { - *self - .pull_request_creation_queue - .lock() - .expect("pull request creation queue lock poisoned") = - PendingPullRequestCreationQueue::default(); - } - #[cfg(test)] pub(super) fn pull_request_creation_queue_len(&self) -> usize { self.pull_request_creation_queue @@ -260,21 +251,18 @@ pub(super) async fn recover_pending_pull_request_creations( .await?; let candidate_count = candidates.len(); let mut pending = Vec::new(); - let mut replay_errors = 0_usize; + let mut load_errors = 0_usize; for run_id in candidates { - let projection = match state.stores.runs.open_run_reader(&run_id).await { - Ok(reader) => match reader.state().await { - Ok(projection) => projection, - Err(error) => { - replay_errors += 1; - warn!(%run_id, %error, "Failed to replay pull request creation candidate"); - continue; - } - }, + if !can_dispatch(&run_id, active, failures) { + continue; + } + let projection = match state.stores.runs.get_cached_projection(&run_id).await { + Ok(Some(projection)) => projection, + Ok(None) => continue, Err(error) => { - replay_errors += 1; - warn!(%run_id, %error, "Failed to open pull request creation candidate"); + load_errors += 1; + warn!(%run_id, %error, "Failed to load pull request creation candidate"); continue; } }; @@ -285,11 +273,6 @@ pub(super) async fn recover_pending_pull_request_creations( else { continue; }; - if active.values().any(|active_id| active_id == &run_id) - || failures.get(&run_id).copied().unwrap_or(0) >= MAX_WORKER_FAILURES_PER_RUN - { - continue; - } pending.push((creation.requested_at, run_id)); } @@ -305,12 +288,23 @@ pub(super) async fn recover_pending_pull_request_creations( candidate_count, pending_count, enqueued, - replay_errors, + load_errors, "Recovered pending pull request creations" ); Ok(()) } +/// Whether the supervisor may hand `run_id` to a worker right now: not +/// already being processed, and not past the store-failure retry cap. +fn can_dispatch( + run_id: &RunId, + active: &HashMap, + failures: &HashMap, +) -> bool { + !active.values().any(|active_id| active_id == run_id) + && failures.get(run_id).copied().unwrap_or(0) < MAX_WORKER_FAILURES_PER_RUN +} + async fn run_pull_request_creation_supervisor(state: Arc) { let shutdown = state.shutdown_token(); let mut workers = JoinSet::new(); @@ -334,9 +328,7 @@ async fn run_pull_request_creation_supervisor(state: Arc) { let Some(run_id) = state.pop_pull_request_creation() else { break; }; - if active.values().any(|active_id| active_id == &run_id) - || failures.get(&run_id).copied().unwrap_or(0) >= MAX_WORKER_FAILURES_PER_RUN - { + if !can_dispatch(&run_id, &active, &failures) { continue; } let handle = workers.spawn( diff --git a/lib/apps/fabro-server/src/server/tests.rs b/lib/apps/fabro-server/src/server/tests.rs index 4faace3fd..25bb99803 100644 --- a/lib/apps/fabro-server/src/server/tests.rs +++ b/lib/apps/fabro-server/src/server/tests.rs @@ -10946,7 +10946,7 @@ async fn pull_request_creation_recovers_durable_request_after_crash_gap() { // Starting the supervisor after the request simulates server recovery: // the durable pending event is enough to resume the operation. - state.clear_pull_request_creation_queue(); + let _ = state.drain_pull_request_creation_queue(); let supervisor = spawn_pull_request_creation_supervisor(Arc::clone(&state)); let creation_body = wait_for_pull_request_creation(&app, run_id).await; diff --git a/lib/components/fabro-store/src/run_summary_store.rs b/lib/components/fabro-store/src/run_summary_store.rs index 30f4890be..18559d5c8 100644 --- a/lib/components/fabro-store/src/run_summary_store.rs +++ b/lib/components/fabro-store/src/run_summary_store.rs @@ -342,14 +342,7 @@ ON CONFLICT(singleton) DO NOTHING .fetch_all(&self.pool) .await? .into_iter() - .map(|stored_id| { - stored_id - .parse::() - .map_err(|_| Error::RunSummaryMismatch { - run_id: stored_id, - field: "id", - }) - }) + .map(parse_stored_run_id) .collect() } @@ -507,7 +500,12 @@ ON CONFLICT(run_id) DO UPDATE SET deleted_at_ms = excluded.deleted_at_ms /// pagination semantics. pub async fn list_all(&self, now: DateTime) -> Result> { let mut query = QueryBuilder::::new(SELECT_RUN_SUMMARIES_SQL); - query.push(" ORDER BY created_at_ms DESC, id DESC"); + push_order( + &mut query, + RunSummarySort::CreatedAt, + RunSummarySortDirection::Desc, + now, + ); let rows = query.build().fetch_all(&self.pool).await?; decode_run_rows(&rows, now) } @@ -527,7 +525,13 @@ ON CONFLICT(run_id) DO UPDATE SET deleted_at_ms = excluded.deleted_at_ms for status in statuses { separated.push_bind(status.to_string()); } - separated.push_unseparated(") ORDER BY created_at_ms DESC, id DESC"); + separated.push_unseparated(")"); + push_order( + &mut query, + RunSummarySort::CreatedAt, + RunSummarySortDirection::Desc, + now, + ); let rows = query.build().fetch_all(&self.pool).await?; decode_run_rows(&rows, now) } @@ -543,14 +547,7 @@ ON CONFLICT(run_id) DO UPDATE SET deleted_at_ms = excluded.deleted_at_ms .fetch_all(&self.pool) .await? .into_iter() - .map(|stored_id| { - stored_id - .parse::() - .map_err(|_| Error::RunSummaryMismatch { - run_id: stored_id, - field: "id", - }) - }) + .map(parse_stored_run_id) .collect() } @@ -569,12 +566,7 @@ FROM runs", rows.iter() .map(|row| { let stored_id: String = row.try_get("id")?; - let id = stored_id - .parse::() - .map_err(|_| Error::RunSummaryMismatch { - run_id: stored_id, - field: "id", - })?; + let id = parse_stored_run_id(stored_id)?; Ok(RunSummaryIdentity { id, workflow_slug: row.try_get("workflow_slug")?, @@ -1462,6 +1454,15 @@ fn push_order( builder.push(", id DESC"); } +fn parse_stored_run_id(stored_id: String) -> Result { + stored_id + .parse::() + .map_err(|_| Error::RunSummaryMismatch { + run_id: stored_id, + field: "id", + }) +} + fn decode_run_row(row: &SqliteRow, now: DateTime) -> Result { let stored_id: String = row.try_get("id")?; let summary_json: String = row.try_get("summary_json")?;