From 560c590eeb1c7eae3a327492a646335ea7a64191 Mon Sep 17 00:00:00 2001 From: CPU-JIA <3144424994@qq.com> Date: Thu, 1 Oct 2026 18:52:23 +0800 Subject: [PATCH 1/4] fix: make bot lease handoff ownership checks atomic --- .github/workflows/ci.yml | 13 ++ docs/docker.md | 15 ++- packages/db/src/worker-fleet.ts | 120 ++++++++++------- packages/db/src/worker-leases.test.ts | 185 ++++++++++++++++++++++++++ 4 files changed, 282 insertions(+), 51 deletions(-) create mode 100644 packages/db/src/worker-leases.test.ts diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index af82aee7..dc10ef20 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -24,6 +24,19 @@ jobs: # past EOL, so the floor here is deliberately above package.json engines. node: [22, 24] name: node ${{ matrix.node }} + services: + redis: + image: redis:7-alpine + ports: + - 6379:6379 + options: >- + --health-cmd "redis-cli ping" + --health-interval 5s + --health-timeout 3s + --health-retries 10 + env: + WECHAT_AI_TEST_REDIS_URL: redis://127.0.0.1:6379/15 + REDIS_URL: redis://127.0.0.1:6379/14 steps: - uses: actions/checkout@v7 diff --git a/docs/docker.md b/docs/docker.md index 135fcb5b..fa3204ba 100644 --- a/docs/docker.md +++ b/docs/docker.md @@ -177,6 +177,19 @@ server { 机器人很多时优先调高 `MAX_BOTS_PER_WORKER` 与系统 `ulimit -n`(注意内存与出站连接数)。 默认 **单副本一体部署** 即可;同镜像多副本已支持(Redis 租约分片 poll)。 +租约续期与释放使用 Redis Lua 原子检查归属,避免旧节点在接管期间覆盖或删除新节点的租约。 +Redis 账号需要允许 `EVAL`(以及脚本内的 `GET` / `EXPIRE` / `DEL` / `SREM`);批量操作仍通过 pipeline 合并往返。 + +### 多节点回归测试 + +CI 在 Node 22 / 24 上使用独立 Redis 服务运行租约接管测试。本地可将 +`WECHAT_AI_TEST_REDIS_URL` 指向**专用测试 Redis** 后运行 `pnpm test`;未设置时这组集成测试会明确跳过。 +测试覆盖续期、单个释放、批量释放及超额认领回收与其他节点接管交错的情况。 + +**交付边界:** 租约协调的是机器人收消息的归属,不是持久化消息队列。 +待回复任务当前仍在进程内存中,进程崩溃或退出可能丢失未完成任务;自动再平衡也可能移走仍有待回复任务的机器人。 +因此节点接管不代表未完成消息一定重试,也不提供 exactly-once 发送保证。 + ## 多节点同构部署(10+ 台) 每台服务器跑**同一镜像**(API + Worker),共用一个 Upstash Redis;用户只访问**主域名**。 @@ -241,4 +254,4 @@ NODE_REGION=cn-east # 可选 目标进程下一轮 reconcile 发现 fence 后停止认领;其他节点 claim 这些 bot。 **解除下线** 后该节点可重新加入。 -这**不会** `docker stop`;若要从 LB 摘流量,还要从 Cloudflare Worker `ORIGINS` 去掉该源站。 \ No newline at end of file +这**不会** `docker stop`;若要从 LB 摘流量,还要从 Cloudflare Worker `ORIGINS` 去掉该源站。 diff --git a/packages/db/src/worker-fleet.ts b/packages/db/src/worker-fleet.ts index 53726305..ca734c3d 100644 --- a/packages/db/src/worker-fleet.ts +++ b/packages/db/src/worker-fleet.ts @@ -798,17 +798,48 @@ export async function claimBotLeases( // Released extras we won but don't need (capacity) const extra = got.slice(keep.length); if (extra.length) { - const drop = db.redis.pipeline(); - for (const botId of extra) { - drop.del(K.botLease(botId)); - } - await drop.exec(); + await releaseOwnedLeasesBatch(db, workerId, extra); } } } return claimed; } +// The ownership check and mutation must execute together. A GET followed by +// SET/DEL (even in pipelines) lets a delayed old node overwrite or delete a +// successor's lease after expiry/rebalance. EXPIRE also cannot resurrect an +// expired lease. These scripts run on the same shared Redis as the bot data. +const RENEW_BOT_LEASE = ` +if redis.call('GET', KEYS[1]) == ARGV[1] then + return redis.call('EXPIRE', KEYS[1], ARGV[3]) +end +redis.call('SREM', KEYS[2], ARGV[2]) +return 0 +`; + +const RELEASE_BOT_LEASE = ` +redis.call('SREM', KEYS[2], ARGV[2]) +if redis.call('GET', KEYS[1]) == ARGV[1] then + return redis.call('DEL', KEYS[1]) +end +return 0 +`; + +/** Do not report a Redis command error or missing reply as a successful write. */ +function leaseResults( + rows: Array<[Error | null, unknown]> | null, + count: number, +): boolean[] { + if (!rows || rows.length !== count) { + throw new Error("Incomplete bot lease results"); + } + return rows.map(([error, value]) => { + if (error) throw error; + if (value !== 0 && value !== 1) throw new Error("Invalid bot lease result"); + return value === 1; + }); +} + /** Renew leases we still own; drop local ownership if stolen/expired. */ export async function renewOwnedLeases( db: RedisStore, @@ -820,30 +851,20 @@ export async function renewOwnedLeases( const lost: string[] = []; if (!botIds.length) return { renewed, lost }; - // Batch GET (one RTT) - const getPipe = db.redis.pipeline(); - for (const botId of botIds) getPipe.get(K.botLease(botId)); - const gets = await getPipe.exec(); - - const setPipe = db.redis.pipeline(); - const lostIds: string[] = []; - botIds.forEach((botId, i) => { - const row = gets?.[i]; - const owner = row && !row[0] ? (row[1] as string | null) : null; - if (owner === workerId) { - setPipe.set(K.botLease(botId), workerId, "EX", ttlSec); - renewed.push(botId); - } else { - lostIds.push(botId); - lost.push(botId); - } - }); - if (renewed.length) await setPipe.exec(); - if (lostIds.length) { - const rem = db.redis.pipeline(); - for (const botId of lostIds) rem.srem(K.workerBots(workerId), botId); - await rem.exec(); + const pipe = db.redis.pipeline(); + for (const botId of botIds) { + pipe.eval( + RENEW_BOT_LEASE, + 2, + K.botLease(botId), + K.workerBots(workerId), + workerId, + botId, + ttlSec, + ); } + const results = leaseResults(await pipe.exec(), botIds.length); + botIds.forEach((botId, i) => (results[i] ? renewed : lost).push(botId)); return { renewed, lost }; } @@ -879,11 +900,14 @@ export async function releaseBotLease( workerId: string, botId: string, ): Promise { - const owner = await db.redis.get(K.botLease(botId)); - if (owner === workerId) { - await db.redis.del(K.botLease(botId)); - } - await db.redis.srem(K.workerBots(workerId), botId); + await db.redis.eval( + RELEASE_BOT_LEASE, + 2, + K.botLease(botId), + K.workerBots(workerId), + workerId, + botId, + ); } /** @@ -897,23 +921,19 @@ export async function releaseOwnedLeasesBatch( botIds: string[], ): Promise { if (!botIds.length) return []; - const getPipe = db.redis.pipeline(); - for (const botId of botIds) getPipe.get(K.botLease(botId)); - const gets = await getPipe.exec(); - - const released: string[] = []; - const delPipe = db.redis.pipeline(); - botIds.forEach((botId, i) => { - const row = gets?.[i]; - const owner = row && !row[0] ? (row[1] as string | null) : null; - if (owner === workerId) { - delPipe.del(K.botLease(botId)); - delPipe.srem(K.workerBots(workerId), botId); - released.push(botId); - } - }); - if (released.length) await delPipe.exec(); - return released; + const pipe = db.redis.pipeline(); + for (const botId of botIds) { + pipe.eval( + RELEASE_BOT_LEASE, + 2, + K.botLease(botId), + K.workerBots(workerId), + workerId, + botId, + ); + } + const results = leaseResults(await pipe.exec(), botIds.length); + return botIds.filter((_, i) => results[i]); } /** diff --git a/packages/db/src/worker-leases.test.ts b/packages/db/src/worker-leases.test.ts new file mode 100644 index 00000000..3a042891 --- /dev/null +++ b/packages/db/src/worker-leases.test.ts @@ -0,0 +1,185 @@ +import { after, before, describe, it } from "node:test"; +import assert from "node:assert/strict"; +import { randomUUID } from "node:crypto"; +import { openDatabase, type RedisStore } from "./client.js"; +import { K } from "./keys.js"; +import { + claimBotLeases, + releaseBotLease, + releaseOwnedLeasesBatch, + renewOwnedLeases, +} from "./worker-fleet.js"; + +// Real Redis is required: mocking EVAL would merely repeat the implementation. +// Use a disposable database. Tests only remove their own keys/set members. +const redisUrl = process.env.WECHAT_AI_TEST_REDIS_URL; + +describe("bot leases across a node handoff", { skip: !redisUrl }, () => { + let db: RedisStore; + const prefix = `lease-test-${randomUUID()}`; + const oldWorker = `${prefix}-old`; + const newWorker = `${prefix}-new`; + const bots: string[] = []; + + before(async () => { + db = openDatabase(redisUrl!); + await db.ping(); + }); + + after(async () => { + if (!db) return; + try { + if (bots.length) { + await db.redis.srem(K.botsPollable, ...bots); + await db.redis.del(...bots.map(K.botLease)); + } + await db.redis.del(K.workerBots(oldWorker), K.workerBots(newWorker)); + } finally { + await db.close(); + } + }); + + async function ownBot() { + const bot = `${prefix}-${bots.length}`; + bots.push(bot); + await db.redis.set(K.botLease(bot), oldWorker, "EX", 45); + await db.redis.sadd(K.workerBots(oldWorker), bot); + return bot; + } + + /** + * Reproduce a successor claiming between the old implementation's GET and + * mutation. For a single atomic operation, hand off before it executes. + * Ownership changes use real Redis; no Lua or Redis semantics are mocked. + */ + function handoffDuringOperation(bot: string): RedisStore { + let transferred = false; + async function transfer() { + if (transferred) return; + transferred = true; + await db.redis.del(K.botLease(bot)); + assert.equal( + await db.redis.set(K.botLease(bot), newWorker, "EX", 90, "NX"), + "OK", + ); + await db.redis.sadd(K.workerBots(newWorker), bot); + } + const redis = new Proxy(db.redis, { + get(target, property) { + if (property === "get") return async (key: string) => { + const result = await target.get(key); + await transfer(); + return result; + }; + if (property === "eval") return async (...args: unknown[]) => { + await transfer(); + return Reflect.apply(target.eval, target, args); + }; + if (property === "pipeline") return () => { + const pipe = target.pipeline(); + let readsOwner = false; + const get = pipe.get.bind(pipe); + pipe.get = ((...args: Parameters) => { + readsOwner = true; + return get(...args); + }) as typeof pipe.get; + const exec = pipe.exec.bind(pipe); + pipe.exec = (async () => { + if (!readsOwner) await transfer(); + const result = await exec(); + if (readsOwner) await transfer(); + return result; + }) as typeof pipe.exec; + return pipe; + }; + const value = Reflect.get(target, property); + return typeof value === "function" ? value.bind(target) : value; + }, + }); + return { redis } as RedisStore; + } + + it("does not overwrite a successor's lease when renewing", async () => { + const bot = await ownBot(); + const result = await renewOwnedLeases(handoffDuringOperation(bot), oldWorker, [bot], 45); + assert.equal(await db.redis.get(K.botLease(bot)), newWorker); + assert.deepEqual(result, { renewed: [], lost: [bot] }); + assert.equal(await db.redis.sismember(K.workerBots(oldWorker), bot), 0); + assert.equal(await db.redis.sismember(K.workerBots(newWorker), bot), 1); + assert.ok((await db.redis.ttl(K.botLease(bot))) > 45); + }); + + it("does not delete a successor's lease on single-bot release", async () => { + const bot = await ownBot(); + await releaseBotLease(handoffDuringOperation(bot), oldWorker, bot); + assert.equal(await db.redis.get(K.botLease(bot)), newWorker); + assert.equal(await db.redis.sismember(K.workerBots(oldWorker), bot), 0); + }); + + it("does not delete a successor's lease during rebalance or shutdown", async () => { + const bot = await ownBot(); + const released = await releaseOwnedLeasesBatch(handoffDuringOperation(bot), oldWorker, [bot]); + assert.equal(await db.redis.get(K.botLease(bot)), newWorker); + assert.deepEqual(released, []); + assert.equal(await db.redis.sismember(K.workerBots(oldWorker), bot), 0); + }); + + it("renews owned leases, reports expired leases, and releases only owned bots", async () => { + const owned = await ownBot(); + const expired = await ownBot(); + await db.redis.del(K.botLease(expired)); + assert.deepEqual(await renewOwnedLeases(db, oldWorker, [owned, expired], 60), { + renewed: [owned], lost: [expired], + }); + assert.ok((await db.redis.ttl(K.botLease(owned))) > 45); + assert.deepEqual(await releaseOwnedLeasesBatch(db, oldWorker, [owned, expired]), [owned]); + assert.equal(await db.redis.get(K.botLease(owned)), null); + assert.equal(await db.redis.sismember(K.workerBots(oldWorker), owned), 0); + assert.equal(await db.redis.sismember(K.workerBots(oldWorker), expired), 0); + }); + + it("propagates Redis renewal errors instead of reporting a successful renewal", async () => { + const bot = await ownBot(); + await assert.rejects(renewOwnedLeases(db, oldWorker, [bot], Number.NaN)); + assert.equal(await db.redis.get(K.botLease(bot)), oldWorker); + }); + + it("propagates Redis release errors without deleting the lease", async () => { + const bot = await ownBot(); + await db.redis.set(K.workerBots(oldWorker), "wrong-type"); + try { + await assert.rejects(releaseOwnedLeasesBatch(db, oldWorker, [bot]), /WRONGTYPE/); + assert.equal(await db.redis.get(K.botLease(bot)), oldWorker); + } finally { + await db.redis.del(K.workerBots(oldWorker)); + } + }); + + it("does not delete a successor's overclaimed lease when trimming capacity", async () => { + const candidates = [await ownBot(), await ownBot()]; + await db.redis.del(...candidates.map(K.botLease)); + await db.redis.srem(K.workerBots(oldWorker), ...candidates); + await db.redis.sadd(K.botsPollable, ...candidates); + let successorBot: string | undefined; + const redis = new Proxy(db.redis, { + get(target, property) { + // Isolate this claim attempt from unrelated bots in the test database. + if (property === "smembers") return async () => [...candidates]; + if (property === "sadd") return async (key: string, ...members: string[]) => { + const result = await target.sadd(key, ...members); + successorBot = candidates.find((id) => !members.includes(id))!; + await target.set(K.botLease(successorBot), newWorker, "EX", 90); + await target.sadd(K.workerBots(newWorker), successorBot); + return result; + }; + const value = Reflect.get(target, property); + return typeof value === "function" ? value.bind(target) : value; + }, + }); + const claimed = await claimBotLeases({ redis } as RedisStore, oldWorker, 1, 45); + assert.equal(claimed.length, 1); + assert.equal(await db.redis.get(K.botLease(claimed[0]!)), oldWorker); + assert.ok(successorBot); + assert.equal(await db.redis.get(K.botLease(successorBot)), newWorker); + }); +}); From b14ac7078e635431ea314fe8b031c99fb4435434 Mon Sep 17 00:00:00 2001 From: CPU-JIA <3144424994@qq.com> Date: Thu, 1 Oct 2026 18:58:00 +0800 Subject: [PATCH 2/4] test: preserve recursive globs so Linux runs top-level cases --- .github/workflows/ci.yml | 4 ++-- apps/api/package.json | 2 +- docs/docker.md | 1 + packages/core/package.json | 2 +- packages/db/package.json | 2 +- packages/ilink/package.json | 2 +- packages/llm/package.json | 2 +- 7 files changed, 8 insertions(+), 7 deletions(-) diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index dc10ef20..bcbb0ee5 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -19,8 +19,8 @@ jobs: strategy: fail-fast: false matrix: - # The `test` scripts pass `src/**/*.test.ts` through to `node --test` - # and rely on Node's own glob expansion, which landed in 22. Node 20 is + # Test scripts quote `src/**/*.test.ts` so the shell cannot expand only + # nested tests and omit top-level files. Node expands the glob (22+). Node 20 is # past EOL, so the floor here is deliberately above package.json engines. node: [22, 24] name: node ${{ matrix.node }} diff --git a/apps/api/package.json b/apps/api/package.json index 1bcdccea..bdcebb15 100644 --- a/apps/api/package.json +++ b/apps/api/package.json @@ -10,7 +10,7 @@ "typecheck": "tsc -p tsconfig.json --noEmit", "ilink:login": "tsx src/cli-login.ts", "diag": "tsx src/cli-doctor.ts", - "test": "node --import tsx --test src/**/*.test.ts" + "test": "node --import tsx --test \"src/**/*.test.ts\"" }, "dependencies": { "@fastify/compress": "^8.0.1", diff --git a/docs/docker.md b/docs/docker.md index fa3204ba..829b7412 100644 --- a/docs/docker.md +++ b/docs/docker.md @@ -185,6 +185,7 @@ Redis 账号需要允许 `EVAL`(以及脚本内的 `GET` / `EXPIRE` / `DEL` / CI 在 Node 22 / 24 上使用独立 Redis 服务运行租约接管测试。本地可将 `WECHAT_AI_TEST_REDIS_URL` 指向**专用测试 Redis** 后运行 `pnpm test`;未设置时这组集成测试会明确跳过。 测试覆盖续期、单个释放、批量释放及超额认领回收与其他节点接管交错的情况。 +各包测试命令的文件通配符须保留引号,交由 Node 展开;否则 Linux shell 可能只匹配子目录,遗漏顶层测试。 **交付边界:** 租约协调的是机器人收消息的归属,不是持久化消息队列。 待回复任务当前仍在进程内存中,进程崩溃或退出可能丢失未完成任务;自动再平衡也可能移走仍有待回复任务的机器人。 diff --git a/packages/core/package.json b/packages/core/package.json index c37dbc09..0217c2f3 100644 --- a/packages/core/package.json +++ b/packages/core/package.json @@ -15,7 +15,7 @@ "scripts": { "build": "tsc -p tsconfig.json", "typecheck": "tsc -p tsconfig.json --noEmit", - "test": "node --import tsx --test src/**/*.test.ts" + "test": "node --import tsx --test \"src/**/*.test.ts\"" }, "dependencies": { "@wechat-ai/db": "workspace:*", diff --git a/packages/db/package.json b/packages/db/package.json index 1e38f486..fdb3df8d 100644 --- a/packages/db/package.json +++ b/packages/db/package.json @@ -16,7 +16,7 @@ "build": "tsc -p tsconfig.json", "typecheck": "tsc -p tsconfig.json --noEmit", "seed": "tsx src/cli-seed.ts", - "test": "node --import tsx --test src/**/*.test.ts" + "test": "node --import tsx --test \"src/**/*.test.ts\"" }, "dependencies": { "ioredis": "^5.6.0" diff --git a/packages/ilink/package.json b/packages/ilink/package.json index baf22e47..96633cae 100644 --- a/packages/ilink/package.json +++ b/packages/ilink/package.json @@ -15,7 +15,7 @@ "scripts": { "build": "tsc -p tsconfig.json", "typecheck": "tsc -p tsconfig.json --noEmit", - "test": "node --import tsx --test src/**/*.test.ts" + "test": "node --import tsx --test \"src/**/*.test.ts\"" }, "devDependencies": { "@types/node": "^22.13.10", diff --git a/packages/llm/package.json b/packages/llm/package.json index 58f110a9..21cfbc96 100644 --- a/packages/llm/package.json +++ b/packages/llm/package.json @@ -15,7 +15,7 @@ "scripts": { "build": "tsc -p tsconfig.json", "typecheck": "tsc -p tsconfig.json --noEmit", - "test": "node --import tsx --test src/**/*.test.ts" + "test": "node --import tsx --test \"src/**/*.test.ts\"" }, "dependencies": { "openai": "^4.91.1" From ca89a4ae3006c4bcdb064b2fc7d37e15bca256a7 Mon Sep 17 00:00:00 2001 From: CPU-JIA <3144424994@qq.com> Date: Thu, 1 Oct 2026 19:55:14 +0800 Subject: [PATCH 3/4] fix: update Fastify and URI parsers with regression coverage --- apps/api/package.json | 2 +- apps/api/src/dependency-regressions.test.ts | 58 +++++++++++++++++++++ pnpm-lock.yaml | 36 ++++++------- 3 files changed, 77 insertions(+), 19 deletions(-) create mode 100644 apps/api/src/dependency-regressions.test.ts diff --git a/apps/api/package.json b/apps/api/package.json index bdcebb15..76e397e0 100644 --- a/apps/api/package.json +++ b/apps/api/package.json @@ -19,7 +19,7 @@ "@wechat-ai/ilink": "workspace:*", "@wechat-ai/llm": "workspace:*", "dotenv": "^16.4.7", - "fastify": "^5.12.3", + "fastify": "^5.12.5", "tsx": "^4.19.3", "zod": "^3.24.2" }, diff --git a/apps/api/src/dependency-regressions.test.ts b/apps/api/src/dependency-regressions.test.ts new file mode 100644 index 00000000..d7fab1de --- /dev/null +++ b/apps/api/src/dependency-regressions.test.ts @@ -0,0 +1,58 @@ +import assert from "node:assert/strict"; +import { it } from "node:test"; +import { createRequire } from "node:module"; +import { connect } from "node:http2"; +import Fastify from "fastify"; + +const require = createRequire(import.meta.url); +const fromFastify = createRequire(require.resolve("fastify")); +const fromCompiler = createRequire(fromFastify.resolve("@fastify/ajv-compiler")); +const fromAjv = createRequire(fromCompiler.resolve("ajv")); + +it("normalizes encoded host case in both installed fast-uri major versions", () => { + for (const uri of [fromCompiler("fast-uri"), fromAjv("fast-uri")]) { + for (const encoded of ["//%41.com", "//%4a.com"]) { + const canonical = encoded.replace("%41", "a").replace("%4a", "j"); + assert.equal(uri.parse(encoded).host, uri.parse(canonical).host); + assert.equal(uri.equal(encoded, canonical), true); + } + assert.equal(uri.equal("//A.com", "//a.com"), true); + } +}); + +it("keeps mailto recipients and reserved headers stable through a round trip", () => { + const uri = fromCompiler("fast-uri"); + for (const query of ["%74o=other@example.com", "%73ubject=hello", "%62ody=hello", "subject=hello%20world"]) { + const parsed = uri.parse(`mailto:user@example.com?${query}`); + const reparsed = uri.parse(uri.serialize(parsed)); + assert.deepEqual(reparsed.to, parsed.to); + assert.equal(reparsed.subject, parsed.subject); + assert.equal(reparsed.body, parsed.body); + } +}); + +it("serves an HTTP/2 trailer response without an uncaught header exception", { timeout: 10_000 }, async () => { + const app = Fastify({ http2: true }); + app.get("/", async (_request, reply) => { + reply.trailer("x-checksum", async () => "ok"); + return "hello"; + }); + const address = await app.listen({ host: "127.0.0.1", port: 0 }); + const session = connect(address); + try { + const body = await new Promise((resolve, reject) => { + session.once("error", reject); + const request = session.request({ ":path": "/" }); + let text = ""; + request.setEncoding("utf8"); + request.on("data", (chunk) => { text += chunk; }); + request.once("error", reject); + request.once("end", () => resolve(text)); + request.end(); + }); + assert.equal(body, "hello"); + } finally { + session.destroy(); + await app.close(); + } +}); diff --git a/pnpm-lock.yaml b/pnpm-lock.yaml index 122f7ed8..fffd9c42 100644 --- a/pnpm-lock.yaml +++ b/pnpm-lock.yaml @@ -29,8 +29,8 @@ importers: specifier: ^16.4.7 version: 16.6.1 fastify: - specifier: ^5.12.3 - version: 5.12.3 + specifier: ^5.12.5 + version: 5.12.5 tsx: specifier: ^4.19.3 version: 4.23.1 @@ -287,8 +287,8 @@ packages: '@fastify/merge-json-schemas@0.2.1': resolution: {integrity: sha512-OA3KGBCy6KtIvLf8DINC5880o5iBlDX4SxzLQS8HorJAbqluzLRn80UXU0bxZn7UOFhFgpRJDasfwn9nG4FG4A==} - '@fastify/proxy-addr@5.1.0': - resolution: {integrity: sha512-INS+6gh91cLUjB+PVHfu1UqcB76Sqtpyp7bnL+FYojhjygvOPA9ctiD/JDKsyD9Xgu4hUhCSJBPig/w7duNajw==} + '@fastify/proxy-addr@5.1.1': + resolution: {integrity: sha512-zv07Y9GEuDsJPegZoDFd4SDWaZOW8N2pa0GSrYmKpId/tjt1Hgo3BjZBVjdVpfVrHaA+Qv5jawtS2O50J5xM9g==} '@ioredis/commands@1.10.0': resolution: {integrity: sha512-UmeW7z4LfctwoQ5wkhVzgq8tXkreED2xZGpX+Bg+zA+WJFZCT6c062AfCK/Dfk81xZnnwdhJCUMkitihRaoC2Q==} @@ -444,17 +444,17 @@ packages: fast-querystring@1.1.2: resolution: {integrity: sha512-g6KuKWmFXc0fID8WWH0jit4g0AGBoJhCkJMb1RmbsSEUNvQ+ZC8D6CUZ+GtF8nMzSPXnhiePyyqqipzNNEnHjg==} - fast-uri@3.1.7: - resolution: {integrity: sha512-dOvZVzjdZdz7phd9v6jCbwxrBW3fK6n8Rc0CtdmM4bumzMnxywBYhuph6J819RRw/ku+rLbelwfMunktuzVVHg==} + fast-uri@3.1.8: + resolution: {integrity: sha512-GZMtZUTNRpOVIECoXwLNZS5xUGE+mVNbTB8h/7Rwh2TFWcBQiPzTgyZi05BF9UMZKkLJv8XBRJTlU7zg8+ZfMg==} - fast-uri@4.1.4: - resolution: {integrity: sha512-dODXrIxlS9JSdgAnhIUKOosKV1oMtU2VtVw87QRaHzyl5jxO290Ii5tEZfCfzfWNHi3jKWwBSdQj0qIyshdZdQ==} + fast-uri@4.2.1: + resolution: {integrity: sha512-TmHQgewjHtMq1E5QKA0tOE0yeYGQs25KZC/ziJpubRtWI15W92e6vPFydWOeZBVROSHBYw/QhD9d5OefdD6LDg==} fastify-plugin@5.1.0: resolution: {integrity: sha512-FAIDA8eovSt5qcDgcBvDuX/v0Cjz0ohGhENZ/wpc3y+oZCY2afZ9Baqql3g/lC+OHRnciQol4ww7tuthOb9idw==} - fastify@5.12.3: - resolution: {integrity: sha512-reZ8wce5VNCcufIt9AVtzZa3L4u1j8esikn7OEgHWLVpRpL5R7Y2+Xzj70OUkv5zDfzUAxXZT6cu4Rt0zr3EKA==} + fastify@5.12.5: + resolution: {integrity: sha512-OB2k1dlxs5/NAABqeKV2FUHkSD2BbENsCak8yULVcymn3fHIPDVa9TI3SDnJSWYSllZmSYuZXy2gTnsT+Sut1A==} fastq@1.20.3: resolution: {integrity: sha512-XKv5nnLs6nLF71NgiKJLIZFLkPyIEuOselLG7ujZnGrRfQK8HpvY+WqKhAJUAdLomwVHErVS4LfxFlPq0/FTAw==} @@ -850,7 +850,7 @@ snapshots: dependencies: ajv: 8.20.0 ajv-formats: 3.0.1(ajv@8.20.0) - fast-uri: 4.1.4 + fast-uri: 4.2.1 '@fastify/compress@8.3.1': dependencies: @@ -875,7 +875,7 @@ snapshots: dependencies: dequal: 2.0.3 - '@fastify/proxy-addr@5.1.0': + '@fastify/proxy-addr@5.1.1': dependencies: '@fastify/forwarded': 3.0.2 ipaddr.js: 2.5.0 @@ -914,7 +914,7 @@ snapshots: ajv@8.20.0: dependencies: fast-deep-equal: 3.1.3 - fast-uri: 3.1.7 + fast-uri: 3.1.8 json-schema-traverse: 1.0.0 require-from-string: 2.0.2 @@ -1044,7 +1044,7 @@ snapshots: '@fastify/merge-json-schemas': 0.2.1 ajv: 8.20.0 ajv-formats: 3.0.1(ajv@8.20.0) - fast-uri: 4.1.4 + fast-uri: 4.2.1 json-schema-ref-resolver: 3.0.0 rfdc: 1.4.1 @@ -1052,18 +1052,18 @@ snapshots: dependencies: fast-decode-uri-component: 1.0.1 - fast-uri@3.1.7: {} + fast-uri@3.1.8: {} - fast-uri@4.1.4: {} + fast-uri@4.2.1: {} fastify-plugin@5.1.0: {} - fastify@5.12.3: + fastify@5.12.5: dependencies: '@fastify/ajv-compiler': 4.0.6 '@fastify/error': 4.2.0 '@fastify/fast-json-stringify-compiler': 5.1.0 - '@fastify/proxy-addr': 5.1.0 + '@fastify/proxy-addr': 5.1.1 abstract-logging: 2.0.1 avvio: 9.3.0 fast-json-stringify: 7.0.1 From 73a05f96b90d8f0500fbb005fbd40292d2f65c5e Mon Sep 17 00:00:00 2001 From: CPU-JIA <3144424994@qq.com> Date: Thu, 1 Oct 2026 19:55:15 +0800 Subject: [PATCH 4/4] fix: persist inbound replies across worker handoff and restart --- .env.example | 2 + apps/api/src/config.ts | 1 + apps/api/src/inbound-admin.test.ts | 90 +++++ apps/api/src/inbound-delivery.test.ts | 93 +++++ apps/api/src/inbound-delivery.ts | 63 ++++ apps/api/src/routes.ts | 23 ++ apps/api/src/worker-recovery.test.ts | 205 +++++++++++ apps/api/src/worker.ts | 500 +++++++++----------------- docs/admin-api.md | 10 +- docs/docker.md | 30 +- docs/runbook.md | 17 +- packages/db/src/inbound-queue.test.ts | 170 +++++++++ packages/db/src/inbound-queue.ts | 196 ++++++++++ packages/db/src/index.ts | 1 + packages/db/src/repos.ts | 19 +- 15 files changed, 1080 insertions(+), 340 deletions(-) create mode 100644 apps/api/src/inbound-admin.test.ts create mode 100644 apps/api/src/inbound-delivery.test.ts create mode 100644 apps/api/src/inbound-delivery.ts create mode 100644 apps/api/src/worker-recovery.test.ts create mode 100644 packages/db/src/inbound-queue.test.ts create mode 100644 packages/db/src/inbound-queue.ts diff --git a/.env.example b/.env.example index 291ae411..ff325b2c 100644 --- a/.env.example +++ b/.env.example @@ -258,6 +258,8 @@ LEASE_RENEW_SEC=15 # WORKER_ID=node-a # Concurrent LLM + WeChat send jobs (poll is decoupled from reply) REPLY_CONCURRENCY=16 +# Retained jobs PER BOT in Redis, including failed jobs awaiting manual retry. +# At capacity, polling retries the batch without advancing its cursor. INBOX_MAX_LEN=20000 # silent | fatal | error | warn | info | debug | trace (anything else → info). # Drives Fastify's logger, which emits one line per request plus framework diff --git a/apps/api/src/config.ts b/apps/api/src/config.ts index 751de27f..3cb0b062 100644 --- a/apps/api/src/config.ts +++ b/apps/api/src/config.ts @@ -117,6 +117,7 @@ export interface AppConfig { workerWeightTtlSec: number; /** Concurrent LLM/reply jobs */ replyConcurrency: number; + /** Retained Redis inbound jobs per bot, including failed jobs. */ inboxMaxLen: number; logLevel: LogLevel; /** Requests slower than this are logged at warn even when they succeed */ diff --git a/apps/api/src/inbound-admin.test.ts b/apps/api/src/inbound-admin.test.ts new file mode 100644 index 00000000..8bdc8789 --- /dev/null +++ b/apps/api/src/inbound-admin.test.ts @@ -0,0 +1,90 @@ +import assert from "node:assert/strict"; +import { after, before, describe, it } from "node:test"; +import { randomUUID } from "node:crypto"; +import Fastify from "fastify"; +import { + openDatabase, K, createLocalUser, createAppSession, deleteUserAccount, + resolveSuperAdminId, persistInbound, claimInbound, inboundQueueKeys, + deleteInboundQueue, type Db, type User, +} from "@wechat-ai/db"; +import { loadConfig } from "./config.js"; +import { registerRoutes, type RouteContext } from "./routes.js"; + +const redisUrl = process.env.WECHAT_AI_TEST_REDIS_URL; +describe("inbound recovery admin API", { skip: !redisUrl }, () => { + let db: Db; + const app = Fastify(); + const id = `inbound-admin-${randomUUID()}`; + const users: User[] = []; + const cookies = new Map(); + let superId: string; + const url = `/api/v1/admin/bots/${id}/inbox`; + const retryUrl = `${url}/${id}/retry`; + before(async () => { + db = openDatabase(redisUrl!); + await db.ping(); + for (let i = 0; i < 3; i++) { + const user = await createLocalUser(db, { + username: `it_${randomUUID().replaceAll("-", "").slice(0, 16)}`, + passwordHash: "unused-test-hash", forceAdmin: i < 2, + }, new Set()); + users.push(user); + cookies.set(user.id, `wa_session=${await createAppSession(db, user.id)}`); + } + superId = (await resolveSuperAdminId(db))!; + assert.ok(users.some((u) => u.id === superId), "use a dedicated test Redis database"); + await db.redis.set(K.bot(id), JSON.stringify({ id, status: "active" })); + await db.redis.set(K.botLease(id), "worker", "EX", 45); + await persistInbound(db, "worker", { + id, botId: id, peerId: "peer", text: "private-message", + contextToken: "private-context", mediaOnly: false, enqueuedAt: new Date().toISOString(), + ...{ mediaRefs: [{ aesKey: "private-media-key" }] }, + }, 100); + await claimInbound(db, id, "worker"); + await db.redis.zadd(inboundQueueKeys(id).active, 0, "peer"); + await claimInbound(db, id, "worker", 120_000, 1); + await registerRoutes(app, { + db, cfg: loadConfig({}), chat: {}, tryChat: {}, worker: {}, loginSessions: {}, + } as RouteContext); + await app.ready(); + }); + after(async () => { + await app.close(); + if (!db) return; + await deleteInboundQueue(db, id); + await db.redis.del(K.bot(id), K.botLease(id)); + for (const row of await db.redis.lrange(K.audit, 0, -1)) { + if (users.some((u) => u.id === JSON.parse(row).actor)) await db.redis.lrem(K.audit, 0, row); + } + for (const user of users) await deleteUserAccount(db, user.id); + await db.close(); + }); + + it("requires the super-admin session for both inspection and retry", async () => { + for (const [method, path] of [["GET", url], ["POST", retryUrl]] as const) { + assert.equal((await app.inject({ method, url: path })).statusCode, 401); + for (const user of users.filter((u) => u.id !== superId)) { + assert.equal((await app.inject({ method, url: path, headers: { cookie: cookies.get(user.id)! } })).statusCode, 403); + } + } + assert.equal(await db.redis.zcard(inboundQueueKeys(id).failed), 1); + }); + + it("returns only failure metadata and retries with an audit record", async () => { + const headers = { cookie: cookies.get(superId)! }; + const response = await app.inject({ method: "GET", url, headers }); + assert.equal(response.statusCode, 200); + assert.match(response.headers["cache-control"]!, /no-store/); + assert.equal(response.json().failed, 1); + assert.equal(response.json().failedJobs[0].id, id); + for (const secret of ["private-message", "private-context", "private-media-key"]) assert.ok(!response.body.includes(secret)); + assert.equal((await app.inject({ method: "POST", url: retryUrl, headers })).statusCode, 200); + assert.equal(await db.redis.zcard(inboundQueueKeys(id).failed), 0); + assert.equal(await db.redis.llen(inboundQueueKeys(id).peers + "peer"), 1); + const audit = (await db.redis.lrange(K.audit, 0, -1)).map((x) => JSON.parse(x)); + const record = audit.find((x) => x.action === "admin_inbound_retry" && x.actor === superId); + assert.deepEqual(JSON.parse(record.meta_json), { botId: id, jobId: id }); + assert.equal((await app.inject({ method: "POST", url: retryUrl, headers })).statusCode, 404); + assert.equal((await app.inject({ method: "GET", url: `/api/v1/admin/bots/${id}-missing/inbox`, headers })).statusCode, 404); + }); +}); diff --git a/apps/api/src/inbound-delivery.test.ts b/apps/api/src/inbound-delivery.test.ts new file mode 100644 index 00000000..0cbcef49 --- /dev/null +++ b/apps/api/src/inbound-delivery.test.ts @@ -0,0 +1,93 @@ +import { after, before, describe, it } from "node:test"; +import assert from "node:assert/strict"; +import { randomUUID } from "node:crypto"; +import { + openDatabase, K, persistInbound, claimInbound, acknowledgeInbound, + inboundQueueKeys, deleteInboundQueue, type Db, type InboundJob, +} from "@wechat-ai/db"; +import type { ILinkClient } from "@wechat-ai/ilink"; +import { InboundDelivery } from "./inbound-delivery.js"; + +const redisUrl = process.env.WECHAT_AI_TEST_REDIS_URL; +describe("reply checkpoints (Redis + controlled iLink sender)", { skip: !redisUrl }, () => { + let db: Db; + const prefix = `delivery-test-${randomUUID()}`; + const jobs: InboundJob[] = []; + before(async () => { db = openDatabase(redisUrl!); await db.ping(); }); + after(async () => { + for (const job of jobs) { + await deleteInboundQueue(db, job.botId); + await db.redis.del(K.bot(job.botId), K.botLease(job.botId), K.inboundSeen(job.id)); + } + await db.close(); + }); + async function setup() { + const id = `${prefix}-${jobs.length}`; + const job = { id, botId: id, peerId: "peer", contextToken: "context", text: "hello", mediaOnly: false, enqueuedAt: new Date().toISOString() }; + jobs.push(job); + await db.redis.set(K.bot(id), "{}"); + await db.redis.set(K.botLease(id), "a", "EX", 45); + await persistInbound(db, "a", job, 100); + return job; + } + + it("does not regenerate or resend confirmed bubbles after a crash before ACK", async () => { + const job = await setup(); + let generated = 0; + const sent: string[] = []; + const client = { async sendText(p: { text: string }) { sent.push(p.text); return { ret: 0 }; } } as unknown as ILinkClient; + const first = (await claimInbound(db, job.botId, "a"))!; + const a = new InboundDelivery(db, first); + const response = await a.step("chat", async () => { generated++; return "reply"; }); + await a.client(client).sendText({ text: response, toUserId: "peer", contextToken: "context" }); + a.close(); // simulate process loss before ACK + await db.redis.zadd(inboundQueueKeys(job.botId).active, 0, job.peerId); + await db.redis.set(K.botLease(job.botId), "b", "EX", 45); + const recovered = (await claimInbound(db, job.botId, "b"))!; + const b = new InboundDelivery(db, recovered); + try { + const replay = await b.step("chat", async () => { generated++; return "wrong regeneration"; }); + await b.client(client).sendText({ text: replay, toUserId: "peer", contextToken: "context" }); + await acknowledgeInbound(db, recovered); + assert.equal(generated, 1); + assert.deepEqual(sent, ["reply"]); + } finally { b.close(); } + }); + + it("reuses client_id when a send succeeded but its acknowledgement was lost", async () => { + const job = await setup(); + const ids: string[] = []; + const client = { async sendText(p: { clientId: string }) { + ids.push(p.clientId); + if (ids.length === 1) throw new Error("connection closed after server accepted send"); + return { ret: 0 }; + } } as unknown as ILinkClient; + const first = (await claimInbound(db, job.botId, "a"))!; + const a = new InboundDelivery(db, first); + try { + await assert.rejects(a.client(client).sendText({ text: "reply", toUserId: "peer", contextToken: "context" })); + } finally { a.close(); } + await db.redis.zadd(inboundQueueKeys(job.botId).active, 0, job.peerId); + const recovered = (await claimInbound(db, job.botId, "a"))!; + const b = new InboundDelivery(db, recovered); + try { + await b.client(client).sendText({ text: "reply", toUserId: "peer", contextToken: "context" }); + assert.equal(ids.length, 2); + assert.equal(ids[0], ids[1]); + await acknowledgeInbound(db, recovered); + } finally { b.close(); } + }); + + it("fences a stale processor before another outbound request", async () => { + const job = await setup(); + const claim = (await claimInbound(db, job.botId, "a"))!; + const delivery = new InboundDelivery(db, claim); + let sends = 0; + const client = { async sendText() { sends++; return { ret: 0 }; } } as unknown as ILinkClient; + await db.redis.zadd(inboundQueueKeys(job.botId).active, 0, job.peerId); + try { + await assert.rejects(delivery.client(client).sendText({ text: "stale", toUserId: "peer", contextToken: "context" }), /lease lost/); + assert.equal(sends, 0); + } finally { delivery.close(); } + }); +}); diff --git a/apps/api/src/inbound-delivery.ts b/apps/api/src/inbound-delivery.ts new file mode 100644 index 00000000..4bead1ce --- /dev/null +++ b/apps/api/src/inbound-delivery.ts @@ -0,0 +1,63 @@ +import { createHash } from "node:crypto"; +import { + INBOUND_LEASE_MS, renewInbound, saveInboundSteps, + type Db, type InboundClaim, +} from "@wechat-ai/db"; +import type { ILinkClient } from "@wechat-ai/ilink"; + +/** Checkpoints belong to the job, not the node that currently processes it. */ +export class InboundDelivery { + private lost = false; + private timer: ReturnType; + + constructor(private db: Db, readonly claim: InboundClaim) { + this.timer = setInterval(() => { + void this.assertActive().catch(() => { this.lost = true; }); + }, INBOUND_LEASE_MS / 3); + this.timer.unref(); + } + + close(): void { clearInterval(this.timer); } + + async assertActive(): Promise { + if (this.lost) throw new Error("Inbound processing lease lost"); + try { + if (!await renewInbound(this.db, this.claim)) throw new Error("Inbound processing lease lost"); + } catch (err) { + this.lost = true; + throw err; + } + } + + clientId(step: string): string { + return `wa-${createHash("sha256").update(`${this.claim.job.id}:${step}`).digest("hex").slice(0, 40)}`; + } + + async step(name: string, run: () => Promise): Promise { + if (Object.hasOwn(this.claim.steps, name)) return this.claim.steps[name]!.value as T; + await this.assertActive(); + const value = await run(); + this.claim.steps[name] = { value }; + await saveInboundSteps(this.db, this.claim); + return value; + } + + /** Already-confirmed bubbles are skipped; ambiguous retries reuse client_id. */ + client(raw: ILinkClient, scope = "send"): ILinkClient { + let sequence = 0; + return new Proxy(raw, { + get: (target, property) => { + if (property === "sendText" || property === "sendImage") { + return (params: Record) => { + const name = `${scope}:${sequence++}:${property}`; + return this.step(name, () => Reflect.apply(target[property], target, [{ + ...params, clientId: this.clientId(name), + }])); + }; + } + const value = Reflect.get(target, property); + return typeof value === "function" ? value.bind(target) : value; + }, + }); + } +} diff --git a/apps/api/src/routes.ts b/apps/api/src/routes.ts index 947f8aa3..d60f40b0 100644 --- a/apps/api/src/routes.ts +++ b/apps/api/src/routes.ts @@ -60,6 +60,9 @@ import { listAuditLogs, listBlockedUserIds, listBotAccounts, + listFailedInbound, + retryFailedInbound, + inboundQueueStats, hasBotCredentials, hasBotCredentialsMany, listBotsByOwner, @@ -3270,6 +3273,26 @@ export async function registerRoutes( }, ); + app.get<{ Params: { botId: string } }>("/api/v1/admin/bots/:botId/inbox", async (req, reply) => { + if (!await requireSuperAdmin(req, reply, ctx)) return; + const botId = req.params.botId; + if (!await getBotAccount(ctx.db, botId)) return reply.code(404).send({ error: "bot_not_found" }); + return { + ...await inboundQueueStats(ctx.db, [botId]), + failedJobs: await listFailedInbound(ctx.db, botId), + }; + }); + + app.post<{ Params: { botId: string; jobId: string } }>("/api/v1/admin/bots/:botId/inbox/:jobId/retry", async (req, reply) => { + const admin = await requireSuperAdmin(req, reply, ctx); + if (!admin) return; + const { botId, jobId } = req.params; + if (!await getBotAccount(ctx.db, botId)) return reply.code(404).send({ error: "bot_not_found" }); + if (!await retryFailedInbound(ctx.db, botId, jobId)) return reply.code(404).send({ error: "failed_job_not_found" }); + await writeAudit(ctx.db, "admin_inbound_retry", admin.id, { botId, jobId }); + return { ok: true }; + }); + /** Start/restart workers for every active bot that has Redis credentials. */ app.post("/api/v1/admin/workers/restart-all", async (req, reply) => { const admin = await requireSuperAdmin(req, reply, ctx); diff --git a/apps/api/src/worker-recovery.test.ts b/apps/api/src/worker-recovery.test.ts new file mode 100644 index 00000000..ae7aaf71 --- /dev/null +++ b/apps/api/src/worker-recovery.test.ts @@ -0,0 +1,205 @@ +import assert from "node:assert/strict"; +import { after, before, describe, it, mock } from "node:test"; +import { randomUUID } from "node:crypto"; +import { + openDatabase, K, inboundQueueKeys, inboundQueueStats, deleteInboundQueue, + type Db, type InboundJob, +} from "@wechat-ai/db"; +import { ILinkClient, type WeixinMessage } from "@wechat-ai/ilink"; +import type { ChatService, InboundChatResult } from "@wechat-ai/core"; +import { BotWorkerManager } from "./worker.js"; + +// Exercise the real worker paths with Redis and controlled external services. +interface Internals { + workerId: string; + loops: Map>; + loopGen: Map; + clients: Map; + replyWorkers: Promise[]; + inboxScanAfter: number; + inboxMaxLen: number; + enqueueFromMessage(bot: string, client: ILinkClient, msg: WeixinMessage): Promise; + shedLeases(n: number, local: number, ctx: { target: number; slack: number }): Promise; + runReplyConsumer(n: number): Promise; + runPollLoop(bot: string, gen: number): Promise; +} +const redisUrl = process.env.WECHAT_AI_TEST_REDIS_URL; +const pause = (ms: number) => new Promise((resolve) => setTimeout(resolve, ms)); +async function until(check: () => boolean | Promise) { + const end = Date.now() + 8000; + while (!await check()) { + assert.ok(Date.now() < end, "worker condition timed out"); + await pause(20); + } +} +function gate() { + let release!: () => void; + const promise = new Promise((resolve) => { release = resolve; }); + return { promise, release }; +} + +describe("worker recovery through real polling/reply paths", { skip: !redisUrl }, () => { + let db: Db; + const prefix = `worker-recovery-${randomUUID()}`; + const bots: string[] = []; + const managers: BotWorkerManager[] = []; + const jobs = new Set(); + before(async () => { db = openDatabase(redisUrl!); await db.ping(); }); + after(async () => { + mock.restoreAll(); + for (const manager of managers) await manager.stopAsync(1000); + for (const bot of bots) { + await deleteInboundQueue(db, bot); + await db.redis.del(K.bot(bot), K.botLease(bot), K.botCreds(bot), K.contextToken(bot, "peer")); + } + for (const id of jobs) await db.redis.del(K.inboundSeen(id)); + await db.close(); + }); + async function bot() { + const id = `${prefix}-${bots.length}`; + bots.push(id); + await db.redis.set(K.bot(id), JSON.stringify({ id, status: "active", updates_cursor: "before" })); + return id; + } + function manager(handle: (p: { text: string }) => Promise) { + const m = new BotWorkerManager({ + db, chat: { handleInbound: handle } as unknown as ChatService, p2pEnabled: false, + replyDelay: { msPerChar: 0, minMs: 0, maxMs: 0, firstMinMs: 0, firstMaxMs: 0, thinkExtraMs: 0 }, + }); + managers.push(m); + const inner = m as unknown as Internals; + inner.workerId = `${prefix}-node-${managers.length}`; + return { m, inner }; + } + function client(send: (p: { text: string; clientId?: string }) => Promise) { + return { + async startTyping() {}, async stopTyping() {}, + async sendText(p: { text: string; clientId?: string }) { await send(p); return { ret: 0 }; }, + } as unknown as ILinkClient; + } + async function own(inner: Internals, id: string, sender: ILinkClient) { + await db.redis.set(K.botLease(id), inner.workerId, "EX", 45); + await db.redis.sadd(K.workerBots(inner.workerId), id); + inner.loops.set(id, Promise.resolve()); + inner.clients.set(id, sender); + } + function message(text: string, time = 1): WeixinMessage { + return { message_type: 1, from_user_id: "peer", context_token: "private-context", create_time_ms: time, item_list: [{ type: 1, text_item: { text } }] }; + } + async function enqueue(inner: Internals, id: string, sender: ILinkClient, text: string, time = 1) { + await inner.enqueueFromMessage(id, sender, message(text, time)); + for (const key of await db.redis.hkeys(inboundQueueKeys(id).jobs)) jobs.add(key); + } + function consume(inner: Internals) { inner.replyWorkers.push(inner.runReplyConsumer(0)); } + async function empty(id: string) { return (await inboundQueueStats(db, [id])).pending === 0; } + + it("replies to persisted messages after rebalance removes the receiving node's client", async () => { + const id = await bot(), sent: string[] = []; + const sender = client(async (p) => { sent.push(p.text); }); + const a = manager(async () => { throw new Error("old node must not generate"); }); + const b = manager(async () => ({ kind: "reply", text: "recovered" })); + await own(a.inner, id, sender); + await enqueue(a.inner, id, sender, "hello"); + await a.inner.shedLeases(1, 1, { target: 0, slack: 0 }); + assert.equal(a.inner.clients.has(id), false); + await own(b.inner, id, sender); + consume(b.inner); + await until(() => empty(id)); + assert.deepEqual(sent, ["recovered"]); + }); + + it("lets an in-flight old owner finish and keeps the successor's next reply ordered", async () => { + const id = await bot(), sent: string[] = [], started = gate(), finish = gate(); + const sender = client(async (p) => { sent.push(p.text); }); + const a = manager(async () => { started.release(); await finish.promise; return { kind: "reply", text: "first" }; }); + const b = manager(async () => ({ kind: "reply", text: "second" })); + try { + await own(a.inner, id, sender); + await enqueue(a.inner, id, sender, "one"); + await enqueue(a.inner, id, sender, "two", 2); + consume(a.inner); + await started.promise; + await a.inner.shedLeases(1, 1, { target: 0, slack: 0 }); + await own(b.inner, id, sender); + consume(b.inner); + await pause(150); + assert.deepEqual(sent, []); + finish.release(); + await until(() => empty(id)); + assert.deepEqual(sent, ["first", "second"]); + } finally { finish.release(); } + }); + + it("retries a failed bubble without regenerating or resending earlier parts after settings change", async () => { + const id = await bot(), sent: { text: string; clientId?: string }[] = []; + let generated = 0; + const sender = client(async (p) => { + sent.push(p); + if (sent.length === 2) throw new Error("ambiguous second send"); + }); + const a = manager(async () => { generated++; return { kind: "reply", bubbles: ["first", "second", "third"] }; }); + await own(a.inner, id, sender); + await enqueue(a.inner, id, sender, "hello"); + consume(a.inner); + await until(() => sent.length === 2); + a.m.applyRuntimeConfig({ splitReply: false }); + await until(() => empty(id)); + assert.equal(generated, 1); + assert.deepEqual(sent.map((p) => p.text), ["first", "second", "second", "third"]); + assert.equal(sent[1]!.clientId, sent[2]!.clientId); + assert.notEqual(sent[0]!.clientId, sent[1]!.clientId); + }); + + it("gracefully ACKs started work and leaves unstarted work for another node", async () => { + const id = await bot(), sent: string[] = [], started = gate(), finish = gate(); + const sender = client(async (p) => { sent.push(p.text); }); + const a = manager(async () => { started.release(); await finish.promise; return { kind: "reply", text: "drained" }; }); + try { + await own(a.inner, id, sender); + await enqueue(a.inner, id, sender, "one"); + await enqueue(a.inner, id, sender, "two", 2); + consume(a.inner); + await started.promise; + const stopping = a.m.stopAsync(3000); + finish.release(); + await stopping; + assert.deepEqual(sent, ["drained"]); + assert.equal((await inboundQueueStats(db, [id])).pending, 1); + const b = manager(async () => ({ kind: "reply", text: "remaining" })); + await own(b.inner, id, sender); + consume(b.inner); + await until(() => empty(id)); + assert.deepEqual(sent, ["drained", "remaining"]); + } finally { finish.release(); } + }); + + it("does not advance the poll cursor when only part of a batch fits in Redis", async () => { + const id = await bot(); + const a = manager(async () => ({ kind: "skip" })); + const sender = client(async () => {}); + await own(a.inner, id, sender); + await db.redis.set(K.botCreds(id), JSON.stringify({ botToken: "test-token" })); + a.inner.inboxMaxLen = 1; + a.inner.loopGen.set(id, 1); + const getUpdates = mock.method(ILinkClient.prototype, "getUpdates", async () => ({ + msgs: [message("first"), message("second", 2)], get_updates_buf: "after", + })); + const startTyping = mock.method(ILinkClient.prototype, "startTyping", async () => {}); + const polling = a.inner.runPollLoop(id, 1); + try { + await until(async () => (await db.redis.hlen(inboundQueueKeys(id).jobs)) === 1); + a.inner.loopGen.set(id, 2); + await polling; + const row = JSON.parse((await db.redis.get(K.bot(id)))!); + assert.equal(row.updates_cursor, "before"); + const pending = (await db.redis.hvals(inboundQueueKeys(id).jobs)).map((x) => JSON.parse(x) as InboundJob); + assert.equal(pending[0]!.text, "first"); + assert.equal(await db.redis.exists(K.inboundSeen(pending[0]!.id)), 0); + jobs.add(pending[0]!.id); + } finally { + a.inner.loopGen.set(id, 2); + getUpdates.mock.restore(); startTyping.mock.restore(); + await polling; + } + }); +}); diff --git a/apps/api/src/worker.ts b/apps/api/src/worker.ts index 339d6a7f..e4349550 100644 --- a/apps/api/src/worker.ts +++ b/apps/api/src/worker.ts @@ -14,6 +14,12 @@ import { import { type Db, type InboundJob, + type InboundClaim, + persistInbound, + claimInbound, + acknowledgeInbound, + retryInbound, + inboundQueueStats, K, claimBotLeases, getBotAccount, @@ -85,6 +91,7 @@ import { type AdminSendResult, } from "./broadcast-runner.js"; import { emitActivity, previewText } from "./activity-stream.js"; +import { InboundDelivery } from "./inbound-delivery.js"; /** * How often a node sweeps the load-weight hash (refresh live entries, delete @@ -130,9 +137,8 @@ export interface BroadcastWorkerConfig { /** * Inbound job plus the attachment coordinates for this message. * - * `mediaRefs` is deliberately NOT on the shared `InboundJob` type: that one has - * a Redis serializer in worker-fleet.ts, and these refs carry the CDN AES key — - * anything holding them must stay in this process's memory. + * Media coordinates contain CDN credentials. The durable Redis inbox has the + * same trust boundary as bot credentials; never expose payloads via status APIs. */ type LocalInboundJob = InboundJob & { mediaRefs?: InboundMediaRef[]; @@ -156,7 +162,7 @@ export interface WorkerOptions { workerWeightTtlSec?: number; /** Concurrent reply consumers (LLM + send) */ replyConcurrency?: number; - /** Max in-process inbox depth before drop */ + /** Max retained inbound jobs per bot, including failed jobs */ inboxMaxLen?: number; /** multi-bubble human-like reply */ splitReply?: boolean; @@ -214,6 +220,7 @@ export interface WorkerRuntimeStats { maxBots: number; atCapacity: boolean; inboxDepth: number; + inboxFailed: number; inboxMaxLen: number; inboxPeak: number; inboxDropped: number; @@ -293,9 +300,9 @@ export class BotWorkerManager { private loopGen = new Map(); private peerLimiter: RateLimiter; - /** In-process job queue (same container; no extra Redis BLPOP connection). */ - private inbox: LocalInboundJob[] = []; - private inboxWaiters: Array<() => void> = []; + private claimingInbox = false; + private inboxScanAfter = 0; + private inboxBotOffset = 0; /** Serialize replies per bot:peer so bubbles stay ordered. */ private peerChains = new Map>(); @@ -324,16 +331,6 @@ export class BotWorkerManager { private p2pEnabled: boolean; /** Prevent concurrent OTA apply on the same process */ private otaInFlight = false; - /** - * Process-local LRU of inbound dedup keys. iLink redelivers the same - * message on a non-advancing cursor / loop restart; without this gate - * each redelivery becomes a full LLM regeneration (near-duplicate replies). - * Redis NX covers multi-replica; this covers the single-process hot path - * without a round trip. - */ - private inboundSeen = new Map(); - private static readonly INBOUND_SEEN_TTL_MS = 10 * 60_000; - private static readonly INBOUND_SEEN_MAX = 2_000; constructor(private opts: WorkerOptions) { this.workerId = @@ -574,6 +571,8 @@ export class BotWorkerManager { const leasedLocal = this.loops.size; const atCapacity = leasedLocal >= this.maxBots && pollable > leasedLocal; + const inbox = await inboundQueueStats(this.opts.db, [...this.loops.keys()]); + this.inboxPeak = Math.max(this.inboxPeak, inbox.pending); return { workerId: this.workerId, @@ -584,7 +583,8 @@ export class BotWorkerManager { pollable, maxBots: this.maxBots, atCapacity, - inboxDepth: this.inbox.length, + inboxDepth: inbox.pending, + inboxFailed: inbox.failed, inboxMaxLen: this.inboxMaxLen, inboxPeak: this.inboxPeak, inboxDropped: this.inboxDropped, @@ -894,6 +894,7 @@ export class BotWorkerManager { botId: string, peerId: string, text: string, + clientId?: string, ): Promise { const body = text?.trim(); if (!body) return { ok: false, reason: "empty" }; @@ -932,6 +933,7 @@ export class BotWorkerManager { toUserId: peerId, text: body, contextToken, + clientId, }); return { ok: true }; } catch (err) { @@ -986,9 +988,6 @@ export class BotWorkerManager { } this.loops.clear(); this.clients.clear(); - // Wake reply consumers so they can exit - for (const w of this.inboxWaiters) w(); - this.inboxWaiters = []; this.unregisterPromise = unregisterWorker( this.opts.db, this.workerId, @@ -999,17 +998,20 @@ export class BotWorkerManager { private unregisterPromise: Promise | null = null; /** - * Stop and wait for the fleet deregistration to land. Without awaiting it, - * SIGTERM races the exit and other nodes wait out LEASE_TTL_SEC before - * re-claiming this node's bots. + * Release polling ownership and drain started replies within the shutdown + * budget. Unacknowledged jobs remain recoverable in the shared inbox. */ - async stopAsync(timeoutMs = 5000): Promise { + async stopAsync(timeoutMs = 25_000): Promise { this.stop(); - if (!this.unregisterPromise) return; - await Promise.race([ - this.unregisterPromise, - new Promise((r) => setTimeout(r, timeoutMs).unref?.()), - ]); + let timer: ReturnType | undefined; + // Queued jobs stay in Redis. Give active jobs a chance to ACK before exit; + // any remaining processing claims expire and are recovered by another node. + try { + await Promise.race([ + Promise.allSettled([this.unregisterPromise, ...this.replyWorkers]), + new Promise((r) => { timer = setTimeout(r, timeoutMs); }), + ]); + } finally { if (timer) clearTimeout(timer); } } /** Subscribe to fleet wake channel so new bots attach without waiting full lease tick. */ @@ -1559,22 +1561,7 @@ export class BotWorkerManager { allowInstall: this.opts.otaAllowInstall !== false, log: this.opts.log, beforeRestart: async () => { - // Drain fleet membership so peers pick up bots quickly - this.stopped = true; - this.runtimeActive = false; - try { - const owned = [...this.loops.keys()]; - if (owned.length) { - await releaseOwnedLeasesBatch( - this.opts.db, - this.workerId, - owned, - ); - } - } catch { - /* */ - } - this.stop(); + await this.stopAsync(); }, }); if (result === "applied") { @@ -1808,7 +1795,7 @@ export class BotWorkerManager { const buf = res.get_updates_buf; if (typeof buf === "string" && buf !== "" && buf !== cursor) { try { - await setBotCursor(this.opts.db, botId, buf); + await setBotCursor(this.opts.db, botId, buf, this.workerId); cursor = buf; } catch (e) { // Still advance in-memory so we don't hot-loop the same batch; @@ -1873,60 +1860,6 @@ export class BotWorkerManager { return createHash("sha1").update(raw).digest("hex").slice(0, 32); } - /** True when this key was already claimed (local LRU and/or Redis NX). */ - private async claimInboundOnce(dedupKey: string): Promise { - const now = Date.now(); - // Expire stale local entries cheaply on the hot path. - const localAt = this.inboundSeen.get(dedupKey); - if (localAt !== undefined) { - if (now - localAt < BotWorkerManager.INBOUND_SEEN_TTL_MS) { - return false; - } - this.inboundSeen.delete(dedupKey); - } - // Bound the map so a long-lived process cannot grow forever. - if (this.inboundSeen.size >= BotWorkerManager.INBOUND_SEEN_MAX) { - const cutoff = now - BotWorkerManager.INBOUND_SEEN_TTL_MS; - for (const [k, t] of this.inboundSeen) { - if (t < cutoff) this.inboundSeen.delete(k); - } - // Still full? drop oldest half. - if (this.inboundSeen.size >= BotWorkerManager.INBOUND_SEEN_MAX) { - const keys = [...this.inboundSeen.keys()].slice( - 0, - Math.floor(BotWorkerManager.INBOUND_SEEN_MAX / 2), - ); - for (const k of keys) this.inboundSeen.delete(k); - } - } - - // Multi-replica: Redis NX is the real gate. Local map is a free filter. - try { - const ok = await this.opts.db.redis.set( - K.inboundSeen(dedupKey), - "1", - "EX", - Math.ceil(BotWorkerManager.INBOUND_SEEN_TTL_MS / 1000), - "NX", - ); - if (ok !== "OK") { - // Another node (or an earlier attempt) already claimed it. - this.inboundSeen.set(dedupKey, now); - return false; - } - } catch (e) { - // Redis blip: fall through to local-only claim so we do not drop the - // message entirely. At-most-once across replicas is best-effort then. - this.opts.log?.( - `[worker] inbound dedup redis failed: ${ - e instanceof Error ? e.message : String(e) - }`, - ); - } - this.inboundSeen.set(dedupKey, now); - return true; - } - private async enqueueFromMessage( botId: string, client: ILinkClient, @@ -1947,124 +1880,85 @@ export class BotWorkerManager { const mediaOnly = !text?.trim(); const trimmed = text?.trim() ?? ""; - // Claim BEFORE startTyping / inbox push so a redelivered inbound never - // burns a typing indicator or an LLM slot. iLink has no msg_id; the key - // is synthesized from create_time_ms + content (see inboundDedupKey). - const dedupKey = this.inboundDedupKey( - botId, - peerId, - msg, - trimmed, - mediaRefs, - ); - const first = await this.claimInboundOnce(dedupKey); - if (!first) { - this.opts.log?.( - `[worker] drop redelivered inbound bot=${botId} peer=${peerId} key=${dedupKey.slice(0, 8)}`, - ); - emitActivity({ - type: "message.dedup", - level: "info", - source: this.workerId, - summary: `dedup drop bot=${botId} peer=${peerId}`, - data: { botId, peerId, dedupKey }, - }); - return; - } - - // Typing ASAP; reply path may be delayed by queue. First call per peer also - // fetches the typing_ticket, which is then cached for ~20h. - void client - .startTyping({ toUserId: peerId, contextToken }) - .catch(() => undefined); - - if (this.inbox.length >= this.inboxMaxLen) { - this.inboxDropped++; - this.lastInboxDropAt = new Date().toISOString(); - this.opts.log?.( - `[worker] inbox full (${this.inboxMaxLen}), drop msg bot=${botId} peer=${peerId} dropped=${this.inboxDropped}`, - ); - emitActivity({ - type: "worker.inbox_drop", - level: "warn", - source: this.workerId, - summary: `inbox full drop bot=${botId} peer=${peerId} total=${this.inboxDropped}`, - data: { - botId, - peerId, - inboxMaxLen: this.inboxMaxLen, - dropped: this.inboxDropped, - }, - }); - return; - } - const job: LocalInboundJob = { - id: `job_${Date.now().toString(36)}_${Math.random().toString(36).slice(2, 10)}`, - botId, - peerId, - contextToken, - text: trimmed, - mediaOnly, + id: this.inboundDedupKey(botId, peerId, msg, trimmed, mediaRefs), + botId, peerId, contextToken, text: trimmed, mediaOnly, enqueuedAt: new Date().toISOString(), ...(mediaRefs.length ? { mediaRefs } : {}), }; + // Persist before advancing the iLink cursor. Failure leaves the entire + // batch replayable; previously accepted jobs deduplicate against Redis. + const result = await persistInbound(this.opts.db, this.workerId, job, this.inboxMaxLen); + if (result === "duplicate") return; + this.inboxScanAfter = 0; + void client.startTyping({ toUserId: peerId, contextToken }).catch(() => undefined); emitActivity({ - type: "message.in", - source: this.workerId, - summary: mediaOnly - ? `in media bot=${botId} peer=${peerId}` - : `in bot=${botId} peer=${peerId} ${(job.text || "").slice(0, 24)}`, - data: streamMessagePayload(job.text || (mediaOnly ? "[media]" : ""), { - botId, - peerId, - role: "user", - jobId: job.id, - mediaOnly: !!mediaOnly, - }), + type: "message.in", source: this.workerId, + summary: `in bot=${botId} peer=${peerId}`, + data: streamMessagePayload(job.text || "[media]", { botId, peerId, jobId: job.id, role: "user" }), }); - this.inbox.push(job); - if (this.inbox.length > this.inboxPeak) { - this.inboxPeak = this.inbox.length; - } - const w = this.inboxWaiters.shift(); - if (w) w(); } - // ── Reply ──────────────────────────────────────────── - - private waitForJob(timeoutMs: number): Promise { - if (this.inbox.length > 0) return Promise.resolve(this.inbox.shift()!); - if (this.stopped) return Promise.resolve(null); - - return new Promise((resolve) => { - let settled = false; - const onReady = () => { - if (settled) return; - settled = true; - clearTimeout(timer); - resolve(this.inbox.shift() ?? null); - }; - this.inboxWaiters.push(onReady); - const timer = setTimeout(() => { - if (settled) return; - settled = true; - const idx = this.inboxWaiters.indexOf(onReady); - if (idx >= 0) this.inboxWaiters.splice(idx, 1); - resolve(this.inbox.shift() ?? null); - }, timeoutMs); - }); + // Only one idle scanner per process, regardless of REPLY_CONCURRENCY. + private async waitForJob(): Promise | null> { + if (this.stopped) return null; + if (this.claimingInbox || Date.now() < this.inboxScanAfter) { + await sleep(100); + return null; + } + this.claimingInbox = true; + try { + const ids = [...this.loops.keys()]; + for (let i = 0; i < ids.length && !this.stopped; i++) { + const botId = ids[this.inboxBotOffset++ % ids.length]!; + const claim = await claimInbound(this.opts.db, botId, this.workerId); + if (claim) return claim; + } + this.inboxScanAfter = Date.now() + 1000; + return null; + } finally { this.claimingInbox = false; } } - private async runReplyConsumer(idx: number): Promise { + private async runReplyConsumer(_idx: number): Promise { while (!this.stopped) { - const job = await this.waitForJob(2000); - if (!job) continue; - await this.enqueuePeerChain(job.botId, job.peerId, () => - this.handleJob(job), - ); + let claim: InboundClaim | null = null; + let delivery: InboundDelivery | undefined; + try { + claim = await this.waitForJob(); + if (!claim) continue; + if (this.stopped) { + await retryInbound(this.opts.db, claim, 0); + break; + } + delivery = new InboundDelivery(this.opts.db, claim); + // A bot can move to another node while this job is in progress. Its + // independent peer claim keeps replies ordered until ACK or expiry. + let rawClient = this.clients.get(claim.job.botId); + if (!rawClient) { + const creds = await this.loadBotToken(claim.job.botId); + if (!creds) throw new Error("No credentials for queued bot"); + rawClient = new ILinkClient({ botToken: creds.botToken, baseUrl: creds.baseUrl ?? undefined }); + } + const active = delivery; + const job = claim.job; + const client = rawClient; + let succeeded = false; + await this.enqueuePeerChain(job.botId, job.peerId, async () => { + await this.handleJob(job, client, active); + succeeded = true; + }); + if (!succeeded) throw new Error("Inbound reply failed"); + await acknowledgeInbound(this.opts.db, claim); + } catch (err) { + delivery?.close(); + this.jobsFailed++; + this.opts.log?.(`[worker] durable reply pending retry: ${err instanceof Error ? err.message : String(err)}`); + if (claim) { + await retryInbound(this.opts.db, claim).catch(() => undefined); + } + await sleep(1000); + } finally { delivery?.close(); } } - void idx; } /** @@ -2074,10 +1968,9 @@ export class BotWorkerManager { * rejection, rate limit, P2P relay hand-off, or crash returned without a stop * the peer would keep seeing "对方正在输入中" until the server timed it out. */ - private async handleJob(job: LocalInboundJob): Promise { - const client = this.clients.get(job.botId); + private async handleJob(job: LocalInboundJob, client: ILinkClient, delivery: InboundDelivery): Promise { try { - await this.handleJobInner(job); + await this.handleJobInner(job, client, delivery); } finally { if (client) { await client @@ -2090,24 +1983,8 @@ export class BotWorkerManager { } } - private async handleJobInner(job: LocalInboundJob): Promise { + private async handleJobInner(job: LocalInboundJob, client: ILinkClient, delivery: InboundDelivery): Promise { const t0 = Date.now(); - const client = this.clients.get(job.botId); - if (!client) { - this.jobsFailed++; - this.opts.log?.( - `[worker] no client for bot=${job.botId}, drop job ${job.id}`, - ); - emitActivity({ - type: "worker.fail", - level: "error", - source: this.workerId, - summary: `no client bot=${job.botId} job=${job.id}`, - data: { botId: job.botId, peerId: job.peerId, jobId: job.id }, - }); - return; - } - // Always refresh context_token so later P2P / proactive push can reach this peer try { await upsertContextToken( @@ -2121,48 +1998,40 @@ export class BotWorkerManager { } const rateKey = `${job.botId}:${job.peerId}`; - if (!this.peerLimiter.tryTake(rateKey)) { - try { - const rateText = "你发得太快啦,请稍等一会儿再聊~"; - await client.sendText({ - toUserId: job.peerId, - text: rateText, - contextToken: job.contextToken, - }); - this.jobsProcessed++; - emitActivity({ - type: "message.reject", - level: "warn", - source: this.workerId, - summary: `rate limit bot=${job.botId} peer=${job.peerId}`, - data: streamMessagePayload(rateText, { - botId: job.botId, - peerId: job.peerId, - role: "system", - reason: "rate_limit", - jobId: job.id, - ms: Date.now() - t0, - }), - }); - } catch { - this.jobsFailed++; - } + if (!(await delivery.step("rate", async () => this.peerLimiter.tryTake(rateKey)))) { + const rateText = "你发得太快啦,请稍等一会儿再聊~"; + await delivery.client(client, "rate-reply").sendText({ + toUserId: job.peerId, + text: rateText, + contextToken: job.contextToken, + }); + this.jobsProcessed++; + emitActivity({ + type: "message.reject", + level: "warn", + source: this.workerId, + summary: `rate limit bot=${job.botId} peer=${job.peerId}`, + data: streamMessagePayload(rateText, { + botId: job.botId, peerId: job.peerId, role: "system", + reason: "rate_limit", jobId: job.id, ms: Date.now() - t0, + }), + }); return; } - // ── P2P intercept (no LLM) ─────────────────────────── - if (this.p2pEnabled && this.p2p) { + // Keep the routing decision stable if P2P settings change between retries. + { try { - const p2pResult = await this.p2p.handleInbound({ + const p2pResult = await delivery.step("p2p", async () => this.p2pEnabled && this.p2p ? this.p2p.handleInbound({ botId: job.botId, peerId: job.peerId, text: job.text, mediaOnly: job.mediaOnly || !job.text.trim(), - }); + }) : { handled: false, localReplies: [], remoteSends: [] }); if (p2pResult.handled) { - for (const text of p2pResult.localReplies) { + for (const [i, text] of p2pResult.localReplies.entries()) { if (!text?.trim()) continue; - await client.sendText({ + await delivery.client(client, `p2p-local:${i}`).sendText({ toUserId: job.peerId, text: text.trim(), contextToken: job.contextToken, @@ -2180,12 +2049,10 @@ export class BotWorkerManager { }), }); } - for (const remote of p2pResult.remoteSends) { - await this.sendToEndpoint( - remote.botId, - remote.peerId, - remote.text, - ); + for (const [i, remote] of p2pResult.remoteSends.entries()) { + await delivery.step(`p2p-remote:${i}`, () => this.sendToEndpoint( + remote.botId, remote.peerId, remote.text, delivery.clientId(`p2p-remote:${i}`), + )); } this.jobsProcessed++; emitActivity({ @@ -2208,31 +2075,8 @@ export class BotWorkerManager { err instanceof Error ? err.message : String(err) }`, ); - // Fall through to roleplay rather than hard-fail chat - } - } - - // Fetch + decrypt attachments here rather than in the poll loop: the long - // poll must stay responsive, and this is already the serialized per-peer - // chain where the LLM call happens. - const attachments = job.mediaRefs?.length - ? await this.buildAttachments(client, job, job.mediaRefs) - : []; - - // Nothing the model could respond to: no text, and no attachment it can - // actually perceive. Answer from a canned line instead of burning a turn. - if (!job.text.trim() && !attachments.some((a) => a.dataUri)) { - try { - await client.sendText({ - toUserId: job.peerId, - text: unreadableMediaReply(job.mediaRefs ?? []), - contextToken: job.contextToken, - }); - this.jobsProcessed++; - } catch { - this.jobsFailed++; + throw err; } - return; } void client @@ -2243,16 +2087,23 @@ export class BotWorkerManager { .catch(() => undefined); try { - const result = await this.opts.chat.handleInbound({ - botAccountId: job.botId, - peerId: job.peerId, - text: job.text.trim(), - contextToken: job.contextToken, - attachments, + // Persist the decision as well as the generated reply. A recovered media + // job must not switch to a canned response after its CDN URL expires. + const result = await delivery.step("chat", async () => { + const attachments = job.mediaRefs?.length + ? await this.buildAttachments(client, job, job.mediaRefs) + : []; + if (!job.text.trim() && !attachments.some((a) => a.dataUri)) { + return { kind: "reject" as const, text: unreadableMediaReply(job.mediaRefs ?? []) }; + } + return this.opts.chat.handleInbound({ + botAccountId: job.botId, peerId: job.peerId, text: job.text.trim(), + contextToken: job.contextToken, attachments, + }); }); if (result.kind === "reject" && result.text) { - await client.sendText({ + await delivery.client(client, "reject").sendText({ toUserId: job.peerId, text: result.text, contextToken: job.contextToken, @@ -2271,33 +2122,36 @@ export class BotWorkerManager { }), }); } else if (result.kind === "reply") { - let parts: ReplyPart[] = - result.parts && result.parts.length > 0 - ? result.parts - : result.bubbles && result.bubbles.length > 0 - ? result.bubbles.map((t) => ({ kind: "text" as const, text: t })) - : result.text - ? [{ kind: "text" as const, text: result.text }] - : []; - // Worker-side validation: never send raw sticker JSON as text - parts = sanitizePartsStripStickerJson(parts, this.maxStickersPerReply()); + const parts = await delivery.step("parts", async () => { + const proposed: ReplyPart[] = + result.parts && result.parts.length > 0 + ? result.parts + : result.bubbles && result.bubbles.length > 0 + ? result.bubbles.map((t) => ({ kind: "text" as const, text: t })) + : result.text + ? [{ kind: "text" as const, text: result.text }] + : []; + // Worker-side validation: never send raw sticker JSON as text. + const sanitized = sanitizePartsStripStickerJson(proposed, this.maxStickersPerReply()); + return this.opts.splitReply === false ? sanitized.slice(0, 1) : sanitized; + }); if (parts.length === 0) { /* empty */ } else { // ChatService already loaded the bot row — it hands the owner back // on the result rather than making us re-read it per message. const ownerUserId = result.ownerUserId || ""; - if (this.opts.splitReply === false || parts.length === 1) { + if (parts.length === 1) { this.opts.log?.( `[worker] send 1 part peer=${job.peerId} kind=${parts[0]!.kind}`, ); - await this.sendReplyPart( - client, + await delivery.step("part:0", () => this.sendReplyPart( + delivery.client(client, "part:0"), job.peerId, job.contextToken, parts[0]!, ownerUserId, - ); + )); } else { this.opts.log?.( `[worker] send ${parts.length} parts peer=${job.peerId}` + @@ -2309,6 +2163,7 @@ export class BotWorkerManager { job.contextToken, parts, ownerUserId, + delivery, ); } const outText = @@ -2363,7 +2218,6 @@ export class BotWorkerManager { }, }); } catch (err) { - this.jobsFailed++; this.opts.log?.( `[worker] chat failed peer=${job.peerId}: ${(err as Error).message}`, ); @@ -2380,15 +2234,7 @@ export class BotWorkerManager { ms: Date.now() - t0, }, }); - try { - await client.sendText({ - toUserId: job.peerId, - text: "抱歉,我这边刚才出了点问题,请稍后再试。", - contextToken: job.contextToken, - }); - } catch { - /* ignore */ - } + throw err; } } @@ -2452,15 +2298,16 @@ export class BotWorkerManager { botId: string, peerId: string, text: string, + clientId?: string, ): Promise { - const result = await this.adminSendText(botId, peerId, text); + const result = await this.adminSendText(botId, peerId, text, clientId); if (!result.ok) { this.opts.log?.( `[worker] p2p send skip bot=${botId} peer=${peerId} reason=${result.reason}${ result.error ? ` err=${result.error}` : "" }`, ); - return; + throw new Error(`P2P delivery failed: ${result.reason}`); } this.opts.log?.( `[worker] p2p sent bot=${botId} peer=${peerId}`, @@ -2473,6 +2320,7 @@ export class BotWorkerManager { contextToken: string, parts: ReplyPart[], ownerUserId: string, + delivery?: InboundDelivery, ): Promise { const list = parts.filter( (p) => @@ -2498,13 +2346,15 @@ export class BotWorkerManager { .catch(() => undefined); await sleep(rand(200, 500)); } - await this.sendReplyPart( - client, + const send = () => this.sendReplyPart( + delivery ? delivery.client(client, `part:${i}`) : client, peerId, contextToken, part, ownerUserId, ); + if (delivery) await delivery.step(`part:${i}`, send); + else await send(); } } @@ -2587,7 +2437,7 @@ export class BotWorkerManager { this.opts.log?.( `[worker] sticker send failed slug=${part.slug}: ${(err as Error).message}`, ); - // Do not fail the whole chain — skip this bubble + throw err; } } } diff --git a/docs/admin-api.md b/docs/admin-api.md index 749e2fdc..906433f5 100644 --- a/docs/admin-api.md +++ b/docs/admin-api.md @@ -130,6 +130,8 @@ Auth: **Cookie 会话**(OAuth 登录后 `wa_session`),`credentials: includ | PATCH | `/api/v1/admin/bots/:botId` | 改名 / `{ status: active\|inactive }` 启停 | | POST | `/api/v1/admin/bots/:botId/stop-worker` | **仅超管** 停止 Worker 轮询 | | POST | `/api/v1/admin/bots/:botId/start-worker` | **仅超管** 启动/重启 Worker(需已有 Redis token) | +| GET | `/api/v1/admin/bots/:botId/inbox` | **仅超管** `pending` / `failed` 计数和最多 100 条失败任务元数据,不含消息正文或凭据 | +| POST | `/api/v1/admin/bots/:botId/inbox/:jobId/retry` | **仅超管** 将失败任务追加回对应聊天对象队尾,保留检查点、重置重试次数;写审计 | | DELETE | `/api/v1/admin/bots/:botId` | 删除机器人 | | GET | `/api/v1/admin/bots/:botId/send-targets` | **仅超管** 该 bot 的 peers + `hasContextToken`(广播勾选) | | POST | `/api/v1/admin/broadcast` | **仅超管** 创建广播任务或 `preview:true` 仅预估人数 | @@ -174,7 +176,7 @@ Auth: **Cookie 会话**(OAuth 登录后 `wa_session`),`credentials: includ | `WORKER_ENABLED` | `true` | 是否在本进程跑 iLink 轮询与回复 | | `MAX_BOTS_PER_WORKER` | `500` | 本进程最多同时 long-poll 的 bot 数 | | `REPLY_CONCURRENCY` | `16` | 进程内 LLM/发送并发 | -| `INBOX_MAX_LEN` | `20000` | 入站队列深度上限 | +| `INBOX_MAX_LEN` | `20000` | 每机器人 Redis 入站保留消息上限(含失败任务);满载不推进游标 | ### 管理后台广播(纯文本推送) @@ -202,7 +204,11 @@ Auth: **Cookie 会话**(OAuth 登录后 `wa_session`),`credentials: includ | `pollable` | 应被轮询的 bot(active + token + 未 pause) | | `nodesOnline` / `nodesTotal` | 在线 / 注册部署节点数 | | `atCapacity` | 任一点触顶或舰队已满 | -| `inboxDepth` 等 | **当前应答节点本机** inbox/任务计数(非全舰队加总) | +| `inboxDepth` / `inboxFailed` | 当前节点所持机器人在 Redis 的待处理 / 失败任务数(非全舰队加总) | +| `inboxPeak` | 当前节点观测到的待处理任务数峰值 | +| `inboxDropped` / `lastInboxDropAt` | 兼容旧客户端的保留字段;持久化队列不主动丢弃,分别为 `0` / `null` | + +失败任务元数据为 `id`、`peerId`、`enqueuedAt`、`error`(通用分类);重试成功为 `{ "ok": true }`,机器人或失败任务不存在返回 404。自动处理最多 3 次,之后保留供人工恢复,详见 [失败消息恢复](./runbook.md#失败消息恢复)。 ### 超管(super admin) diff --git a/docs/docker.md b/docs/docker.md index 829b7412..67b81ca7 100644 --- a/docs/docker.md +++ b/docs/docker.md @@ -178,18 +178,34 @@ server { 默认 **单副本一体部署** 即可;同镜像多副本已支持(Redis 租约分片 poll)。 租约续期与释放使用 Redis Lua 原子检查归属,避免旧节点在接管期间覆盖或删除新节点的租约。 -Redis 账号需要允许 `EVAL`(以及脚本内的 `GET` / `EXPIRE` / `DEL` / `SREM`);批量操作仍通过 pipeline 合并往返。 +Redis 账号需要允许 `EVAL` 及脚本使用的字符串、Hash、List、Set、Sorted Set、`TIME` 命令;批量操作仍通过 pipeline 合并往返。 ### 多节点回归测试 -CI 在 Node 22 / 24 上使用独立 Redis 服务运行租约接管测试。本地可将 -`WECHAT_AI_TEST_REDIS_URL` 指向**专用测试 Redis** 后运行 `pnpm test`;未设置时这组集成测试会明确跳过。 -测试覆盖续期、单个释放、批量释放及超额认领回收与其他节点接管交错的情况。 +CI 在 Node 22 / 24 上使用独立 Redis 服务运行租约和消息恢复测试。本地将 +`WECHAT_AI_TEST_REDIS_URL` 指向**专用测试 Redis**,并将 `REDIS_URL` 指向另一个测试库(旧有核心测试使用),然后运行 `pnpm -r test`;未设置时对应集成测试会明确跳过。 +测试覆盖租约竞态、消息再平衡交接、处理超时恢复、逐片段重试、关闭排空、入队失败不推进游标,以及失败任务接口的鉴权和脱敏。 各包测试命令的文件通配符须保留引号,交由 Node 展开;否则 Linux shell 可能只匹配子目录,遗漏顶层测试。 -**交付边界:** 租约协调的是机器人收消息的归属,不是持久化消息队列。 -待回复任务当前仍在进程内存中,进程崩溃或退出可能丢失未完成任务;自动再平衡也可能移走仍有待回复任务的机器人。 -因此节点接管不代表未完成消息一定重试,也不提供 exactly-once 发送保证。 +### 持久化回复与交付边界 + +入站消息先写 Redis,再推进 iLink 游标。队列按机器人保存,同一聊天对象 FIFO,不同聊天对象可并发。 +机器人迁移后,新节点从共享队列继续处理;旧节点已开始的回复使用独立的 120 秒处理租约(每 40 秒续期),可以完成并 ACK,后续同对象消息等待它结束或超时。 +正常退出和 OTA 最多等待 25 秒排空在途回复;未开始或未 ACK 的消息仍保留,崩溃后由当前机器人归属节点恢复。容器停止宽限期建议至少 30 秒。 + +每条消息最多自动处理 3 次;仍失败则保留原消息和检查点供超管查看、重试(见 [运维手册](./runbook.md#失败消息恢复))。 +`INBOX_MAX_LEN` 现在是**每机器人**保留消息总数,包括失败队列;满载时保留游标并重试收取,不丢弃新消息。 +已保存的生成结果和成功发送片段会复用;重试发送使用稳定的 `client_id`。这是有重试上限的至少一次处理: +如果 iLink 已接受请求但确认丢失,是否去重仍取决于服务端;生成、P2P 状态变更完成但检查点尚未保存的崩溃窗口也可能重做。**不承诺端到端 exactly-once**。 + +各节点必须使用同一份支持 Lua 的共享 Redis(当前实现不支持 Redis Cluster 分片)。 +队列包含消息正文、上下文 token、媒体下载凭据与回复检查点,属于与 Bot token 相同的敏感数据;限制 Redis 访问权限,远程连接使用 TLS,保护备份。 +消息 ACK 后删除队列 payload;失败消息保留直到成功重试或删除机器人,删除机器人也清理其队列并废止处理 token。 +为 Redis 启用符合业务容灾要求的持久化与备份,并配置 `noeviction` 或托管服务等效的禁止逐出策略;恢复保证以 Redis 已保存数据仍在为前提。 +容量不足时应先处理失败消息或增加 Redis 容量,不要使用会随机逐出队列键的策略。 + +**升级:** 旧进程内存中的消息不能由新队列恢复。先停止新的入站流量并等待旧进程在途工作结束,再统一升级所有节点;避免与旧租约/内存队列版本混跑。 +回滚前先排空新队列,旧版本不会消费这些持久化任务。 ## 多节点同构部署(10+ 台) diff --git a/docs/runbook.md b/docs/runbook.md index 1c2eacae..fd3688dc 100644 --- a/docs/runbook.md +++ b/docs/runbook.md @@ -253,9 +253,23 @@ pnpm dev API 与 iLink Worker **同进程**(单镜像 / 单容器)。 - 收消息:`getUpdates` 长轮询(每 bot 一路,有 `MAX_BOTS_PER_WORKER` 上限) -- 回消息:进程内 inbox 队列 + `REPLY_CONCURRENCY` 并发,避免 LLM 堵住轮询 +- 回消息:Redis 持久化 inbox + `REPLY_CONCURRENCY` 并发,同一机器人/聊天对象按序处理,避免 LLM 堵住轮询 - 日志出现 `at capacity`:提高 `MAX_BOTS_PER_WORKER`,或同镜像多副本分片 +### 失败消息恢复 + +入队成功后才推进收取游标。处理租约默认 120 秒、每 40 秒续期;进程丢失后,当前机器人归属节点在租约到期后重新处理,已确认完成的生成/发送步骤会复用。 +连续 3 次处理失败或中断的消息移入失败队列,后续同对象消息可继续。人工重试会把失败消息追加到该聊天对象队尾,重置尝试次数并保留检查点,因此较新的消息可能已先回复。 + +1. 用超管登录会话调用 `GET /api/v1/admin/bots/:botId/inbox`,读取 `pending`、`failed` 和最多 100 条 `failedJobs` 元数据。 +2. 检查应用日志并修复实际原因(如 LLM/iLink 不可用、账号会话过期)。查询接口不返回正文、context token 或媒体密钥。 +3. 对选中的任务调用 `POST /api/v1/admin/bots/:botId/inbox/:jobId/retry`,成功返回 `{ "ok": true }` 并记录 `admin_inbound_retry` 审计。 + +`INBOX_MAX_LEN` 按机器人计算,包含失败消息。满载时收取批次重试、游标不前进;留意 Redis 内存和失败数量。 +队列依赖共享 Redis 的持久化且禁止逐出,ACK 后清理 payload,失败消息不会自动到期;删除机器人会一并清理队列。 +发送确认丢失时会用相同 `client_id` 重试,但未经真实 iLink 去重验证,不能视为 exactly-once。 +节点升级、回滚和敏感数据保护要求见 [Docker 部署文档](./docker.md#持久化回复与交付边界)。 + ## 5. 故障 | 现象 | 处理 | @@ -278,6 +292,7 @@ API 与 iLink Worker **同进程**(单镜像 / 单容器)。 - Upstash:控制台备份 / 导出(按套餐) - Redis `wa:bot:{id}:creds`(Bot token,敏感) +- Redis `wa:bot:{id}:inbound:*`(待处理/失败消息、context token、媒体凭据、回复检查点,敏感;备份须与其余机器人数据保持一致) - `.env`(勿提交 Git) ## 7. 文档 diff --git a/packages/db/src/inbound-queue.test.ts b/packages/db/src/inbound-queue.test.ts new file mode 100644 index 00000000..6b431742 --- /dev/null +++ b/packages/db/src/inbound-queue.test.ts @@ -0,0 +1,170 @@ +import { after, before, describe, it } from "node:test"; +import assert from "node:assert/strict"; +import { randomUUID } from "node:crypto"; +import { openDatabase, type RedisStore } from "./client.js"; +import { K } from "./keys.js"; +import { deleteBotAccount, setBotCursor } from "./repos.js"; +import type { InboundJob } from "./worker-fleet.js"; +import { + persistInbound, claimInbound, renewInbound, saveInboundSteps, + acknowledgeInbound, inboundQueueKeys, inboundQueueStats, + listFailedInbound, retryFailedInbound, type InboundClaim, +} from "./inbound-queue.js"; + +const redisUrl = process.env.WECHAT_AI_TEST_REDIS_URL; +describe("durable inbound delivery (real Redis)", { skip: !redisUrl }, () => { + let db: RedisStore; + const prefix = `inbox-test-${randomUUID()}`; + const jobs: InboundJob[] = []; + const bots: string[] = []; + before(async () => { db = openDatabase(redisUrl!); await db.ping(); }); + after(async () => { + if (!db) return; + try { + for (const bot of bots) await db.redis.del(K.bot(bot), K.botLease(bot), ...Object.values(inboundQueueKeys(bot))); + for (const job of jobs) await db.redis.del(K.inboundSeen(job.id), inboundQueueKeys(job.botId).peers + job.peerId); + } finally { await db.close(); } + }); + async function bot() { + const id = `${prefix}-${bots.length}`; + bots.push(id); + await db.redis.set(K.bot(id), "{}"); + await db.redis.set(K.botLease(id), "a", "EX", 45); + return id; + } + function job(botId: string, peerId = "peer"): InboundJob { + const value = { id: `${prefix}-job-${jobs.length}`, botId, peerId, text: "hello", contextToken: "private-context", mediaOnly: false, enqueuedAt: new Date().toISOString() }; + jobs.push(value); + return value; + } + async function expire(claim: InboundClaim) { + await db.redis.zadd(inboundQueueKeys(claim.job.botId).active, 0, claim.job.peerId); + } + + it("atomically deduplicates acceptance and leaves rejected messages replayable", async () => { + const id = await bot(), first = job(id), second = job(id); + const results = await Promise.all([persistInbound(db, "a", first, 1), persistInbound(db, "a", first, 1)]); + assert.deepEqual(results.sort(), ["accepted", "duplicate"]); + await assert.rejects(persistInbound(db, "a", second, 1), /full/); + assert.equal(await db.redis.exists(K.inboundSeen(second.id)), 0); + const claim = (await claimInbound(db, id, "a"))!; + await acknowledgeInbound(db, claim); + assert.equal(await persistInbound(db, "a", second, 1), "accepted"); + assert.equal(await persistInbound(db, "a", first, 1), "duplicate"); + }); + + it("keeps same-peer order across handoff while allowing other peers to progress", async () => { + const id = await bot(), first = job(id), second = job(id), other = job(id, "other"); + for (const item of [first, second, other]) await persistInbound(db, "a", item, 10); + const old = (await claimInbound(db, id, "a"))!; + assert.equal(old.job.id, first.id); + await db.redis.set(K.botLease(id), "b", "EX", 45); + assert.equal(await claimInbound(db, id, "a"), null); + const parallel = (await claimInbound(db, id, "b"))!; + assert.equal(parallel.job.id, other.id); + assert.equal(await claimInbound(db, id, "b"), null); + assert.equal(await renewInbound(db, old), true, "old in-flight job may finish independently of polling ownership"); + await expire(old); + const recovered = (await claimInbound(db, id, "b"))!; + assert.equal(recovered.job.id, first.id); + assert.equal(recovered.attempts, 2); + await assert.rejects(acknowledgeInbound(db, old), /lease lost/); + await assert.rejects(saveInboundSteps(db, old), /lease lost/); + await acknowledgeInbound(db, recovered); + const next = (await claimInbound(db, id, "b"))!; + assert.equal(next.job.id, second.id); + await acknowledgeInbound(db, next); + await acknowledgeInbound(db, parallel); + assert.deepEqual(await inboundQueueStats(db, [id]), { pending: 0, failed: 0 }); + }); + + it("recovers media coordinates and completed generation/send checkpoints after a lost worker", async () => { + const id = await bot(); + const item = { ...job(id), mediaRefs: [{ kind: "image", aesKey: "private-media-key", encryptQueryParam: "private-query" }] }; + await persistInbound(db, "a", item, 10); + const old = (await claimInbound(db, id, "a"))!; + old.steps.chat = { value: { text: "generated once" } }; + old.steps["send:0:sendText"] = { value: { ret: 0 } }; + await saveInboundSteps(db, old); + await expire(old); + await db.close(); + db = openDatabase(redisUrl!); + await db.redis.set(K.botLease(id), "b", "EX", 45); + const recovered = (await claimInbound(db, id, "b"))!; + assert.deepEqual(recovered.job, item); + assert.deepEqual(recovered.steps, old.steps); + assert.equal(await persistInbound(db, "b", item, 10), "duplicate"); + await acknowledgeInbound(db, recovered); + }); + + it("retains exhausted work for inspection and explicit retry without exposing payloads", async () => { + const id = await bot(), item = job(id); + await persistInbound(db, "a", item, 10); + for (let n = 0; n < 3; n++) { + const claim = (await claimInbound(db, id, "a"))!; + assert.equal(claim.attempts, n + 1); + await expire(claim); + } + assert.equal(await claimInbound(db, id, "a"), null); + assert.deepEqual(await inboundQueueStats(db, [id]), { pending: 0, failed: 1 }); + const failed = await listFailedInbound(db, id); + assert.equal(failed[0]!.id, item.id); + assert.ok(!JSON.stringify(failed).includes(item.contextToken)); + assert.ok(!Object.hasOwn(failed[0]!, "text")); + assert.equal(await retryFailedInbound(db, id, item.id), true); + assert.equal(await retryFailedInbound(db, id, item.id), false); + const retried = (await claimInbound(db, id, "a"))!; + assert.equal(retried.attempts, 1); + await acknowledgeInbound(db, retried); + assert.equal(await db.redis.hlen(inboundQueueKeys(id).jobs), 0); + }); + + it("never renews an expired processing claim or accepts from a stale polling node", async () => { + const id = await bot(), item = job(id); + await persistInbound(db, "a", item, 10); + const claim = (await claimInbound(db, id, "a"))!; + await expire(claim); + assert.equal(await renewInbound(db, claim), false); + await db.redis.set(K.botLease(id), "b", "EX", 45); + const next = job(id); + await assert.rejects(persistInbound(db, "a", next, 10), /lease lost/); + assert.equal(await db.redis.hexists(inboundQueueKeys(id).jobs, next.id), 0); + }); + + it("deletes queued, active and failed payloads with the bot and fences old processors", async () => { + const id = await bot(); + const active = job(id, "active"), queued = job(id, "queued"), failed = job(id, "failed"); + for (const item of [active, queued, failed]) await persistInbound(db, "a", item, 10); + const first = (await claimInbound(db, id, "a"))!; + const second = (await claimInbound(db, id, "a"))!; + const third = (await claimInbound(db, id, "a"))!; + third.steps.chat = { value: "sensitive generated reply" }; + await saveInboundSteps(db, third); + await expire(third); + assert.equal(await claimInbound(db, id, "a", 120_000, 1), null); + assert.equal((await listFailedInbound(db, id)).length, 1); + await persistInbound(db, "a", job(id, "unstarted"), 10); + assert.equal(await deleteBotAccount(db, id), true); + assert.equal(await renewInbound(db, first), false); + await assert.rejects(acknowledgeInbound(db, second), /lease lost/); + for (const item of jobs.filter((x) => x.botId === id)) { + assert.equal(await db.redis.exists(inboundQueueKeys(id).peers + item.peerId), 0); + } + for (const key of Object.values(inboundQueueKeys(id))) assert.equal(await db.redis.exists(key), 0); + await db.redis.set(K.botLease(id), "a", "EX", 45); // delayed stale owner + await assert.rejects(persistInbound(db, "a", active, 10), /lease lost/); + assert.equal(await claimInbound(db, id, "a"), null); + await assert.rejects(setBotCursor(db, id, "stale", "a"), /lease lost/); + assert.equal(await db.redis.exists(K.bot(id)), 0); + }); + + it("fences cursor commits after polling ownership changes", async () => { + const id = await bot(); + await setBotCursor(db, id, "accepted", "a"); + await db.redis.set(K.botLease(id), "b", "EX", 45); + await assert.rejects(setBotCursor(db, id, "stale", "a"), /lease lost/); + assert.equal((await db.getJson<{ updates_cursor: string }>(K.bot(id)))!.updates_cursor, "accepted"); + await setBotCursor(db, id, "successor", "b"); + assert.equal((await db.getJson<{ updates_cursor: string }>(K.bot(id)))!.updates_cursor, "successor"); + }); +}); diff --git a/packages/db/src/inbound-queue.ts b/packages/db/src/inbound-queue.ts new file mode 100644 index 00000000..9ca0b520 --- /dev/null +++ b/packages/db/src/inbound-queue.ts @@ -0,0 +1,196 @@ +import { randomUUID } from "node:crypto"; +import type { RedisStore } from "./client.js"; +import { K } from "./keys.js"; +import type { InboundJob } from "./worker-fleet.js"; + +export interface InboundClaim { + job: T; + token: string; + attempts: number; + steps: Record; +} + +export const INBOUND_LEASE_MS = 120_000; +export const INBOUND_MAX_ATTEMPTS = 3; +const DONE_TTL_SEC = 600; + +// All keys live in the existing shared, non-cluster Redis database. Peer lists +// are derived from this bot's namespace inside Lua, never from caller key names. +export function inboundQueueKeys(botId: string) { + const base = `wa:bot:${botId}:inbound:`; + return { + jobs: `${base}jobs`, ready: `${base}ready`, active: `${base}active`, + claims: `${base}claims`, attempts: `${base}attempts`, steps: `${base}steps`, + failed: `${base}failed`, errors: `${base}errors`, sequence: `${base}seq`, + peers: `${base}peer:`, + }; +} + +const TIME = "local t=redis.call('TIME'); local now=t[1]*1000+math.floor(t[2]/1000)\n"; + +/** Atomically accept a message; a full/unavailable inbox never marks it seen. */ +export async function persistInbound( + db: RedisStore, workerId: string, job: InboundJob, maxLen: number, +): Promise<"accepted" | "duplicate"> { + const q = inboundQueueKeys(job.botId); + const result = Number(await db.redis.eval(` +if redis.call('GET',KEYS[1]) ~= ARGV[1] or redis.call('EXISTS',KEYS[8]) == 0 then return -2 end +if redis.call('HEXISTS',KEYS[2],ARGV[2]) == 1 or redis.call('EXISTS',KEYS[3]) == 1 then return 0 end +if redis.call('HLEN',KEYS[2]) >= tonumber(ARGV[5]) then return -1 end +redis.call('HSET',KEYS[2],ARGV[2],ARGV[4]) +local n=redis.call('RPUSH',KEYS[4],ARGV[2]) +if n == 1 and not redis.call('ZSCORE',KEYS[6],ARGV[3]) then + redis.call('ZADD',KEYS[5],redis.call('INCR',KEYS[7]),ARGV[3]) +end +return 1`, 8, K.botLease(job.botId), q.jobs, K.inboundSeen(job.id), + q.peers + job.peerId, q.ready, q.active, q.sequence, + K.bot(job.botId), + workerId, job.id, job.peerId, JSON.stringify(job), maxLen)); + if (result === -2) throw new Error("Inbound bot lease lost"); + if (result === -1) throw new Error("Durable inbox full; cursor not advanced"); + return result === 1 ? "accepted" : "duplicate"; +} + +export async function claimInbound( + db: RedisStore, botId: string, workerId: string, + leaseMs = INBOUND_LEASE_MS, maxAttempts = INBOUND_MAX_ATTEMPTS, +): Promise | null> { + const q = inboundQueueKeys(botId); + const token = randomUUID(); + const raw = await db.redis.eval(TIME + ` +if redis.call('GET',KEYS[1]) ~= ARGV[1] or redis.call('EXISTS',KEYS[11]) == 0 then return nil end +local expired=redis.call('ZRANGEBYSCORE',KEYS[3],'-inf',now,'LIMIT',0,32) +for _,peer in ipairs(expired) do + redis.call('ZREM',KEYS[3],peer) + redis.call('HDEL',KEYS[4],peer) + redis.call('ZADD',KEYS[2],redis.call('INCR',KEYS[10]),peer) +end +local peers=redis.call('ZRANGE',KEYS[2],0,0) +if #peers == 0 then return nil end +local peer=peers[1] +local list=ARGV[4]..peer +local id=redis.call('LINDEX',list,0) +if not id then return redis.error_reply('Missing inbound peer queue') end +local payload=redis.call('HGET',KEYS[5],id) +if not payload then return redis.error_reply('Missing inbound payload') end +redis.call('ZREM',KEYS[2],peer) +local attempts=redis.call('HINCRBY',KEYS[6],id,1) +if attempts > tonumber(ARGV[5]) then + redis.call('LPOP',list) + redis.call('ZADD',KEYS[8],now,id) + redis.call('HSET',KEYS[9],id,'retry_exhausted') + if redis.call('LLEN',list) > 0 then redis.call('ZADD',KEYS[2],redis.call('INCR',KEYS[10]),peer) end + return nil +end +redis.call('ZADD',KEYS[3],now+tonumber(ARGV[3]),peer) +redis.call('HSET',KEYS[4],peer,ARGV[2]) +return {payload,tostring(attempts),redis.call('HGET',KEYS[7],id) or '{}'}`, + 11, K.botLease(botId), q.ready, q.active, q.claims, q.jobs, q.attempts, + q.steps, q.failed, q.errors, q.sequence, K.bot(botId), + workerId, token, leaseMs, q.peers, maxAttempts) as string[] | null; + if (!raw) return null; + return { job: JSON.parse(raw[0]!) as T, token, attempts: Number(raw[1]), steps: JSON.parse(raw[2]!) }; +} + +const OWNED = ` +if redis.call('HGET',KEYS[1],ARGV[1]) ~= ARGV[2] then return 0 end +local deadline=redis.call('ZSCORE',KEYS[2],ARGV[1]) +if not deadline or tonumber(deadline) <= now then return 0 end +`; + +export async function renewInbound(db: RedisStore, claim: InboundClaim, leaseMs = INBOUND_LEASE_MS): Promise { + const q = inboundQueueKeys(claim.job.botId); + return Number(await db.redis.eval(TIME + OWNED + ` +redis.call('ZADD',KEYS[2],now+tonumber(ARGV[3]),ARGV[1]); return 1`, + 2, q.claims, q.active, claim.job.peerId, claim.token, leaseMs)) === 1; +} + +export async function saveInboundSteps(db: RedisStore, claim: InboundClaim): Promise { + const q = inboundQueueKeys(claim.job.botId); + const ok = await db.redis.eval(TIME + OWNED + ` +redis.call('HSET',KEYS[3],ARGV[3],ARGV[4]); return 1`, + 3, q.claims, q.active, q.steps, claim.job.peerId, claim.token, + claim.job.id, JSON.stringify(claim.steps)); + if (Number(ok) !== 1) throw new Error("Inbound processing lease lost"); +} + +/** ACK is the only normal path that deletes a pending payload. */ +export async function acknowledgeInbound(db: RedisStore, claim: InboundClaim): Promise { + const q = inboundQueueKeys(claim.job.botId); + const ok = await db.redis.eval(TIME + OWNED + ` +if redis.call('LINDEX',KEYS[3],0) ~= ARGV[3] then return 0 end +redis.call('LPOP',KEYS[3]) +redis.call('HDEL',KEYS[1],ARGV[1]); redis.call('ZREM',KEYS[2],ARGV[1]) +for i=4,7 do redis.call('HDEL',KEYS[i],ARGV[3]) end +redis.call('SET',KEYS[8],'1','EX',ARGV[4]) +if redis.call('LLEN',KEYS[3]) > 0 then redis.call('ZADD',KEYS[9],redis.call('INCR',KEYS[10]),ARGV[1]) end +return 1`, 10, q.claims, q.active, q.peers + claim.job.peerId, + q.jobs, q.steps, q.attempts, q.errors, K.inboundSeen(claim.job.id), q.ready, q.sequence, + claim.job.peerId, claim.token, claim.job.id, DONE_TTL_SEC); + if (Number(ok) !== 1) throw new Error("Inbound acknowledgement lease lost"); +} + +/** Retain payload/checkpoints, release the peer after a bounded retry delay. */ +export async function retryInbound(db: RedisStore, claim: InboundClaim, delayMs = 2000): Promise { + const q = inboundQueueKeys(claim.job.botId); + await db.redis.eval(TIME + OWNED + ` +redis.call('HSET',KEYS[3],ARGV[3],'processing_failed') +redis.call('ZADD',KEYS[2],now+tonumber(ARGV[4]),ARGV[1]); return 1`, + 3, q.claims, q.active, q.errors, claim.job.peerId, claim.token, claim.job.id, delayMs); +} + +export async function inboundQueueStats(db: RedisStore, botIds: string[]): Promise<{ pending: number; failed: number }> { + if (!botIds.length) return { pending: 0, failed: 0 }; + const pipe = db.redis.pipeline(); + for (const id of botIds) { + const q = inboundQueueKeys(id); + pipe.hlen(q.jobs).zcard(q.failed); + } + const rows = await pipe.exec(); + if (!rows) throw new Error("Missing inbox stats"); + for (const [err] of rows) if (err) throw err; + let pending = 0, failed = 0; + for (let i = 0; i < rows.length; i += 2) { + const count = Number(rows[i + 1]![1]); + pending += Number(rows[i]![1]) - count; + failed += count; + } + return { pending, failed }; +} + +export async function listFailedInbound(db: RedisStore, botId: string) { + const q = inboundQueueKeys(botId); + const ids = await db.redis.zrange(q.failed, 0, 99); + if (!ids.length) return []; + const [jobs, errors] = await Promise.all([db.redis.hmget(q.jobs, ...ids), db.redis.hmget(q.errors, ...ids)]); + return ids.map((id, i) => { + const job = jobs[i] ? JSON.parse(jobs[i]!) as InboundJob : null; + return { id, peerId: job?.peerId, enqueuedAt: job?.enqueuedAt, error: errors[i] }; + }); +} + +export async function retryFailedInbound(db: RedisStore, botId: string, id: string): Promise { + const q = inboundQueueKeys(botId); + return Number(await db.redis.eval(` +local raw=redis.call('HGET',KEYS[1],ARGV[1]) +if not raw or not redis.call('ZSCORE',KEYS[2],ARGV[1]) then return 0 end +local job=cjson.decode(raw); local list=ARGV[2]..job.peerId +redis.call('ZREM',KEYS[2],ARGV[1]); redis.call('HDEL',KEYS[3],ARGV[1]); redis.call('HDEL',KEYS[4],ARGV[1]) +local n=redis.call('RPUSH',list,ARGV[1]) +if n == 1 and not redis.call('ZSCORE',KEYS[6],job.peerId) then + redis.call('ZADD',KEYS[5],redis.call('INCR',KEYS[7]),job.peerId) +end +return 1`, 7, q.jobs, q.failed, q.attempts, q.errors, q.ready, q.active, q.sequence, id, q.peers)) === 1; +} + +/** Called after deleting the bot record; also invalidates in-flight claim tokens. */ +export async function deleteInboundQueue(db: RedisStore, botId: string): Promise { + const q = inboundQueueKeys(botId); + const keys = [q.ready, q.active, q.jobs, q.claims, q.attempts, q.steps, q.failed, q.errors, q.sequence]; + await db.redis.eval(` +for i=1,2 do + for _,peer in ipairs(redis.call('ZRANGE',KEYS[i],0,-1)) do redis.call('DEL',ARGV[1]..peer) end +end +for _,key in ipairs(KEYS) do redis.call('DEL',key) end +return 1`, keys.length, ...keys, q.peers); +} diff --git a/packages/db/src/index.ts b/packages/db/src/index.ts index c9c0162a..7892f3ce 100644 --- a/packages/db/src/index.ts +++ b/packages/db/src/index.ts @@ -12,5 +12,6 @@ export * from "./broadcast-repos.js"; export * from "./bot-login-repos.js"; export * from "./seed.js"; export * from "./worker-fleet.js"; +export * from "./inbound-queue.js"; export * from "./ota-paths.js"; export * from "./ota-repos.js"; diff --git a/packages/db/src/repos.ts b/packages/db/src/repos.ts index 67faf8bb..caa38126 100644 --- a/packages/db/src/repos.ts +++ b/packages/db/src/repos.ts @@ -2,6 +2,7 @@ import crypto from "node:crypto"; import type { RedisStore } from "./client.js"; import { newId, nowIso, dayKey } from "./client.js"; import { K } from "./keys.js"; +import { deleteInboundQueue } from "./inbound-queue.js"; import { defaultPersonaIdCache, invalidateDefaultPersonaCache, @@ -938,6 +939,7 @@ export async function deleteBotAccount( await forceReleaseBotLease(db, botId); await db.del(K.bot(botId)); await deleteBotCredentials(db, botId); + await deleteInboundQueue(db, botId); await db.redis.srem(K.botsAll, botId); if (bot.owner_user_id) { await db.redis.srem(K.botsByOwner(bot.owner_user_id), botId); @@ -1044,12 +1046,19 @@ export async function setBotCursor( db: RedisStore, botId: string, cursor: string, + workerId?: string, ): Promise { - const bot = await getBotAccount(db, botId); - if (!bot) return; - bot.updates_cursor = cursor; - bot.updated_at = nowIso(); - await db.setJson(K.bot(botId), bot); + // A delayed poller must not overwrite the successor's cursor or resurrect + // a deleted bot row after its batch has been persisted. + const ok = await db.redis.eval(` +local raw=redis.call('GET',KEYS[1]) +if not raw then return 0 end +if ARGV[1] ~= '' and redis.call('GET',KEYS[2]) ~= ARGV[1] then return 0 end +local bot=cjson.decode(raw) +bot.updates_cursor=ARGV[2]; bot.updated_at=ARGV[3] +redis.call('SET',KEYS[1],cjson.encode(bot)); return 1`, + 2, K.bot(botId), K.botLease(botId), workerId ?? "", cursor, nowIso()); + if (workerId && Number(ok) !== 1) throw new Error("Inbound bot lease lost"); } export async function setBotStatus(