diff --git a/CHANGES.md b/CHANGES.md index 815fa5a1..f7a384d5 100644 --- a/CHANGES.md +++ b/CHANGES.md @@ -67,6 +67,12 @@ To be released. successor account at the top of the profile card and muting the old profile's visual treatment. + - Limited pending poll notification messages to 100 across the shared task + queue, including delayed wakeups and retries. Repeated schedules coalesce, + earlier wakeups take priority at the limit, and deferred work is recovered + after expiry. Delayed poll schedules no longer allocate worker timers. + [[#653], [#663]] + - Replaced the undocumented `SEONBI_URL` integration with the embedded [Gukhanmun] Node-API binding. Set `GUKHANMUN` to a comma-separated list of locale patterns, such as `ko,ko-*`, to add Hangul readings to Hanja in @@ -155,7 +161,9 @@ To be released. [#648]: https://github.com/fedify-dev/hollo/pull/648 [#649]: https://github.com/fedify-dev/hollo/pull/649 [#650]: https://github.com/fedify-dev/hollo/pull/650 +[#653]: https://github.com/fedify-dev/hollo/issues/653 [#660]: https://github.com/fedify-dev/hollo/pull/660 +[#663]: https://github.com/fedify-dev/hollo/pull/663 Version 0.9.22 diff --git a/docs/src/content/docs/install/workers.mdx b/docs/src/content/docs/install/workers.mdx index d0634ce3..ecdde005 100644 --- a/docs/src/content/docs/install/workers.mdx +++ b/docs/src/content/docs/install/workers.mdx @@ -120,6 +120,40 @@ interrupted work remains recoverable. Stop older nodes before upgrading and run compatible task definitions on every node. +The shared database stores at most 100 pending poll messages, including delayed +wakeups and retries, with one wakeup per poll for new schedules. Repeated votes +and expiry updates coalesce: an earlier expiry advances the existing wakeup; +a later expiry is reloaded when that wakeup runs. At the limit, an earlier +wakeup replaces the latest one. Deferred or displaced polls are recovered after +expiry. Bursts can also defer admission when its database lock cannot be +acquired within one second. These limits apply across nodes, rather than per +worker. They bound stored work, not notifications or running handlers. + +Delayed poll wakeups use the queue's five-second polling interval rather than +allocating a timer in each worker. Recovery still pauses at 200 shared ready +messages. The limit of 100 matches one recovery page; it is not a throughput +guarantee. Stop all older producers before upgrading. Legacy queued messages +count toward the limit, but an existing backlog above 100 drains gradually; +new admissions do not increase that backlog. + +The `hollo.poll-notifications.queue` logger reports capacity deferrals, +displaced wakeups and lock timeouts at warning level, at most once per reason +per minute per process. Inspect storage pressure with the query below. +Enqueue counters include coalesced and deferred attempts, so they do not +measure stored messages. + +~~~~ sql +SELECT count(*) AS queued, + count(*) FILTER (WHERE created + delay <= CURRENT_TIMESTAMP) AS ready, + count(*) FILTER (WHERE created + delay > CURRENT_TIMESTAMP) AS delayed, + count(*) FILTER ( + WHERE message->>'type' = 'task' + AND message->>'taskName' = 'hollo.poll-notification.v1' + ) AS pending_polls +FROM hollo_task_message_v1; +~~~~ + + Docker Compose setup -------------------- diff --git a/docs/src/content/docs/ja/install/workers.mdx b/docs/src/content/docs/ja/install/workers.mdx index fede8d37..308d343e 100644 --- a/docs/src/content/docs/ja/install/workers.mdx +++ b/docs/src/content/docs/ja/install/workers.mdx @@ -110,6 +110,37 @@ Holloは6つの主要コンポーネントで構成されています: タスク定義を使ってください。 +共有データベースの待機中投票メッセージは、遅延と再試行を含め最大100件です。 +新しい予約は投票ごとに1件にまとめます。繰り返しの投票や期限更新を統合し、 +期限が早まると既存の予約を前倒しします。期限の延長は実行時に再取得します。 +上限では、早い予約が最も遅い予約を置き換えます。延期・置換された投票は期限後に +復旧します。集中した負荷で入場ロックを1秒以内に取得できない場合も延期します。 +この上限は全ノードで共有し、実行中のハンドラーや通知件数ではなく保存件数を +制限します。 + +遅延予約はワーカーごとのタイマーではなく、キューの5秒間隔の確認を使います。 +復旧は共有キューの即時実行メッセージが200件になると引き続き停止します。 +100件は復旧1ページ分であり、処理能力の保証ではありません。更新前に旧版の +生成ノードをすべて停止してください。旧メッセージも上限に数えますが、既に100件を +超えるバックログは徐々に減少します。新規予約はその件数を増やしません。 + +`hollo.poll-notifications.queue` ロガーは容量による延期、予約の置換、ロックの +タイムアウトを warning レベルで記録します。理由ごと・プロセスごとに毎分最大1回 +です。保存負荷は以下のSQLで確認できます。enqueue カウンターは統合・延期された +試行も含むため、保存メッセージ数ではありません。 + +~~~~ sql +SELECT count(*) AS queued, + count(*) FILTER (WHERE created + delay <= CURRENT_TIMESTAMP) AS ready, + count(*) FILTER (WHERE created + delay > CURRENT_TIMESTAMP) AS delayed, + count(*) FILTER ( + WHERE message->>'type' = 'task' + AND message->>'taskName' = 'hollo.poll-notification.v1' + ) AS pending_polls +FROM hollo_task_message_v1; +~~~~ + + Docker Compose設定 ------------------ diff --git a/docs/src/content/docs/ko/install/workers.mdx b/docs/src/content/docs/ko/install/workers.mdx index 786761f7..3cff4516 100644 --- a/docs/src/content/docs/ko/install/workers.mdx +++ b/docs/src/content/docs/ko/install/workers.mdx @@ -109,6 +109,38 @@ Hollo는 여섯 가지 주요 구성 요소로 이루어져 있습니다: 업그레이드 전 기존 노드를 중지하고 모든 노드에 호환되는 태스크 정의를 배포하세요. +공유 데이터베이스에는 지연 wakeup과 재시도를 포함해 대기 투표 메시지를 최대 +100개 저장하며, 새 예약은 투표당 하나의 wakeup으로 합칩니다. 반복 참여와 만료 +변경도 합쳐지고, 만료가 앞당겨지면 기존 wakeup을 앞당깁니다. 만료가 늦춰지면 +wakeup 실행 시 현재 만료 시각을 다시 읽습니다. 상한에서는 더 이른 wakeup이 +가장 늦은 메시지를 대체합니다. 보류되거나 밀려난 투표는 만료 후 복구됩니다. 부하가 +몰려 입장 잠금을 1초 안에 얻지 못해도 예약을 보류할 수 있습니다. 이 상한은 +워커별이 아니라 모든 노드에 걸쳐 적용되며, 실행 중 handler나 알림 수가 아닌 +저장된 작업 수를 제한합니다. + +지연 예약은 워커마다 타이머를 할당하는 대신 큐의 5초 주기 확인을 사용합니다. +복구는 공유 ready 메시지가 200개이면 여전히 일시 중지합니다. 100개 상한은 +복구 한 페이지 크기이며 처리량 보장은 아닙니다. 업그레이드 전에 모든 기존 +생산자를 중지하세요. 기존 메시지도 상한에 포함되지만 이미 100개를 넘은 backlog는 +점차 줄어듭니다. 새 예약으로 그 backlog가 더 커지지는 않습니다. + +`hollo.poll-notifications.queue` 로거는 상한에 따른 보류, wakeup 대체와 잠금 +시간 초과를 warning 수준으로 원인별·프로세스별 매분 최대 한 번 기록합니다. +아래 SQL로 저장 부하를 확인할 수 있습니다. enqueue 카운터에는 합쳐지거나 +보류된 시도도 포함되므로 저장 메시지 수를 나타내지 않습니다. + +~~~~ sql +SELECT count(*) AS queued, + count(*) FILTER (WHERE created + delay <= CURRENT_TIMESTAMP) AS ready, + count(*) FILTER (WHERE created + delay > CURRENT_TIMESTAMP) AS delayed, + count(*) FILTER ( + WHERE message->>'type' = 'task' + AND message->>'taskName' = 'hollo.poll-notification.v1' + ) AS pending_polls +FROM hollo_task_message_v1; +~~~~ + + Docker Compose 설정 ------------------- diff --git a/docs/src/content/docs/zh-cn/install/workers.mdx b/docs/src/content/docs/zh-cn/install/workers.mdx index 9680b926..6fb574be 100644 --- a/docs/src/content/docs/zh-cn/install/workers.mdx +++ b/docs/src/content/docs/zh-cn/install/workers.mdx @@ -95,6 +95,33 @@ Hollo由六个主要组件组成: 事务结束,中断的工作可恢复。升级前停止旧节点,并在所有节点使用兼容的任务定义。 +共享数据库最多保存 100 条待处理投票消息,包括延迟唤醒和重试。新预约按每个 +投票合并为一条消息。重复投票和到期更新会合并;到期提前时,现有唤醒也提前; +到期延后时,执行时重新读取。达到上限后,更早的唤醒替换最晚的消息。被延期或 +替换的投票在到期后恢复。突发负载导致入队锁无法在一秒内取得时也可能延期。 +上限由所有节点共享,而非按工作器计算,限制保存的任务数而非通知或运行的处理器。 + +延迟唤醒使用队列每五秒的轮询,不为每个工作器分配计时器。共享就绪消息达到 +200 条时,恢复仍会暂停。100 条上限等于一页恢复批次,并不保证吞吐量。升级前 +停止所有旧的生产节点。旧消息也计入上限,但已超过 100 条的积压会逐渐减少; +新预约不会增加该积压。 + +`hollo.poll-notifications.queue` 日志以 warning 级别报告容量延期、唤醒替换和 +锁超时,每个原因、每个进程每分钟最多一次。用下面的 SQL 观察存储压力。入队计数 +包含合并及延期的尝试,不代表保存消息数。 + +~~~~ sql +SELECT count(*) AS queued, + count(*) FILTER (WHERE created + delay <= CURRENT_TIMESTAMP) AS ready, + count(*) FILTER (WHERE created + delay > CURRENT_TIMESTAMP) AS delayed, + count(*) FILTER ( + WHERE message->>'type' = 'task' + AND message->>'taskName' = 'hollo.poll-notification.v1' + ) AS pending_polls +FROM hollo_task_message_v1; +~~~~ + + Docker Compose设置 ------------------ diff --git a/docs/src/content/docs/zh-tw/install/workers.mdx b/docs/src/content/docs/zh-tw/install/workers.mdx index af964284..dce61bad 100644 --- a/docs/src/content/docs/zh-tw/install/workers.mdx +++ b/docs/src/content/docs/zh-tw/install/workers.mdx @@ -95,6 +95,33 @@ Hollo 由六個主要元件組成: 交易結束,中斷的工作可復原。升級前停止舊節點,並在所有節點使用相容的任務定義。 +共享資料庫最多儲存 100 條待處理投票訊息,包括延遲喚醒與重試。新預約按每個 +投票合併為一條訊息。重複投票與到期更新會合併;到期提前時,現有喚醒也提前; +到期延後時,執行時重新讀取。達到上限後,更早的喚醒取代最晚的訊息。被延後或 +取代的投票在到期後復原。突發負載導致入列鎖無法在一秒內取得時也可能延後。 +上限由所有節點共享,而非按工作節點計算,限制儲存的任務數而非通知或執行的處理器。 + +延遲喚醒使用佇列每五秒的輪詢,不為每個工作節點分配計時器。共享就緒訊息達到 +200 條時,復原仍會暫停。100 條上限等於一頁復原批次,並不保證吞吐量。升級前 +停止所有舊的生產節點。舊訊息也計入上限,但已超過 100 條的積壓會逐漸減少; +新預約不會增加該積壓。 + +`hollo.poll-notifications.queue` 日誌以 warning 級別報告容量延後、喚醒取代與 +鎖逾時,每個原因、每個程序每分鐘最多一次。用下面的 SQL 觀察儲存壓力。入列計數 +包含合併及延後的嘗試,不代表儲存訊息數。 + +~~~~ sql +SELECT count(*) AS queued, + count(*) FILTER (WHERE created + delay <= CURRENT_TIMESTAMP) AS ready, + count(*) FILTER (WHERE created + delay > CURRENT_TIMESTAMP) AS delayed, + count(*) FILTER ( + WHERE message->>'type' = 'task' + AND message->>'taskName' = 'hollo.poll-notification.v1' + ) AS pending_polls +FROM hollo_task_message_v1; +~~~~ + + Docker Compose 設定 ------------------- diff --git a/src/background/poll-queue.test.ts b/src/background/poll-queue.test.ts new file mode 100644 index 00000000..afe24d69 --- /dev/null +++ b/src/background/poll-queue.test.ts @@ -0,0 +1,701 @@ +import { + createFederation, + MemoryKvStore, + type Message, + type MessageQueueEnqueueOptions, +} from "@fedify/fedify"; +import { Temporal } from "@js-temporal/polyfill"; +import { eq } from "drizzle-orm"; +import { + afterAll, + afterEach, + beforeEach, + describe, + expect, + it, + vi, +} from "vitest"; + +import { cleanDatabase } from "../../tests/helpers"; +import { createAccount } from "../../tests/helpers/oauth"; +import { createExpiredPollPost } from "../../tests/helpers/poll"; +import * as cleanupWorker from "../cleanup/worker"; +import db, { postgres } from "../db"; +import { registerRemoteReplyScrapes } from "../federation/replies-tasks"; +import * as importWorker from "../import/worker"; +import { registerPollNotifications } from "../poll-notification-tasks"; +import * as schema from "../schema"; +import { type Uuid, uuidv7 } from "../uuid"; +import { registerBackgroundJobs } from "./jobs"; +import { itemLeases } from "./lease"; +import { + POLL_ADMISSION_LOCK, + POLL_QUEUE_LIMIT, + PollMessageQueue, +} from "./poll-queue"; +import { TaskMessageQueue } from "./queue"; + +const tableName = "hollo_poll_queue_test"; +const channelName = "hollo_poll_queue_test"; +const taskName = "hollo.poll-notification.v1"; +const key = (id: string) => `hollo.poll-notification:${id}`; +const delay = (seconds: number): MessageQueueEnqueueOptions => ({ + // The installed polyfills share the queue's total()/toString() interface. + delay: Temporal.Duration.from({ + milliseconds: seconds * 1000, + }) as unknown as NonNullable, +}); + +function fixture(raw = false) { + const queue = new PollMessageQueue(postgres, { + tableName, + channelName, + pollInterval: { milliseconds: 20 }, + handlerTimeout: { seconds: 0 }, + }); + const source = raw ? queue.queue : queue; + const federation = createFederation({ + kv: new MemoryKvStore(), + queue: { task: new TaskMessageQueue(source) }, + taskQueueResolution: "strict", + manuallyStartQueue: true, + }); + let now: Date | undefined; + const tasks = registerPollNotifications( + federation, + async () => (await source.getDepth()).ready!, + { clock: () => now ?? new Date() }, + ); + const ctx = federation.createContext( + new URL("https://hollo.test"), + undefined, + ); + return { + queue, + federation, + tasks, + ctx, + setNow(value: Date) { + now = value; + }, + async run() { + const [row] = await postgres` + DELETE FROM ${postgres(tableName)} WHERE id = ( + SELECT id FROM ${postgres(tableName)} ORDER BY created LIMIT 1 + ) RETURNING message + `; + await federation.processQueuedTask(undefined, row.message); + }, + }; +} + +async function seed(expires = new Date(Date.now() + 60_000), author?: Uuid) { + const authorId = author ?? ((await createAccount()).id as Uuid); + return { + ...(await createExpiredPollPost(authorId, expires)), + authorId, + expires, + }; +} + +async function rows() { + return await postgres` + SELECT *, (created + delay)::text AS due, xmin::text AS version, + extract(day FROM delay) AS days, extract(month FROM delay) AS months, + jsonb_typeof(message) AS shape, + extract(microseconds FROM created)::integer % 1000 AS micros + FROM ${postgres(tableName)} ORDER BY created + `; +} + +async function notifications() { + return await db.query.notifications.findMany({ + where: { type: { eq: "poll" } }, + }); +} + +function envelope( + id = uuidv7(), + orderingKey: string | undefined = key(id), +): Extract { + return { + type: "task", + id: crypto.randomUUID(), + taskName, + data: "unused", + baseUrl: "https://hollo.test", + started: new Date().toISOString(), + attempt: 0, + orderingKey, + traceContext: {}, + }; +} + +beforeEach(cleanDatabase); +afterEach(() => vi.restoreAllMocks()); +afterAll(() => itemLeases.close()); + +describe("PostgreSQL poll admission", () => { + it("measures repeated scheduling against the raw backend and emits no delayed NOTIFYs", async () => { + const p = await seed(); + const raw = fixture(true); + let notices = 0; + const listener = await postgres.listen(channelName, () => { + notices++; + }); + try { + for (let i = 0; i < 1000; i++) await raw.tasks.enqueue(raw.ctx, p.pollId); + expect(await raw.queue.getPollDepth()).toBe(1000); + await vi.waitFor(() => expect(notices).toBe(1000)); + await postgres`TRUNCATE ${postgres(tableName)}`; + notices = 0; + const bounded = fixture(); + const started = performance.now(); + for (let i = 0; i < 1000; i++) + await bounded.tasks.enqueue(bounded.ctx, p.pollId); + expect(await bounded.queue.getDepth()).toMatchObject({ + queued: 1, + ready: 0, + delayed: 1, + }); + expect(notices).toBe(0); + console.info("poll queue repeated schedules", { + attempts: 1000, + rawDepth: 1000, + boundedDepth: 1, + rawDelayedNotifies: 1000, + boundedDelayedNotifies: notices, + boundedMilliseconds: Math.round(performance.now() - started), + }); + } finally { + await listener.unlisten(); + } + }, 30_000); + + it("coalesces votes and expiry edits across producers, sampling both bounds", async () => { + const a = fixture(); + const b = fixture(); + const p = await seed(); + const voter = await createAccount({ username: "voter" }); + await db.insert(schema.pollVotes).values({ + pollId: p.pollId, + accountId: voter.id as Uuid, + optionIndex: 0, + }); + await a.tasks.enqueue(a.ctx, p.pollId); + let maximum = 0; + for (let batch = 0; batch < 20; batch++) { + await db + .update(schema.polls) + .set({ expires: new Date(Date.now() + (batch % 2 ? 120_000 : 60_000)) }) + .where(eq(schema.polls.id, p.pollId)); + await Promise.all( + Array.from({ length: 50 }, (_, i) => + (i % 2 ? a : b).tasks.enqueue(a.ctx, p.pollId), + ), + ); + const stored = await rows(); + maximum = Math.max(maximum, stored.length); + expect(stored).toHaveLength(1); + expect(stored[0].ordering_key).toBeNull(); + } + await db + .update(schema.polls) + .set({ expires: new Date(0) }) + .where(eq(schema.polls.id, p.pollId)); + await b.tasks.enqueue(b.ctx, p.pollId); + await a.run(); + expect(await notifications()).toHaveLength(2); + expect( + (await db.query.notificationGroups.findMany()).map( + (g) => g.notificationsCount, + ), + ).toEqual([1, 1]); + console.info("poll queue concurrent edits/votes", { + attempts: 1000, + producers: 2, + maximum, + notifications: 2, + }); + }, 30_000); + + it("bounds distinct polls and retries globally across concurrent producers", async () => { + const a = fixture(); + const b = fixture(); + let maximum = 0; + for (let batch = 0; batch < 15; batch++) { + await Promise.all( + Array.from({ length: 20 }, (_, i) => { + const message = { ...envelope(), attempt: i % 3 }; + return (i % 2 ? a : b).queue.enqueue(message, delay(3600)); + }), + ); + maximum = Math.max(maximum, await a.queue.getPollDepth()); + expect(maximum).toBeLessThanOrEqual(POLL_QUEUE_LIMIT); + const [duplicates] = await postgres` + SELECT count(*) AS count FROM ( + SELECT message->>'orderingKey' FROM ${postgres(tableName)} + GROUP BY message->>'orderingKey' HAVING count(*) > 1 + ) duplicates + `; + expect(Number(duplicates.count)).toBe(0); + } + expect(maximum).toBe(POLL_QUEUE_LIMIT); + console.info("poll queue distinct concurrent schedules", { + attempts: 300, + producers: 2, + maximum, + }); + }, 15_000); + + it("preserves envelopes, microseconds, FIFO priority and no-op rows", async () => { + const f = fixture(); + const messages = Array.from({ length: 10 }, () => envelope()); + for (const message of messages) + await f.queue.enqueue(message, delay(0.125)); + const before = await rows(); + expect(before.some((row) => Number(row.micros) !== 0)).toBe(true); + expect(before.map((row) => row.message)).toEqual(messages); + expect( + before.every( + (row) => + row.shape === "object" && + row.ordering_key === null && + Number(row.days) === 0 && + Number(row.months) === 0, + ), + ).toBe(true); + for (const message of messages) await f.queue.enqueue(message, delay(100)); + expect((await rows()).map((row) => row.version)).toEqual( + before.map((row) => row.version), + ); + await f.queue.queue.enqueue(messages[0], delay(0.125)); + expect((await rows()).at(-1)?.message).toEqual(before[0].message); + const columns = await postgres` + SELECT column_name FROM information_schema.columns + WHERE table_name = ${tableName} ORDER BY ordinal_position + `; + expect(columns.map((row) => row.column_name)).toEqual([ + "id", + "message", + "delay", + "created", + "ordering_key", + ]); + }); + + it("pulls earlier expiry to now once, retains FIFO time and seconds-only delays", async () => { + const f = fixture(); + const p = await seed(); + await f.tasks.enqueue(f.ctx, p.pollId); + const original = (await rows())[0]; + let notices = 0; + const listener = await postgres.listen(channelName, () => { + notices++; + }); + try { + await db + .update(schema.polls) + .set({ expires: new Date(0) }) + .where(eq(schema.polls.id, p.pollId)); + await f.tasks.enqueue(f.ctx, p.pollId); + await vi.waitFor(() => expect(notices).toBe(1)); + const pulled = (await rows())[0]; + expect(pulled.created).toEqual(original.created); + expect(Number(pulled.days)).toBe(0); + expect(Number(pulled.months)).toBe(0); + await f.tasks.enqueue(f.ctx, p.pollId); + expect((await rows())[0].version).toBe(pulled.version); + await f.run(); + expect(await notifications()).toHaveLength(1); + expect(notices).toBe(1); + await f.queue.enqueue(envelope()); + await vi.waitFor(() => expect(notices).toBe(2)); + } finally { + await listener.unlisten(); + } + }); + + it("admits expired work by replacing only later polls, even with legacy overload", async () => { + const f = fixture(); + for (let i = 0; i < 120; i++) + await f.queue.queue.enqueue( + { ...envelope(), orderingKey: undefined }, + delay(604800), + ); + const p = await seed(new Date("2026-01-01")); + await f.tasks.recover(f.ctx); + expect(await f.queue.getPollDepth()).toBe(120); + const [row] = await postgres` + DELETE FROM ${postgres(tableName)} WHERE message->>'orderingKey' = ${key(p.pollId)} RETURNING message + `; + expect(row).toBeDefined(); + await f.federation.processQueuedTask(undefined, row.message); + expect(await notifications()).toHaveLength(1); + expect(await f.queue.getPollDepth()).toBe(119); + // Consumption drains the legacy surplus; admissions never add to it. + await postgres`DELETE FROM ${postgres(tableName)} WHERE id IN (SELECT id FROM ${postgres(tableName)} LIMIT 20)`; + await f.queue.enqueue(envelope(), delay(1)); + expect(await f.queue.getPollDepth()).toBe(POLL_QUEUE_LIMIT); + }); + + it("defers later work at the cap and rolls eviction back when insertion fails", async () => { + const f = fixture(); + for (let i = 0; i < POLL_QUEUE_LIMIT; i++) + await f.queue.enqueue(envelope(), delay(3600)); + const before = (await rows()).map((row) => row.id); + await f.queue.enqueue(envelope(), delay(7200)); + expect((await rows()).map((row) => row.id)).toEqual(before); + const rejected = { ...envelope(), data: "reject" }; + let notices = 0; + const listener = await postgres.listen(channelName, () => { + notices++; + }); + await postgres.unsafe(`CREATE OR REPLACE FUNCTION hollo_poll_queue_reject() RETURNS trigger LANGUAGE plpgsql AS $$ + BEGIN IF NEW.message->>'data' = 'reject' THEN RAISE EXCEPTION 'forced insert failure'; END IF; RETURN NEW; END $$`); + await postgres.unsafe( + `CREATE TRIGGER reject_poll BEFORE INSERT ON "${tableName}" FOR EACH ROW EXECUTE FUNCTION hollo_poll_queue_reject()`, + ); + try { + await expect(f.queue.enqueue(rejected)).rejects.toThrow( + "forced insert failure", + ); + expect((await rows()).map((row) => row.id)).toEqual(before); + expect(notices).toBe(0); + } finally { + await postgres.unsafe(`DROP TRIGGER reject_poll ON "${tableName}"`); + await postgres.unsafe("DROP FUNCTION hollo_poll_queue_reject()"); + await listener.unlisten(); + } + }); + + it("defers a held admission lock within a second and recovers the missed enqueue", async () => { + const f = fixture(); + const p = await seed(new Date("2026-01-01")); + await f.queue.queue.initialize(); + const held = await postgres.reserve(); + await held`SELECT pg_advisory_lock(${POLL_ADMISSION_LOCK}::bigint)`; + try { + const started = performance.now(); + await f.tasks.enqueue(f.ctx, p.pollId); + expect(performance.now() - started).toBeGreaterThanOrEqual(900); + expect(performance.now() - started).toBeLessThan(2500); + expect(await f.queue.getPollDepth()).toBe(0); + } finally { + await held`SELECT pg_advisory_unlock(${POLL_ADMISSION_LOCK}::bigint)`; + held.release(); + } + await f.tasks.recover(f.ctx); + await f.run(); + expect(await notifications()).toHaveLength(1); + }); + + it("recovers 250 expired polls through bounded pages without duplicate counts", async () => { + const f = fixture(); + const author = (await createAccount()).id as Uuid; + for (let i = 0; i < 250; i++) + await seed(new Date(1767225600000 + i), author); + let passes = 0; + let maximum = 0; + while ((await notifications()).length < 250 && passes < 10) { + await f.tasks.recover(f.ctx); + const depth = await f.queue.getPollDepth(); + maximum = Math.max(maximum, depth); + expect(depth).toBeLessThanOrEqual(POLL_QUEUE_LIMIT); + for (let i = 0; i < depth; i++) await f.run(); + passes++; + } + expect(await notifications()).toHaveLength(250); + expect( + (await db.query.notificationGroups.findMany()).every( + (g) => g.notificationsCount === 1, + ), + ).toBe(true); + console.info("poll queue recovery backlog", { + polls: 250, + passes, + maximum, + }); + }, 30_000); + + it("does not lose wakeups when consumption races with coalescing", async () => { + const a = fixture(); + const b = fixture(); + const p = await seed(); + await a.tasks.enqueue(a.ctx, p.pollId); + for (let i = 0; i < 50; i++) { + await Promise.all([ + b.tasks.enqueue(b.ctx, p.pollId), + postgres`DELETE FROM ${postgres(tableName)} WHERE message->>'orderingKey' = ${key(p.pollId)}`, + ]); + expect(await a.queue.getPollDepth()).toBeLessThanOrEqual(1); + // A lost handler is repaired by a subsequent schedule/recovery, rather + // than suppressed by a reservation left behind after consumption. + await a.tasks.enqueue(a.ctx, p.pollId); + expect(await a.queue.getPollDepth()).toBe(1); + } + await db + .update(schema.polls) + .set({ expires: new Date("2026-01-01") }) + .where(eq(schema.polls.id, p.pollId)); + await b.tasks.enqueue(b.ctx, p.pollId); + await a.run(); + expect(await notifications()).toHaveLength(1); + }); + + it("coalesces with elapsed seconds across a daylight-saving transition", async () => { + const f = fixture(); + const message = envelope(); + await f.queue.enqueue(message, delay(604800)); + // Keep the original FIFO timestamp, spanning the New York spring change. + await postgres` + UPDATE ${postgres(tableName)} SET created = '2026-03-07 12:00:00+00'::timestamptz, + delay = make_interval(secs => extract(epoch FROM + (clock_timestamp() + interval '1 hour' - '2026-03-07 12:00:00+00'::timestamptz))::double precision) + `; + await f.queue.enqueue(message, delay(60)); + const [row] = await postgres.begin(async (tx) => { + await tx`SET LOCAL TIME ZONE 'America/New_York'`; + return await tx` + SELECT extract(day FROM delay) AS days, + extract(month FROM delay) AS months, + extract(epoch FROM (created + delay - clock_timestamp())) AS remaining + FROM ${tx(tableName)} + `; + }); + expect(Number(row.days)).toBe(0); + expect(Number(row.months)).toBe(0); + expect(Number(row.remaining)).toBeGreaterThan(55); + expect(Number(row.remaining)).toBeLessThanOrEqual(60); + }); + + it("keeps retries capped, recovers exhaustion and restarts after consumed-message loss", async () => { + const f = fixture(); + const p = await seed(new Date("2026-01-01")); + await f.tasks.enqueue(f.ctx, p.pollId); + const failure = vi + .spyOn(db, "transaction") + .mockRejectedValue(new Error("notification database offline")); + for (let i = 0; i < 3; i++) { + await f.run(); + expect(await f.queue.getPollDepth()).toBeLessThanOrEqual(1); + } + expect(failure).toHaveBeenCalledTimes(3); + expect(await f.queue.getPollDepth()).toBe(0); + failure.mockRestore(); + await f.tasks.recover(f.ctx); + // PostgreSQL deletes before calling the handler; emulate exit in that gap. + await postgres`DELETE FROM ${postgres(tableName)}`; + const restarted = fixture(); + await restarted.tasks.recover(restarted.ctx); + await restarted.run(); + await restarted.tasks.recover(restarted.ctx); + expect(await restarted.queue.getPollDepth()).toBe(0); + expect(await notifications()).toHaveLength(1); + }); + + it.each(["earlier", "later"])( + "delivers once at the current %s expiry through a real listener", + async (direction) => { + const f = fixture(); + const p = await seed(new Date(Date.now() + 400)); + await f.tasks.enqueue(f.ctx, p.pollId); + const expires = new Date( + Date.now() + (direction === "earlier" ? 100 : 900), + ); + await db + .update(schema.polls) + .set({ expires }) + .where(eq(schema.polls.id, p.pollId)); + await f.tasks.enqueue(f.ctx, p.pollId); + const abort = new AbortController(); + const listen = f.federation.startQueue(undefined, { + queue: "task", + signal: abort.signal, + }); + try { + expect(await notifications()).toHaveLength(0); + await vi.waitFor( + async () => expect(await notifications()).toHaveLength(1), + { timeout: 3000 }, + ); + expect((await notifications())[0].created).toEqual(expires); + expect(Date.now()).toBeGreaterThanOrEqual(+expires); + } finally { + abort.abort(); + await listen; + } + expect(await f.queue.getPollDepth()).toBe(0); + }, + ); + + it("keeps admission bounded with a large unrelated delayed backlog", async () => { + const f = fixture(); + await f.queue.queue.initialize(); + await postgres` + INSERT INTO ${postgres(tableName)} (message, delay) + SELECT '{"type":"task","taskName":"hollo.cleanup-item.v1"}'::jsonb, + interval '1 hour' FROM generate_series(1, 10000) + `; + const started = performance.now(); + await f.queue.enqueue(envelope(), delay(60)); + const milliseconds = Math.round(performance.now() - started); + expect(await f.queue.getPollDepth()).toBe(1); + expect((await f.queue.getDepth()).queued).toBe(10001); + console.info("poll admission with unrelated backlog", { + unrelated: 10000, + milliseconds, + }); + }); + + it("registers the bounded adapter in the real federation queue", async () => { + const { taskQueue } = await import("../federation/federation"); + expect(taskQueue.queue).toBeInstanceOf(PollMessageQueue); + expect("enqueueMany" in taskQueue.queue).toBe(false); + }); + + it.each(["delayed", "ready"])( + "runs imports, cleanup and replies with two consumers and saturated %s poll work", + async (load) => { + const a = fixture(); + const b = fixture(); + vi.spyOn(importWorker, "executeImportItem").mockResolvedValue(undefined); + vi.spyOn(cleanupWorker, "executeCleanupItem").mockResolvedValue( + undefined, + ); + const author = (await createAccount()).id as Uuid; + // Delayed reservations occupy the cap; expired work displaces them. + for (let i = 0; i < POLL_QUEUE_LIMIT; i++) + await a.queue.enqueue(envelope(), delay(604800)); + const jobs = [a, b].map((f) => + registerBackgroundJobs( + f.federation, + async () => (await f.queue.getDepth()).ready!, + ), + ); + let fetched = 0; + const scrapeOptions = { + maxDepth: 1, + intervalSeconds: 0, + documentLoader: async (url: string) => { + fetched++; + return { + contextUrl: null, + documentUrl: url, + document: { + "@context": "https://www.w3.org/ns/activitystreams", + id: url, + type: "OrderedCollection", + totalItems: 0, + orderedItems: [], + }, + }; + }, + }; + const scrapes = [a, b].map((f) => + registerRemoteReplyScrapes( + f.federation, + async () => (await f.queue.getDepth()).ready!, + scrapeOptions, + ), + ); + const importId = uuidv7(); + const cleanupId = uuidv7(); + await db.insert(schema.importJobs).values({ + id: importId, + accountOwnerId: author, + category: "bookmarks", + totalItems: 1, + }); + await db + .insert(schema.importJobItems) + .values({ id: uuidv7(), jobId: importId, data: {} }); + await db.insert(schema.cleanupJobs).values({ + id: cleanupId, + category: "cleanup_thumbnails", + totalItems: 1, + }); + await db.insert(schema.cleanupJobItems).values({ + id: uuidv7(), + jobId: cleanupId, + data: { kind: "enumerate_proxy_cache" }, + }); + const p = await seed(new Date("2026-01-01"), author); + const scrapeId = uuidv7(); + await db + .insert(schema.remoteReplyScrapeOrigins) + .values({ originHost: "remote.test", nextRequestAt: new Date(0) }); + await db.insert(schema.remoteReplyScrapeJobs).values({ + id: scrapeId, + postId: p.postId, + postIri: `https://hollo.test/posts/${p.postId}`, + repliesIri: "https://remote.test/replies", + baseUrl: "https://hollo.test", + originHost: "remote.test", + nextAttemptAt: new Date(0), + nextDispatchAt: new Date(0), + }); + await jobs[0].enqueueJob(a.ctx, "import", importId); + await jobs[1].enqueueJob(b.ctx, "cleanup", cleanupId); + await scrapes[0].enqueue(a.ctx, scrapeId); + if (load === "ready") { + for (let i = 0; i < POLL_QUEUE_LIMIT - 1; i++) { + const due = await seed(new Date("2026-01-01"), author); + await b.tasks.enqueue(b.ctx, due.pollId); + } + } + await a.tasks.enqueue(a.ctx, p.pollId); + const expectedNotifications = load === "ready" ? POLL_QUEUE_LIMIT : 1; + const startDepth = await a.queue.getDepth(); + expect(startDepth.queued).toBe(103); + expect(await notifications()).toHaveLength(0); // web-only registration/enqueue + const abort = new AbortController(); + const started = performance.now(); + const workers = [a, b].map((f) => + f.federation.startQueue(undefined, { + queue: "task", + signal: abort.signal, + }), + ); + try { + await vi.waitFor( + async () => { + expect( + ( + await db.query.importJobs.findFirst({ + where: { id: { eq: importId } }, + }) + )?.status, + ).toBe("completed"); + expect( + ( + await db.query.cleanupJobs.findFirst({ + where: { id: { eq: cleanupId } }, + }) + )?.status, + ).toBe("completed"); + expect( + (await db.query.remoteReplyScrapeJobs.findFirst())?.status, + ).toBe("completed"); + expect(await notifications()).toHaveLength(expectedNotifications); + }, + { timeout: 5000 }, + ); + expect(importWorker.executeImportItem).toHaveBeenCalledTimes(1); + expect(cleanupWorker.executeCleanupItem).toHaveBeenCalledTimes(1); + expect(fetched).toBe(1); + expect(await a.queue.getPollDepth()).toBe(load === "ready" ? 0 : 99); + console.info("poll queue mixed workload", { + consumers: 2, + startDepth, + endDepth: await a.queue.getDepth(), + milliseconds: Math.round(performance.now() - started), + }); + } finally { + abort.abort(); + await Promise.all(workers); + } + }, + 15_000, + ); +}); diff --git a/src/background/poll-queue.ts b/src/background/poll-queue.ts new file mode 100644 index 00000000..7ac543d0 --- /dev/null +++ b/src/background/poll-queue.ts @@ -0,0 +1,180 @@ +import type { + MessageQueue, + MessageQueueEnqueueOptions, + MessageQueueListenOptions, +} from "@fedify/fedify"; +import { + PostgresMessageQueue, + type PostgresMessageQueueOptions, +} from "@fedify/postgres"; +import { getLogger } from "@logtape/logtape"; +import type { Sql } from "postgres"; + +const logger = getLogger(["hollo", "poll-notifications", "queue"]); +const TASK_NAME = "hollo.poll-notification.v1"; +export const POLL_QUEUE_LIMIT = 100; +// Separate from the two-int keyspace used by Fedify's ordering locks and +// from Hollo's passkey lock (7626128400). +export const POLL_ADMISSION_LOCK = 7626128653; + +/** + * Bounds poll messages in Fedify's PostgreSQL task table, including retries. + * Queue rows are the reservation: consumption and rollback cannot leave a + * detached marker suppressing a needed wakeup. Other workloads pass through. + */ +export class PollMessageQueue implements MessageQueue { + readonly queue: PostgresMessageQueue; + readonly tableName: string; + readonly channelName: string; + private readonly warnings = new Map(); + + constructor( + private readonly sql: Sql, + options: PostgresMessageQueueOptions = {}, + ) { + this.tableName = options.tableName ?? "hollo_task_message_v1"; + this.channelName = options.channelName ?? "hollo_task_channel_v1"; + this.queue = new PostgresMessageQueue(sql, { + ...options, + tableName: this.tableName, + channelName: this.channelName, + }); + } + + getDepth() { + return this.queue.getDepth(); + } + + async getPollDepth(): Promise { + await this.queue.initialize(); + const [row] = await this.sql` + SELECT count(*) AS count FROM ${this.sql(this.tableName)} + WHERE message->>'type' = 'task' AND message->>'taskName' = ${TASK_NAME} + `; + return Number(row.count); + } + + listen( + handler: (message: unknown) => void | Promise, + options?: MessageQueueListenOptions, + ) { + return this.queue.listen(handler, options); + } + + private warn(reason: string) { + const now = Date.now(); + if (now - (this.warnings.get(reason) ?? -Infinity) < 60_000) return; + this.warnings.set(reason, now); + logger.warning( + "Poll queue admission {reason}; displaced or deferred polls remain " + + "recoverable after expiry (stored limit {limit}).", + { reason, limit: POLL_QUEUE_LIMIT }, + ); + } + + async enqueue(message: unknown, options?: MessageQueueEnqueueOptions) { + if ( + message == null || + typeof message !== "object" || + !("type" in message) || + message.type !== "task" || + !("taskName" in message) || + message.taskName !== TASK_NAME + ) { + await this.queue.enqueue(message, options); + return; + } + await this.queue.initialize(); + const key = + "orderingKey" in message && typeof message.orderingKey === "string" + ? message.orderingKey + : undefined; + const seconds = Math.max(0, options?.delay?.total("seconds") ?? 0); + try { + const reason = await this.sql.begin( + "isolation level read committed", + async (tx) => { + await tx`SET LOCAL lock_timeout = '1s'`; + // Take the lock in its own statement: subsequent statements must + // see producers that committed while we waited for it. + await tx`SELECT pg_advisory_xact_lock(${POLL_ADMISSION_LOCK}::bigint)`; + const [clock] = await tx` + WITH clock AS (SELECT clock_timestamp() AS now) + SELECT now::text AS now, + (now + make_interval(secs => ${seconds}::double precision))::text AS due + FROM clock + `; + if (key != null) { + const [existing] = await tx` + SELECT id, + created + delay > ${clock.due}::text::timestamptz AS later, + created + delay > ${clock.now}::text::timestamptz AS delayed + FROM ${tx(this.tableName)} + WHERE message->>'type' = 'task' AND message->>'taskName' = ${TASK_NAME} + AND message->>'orderingKey' = ${key} + LIMIT 1 FOR UPDATE + `; + if (existing) { + if (existing.later) { + await tx` + UPDATE ${tx(this.tableName)} + SET delay = make_interval(secs => extract(epoch FROM + (${clock.due}::text::timestamptz - created))::double precision) + WHERE id = ${existing.id} + `; + if (existing.delayed && seconds === 0) + await tx`SELECT pg_notify(${this.channelName}, 'PT0S')`; + } + return; + } + } + const [depth] = await tx` + SELECT count(*) AS count FROM ${tx(this.tableName)} + WHERE message->>'type' = 'task' AND message->>'taskName' = ${TASK_NAME} + `; + const full = Number(depth.count) >= POLL_QUEUE_LIMIT; + if (full) { + const removed = await tx` + WITH latest AS ( + SELECT id FROM ${tx(this.tableName)} + WHERE message->>'type' = 'task' AND message->>'taskName' = ${TASK_NAME} + AND created + delay > ${clock.due}::text::timestamptz + ORDER BY created + delay DESC LIMIT 1 FOR UPDATE + ) + DELETE FROM ${tx(this.tableName)} + WHERE id IN (SELECT id FROM latest) RETURNING id + `; + if (removed.length === 0) return "deferred-capacity"; + } + // Explicit text casts bypass postgres.js's JSON/date parameter + // serializers. Elapsed seconds avoid calendar-day/DST arithmetic. + // Keep the ordering column null to preserve the backend's FIFO + // priority; the envelope key is only a coalescing identity. + await tx` + INSERT INTO ${tx(this.tableName)} (message, created, delay, ordering_key) + VALUES (${JSON.stringify(message)}::text::jsonb, + ${clock.now}::text::timestamptz, + make_interval(secs => ${seconds}::double precision), NULL) + `; + // Delayed NOTIFYs allocate a timer in every worker. Polling discovers + // future rows within the backend's five-second interval instead. + if (seconds === 0) + await tx`SELECT pg_notify(${this.channelName}, 'PT0S')`; + return full ? "evicted-later-wakeup" : undefined; + }, + ); + if (reason) this.warn(reason); + } catch (error) { + if ( + error != null && + typeof error === "object" && + "code" in error && + error.code === "55P03" + ) { + this.warn("lock-timeout"); + return; + } + throw error; + } + } +} diff --git a/src/federation/federation.ts b/src/federation/federation.ts index 2a603844..1ae4310b 100644 --- a/src/federation/federation.ts +++ b/src/federation/federation.ts @@ -17,6 +17,7 @@ import { import metadata from "../../package.json" with { type: "json" }; import { registerBackgroundJobs } from "../background/jobs"; +import { PollMessageQueue } from "../background/poll-queue"; import { TaskMessageQueue } from "../background/queue"; import { postgres } from "../db"; import { FEDIFY_ORIGIN } from "../env"; @@ -46,9 +47,7 @@ const activityQueue = new ParallelMessageQueue( 10, ); export const taskQueue = new TaskMessageQueue( - new PostgresMessageQueue(postgres, { - tableName: "hollo_task_message_v1", - channelName: "hollo_task_channel_v1", + new PollMessageQueue(postgres, { handlerTimeout: { seconds: 0 }, }), ); diff --git a/src/poll-notification-tasks.ts b/src/poll-notification-tasks.ts index 4435a693..28d89a54 100644 --- a/src/poll-notification-tasks.ts +++ b/src/poll-notification-tasks.ts @@ -59,6 +59,7 @@ export function registerPollNotifications( task, { pollId }, { + orderingKey: `hollo.poll-notification:${pollId}`, delay: { milliseconds: Math.min( MAX_DELAY_MS,