From f97fe72b094a38651360751cf098670f4bb40b0c Mon Sep 17 00:00:00 2001 From: Hong Minhee Date: Tue, 6 Oct 2026 03:35:56 +0900 Subject: [PATCH 1/2] Limit pending poll notification tasks Coalesce poll wakeups in PostgreSQL and cap stored poll messages at 100. Prefer earlier wakeups at capacity and recover displaced or deferred polls after expiry. Keep delayed schedules from allocating worker timers. Cover concurrent scheduling, expiry changes, recovery and shared worker progress, and document queue limits and SQL pressure checks. Fixes https://github.com/fedify-dev/hollo/issues/653 Assisted-by: Codex:gpt-6.1-sol Assisted-by: Claude Code:claude-fable-5-1 --- CHANGES.md | 8 + docs/src/content/docs/install/workers.mdx | 34 + docs/src/content/docs/ja/install/workers.mdx | 31 + docs/src/content/docs/ko/install/workers.mdx | 32 + .../content/docs/zh-cn/install/workers.mdx | 27 + .../content/docs/zh-tw/install/workers.mdx | 27 + src/background/poll-queue.test.ts | 701 ++++++++++++++++++ src/background/poll-queue.ts | 180 +++++ src/federation/federation.ts | 5 +- src/poll-notification-tasks.ts | 1 + 10 files changed, 1043 insertions(+), 3 deletions(-) create mode 100644 src/background/poll-queue.test.ts create mode 100644 src/background/poll-queue.ts 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..86237975 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, From f2221a3a8c8890e00820a814242e4e73cbe627e2 Mon Sep 17 00:00:00 2001 From: Hong Minhee Date: Tue, 6 Oct 2026 12:54:46 +0900 Subject: [PATCH 2/2] Clarify poll expiry reload in Korean docs Name the current expiry time as the value reloaded by a wakeup, so postponed polls are not described as rereading the current time. https://github.com/fedify-dev/hollo/pull/663#discussion_r4187847612 Assisted-by: Codex:gpt-6.1-sol --- docs/src/content/docs/ko/install/workers.mdx | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/docs/src/content/docs/ko/install/workers.mdx b/docs/src/content/docs/ko/install/workers.mdx index 86237975..3cff4516 100644 --- a/docs/src/content/docs/ko/install/workers.mdx +++ b/docs/src/content/docs/ko/install/workers.mdx @@ -112,8 +112,8 @@ Hollo는 여섯 가지 주요 구성 요소로 이루어져 있습니다: 공유 데이터베이스에는 지연 wakeup과 재시도를 포함해 대기 투표 메시지를 최대 100개 저장하며, 새 예약은 투표당 하나의 wakeup으로 합칩니다. 반복 참여와 만료 변경도 합쳐지고, 만료가 앞당겨지면 기존 wakeup을 앞당깁니다. 만료가 늦춰지면 -wakeup 실행 시 현재 시각을 다시 읽습니다. 상한에서는 더 이른 wakeup이 가장 -늦은 메시지를 대체합니다. 보류되거나 밀려난 투표는 만료 후 복구됩니다. 부하가 +wakeup 실행 시 현재 만료 시각을 다시 읽습니다. 상한에서는 더 이른 wakeup이 +가장 늦은 메시지를 대체합니다. 보류되거나 밀려난 투표는 만료 후 복구됩니다. 부하가 몰려 입장 잠금을 1초 안에 얻지 못해도 예약을 보류할 수 있습니다. 이 상한은 워커별이 아니라 모든 노드에 걸쳐 적용되며, 실행 중 handler나 알림 수가 아닌 저장된 작업 수를 제한합니다.