diff --git a/src/core/car-fetch.ts b/src/core/car-fetch.ts index 7acf3f2..65597c4 100644 --- a/src/core/car-fetch.ts +++ b/src/core/car-fetch.ts @@ -90,6 +90,22 @@ async function walkMST( // ─── public API ────────────────────────────────────────────────────────────── +/** + * Thrown when the PDS returns 401 on com.atproto.sync.getRepo. + * Callers can catch this specifically to surface a re-auth prompt rather than + * treating it as a generic network error. + */ +export class CARFetchUnauthorizedError extends Error { + constructor(pdsUrl: string, did: string) { + super( + `CAR fetch returned 401 Unauthorized for ${did} at ${pdsUrl}. ` + + `The PDS requires authentication but a valid token could not be obtained. ` + + `Try signing out and back in to refresh your session.` + ); + this.name = 'CARFetchUnauthorizedError'; + } +} + export interface CARRecord { rkey: string; uri: string; @@ -120,6 +136,9 @@ export async function fetchRepoViaCAR( const response = await fetch(url, { headers, signal }); if (!response.ok) { + if (response.status === 401) { + throw new CARFetchUnauthorizedError(pdsUrl, did); + } throw new Error(`CAR fetch failed: ${response.status} ${response.statusText}`); } @@ -180,7 +199,9 @@ export function getPdsUrlFromAgent(agent: unknown): string { export async function getAgentToken(agent: unknown): Promise { const a = agent as Record; - // Password-auth AtpAgent: session carries a plain JWT. + // Password-auth CredentialSession (AtpAgent): + // session.accessJwt holds the current JWT. It may be expired — callers + // should handle CARFetchUnauthorizedError and retry after refreshing. const jwt = (a['session'] as any)?.accessJwt; if (jwt) return jwt as string; @@ -191,7 +212,19 @@ export async function getAgentToken(agent: unknown): Promise const tokens = await sm.getTokens() as { accessToken?: string } | null; if (tokens?.accessToken) return tokens.accessToken; } catch { - // If the OAuth token is expired and can't be refreshed silently, fall through. + // Token read failed — try a silent refresh before giving up. + } + + // If getTokens() returned nothing (expired session), attempt a silent + // refresh via the session manager and retry once. + if (typeof sm?.refresh === 'function') { + try { + await sm.refresh(); + const refreshed = await sm.getTokens() as { accessToken?: string } | null; + if (refreshed?.accessToken) return refreshed.accessToken; + } catch { + // Refresh failed — fall through and return undefined. + } } } diff --git a/src/core/sync.ts b/src/core/sync.ts index d4cf757..5cb0ec7 100644 --- a/src/core/sync.ts +++ b/src/core/sync.ts @@ -7,7 +7,7 @@ import type { Agent } from '@atproto/api'; import type { PlayRecord } from './types.js'; import { RECORD_TYPE } from './config.js'; -import { fetchRepoViaCAR, getPdsUrlFromAgent, getAgentToken } from './car-fetch.js'; +import { fetchRepoViaCAR, getPdsUrlFromAgent, getAgentToken, CARFetchUnauthorizedError } from './car-fetch.js'; export interface ExistingRecord { uri: string; @@ -44,11 +44,43 @@ export async function fetchExistingRecords( signal?.throwIfAborted(); const pdsUrl = getPdsUrlFromAgent(agent); - const token = await getAgentToken(agent); - const carRecords = await fetchRepoViaCAR(pdsUrl, did, RECORD_TYPE, signal, token); + let token = await getAgentToken(agent); + let carRecords; + try { + carRecords = await fetchRepoViaCAR(pdsUrl, did, RECORD_TYPE, signal, token); + } catch (err) { + if (err instanceof CARFetchUnauthorizedError) { + // The token we sent was invalid or expired. Try to silently refresh the + // session (works for both CredentialSession / AtpAgent and OAuth agents + // that expose a refreshSession method on their session manager) then + // retry the CAR fetch exactly once before giving up. + const sm = (agent as any)?.sessionManager; + let retried = false; + if (typeof sm?.refreshSession === 'function') { + try { + await sm.refreshSession(); + const freshToken = await getAgentToken(agent); + if (freshToken && freshToken !== token) { + carRecords = await fetchRepoViaCAR(pdsUrl, did, RECORD_TYPE, signal, freshToken); + token = freshToken; + retried = true; + } + } catch { + // Refresh or second fetch failed — fall through and throw below. + } + } + if (!retried) { + // Clear the stale session cache so the next call starts clean. + sessionCache.delete(did); + throw err; + } + } else { + throw err; + } + } const map = new Map(); - for (const rec of carRecords) { + for (const rec of carRecords!) { const value = rec.value as unknown as PlayRecord; map.set(recordKey(value), { uri: rec.uri, cid: rec.cid, value }); } @@ -76,10 +108,37 @@ export async function fetchAllRecordsForDedup( signal?.throwIfAborted(); const pdsUrl = getPdsUrlFromAgent(agent); - const token = await getAgentToken(agent); - const carRecords = await fetchRepoViaCAR(pdsUrl, did, RECORD_TYPE, signal, token); + let token = await getAgentToken(agent); + let carRecords; + try { + carRecords = await fetchRepoViaCAR(pdsUrl, did, RECORD_TYPE, signal, token); + } catch (err) { + if (err instanceof CARFetchUnauthorizedError) { + const sm = (agent as any)?.sessionManager; + let retried = false; + if (typeof sm?.refreshSession === 'function') { + try { + await sm.refreshSession(); + const freshToken = await getAgentToken(agent); + if (freshToken && freshToken !== token) { + carRecords = await fetchRepoViaCAR(pdsUrl, did, RECORD_TYPE, signal, freshToken); + token = freshToken; + retried = true; + } + } catch { + // fall through + } + } + if (!retried) { + sessionCache.delete(did); + throw err; + } + } else { + throw err; + } + } - const all: ExistingRecord[] = carRecords.map((rec) => ({ + const all: ExistingRecord[] = carRecords!.map((rec) => ({ uri: rec.uri, cid: rec.cid, value: rec.value as unknown as PlayRecord, diff --git a/src/lib/publisher.ts b/src/lib/publisher.ts index 06c4889..11be410 100644 --- a/src/lib/publisher.ts +++ b/src/lib/publisher.ts @@ -320,25 +320,24 @@ export async function publishRecordsWithApplyWrites( const rateLimitError = isRateLimitError(err); if (rateLimitError) { - log.warn('⚠️ Rate limit hit (unexpected with proactive pacing) - updating from error headers...'); - - // Extract and update from error headers - let headers: Record | undefined; - if (err?.response?.headers) { - headers = err.response.headers; - } else if (err?.headers) { - headers = err.headers; - } - - if (headers && Object.keys(headers).length > 0) { - const normalized = normalizeHeaders(headers); - const hasRateLimitHeaders = Object.keys(normalized).some(k => k.includes('ratelimit')); - if (hasRateLimitHeaders) { - rl.updateFromHeaders(normalized); - } - } - - // Wait for permit and retry + log.warn('⚠️ Rate limit hit — pausing until quota resets…'); + + // XRPCError (from @atproto/xrpc) carries response headers directly on + // err.headers as a plain Record built via + // Object.fromEntries(response.headers.entries()), so keys are already + // lowercase. There is no err.response property on this error class. + const errHeaders: Record | undefined = + err?.headers && typeof err.headers === 'object' + ? (err.headers as Record) + : undefined; + + // handleRateLimitHit zeroes remaining unconditionally — this is the + // critical fix. Previously we only called updateFromHeaders when + // headers were present, meaning a headerless 429 left state untouched + // and waitForPermit returned immediately, sending another request. + rl.handleRateLimitHit(errHeaders ? normalizeHeaders(errHeaders) : undefined); + + // Now waitForPermit will block until the window resets. await rl.waitForPermit(batchPoints); continue; diff --git a/src/utils/rate-limiter.ts b/src/utils/rate-limiter.ts index 1d49bb0..46f5803 100644 --- a/src/utils/rate-limiter.ts +++ b/src/utils/rate-limiter.ts @@ -685,6 +685,45 @@ export class RateLimiter { log.info('[RateLimiter] ✅ Wait complete - quota restored'); } + /** + * Called when the server returns a 429. + * + * Zeroes `remaining` unconditionally so the next `waitForPermit` call + * actually blocks until the window resets, regardless of whether the 429 + * response included rate-limit headers. If headers ARE present they are + * applied first (so we get an accurate `resetAt`), then remaining is + * forced to 0. + * + * @param errHeaders Optional normalised headers from the 429 error response. + */ + handleRateLimitHit(errHeaders?: Record): void { + // Apply header info first so resetAt is as accurate as possible. + if (errHeaders && Object.keys(errHeaders).length > 0) { + this.updateFromHeaders(errHeaders); + } + + const now = Math.floor(Date.now() / 1000); + const state = this.readState(); + if (state) { + state.remaining = 0; + state.updatedAt = now; + this.writeState(state); + log.warn(`[RateLimiter] 🛑 429 received — zeroed remaining, will wait until ${new Date(state.resetAt * 1000).toISOString()}`); + } else { + // No state yet — create a blocking stub with a 60-second reset. + this.writeState({ + limit: 5000, + remaining: 0, + resetAt: now + 60, + windowSeconds: 3600, + updatedAt: now, + headroomThreshold: this.headroomThreshold, + }); + this.hasLearnedFromServer = true; + log.warn('[RateLimiter] 🛑 429 received (no prior state) — blocking for 60s'); + } + } + /** * Wait for a permit with the given number of points. * Combines reserveQuota and waitForReset - loops until permit granted.