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 <noreply@anthropic.com>
This commit is contained in:
Scott Werner 2026-08-31 14:58:23 -04:00
parent 1a2d8a9056
commit f2a2630ad0
4 changed files with 49 additions and 66 deletions

View file

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

View file

@ -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<task::Id, RunId>,
failures: &HashMap<RunId, u32>,
) -> 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<AppState>) {
let shutdown = state.shutdown_token();
let mut workers = JoinSet::new();
@ -334,9 +328,7 @@ async fn run_pull_request_creation_supervisor(state: Arc<AppState>) {
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(

View file

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

View file

@ -342,14 +342,7 @@ ON CONFLICT(singleton) DO NOTHING
.fetch_all(&self.pool)
.await?
.into_iter()
.map(|stored_id| {
stored_id
.parse::<RunId>()
.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<Utc>) -> Result<Vec<Run>> {
let mut query = QueryBuilder::<Sqlite>::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::<RunId>()
.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::<RunId>()
.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<RunId> {
stored_id
.parse::<RunId>()
.map_err(|_| Error::RunSummaryMismatch {
run_id: stored_id,
field: "id",
})
}
fn decode_run_row(row: &SqliteRow, now: DateTime<Utc>) -> Result<Run> {
let stored_id: String = row.try_get("id")?;
let summary_json: String = row.try_get("summary_json")?;