diff --git a/CLAUDE.md b/CLAUDE.md index 1bb3d94..6a3ee22 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -91,6 +91,12 @@ regression is pinned). Extend it, don't prune it. - worker tests: `bun test`. zig tests: `just test-services` (or `cd services && zig build test`). - never run work proportional to the corpus on the request/cron path (see the invariant above). +- **the cron contract**: every cron selector must be an index SEARCH (smoke's + `test_cron_selector_plans` EXPLAIN-guards this live against Turso), every batch size must + state what it honestly covers against measured arrival rates (~2,900 new actors/hour), and + the cron NEVER claims to drain a corpus-scale backlog — backlogs belong to the bulk paths + (`scripts/bulk-enrich.py` for profiles, `scripts/plc-identity-sync.py` for identity), run + off-worker. A cron sized within ~2× of arrivals is treading water, not draining. - every equality query on `actors.handle` needs `COLLATE NOCASE`, or it scans the whole table — `idx_actors_handle` is declared with that collation. Paginate on `rowid`. - Turso is single-writer and the live ingester shares it. Bulk writes must batch (one diff --git a/scripts/smoke.py b/scripts/smoke.py index 14a727e..d71e62d 100755 --- a/scripts/smoke.py +++ b/scripts/smoke.py @@ -528,6 +528,61 @@ def test_sync_keyset_plan(): plan) +def test_cron_selector_plans(): + """THE CRON CONTRACT: every cron selector must plan as an index SEARCH on + Turso, never a corpus scan. The moderation-refresh selector used to run two + full-table ORDER BY scans over ~11M rows every hour; these EXPLAINs pin its + replacement (quality-keyset walk + firehose event queue, src/cron.ts) to + the indexes that serve them. Like test_sync_keyset_plan, only a live + EXPLAIN against Turso is trustworthy — local SQLite plans differently.""" + print("\n--- turso cron selector plans ---") + turso_url = os.environ.get("TURSO_URL", "") + token = os.environ.get("TURSO_AUTH_TOKEN", "") + if not turso_url or not token: + print(f" [{SKIP}] skipped (set TURSO_URL + TURSO_AUTH_TOKEN)") + return + if turso_url.startswith("libsql://"): + turso_url = "https://" + turso_url.removeprefix("libsql://") + + selectors = { + # quality walk: keyset literals inlined exactly as src/cron.ts builds them + "quality walk uses idx_actors_visible_quality_score": ( + "EXPLAIN QUERY PLAN SELECT did FROM actors " + "WHERE handle != '' AND hidden = 0 AND quality_score > 0 " + "AND (quality_score < 123.45 OR (quality_score = 123.45 AND rowid < 678)) " + "ORDER BY quality_score DESC, rowid DESC LIMIT 800" + ), + "event queue uses idx_actors_refresh_queue": ( + "EXPLAIN QUERY PLAN SELECT did FROM actors " + "WHERE profile_checked_at = 0 AND handle != '' " + "ORDER BY updated_at DESC LIMIT 200" + ), + } + for label, sql in selectors.items(): + body = json.dumps({ + "requests": [ + {"type": "execute", "stmt": {"sql": sql}}, + {"type": "close"}, + ] + }).encode() + try: + req = urllib.request.Request( + f"{turso_url}/v2/pipeline", data=body, method="POST", + headers={"Authorization": f"Bearer {token}", "Content-Type": "application/json"}, + ) + with urllib.request.urlopen(req, timeout=15) as resp: + result = json.loads(resp.read()) + rows = result["results"][0]["response"]["result"]["rows"] + details = [r[3]["value"] for r in rows] + except Exception as e: + check(f"{label} (plan fetch)", False, str(e)) + continue + plan = "; ".join(details) + check(label, + any("SEARCH" in d for d in details) and not any("SCAN" in d for d in details), + plan) + + def main(): parser = argparse.ArgumentParser(description="typeahead smoke tests") parser.add_argument("--url", required=True, help="typeahead service URL") @@ -567,6 +622,7 @@ def main(): test_page_scripts(args.url) test_sync_keyset_plan() + test_cron_selector_plans() if args.compare: test_comparison(args.url, args.queries) diff --git a/src/cron.test.ts b/src/cron.test.ts index 1aa3d0b..43009e6 100644 --- a/src/cron.test.ts +++ b/src/cron.test.ts @@ -99,8 +99,8 @@ describe("refreshModeration advances the freshness cursor", () => { bind(...args: any[]) { s.binds = args; return s; }, async all() { calls.push({ sql, binds: s.binds, method: "all" }); - if (sql.includes("ORDER BY (unixepoch()")) return { results: [row] }; // priority pool - return { results: [] }; // floor pool, overrides, domains + if (sql.includes("ORDER BY quality_score DESC")) return { results: [row] }; // quality walk + return { results: [] }; // event queue, overrides, domains }, async run() { calls.push({ sql, binds: s.binds, method: "run" }); return { meta: { changes: 1 } }; }, async first() { return null; }, @@ -119,6 +119,18 @@ describe("refreshModeration advances the freshness cursor", () => { }; } + // KV mock capturing cursor movement + function mockEnv(cursor: string | null = null) { + const kv = { + puts: [] as { key: string; value: string }[], + deletes: [] as string[], + async get(_key: string) { return cursor; }, + async put(key: string, value: string) { kv.puts.push({ key, value }); }, + async delete(key: string) { kv.deletes.push(key); }, + }; + return { env: { KV: kv } as any, kv }; + } + function mockGetProfiles(profile: any) { globalThis.fetch = (async () => ({ status: 200, ok: true, json: async () => ({ profiles: [profile] }), @@ -148,7 +160,7 @@ describe("refreshModeration advances the freshness cursor", () => { labels: [], createdAt: "2024-01-01T00:00:00.000Z", associated: {}, followersCount: 3041, followsCount: 100, postsCount: 500, }); - await refreshModeration(db, {} as any); + await refreshModeration(db, mockEnv().env); expect(cursorAdvances(calls)).toBe(true); expect(materialUpdateAdvancesReplica(calls)).toBe(true); }); @@ -165,10 +177,79 @@ describe("refreshModeration advances the freshness cursor", () => { const row = staleRow(f.qualityScore); // align score so identical holds const { db, calls } = mockDb(row); mockGetProfiles(profile); - await refreshModeration(db, {} as any); + await refreshModeration(db, mockEnv().env); // even though nothing changed, we must record that we checked it — else the - // priority pool re-selects this same row forever and never walks onward. + // event queue re-selects this same row forever and never drains. expect(cursorAdvances(calls)).toBe(true); expect(materialUpdateAdvancesReplica(calls)).toBe(false); }); }); + +// THE CRON CONTRACT: selection must be O(batch) and index-served — never a +// full-corpus scan or sort. The previous selector ran two full-table ORDER BY +// scans over ~11M rows every hour; these tests pin the replacement's shape. +describe("refreshModeration quality-walk cursor", () => { + const origFetch = globalThis.fetch; + afterEach(() => { globalThis.fetch = origFetch; }); + + function scanTrap() { + // db that returns empty pools and captures the walk SQL + const sqls: string[] = []; + const stmt = (sql: string) => { + const s: any = { + sql, + bind() { return s; }, + async all() { sqls.push(sql); return { results: [] }; }, + async run() { return { meta: { changes: 0 } }; }, + }; + return s; + }; + return { db: { prepare: stmt, async batch() { return []; } } as any, sqls }; + } + + test("keyset cursor is inlined as literals, not bound params (Turso plan degradation)", async () => { + const { db, sqls } = scanTrap(); + const { env } = mockEnvWithCursor(JSON.stringify({ qs: 123.45, rowid: 678 })); + await refreshModeration(db, env); + const walk = sqls.find((s) => s.includes("ORDER BY quality_score DESC"))!; + expect(walk).toContain("quality_score < 123.45"); + expect(walk).toContain("rowid < 678"); + expect(walk).not.toContain("?2"); // keyset values never arrive as params + }); + + test("a poisoned cursor cannot inject SQL — non-numeric values are dropped", async () => { + const { db, sqls } = scanTrap(); + const { env } = mockEnvWithCursor(JSON.stringify({ qs: "1; DROP TABLE actors", rowid: 678 })); + await refreshModeration(db, env); + const walk = sqls.find((s) => s.includes("ORDER BY quality_score DESC"))!; + expect(walk).not.toContain("DROP"); + expect(walk).not.toContain("quality_score <"); // no keyset at all — walk restarts from top + }); + + test("short walk page wraps the cursor (pool exhausted)", async () => { + const { db } = scanTrap(); + const kvOps: string[] = []; + const env = { + KV: { + async get() { return JSON.stringify({ qs: 0.01, rowid: 1 }); }, + async put(key: string) { kvOps.push(`put:${key}`); }, + async delete(key: string) { kvOps.push(`delete:${key}`); }, + }, + } as any; + await refreshModeration(db, env); + expect(kvOps).toContain("delete:mod_walk_cursor"); + expect(kvOps).not.toContain("put:mod_walk_cursor"); + }); + + function mockEnvWithCursor(cursor: string) { + return { + env: { + KV: { + async get() { return cursor; }, + async put() {}, + async delete() {}, + }, + } as any, + }; + } +}); diff --git a/src/cron.ts b/src/cron.ts index 7fd57b0..8d3a146 100644 --- a/src/cron.ts +++ b/src/cron.ts @@ -34,51 +34,85 @@ export function isProfileIdentical( ); } -/** refresh moderation labels, walking the full index over multiple cron runs */ +/** refresh moderation labels and profile aggregates for a bounded batch. + * + * THE CRON CONTRACT (see CLAUDE.md): every selector here must be an index + * SEARCH, never a corpus scan, and every batch size must state what it can + * honestly cover. The previous two-pool selector broke both rules: it ran + * two full-table ORDER BY scans over ~11M rows every hour to pick 1,000, and + * its doc claimed it was "walking the full index over multiple cron runs" — + * at 24k actors/day against an 11M corpus growing 70k/day, that walk was + * never going to happen. What it could honestly do is what this now says + * it does: + * + * QUALITY WALK (800/run): a keyset cursor (KV `mod_walk_cursor`) stepping + * down idx_actors_visible_quality_score. The scored-visible pool — every + * actor a search can actually surface with ranking weight — measured + * 41,355 rows on 2026-08-16, so a full sweep completes in ~52 runs: + * every rankable actor's labels/avatar/counts re-checked every ~2-3 days, + * deterministically, with the sweep position observable in KV. The keyset + * values are inlined as numeric literals, NOT bound params — the OR-keyset + * shape degrades to a full index scan on Turso when the values arrive as + * hrana params (the Aug 2026 staleness incident; smoke EXPLAIN-guards it). + * + * EVENT QUEUE (200/run): actors the firehose re-touched — handlers/ingest.ts + * resets profile_checked_at=0 on profile/identity events so appview + * aggregates get re-fetched. Served newest-activity-first by the partial + * index idx_actors_refresh_queue. The same index also holds the historical + * never-checked backlog (profile_checked_at=0 since insert); reading DESC + * means live activity always wins, and draining that backlog is explicitly + * NOT this job — it belongs to scripts/bulk-enrich.py, run off-worker. + * + * Score-0 / hidden actors outside both queues are reached on firehose + * activity (event queue) or by bulk runs — the cron makes no pretense of + * touching all 11M rows on any schedule. */ export async function refreshModeration(db: TursoDB, env: Env): Promise { - // Two-pool selector — replaces the old rowid cursor that took ~245 days to - // walk the corpus and treated viral handles the same as 2022-dormant ones. - // - // PRIORITY POOL (800/run): ORDER BY staleness × value × activity DESC. - // - staleness: now - profile_checked_at (NULL → very stale) - // - value: log(1 + quality_score), so popular accounts refresh more - // - activity: recent firehose touch via last_activity_at boosts further - // COLD FLOOR (200/run): ORDER BY profile_checked_at ASC NULLS FIRST. - // Guarantees the long tail gets touched without starving on score=0 - // accounts that never get a priority bump. Also picks up actors whose - // firehose enqueue set profile_checked_at=0 (see handlers/ingest.ts). - // - // Per Turso Scaler (100B reads/mo): one full scan per tick ~ 6M reads × 24/day - // × 30 = 4.3B/mo ≈ 4% of budget. Will add a functional index later. - const PRIORITY_N = 800; - const FLOOR_N = 200; - - const priorityRows = (await db.prepare( - `SELECT rowid, did, handle, hidden, labels, pds, + const WALK_N = 800; + const EVENT_N = 200; + const ACTOR_COLS = `rowid, did, handle, hidden, labels, pds, display_name, avatar_url, created_at, associated, - followers_count, follows_count, posts_count, quality_score - FROM actors - ORDER BY (unixepoch() - COALESCE(profile_checked_at, 0)) - * (1.0 + ln(1.0 + COALESCE(quality_score, 0))) - * (CASE WHEN last_activity_at > unixepoch() - 86400 THEN 3.0 - WHEN last_activity_at > unixepoch() - 604800 THEN 1.5 - ELSE 1.0 END) - DESC - LIMIT ?1` - ).bind(PRIORITY_N).all()).results ?? []; + followers_count, follows_count, posts_count, quality_score`; - const floorRows = (await db.prepare( - `SELECT rowid, did, handle, hidden, labels, pds, - display_name, avatar_url, created_at, associated, - followers_count, follows_count, posts_count, quality_score - FROM actors - ORDER BY COALESCE(profile_checked_at, 0) ASC + let cursor: { qs: number; rowid: number } | null = null; + try { + cursor = JSON.parse((await env.KV.get("mod_walk_cursor")) || "null"); + } catch {} + const qs = Number(cursor?.qs); + const rid = Number(cursor?.rowid); + const keyset = Number.isFinite(qs) && Number.isFinite(rid) + ? `AND (quality_score < ${qs} OR (quality_score = ${qs} AND rowid < ${rid}))` + : ""; + // predicates must match idx_actors_visible_quality_score's WHERE clause + // exactly or the partial index is unusable and this reverts to a scan. + const walkRows = (await db.prepare( + `SELECT ${ACTOR_COLS} FROM actors + WHERE handle != '' AND hidden = 0 AND quality_score > 0 ${keyset} + ORDER BY quality_score DESC, rowid DESC + LIMIT ${WALK_N}` + ).all()).results ?? []; + + if (walkRows.length < WALK_N) { + // pool exhausted — wrap to the top next run + await env.KV.delete("mod_walk_cursor").catch(() => {}); + } else { + const last = walkRows[walkRows.length - 1]; + await env.KV.put( + "mod_walk_cursor", + JSON.stringify({ qs: last.quality_score, rowid: last.rowid }) + ).catch(() => {}); + } + + const eventRows = (await db.prepare( + `SELECT ${ACTOR_COLS} FROM actors + WHERE profile_checked_at = 0 AND handle != '' + ORDER BY updated_at DESC LIMIT ?1` - ).bind(FLOOR_N).all()).results ?? []; + ).bind(EVENT_N).all()).results ?? []; - // dedup: priority pool wins (cold floor is a backstop, not a duplicate-fetcher). - const seen = new Set(priorityRows.map((r) => r.did)); - const results = [...priorityRows, ...floorRows.filter((r) => !seen.has(r.did))]; + // dedup: event queue wins — a firehose-touched actor is fresher signal than + // its scheduled sweep slot, and the walk cursor has already moved past it. + const seen = new Set(eventRows.map((r) => r.did)); + const results = [...eventRows, ...walkRows.filter((r) => !seen.has(r.did))]; if (results.length === 0) { console.log(JSON.stringify({ event: "moderation_refresh", status: "empty" })); @@ -244,7 +278,8 @@ export async function refreshModeration(db: TursoDB, env: Env): Promise { console.log(JSON.stringify({ event: "moderation_refresh", checked, changed, skipped, deleted, - priority: priorityRows.length, floor: floorRows.length, + walk: walkRows.length, events: eventRows.length, + walk_cursor_qs: walkRows.length === WALK_N ? walkRows[walkRows.length - 1].quality_score : null, ms: Date.now() - t0, })); } diff --git a/src/enrichment.ts b/src/enrichment.ts index 55af2c1..e593455 100644 --- a/src/enrichment.ts +++ b/src/enrichment.ts @@ -234,6 +234,13 @@ export async function enrichActors(db: TursoDB, env: Env): Promise<{ resolved: n // profile-checked, which is what a 61.9%-and-flat avatar coverage and a // near-zero hidden count on /stats actually measured. avatar coverage // among CHECKED actors is ~95%. + // + // Draining that multi-million backlog is NOT this job. At 5,000/run the + // cron nets ~2,100/hour over arrivals — months to converge, and only if + // every run completes (a single 429 forfeits the rest of the run). The + // backlog path is scripts/bulk-enrich.py, run off-worker; this cron's + // honest contract is keeping pace with arrivals so the backlog never + // grows again. const { results: profileRows } = await db.prepare( `SELECT did, handle, pds FROM actors WHERE handle != ''