Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
65 changes: 42 additions & 23 deletions src/orb/apr-repo-transfer.ts
Original file line number Diff line number Diff line change
Expand Up @@ -210,16 +210,22 @@ export async function probeAprRepoTransfer(
env: Env,
transfer: Pick<PendingAprRepoTransfer, "repoFullName" | "newOwner" | "installationId">,
): Promise<AprRepoTransferProbe> {
const token = await createInstallationToken(env, transfer.installationId);
const response = await timeoutFetch(`https://api.github.com/repos/${transfer.repoFullName}`, {
headers: githubHeaders({ token }),
});
if (response.status === 404) return { state: "access_departed" };
if (!response.ok) return { state: "pending" };
const body = (await response.json().catch(() => null)) as { owner?: { login?: string } } | null;
const owner = body?.owner?.login;
if (owner && owner.toLowerCase() === transfer.newOwner.toLowerCase()) return { state: "resolved_under_target" };
return { state: "pending" };
// #8331: token mint / network / AbortSignal failures must not escape — treat them as still-pending
// so the next poll retries, matching this function's documented "Never throws" contract.
try {
const token = await createInstallationToken(env, transfer.installationId);
const response = await timeoutFetch(`https://api.github.com/repos/${transfer.repoFullName}`, {
headers: githubHeaders({ token }),
});
if (response.status === 404) return { state: "access_departed" };
if (!response.ok) return { state: "pending" };
const body = (await response.json().catch(() => null)) as { owner?: { login?: string } } | null;
const owner = body?.owner?.login;
if (owner && owner.toLowerCase() === transfer.newOwner.toLowerCase()) return { state: "resolved_under_target" };
return { state: "pending" };
} catch {
return { state: "pending" };
}
}

