Narrow pull request recovery to unresolved creation requests

The supervisor's periodic recovery scan asked SQLite for every run that
had ever recorded an explicit pull request creation request and then
replayed each inactive candidate's full history to learn whether the
request was still pending. With projections now loaded on demand that
set grows without bound and was replayed every scan.

The candidate query now mirrors the projection reducer: a run is a
candidate only when its latest creation request has no later request,
created, linked, or unlinked event, and no later failure naming the same
creation id. Callers still replay each candidate to confirm, so the query
only has to avoid omitting a pending run, and the replayed set is bounded
by in-flight requests.

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
This commit is contained in:
Scott Werner 2026-09-02 12:48:21 -04:00
parent 13f3ff867b
commit 064074233c

View file

@ -536,13 +536,36 @@ ON CONFLICT(run_id) DO UPDATE SET deleted_at_ms = excluded.deleted_at_ms
decode_run_rows(&rows, now)
}
/// Run ids that have ever recorded an explicit pull request creation
/// request. Callers replay these candidate histories to determine whether
/// their latest request is still pending.
/// Run ids whose latest explicit pull request creation request has no
/// later event that would resolve it. This mirrors the projection reducer:
/// a newer request supersedes the old one; `created`, `linked`, and
/// `unlinked` resolve any pending request; `failed` resolves only the
/// request whose creation id it names. Callers still replay each candidate
/// to confirm, so this must never omit a genuinely pending run, but it
/// keeps the replayed set bounded by in-flight requests rather than by
/// every run that ever asked for a pull request.
pub async fn list_pull_request_creation_candidate_run_ids(&self) -> Result<Vec<RunId>> {
sqlx::query_scalar::<_, String>(
"SELECT DISTINCT run_id FROM run_events \
WHERE event_name = 'pull_request.creation_requested'",
"SELECT DISTINCT requested.run_id FROM run_events AS requested \
WHERE requested.event_name = 'pull_request.creation_requested' \
AND NOT EXISTS ( \
SELECT 1 FROM run_events AS later \
WHERE later.run_id = requested.run_id \
AND later.seq > requested.seq \
AND ( \
later.event_name IN ( \
'pull_request.creation_requested', \
'pull_request.created', \
'pull_request.linked', \
'pull_request.unlinked' \
) \
OR ( \
later.event_name = 'pull_request.failed' \
AND json_extract(later.event_json, '$.properties.creation_id') \
= json_extract(requested.event_json, '$.properties.creation_id') \
) \
) \
)",
)
.fetch_all(&self.pool)
.await?
@ -2589,6 +2612,213 @@ mod tests {
assert_eq!(candidates, expected);
}
/// The candidate query must agree with the projection reducer about which
/// later events resolve a creation request, so recovery replays only runs
/// that are still plausibly pending and never skips one that is.
#[tokio::test]
async fn pull_request_creation_candidates_exclude_requests_resolved_by_later_events() {
let (_directory, store) = store().await;
let created_at = dt("2026-08-27T12:00:00Z");
let base = created_at.timestamp_millis().cast_unsigned();
let link = serde_json::json!({ "owner": "acme", "repo": "widgets", "number": 7 });
let created = serde_json::json!({
"pr_url": "https://github.com/acme/widgets/pull/7",
"pr_number": 7,
"owner": "acme",
"repo": "widgets",
"base_branch": "main",
"head_branch": "fabro/run/7",
"title": "Widgets",
"draft": false,
});
let request = |creation_id: PullRequestCreationId| {
serde_json::json!({
"creation_id": creation_id,
"model": "test-model",
"force": false,
})
};
let resolved_by_created = run_id(base, 1);
let resolved_by_linked = run_id(base + 1, 2);
let resolved_by_unlinked = run_id(base + 2, 3);
let resolved_by_matching_failure = run_id(base + 3, 4);
let failed_for_other_request = run_id(base + 4, 5);
let failed_without_creation_id = run_id(base + 5, 6);
let retried_after_failure = run_id(base + 6, 7);
let superseded_then_resolved = run_id(base + 7, 8);
let all_ids = [
resolved_by_created,
resolved_by_linked,
resolved_by_unlinked,
resolved_by_matching_failure,
failed_for_other_request,
failed_without_creation_id,
retried_after_failure,
superseded_then_resolved,
];
for id in all_ids {
store
.upsert_projection(&entry(projection(id, "candidate", id.created_at()), 9))
.await
.unwrap();
}
let seed = async |id: RunId, seq: u32, event: &str, properties: serde_json::Value| {
let payload = sql_event_payload(&id, event, None, None, None, properties);
seed_sql_event(&store, &id, seq, &payload).await;
};
let id_a = PullRequestCreationId::new();
seed(
resolved_by_created,
2,
"pull_request.creation_requested",
request(id_a),
)
.await;
seed(
resolved_by_created,
3,
"pull_request.created",
created.clone(),
)
.await;
let id_b = PullRequestCreationId::new();
seed(
resolved_by_linked,
2,
"pull_request.creation_requested",
request(id_b),
)
.await;
seed(
resolved_by_linked,
3,
"pull_request.linked",
serde_json::json!({ "pull_request": link }),
)
.await;
let id_c = PullRequestCreationId::new();
seed(
resolved_by_unlinked,
2,
"pull_request.creation_requested",
request(id_c),
)
.await;
seed(
resolved_by_unlinked,
3,
"pull_request.unlinked",
serde_json::json!({ "pull_request": link }),
)
.await;
let id_d = PullRequestCreationId::new();
seed(
resolved_by_matching_failure,
2,
"pull_request.creation_requested",
request(id_d),
)
.await;
seed(
resolved_by_matching_failure,
3,
"pull_request.failed",
serde_json::json!({ "creation_id": id_d, "error": "boom" }),
)
.await;
let id_e = PullRequestCreationId::new();
seed(
failed_for_other_request,
2,
"pull_request.creation_requested",
request(id_e),
)
.await;
seed(
failed_for_other_request,
3,
"pull_request.failed",
serde_json::json!({ "creation_id": PullRequestCreationId::new(), "error": "other" }),
)
.await;
let id_f = PullRequestCreationId::new();
seed(
failed_without_creation_id,
2,
"pull_request.creation_requested",
request(id_f),
)
.await;
seed(
failed_without_creation_id,
3,
"pull_request.failed",
serde_json::json!({ "error": "publish stage failure" }),
)
.await;
let id_g = PullRequestCreationId::new();
seed(
retried_after_failure,
2,
"pull_request.creation_requested",
request(id_g),
)
.await;
seed(
retried_after_failure,
3,
"pull_request.failed",
serde_json::json!({ "creation_id": id_g, "error": "boom" }),
)
.await;
seed(
retried_after_failure,
4,
"pull_request.creation_requested",
request(PullRequestCreationId::new()),
)
.await;
let id_h = PullRequestCreationId::new();
seed(
superseded_then_resolved,
2,
"pull_request.creation_requested",
request(id_h),
)
.await;
seed(
superseded_then_resolved,
3,
"pull_request.creation_requested",
request(PullRequestCreationId::new()),
)
.await;
seed(superseded_then_resolved, 4, "pull_request.created", created).await;
let mut candidates = store
.list_pull_request_creation_candidate_run_ids()
.await
.unwrap();
candidates.sort_unstable();
let mut expected = vec![
failed_for_other_request,
failed_without_creation_id,
retried_after_failure,
];
expected.sort_unstable();
assert_eq!(candidates, expected);
}
#[tokio::test]
async fn projection_persists_billing_diff_and_derived_size() {
let (_directory, store) = store().await;