Repository navigation
Back off checkpoint capture after repeated failures #4517
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Closed
yashranaway
wants to merge
5
commits into
pingdotgg:main
from
yashranaway:checkpoint-capture-backoff
+542
−2
Closed
Changes from 3 commits
Commits
Show all changes
5 commits
Select commit
Hold shift + click to select a range
361ee4a
Back off checkpoint capture after repeated failures
yashranaway a92a481
Read the backoff clock from the Effect clock
yashranaway fa28868
Bound the set of workspaces tracked for capture backoff
yashranaway b2c9ddb
Track recency on skips and reserve the retry after a cooldown
yashranaway 5a76e35
Reserve a capture retry only for a workspace in cooldown
yashranaway File filter
Filter by extension
Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Some comments aren't visible on the classic Files Changed page.
There are no files selected for viewing
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,127 @@ | ||
| import { describe, expect, it } from "vite-plus/test"; | ||
|
|
||
| import { cooldownForFailureCount, makeCaptureBackoff } from "./CaptureBackoff.ts"; | ||
|
|
||
| const MINUTE = 60_000; | ||
| const CWD = "/repo/workspace"; | ||
|
|
||
| describe("cooldownForFailureCount", () => { | ||
| it("tolerates transient failures before opening a cooldown", () => { | ||
| expect(cooldownForFailureCount(1)).toBe(0); | ||
| expect(cooldownForFailureCount(2)).toBe(0); | ||
| }); | ||
|
|
||
| it("backs off further the longer capture keeps failing", () => { | ||
| expect(cooldownForFailureCount(3)).toBe(5 * MINUTE); | ||
| expect(cooldownForFailureCount(4)).toBe(10 * MINUTE); | ||
| expect(cooldownForFailureCount(5)).toBe(20 * MINUTE); | ||
| }); | ||
|
|
||
| it("caps the cooldown so a workspace is retried eventually", () => { | ||
| expect(cooldownForFailureCount(20)).toBe(60 * MINUTE); | ||
| expect(cooldownForFailureCount(500)).toBe(60 * MINUTE); | ||
| }); | ||
| }); | ||
|
|
||
| describe("makeCaptureBackoff", () => { | ||
| it("does not skip before the failure threshold is reached", () => { | ||
| const backoff = makeCaptureBackoff<string>(); | ||
|
|
||
| backoff.recordFailure(CWD, 0, "timeout"); | ||
| backoff.recordFailure(CWD, 1_000, "timeout"); | ||
|
|
||
| expect(backoff.evaluate(CWD, 2_000).skip).toBe(false); | ||
| }); | ||
|
|
||
| it("skips and replays the recorded failure once capture keeps failing", () => { | ||
| const backoff = makeCaptureBackoff<string>(); | ||
|
|
||
| for (const at of [0, 1_000, 2_000]) { | ||
| backoff.recordFailure(CWD, at, "git add timed out"); | ||
| } | ||
|
|
||
| const decision = backoff.evaluate(CWD, 3_000); | ||
| expect(decision.skip).toBe(true); | ||
| expect(decision.lastError).toBe("git add timed out"); | ||
| expect(decision.remainingMs).toBeGreaterThan(0); | ||
| }); | ||
|
|
||
| it("retries again once the cooldown elapses", () => { | ||
| const backoff = makeCaptureBackoff<string>(); | ||
|
|
||
| for (const at of [0, 0, 0]) { | ||
| backoff.recordFailure(CWD, at, "timeout"); | ||
| } | ||
|
|
||
| expect(backoff.evaluate(CWD, 5 * MINUTE - 1).skip).toBe(true); | ||
| expect(backoff.evaluate(CWD, 5 * MINUTE).skip).toBe(false); | ||
| }); | ||
|
|
||
| it("clears the record after a capture succeeds", () => { | ||
| const backoff = makeCaptureBackoff<string>(); | ||
|
|
||
| for (const at of [0, 0, 0]) { | ||
| backoff.recordFailure(CWD, at, "timeout"); | ||
| } | ||
| expect(backoff.evaluate(CWD, 1_000).skip).toBe(true); | ||
|
|
||
| backoff.recordSuccess(CWD); | ||
|
|
||
| expect(backoff.evaluate(CWD, 1_000).skip).toBe(false); | ||
| expect(backoff.trackedWorkspaceCount).toBe(0); | ||
| }); | ||
|
|
||
| it("tracks each workspace independently", () => { | ||
| const backoff = makeCaptureBackoff<string>(); | ||
| const healthy = "/repo/other"; | ||
|
|
||
| for (const at of [0, 0, 0]) { | ||
| backoff.recordFailure(CWD, at, "timeout"); | ||
| } | ||
|
|
||
| expect(backoff.evaluate(CWD, 1_000).skip).toBe(true); | ||
| expect(backoff.evaluate(healthy, 1_000).skip).toBe(false); | ||
| }); | ||
|
|
||
| it("bounds how many workspaces it retains", () => { | ||
| const backoff = makeCaptureBackoff<string>(); | ||
|
|
||
| for (let index = 0; index < 400; index += 1) { | ||
| backoff.recordFailure(`/repo/workspace-${index}`, index, "timeout"); | ||
| } | ||
|
|
||
| expect(backoff.trackedWorkspaceCount).toBe(256); | ||
| // The oldest entries are the ones dropped: an evicted workspace starts | ||
| // counting again from one, a retained one continues. | ||
| expect(backoff.recordFailure("/repo/workspace-0", 1_000, "timeout")).toBe(1); | ||
| expect(backoff.recordFailure("/repo/workspace-399", 1_000, "timeout")).toBe(2); | ||
| }); | ||
|
|
||
| it("keeps a repeatedly failing workspace alive through eviction pressure", () => { | ||
| const backoff = makeCaptureBackoff<string>(); | ||
|
|
||
| // A workspace in active use fails on every turn while many one-off | ||
| // workspaces churn past it. | ||
| for (let index = 0; index < 400; index += 1) { | ||
| backoff.recordFailure(CWD, index, "timeout"); | ||
| backoff.recordFailure(`/repo/other-${index}`, index, "timeout"); | ||
| } | ||
|
|
||
| expect(backoff.evaluate(CWD, 400).skip).toBe(true); | ||
| }); | ||
|
|
||
| it("keeps extending the cooldown while failures continue", () => { | ||
| const backoff = makeCaptureBackoff<string>(); | ||
|
|
||
| for (const at of [0, 0, 0]) { | ||
| backoff.recordFailure(CWD, at, "timeout"); | ||
| } | ||
| const firstRemaining = backoff.evaluate(CWD, 0).remainingMs; | ||
|
|
||
| // The next attempt after the cooldown fails again, so the wait grows. | ||
| backoff.recordFailure(CWD, 5 * MINUTE, "timeout"); | ||
| const secondRemaining = backoff.evaluate(CWD, 5 * MINUTE).remainingMs; | ||
|
|
||
| expect(secondRemaining).toBeGreaterThan(firstRemaining); | ||
| }); | ||
| }); |
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,108 @@ | ||
| /** | ||
| * Consecutive-failure backoff for workspace checkpoint capture. | ||
| * | ||
| * Capture runs a full-tree `git add -A` under a temporary index after every | ||
| * completed turn. On a repository large enough to exceed the VCS process | ||
| * timeout, that capture can never succeed, so the unguarded retry pegs a CPU | ||
| * core for as long as any thread is in use and litters `.git/objects/pack` | ||
| * with `tmp_pack_*` files from each killed process. | ||
| * | ||
| * After a few consecutive failures for a workspace, capture is skipped for a | ||
| * growing cooldown instead of being retried every turn. Any success clears | ||
| * the record, so a transient failure (a lock held by a concurrent git | ||
| * command) costs nothing. | ||
| * | ||
| * @module CaptureBackoff | ||
| */ | ||
|
|
||
| /** Failures tolerated before a workspace enters cooldown. */ | ||
| const FAILURE_THRESHOLD = 3; | ||
| const BASE_COOLDOWN_MS = 5 * 60_000; | ||
| const MAX_COOLDOWN_MS = 60 * 60_000; | ||
|
|
||
| /** | ||
| * Only a workspace that later succeeds clears its own record, so a workspace | ||
| * that fails once and is then abandoned would otherwise be retained for the | ||
| * lifetime of the server. Bound the tracked set and evict least-recently | ||
| * touched entries; dropping one only costs a workspace its failure history. | ||
| */ | ||
| const MAX_TRACKED_WORKSPACES = 256; | ||
|
|
||
| export interface CaptureBackoffDecision<E> { | ||
| readonly skip: boolean; | ||
| /** Milliseconds left in the cooldown, for logging. Zero when not skipping. */ | ||
| readonly remainingMs: number; | ||
| /** | ||
| * The failure that opened the cooldown. Replayed instead of inventing a new | ||
| * error, so callers keep seeing the real reason capture is unavailable and | ||
| * the error channel is unchanged. | ||
| */ | ||
| readonly lastError: E | null; | ||
| } | ||
|
|
||
| export function cooldownForFailureCount(consecutiveFailures: number): number { | ||
| if (consecutiveFailures < FAILURE_THRESHOLD) return 0; | ||
| const doublings = consecutiveFailures - FAILURE_THRESHOLD; | ||
| // Clamp the exponent before shifting so a long-lived workspace cannot | ||
| // overflow into a negative or infinite cooldown. | ||
| const scale = 2 ** Math.min(doublings, 10); | ||
| return Math.min(BASE_COOLDOWN_MS * scale, MAX_COOLDOWN_MS); | ||
| } | ||
|
|
||
| interface WorkspaceRecord<E> { | ||
| consecutiveFailures: number; | ||
| skipUntilMs: number; | ||
| lastError: E; | ||
| } | ||
|
|
||
| /** | ||
| * Tracks capture health per workspace. Callers ask whether to skip, then | ||
| * report the outcome of any capture they actually ran. | ||
| */ | ||
| export function makeCaptureBackoff<E>() { | ||
| const recordByCwd = new Map<string, WorkspaceRecord<E>>(); | ||
|
|
||
| return { | ||
| evaluate(cwd: string, nowMs: number): CaptureBackoffDecision<E> { | ||
| const record = recordByCwd.get(cwd); | ||
| if (!record || nowMs >= record.skipUntilMs) { | ||
|
macroscopeapp[bot] marked this conversation as resolved.
Outdated
|
||
| return { skip: false, remainingMs: 0, lastError: null }; | ||
| } | ||
|
macroscopeapp[bot] marked this conversation as resolved.
Outdated
|
||
| return { | ||
| skip: true, | ||
| remainingMs: record.skipUntilMs - nowMs, | ||
| lastError: record.lastError, | ||
| }; | ||
| }, | ||
|
|
||
| recordSuccess(cwd: string): void { | ||
| recordByCwd.delete(cwd); | ||
| }, | ||
|
|
||
| /** Returns the resulting consecutive failure count for this workspace. */ | ||
| recordFailure(cwd: string, nowMs: number, error: E): number { | ||
| const consecutiveFailures = (recordByCwd.get(cwd)?.consecutiveFailures ?? 0) + 1; | ||
| const cooldownMs = cooldownForFailureCount(consecutiveFailures); | ||
| // Re-insert so Map iteration order tracks recency, keeping the | ||
| // workspaces actually in use ahead of stale ones during eviction. | ||
| recordByCwd.delete(cwd); | ||
| recordByCwd.set(cwd, { | ||
| consecutiveFailures, | ||
| skipUntilMs: cooldownMs === 0 ? 0 : nowMs + cooldownMs, | ||
| lastError: error, | ||
| }); | ||
|
|
||
| while (recordByCwd.size > MAX_TRACKED_WORKSPACES) { | ||
| const oldestCwd = recordByCwd.keys().next().value; | ||
| if (oldestCwd === undefined) break; | ||
| recordByCwd.delete(oldestCwd); | ||
| } | ||
|
|
||
| return consecutiveFailures; | ||
| }, | ||
|
|
||
| get trackedWorkspaceCount(): number { | ||
| return recordByCwd.size; | ||
| }, | ||
| }; | ||
| } | ||
130 changes: 130 additions & 0 deletions
130
apps/server/src/checkpointing/CheckpointCaptureBackoff.test.ts
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,130 @@ | ||
| import { it } from "@effect/vitest"; | ||
| import { ThreadId, VcsProcessTimeoutError } from "@t3tools/contracts"; | ||
| import * as Effect from "effect/Effect"; | ||
| import * as Layer from "effect/Layer"; | ||
| import { describe, expect } from "vite-plus/test"; | ||
|
|
||
| import * as CheckpointStore from "./CheckpointStore.ts"; | ||
| import { checkpointRefForThreadTurn } from "./Utils.ts"; | ||
| import * as VcsDriverRegistry from "../vcs/VcsDriverRegistry.ts"; | ||
| import type * as VcsDriver from "../vcs/VcsDriver.ts"; | ||
|
|
||
| const CWD = "/repo/huge-monorepo"; | ||
| const THREAD_ID = ThreadId.make("thread-checkpoint-backoff"); | ||
|
|
||
| const captureTimeout = new VcsProcessTimeoutError({ | ||
| operation: "GitVcsDriver.checkpoints.captureCheckpoint", | ||
| command: "git add", | ||
| cwd: CWD, | ||
| timeoutMs: 30_000, | ||
| }); | ||
|
|
||
| /** | ||
| * A driver whose capture always exceeds the process timeout, matching a | ||
| * repository too large for a full-tree `git add -A` to finish in time. | ||
| */ | ||
| function makeAlwaysTimingOutRegistry( | ||
| captureAttempts: { count: number }, | ||
| options: { readonly succeedOnAttempt?: number } = {}, | ||
| ) { | ||
| const checkpoints = { | ||
| captureCheckpoint: () => | ||
| Effect.suspend(() => { | ||
| captureAttempts.count += 1; | ||
| return captureAttempts.count === options.succeedOnAttempt | ||
| ? Effect.void | ||
| : Effect.fail(captureTimeout); | ||
| }), | ||
| hasCheckpointRef: () => Effect.succeed(false), | ||
| restoreCheckpoint: () => Effect.succeed(false), | ||
| diffCheckpoints: () => Effect.succeed(""), | ||
| deleteCheckpointRefs: () => Effect.void, | ||
| } as unknown as VcsDriver.VcsCheckpointOps; | ||
|
|
||
| const handle = { | ||
| kind: "git" as const, | ||
| repository: { | ||
| kind: "git" as const, | ||
| rootPath: CWD, | ||
| metadataPath: `${CWD}/.git`, | ||
| freshness: { source: "cache" as const, checkedAt: 0 }, | ||
| }, | ||
| driver: { checkpoints } as unknown as VcsDriver.VcsDriver["Service"], | ||
| } as unknown as VcsDriverRegistry.VcsDriverHandle; | ||
|
|
||
| return Layer.succeed( | ||
| VcsDriverRegistry.VcsDriverRegistry, | ||
| VcsDriverRegistry.VcsDriverRegistry.of({ | ||
| get: () => Effect.succeed(handle.driver), | ||
| detect: () => Effect.succeed(handle), | ||
| resolve: () => Effect.succeed(handle), | ||
| }), | ||
| ); | ||
| } | ||
|
|
||
| describe("checkpoint capture backoff", () => { | ||
| it.effect("stops re-running a capture that keeps timing out", () => { | ||
| const captureAttempts = { count: 0 }; | ||
|
|
||
| return Effect.gen(function* () { | ||
| const store = yield* CheckpointStore.CheckpointStore; | ||
| const capture = (turn: number) => | ||
| store | ||
| .captureCheckpoint({ | ||
| cwd: CWD, | ||
| checkpointRef: checkpointRefForThreadTurn(THREAD_ID, turn), | ||
| }) | ||
| .pipe(Effect.flip); | ||
|
|
||
| // The first three turns each pay for a real capture attempt. | ||
| for (let turn = 1; turn <= 3; turn += 1) { | ||
| const error = yield* capture(turn); | ||
| expect(error._tag).toBe("VcsProcessTimeoutError"); | ||
| } | ||
| expect(captureAttempts.count).toBe(3); | ||
|
|
||
| // Every later turn inside the cooldown replays the recorded failure | ||
| // without spawning git again. The cooldown is minutes long, so the rest | ||
| // of this test runs well inside it. | ||
| for (let turn = 4; turn <= 20; turn += 1) { | ||
| const error = yield* capture(turn); | ||
| expect(error).toBe(captureTimeout); | ||
| } | ||
| expect(captureAttempts.count).toBe(3); | ||
| }).pipe( | ||
| Effect.provide( | ||
| CheckpointStore.layer.pipe(Layer.provide(makeAlwaysTimingOutRegistry(captureAttempts))), | ||
| ), | ||
| ); | ||
| }); | ||
|
|
||
| it.effect("keeps capturing for a workspace that recovers", () => { | ||
| const captureAttempts = { count: 0 }; | ||
|
|
||
| return Effect.gen(function* () { | ||
| const store = yield* CheckpointStore.CheckpointStore; | ||
| const capture = (turn: number) => | ||
| store.captureCheckpoint({ | ||
| cwd: CWD, | ||
| checkpointRef: checkpointRefForThreadTurn(THREAD_ID, turn), | ||
| }); | ||
|
|
||
| // Two failures stay under the threshold, and the success that follows | ||
| // clears them, so a later isolated failure does not open a cooldown. | ||
| yield* capture(1).pipe(Effect.flip); | ||
| yield* capture(2).pipe(Effect.flip); | ||
| yield* capture(3); | ||
| yield* capture(4).pipe(Effect.flip); | ||
| yield* capture(5).pipe(Effect.flip); | ||
| yield* capture(6).pipe(Effect.flip); | ||
|
|
||
| expect(captureAttempts.count).toBe(6); | ||
| }).pipe( | ||
| Effect.provide( | ||
| CheckpointStore.layer.pipe( | ||
| Layer.provide(makeAlwaysTimingOutRegistry(captureAttempts, { succeedOnAttempt: 3 })), | ||
| ), | ||
| ), | ||
| ); | ||
| }); | ||
| }); |
Oops, something went wrong.
Oops, something went wrong.
Add this suggestion to a batch that can be applied as a single commit.
This suggestion is invalid because no changes were made to the code.
Suggestions cannot be applied while the pull request is closed.
Suggestions cannot be applied while viewing a subset of changes.
Only one suggestion per line can be applied in a batch.
Add this suggestion to a batch that can be applied as a single commit.
Applying suggestions on deleted lines is not supported.
You must change the existing code in this line in order to create a valid suggestion.
Outdated suggestions cannot be applied.
This suggestion has been applied or marked resolved.
Suggestions cannot be applied from pending reviews.
Suggestions cannot be applied on multi-line comments.
Suggestions cannot be applied while the pull request is queued to merge.
Suggestion cannot be applied right now. Please check back later.
Uh oh!
There was an error while loading. Please reload this page.