diff --git a/src/core/transport.ts b/src/core/transport.ts index f2565bf..d590b26 100644 --- a/src/core/transport.ts +++ b/src/core/transport.ts @@ -222,9 +222,12 @@ export class ApiClient { } private async runSessionRefresh(usedToken: string, flight: RefreshFlight): Promise { + // 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", @@ -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 @@ -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; } diff --git a/test/auth_401.test.ts b/test/auth_401.test.ts index 66c2d26..d9134e3 100644 --- a/test/auth_401.test.ts +++ b/test/auth_401.test.ts @@ -299,6 +299,50 @@ 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((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((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"; @@ -306,7 +350,9 @@ test("ApiClient: one cancelled 401 waiter does not abort a shared refresh needed let refreshRequests = 0; let releaseRefresh: (() => void) | undefined; let announceBothOld: (() => void) | undefined; + let announceBodyRead: (() => void) | undefined; const bothOld = new Promise((resolve) => { announceBothOld = resolve; }); + const bodyRead = new Promise((resolve) => { announceBodyRead = resolve; }); const refreshGate = new Promise((resolve) => { releaseRefresh = resolve; }); const store: import("../src/core/auth.js").TokenStore = { async get() { return token; }, @@ -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; @@ -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((resolve) => setImmediate(resolve)); firstController.abort(); await assert.rejects(first, (error: unknown) => (error as Error).name === "AbortError");