mirror of
https://github.com/toeverything/AFFiNE.git
synced 2026-07-24 13:58:50 +08:00
feat(server): support refresh token (#15218)
This commit is contained in:
@@ -0,0 +1 @@
|
||||
export * from './token-broker';
|
||||
@@ -0,0 +1,409 @@
|
||||
import { describe, expect, test, vi } from 'vitest';
|
||||
|
||||
import {
|
||||
AuthSessionError,
|
||||
AuthTokenBroker,
|
||||
type AuthTokenPair,
|
||||
type AuthTokenResponse,
|
||||
type AuthTokenStorage,
|
||||
createRealtimeAuthAdapter,
|
||||
createRequestAuthAdapter,
|
||||
withAuthRetry,
|
||||
} from './token-broker';
|
||||
|
||||
const now = Date.parse('2026-07-11T00:00:00.000Z');
|
||||
|
||||
function deferred<T>() {
|
||||
let resolve!: (value: T) => void;
|
||||
let reject!: (error: unknown) => void;
|
||||
const promise = new Promise<T>((done, fail) => {
|
||||
resolve = done;
|
||||
reject = fail;
|
||||
});
|
||||
return { promise, reject, resolve };
|
||||
}
|
||||
|
||||
function response(sequence: number): AuthTokenResponse {
|
||||
return {
|
||||
tokenType: 'Bearer',
|
||||
accessToken: `access-${sequence}`,
|
||||
expiresIn: 900,
|
||||
refreshToken: `refresh-${sequence}`,
|
||||
refreshExpiresAt: '2026-08-11T00:00:00.000Z',
|
||||
session: {
|
||||
id: 'session-1',
|
||||
absoluteExpiresAt: '2027-01-11T00:00:00.000Z',
|
||||
},
|
||||
};
|
||||
}
|
||||
|
||||
function pair(sequence: number, accessExpiresAt: number): AuthTokenPair {
|
||||
return {
|
||||
version: 1,
|
||||
...response(sequence),
|
||||
accessExpiresAt: new Date(accessExpiresAt).toISOString(),
|
||||
};
|
||||
}
|
||||
|
||||
function setup(initial: AuthTokenPair | null) {
|
||||
let persisted = initial;
|
||||
const storage: AuthTokenStorage = {
|
||||
load: vi.fn(async () => persisted),
|
||||
save: vi.fn(async next => {
|
||||
persisted = next;
|
||||
}),
|
||||
clear: vi.fn(async () => {
|
||||
persisted = null;
|
||||
}),
|
||||
};
|
||||
const transport = { refresh: vi.fn(async () => response(2)) };
|
||||
return {
|
||||
broker: new AuthTokenBroker(storage, transport, {
|
||||
now: () => now,
|
||||
retryDelays: [],
|
||||
}),
|
||||
storage,
|
||||
transport,
|
||||
persisted: () => persisted,
|
||||
};
|
||||
}
|
||||
|
||||
describe('AuthTokenBroker', () => {
|
||||
test('proactively refreshes and atomically persists before publishing', async () => {
|
||||
const { broker, storage, persisted } = setup(pair(1, now + 30_000));
|
||||
const states: string[] = [];
|
||||
broker.observeAuthState(state => states.push(state.status));
|
||||
|
||||
await expect(broker.getValidAccessToken()).resolves.toBe('access-2');
|
||||
expect(storage.save).toHaveBeenCalledOnce();
|
||||
expect(persisted()?.accessToken).toBe('access-2');
|
||||
expect(states).toEqual([
|
||||
'initializing',
|
||||
'authenticated',
|
||||
'refreshing',
|
||||
'authenticated',
|
||||
]);
|
||||
});
|
||||
|
||||
test('collapses concurrent refresh into one transport request', async () => {
|
||||
const { broker, transport } = setup(pair(1, now + 30_000));
|
||||
const tokens = await Promise.all(
|
||||
Array.from({ length: 100 }, () => broker.getValidAccessToken())
|
||||
);
|
||||
|
||||
expect(new Set(tokens)).toEqual(new Set(['access-2']));
|
||||
expect(transport.refresh).toHaveBeenCalledOnce();
|
||||
});
|
||||
|
||||
test('keeps credentials for transient failures', async () => {
|
||||
const { broker, storage, transport } = setup(pair(1, now + 30_000));
|
||||
transport.refresh.mockRejectedValueOnce({
|
||||
code: 'NETWORK_ERROR',
|
||||
config: { body: 'refresh-1' },
|
||||
});
|
||||
const states: unknown[] = [];
|
||||
broker.observeAuthState(state => states.push(state));
|
||||
|
||||
const error = await broker.getValidAccessToken().catch(error => error);
|
||||
expect(error).toMatchObject({
|
||||
transient: true,
|
||||
});
|
||||
expect(JSON.stringify(error)).not.toContain('refresh-1');
|
||||
expect(JSON.stringify(states)).not.toContain('refresh-1');
|
||||
expect(storage.clear).not.toHaveBeenCalled();
|
||||
expect(states.at(-1)).toMatchObject({
|
||||
status: 'offline-authenticated',
|
||||
code: 'NETWORK_ERROR',
|
||||
});
|
||||
});
|
||||
|
||||
test('retries transient refresh failures with bounded jittered backoff', async () => {
|
||||
const { storage, transport } = setup(pair(1, now + 30_000));
|
||||
transport.refresh
|
||||
.mockRejectedValueOnce(new TypeError('offline'))
|
||||
.mockResolvedValueOnce(response(2));
|
||||
const sleep = vi.fn(async () => {});
|
||||
const broker = new AuthTokenBroker(storage, transport, {
|
||||
now: () => now,
|
||||
random: () => 0.5,
|
||||
retryDelays: [250],
|
||||
sleep,
|
||||
});
|
||||
|
||||
await expect(broker.getValidAccessToken()).resolves.toBe('access-2');
|
||||
expect(transport.refresh).toHaveBeenCalledTimes(2);
|
||||
expect(sleep).toHaveBeenCalledWith(250);
|
||||
});
|
||||
|
||||
test('clears credentials only for permanent auth failures', async () => {
|
||||
const { broker, storage, transport } = setup(pair(1, now + 30_000));
|
||||
transport.refresh.mockRejectedValueOnce({ code: 'AUTH_SESSION_REVOKED' });
|
||||
|
||||
await expect(broker.getValidAccessToken()).rejects.toMatchObject({
|
||||
code: 'AUTH_SESSION_REVOKED',
|
||||
transient: false,
|
||||
});
|
||||
expect(storage.clear).toHaveBeenCalledOnce();
|
||||
});
|
||||
|
||||
test('retries a request once only for access-token expiration', async () => {
|
||||
const { broker, transport } = setup(pair(1, now + 600_000));
|
||||
const request = vi
|
||||
.fn<(token: string) => Promise<string>>()
|
||||
.mockRejectedValueOnce({ code: 'ACCESS_TOKEN_EXPIRED' })
|
||||
.mockResolvedValueOnce('ok');
|
||||
|
||||
await expect(withAuthRetry(broker, request)).resolves.toBe('ok');
|
||||
expect(request).toHaveBeenCalledTimes(2);
|
||||
expect(transport.refresh).toHaveBeenCalledOnce();
|
||||
});
|
||||
|
||||
test('does not retry permission failures', async () => {
|
||||
const { broker, transport } = setup(pair(1, now + 600_000));
|
||||
await expect(
|
||||
withAuthRetry(broker, async () => {
|
||||
throw new AuthSessionError('FORBIDDEN', false);
|
||||
})
|
||||
).rejects.toMatchObject({ code: 'FORBIDDEN' });
|
||||
await expect(
|
||||
createRealtimeAuthAdapter(broker).recover({ code: 'FORBIDDEN' })
|
||||
).rejects.toMatchObject({ code: 'FORBIDDEN' });
|
||||
expect(transport.refresh).not.toHaveBeenCalled();
|
||||
});
|
||||
|
||||
test('serializes initialization before set and clear mutations', async () => {
|
||||
const loaded = deferred<AuthTokenPair | null>();
|
||||
let persisted: AuthTokenPair | null = pair(1, now + 600_000);
|
||||
const storage: AuthTokenStorage = {
|
||||
load: vi.fn(() => loaded.promise),
|
||||
save: vi.fn(async next => {
|
||||
persisted = next;
|
||||
}),
|
||||
clear: vi.fn(async () => {
|
||||
persisted = null;
|
||||
}),
|
||||
};
|
||||
const broker = new AuthTokenBroker(
|
||||
storage,
|
||||
{ refresh: vi.fn() },
|
||||
{
|
||||
now: () => now,
|
||||
}
|
||||
);
|
||||
const setting = broker.set(response(2));
|
||||
loaded.resolve(persisted);
|
||||
await setting;
|
||||
|
||||
expect(await broker.getValidAccessToken(0)).toBe('access-2');
|
||||
await broker.clear('logout');
|
||||
expect(persisted).toBeNull();
|
||||
expect(await broker.getValidAccessToken()).toBeNull();
|
||||
});
|
||||
|
||||
test('does not let delayed initialization undo clear', async () => {
|
||||
const loaded = deferred<AuthTokenPair | null>();
|
||||
let persisted: AuthTokenPair | null = pair(1, now + 600_000);
|
||||
const storage: AuthTokenStorage = {
|
||||
load: vi.fn(() => loaded.promise),
|
||||
save: vi.fn(),
|
||||
clear: vi.fn(async () => {
|
||||
persisted = null;
|
||||
}),
|
||||
};
|
||||
const broker = new AuthTokenBroker(
|
||||
storage,
|
||||
{ refresh: vi.fn() },
|
||||
{
|
||||
now: () => now,
|
||||
}
|
||||
);
|
||||
const clearing = broker.clear('logout');
|
||||
loaded.resolve(persisted);
|
||||
await clearing;
|
||||
|
||||
expect(persisted).toBeNull();
|
||||
expect(await broker.getValidAccessToken()).toBeNull();
|
||||
});
|
||||
|
||||
test('retries pending secure-store persistence without rotating again', async () => {
|
||||
const { broker, storage, transport, persisted } = setup(
|
||||
pair(1, now + 30_000)
|
||||
);
|
||||
vi.mocked(storage.save).mockRejectedValueOnce(new Error('locked'));
|
||||
|
||||
await expect(broker.getValidAccessToken()).rejects.toMatchObject({
|
||||
code: 'AUTH_TOKEN_STORAGE_UNAVAILABLE',
|
||||
});
|
||||
const recovered = await broker.refresh('storage-retry');
|
||||
expect(recovered).toMatchObject({ accessToken: 'access-2' });
|
||||
expect(JSON.stringify(recovered)).not.toContain('refresh-2');
|
||||
expect(transport.refresh).toHaveBeenCalledOnce();
|
||||
expect(persisted()?.accessToken).toBe('access-2');
|
||||
});
|
||||
|
||||
test('does not resurrect a session when clear races refresh', async () => {
|
||||
const { broker, transport, persisted } = setup(pair(1, now + 600_000));
|
||||
await broker.getValidAccessToken(0);
|
||||
const rotated = deferred<AuthTokenResponse>();
|
||||
transport.refresh.mockReturnValueOnce(rotated.promise);
|
||||
const refreshing = broker.refresh('manual');
|
||||
await Promise.resolve();
|
||||
const clearing = broker.clear('logout');
|
||||
rotated.resolve(response(2));
|
||||
|
||||
await expect(refreshing).rejects.toMatchObject({
|
||||
code: 'AUTH_OPERATION_CANCELLED',
|
||||
});
|
||||
await clearing;
|
||||
expect(persisted()).toBeNull();
|
||||
expect(await broker.getValidAccessToken()).toBeNull();
|
||||
});
|
||||
|
||||
test('isolates observers and rejects malformed persisted records', async () => {
|
||||
const { broker } = setup(pair(1, now + 600_000));
|
||||
const states: unknown[] = [];
|
||||
broker.observeAuthState(() => {
|
||||
throw new Error('observer failure');
|
||||
});
|
||||
broker.observeAuthState(state => states.push(state));
|
||||
await expect(broker.getValidAccessToken()).resolves.toBe('access-1');
|
||||
expect(JSON.stringify(states)).not.toContain('refresh-1');
|
||||
|
||||
const storage: AuthTokenStorage = {
|
||||
load: vi.fn(async () => ({ version: 2 }) as never),
|
||||
save: vi.fn(),
|
||||
clear: vi.fn(),
|
||||
};
|
||||
const malformed = new AuthTokenBroker(storage, { refresh: vi.fn() });
|
||||
await expect(malformed.getValidAccessToken()).resolves.toBeNull();
|
||||
expect(storage.clear).toHaveBeenCalledOnce();
|
||||
});
|
||||
|
||||
test('shares one refresh across HTTP replay and realtime recovery', async () => {
|
||||
const { broker, transport } = setup(pair(1, now + 600_000));
|
||||
const rotated = deferred<AuthTokenResponse>();
|
||||
transport.refresh.mockReturnValueOnce(rotated.promise);
|
||||
const request = vi
|
||||
.fn<(token: string) => Promise<string>>()
|
||||
.mockRejectedValueOnce({ code: 'ACCESS_TOKEN_EXPIRED' })
|
||||
.mockResolvedValueOnce('ok');
|
||||
const realtime = createRealtimeAuthAdapter(broker);
|
||||
const http = createRequestAuthAdapter(broker).execute(request);
|
||||
const socket = realtime.recover({ code: 'ACCESS_TOKEN_EXPIRED' });
|
||||
await Promise.resolve();
|
||||
rotated.resolve(response(2));
|
||||
|
||||
await expect(Promise.all([http, socket])).resolves.toEqual([
|
||||
'ok',
|
||||
'access-2',
|
||||
]);
|
||||
expect(transport.refresh).toHaveBeenCalledOnce();
|
||||
});
|
||||
|
||||
test.each([
|
||||
{ name: 'permanent', error: { code: 'AUTH_SESSION_REVOKED' } },
|
||||
{ name: 'transient', error: new TypeError('offline') },
|
||||
])('ignores a stale $name refresh failure after a new login', async item => {
|
||||
const { broker, transport, persisted } = setup(pair(1, now + 600_000));
|
||||
await broker.getValidAccessToken(0);
|
||||
const oldRefresh = deferred<AuthTokenResponse>();
|
||||
transport.refresh.mockReturnValueOnce(oldRefresh.promise);
|
||||
const refreshing = broker.refresh('manual');
|
||||
await Promise.resolve();
|
||||
await broker.set(response(3));
|
||||
oldRefresh.reject(item.error);
|
||||
|
||||
await expect(refreshing).rejects.toMatchObject({
|
||||
code: 'AUTH_OPERATION_CANCELLED',
|
||||
});
|
||||
expect(await broker.getValidAccessToken(0)).toBe('access-3');
|
||||
expect(persisted()?.accessToken).toBe('access-3');
|
||||
});
|
||||
|
||||
test('retries initialization after a transient storage failure', async () => {
|
||||
const stored = pair(1, now + 600_000);
|
||||
const storage: AuthTokenStorage = {
|
||||
load: vi
|
||||
.fn<() => Promise<AuthTokenPair | null>>()
|
||||
.mockRejectedValueOnce(new Error('locked'))
|
||||
.mockResolvedValueOnce(stored),
|
||||
save: vi.fn(),
|
||||
clear: vi.fn(),
|
||||
};
|
||||
const broker = new AuthTokenBroker(
|
||||
storage,
|
||||
{ refresh: vi.fn() },
|
||||
{
|
||||
now: () => now,
|
||||
}
|
||||
);
|
||||
const states: string[] = [];
|
||||
broker.observeAuthState(state => states.push(state.status));
|
||||
|
||||
await expect(broker.getValidAccessToken()).rejects.toMatchObject({
|
||||
code: 'AUTH_TOKEN_STORAGE_UNAVAILABLE',
|
||||
transient: true,
|
||||
});
|
||||
await expect(broker.getValidAccessToken()).resolves.toBe('access-1');
|
||||
expect(storage.load).toHaveBeenCalledTimes(2);
|
||||
expect(states).toEqual(['initializing', 'authenticated']);
|
||||
});
|
||||
|
||||
test.each([0, Number.NaN, Number.POSITIVE_INFINITY, Number.MAX_VALUE])(
|
||||
'rejects malformed expiresIn %s before persistence',
|
||||
async expiresIn => {
|
||||
const { broker, storage, transport } = setup(pair(1, now + 30_000));
|
||||
transport.refresh.mockResolvedValueOnce({ ...response(2), expiresIn });
|
||||
|
||||
await expect(broker.getValidAccessToken()).rejects.toMatchObject({
|
||||
code: 'AUTH_TOKEN_RESPONSE_INVALID',
|
||||
});
|
||||
expect(storage.save).not.toHaveBeenCalled();
|
||||
}
|
||||
);
|
||||
|
||||
test('normalizes clear storage failures after clearing memory', async () => {
|
||||
const { broker, storage } = setup(pair(1, now + 600_000));
|
||||
await broker.getValidAccessToken(0);
|
||||
vi.mocked(storage.clear).mockRejectedValueOnce(new Error('locked'));
|
||||
|
||||
await expect(broker.clear('logout')).rejects.toMatchObject({
|
||||
code: 'AUTH_TOKEN_STORAGE_UNAVAILABLE',
|
||||
});
|
||||
await expect(broker.getValidAccessToken()).resolves.toBeNull();
|
||||
});
|
||||
|
||||
test('revokes captured credential even when local clear fails', async () => {
|
||||
const { broker, storage } = setup(pair(1, now + 600_000));
|
||||
await broker.getValidAccessToken(0);
|
||||
vi.mocked(storage.clear).mockRejectedValueOnce(new Error('locked'));
|
||||
const revoke = vi.fn(async () => {});
|
||||
|
||||
await expect(broker.revoke('logout', revoke)).rejects.toMatchObject({
|
||||
code: 'AUTH_TOKEN_STORAGE_UNAVAILABLE',
|
||||
});
|
||||
expect(revoke).toHaveBeenCalledWith('refresh-1');
|
||||
await expect(broker.getValidAccessToken()).resolves.toBeNull();
|
||||
});
|
||||
|
||||
test('initializes when the first consumer only observes state', async () => {
|
||||
const { broker, storage } = setup(pair(1, now + 600_000));
|
||||
const states: string[] = [];
|
||||
broker.observeAuthState(state => states.push(state.status));
|
||||
await vi.waitFor(() => expect(states.at(-1)).toBe('authenticated'));
|
||||
|
||||
expect(storage.load).toHaveBeenCalledOnce();
|
||||
expect(states).toEqual(['initializing', 'authenticated']);
|
||||
});
|
||||
|
||||
test('redacts unknown external error codes', async () => {
|
||||
const { broker, transport } = setup(pair(1, now + 30_000));
|
||||
transport.refresh.mockRejectedValueOnce({
|
||||
code: 'secret-refresh-1',
|
||||
});
|
||||
|
||||
await expect(broker.getValidAccessToken()).rejects.toMatchObject({
|
||||
code: 'AUTH_SESSION_TEMPORARILY_UNAVAILABLE',
|
||||
});
|
||||
});
|
||||
});
|
||||
@@ -0,0 +1,464 @@
|
||||
export const PERMANENT_AUTH_CODES = new Set([
|
||||
'ACCESS_TOKEN_INVALID',
|
||||
'AUTH_SESSION_EXPIRED',
|
||||
'AUTH_SESSION_REVOKED',
|
||||
'REFRESH_TOKEN_INVALID',
|
||||
'REFRESH_TOKEN_REUSED',
|
||||
'UNSUPPORTED_CLIENT_VERSION',
|
||||
]);
|
||||
const PUBLIC_AUTH_CODES = new Set([
|
||||
...PERMANENT_AUTH_CODES,
|
||||
'ACCESS_TOKEN_EXPIRED',
|
||||
'AUTH_SESSION_TEMPORARILY_UNAVAILABLE',
|
||||
'AUTH_TOKEN_RESPONSE_INVALID',
|
||||
'AUTH_TOKEN_STORAGE_UNAVAILABLE',
|
||||
'FORBIDDEN',
|
||||
'NETWORK_ERROR',
|
||||
'TOO_MANY_REQUESTS',
|
||||
]);
|
||||
const SAFE_AUTH_CODES = new Set([
|
||||
...PUBLIC_AUTH_CODES,
|
||||
'AUTH_OPERATION_CANCELLED',
|
||||
'AUTH_SESSION_EMPTY',
|
||||
]);
|
||||
|
||||
export interface AuthTokenPair {
|
||||
version: 1;
|
||||
tokenType: 'Bearer';
|
||||
accessToken: string;
|
||||
accessExpiresAt: string;
|
||||
refreshToken: string;
|
||||
refreshExpiresAt: string;
|
||||
session: {
|
||||
id: string;
|
||||
absoluteExpiresAt: string;
|
||||
};
|
||||
}
|
||||
|
||||
export interface AuthTokenResponse {
|
||||
tokenType: 'Bearer';
|
||||
accessToken: string;
|
||||
expiresIn: number;
|
||||
refreshToken: string;
|
||||
refreshExpiresAt: string;
|
||||
session: AuthTokenPair['session'];
|
||||
}
|
||||
|
||||
export interface AuthSessionSnapshot {
|
||||
accessExpiresAt: string;
|
||||
refreshExpiresAt: string;
|
||||
session: AuthTokenPair['session'];
|
||||
}
|
||||
|
||||
export interface AuthAccessToken extends AuthSessionSnapshot {
|
||||
accessToken: string;
|
||||
}
|
||||
|
||||
export type AuthState =
|
||||
| { status: 'initializing' }
|
||||
| { status: 'empty' }
|
||||
| { status: 'authenticated'; session: AuthSessionSnapshot }
|
||||
| { status: 'refreshing'; session: AuthSessionSnapshot; reason: string }
|
||||
| {
|
||||
status: 'offline-authenticated';
|
||||
session: AuthSessionSnapshot;
|
||||
code: string;
|
||||
}
|
||||
| { status: 'revoked'; code: string };
|
||||
|
||||
export interface AuthTokenStorage {
|
||||
load(): Promise<AuthTokenPair | null>;
|
||||
save(pair: AuthTokenPair): Promise<void>;
|
||||
clear(): Promise<void>;
|
||||
}
|
||||
|
||||
export interface AuthTokenTransport {
|
||||
refresh(refreshToken: string): Promise<AuthTokenResponse>;
|
||||
}
|
||||
|
||||
export interface AuthTokenBrokerContract {
|
||||
getValidAccessToken(minValidity?: number): Promise<string | null>;
|
||||
refresh(reason: string): Promise<AuthAccessToken>;
|
||||
clear(reason: string): Promise<void>;
|
||||
observeAuthState(listener: (state: AuthState) => void): () => void;
|
||||
}
|
||||
|
||||
export interface AuthTokenBrokerOptions {
|
||||
now?: () => number;
|
||||
random?: () => number;
|
||||
retryDelays?: number[];
|
||||
sleep?: (delay: number) => Promise<void>;
|
||||
}
|
||||
|
||||
export class AuthSessionError extends Error {
|
||||
constructor(
|
||||
readonly code: string,
|
||||
readonly transient: boolean
|
||||
) {
|
||||
super(code);
|
||||
}
|
||||
}
|
||||
|
||||
export class AuthTokenBroker implements AuthTokenBrokerContract {
|
||||
private pair: AuthTokenPair | null = null;
|
||||
private initialized: Promise<void> | null = null;
|
||||
private refreshPromise: Promise<AuthTokenPair> | null = null;
|
||||
private pending: { pair: AuthTokenPair; epoch: number } | null = null;
|
||||
private mutationEpoch = 0;
|
||||
private storageMutation: Promise<void> = Promise.resolve();
|
||||
private readonly listeners = new Set<(state: AuthState) => void>();
|
||||
private state: AuthState = { status: 'initializing' };
|
||||
private readonly now: () => number;
|
||||
private readonly random: () => number;
|
||||
private readonly retryDelays: number[];
|
||||
private readonly sleep: (delay: number) => Promise<void>;
|
||||
|
||||
constructor(
|
||||
private readonly storage: AuthTokenStorage,
|
||||
private readonly transport: AuthTokenTransport,
|
||||
options: AuthTokenBrokerOptions = {}
|
||||
) {
|
||||
this.now = options.now ?? Date.now;
|
||||
this.random = options.random ?? Math.random;
|
||||
this.retryDelays = options.retryDelays ?? [250, 1000];
|
||||
this.sleep =
|
||||
options.sleep ??
|
||||
(delay => new Promise(resolve => setTimeout(resolve, delay)));
|
||||
}
|
||||
|
||||
observeAuthState(listener: (state: AuthState) => void) {
|
||||
this.listeners.add(listener);
|
||||
try {
|
||||
listener(this.state);
|
||||
} catch {
|
||||
// Observers cannot participate in credential lifecycle decisions.
|
||||
}
|
||||
void this.initialize().catch(() => {});
|
||||
return () => {
|
||||
this.listeners.delete(listener);
|
||||
};
|
||||
}
|
||||
|
||||
async set(response: AuthTokenResponse) {
|
||||
await this.initialize();
|
||||
const epoch = ++this.mutationEpoch;
|
||||
const pair = this.toPair(response);
|
||||
this.pending = { pair, epoch };
|
||||
return this.toAccessToken(await this.persistPending());
|
||||
}
|
||||
|
||||
async getValidAccessToken(minValidity = 60_000) {
|
||||
await this.initialize();
|
||||
if (!this.pair) return null;
|
||||
if (Date.parse(this.pair.accessExpiresAt) - this.now() <= minValidity) {
|
||||
await this.refresh('proactive');
|
||||
}
|
||||
return this.pair?.accessToken ?? null;
|
||||
}
|
||||
|
||||
async refresh(reason: string) {
|
||||
await this.initialize();
|
||||
if (!this.pair) throw new AuthSessionError('AUTH_SESSION_EMPTY', false);
|
||||
if (!this.refreshPromise) {
|
||||
this.refreshPromise = (
|
||||
this.pending ? this.persistPending() : this.performRefresh(reason)
|
||||
).finally(() => {
|
||||
this.refreshPromise = null;
|
||||
});
|
||||
}
|
||||
return this.toAccessToken(await this.refreshPromise);
|
||||
}
|
||||
|
||||
async clear(_reason: string) {
|
||||
await this.initialize();
|
||||
++this.mutationEpoch;
|
||||
this.pending = null;
|
||||
let failed = false;
|
||||
try {
|
||||
await this.mutateStorage(() => this.storage.clear());
|
||||
} catch {
|
||||
failed = true;
|
||||
} finally {
|
||||
this.pair = null;
|
||||
this.publish({ status: 'empty' });
|
||||
}
|
||||
if (failed) {
|
||||
throw new AuthSessionError('AUTH_TOKEN_STORAGE_UNAVAILABLE', true);
|
||||
}
|
||||
}
|
||||
|
||||
async revoke(
|
||||
_reason: string,
|
||||
revokeToken: (refreshToken: string) => Promise<void>
|
||||
) {
|
||||
await this.initialize();
|
||||
++this.mutationEpoch;
|
||||
this.pending = null;
|
||||
const refreshToken = this.pair?.refreshToken;
|
||||
let storageError: unknown;
|
||||
try {
|
||||
await this.mutateStorage(() => this.storage.clear());
|
||||
} catch (error) {
|
||||
storageError = error;
|
||||
} finally {
|
||||
this.pair = null;
|
||||
this.publish({ status: 'empty' });
|
||||
}
|
||||
if (refreshToken) await revokeToken(refreshToken);
|
||||
if (storageError) {
|
||||
throw new AuthSessionError('AUTH_TOKEN_STORAGE_UNAVAILABLE', true);
|
||||
}
|
||||
}
|
||||
|
||||
private async initialize() {
|
||||
if (!this.initialized) {
|
||||
const attempt = Promise.resolve()
|
||||
.then(() => this.storage.load())
|
||||
.then(async pair => {
|
||||
if (pair && !isAuthTokenPair(pair)) {
|
||||
await this.mutateStorage(() => this.storage.clear());
|
||||
pair = null;
|
||||
}
|
||||
this.pair = pair;
|
||||
this.publish(
|
||||
pair
|
||||
? { status: 'authenticated', session: this.toSnapshot(pair) }
|
||||
: { status: 'empty' }
|
||||
);
|
||||
})
|
||||
.catch(() => {
|
||||
throw new AuthSessionError('AUTH_TOKEN_STORAGE_UNAVAILABLE', true);
|
||||
});
|
||||
this.initialized = attempt;
|
||||
try {
|
||||
await attempt;
|
||||
} catch (error) {
|
||||
if (this.initialized === attempt) this.initialized = null;
|
||||
throw error;
|
||||
}
|
||||
return;
|
||||
}
|
||||
await this.initialized;
|
||||
}
|
||||
|
||||
private async performRefresh(reason: string) {
|
||||
const current = this.pair;
|
||||
if (!current) throw new AuthSessionError('AUTH_SESSION_EMPTY', false);
|
||||
const epoch = this.mutationEpoch;
|
||||
this.publish({
|
||||
status: 'refreshing',
|
||||
session: this.toSnapshot(current),
|
||||
reason,
|
||||
});
|
||||
try {
|
||||
const pair = this.toPair(await this.requestRefresh(current.refreshToken));
|
||||
if (epoch !== this.mutationEpoch) {
|
||||
throw new AuthSessionError('AUTH_OPERATION_CANCELLED', true);
|
||||
}
|
||||
this.pending = { pair, epoch };
|
||||
return await this.persistPending();
|
||||
} catch (error) {
|
||||
const classified = classifyAuthError(error);
|
||||
if (
|
||||
classified.code === 'AUTH_OPERATION_CANCELLED' ||
|
||||
epoch !== this.mutationEpoch
|
||||
) {
|
||||
throw new AuthSessionError('AUTH_OPERATION_CANCELLED', true);
|
||||
}
|
||||
if (!classified.transient) {
|
||||
try {
|
||||
await this.mutateStorage(() => this.storage.clear());
|
||||
} catch {
|
||||
// The in-memory credential must still become unusable immediately.
|
||||
} finally {
|
||||
this.pair = null;
|
||||
this.publish({ status: 'revoked', code: classified.code });
|
||||
}
|
||||
} else {
|
||||
this.publish({
|
||||
status: 'offline-authenticated',
|
||||
session: this.toSnapshot(current),
|
||||
code: classified.code,
|
||||
});
|
||||
}
|
||||
throw classified;
|
||||
}
|
||||
}
|
||||
|
||||
private async persistPending() {
|
||||
const pending = this.pending;
|
||||
if (!pending) throw new AuthSessionError('AUTH_SESSION_EMPTY', false);
|
||||
try {
|
||||
await this.mutateStorage(() => this.storage.save(pending.pair));
|
||||
} catch {
|
||||
throw new AuthSessionError('AUTH_TOKEN_STORAGE_UNAVAILABLE', true);
|
||||
}
|
||||
if (pending.epoch !== this.mutationEpoch) {
|
||||
throw new AuthSessionError('AUTH_OPERATION_CANCELLED', true);
|
||||
}
|
||||
this.pending = null;
|
||||
this.pair = pending.pair;
|
||||
this.publish({
|
||||
status: 'authenticated',
|
||||
session: this.toSnapshot(pending.pair),
|
||||
});
|
||||
return pending.pair;
|
||||
}
|
||||
|
||||
private async mutateStorage<T>(operation: () => Promise<T>) {
|
||||
const result = this.storageMutation.then(operation, operation);
|
||||
this.storageMutation = result.then(
|
||||
() => {},
|
||||
() => {}
|
||||
);
|
||||
return await result;
|
||||
}
|
||||
|
||||
private async requestRefresh(refreshToken: string) {
|
||||
for (let attempt = 0; ; attempt++) {
|
||||
try {
|
||||
return await this.transport.refresh(refreshToken);
|
||||
} catch (error) {
|
||||
const classified = classifyAuthError(error);
|
||||
const delay = this.retryDelays[attempt];
|
||||
if (!classified.transient || delay === undefined) throw classified;
|
||||
await this.sleep(delay * (0.75 + this.random() * 0.5));
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private toPair(response: AuthTokenResponse): AuthTokenPair {
|
||||
const accessExpiresAt = this.now() + response.expiresIn * 1000;
|
||||
if (
|
||||
!Number.isFinite(response.expiresIn) ||
|
||||
response.expiresIn <= 0 ||
|
||||
!Number.isFinite(accessExpiresAt) ||
|
||||
Math.abs(accessExpiresAt) > 8.64e15
|
||||
) {
|
||||
throw new AuthSessionError('AUTH_TOKEN_RESPONSE_INVALID', true);
|
||||
}
|
||||
const pair: AuthTokenPair = {
|
||||
version: 1,
|
||||
tokenType: response.tokenType,
|
||||
accessToken: response.accessToken,
|
||||
accessExpiresAt: new Date(accessExpiresAt).toISOString(),
|
||||
refreshToken: response.refreshToken,
|
||||
refreshExpiresAt: response.refreshExpiresAt,
|
||||
session: response.session,
|
||||
};
|
||||
if (!isAuthTokenPair(pair)) {
|
||||
throw new AuthSessionError('AUTH_TOKEN_RESPONSE_INVALID', true);
|
||||
}
|
||||
return pair;
|
||||
}
|
||||
|
||||
private toSnapshot(pair: AuthTokenPair): AuthSessionSnapshot {
|
||||
return {
|
||||
accessExpiresAt: pair.accessExpiresAt,
|
||||
refreshExpiresAt: pair.refreshExpiresAt,
|
||||
session: { ...pair.session },
|
||||
};
|
||||
}
|
||||
|
||||
private toAccessToken(pair: AuthTokenPair): AuthAccessToken {
|
||||
return {
|
||||
...this.toSnapshot(pair),
|
||||
accessToken: pair.accessToken,
|
||||
};
|
||||
}
|
||||
|
||||
private publish(state: AuthState) {
|
||||
this.state = state;
|
||||
for (const listener of this.listeners) {
|
||||
try {
|
||||
listener(state);
|
||||
} catch {
|
||||
// Observers cannot participate in credential lifecycle decisions.
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
export function classifyAuthError(error: unknown) {
|
||||
if (error instanceof AuthSessionError) {
|
||||
return SAFE_AUTH_CODES.has(error.code)
|
||||
? error
|
||||
: new AuthSessionError('AUTH_SESSION_TEMPORARILY_UNAVAILABLE', true);
|
||||
}
|
||||
const code =
|
||||
typeof error === 'object' &&
|
||||
error !== null &&
|
||||
'code' in error &&
|
||||
typeof error.code === 'string'
|
||||
? error.code
|
||||
: 'AUTH_SESSION_TEMPORARILY_UNAVAILABLE';
|
||||
const publicCode = PUBLIC_AUTH_CODES.has(code)
|
||||
? code
|
||||
: 'AUTH_SESSION_TEMPORARILY_UNAVAILABLE';
|
||||
return new AuthSessionError(
|
||||
publicCode,
|
||||
!PERMANENT_AUTH_CODES.has(publicCode)
|
||||
);
|
||||
}
|
||||
|
||||
export async function withAuthRetry<T>(
|
||||
broker: AuthTokenBrokerContract,
|
||||
request: (accessToken: string) => Promise<T>
|
||||
) {
|
||||
const token = await broker.getValidAccessToken();
|
||||
if (!token) throw new AuthSessionError('AUTH_SESSION_EMPTY', false);
|
||||
try {
|
||||
return await request(token);
|
||||
} catch (error) {
|
||||
const classified = classifyAuthError(error);
|
||||
if (classified.code !== 'ACCESS_TOKEN_EXPIRED') throw classified;
|
||||
const pair = await broker.refresh('access-token-expired');
|
||||
return await request(pair.accessToken);
|
||||
}
|
||||
}
|
||||
|
||||
export function createRealtimeAuthAdapter(broker: AuthTokenBrokerContract) {
|
||||
return {
|
||||
getAccessToken: () => broker.getValidAccessToken(120_000),
|
||||
recover: async (error: unknown) => {
|
||||
const classified = classifyAuthError(error);
|
||||
if (classified.code !== 'ACCESS_TOKEN_EXPIRED') throw classified;
|
||||
return (await broker.refresh('realtime-access-token-expired'))
|
||||
.accessToken;
|
||||
},
|
||||
};
|
||||
}
|
||||
|
||||
export function createRequestAuthAdapter(broker: AuthTokenBrokerContract) {
|
||||
return {
|
||||
execute: <T>(request: (accessToken: string) => Promise<T>) =>
|
||||
withAuthRetry(broker, request),
|
||||
};
|
||||
}
|
||||
|
||||
export function createWorkerAuthAdapter(broker: AuthTokenBrokerContract) {
|
||||
return {
|
||||
getAccessToken: () => broker.getValidAccessToken(60_000),
|
||||
};
|
||||
}
|
||||
|
||||
export function isAuthTokenPair(value: unknown): value is AuthTokenPair {
|
||||
if (!value || typeof value !== 'object') return false;
|
||||
const pair = value as Partial<AuthTokenPair>;
|
||||
return (
|
||||
pair.version === 1 &&
|
||||
pair.tokenType === 'Bearer' &&
|
||||
typeof pair.accessToken === 'string' &&
|
||||
pair.accessToken.length > 0 &&
|
||||
typeof pair.refreshToken === 'string' &&
|
||||
pair.refreshToken.length > 0 &&
|
||||
typeof pair.accessExpiresAt === 'string' &&
|
||||
Number.isFinite(Date.parse(pair.accessExpiresAt)) &&
|
||||
typeof pair.refreshExpiresAt === 'string' &&
|
||||
Number.isFinite(Date.parse(pair.refreshExpiresAt)) &&
|
||||
!!pair.session &&
|
||||
typeof pair.session.id === 'string' &&
|
||||
typeof pair.session.absoluteExpiresAt === 'string' &&
|
||||
Number.isFinite(Date.parse(pair.session.absoluteExpiresAt))
|
||||
);
|
||||
}
|
||||
Reference in New Issue
Block a user