Skip to content

Commit d2e61b1

Browse files
committed
fix: harden supervisor reset edge cases
1 parent 63ad12a commit d2e61b1

4 files changed

Lines changed: 189 additions & 27 deletions

File tree

packages/server/src/__tests__/supervisor-manager.test.ts

Lines changed: 50 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -798,6 +798,56 @@ describe("SupervisorManager cycle triggers", () => {
798798
expect(deps.supervisorRepo.findById(supervisor.id)?.objective).toBe("Initial objective");
799799
});
800800

801+
it("preserves a pause request while an objective change abort is in flight", async () => {
802+
const supervisor = await manager.create({
803+
sessionId: "sess-objective-pause-race",
804+
workspaceId: "ws-1",
805+
objective: "Initial objective",
806+
evaluatorProviderId: "codex",
807+
});
808+
809+
let observedSignal: AbortSignal | undefined;
810+
vi.spyOn(getManagerInternals().evaluator, "evaluate").mockImplementationOnce(
811+
async (_supervisor, _context, options) =>
812+
await new Promise<SupervisorEvaluationResult>((_resolve, reject) => {
813+
observedSignal = options?.signal;
814+
options?.signal?.addEventListener(
815+
"abort",
816+
() => {
817+
reject({
818+
code: "supervisor_eval_aborted",
819+
message: "Supervisor evaluator aborted",
820+
});
821+
},
822+
{ once: true }
823+
);
824+
})
825+
);
826+
827+
await manager.triggerEvaluation(supervisor.id);
828+
829+
await waitFor(() => {
830+
expect(observedSignal).toBeDefined();
831+
expect(manager.get(supervisor.id)?.state).toBe("evaluating");
832+
});
833+
834+
const updatedPromise = manager.update(supervisor.id, {
835+
objective: "New objective",
836+
});
837+
const paused = await manager.pause(supervisor.id);
838+
839+
await waitFor(() => {
840+
expect(observedSignal?.aborted).toBe(true);
841+
});
842+
843+
const updated = await updatedPromise;
844+
845+
expect(paused.state).toBe("paused");
846+
expect(updated.state).toBe("paused");
847+
expect(updated.objective).toBe("New objective");
848+
expect(manager.get(supervisor.id)?.state).toBe("paused");
849+
});
850+
801851
it("retries evaluator timeout up to the global retry budget", async () => {
802852
vi.useFakeTimers();
803853
deps.settingsRepo.get = vi.fn((key: string) => {

packages/server/src/supervisor/manager.ts

Lines changed: 4 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -824,9 +824,12 @@ export class SupervisorManager {
824824
}
825825

826826
if (this.pendingObjectiveUpdates.has(supervisorId)) {
827+
const nextState: SupervisorState = this.pendingPauses.has(supervisorId)
828+
? "paused"
829+
: "idle";
827830
const recoveredSupervisor = this.attachCycles(
828831
this.deps.supervisorRepo.update(supervisorId, {
829-
state: "idle",
832+
state: nextState,
830833
stopReason: null,
831834
errorReason: null,
832835
updatedAt: Date.now(),
@@ -838,7 +841,6 @@ export class SupervisorManager {
838841
this.broadcastState(recoveredSupervisor, "state_changed");
839842
this.deps.cycleRepo.pruneOldest(supervisorId, this.config.maxCyclesPerSession);
840843
this.scheduler.refresh();
841-
this.pendingPauses.delete(supervisorId);
842844

843845
return abortedCycle;
844846
}

packages/server/src/supervisor/target-store.atomic.test.ts

Lines changed: 96 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -1,24 +1,33 @@
1-
import { mkdtempSync, rmSync } from "node:fs";
1+
import { existsSync, mkdtempSync, readdirSync, readFileSync, rmSync } from "node:fs";
22
import { tmpdir } from "node:os";
33
import { join } from "node:path";
44
import { afterEach, beforeEach, describe, expect, it, vi } from "vitest";
55

6-
const renameState = vi.hoisted(() => ({
7-
callCount: 0,
8-
failOnCall: 0,
6+
const fsState = vi.hoisted(() => ({
7+
renameCallCount: 0,
8+
failRenameOnCalls: [] as number[],
9+
rmCallCount: 0,
10+
failRmOnCalls: [] as number[],
911
}));
1012

1113
vi.mock("node:fs/promises", async () => {
1214
const actual = await vi.importActual<typeof import("node:fs/promises")>("node:fs/promises");
1315
return {
1416
...actual,
1517
rename: vi.fn(async (from: string, to: string) => {
16-
renameState.callCount += 1;
17-
if (renameState.failOnCall !== 0 && renameState.callCount === renameState.failOnCall) {
18+
fsState.renameCallCount += 1;
19+
if (fsState.failRenameOnCalls.includes(fsState.renameCallCount)) {
1820
throw new Error("promote failed");
1921
}
2022
return actual.rename(from, to);
2123
}),
24+
rm: vi.fn(async (...args: Parameters<typeof actual.rm>) => {
25+
fsState.rmCallCount += 1;
26+
if (fsState.failRmOnCalls.includes(fsState.rmCallCount)) {
27+
throw new Error("cleanup failed");
28+
}
29+
return actual.rm(...args);
30+
}),
2231
};
2332
});
2433

@@ -37,8 +46,10 @@ describe("target store atomic reset", () => {
3746

3847
beforeEach(() => {
3948
workspacePath = mkdtempSync(join(tmpdir(), "supervisor-target-store-atomic-"));
40-
renameState.callCount = 0;
41-
renameState.failOnCall = 0;
49+
fsState.renameCallCount = 0;
50+
fsState.failRenameOnCalls = [];
51+
fsState.rmCallCount = 0;
52+
fsState.failRmOnCalls = [];
4253
});
4354

4455
afterEach(() => {
@@ -77,7 +88,7 @@ describe("target store atomic reset", () => {
7788
attemptCount: 1,
7889
});
7990

80-
renameState.failOnCall = 2;
91+
fsState.failRenameOnCalls = [2];
8192

8293
await expect(
8394
resetTargetFiles(workspacePath, {
@@ -110,4 +121,80 @@ describe("target store atomic reset", () => {
110121
},
111122
]);
112123
});
124+
125+
it("keeps the promoted target live when backup cleanup fails", async () => {
126+
await createTargetFiles(workspacePath, {
127+
targetId: "tgt-1",
128+
sessionId: "sess-1",
129+
workspaceId: "ws-1",
130+
objective: "Old objective",
131+
createdAt: 1,
132+
});
133+
134+
fsState.failRmOnCalls = [1];
135+
136+
await expect(
137+
resetTargetFiles(workspacePath, {
138+
targetId: "tgt-1",
139+
sessionId: "sess-1",
140+
workspaceId: "ws-1",
141+
objective: "New objective",
142+
createdAt: 3,
143+
})
144+
).resolves.toBeUndefined();
145+
146+
const meta = await readTargetMeta(workspacePath, "tgt-1");
147+
const memory = await loadTargetMemory(workspacePath, "tgt-1");
148+
const cycles = await readTargetCycleRecords(workspacePath, "tgt-1");
149+
150+
expect(meta.objective).toBe("New objective");
151+
expect(memory).toMatchObject({
152+
targetId: "tgt-1",
153+
planGenerated: false,
154+
stalledCount: 0,
155+
updatedAt: 3,
156+
});
157+
expect(cycles).toEqual([]);
158+
});
159+
160+
it("preserves the backup target when both promotion and restore fail", async () => {
161+
await createTargetFiles(workspacePath, {
162+
targetId: "tgt-1",
163+
sessionId: "sess-1",
164+
workspaceId: "ws-1",
165+
objective: "Old objective",
166+
createdAt: 1,
167+
});
168+
169+
fsState.failRenameOnCalls = [2, 3];
170+
171+
await expect(
172+
resetTargetFiles(workspacePath, {
173+
targetId: "tgt-1",
174+
sessionId: "sess-1",
175+
workspaceId: "ws-1",
176+
objective: "New objective",
177+
createdAt: 3,
178+
})
179+
).rejects.toThrow("promote failed");
180+
181+
const targetsRoot = join(workspacePath, ".coder-studio", "supervisor", "targets");
182+
const entries = readdirSync(targetsRoot);
183+
const backupEntry = entries.find((entry) => entry.startsWith("tgt-1.backup-"));
184+
const stagingEntry = entries.find((entry) => entry.startsWith("tgt-1.reset-"));
185+
186+
expect(existsSync(join(targetsRoot, "tgt-1"))).toBe(false);
187+
expect(backupEntry).toBeDefined();
188+
expect(stagingEntry).toBeDefined();
189+
190+
const backupMeta = JSON.parse(
191+
readFileSync(join(targetsRoot, backupEntry!, "meta.json"), "utf-8")
192+
) as { objective: string };
193+
const stagedMeta = JSON.parse(
194+
readFileSync(join(targetsRoot, stagingEntry!, "meta.json"), "utf-8")
195+
) as { objective: string };
196+
197+
expect(backupMeta.objective).toBe("Old objective");
198+
expect(stagedMeta.objective).toBe("New objective");
199+
});
113200
});

packages/server/src/supervisor/target-store.ts

Lines changed: 39 additions & 16 deletions
Original file line numberDiff line numberDiff line change
@@ -51,6 +51,19 @@ function hasCode(error: unknown, code: string): boolean {
5151
);
5252
}
5353

54+
function errorMessage(error: unknown, fallback: string): string {
55+
if (error instanceof Error) {
56+
return error.message;
57+
}
58+
if (error && typeof error === "object" && "message" in error) {
59+
const value = (error as { message?: unknown }).message;
60+
if (typeof value === "string") {
61+
return value;
62+
}
63+
}
64+
return fallback;
65+
}
66+
5467
async function writeJsonIfMissing(path: string, value: unknown): Promise<void> {
5568
try {
5669
await writeFile(path, JSON.stringify(value, null, 2) + "\n", {
@@ -154,15 +167,16 @@ export async function resetTargetFiles(
154167
await mkdir(parentDir, { recursive: true });
155168
const stagingDir = await mkdtemp(join(parentDir, `${input.targetId}.reset-`));
156169

157-
let movedExisting = false;
170+
let backupCreated = false;
158171
let promoted = false;
172+
let restored = false;
159173

160174
try {
161175
await writeResetTargetFiles(stagingDir, input);
162176

163177
try {
164178
await rename(dir, backupDir);
165-
movedExisting = true;
179+
backupCreated = true;
166180
} catch (error) {
167181
if (!hasCode(error, "ENOENT")) {
168182
throw error;
@@ -172,27 +186,36 @@ export async function resetTargetFiles(
172186
try {
173187
await rename(stagingDir, dir);
174188
promoted = true;
175-
} catch (error) {
176-
if (movedExisting) {
177-
await rename(backupDir, dir);
178-
movedExisting = false;
189+
} catch (promoteError) {
190+
if (backupCreated) {
191+
try {
192+
await rename(backupDir, dir);
193+
backupCreated = false;
194+
restored = true;
195+
} catch (restoreError) {
196+
throw new Error(
197+
`Failed to promote target reset (${errorMessage(
198+
promoteError,
199+
"unknown promote error"
200+
)}); restore also failed (${errorMessage(restoreError, "unknown restore error")})`
201+
);
202+
}
179203
}
180-
throw error;
181-
}
182-
183-
if (movedExisting) {
184-
await rm(backupDir, { recursive: true, force: true });
185-
movedExisting = false;
204+
throw promoteError;
186205
}
187206
} catch (error) {
188-
if (!promoted) {
207+
if (restored || !backupCreated) {
189208
await rm(stagingDir, { recursive: true, force: true }).catch(() => {});
190209
}
191-
if (movedExisting) {
192-
await rm(backupDir, { recursive: true, force: true }).catch(() => {});
193-
}
194210
throw error;
195211
}
212+
213+
if (backupCreated) {
214+
await rm(backupDir, { recursive: true, force: true }).catch(() => {});
215+
}
216+
if (!promoted) {
217+
await rm(stagingDir, { recursive: true, force: true }).catch(() => {});
218+
}
196219
}
197220

198221
export async function readTargetMeta(

0 commit comments

Comments
 (0)