From f229b40d59d1c8a5b88631e9f5be482b6ec17a59 Mon Sep 17 00:00:00 2001 From: Hunter B Date: Tue, 29 Sep 2026 04:18:43 -0700 Subject: [PATCH 1/8] fix(web): gate the public digest on maintainer approval; keep resolved drafts resolved - runDigest stages the structured weekly record with approved:false and picks its fields explicitly; /digest renders only approved, well-formed records, so the "maintainer-approved" copy on the page and in llms.txt is accurate. Posting the digest from /admin approves it (unedited text only); discarding deletes the staged record. - Posting or discarding writes a draft-resolved marker (90 days). saveDraft and hasFreshDraft honor it and never overwrite a posted draft, so the triage / PR-review / stale / dupes / digest / content-watch runs no longer bring a resolved draft back as pending or spend a model call on it. - /api/admin/post rejects already-posted drafts (409), claims the draft before the GitHub call, returns ok with a warning when bookkeeping fails after GitHub accepted the post, reads the body with readBoundedBody, and answers malformed JSON with a JSON 400. - listDrafts follows the KV list cursor instead of stopping at 100 keys. Tests: web vitest 57 files / 515 passed; new lib/community-agent-review-state.test.ts 14 passed (13 fail on the unfixed source); tsc --noEmit exit 0; eslint on touched files exit 0. Co-Authored-By: Claude Opus 5.5 (1M context) Claude-Session: https://claude.ai/code/session_014ZwqatxgVFxHvovngywnks --- web/app/[locale]/digest/page.tsx | 29 +- web/app/api/admin/post/route.ts | 194 ++++++--- web/lib/community-agent-review-state.test.ts | 391 +++++++++++++++++++ web/lib/community-agent-tasks.ts | 50 ++- web/lib/community-agent.ts | 174 ++++++++- web/lib/public-api-security.test.ts | 14 +- 6 files changed, 746 insertions(+), 106 deletions(-) create mode 100644 web/lib/community-agent-review-state.test.ts diff --git a/web/app/[locale]/digest/page.tsx b/web/app/[locale]/digest/page.tsx index 69c63c989d..c5a658d6ed 100644 --- a/web/app/[locale]/digest/page.tsx +++ b/web/app/[locale]/digest/page.tsx @@ -1,24 +1,15 @@ import { PageHeader } from "@/components/page-header"; import { EmptyState } from "@/components/surface-state"; +import { + DIGEST_RECORD_PREFIX, + isPublishedDigest, + type WeeklyDigestRecord as WeeklyDigest, +} from "@/lib/community-agent"; import { getDigest, pickTextLocale } from "@/lib/i18n/dictionaries"; import { getEnv } from "@/lib/kv"; import { buildPageMetadata } from "@/lib/page-meta"; -// Define the exact structure of the Digest data to fix all the 'any' type errors -interface DigestSection { - heading: string; - items: string[]; -} - -interface WeeklyDigest { - weekId: string; - titleEn: string; - titleZh: string; - summaryEn: string; - summaryZh: string; - sections: DigestSection[]; - generatedAt: string; -} +type DigestSection = WeeklyDigest["sections"][number]; export const revalidate = 3600; // Cache page updates hourly @@ -47,7 +38,7 @@ export default async function DigestArchivePage({ params }: { params: Promise<{ if (kv) { try { // Fetch all weekly digest keys generated by the agent tasks - const { keys } = await kv.list({ prefix: "digest:weekly-" }); + const { keys } = await kv.list({ prefix: DIGEST_RECORD_PREFIX }); if (keys && keys.length > 0) { const digestsRaw = await Promise.all( @@ -58,12 +49,14 @@ export default async function DigestArchivePage({ params }: { params: Promise<{ ); // Parse each entry independently so a single malformed record can't - // blank the entire archive — skip only the bad ones. + // blank the entire archive — skip only the bad ones. The cron stages + // records unapproved; only maintainer-approved ones are public. digests = digestsRaw .filter((item: string | null): item is string => Boolean(item)) .flatMap((item: string) => { try { - return [JSON.parse(item) as WeeklyDigest]; + const parsed: unknown = JSON.parse(item); + return isPublishedDigest(parsed) ? [parsed] : []; } catch (e) { console.error("Skipping malformed digest entry:", e); return []; diff --git a/web/app/api/admin/post/route.ts b/web/app/api/admin/post/route.ts index 2b917e442d..76f6e3fad5 100644 --- a/web/app/api/admin/post/route.ts +++ b/web/app/api/admin/post/route.ts @@ -1,8 +1,14 @@ import { NextResponse } from "next/server"; +import { BodyReadError, readBoundedBody } from "@/lib/bounded-body"; import { + approveDigestRecord, + clearDraftResolution, + deleteDigestRecord, deleteDraft, getAgentEnv, getDraft, + getDraftResolution, + markDraftResolved, parseDraftKey, validateSession, type CommunityAgentEnv, @@ -51,24 +57,42 @@ export async function POST(req: Request) { ); } - const contentLength = Number(req.headers.get("content-length") ?? "0"); - if (contentLength > MAX_BODY_BYTES) { - return NextResponse.json({ error: "payload too large" }, { status: 413 }); + // Count the real body bytes (Content-Length is only an early rejection) + // and answer malformed JSON with a 400 instead of an unhandled 500. + let body: unknown; + try { + const bytes = await readBoundedBody(req, MAX_BODY_BYTES); + body = JSON.parse(new TextDecoder().decode(bytes)); + } catch (e) { + if (e instanceof BodyReadError) { + return NextResponse.json({ error: e.message }, { status: e.status }); + } + return NextResponse.json({ error: "invalid JSON body" }, { status: 400 }); } - - const body = await req.json() as { action: string; draftKey: string; editedBody?: string; lang?: "en" | "zh" }; - const { action, draftKey, editedBody, lang } = body; - - if (!ALLOWED_ACTIONS.has(action)) { + if (!body || typeof body !== "object" || Array.isArray(body)) { + return NextResponse.json({ error: "invalid JSON body" }, { status: 400 }); + } + const { action, draftKey, editedBody, lang } = body as { + action?: unknown; + draftKey?: unknown; + editedBody?: unknown; + lang?: unknown; + }; + + if (typeof action !== "string" || !ALLOWED_ACTIONS.has(action)) { return NextResponse.json({ error: "unknown action" }, { status: 400 }); } if (typeof draftKey !== "string" || !draftKey || draftKey.length > 256) { return NextResponse.json({ error: "missing or invalid draftKey" }, { status: 400 }); } - if (!parseDraftKey(draftKey)) { + const parsedKey = parseDraftKey(draftKey); + if (!parsedKey) { return NextResponse.json({ error: "invalid draftKey namespace" }, { status: 400 }); } - if (editedBody !== undefined && (typeof editedBody !== "string" || editedBody.length > MAX_BODY_BYTES)) { + if (editedBody !== undefined && typeof editedBody !== "string") { + return NextResponse.json({ error: "invalid editedBody" }, { status: 400 }); + } + if (editedBody !== undefined && editedBody.length > MAX_BODY_BYTES) { return NextResponse.json({ error: "editedBody too long" }, { status: 413 }); } if (lang !== undefined && lang !== "en" && lang !== "zh") { @@ -81,7 +105,16 @@ export async function POST(req: Request) { } if (action === "discard") { - await deleteDraft(env.CURATED_KV, draftKey); + try { + // The marker stops the next cron run from regenerating this draft. + await markDraftResolved(env.CURATED_KV, parsedKey.type, parsedKey.id, "discarded"); + await deleteDraft(env.CURATED_KV, draftKey); + if (draft.type === "digest") { + await deleteDigestRecord(env.CURATED_KV, draft.id); + } + } catch (e) { + return NextResponse.json({ error: `discard failed: ${String(e)}` }, { status: 500 }); + } return NextResponse.json({ ok: true, action: "discarded" }); } @@ -89,8 +122,72 @@ export async function POST(req: Request) { if (!env.MAINTAINER_GITHUB_PAT) { return NextResponse.json({ error: "MAINTAINER_GITHUB_PAT not configured" }, { status: 500 }); } + if (draft.type !== "digest" && !draft.targetNumber) { + return NextResponse.json({ error: "no target number" }, { status: 400 }); + } + + // Posting is not idempotent on GitHub: refuse a draft that is already + // posted or has a post in flight (second tab, retry after a partial + // failure), then claim it before the GitHub call. + if (draft.posted) { + return NextResponse.json({ error: "draft already posted" }, { status: 409 }); + } + const resolution = await getDraftResolution(env.CURATED_KV, parsedKey.type, parsedKey.id); + if (resolution) { + return NextResponse.json({ error: `draft already ${resolution.state}` }, { status: 409 }); + } + try { + await markDraftResolved(env.CURATED_KV, parsedKey.type, parsedKey.id, "posting"); + } catch (e) { + return NextResponse.json({ error: `could not claim draft: ${String(e)}` }, { status: 500 }); + } - const commentBody = editedBody ?? (lang === "zh" ? draft.bodyZh : draft.bodyEn); + const originalBody = lang === "zh" ? draft.bodyZh : draft.bodyEn; + const commentBody = editedBody ?? originalBody; + + // After GitHub accepted the post, bookkeeping failures must not turn into + // an error the maintainer would "fix" by posting again. + const recordPosted = async (): Promise => { + try { + await markDraftResolved(env.CURATED_KV, parsedKey.type, parsedKey.id, "posted"); + await env.CURATED_KV?.put(draftKey, JSON.stringify(draft), { expirationTtl: 60 * 60 * 24 * 7 }); + return undefined; + } catch (e) { + return `posted, but saving draft state failed: ${String(e)}`; + } + }; + + const postToGitHub = async (url: string, payload: unknown, authScheme: "token" | "Bearer") => { + try { + return await fetch(url, { + method: "POST", + headers: { + Accept: "application/vnd.github+json", + Authorization: `${authScheme} ${env.MAINTAINER_GITHUB_PAT}`, + "X-GitHub-Api-Version": "2022-11-28", + "Content-Type": "application/json", + }, + body: JSON.stringify(payload), + }); + } catch { + // Outcome unknown: keep the short-lived claim so an immediate retry + // cannot double-post. + return null; + } + }; + const unknownOutcome = () => + NextResponse.json( + { error: "GitHub request failed; check GitHub before retrying (retry unlocks in 15 minutes)" }, + { status: 502 } + ); + const githubFailed = async (res: Response) => { + // GitHub definitively rejected the post, so release the claim. + try { + await clearDraftResolution(env.CURATED_KV, parsedKey.type, parsedKey.id); + } catch { /* the claim expires on its own */ } + const text = await res.text(); + return NextResponse.json({ error: `GitHub ${res.status}: ${text}` }, { status: 502 }); + }; if (draft.type === "digest") { const digestBody = commentBody; @@ -100,60 +197,51 @@ export async function POST(req: Request) { const digestRepo = env.GITHUB_REPO ?? "Hmbown/CodeWhale"; const issuesUrl = `https://api.github.com/repos/${digestRepo}/issues`; - const digestRes = await fetch(issuesUrl, { - method: "POST", - headers: { - Accept: "application/vnd.github+json", - Authorization: `token ${env.MAINTAINER_GITHUB_PAT}`, - "X-GitHub-Api-Version": "2022-11-28", - "Content-Type": "application/json", - }, - body: JSON.stringify({ title, body: digestBody, labels: ["digest"] }), - }); - - if (!digestRes.ok) { - const text = await digestRes.text(); - return NextResponse.json({ error: `GitHub ${digestRes.status}: ${text}` }, { status: 502 }); - } + const digestRes = await postToGitHub(issuesUrl, { title, body: digestBody, labels: ["digest"] }, "token"); + if (!digestRes) return unknownOutcome(); + if (!digestRes.ok) return githubFailed(digestRes); - const issue = await digestRes.json() as { number: number; html_url: string }; + const issue = await digestRes.json().catch(() => ({})) as { number?: number; html_url?: string }; draft.posted = true; - draft.targetNumber = issue.number; - draft.targetUrl = issue.html_url; - await env.CURATED_KV?.put(draftKey, JSON.stringify(draft), { expirationTtl: 60 * 60 * 24 * 7 }); - - return NextResponse.json({ ok: true, action: "posted", number: issue.number, url: issue.html_url }); - } + if (typeof issue.number === "number") draft.targetNumber = issue.number; + if (typeof issue.html_url === "string") draft.targetUrl = issue.html_url; + let warning = await recordPosted(); + + // Publishing to /digest is the approval. The structured record holds + // the unedited model text, so an edited digest is posted to GitHub but + // not published there. + let published = false; + if (editedBody === undefined || editedBody === originalBody) { + try { + published = await approveDigestRecord(env.CURATED_KV, draft.id); + } catch (e) { + warning ??= `posted, but publishing the digest page failed: ${String(e)}`; + } + } - if (!draft.targetNumber) { - return NextResponse.json({ error: "no target number" }, { status: 400 }); + return NextResponse.json({ + ok: true, + action: "posted", + number: issue.number, + url: issue.html_url, + published, + ...(warning ? { warning } : {}), + }); } const repo = env.GITHUB_REPO ?? "Hmbown/CodeWhale"; const commentUrl = `https://api.github.com/repos/${repo}/issues/${draft.targetNumber}/comments`; - const ghRes = await fetch(commentUrl, { - method: "POST", - headers: { - Accept: "application/vnd.github+json", - Authorization: `Bearer ${env.MAINTAINER_GITHUB_PAT}`, - "X-GitHub-Api-Version": "2022-11-28", - "Content-Type": "application/json", - }, - body: JSON.stringify({ body: commentBody }), - }); - - if (!ghRes.ok) { - const text = await ghRes.text(); - return NextResponse.json({ error: `GitHub ${ghRes.status}: ${text}` }, { status: 502 }); - } + const ghRes = await postToGitHub(commentUrl, { body: commentBody }, "Bearer"); + if (!ghRes) return unknownOutcome(); + if (!ghRes.ok) return githubFailed(ghRes); // Mark as posted draft.posted = true; - await env.CURATED_KV?.put(draftKey, JSON.stringify(draft), { expirationTtl: 60 * 60 * 24 * 7 }); + const warning = await recordPosted(); - return NextResponse.json({ ok: true, action: "posted" }); + return NextResponse.json({ ok: true, action: "posted", ...(warning ? { warning } : {}) }); } // ALLOWED_ACTIONS guard above means this is unreachable. diff --git a/web/lib/community-agent-review-state.test.ts b/web/lib/community-agent-review-state.test.ts new file mode 100644 index 0000000000..60ed6c310d --- /dev/null +++ b/web/lib/community-agent-review-state.test.ts @@ -0,0 +1,391 @@ +import { afterEach, beforeEach, describe, expect, it, vi } from "vitest"; + +const mocks = vi.hoisted(() => ({ + agentChat: vi.fn(), + getAgentEnv: vi.fn(), + validateSession: vi.fn(), + fetchRepoStats: vi.fn(), +})); + +vi.mock("@/lib/community-agent", async (importOriginal) => { + const actual = await importOriginal(); + return { + ...actual, + agentChat: mocks.agentChat, + getAgentEnv: mocks.getAgentEnv, + validateSession: mocks.validateSession, + }; +}); + +vi.mock("@/lib/github", async (importOriginal) => { + const actual = await importOriginal(); + return { ...actual, fetchRepoStats: mocks.fetchRepoStats }; +}); + +import { POST as adminPost } from "../app/api/admin/post/route"; +import { isPublishedDigest, listDrafts, saveDraft, type AgentDraft } from "./community-agent"; +import { runDigest, runDupes, runTriage } from "./community-agent-tasks"; + +/** In-memory KV that pages like Cloudflare KV (max 1000 keys per list call). */ +class FakeKv { + readonly values = new Map(); + failPutsMatching: RegExp | null = null; + + async get(key: string): Promise { + return this.values.get(key) ?? null; + } + + async put(key: string, value: string): Promise { + if (this.failPutsMatching?.test(key)) throw new Error("kv put failed"); + this.values.set(key, value); + } + + async list(options?: { prefix?: string; limit?: number; cursor?: string }) { + const prefix = options?.prefix ?? ""; + const limit = Math.min(options?.limit ?? 1000, 1000); + const start = options?.cursor ? Number(options.cursor) : 0; + const all = [...this.values.keys()].filter((k) => k.startsWith(prefix)).sort(); + const page = all.slice(start, start + limit); + const next = start + page.length; + const complete = next >= all.length; + return { + keys: page.map((name) => ({ name })), + list_complete: complete, + ...(complete ? {} : { cursor: String(next) }), + }; + } + + async delete(key: string): Promise { + this.values.delete(key); + } +} + +function draft(overrides: Partial = {}): AgentDraft { + return { + id: "42", + type: "triage", + targetNumber: 42, + bodyEn: "English body", + bodyZh: "中文正文", + generatedAt: "2026-01-01T00:00:00.000Z", + posted: false, + ...overrides, + }; +} + +function jsonResponse(value: unknown, status = 200): Response { + return new Response(JSON.stringify(value), { status, headers: { "content-type": "application/json" } }); +} + +function inputUrl(input: string | URL | Request): string { + if (typeof input === "string") return input; + return input instanceof URL ? input.toString() : input.url; +} + +function adminRequest(body: BodyInit, headers: Record = {}): Request { + return new Request("https://codewhale.net/api/admin/post", { + method: "POST", + headers: { + "content-type": "application/json", + cookie: "mt_sid=test-session", + origin: "https://codewhale.net", + ...headers, + }, + body, + // Required by undici for a streamed request body. + ...(body instanceof ReadableStream ? { duplex: "half" } : {}), + } as RequestInit); +} + +function postBody(value: Record): string { + return JSON.stringify(value); +} + +function useAdminEnv(kv: FakeKv) { + mocks.getAgentEnv.mockResolvedValue({ + CURATED_KV: kv, + MAINTAINER_TOKEN: "configured", + MAINTAINER_GITHUB_PAT: "ghp_test", + GITHUB_REPO: "Hmbown/CodeWhale", + }); +} + +/** GitHub stub that records every comment/issue creation. */ +function stubGitHub(status = 201) { + const posts: string[] = []; + const fetchMock = vi.fn(async (input: string | URL | Request) => { + const url = inputUrl(input); + posts.push(url); + if (url.endsWith("/issues")) { + return jsonResponse({ number: 900, html_url: "https://github.com/Hmbown/CodeWhale/issues/900" }, status); + } + return jsonResponse({ id: 1 }, status); + }); + vi.stubGlobal("fetch", fetchMock); + return posts; +} + +const DIGEST_MODEL_OUTPUT = { + titleEn: "Weekly Digest", + titleZh: "每周摘要", + summaryEn: "A quiet week.", + summaryZh: "平静的一周。", + sections: [{ heading: "Shipped", items: ["PR #1: fix"] }], +}; + +function stubDigestSources() { + mocks.fetchRepoStats.mockResolvedValue({ stars: 1, forks: 1 }); + vi.stubGlobal("fetch", vi.fn(async (input: string | URL | Request) => { + const url = inputUrl(input); + if (url.includes("/issues?") || url.includes("/pulls?")) return jsonResponse([]); + throw new Error(`unexpected URL: ${url}`); + })); +} + +function onlyKey(kv: FakeKv, prefix: string): string { + const keys = [...kv.values.keys()].filter((k) => k.startsWith(prefix)); + expect(keys).toHaveLength(1); + return keys[0]; +} + +beforeEach(() => { + mocks.agentChat.mockReset(); + mocks.getAgentEnv.mockReset(); + mocks.fetchRepoStats.mockReset(); + mocks.validateSession.mockReset(); + mocks.validateSession.mockResolvedValue(true); +}); + +afterEach(() => { + vi.unstubAllGlobals(); +}); + +describe("weekly digest publication requires maintainer approval", () => { + it("stages the cron digest unapproved and publishes it only when the maintainer posts it", async () => { + const kv = new FakeKv(); + stubDigestSources(); + // Model output cannot self-approve. + mocks.agentChat.mockResolvedValue({ + content: JSON.stringify({ ...DIGEST_MODEL_OUTPUT, approved: true }), + usage: { input: 1, output: 1 }, + }); + + await expect(runDigest({ CURATED_KV: kv, DEEPSEEK_API_KEY: "k" })).resolves.toMatchObject({ ok: true }); + const recordKey = onlyKey(kv, "digest:weekly-"); + const staged: unknown = JSON.parse(kv.values.get(recordKey)!); + expect(staged).toMatchObject({ approved: false }); + expect(isPublishedDigest(staged)).toBe(false); + + useAdminEnv(kv); + const posts = stubGitHub(); + const draftKey = onlyKey(kv, "draft:digest:"); + const res = await adminPost(adminRequest(postBody({ action: "post", draftKey, lang: "en" }))); + await expect(res.json()).resolves.toMatchObject({ ok: true, published: true, number: 900 }); + expect(posts).toHaveLength(1); + expect(isPublishedDigest(JSON.parse(kv.values.get(recordKey)!))).toBe(true); + }); + + it("does not publish an edited digest's unedited model text", async () => { + const kv = new FakeKv(); + stubDigestSources(); + mocks.agentChat.mockResolvedValue({ content: JSON.stringify(DIGEST_MODEL_OUTPUT), usage: { input: 1, output: 1 } }); + await runDigest({ CURATED_KV: kv, DEEPSEEK_API_KEY: "k" }); + + useAdminEnv(kv); + stubGitHub(); + const draftKey = onlyKey(kv, "draft:digest:"); + const res = await adminPost(adminRequest(postBody({ action: "post", draftKey, lang: "en", editedBody: "# Edited" }))); + await expect(res.json()).resolves.toMatchObject({ ok: true, published: false }); + expect(isPublishedDigest(JSON.parse(kv.values.get(onlyKey(kv, "digest:weekly-"))!))).toBe(false); + }); + + it("discarding a digest removes the staged record and the cron does not regenerate it", async () => { + const kv = new FakeKv(); + stubDigestSources(); + mocks.agentChat.mockResolvedValue({ content: JSON.stringify(DIGEST_MODEL_OUTPUT), usage: { input: 1, output: 1 } }); + await runDigest({ CURATED_KV: kv, DEEPSEEK_API_KEY: "k" }); + + useAdminEnv(kv); + const draftKey = onlyKey(kv, "draft:digest:"); + const res = await adminPost(adminRequest(postBody({ action: "discard", draftKey }))); + await expect(res.json()).resolves.toMatchObject({ ok: true, action: "discarded" }); + expect([...kv.values.keys()].filter((k) => k.startsWith("digest:weekly-"))).toEqual([]); + + mocks.agentChat.mockClear(); + await expect(runDigest({ CURATED_KV: kv, DEEPSEEK_API_KEY: "k" })).resolves.toMatchObject({ skipped: true }); + expect(mocks.agentChat).not.toHaveBeenCalled(); + expect(kv.values.has(draftKey)).toBe(false); + }); + + it("surfaces a GitHub failure on a digest post as a 502 and does not publish", async () => { + const kv = new FakeKv(); + stubDigestSources(); + mocks.agentChat.mockResolvedValue({ content: JSON.stringify(DIGEST_MODEL_OUTPUT), usage: { input: 1, output: 1 } }); + await runDigest({ CURATED_KV: kv, DEEPSEEK_API_KEY: "k" }); + + useAdminEnv(kv); + stubGitHub(422); + const draftKey = onlyKey(kv, "draft:digest:"); + const res = await adminPost(adminRequest(postBody({ action: "post", draftKey, lang: "en" }))); + expect(res.status).toBe(502); + const payload = await res.json(); + expect(payload.ok).toBeUndefined(); + expect(payload.error).toMatch(/^GitHub 422/); + expect(isPublishedDigest(JSON.parse(kv.values.get(onlyKey(kv, "digest:weekly-"))!))).toBe(false); + }); + + it("hides legacy records that were never approved", () => { + const legacy = { ...DIGEST_MODEL_OUTPUT, weekId: "2026-W01", generatedAt: "2026-01-05T00:00:00.000Z" }; + expect(isPublishedDigest(legacy)).toBe(false); + expect(isPublishedDigest({ ...legacy, approved: true })).toBe(true); + expect(isPublishedDigest({ ...legacy, approved: true, sections: [{ heading: "x", items: [1] }] })).toBe(false); + }); +}); + +describe("resolved drafts are not resurrected by the cron", () => { + function stubIssues(updatedAt: string) { + vi.stubGlobal("fetch", vi.fn(async (input: string | URL | Request) => { + const url = inputUrl(input); + if (!url.includes("/issues?")) throw new Error(`unexpected URL: ${url}`); + return jsonResponse([{ + number: 42, + title: "Issue", + body: "body", + updated_at: updatedAt, + html_url: "https://github.com/Hmbown/CodeWhale/issues/42", + labels: [], + }]); + })); + } + + beforeEach(() => { + mocks.agentChat.mockResolvedValue({ + content: JSON.stringify({ bodyEn: "review", bodyZh: "审阅" }), + usage: { input: 1, output: 1 }, + }); + }); + + it("does not redraft a discarded triage item", async () => { + const kv = new FakeKv(); + await saveDraft(kv, draft()); + useAdminEnv(kv); + const res = await adminPost(adminRequest(postBody({ action: "discard", draftKey: "draft:triage:42" }))); + expect(res.status).toBe(200); + + stubIssues("2020-01-01T00:00:00.000Z"); + await expect(runTriage({ CURATED_KV: kv, DEEPSEEK_API_KEY: "k" })).resolves.toMatchObject({ processed: 0, skipped: 1 }); + expect(mocks.agentChat).not.toHaveBeenCalled(); + expect(kv.values.has("draft:triage:42")).toBe(false); + }); + + it("does not turn a posted triage draft back into a pending one after the issue updates", async () => { + const kv = new FakeKv(); + await saveDraft(kv, draft()); + useAdminEnv(kv); + stubGitHub(); + const res = await adminPost(adminRequest(postBody({ action: "post", draftKey: "draft:triage:42", lang: "en" }))); + expect(res.status).toBe(200); + + // Our own comment bumps updated_at past the draft's generatedAt. + stubIssues("2099-01-01T00:00:00.000Z"); + await runTriage({ CURATED_KV: kv, DEEPSEEK_API_KEY: "k" }); + expect(mocks.agentChat).not.toHaveBeenCalled(); + expect(JSON.parse(kv.values.get("draft:triage:42")!)).toMatchObject({ posted: true }); + }); + + it("does not let the dupes run overwrite a posted dupes draft", async () => { + const kv = new FakeKv(); + kv.values.set("draft:dupes:7", JSON.stringify(draft({ id: "7", type: "dupes", targetNumber: 7, posted: true }))); + vi.stubGlobal("fetch", vi.fn(async () => jsonResponse( + [1, 2, 7].map((n) => ({ number: n, title: `t${n}`, updated_at: "2020-01-01T00:00:00.000Z", html_url: `u${n}` })) + ))); + mocks.agentChat.mockResolvedValue({ + content: JSON.stringify({ suggestions: [{ targetNumber: 1, duplicateNumber: 7, reason: "r", bodyEn: "dup", bodyZh: "重复" }] }), + usage: { input: 1, output: 1 }, + }); + + await expect(runDupes({ CURATED_KV: kv, DEEPSEEK_API_KEY: "k" })).resolves.toMatchObject({ ok: true, processed: 0 }); + expect(JSON.parse(kv.values.get("draft:dupes:7")!)).toMatchObject({ posted: true }); + }); +}); + +describe("admin post action is idempotent and bounded", () => { + it("refuses to post a draft that is already posted", async () => { + const kv = new FakeKv(); + kv.values.set("draft:triage:42", JSON.stringify(draft({ posted: true }))); + useAdminEnv(kv); + const posts = stubGitHub(); + + const res = await adminPost(adminRequest(postBody({ action: "post", draftKey: "draft:triage:42" }))); + expect(res.status).toBe(409); + expect(posts).toEqual([]); + }); + + it("returns ok when saving state fails after GitHub accepted the comment, and a retry does not post twice", async () => { + const kv = new FakeKv(); + await saveDraft(kv, draft()); + useAdminEnv(kv); + const posts = stubGitHub(); + kv.failPutsMatching = /^draft:triage:42$/; + + const first = await adminPost(adminRequest(postBody({ action: "post", draftKey: "draft:triage:42" }))); + expect(first.status).toBe(200); + await expect(first.json()).resolves.toMatchObject({ ok: true, action: "posted", warning: expect.any(String) }); + + const retry = await adminPost(adminRequest(postBody({ action: "post", draftKey: "draft:triage:42" }))); + expect(retry.status).toBe(409); + expect(posts).toHaveLength(1); + }); + + it("releases the claim when GitHub rejects the post so the maintainer can retry", async () => { + const kv = new FakeKv(); + await saveDraft(kv, draft()); + useAdminEnv(kv); + stubGitHub(500); + + const failed = await adminPost(adminRequest(postBody({ action: "post", draftKey: "draft:triage:42" }))); + expect(failed.status).toBe(502); + + const posts = stubGitHub(); + const retry = await adminPost(adminRequest(postBody({ action: "post", draftKey: "draft:triage:42" }))); + expect(retry.status).toBe(200); + expect(posts).toHaveLength(1); + }); + + it("answers malformed JSON with a JSON 400", async () => { + const kv = new FakeKv(); + useAdminEnv(kv); + + const res = await adminPost(adminRequest("{not json")); + expect(res.status).toBe(400); + await expect(res.json()).resolves.toEqual({ error: "invalid JSON body" }); + }); + + it("counts streamed body bytes instead of trusting a missing Content-Length", async () => { + const kv = new FakeKv(); + useAdminEnv(kv); + const chunk = new TextEncoder().encode("x".repeat(40_000)); + const stream = new ReadableStream({ + start(controller) { + controller.enqueue(chunk); + controller.enqueue(chunk); + controller.close(); + }, + }); + + const res = await adminPost(adminRequest(stream)); + expect(res.status).toBe(413); + }); +}); + +describe("admin draft queue", () => { + it("lists every draft across KV list pages", async () => { + const kv = new FakeKv(); + for (let i = 1; i <= 1_500; i++) { + kv.values.set(`draft:triage:${i}`, JSON.stringify(draft({ id: String(i), targetNumber: i }))); + } + + const drafts = await listDrafts(kv); + expect(drafts).toHaveLength(1_500); + }); +}); diff --git a/web/lib/community-agent-tasks.ts b/web/lib/community-agent-tasks.ts index 11bf69ef28..b06d3a1703 100644 --- a/web/lib/community-agent-tasks.ts +++ b/web/lib/community-agent-tasks.ts @@ -10,7 +10,10 @@ import { DIGEST_PROMPT, saveDraft, hasFreshDraft, + getDraftResolution, + digestRecordKey, logUsage, + type WeeklyDigestRecord, type AgentDraft, type DeepSeekEnv, } from "@/lib/community-agent"; @@ -342,8 +345,8 @@ export async function runDupes(env: AgentEnv): Promise> generatedAt: new Date().toISOString(), posted: false, }; - await saveDraft(env.CURATED_KV, draft); - processed++; + // saveDraft refuses identities the maintainer already posted or discarded. + if (await saveDraft(env.CURATED_KV, draft)) processed++; } await logUsage(env.CURATED_KV, usage.input, usage.output); @@ -357,7 +360,17 @@ export async function runDigest(env: AgentEnv): Promise> const repo = env.GITHUB_REPO ?? "Hmbown/CodeWhale"; const weekAgo = new Date(Date.now() - 7 * 24 * 60 * 60 * 1000).toISOString(); + // Compute week ID + const now = new Date(); + const startOfYear = new Date(now.getFullYear(), 0, 1); + const weekNum = Math.ceil(((now.getTime() - startOfYear.getTime()) / 86400000 + startOfYear.getDay() + 1) / 7); + const weekId = `${now.getFullYear()}-W${String(weekNum).padStart(2, "0")}`; + try { + if (await getDraftResolution(env.CURATED_KV, "digest", weekId)) { + return { ok: true, skipped: true, weekId, reason: "digest already reviewed" }; + } + const [issuesRes, pullsRes, stats] = await Promise.all([ fetch( `https://api.github.com/repos/${repo}/issues?state=all&since=${weekAgo}&per_page=50&sort=updated&direction=desc`, @@ -414,12 +427,6 @@ export async function runDigest(env: AgentEnv): Promise> const parsed = JSON.parse(content) as { titleEn: string; titleZh: string; summaryEn: string; summaryZh: string; sections: { heading: string; items: string[] }[] }; - // Compute week ID - const now = new Date(); - const startOfYear = new Date(now.getFullYear(), 0, 1); - const weekNum = Math.ceil(((now.getTime() - startOfYear.getTime()) / 86400000 + startOfYear.getDay() + 1) / 7); - const weekId = `${now.getFullYear()}-W${String(weekNum).padStart(2, "0")}`; - const draft: AgentDraft = { id: weekId, type: "digest", @@ -429,14 +436,27 @@ export async function runDigest(env: AgentEnv): Promise> posted: false, }; - await saveDraft(env.CURATED_KV, draft); + if (!(await saveDraft(env.CURATED_KV, draft))) { + return { ok: true, skipped: true, weekId, reason: "digest already reviewed" }; + } - // Also save the structured digest for the weekly page - await env.CURATED_KV?.put( - `digest:weekly-${weekId}`, - JSON.stringify({ ...parsed, weekId, generatedAt: draft.generatedAt }), - { expirationTtl: 60 * 60 * 24 * 90 } - ); + // Stage the structured digest for the weekly page, unapproved. The page + // renders it only after the maintainer posts the draft from /admin. + // Pick fields explicitly: model output must never be able to set + // `approved` or any other record key. + const record: WeeklyDigestRecord = { + titleEn: parsed.titleEn, + titleZh: parsed.titleZh, + summaryEn: parsed.summaryEn, + summaryZh: parsed.summaryZh, + sections: parsed.sections, + weekId, + generatedAt: draft.generatedAt, + approved: false, + }; + await env.CURATED_KV?.put(digestRecordKey(weekId), JSON.stringify(record), { + expirationTtl: 60 * 60 * 24 * 90, + }); await logUsage(env.CURATED_KV, usage.input, usage.output); return { ok: true, weekId }; diff --git a/web/lib/community-agent.ts b/web/lib/community-agent.ts index 594d5425c6..429e7a54b6 100644 --- a/web/lib/community-agent.ts +++ b/web/lib/community-agent.ts @@ -234,7 +234,11 @@ ${VOICE_CONSTRAINTS}`; interface KVNamespace { get(key: string): Promise; put(key: string, value: string, opts?: { expirationTtl?: number }): Promise; - list(opts?: { prefix?: string; limit?: number }): Promise<{ keys: { name: string }[] }>; + list(opts?: { prefix?: string; limit?: number; cursor?: string }): Promise<{ + keys: { name: string }[]; + list_complete?: boolean; + cursor?: string; + }>; delete(key: string): Promise; } @@ -270,10 +274,151 @@ export async function getAgentEnv(): Promise { } } -export async function saveDraft(kv: KVNamespace | undefined, draft: AgentDraft): Promise { - if (!kv) return; +/** + * Persist a generated draft for maintainer review. Returns false without + * writing when the maintainer already posted or discarded this draft + * identity, so a cron run can never resurrect a resolved draft as pending. + */ +export async function saveDraft(kv: KVNamespace | undefined, draft: AgentDraft): Promise { + if (!kv) return false; const key = draftKey(draft.type, draft.id); + if (await getDraftResolution(kv, draft.type, draft.id)) return false; + const existing = await getDraft(kv, key); + if (existing?.posted) return false; await kv.put(key, JSON.stringify(draft), { expirationTtl: 60 * 60 * 24 * 30 }); // 30 days + return true; +} + +// --- Draft resolution markers --- +// +// Posting or discarding deletes/expires the draft itself, so the marker is +// what tells the generators and the post action that a maintainer already +// acted on this identity. It outlives the draft TTLs. + +export type DraftResolutionState = "posting" | "posted" | "discarded"; + +export interface DraftResolution { + state: DraftResolutionState; + at: string; +} + +const RESOLUTION_PREFIX = "draft-resolved:"; +const RESOLUTION_TTL_SEC = 60 * 60 * 24 * 90; // 90 days +// A claim taken before the GitHub call. It is short-lived so an unknown +// outcome (network error mid-request) blocks immediate retries without +// locking the draft forever. +const POSTING_CLAIM_TTL_SEC = 60 * 15; + +function resolutionKey(type: AgentDraftType, id: string): string { + return RESOLUTION_PREFIX + draftKey(type, id).slice("draft:".length); +} + +export async function getDraftResolution( + kv: KVNamespace | undefined, + type: AgentDraftType, + id: string +): Promise { + if (!kv) return null; + const raw = await kv.get(resolutionKey(type, id)); + if (!raw) return null; + try { + const parsed = JSON.parse(raw) as Partial; + if (parsed.state === "posting" || parsed.state === "posted" || parsed.state === "discarded") { + return { state: parsed.state, at: typeof parsed.at === "string" ? parsed.at : "" }; + } + } catch { + /* fall through: treat an unreadable marker as present */ + } + return { state: "posted", at: "" }; +} + +export async function markDraftResolved( + kv: KVNamespace | undefined, + type: AgentDraftType, + id: string, + state: DraftResolutionState +): Promise { + if (!kv) return; + const value: DraftResolution = { state, at: new Date().toISOString() }; + await kv.put(resolutionKey(type, id), JSON.stringify(value), { + expirationTtl: state === "posting" ? POSTING_CLAIM_TTL_SEC : RESOLUTION_TTL_SEC, + }); +} + +export async function clearDraftResolution( + kv: KVNamespace | undefined, + type: AgentDraftType, + id: string +): Promise { + if (!kv) return; + await kv.delete(resolutionKey(type, id)); +} + +// --- Public weekly digest records --- +// +// The cron writes the structured digest unapproved; only the maintainer's +// post action flips `approved`, and the public /digest page renders only +// approved records. + +export const DIGEST_RECORD_PREFIX = "digest:weekly-"; +const DIGEST_RECORD_TTL_SEC = 60 * 60 * 24 * 90; + +export interface WeeklyDigestRecord { + weekId: string; + titleEn: string; + titleZh: string; + summaryEn: string; + summaryZh: string; + sections: { heading: string; items: string[] }[]; + generatedAt: string; + approved?: boolean; + approvedAt?: string; +} + +export function digestRecordKey(weekId: string): string { + if (!DRAFT_ID_PATTERN.test(weekId)) throw new Error("invalid digest id"); + return DIGEST_RECORD_PREFIX + weekId; +} + +/** True only for a well-formed record a maintainer approved for publication. */ +export function isPublishedDigest(value: unknown): value is WeeklyDigestRecord { + if (!value || typeof value !== "object") return false; + const d = value as Record; + return ( + d.approved === true && + typeof d.weekId === "string" && + typeof d.titleEn === "string" && + typeof d.titleZh === "string" && + typeof d.summaryEn === "string" && + typeof d.summaryZh === "string" && + typeof d.generatedAt === "string" && + Array.isArray(d.sections) && + d.sections.every( + (s: unknown) => + !!s && + typeof (s as { heading?: unknown }).heading === "string" && + Array.isArray((s as { items?: unknown }).items) && + (s as { items: unknown[] }).items.every((i) => typeof i === "string") + ) + ); +} + +/** Mark the stored digest for `weekId` approved. Returns false if it is gone. */ +export async function approveDigestRecord(kv: KVNamespace | undefined, weekId: string): Promise { + if (!kv) return false; + const key = digestRecordKey(weekId); + const raw = await kv.get(key); + if (!raw) return false; + const record = JSON.parse(raw) as WeeklyDigestRecord; + const approved: WeeklyDigestRecord = { ...record, approved: true, approvedAt: new Date().toISOString() }; + if (!isPublishedDigest(approved)) return false; + await kv.put(key, JSON.stringify(approved), { expirationTtl: DIGEST_RECORD_TTL_SEC }); + return true; +} + +export async function deleteDigestRecord(kv: KVNamespace | undefined, weekId: string): Promise { + if (!kv) return; + await kv.delete(digestRecordKey(weekId)); } /** @@ -303,11 +448,18 @@ export async function getDraft(kv: KVNamespace | undefined, key: string): Promis export async function listDrafts(kv: KVNamespace | undefined, prefix = "draft:"): Promise { if (!kv) return []; - const listed = await kv.list({ prefix, limit: 100 }); const drafts: AgentDraft[] = []; - for (const k of listed.keys) { - const draft = await getDraft(kv, k.name); - if (draft) drafts.push(draft); + let cursor: string | undefined; + // KV returns at most 1000 keys per call; follow the cursor so the admin + // queue is never silently truncated. The page cap only bounds a runaway. + for (let page = 0; page < 20; page++) { + const listed = await kv.list({ prefix, limit: 1000, ...(cursor ? { cursor } : {}) }); + for (const k of listed.keys) { + const draft = await getDraft(kv, k.name); + if (draft) drafts.push(draft); + } + if (listed.list_complete !== false || !listed.cursor) break; + cursor = listed.cursor; } return drafts; } @@ -379,6 +531,11 @@ export async function logUsage( await kv.put(key, JSON.stringify(existing), { expirationTtl: 60 * 60 * 24 * 90 }); // 90 days } +/** + * True when a generator should not spend a model call drafting this item: + * the maintainer already posted or discarded it, or the stored draft is newer + * than the item's last update. + */ export async function hasFreshDraft( kv: KVNamespace | undefined, type: string, @@ -387,8 +544,9 @@ export async function hasFreshDraft( ): Promise { if (!kv) return false; if (!AGENT_DRAFT_TYPE_SET.has(type)) return false; + if (await getDraftResolution(kv, type as AgentDraftType, id)) return true; const existing = await getDraft(kv, draftKey(type as AgentDraftType, id)); if (!existing) return false; - // Skip if draft is newer than the item's last update + if (existing.posted) return true; return new Date(existing.generatedAt) > new Date(updatedAt); } diff --git a/web/lib/public-api-security.test.ts b/web/lib/public-api-security.test.ts index e56fc7c896..e21109361b 100644 --- a/web/lib/public-api-security.test.ts +++ b/web/lib/public-api-security.test.ts @@ -61,16 +61,6 @@ describe("public API security contracts", () => { ); }); - it("digest post: GitHub API failure surfaces a 502 error, not ok:true", () => { - const source = routeSource("admin/post"); - // On a failed digest GitHub call the handler must return a non-ok error payload - expect(source).toContain("digestRes.ok"); - // Must propagate the GitHub status rather than swallowing it - const digestErrorPath = source.slice( - source.indexOf("digestRes.ok"), - source.indexOf("digestRes.ok") + 300, - ); - expect(digestErrorPath).toContain("status: 502"); - expect(digestErrorPath).not.toContain('ok: true'); - }); + // "digest post: GitHub API failure surfaces a 502" is covered behaviorally + // in community-agent-review-state.test.ts. }); From 1d83478bd3cee4aa9012ba198016c56675bb5671 Mon Sep 17 00:00:00 2001 From: Hunter B Date: Tue, 29 Sep 2026 04:26:57 -0700 Subject: [PATCH 2/8] fix(web): tighten community draft review follow-ups Review of the first commit on this branch found gaps; this addresses them. - Posting claim: the claim now carries a random token that is read back before the GitHub call, so two concurrent posts of one draft that meet in the same KV location post once. KV has no compare-and-set and is eventually consistent across locations, so this is documented as best effort, not a lock. - GitHub 5xx, 408 and 429 answers keep the claim (the post may exist); only other 4xx answers release it for an immediate retry. - Resolution markers no longer suppress triage / PR-review / stale drafting for 90 days regardless of activity: an item updated more than 10 minutes after the maintainer's decision can be drafted again (our own post's updated_at bump stays covered). Digests, dupes and content-watch findings keep the 90-day suppression; their identity already encodes the content. - /digest publishes only the language the maintainer reviewed (approvedLang), revalidates the page after approval or discard, and a posted draft can no longer be discarded (409), which would have silently unpublished a digest. - The admin client shows the route's warning (state not saved, or digest not published because it was edited or its record is gone). - listDrafts reads at most 500 drafts in parallel batches of 50 and keeps what it read when a later list page fails. - runSemanticDrift and the triage / PR-review / stale runs count a draft only when saveDraft wrote it. Tests: web vitest 57 files / 526 passed; targeted community-agent-review-state, content-watch, community-agent-security, public-api-security 4 files / 45 passed; 14 of the new or changed tests fail on the previous source; tsc --noEmit exit 0; eslint on touched files exit 0; web npm run check ran through next build. Co-Authored-By: Claude Opus 5.5 (1M context) Claude-Session: https://claude.ai/code/session_014ZwqatxgVFxHvovngywnks --- web/app/[locale]/admin/admin-client.tsx | 3 + web/app/[locale]/digest/page.tsx | 15 +- web/app/api/admin/post/route.ts | 63 ++++++-- web/lib/community-agent-review-state.test.ts | 160 +++++++++++++++++-- web/lib/community-agent-tasks.ts | 12 +- web/lib/community-agent.ts | 146 ++++++++++++++--- web/lib/content-watch.test.ts | 15 +- web/lib/content-watch.ts | 4 +- 8 files changed, 356 insertions(+), 62 deletions(-) diff --git a/web/app/[locale]/admin/admin-client.tsx b/web/app/[locale]/admin/admin-client.tsx index 51275d493e..58c51c4f54 100644 --- a/web/app/[locale]/admin/admin-client.tsx +++ b/web/app/[locale]/admin/admin-client.tsx @@ -37,6 +37,9 @@ export function AdminClient({ drafts, posted, isZh, typeLabels }: Props) { } } setEditing(null); + // The post went through but something after it did not (draft state + // not saved, or the digest was not published on /digest). + if (typeof data.warning === "string") alert(data.warning); } else { alert(`Error: ${data.error}`); } diff --git a/web/app/[locale]/digest/page.tsx b/web/app/[locale]/digest/page.tsx index c5a658d6ed..ff5edac46d 100644 --- a/web/app/[locale]/digest/page.tsx +++ b/web/app/[locale]/digest/page.tsx @@ -5,7 +5,7 @@ import { isPublishedDigest, type WeeklyDigestRecord as WeeklyDigest, } from "@/lib/community-agent"; -import { getDigest, pickTextLocale } from "@/lib/i18n/dictionaries"; +import { getDigest } from "@/lib/i18n/dictionaries"; import { getEnv } from "@/lib/kv"; import { buildPageMetadata } from "@/lib/page-meta"; @@ -71,9 +71,9 @@ export default async function DigestArchivePage({ params }: { params: Promise<{ } } - // The records are bilingual; show the reader's language (English for - // every locale without a Chinese record). - const zh = pickTextLocale(locale) === "zh"; + // A maintainer reviews a digest in one language (the admin locale), and + // only that language is published, so each digest renders in the language + // it was reviewed in whatever the reader's locale. if (digests.length === 0) { return ( @@ -91,7 +91,9 @@ export default async function DigestArchivePage({ params }: { params: Promise<{
- {digests.map((digest: WeeklyDigest) => ( + {digests.map((digest: WeeklyDigest) => { + const zh = digest.approvedLang === "zh"; + return (
{digest.weekId} @@ -111,7 +113,8 @@ export default async function DigestArchivePage({ params }: { params: Promise<{ ))}
- ))} + ); + })}
diff --git a/web/app/api/admin/post/route.ts b/web/app/api/admin/post/route.ts index 76f6e3fad5..1e8fb56c05 100644 --- a/web/app/api/admin/post/route.ts +++ b/web/app/api/admin/post/route.ts @@ -1,7 +1,9 @@ +import { revalidatePath } from "next/cache"; import { NextResponse } from "next/server"; import { BodyReadError, readBoundedBody } from "@/lib/bounded-body"; import { approveDigestRecord, + claimDraftForPosting, clearDraftResolution, deleteDigestRecord, deleteDraft, @@ -41,6 +43,25 @@ const ALLOWED_ACTIONS = new Set(["post", "discard"]); const ALLOWED_ORIGINS = new Set(["https://codewhale.net", "https://www.codewhale.net"]); const MAX_BODY_BYTES = 65_536; +/** Refresh the ISR copy of /digest after its published set changes. */ +function revalidateDigestPage() { + try { + revalidatePath("/[locale]/digest", "page"); + } catch (e) { + // The page still refreshes on its hourly revalidate. + console.error("digest revalidation failed", e); + } +} + +/** + * A GitHub 4xx other than 408/429 means the post was not created, so the + * claim can be released. A 5xx, 408 or 429 can come back after GitHub already + * created the comment or issue, so its outcome is unknown. + */ +function githubDefinitelyRejected(status: number): boolean { + return status >= 400 && status < 500 && status !== 408 && status !== 429; +} + export async function POST(req: Request) { const env = await getAgentEnv(); @@ -105,6 +126,12 @@ export async function POST(req: Request) { } if (action === "discard") { + // A posted draft is already public on GitHub (and, for a digest, on + // /digest); discarding it would silently unpublish or relabel it. + const resolution = await getDraftResolution(env.CURATED_KV, parsedKey.type, parsedKey.id); + if (draft.posted || resolution?.state === "posted" || resolution?.state === "posting") { + return NextResponse.json({ error: "draft already posted" }, { status: 409 }); + } try { // The marker stops the next cron run from regenerating this draft. await markDraftResolved(env.CURATED_KV, parsedKey.type, parsedKey.id, "discarded"); @@ -115,6 +142,7 @@ export async function POST(req: Request) { } catch (e) { return NextResponse.json({ error: `discard failed: ${String(e)}` }, { status: 500 }); } + if (draft.type === "digest") revalidateDigestPage(); return NextResponse.json({ ok: true, action: "discarded" }); } @@ -137,7 +165,9 @@ export async function POST(req: Request) { return NextResponse.json({ error: `draft already ${resolution.state}` }, { status: 409 }); } try { - await markDraftResolved(env.CURATED_KV, parsedKey.type, parsedKey.id, "posting"); + if (!(await claimDraftForPosting(env.CURATED_KV, parsedKey.type, parsedKey.id))) { + return NextResponse.json({ error: "draft already posting" }, { status: 409 }); + } } catch (e) { return NextResponse.json({ error: `could not claim draft: ${String(e)}` }, { status: 500 }); } @@ -153,7 +183,7 @@ export async function POST(req: Request) { await env.CURATED_KV?.put(draftKey, JSON.stringify(draft), { expirationTtl: 60 * 60 * 24 * 7 }); return undefined; } catch (e) { - return `posted, but saving draft state failed: ${String(e)}`; + return `Posted to GitHub, but saving the draft state failed (${String(e)}). Do not post it again; it may reappear as pending.`; } }; @@ -175,17 +205,20 @@ export async function POST(req: Request) { return null; } }; - const unknownOutcome = () => + const unknownOutcome = (detail: string) => NextResponse.json( - { error: "GitHub request failed; check GitHub before retrying (retry unlocks in 15 minutes)" }, + { error: `${detail}; check GitHub before retrying (retry unlocks in 15 minutes)` }, { status: 502 } ); const githubFailed = async (res: Response) => { - // GitHub definitively rejected the post, so release the claim. + const text = await res.text().catch(() => ""); + if (!githubDefinitelyRejected(res.status)) { + // Keep the claim: GitHub may have created the post anyway. + return unknownOutcome(`GitHub ${res.status}: ${text}`); + } try { await clearDraftResolution(env.CURATED_KV, parsedKey.type, parsedKey.id); } catch { /* the claim expires on its own */ } - const text = await res.text(); return NextResponse.json({ error: `GitHub ${res.status}: ${text}` }, { status: 502 }); }; @@ -198,7 +231,7 @@ export async function POST(req: Request) { const issuesUrl = `https://api.github.com/repos/${digestRepo}/issues`; const digestRes = await postToGitHub(issuesUrl, { title, body: digestBody, labels: ["digest"] }, "token"); - if (!digestRes) return unknownOutcome(); + if (!digestRes) return unknownOutcome("GitHub request failed"); if (!digestRes.ok) return githubFailed(digestRes); const issue = await digestRes.json().catch(() => ({})) as { number?: number; html_url?: string }; @@ -208,16 +241,20 @@ export async function POST(req: Request) { if (typeof issue.html_url === "string") draft.targetUrl = issue.html_url; let warning = await recordPosted(); - // Publishing to /digest is the approval. The structured record holds - // the unedited model text, so an edited digest is posted to GitHub but - // not published there. + // Publishing to /digest is the approval, for the language shown to the + // maintainer only. The structured record holds the unedited model + // text, so an edited digest is posted to GitHub but not published there. let published = false; if (editedBody === undefined || editedBody === originalBody) { try { - published = await approveDigestRecord(env.CURATED_KV, draft.id); + published = await approveDigestRecord(env.CURATED_KV, draft.id, lang === "zh" ? "zh" : "en"); + if (published) revalidateDigestPage(); } catch (e) { - warning ??= `posted, but publishing the digest page failed: ${String(e)}`; + warning ??= `Posted to GitHub, but publishing the digest page failed (${String(e)}).`; } + if (!published) warning ??= "Posted to GitHub, but the weekly record is gone, so /digest does not show it."; + } else { + warning ??= "Posted to GitHub. The text was edited, so /digest does not show this week."; } return NextResponse.json({ @@ -234,7 +271,7 @@ export async function POST(req: Request) { const commentUrl = `https://api.github.com/repos/${repo}/issues/${draft.targetNumber}/comments`; const ghRes = await postToGitHub(commentUrl, { body: commentBody }, "Bearer"); - if (!ghRes) return unknownOutcome(); + if (!ghRes) return unknownOutcome("GitHub request failed"); if (!ghRes.ok) return githubFailed(ghRes); // Mark as posted diff --git a/web/lib/community-agent-review-state.test.ts b/web/lib/community-agent-review-state.test.ts index 60ed6c310d..e50d1528d4 100644 --- a/web/lib/community-agent-review-state.test.ts +++ b/web/lib/community-agent-review-state.test.ts @@ -23,26 +23,43 @@ vi.mock("@/lib/github", async (importOriginal) => { }); import { POST as adminPost } from "../app/api/admin/post/route"; -import { isPublishedDigest, listDrafts, saveDraft, type AgentDraft } from "./community-agent"; +import { + isPublishedDigest, + listDrafts, + MAX_LISTED_DRAFTS, + saveDraft, + type AgentDraft, +} from "./community-agent"; import { runDigest, runDupes, runTriage } from "./community-agent-tasks"; /** In-memory KV that pages like Cloudflare KV (max 1000 keys per list call). */ class FakeKv { readonly values = new Map(); failPutsMatching: RegExp | null = null; + /** Yield to the event loop inside every call, like a network round trip. */ + yieldEachOp = false; + pageSize = 1000; + failListAtCursor: string | null = null; + + private async tick() { + if (this.yieldEachOp) await new Promise((resolve) => setTimeout(resolve, 0)); + } async get(key: string): Promise { + await this.tick(); return this.values.get(key) ?? null; } async put(key: string, value: string): Promise { + await this.tick(); if (this.failPutsMatching?.test(key)) throw new Error("kv put failed"); this.values.set(key, value); } async list(options?: { prefix?: string; limit?: number; cursor?: string }) { + if (options?.cursor && options.cursor === this.failListAtCursor) throw new Error("kv list failed"); const prefix = options?.prefix ?? ""; - const limit = Math.min(options?.limit ?? 1000, 1000); + const limit = Math.min(options?.limit ?? 1000, this.pageSize); const start = options?.cursor ? Number(options.cursor) : 0; const all = [...this.values.keys()].filter((k) => k.startsWith(prefix)).sort(); const page = all.slice(start, start + limit); @@ -182,7 +199,40 @@ describe("weekly digest publication requires maintainer approval", () => { const res = await adminPost(adminRequest(postBody({ action: "post", draftKey, lang: "en" }))); await expect(res.json()).resolves.toMatchObject({ ok: true, published: true, number: 900 }); expect(posts).toHaveLength(1); - expect(isPublishedDigest(JSON.parse(kv.values.get(recordKey)!))).toBe(true); + const published = JSON.parse(kv.values.get(recordKey)!); + expect(isPublishedDigest(published)).toBe(true); + // Only the reviewed language is published. + expect(published).toMatchObject({ approvedLang: "en" }); + }); + + it("records the Chinese review as the published language when posted from the zh admin", async () => { + const kv = new FakeKv(); + stubDigestSources(); + mocks.agentChat.mockResolvedValue({ content: JSON.stringify(DIGEST_MODEL_OUTPUT), usage: { input: 1, output: 1 } }); + await runDigest({ CURATED_KV: kv, DEEPSEEK_API_KEY: "k" }); + + useAdminEnv(kv); + stubGitHub(); + const draftKey = onlyKey(kv, "draft:digest:"); + const res = await adminPost(adminRequest(postBody({ action: "post", draftKey, lang: "zh" }))); + await expect(res.json()).resolves.toMatchObject({ ok: true, published: true }); + expect(JSON.parse(kv.values.get(onlyKey(kv, "digest:weekly-"))!)).toMatchObject({ approvedLang: "zh" }); + }); + + it("refuses to discard a posted digest, which would silently unpublish it", async () => { + const kv = new FakeKv(); + stubDigestSources(); + mocks.agentChat.mockResolvedValue({ content: JSON.stringify(DIGEST_MODEL_OUTPUT), usage: { input: 1, output: 1 } }); + await runDigest({ CURATED_KV: kv, DEEPSEEK_API_KEY: "k" }); + + useAdminEnv(kv); + stubGitHub(); + const draftKey = onlyKey(kv, "draft:digest:"); + expect((await adminPost(adminRequest(postBody({ action: "post", draftKey, lang: "en" })))).status).toBe(200); + + const discard = await adminPost(adminRequest(postBody({ action: "discard", draftKey }))); + expect(discard.status).toBe(409); + expect(isPublishedDigest(JSON.parse(kv.values.get(onlyKey(kv, "digest:weekly-"))!))).toBe(true); }); it("does not publish an edited digest's unedited model text", async () => { @@ -195,7 +245,12 @@ describe("weekly digest publication requires maintainer approval", () => { stubGitHub(); const draftKey = onlyKey(kv, "draft:digest:"); const res = await adminPost(adminRequest(postBody({ action: "post", draftKey, lang: "en", editedBody: "# Edited" }))); - await expect(res.json()).resolves.toMatchObject({ ok: true, published: false }); + // The maintainer is told why the digest is not on /digest. + await expect(res.json()).resolves.toMatchObject({ + ok: true, + published: false, + warning: expect.stringContaining("/digest"), + }); expect(isPublishedDigest(JSON.parse(kv.values.get(onlyKey(kv, "digest:weekly-"))!))).toBe(false); }); @@ -237,8 +292,11 @@ describe("weekly digest publication requires maintainer approval", () => { it("hides legacy records that were never approved", () => { const legacy = { ...DIGEST_MODEL_OUTPUT, weekId: "2026-W01", generatedAt: "2026-01-05T00:00:00.000Z" }; expect(isPublishedDigest(legacy)).toBe(false); - expect(isPublishedDigest({ ...legacy, approved: true })).toBe(true); - expect(isPublishedDigest({ ...legacy, approved: true, sections: [{ heading: "x", items: [1] }] })).toBe(false); + expect(isPublishedDigest({ ...legacy, approved: true })).toBe(false); + expect(isPublishedDigest({ ...legacy, approved: true, approvedLang: "en" })).toBe(true); + expect( + isPublishedDigest({ ...legacy, approved: true, approvedLang: "en", sections: [{ heading: "x", items: [1] }] }) + ).toBe(false); }); }); @@ -287,12 +345,30 @@ describe("resolved drafts are not resurrected by the cron", () => { expect(res.status).toBe(200); // Our own comment bumps updated_at past the draft's generatedAt. - stubIssues("2099-01-01T00:00:00.000Z"); + stubIssues(new Date().toISOString()); await runTriage({ CURATED_KV: kv, DEEPSEEK_API_KEY: "k" }); expect(mocks.agentChat).not.toHaveBeenCalled(); expect(JSON.parse(kv.values.get("draft:triage:42")!)).toMatchObject({ posted: true }); }); + it("drafts again when the item has new activity well after the maintainer's decision", async () => { + const kv = new FakeKv(); + await saveDraft(kv, draft()); + useAdminEnv(kv); + stubGitHub(); + expect((await adminPost(adminRequest(postBody({ action: "post", draftKey: "draft:triage:42" })))).status).toBe(200); + + // New commits or a reply a day later. + stubIssues(new Date(Date.now() + 24 * 60 * 60 * 1000).toISOString()); + await expect(runTriage({ CURATED_KV: kv, DEEPSEEK_API_KEY: "k" })).resolves.toMatchObject({ processed: 1 }); + expect(JSON.parse(kv.values.get("draft:triage:42")!)).toMatchObject({ posted: false, bodyEn: "review" }); + + // The fresh draft can be posted. + const posts = stubGitHub(); + expect((await adminPost(adminRequest(postBody({ action: "post", draftKey: "draft:triage:42" })))).status).toBe(200); + expect(posts).toHaveLength(1); + }); + it("does not let the dupes run overwrite a posted dupes draft", async () => { const kv = new FakeKv(); kv.values.set("draft:dupes:7", JSON.stringify(draft({ id: "7", type: "dupes", targetNumber: 7, posted: true }))); @@ -341,7 +417,7 @@ describe("admin post action is idempotent and bounded", () => { const kv = new FakeKv(); await saveDraft(kv, draft()); useAdminEnv(kv); - stubGitHub(500); + stubGitHub(422); const failed = await adminPost(adminRequest(postBody({ action: "post", draftKey: "draft:triage:42" }))); expect(failed.status).toBe(502); @@ -352,6 +428,41 @@ describe("admin post action is idempotent and bounded", () => { expect(posts).toHaveLength(1); }); + it.each([500, 502, 504, 429])( + "keeps the claim when GitHub answers %i, since the post may have been created", + async (status) => { + const kv = new FakeKv(); + await saveDraft(kv, draft()); + useAdminEnv(kv); + const posts = stubGitHub(status); + + const failed = await adminPost(adminRequest(postBody({ action: "post", draftKey: "draft:triage:42" }))); + expect(failed.status).toBe(502); + await expect(failed.json()).resolves.toMatchObject({ error: expect.stringContaining("check GitHub before retrying") }); + + stubGitHub(); + const retry = await adminPost(adminRequest(postBody({ action: "post", draftKey: "draft:triage:42" }))); + expect(retry.status).toBe(409); + expect(posts).toHaveLength(1); + } + ); + + it("posts once when two requests for the same draft race", async () => { + const kv = new FakeKv(); + await saveDraft(kv, draft()); + kv.yieldEachOp = true; + useAdminEnv(kv); + const posts = stubGitHub(); + + const results = await Promise.all([ + adminPost(adminRequest(postBody({ action: "post", draftKey: "draft:triage:42" }))), + adminPost(adminRequest(postBody({ action: "post", draftKey: "draft:triage:42" }))), + ]); + + expect(results.map((r) => r.status).sort()).toEqual([200, 409]); + expect(posts).toHaveLength(1); + }); + it("answers malformed JSON with a JSON 400", async () => { const kv = new FakeKv(); useAdminEnv(kv); @@ -379,13 +490,36 @@ describe("admin post action is idempotent and bounded", () => { }); describe("admin draft queue", () => { - it("lists every draft across KV list pages", async () => { - const kv = new FakeKv(); - for (let i = 1; i <= 1_500; i++) { + function fill(kv: FakeKv, count: number) { + for (let i = 1; i <= count; i++) { kv.values.set(`draft:triage:${i}`, JSON.stringify(draft({ id: String(i), targetNumber: i }))); } + } + + it("follows the KV list cursor past the first page", async () => { + const kv = new FakeKv(); + kv.pageSize = 100; + fill(kv, 250); + + await expect(listDrafts(kv)).resolves.toHaveLength(250); + }); + + it("reads at most MAX_LISTED_DRAFTS drafts in one request", async () => { + const kv = new FakeKv(); + fill(kv, MAX_LISTED_DRAFTS + 1_000); + const get = vi.spyOn(kv, "get"); + + await expect(listDrafts(kv)).resolves.toHaveLength(MAX_LISTED_DRAFTS); + expect(get).toHaveBeenCalledTimes(MAX_LISTED_DRAFTS); + }); + + it("keeps the drafts already read when a later list page fails", async () => { + const kv = new FakeKv(); + kv.pageSize = 100; + kv.failListAtCursor = "200"; + fill(kv, 250); + vi.spyOn(console, "error").mockImplementation(() => {}); - const drafts = await listDrafts(kv); - expect(drafts).toHaveLength(1_500); + await expect(listDrafts(kv)).resolves.toHaveLength(200); }); }); diff --git a/web/lib/community-agent-tasks.ts b/web/lib/community-agent-tasks.ts index b06d3a1703..89e39aa691 100644 --- a/web/lib/community-agent-tasks.ts +++ b/web/lib/community-agent-tasks.ts @@ -132,9 +132,9 @@ export async function runTriage(env: AgentEnv): Promise> generatedAt: new Date().toISOString(), posted: false, }; - await saveDraft(env.CURATED_KV, draft); await logUsage(env.CURATED_KV, usage.input, usage.output); - processed++; + if (await saveDraft(env.CURATED_KV, draft, issue.updated_at)) processed++; + else skipped++; } catch { skipped++; } @@ -215,9 +215,9 @@ export async function runPrReview(env: AgentEnv): Promise> generatedAt: new Date().toISOString(), posted: false, }; - await saveDraft(env.CURATED_KV, draft); await logUsage(env.CURATED_KV, usage.input, usage.output); - processed++; + if (await saveDraft(env.CURATED_KV, draft, issue.updated_at)) processed++; + else skipped++; } catch { skipped++; } diff --git a/web/lib/community-agent.ts b/web/lib/community-agent.ts index 429e7a54b6..8555c9d33a 100644 --- a/web/lib/community-agent.ts +++ b/web/lib/community-agent.ts @@ -278,13 +278,30 @@ export async function getAgentEnv(): Promise { * Persist a generated draft for maintainer review. Returns false without * writing when the maintainer already posted or discarded this draft * identity, so a cron run can never resurrect a resolved draft as pending. + * + * `sourceUpdatedAt` is the GitHub item's `updated_at` for drafts about an + * issue or PR. Activity well after the maintainer's decision (new commits, a + * reply) reopens the identity; see `resolutionCovers`. */ -export async function saveDraft(kv: KVNamespace | undefined, draft: AgentDraft): Promise { +export async function saveDraft( + kv: KVNamespace | undefined, + draft: AgentDraft, + sourceUpdatedAt?: string +): Promise { if (!kv) return false; const key = draftKey(draft.type, draft.id); - if (await getDraftResolution(kv, draft.type, draft.id)) return false; - const existing = await getDraft(kv, key); - if (existing?.posted) return false; + const resolution = await getDraftResolution(kv, draft.type, draft.id); + if (resolution) { + if (resolutionCovers(resolution, sourceUpdatedAt)) return false; + // Reopened by new activity. Clear the marker first: if the put below + // then fails, the old posted draft still blocks a redraft, which is the + // safe side. + await clearDraftResolution(kv, draft.type, draft.id); + } else { + // A posted draft without a marker predates the markers; keep it. + const existing = await getDraft(kv, key); + if (existing?.posted) return false; + } await kv.put(key, JSON.stringify(draft), { expirationTtl: 60 * 60 * 24 * 30 }); // 30 days return true; } @@ -300,6 +317,8 @@ export type DraftResolutionState = "posting" | "posted" | "discarded"; export interface DraftResolution { state: DraftResolutionState; at: string; + /** Random token of the request holding a "posting" claim. */ + claim?: string; } const RESOLUTION_PREFIX = "draft-resolved:"; @@ -308,6 +327,25 @@ const RESOLUTION_TTL_SEC = 60 * 60 * 24 * 90; // 90 days // outcome (network error mid-request) blocks immediate retries without // locking the draft forever. const POSTING_CLAIM_TTL_SEC = 60 * 15; +// Our own post bumps the GitHub item's updated_at moments before the marker +// is written; only activity later than this after the decision reopens it. +const RESOLUTION_REOPEN_GRACE_MS = 10 * 60 * 1000; + +/** + * True while a maintainer decision still applies to the item as it is now. + * A decision covers everything up to shortly after it was made; an item + * updated later (new commits, a reply) can be drafted again. An in-flight + * post, or an item without an update time (digests, dupes, content-watch + * findings, whose identity already encodes their content), stays covered + * for the marker's lifetime. + */ +export function resolutionCovers(resolution: DraftResolution, sourceUpdatedAt?: string): boolean { + if (resolution.state === "posting" || !sourceUpdatedAt) return true; + const decidedAt = Date.parse(resolution.at); + const updatedAt = Date.parse(sourceUpdatedAt); + if (!Number.isFinite(decidedAt) || !Number.isFinite(updatedAt)) return true; + return updatedAt <= decidedAt + RESOLUTION_REOPEN_GRACE_MS; +} function resolutionKey(type: AgentDraftType, id: string): string { return RESOLUTION_PREFIX + draftKey(type, id).slice("draft:".length); @@ -324,7 +362,11 @@ export async function getDraftResolution( try { const parsed = JSON.parse(raw) as Partial; if (parsed.state === "posting" || parsed.state === "posted" || parsed.state === "discarded") { - return { state: parsed.state, at: typeof parsed.at === "string" ? parsed.at : "" }; + return { + state: parsed.state, + at: typeof parsed.at === "string" ? parsed.at : "", + ...(typeof parsed.claim === "string" ? { claim: parsed.claim } : {}), + }; } } catch { /* fall through: treat an unreadable marker as present */ @@ -345,6 +387,29 @@ export async function markDraftResolved( }); } +/** + * Claim a draft before posting it to GitHub. Returns false when another + * request holds or wins the claim. + * + * KV has no compare-and-set, so this is best effort, not a lock: it writes a + * random token and reads it back, which stops concurrent posts that meet in + * one location (two tabs, a double submit) but cannot fully order writers in + * different edge locations, because KV is eventually consistent across them. + */ +export async function claimDraftForPosting( + kv: KVNamespace | undefined, + type: AgentDraftType, + id: string +): Promise { + if (!kv) return true; + const key = resolutionKey(type, id); + const claim = crypto.randomUUID(); + const value: DraftResolution = { state: "posting", at: new Date().toISOString(), claim }; + await kv.put(key, JSON.stringify(value), { expirationTtl: POSTING_CLAIM_TTL_SEC }); + const current = await getDraftResolution(kv, type, id); + return current?.state === "posting" && current.claim === claim; +} + export async function clearDraftResolution( kv: KVNamespace | undefined, type: AgentDraftType, @@ -358,7 +423,7 @@ export async function clearDraftResolution( // // The cron writes the structured digest unapproved; only the maintainer's // post action flips `approved`, and the public /digest page renders only -// approved records. +// approved records, in the one language the maintainer reviewed. export const DIGEST_RECORD_PREFIX = "digest:weekly-"; const DIGEST_RECORD_TTL_SEC = 60 * 60 * 24 * 90; @@ -373,6 +438,8 @@ export interface WeeklyDigestRecord { generatedAt: string; approved?: boolean; approvedAt?: string; + /** The language the maintainer reviewed; the only one /digest shows. */ + approvedLang?: "en" | "zh"; } export function digestRecordKey(weekId: string): string { @@ -386,6 +453,7 @@ export function isPublishedDigest(value: unknown): value is WeeklyDigestRecord { const d = value as Record; return ( d.approved === true && + (d.approvedLang === "en" || d.approvedLang === "zh") && typeof d.weekId === "string" && typeof d.titleEn === "string" && typeof d.titleZh === "string" && @@ -403,14 +471,26 @@ export function isPublishedDigest(value: unknown): value is WeeklyDigestRecord { ); } -/** Mark the stored digest for `weekId` approved. Returns false if it is gone. */ -export async function approveDigestRecord(kv: KVNamespace | undefined, weekId: string): Promise { +/** + * Mark the stored digest for `weekId` approved in the language the maintainer + * reviewed. Returns false if it is gone. + */ +export async function approveDigestRecord( + kv: KVNamespace | undefined, + weekId: string, + lang: "en" | "zh" +): Promise { if (!kv) return false; const key = digestRecordKey(weekId); const raw = await kv.get(key); if (!raw) return false; const record = JSON.parse(raw) as WeeklyDigestRecord; - const approved: WeeklyDigestRecord = { ...record, approved: true, approvedAt: new Date().toISOString() }; + const approved: WeeklyDigestRecord = { + ...record, + approved: true, + approvedAt: new Date().toISOString(), + approvedLang: lang, + }; if (!isPublishedDigest(approved)) return false; await kv.put(key, JSON.stringify(approved), { expirationTtl: DIGEST_RECORD_TTL_SEC }); return true; @@ -446,20 +526,42 @@ export async function getDraft(kv: KVNamespace | undefined, key: string): Promis } } +// One draft read is one KV operation, and a Worker invocation gets a bounded +// number of them, so the admin queue reads at most this many drafts. +export const MAX_LISTED_DRAFTS = 500; +const DRAFT_READ_BATCH = 50; + +/** + * Read up to MAX_LISTED_DRAFTS drafts, following the KV list cursor and + * reading drafts in parallel batches. A failure after the first list call + * returns what was read so far instead of discarding it. + */ export async function listDrafts(kv: KVNamespace | undefined, prefix = "draft:"): Promise { if (!kv) return []; const drafts: AgentDraft[] = []; + let seen = 0; let cursor: string | undefined; - // KV returns at most 1000 keys per call; follow the cursor so the admin - // queue is never silently truncated. The page cap only bounds a runaway. - for (let page = 0; page < 20; page++) { - const listed = await kv.list({ prefix, limit: 1000, ...(cursor ? { cursor } : {}) }); - for (const k of listed.keys) { - const draft = await getDraft(kv, k.name); - if (draft) drafts.push(draft); + try { + while (seen < MAX_LISTED_DRAFTS) { + const listed = await kv.list({ + prefix, + limit: Math.min(1000, MAX_LISTED_DRAFTS - seen), + ...(cursor ? { cursor } : {}), + }); + const names = listed.keys.map((k) => k.name).slice(0, MAX_LISTED_DRAFTS - seen); + seen += names.length; + for (let i = 0; i < names.length; i += DRAFT_READ_BATCH) { + const batch = await Promise.all( + names.slice(i, i + DRAFT_READ_BATCH).map((name) => getDraft(kv, name).catch(() => null)) + ); + for (const draft of batch) if (draft) drafts.push(draft); + } + if (names.length === 0 || listed.list_complete !== false || !listed.cursor) break; + cursor = listed.cursor; } - if (listed.list_complete !== false || !listed.cursor) break; - cursor = listed.cursor; + } catch (e) { + if (seen === 0) throw e; + console.error("listDrafts: returning a partial queue", e); } return drafts; } @@ -533,8 +635,8 @@ export async function logUsage( /** * True when a generator should not spend a model call drafting this item: - * the maintainer already posted or discarded it, or the stored draft is newer - * than the item's last update. + * a maintainer decision still covers it (see `resolutionCovers`), or the + * stored draft is newer than the item's last update. */ export async function hasFreshDraft( kv: KVNamespace | undefined, @@ -544,9 +646,11 @@ export async function hasFreshDraft( ): Promise { if (!kv) return false; if (!AGENT_DRAFT_TYPE_SET.has(type)) return false; - if (await getDraftResolution(kv, type as AgentDraftType, id)) return true; + const resolution = await getDraftResolution(kv, type as AgentDraftType, id); + if (resolution) return resolutionCovers(resolution, updatedAt); const existing = await getDraft(kv, draftKey(type as AgentDraftType, id)); if (!existing) return false; + // A posted draft without a marker predates the markers. if (existing.posted) return true; return new Date(existing.generatedAt) > new Date(updatedAt); } diff --git a/web/lib/content-watch.test.ts b/web/lib/content-watch.test.ts index 6da9c0d708..9c40b073dd 100644 --- a/web/lib/content-watch.test.ts +++ b/web/lib/content-watch.test.ts @@ -6,7 +6,7 @@ vi.mock("./community-agent", async (importOriginal) => { }); import { runLinkCheck, runSemanticDrift, watchDraftId } from "./content-watch"; -import { agentChat, draftStorageKey, type AgentDraft } from "./community-agent"; +import { agentChat, deleteDraft, draftStorageKey, markDraftResolved, type AgentDraft } from "./community-agent"; const agentChatMock = agentChat as Mock; @@ -157,6 +157,19 @@ describe("runSemanticDrift draft identity", () => { expect(draftEntries(kv)).toHaveLength(1); }); + it("does not count a finding the maintainer already discarded as drafted", async () => { + mockDrifts([finding]); + await runSemanticDrift(env()); + const [[key, draft]] = draftEntries(kv); + await markDraftResolved(kv, "semantic-drift", draft.id, "discarded"); + await deleteDraft(kv, key); + + const second = await runSemanticDrift(env()); + + expect(second).toEqual({ ok: true, drafted: 0 }); + expect(draftEntries(kv)).toHaveLength(0); + }); + it("creates a new draft when the finding's evidence changes", async () => { mockDrifts([finding]); await runSemanticDrift(env()); diff --git a/web/lib/content-watch.ts b/web/lib/content-watch.ts index 7c168784ec..46680917dc 100644 --- a/web/lib/content-watch.ts +++ b/web/lib/content-watch.ts @@ -363,8 +363,8 @@ ${docsText}`; const existing = await getDraft(env.CURATED_KV, draftStorageKey(draft)); if (existing) continue; - await saveDraft(env.CURATED_KV, draft); - drafted++; + // saveDraft refuses findings the maintainer already discarded. + if (await saveDraft(env.CURATED_KV, draft)) drafted++; } return { ok: true, drafted }; From 989081cb832a30949cd5a72c09979b8b9e4da1d3 Mon Sep 17 00:00:00 2001 From: Hunter B Date: Tue, 29 Sep 2026 11:26:21 -0700 Subject: [PATCH 3/8] fix(web): bind digest approval to the reviewed draft; exclusive per-request claims Independent review of 1d83478bd held two P1s. Approval binding: approveDigestRecord took only weekId/lang and approved whatever record sat at the key, so draft B + record A (B's record put failed after its draft put) posted B to GitHub and published A. The route also reloaded and posted current storage, so a draft regenerated after the admin page loaded was posted unseen. - The admin page sends reviewedSha256 (SHA-256 of the text shown in the selected language) with every post/discard; the route answers 409 before any claim or GitHub call when the stored draft no longer hashes to it. - approveDigestRecord(kv, draft, lang) publishes only a well-formed record of the same weekId and generatedAt that re-renders (renderDigestBody, now the single renderer shared with runDigest) to exactly the reviewed text, copying fields explicitly in one put. Otherwise it stays unpublished and the maintainer gets a mismatch/missing warning. Claim serialization: the token claim overwrote one shared marker key, so two handlers that both passed the absence check could each write and read back their own token and both post. Claims now live in per-request keys (draft-claim:::); a request proceeds only if a list shows its key alone and no decision marker exists, otherwise it deletes its key and answers 409. Overlapping claimants can both refuse, never both proceed. Discard takes the same claim, so a discard cannot overlap a post. The claim is released on a definite GitHub rejection or once the posted marker is written; on an unknown outcome it is kept until it expires. The code comment states the limit: this relies on a list seeing finished puts, which Workers KV guarantees within one location, not across locations. Tests: community-agent-review-state, content-watch, community-agent-security, public-api-security: 4 files, 51 passed. Against the previous source 6 of the new tests fail (draft B/record A publication, explicit-field publish, barrier race post/post, post/discard overlap, stale reviewed text, missing hash). tsc --noEmit exit 0; eslint on touched files exit 0. Co-Authored-By: Claude Opus 5.5 (1M context) Claude-Session: https://claude.ai/code/session_014ZwqatxgVFxHvovngywnks --- web/app/[locale]/admin/admin-client.tsx | 16 +- web/app/api/admin/post/route.ts | 64 ++++- web/lib/community-agent-review-state.test.ts | 232 +++++++++++++++++-- web/lib/community-agent-tasks.ts | 5 +- web/lib/community-agent.ts | 166 ++++++++++--- 5 files changed, 410 insertions(+), 73 deletions(-) diff --git a/web/app/[locale]/admin/admin-client.tsx b/web/app/[locale]/admin/admin-client.tsx index 58c51c4f54..273069fbae 100644 --- a/web/app/[locale]/admin/admin-client.tsx +++ b/web/app/[locale]/admin/admin-client.tsx @@ -1,7 +1,7 @@ "use client"; import { useState } from "react"; -import { draftStorageKey, type AgentDraft } from "@/lib/community-agent"; +import { draftStorageKey, reviewedBodyHash, type AgentDraft } from "@/lib/community-agent"; interface Props { drafts: AgentDraft[]; @@ -17,13 +17,17 @@ export function AdminClient({ drafts, posted, isZh, typeLabels }: Props) { const [editBody, setEditBody] = useState(""); const [loading, setLoading] = useState(null); - const handleAction = async (draftKey: string, action: "post" | "discard", editedBody?: string) => { + const handleAction = async (draft: AgentDraft, action: "post" | "discard", editedBody?: string) => { + const draftKey = draftStorageKey(draft); setLoading(draftKey); try { + // Binds the action to the text shown here; the route refuses it if the + // stored draft was regenerated since this page loaded. + const reviewedSha256 = await reviewedBodyHash(isZh ? draft.bodyZh : draft.bodyEn); const res = await fetch("/api/admin/post", { method: "POST", headers: { "Content-Type": "application/json" }, - body: JSON.stringify({ action, draftKey, editedBody, lang: isZh ? "zh" : "en" }), + body: JSON.stringify({ action, draftKey, editedBody, lang: isZh ? "zh" : "en", reviewedSha256 }), }); const data = await res.json(); if (data.ok) { @@ -93,7 +97,7 @@ export function AdminClient({ drafts, posted, isZh, typeLabels }: Props) { />