Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 2 additions & 0 deletions .env.example
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
17 changes: 15 additions & 2 deletions .github/workflows/ci.yml
Original file line number Diff line number Diff line change
Expand Up @@ -19,11 +19,24 @@ 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 }}
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

Expand Down
4 changes: 2 additions & 2 deletions apps/api/package.json
Original file line number Diff line number Diff line change
Expand Up @@ -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",
Expand All @@ -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"
},
Expand Down
1 change: 1 addition & 0 deletions apps/api/src/config.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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 */
Expand Down
58 changes: 58 additions & 0 deletions apps/api/src/dependency-regressions.test.ts
Original file line number Diff line number Diff line change
@@ -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<string>((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();
}
});
90 changes: 90 additions & 0 deletions apps/api/src/inbound-admin.test.ts
Original file line number Diff line number Diff line change
@@ -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<string, string>();
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);
});
});
93 changes: 93 additions & 0 deletions apps/api/src/inbound-delivery.test.ts
Original file line number Diff line number Diff line change
@@ -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(); }
});
});
63 changes: 63 additions & 0 deletions apps/api/src/inbound-delivery.ts
Original file line number Diff line number Diff line change
@@ -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<typeof setInterval>;

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<void> {
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<T>(name: string, run: () => Promise<T>): Promise<T> {
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<string, unknown>) => {
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;
},
});
}
}
Loading