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); }, });