From f9f6bb761552e4c23c632542e3386a906c298db6 Mon Sep 17 00:00:00 2001 From: Ewan Croft Date: Thu, 13 Aug 2026 07:48:05 +0100 Subject: [PATCH] fix(croft-click-core): pace against observed cost, not assumed bucket MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The proactive rate pacer assumed every ATProto write costs a fixed 3 points, and applied that against whatever RateLimit-* headers the server happened to report. But a PDS can report capacity for several differently-scoped buckets (e.g. a broad per-IP request cap vs. a write-specific points budget) under the same generic header names, with no way to tell which one it is from the headers alone. Applying the write-cost constant to a request-counted bucket wildly overestimates cost per record and can stall batched writes for multiples of the pacer's 5-minute max delay, for no real reason. RateLimiter now tracks the empirically observed points-consumed-per- record between consecutive same-window readings and exposes it via getPointsPerRecord()/hasObservedPointsPerRecord(), self-correcting to whichever bucket is actually being reported instead of assuming. The pacer accepts this as an optional parameter (defaulting to the old constant when unset). Callers skip proactive pacing entirely until an empirical sample exists, relying on the existing reactive 429 handling in the meantime — which is safe regardless of which bucket is real. Also: malachite's polish command relied entirely on a cli-progress bar for output, which renders nothing meaningful when stdout isn't a TTY (CI, piped output, --non-interactive runs). Added throttled explicit log lines as a fallback so progress is actually visible in that case. Co-Authored-By: Claude Sonnet 5 --- packages/croft-click-core/src/polish.ts | 24 ++++-- .../src/proactive-rate-pacer.ts | 39 ++++++--- packages/croft-click-core/src/publisher.ts | 16 ++-- packages/croft-click-core/src/rate-limiter.ts | 79 +++++++++++++++++-- packages/malachite/src/lib/polish.ts | 17 ++++ 5 files changed, 146 insertions(+), 29 deletions(-) diff --git a/packages/croft-click-core/src/polish.ts b/packages/croft-click-core/src/polish.ts index c1166c2..bde4e48 100644 --- a/packages/croft-click-core/src/polish.ts +++ b/packages/croft-click-core/src/polish.ts @@ -223,7 +223,7 @@ export async function migrateLegacyRecords( try { const respHeaders = (response as any)?.headers as Record | undefined; if (respHeaders && Object.keys(respHeaders).length > 0) { - rl.updateFromHeaders(normalizeHeaders(respHeaders)); + rl.updateFromHeaders(normalizeHeaders(respHeaders), batch.length); } } catch { // ignore header parse errors @@ -260,9 +260,15 @@ export async function migrateLegacyRecords( if (i + MAX_WRITES_PER_BATCH < plan.toBackfill.length) { const cap = rl.getServerCapacity(); - if (cap) { + if (cap && rl.hasObservedPointsPerRecord()) { const actualQuota = rl.getActualRemaining(); - const pacing = pacer.calculateDelay(batch.length, cap.limit, cap.windowSeconds, actualQuota); + const pacing = pacer.calculateDelay( + batch.length, + cap.limit, + cap.windowSeconds, + actualQuota, + rl.getPointsPerRecord(POINTS_PER_RECORD) + ); await sleep(pacing.delayMs, signal); } } @@ -322,7 +328,7 @@ export async function migrateLegacyRecords( try { const respHeaders = (response as any)?.headers as Record | undefined; if (respHeaders && Object.keys(respHeaders).length > 0) { - rl.updateFromHeaders(normalizeHeaders(respHeaders)); + rl.updateFromHeaders(normalizeHeaders(respHeaders), batch.length); } } catch { // ignore header parse errors @@ -355,9 +361,15 @@ export async function migrateLegacyRecords( if (i + MAX_WRITES_PER_BATCH < toDelete.length) { const cap = rl.getServerCapacity(); - if (cap) { + if (cap && rl.hasObservedPointsPerRecord()) { const actualQuota = rl.getActualRemaining(); - const pacing = pacer.calculateDelay(batch.length, cap.limit, cap.windowSeconds, actualQuota); + const pacing = pacer.calculateDelay( + batch.length, + cap.limit, + cap.windowSeconds, + actualQuota, + rl.getPointsPerRecord(POINTS_PER_RECORD) + ); await sleep(pacing.delayMs, signal); } } diff --git a/packages/croft-click-core/src/proactive-rate-pacer.ts b/packages/croft-click-core/src/proactive-rate-pacer.ts index da56dc0..ad44d27 100644 --- a/packages/croft-click-core/src/proactive-rate-pacer.ts +++ b/packages/croft-click-core/src/proactive-rate-pacer.ts @@ -69,8 +69,15 @@ export interface PacingCalculation { * - This creates a steady-state where we never exhaust quota */ export class ProactiveRatePacer { - /** Points per record in ATProto (constant) */ - private readonly POINTS_PER_RECORD = 3; + /** + * Default points-per-record assumption, used only when the caller hasn't + * supplied an empirically observed value (see RateLimiter.getPointsPerRecord). + * The server may report capacity for a bucket with entirely different + * per-request semantics (e.g. a flat per-IP request cap rather than a + * write-cost budget) under the same header names, so this constant is a + * starting assumption to converge from, not a fact about every bucket. + */ + private readonly DEFAULT_POINTS_PER_RECORD = 3; /** Target utilization of maximum rate (80% = comfortable margin) */ private readonly TARGET_UTILIZATION = 0.80; @@ -104,17 +111,21 @@ export class ProactiveRatePacer { * @param serverLimit Total server capacity (e.g., 5000) * @param windowSeconds Window duration (e.g., 3600 = 1 hour) * @param currentRemaining Current ACTUAL quota remaining (not safe quota) + * @param pointsPerRecord Cost per record against the reported bucket — + * pass an empirically observed value when available (see + * RateLimiter.getPointsPerRecord); defaults to the write-cost assumption. * @returns Optimal delay calculation */ calculateDelay( batchSize: number, serverLimit: number, windowSeconds: number, - currentRemaining: number + currentRemaining: number, + pointsPerRecord: number = this.DEFAULT_POINTS_PER_RECORD ): PacingCalculation { // Calculate maximum sustainable rate (records per second) const pointsPerSecond = serverLimit / windowSeconds; - const maxRecordsPerSecond = pointsPerSecond / this.POINTS_PER_RECORD; + const maxRecordsPerSecond = pointsPerSecond / pointsPerRecord; // Calculate quota health (percentage of limit remaining) const quotaHealthPercent = (currentRemaining / serverLimit) * 100; @@ -184,17 +195,21 @@ export class ProactiveRatePacer { * @param windowSeconds Window duration * @param currentRemaining Current ACTUAL remaining (not safe quota) * @param maxBatchSize Hard limit (PDS max = 200) + * @param pointsPerRecord Cost per record against the reported bucket — + * pass an empirically observed value when available; defaults to the + * write-cost assumption. * @returns Recommended batch size */ calculateOptimalBatchSize( serverLimit: number, windowSeconds: number, currentRemaining: number, - maxBatchSize: number = 200 + maxBatchSize: number = 200, + pointsPerRecord: number = this.DEFAULT_POINTS_PER_RECORD ): number { // Calculate sustainable rate const pointsPerSecond = serverLimit / windowSeconds; - const maxRecordsPerSecond = pointsPerSecond / this.POINTS_PER_RECORD; + const maxRecordsPerSecond = pointsPerSecond / pointsPerRecord; // Quota health determines target rate const quotaHealthPercent = (currentRemaining / serverLimit) * 100; @@ -203,12 +218,12 @@ export class ProactiveRatePacer { // OPTIMIZED: Progressive throttling for batch size calculation if (quotaHealthPercent < 5) { // Critical: tiny batches (1-10 records) to maintain minimal progress - const criticalSize = Math.max(1, Math.min(10, Math.floor(currentRemaining / this.POINTS_PER_RECORD))); + const criticalSize = Math.max(1, Math.min(10, Math.floor(currentRemaining / pointsPerRecord))); return criticalSize; } else if (quotaHealthPercent < 15) { // Very Low: small batches (10-20 records) to gradually rebuild targetUtilization = 0.10; - const lowSize = Math.max(10, Math.min(20, Math.floor(currentRemaining / this.POINTS_PER_RECORD))); + const lowSize = Math.max(10, Math.min(20, Math.floor(currentRemaining / pointsPerRecord))); return lowSize; } else if (quotaHealthPercent < 30) { // Low: conservative batches @@ -243,17 +258,21 @@ export class ProactiveRatePacer { * @param serverLimit Server capacity * @param windowSeconds Window duration * @param currentQuota Current quota remaining + * @param pointsPerRecord Cost per record against the reported bucket — + * pass an empirically observed value when available; defaults to the + * write-cost assumption. * @returns Estimated seconds to completion */ estimateTimeToCompletion( remainingRecords: number, serverLimit: number, windowSeconds: number, - currentQuota: number + currentQuota: number, + pointsPerRecord: number = this.DEFAULT_POINTS_PER_RECORD ): number { // Calculate sustainable rate const pointsPerSecond = serverLimit / windowSeconds; - const maxRecordsPerSecond = pointsPerSecond / this.POINTS_PER_RECORD; + const maxRecordsPerSecond = pointsPerSecond / pointsPerRecord; // Use current quota health to estimate average utilization const quotaHealthPercent = (currentQuota / serverLimit) * 100; diff --git a/packages/croft-click-core/src/publisher.ts b/packages/croft-click-core/src/publisher.ts index 40ab1e4..647369f 100644 --- a/packages/croft-click-core/src/publisher.ts +++ b/packages/croft-click-core/src/publisher.ts @@ -98,7 +98,8 @@ export async function publishRecords( serverCapacity.limit, serverCapacity.windowSeconds, actualRemaining, - MAX_PDS_BATCH_SIZE + MAX_PDS_BATCH_SIZE, + rl.getPointsPerRecord(POINTS_PER_RECORD) ); currentDelay = 500; onLog('info', `Using saved server info: ${serverCapacity.limit} pts/${serverCapacity.windowSeconds}s`); @@ -135,7 +136,8 @@ export async function publishRecords( capacity.limit, capacity.windowSeconds, actualRemaining, - MAX_PDS_BATCH_SIZE + MAX_PDS_BATCH_SIZE, + rl.getPointsPerRecord(POINTS_PER_RECORD) ); // Apply adaptive scaling from performance metrics @@ -222,7 +224,7 @@ export async function publishRecords( const rawHeaders = extractHeaders(response); if (Object.keys(rawHeaders).length > 0) { const norm = normalizeHeaders(rawHeaders); - rl.updateFromHeaders(norm); + rl.updateFromHeaders(norm, batch.length); // After first response, log the learned capacity and recalculate if (batchCounter === 1 && rl.hasServerInfo()) { @@ -233,7 +235,8 @@ export async function publishRecords( cap.limit, cap.windowSeconds, remaining, - MAX_PDS_BATCH_SIZE + MAX_PDS_BATCH_SIZE, + rl.getPointsPerRecord(POINTS_PER_RECORD) ); const quotaPercent = ((remaining / cap.limit) * 100).toFixed(1); @@ -253,13 +256,14 @@ export async function publishRecords( // PROACTIVE PACING: Calculate optimal delay for next batch if (i < total) { const cap = rl.getServerCapacity(); - if (cap) { + if (cap && rl.hasObservedPointsPerRecord()) { const actualQuota = rl.getActualRemaining(); const pacing = pacer.calculateDelay( currentBatchSize, cap.limit, cap.windowSeconds, - actualQuota + actualQuota, + rl.getPointsPerRecord(POINTS_PER_RECORD) ); // Update delay if changed significantly diff --git a/packages/croft-click-core/src/rate-limiter.ts b/packages/croft-click-core/src/rate-limiter.ts index 982261e..6666d9e 100644 --- a/packages/croft-click-core/src/rate-limiter.ts +++ b/packages/croft-click-core/src/rate-limiter.ts @@ -16,11 +16,32 @@ export class RateLimiter { private state: State | null = null; private readonly headroom: number; + /** + * Empirically observed points consumed per record, derived from successive + * quota readings rather than assumed. The server reports whichever bucket + * is currently tightest (e.g. a broad per-IP request cap vs. a write-specific + * points budget) under the same generic header names, with no way to tell + * which one it is from the headers alone. A hardcoded "points per record" + * constant is only correct when the reported bucket happens to be the + * write-specific one; against a request-counted bucket it wildly + * overestimates cost per record. Tracking the observed delta instead + * self-corrects regardless of which bucket is being reported. + */ + private observedPointsPerRecord: number | null = null; + private lastRemaining: number | null = null; + private lastResetAt: number | null = null; + constructor(opts?: { headroom?: number }) { this.headroom = opts?.headroom ?? 0.15; } - updateFromHeaders(headers: Record): void { + /** + * @param batchSize Number of records the just-completed request covered. + * Pass this to let the limiter learn the real points-per-record cost of + * whichever bucket the server is reporting; omit for reads/other calls + * that shouldn't feed the estimate. + */ + updateFromHeaders(headers: Record, batchSize?: number): void { const h = normalizeHeaders(headers); const get = (k: string) => h[k] ?? h[`x-${k}`] ?? ''; @@ -36,12 +57,56 @@ export class RateLimiter { if (m) windowSeconds = parseInt(m[1], 10); const now = Math.floor(Date.now() / 1000); - this.state = { - limit, - remaining, - resetAt: isNaN(reset) ? now + windowSeconds : reset, - windowSeconds, - }; + const resetAt = isNaN(reset) ? now + windowSeconds : reset; + + // Only learn from this reading if it's a same-window continuation of the + // previous one — a window reset (or first-ever reading) makes the delta + // meaningless, and a jump to a differently-shaped bucket (different + // limit/window) resets our sample rather than corrupting the average. + if ( + batchSize && + batchSize > 0 && + this.state && + this.state.limit === limit && + this.state.windowSeconds === windowSeconds && + this.lastRemaining !== null && + this.lastResetAt === resetAt + ) { + const consumed = this.lastRemaining - remaining; + if (consumed > 0) { + const observed = consumed / batchSize; + this.observedPointsPerRecord = + this.observedPointsPerRecord === null + ? observed + : this.observedPointsPerRecord * 0.5 + observed * 0.5; + } + } + + this.lastRemaining = remaining; + this.lastResetAt = resetAt; + this.state = { limit, remaining, resetAt, windowSeconds }; + } + + /** + * Points-per-record cost to use for pacing math: the empirically observed + * value once we have one, otherwise the caller's assumed default. + */ + getPointsPerRecord(fallback: number): number { + return this.observedPointsPerRecord ?? fallback; + } + + /** + * Whether an empirical points-per-record estimate exists yet. Before the + * first same-window reading pair, any capacity/window the server reports + * could belong to a bucket with completely different per-request cost + * semantics than the caller's assumed default — proactively pacing against + * it is as likely to be needlessly conservative as it is accurate. Callers + * should skip proactive delay until this is true and rely on the reactive + * 429 path (`handleRateLimitHit` + `waitForPermit`) in the meantime, which + * is safe regardless of which bucket is actually being reported. + */ + hasObservedPointsPerRecord(): boolean { + return this.observedPointsPerRecord !== null; } getActualRemaining(): number { diff --git a/packages/malachite/src/lib/polish.ts b/packages/malachite/src/lib/polish.ts index f72d0a1..814ca76 100644 --- a/packages/malachite/src/lib/polish.ts +++ b/packages/malachite/src/lib/polish.ts @@ -114,6 +114,22 @@ export async function migrateLegacyRecords( log.section('Backfilling Production Records'); } + // cli-progress renders nothing meaningful when stdout isn't a TTY (CI, + // piped/redirected output, --non-interactive runs) — fall back to explicit, + // throttled log lines so progress is actually visible in that context. + const isTTY = Boolean(process.stdout.isTTY); + let lastLogAt = 0; + const logProgressLine = (phase: 'backfill' | 'delete', done: number, total: number) => { + if (isTTY || total === 0) return; + const now = Date.now(); + if (done < total && now - lastLogAt < 5000) return; // throttle to ~5s; always emit the final line + lastLogAt = now; + const pct = ((done / total) * 100).toFixed(1); + const elapsed = ((now - start) / 1000).toFixed(0); + const label = phase === 'backfill' ? 'Backfilling' : 'Removing legacy'; + log.info(`${label}: ${done.toLocaleString()}/${total.toLocaleString()} (${pct}%) — ${elapsed}s elapsed`); + }; + const result = await coreMigrate(agent, plan, { dryRun, onProgress: (phase, done, total) => { @@ -126,6 +142,7 @@ export async function migrateLegacyRecords( } bars.delete.update(done, {}); } + logProgressLine(phase, done, total); }, }); -- 2.51.2