/**
Expand Down Expand Up @@ -284,20 +290,33 @@ export async function pollPendingAprRepoTransfers(
const now = deps.now();
const results: AprRepoTransferPollResult[] = [];
for (const transfer of pending) {
const probe = await deps.probe(env, transfer);
const outcome = classifyAprRepoTransferOutcome({
probe,
initiatedAt: transfer.initiatedAt,
now,
...(deps.expiryMs !== undefined ? { expiryMs: deps.expiryMs } : {}),
});
if (outcome === "pending") {
await deps.setDispatchPaused(env, transfer.repoFullName, true);
} else {
await deps.markResolved(env, transfer, outcome);
if (outcome !== "accepted_departed") await deps.setDispatchPaused(env, transfer.repoFullName, false);
// #8331: isolate each transfer so one probe/dependency throw cannot abort the rest of the batch
// (mirrors retryFailedRelays' per-row independence — a bad row must never starve siblings).
try {
const probe = await deps.probe(env, transfer);
const outcome = classifyAprRepoTransferOutcome({
probe,
initiatedAt: transfer.initiatedAt,
now,
...(deps.expiryMs !== undefined ? { expiryMs: deps.expiryMs } : {}),
});
if (outcome === "pending") {
await deps.setDispatchPaused(env, transfer.repoFullName, true);
} else {
await deps.markResolved(env, transfer, outcome);
if (outcome !== "accepted_departed") await deps.setDispatchPaused(env, transfer.repoFullName, false);
}
results.push({ repoFullName: transfer.repoFullName, outcome });
} catch (error) {
console.error(
JSON.stringify({
level: "error",
event: "apr_repo_transfer_poll_item_failed",
repoFullName: transfer.repoFullName,
message: String(error).slice(0, 200),
}),
);
}
results.push({ repoFullName: transfer.repoFullName, outcome });
}
return results;
}
78 changes: 78 additions & 0 deletions test/unit/orb-apr-repo-transfer.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -336,6 +336,19 @@ describe("probeAprRepoTransfer (#7741)", () => {
stubFetch(() => new Response("", { status: 500 }));
expect(await probeAprRepoTransfer(createTestEnv(), transfer)).toEqual({ state: "pending" });
});

// #8331: "Never throws" must hold for token-mint and fetch rejections, not only HTTP status codes.
it("stays pending (does not throw) when createInstallationToken rejects", async () => {
mockedProbeToken.mockRejectedValue(new Error("token mint failed"));
await expect(probeAprRepoTransfer(createTestEnv(), transfer)).resolves.toEqual({ state: "pending" });
});

it("stays pending (does not throw) when the outbound fetch rejects", async () => {
vi.stubGlobal("fetch", async () => {
throw new Error("network down");
});
await expect(probeAprRepoTransfer(createTestEnv(), transfer)).resolves.toEqual({ state: "pending" });
});
});

describe("setAprRepoDispatchPaused (#7741 deliverable 2)", () => {
Expand Down Expand Up @@ -439,4 +452,69 @@ describe("pollPendingAprRepoTransfers (#7741 deliverables 1+2)", () => {
expect(resolved).toEqual([{ repoFullName: "loopover-repos/stale", outcome: "expired" }]);
expect(paused).toEqual([{ repoFullName: "loopover-repos/stale", paused: false }]);
});

// #8331: one rejecting probe must not starve the rest of the batch.
it("continues processing later transfers when one probe rejects", async () => {
const pending: PendingAprRepoTransfer[] = [
{ ...base, repoFullName: "loopover-repos/broken", initiatedAt: now - 1000 },
{ ...base, repoFullName: "loopover-repos/ok", initiatedAt: now - 1000 },
];
const paused: Array<{ repoFullName: string; paused: boolean }> = [];
const resolved: Array<{ repoFullName: string; outcome: string }> = [];
const errorSpy = vi.spyOn(console, "error").mockImplementation(() => undefined);
const deps: AprRepoTransferPollDeps = {
listPending: async () => pending,
probe: async (_env, t) => {
if (t.repoFullName === "loopover-repos/broken") throw new Error("probe boom");
return { state: "resolved_under_target" };
},
now: () => now,
markResolved: async (_env, t, outcome) => {
resolved.push({ repoFullName: t.repoFullName, outcome });
},
setDispatchPaused: async (_env, repoFullName, p) => {
paused.push({ repoFullName, paused: p });
},
expiryMs: 5000,
};

const results = await pollPendingAprRepoTransfers(createTestEnv(), deps);

expect(results).toEqual([{ repoFullName: "loopover-repos/ok", outcome: "accepted" }]);
expect(resolved).toEqual([{ repoFullName: "loopover-repos/ok", outcome: "accepted" }]);
expect(paused).toEqual([{ repoFullName: "loopover-repos/ok", paused: false }]);
expect(errorSpy).toHaveBeenCalled();
const logged = String(errorSpy.mock.calls[0]?.[0] ?? "");
expect(logged).toContain("apr_repo_transfer_poll_item_failed");
expect(logged).toContain("loopover-repos/broken");
errorSpy.mockRestore();
});

it("continues when a later dependency call (markResolved) rejects after a successful probe", async () => {
const pending: PendingAprRepoTransfer[] = [
{ ...base, repoFullName: "loopover-repos/write-fail", initiatedAt: now - 1000 },
{ ...base, repoFullName: "loopover-repos/ok", initiatedAt: now - 1000 },
];
const paused: Array<{ repoFullName: string; paused: boolean }> = [];
const errorSpy = vi.spyOn(console, "error").mockImplementation(() => undefined);
const deps: AprRepoTransferPollDeps = {
listPending: async () => pending,
probe: async () => ({ state: "resolved_under_target" }),
now: () => now,
markResolved: async (_env, t) => {
if (t.repoFullName === "loopover-repos/write-fail") throw new Error("persist failed");
},
setDispatchPaused: async (_env, repoFullName, p) => {
paused.push({ repoFullName, paused: p });
},
expiryMs: 5000,
};

const results = await pollPendingAprRepoTransfers(createTestEnv(), deps);

expect(results).toEqual([{ repoFullName: "loopover-repos/ok", outcome: "accepted" }]);
expect(paused).toEqual([{ repoFullName: "loopover-repos/ok", paused: false }]);
expect(errorSpy).toHaveBeenCalled();
errorSpy.mockRestore();
});
});