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
19 changes: 11 additions & 8 deletions src/core/transport.ts
Original file line number Diff line number Diff line change
Expand Up @@ -222,9 +222,12 @@ export class ApiClient {
}

private async runSessionRefresh(usedToken: string, flight: RefreshFlight): Promise<boolean> {
// Also race the header and body promises against this deadline. A fetch
// implementation can return 200 headers while its JSON body never ends,
// even after the passed signal aborts; that must not pin the shared flight.
const signal = AbortSignal.any([flight.controller.signal, AbortSignal.timeout(10_000)]);
try {
const res = await fetch(this.url(REFRESH_PATH), {
const res = await raceAgainst(fetch(this.url(REFRESH_PATH), {
method: "POST",
headers: {
"Content-Type": "application/json",
Expand All @@ -233,21 +236,21 @@ export class ApiClient {
},
body: "{}",
signal,
});
if (flight.controller.signal.aborted || !res.ok) return false;
}), signal, 0);
if (signal.aborted || !res.ok) return false;
let body: { session_token?: string } | undefined;
try {
body = (await res.json()) as { session_token?: string };
body = (await raceAgainst(res.json(), signal, 0)) as { session_token?: string };
} catch {
return false;
}
if (flight.controller.signal.aborted || !body?.session_token) return false;
if (signal.aborted || !body?.session_token) return false;

// Do not overwrite a rotation performed by another process while this
// network request was in flight. A matching fresh value already means
// callers may retry; any other value belongs to a newer authority.
const current = await this.tokens.get();
if (flight.controller.signal.aborted) return false;
const current = await raceAgainst(this.tokens.get(), signal, 0);
if (signal.aborted) return false;
if (current !== usedToken) return current === body.session_token;

// update() (when the store distinguishes it) swaps the ACTIVE token
Expand All @@ -258,7 +261,7 @@ export class ApiClient {
// never return to the caller before a late credential mutation occurs.
flight.committing = true;
await (this.tokens.update?.(body.session_token) ?? this.tokens.set(body.session_token));
return !flight.controller.signal.aborted;
return !signal.aborted;
} catch {
return false;
}
Expand Down
54 changes: 52 additions & 2 deletions test/auth_401.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -299,14 +299,60 @@ test("ApiClient: cancelling the sole 401 waiter aborts refresh and prevents a la
}
});

test("ApiClient: a stalled 200 refresh body releases a cancelled flight so a later 401 can retry", async () => {
const real = globalThis.fetch;
const store = new StaticTokenStore("sess_old");
let refreshRequests = 0;
let announceBodyRead: (() => void) | undefined;
const bodyRead = new Promise<void>((resolve) => { announceBodyRead = resolve; });
let releaseBody: ((value: { session_token: string }) => void) | undefined;
const bodyGate = new Promise<{ session_token: string }>((resolve) => { releaseBody = resolve; });
globalThis.fetch = (async (url: unknown) => {
if (String(url).endsWith("/auth/refresh")) {
refreshRequests += 1;
if (refreshRequests > 1) return jsonRes(401, { detail: "refresh unavailable" });
const stalled = jsonRes(200, { session_token: "ignored" });
// Model a fetch implementation that returned 200 headers but leaves
// res.json() pending even after the request's signal aborts.
Object.defineProperty(stalled, "json", {
value: () => { announceBodyRead?.(); return bodyGate; },
});
return stalled;
}
return jsonRes(401, { detail: "expired" });
}) as typeof globalThis.fetch;
try {
const api = new ApiClient("https://api.example", store);
const firstController = new AbortController();
const first = api.getJson("/models", firstController.signal, 1_000);
await bodyRead;
firstController.abort();
await assert.rejects(first, (error: unknown) => (error as Error).name === "AbortError");
await new Promise<void>((resolve) => setImmediate(resolve));

await assert.rejects(
() => api.getJson("/models", AbortSignal.timeout(500), 1_000),
(error: unknown) => error instanceof HttpError && error.status === 401,
"a later 401 must start a new refresh instead of joining the dead body reader",
);
assert.equal(refreshRequests, 2);
assert.equal(await store.get(), "sess_old", "aborted refresh cannot rotate the token");
} finally {
releaseBody?.({ session_token: "sess_late" });
globalThis.fetch = real;
}
});

test("ApiClient: one cancelled 401 waiter does not abort a shared refresh needed by another request", async () => {
const real = globalThis.fetch;
let token = "sess_old";
let oldRequests = 0;
let refreshRequests = 0;
let releaseRefresh: (() => void) | undefined;
let announceBothOld: (() => void) | undefined;
let announceBodyRead: (() => void) | undefined;
const bothOld = new Promise<void>((resolve) => { announceBothOld = resolve; });
const bodyRead = new Promise<void>((resolve) => { announceBodyRead = resolve; });
const refreshGate = new Promise<void>((resolve) => { releaseRefresh = resolve; });
const store: import("../src/core/auth.js").TokenStore = {
async get() { return token; },
Expand All @@ -318,8 +364,11 @@ test("ApiClient: one cancelled 401 waiter does not abort a shared refresh needed
const target = String(url);
if (target.endsWith("/auth/refresh")) {
refreshRequests += 1;
await refreshGate;
return jsonRes(200, { session_token: "sess_new" });
const response = jsonRes(200, { session_token: "ignored" });
Object.defineProperty(response, "json", {
value: () => { announceBodyRead?.(); return refreshGate.then(() => ({ session_token: "sess_new" })); },
});
return response;
}
if (bearer(init ?? {}) === "Bearer sess_new") return jsonRes(200, { ok: true });
oldRequests += 1;
Expand All @@ -332,6 +381,7 @@ test("ApiClient: one cancelled 401 waiter does not abort a shared refresh needed
const first = api.getJson("/models", firstController.signal, 1_000);
const second = api.getJson<{ ok: boolean }>("/models", undefined, 1_000);
await bothOld;
await bodyRead;
await new Promise<void>((resolve) => setImmediate(resolve));
firstController.abort();
await assert.rejects(first, (error: unknown) => (error as Error).name === "AbortError");
Expand Down
Loading