fix: use passive QQ event replies
This commit is contained in:
@@ -437,6 +437,85 @@ test("envelope parsers accept valid tails and reject invalid ones", () => {
|
||||
assert.equal(parseWorkerResult(`x\n<GORI_WORKER_RESULT_V1>{"status":"DONE","summary":"s"}</GORI_WORKER_RESULT_V1>`), undefined);
|
||||
});
|
||||
|
||||
test("start_next blocked by another user's pending confirmation appends a correction and exposes no details", async () => {
|
||||
const harness = await createHarness();
|
||||
try {
|
||||
await harness.manager.prompt(request("create proposal: succeed", "user-a"));
|
||||
const first = harness.proposals.list()[0]!;
|
||||
await harness.manager.prompt(request("confirm", "user-a"));
|
||||
await waitFor(() => harness.proposals.get(first.id)!.status === "awaiting_user_confirmation");
|
||||
|
||||
await harness.manager.prompt(request("create proposal: succeed", "user-b"));
|
||||
const second = harness.proposals.list().find((proposal) => proposal.id !== first.id)!;
|
||||
await harness.manager.prompt(request("confirm", "user-b"));
|
||||
assert.equal(harness.proposals.get(second.id)!.status, "queued");
|
||||
|
||||
const blocked = await harness.manager.prompt(request("start next", "user-b"));
|
||||
assert.match(blocked.text, /还不能开始/);
|
||||
assert.match(blocked.text, /发起人确认/);
|
||||
assert.match(blocked.text, /保留在队列里/);
|
||||
assert.equal(harness.proposals.get(second.id)!.status, "queued");
|
||||
assert.equal(harness.proposals.get(first.id)!.status, "awaiting_user_confirmation");
|
||||
|
||||
const bPrompt = readLog(harness.logFile).find((entry) => entry.method === "session/prompt" && entry.text?.startsWith("[User message]\nstart next"));
|
||||
assert.ok(bPrompt?.text);
|
||||
assert.match(bPrompt.text!, /blocked: another proposal is awaiting owner confirmation/);
|
||||
assert.doesNotMatch(bPrompt.text!, /goal done/);
|
||||
|
||||
const bStatus = harness.manager.status("webhook", "chat-1", "user-b");
|
||||
assert.equal(bStatus.schedulerState, "awaiting_confirmation");
|
||||
assert.equal(bStatus.blockedReason, "another proposal is awaiting owner confirmation");
|
||||
assert.equal(bStatus.myQueuedProposals, 1);
|
||||
assert.equal(bStatus.myPendingConfirmations, 0);
|
||||
assert.equal(bStatus.nextAction, "wait for the current proposal to settle");
|
||||
|
||||
const aStatus = harness.manager.status("webhook", "chat-1", "user-a");
|
||||
assert.equal(aStatus.myPendingConfirmations, 1);
|
||||
assert.equal(aStatus.blockedReason, "your proposal is awaiting your confirmation");
|
||||
assert.equal(aStatus.nextAction, "confirm your pending proposal");
|
||||
} finally {
|
||||
await closeHarness(harness);
|
||||
}
|
||||
});
|
||||
|
||||
test("confirm and stop actions with no matching proposal append an ownership correction", async () => {
|
||||
const harness = await createHarness();
|
||||
try {
|
||||
await harness.manager.prompt(request("create proposal: hang", "user-a"));
|
||||
|
||||
const confirmed = await harness.manager.prompt(request("confirm", "user-b"));
|
||||
assert.match(confirmed.text, /没有可确认的提案/);
|
||||
assert.match(confirmed.text, /自己发起/);
|
||||
|
||||
const stopped = await harness.manager.prompt(request("stop", "user-b"));
|
||||
assert.match(stopped.text, /没有可停止的任务/);
|
||||
assert.match(stopped.text, /自己发起/);
|
||||
|
||||
const cancelled = await harness.manager.prompt(request("cancel proposal", "user-b"));
|
||||
assert.match(cancelled.text, /没有可取消的提案/);
|
||||
assert.match(cancelled.text, /自己发起/);
|
||||
|
||||
assert.equal(harness.proposals.list()[0]!.status, "proposed");
|
||||
} finally {
|
||||
await closeHarness(harness);
|
||||
}
|
||||
});
|
||||
|
||||
test("start_next reports when the queue is empty instead of claiming a start", async () => {
|
||||
const harness = await createHarness();
|
||||
try {
|
||||
const reply = await harness.manager.prompt(request("start next", "user-a"));
|
||||
assert.match(reply.text, /还不能开始/);
|
||||
assert.match(reply.text, /没有已确认并排队/);
|
||||
const status = harness.manager.status("webhook", "chat-1", "user-a");
|
||||
assert.equal(status.schedulerState, "idle");
|
||||
assert.equal(status.blockedReason, "none");
|
||||
assert.equal(status.nextAction, "none");
|
||||
} finally {
|
||||
await closeHarness(harness);
|
||||
}
|
||||
});
|
||||
|
||||
function processGroupExists(pgid: number): boolean {
|
||||
try { process.kill(-pgid, 0); return true; }
|
||||
catch (error) { return (error as NodeJS.ErrnoException).code === "EPERM"; }
|
||||
|
||||
+208
-8
@@ -68,6 +68,9 @@ const policy = { allowedUsers: [] as string[], allowedChats: [] as string[], req
|
||||
const message = (text: string, messageId = text, chatId = "chat"): IncomingMessage => ({
|
||||
platform: "qq", chatId, userId: "user", text, messageId
|
||||
});
|
||||
const groupMessage = (text: string, messageId: string): IncomingMessage => ({
|
||||
platform: "qq", chatId: "group:g1", userId: "user", text, messageId, isGroup: true
|
||||
});
|
||||
const flush = async () => { await Promise.resolve(); await Promise.resolve(); await new Promise<void>((resolve) => setImmediate(resolve)); };
|
||||
const waitFor = async (condition: () => boolean) => {
|
||||
for (let count = 0; count < 50 && !condition(); count++) await flush();
|
||||
@@ -78,6 +81,25 @@ function recordingAdapter(sent: OutgoingMessage[] = []): PlatformAdapter {
|
||||
return { name: "test", async handleWebhook() { return {}; }, async sendMessage(outgoing) { sent.push(outgoing); } };
|
||||
}
|
||||
|
||||
// Simulates QQ passive-reply dedup: a repeated msg_id + msg_seq pair is rejected.
|
||||
function qqDedupAdapter(sent: OutgoingMessage[] = [], fail?: (outgoing: OutgoingMessage) => string | undefined): PlatformAdapter {
|
||||
const used = new Set<string>();
|
||||
return {
|
||||
name: "test",
|
||||
async handleWebhook() { return {}; },
|
||||
async sendMessage(outgoing) {
|
||||
const failure = fail?.(outgoing);
|
||||
if (failure) throw new Error(failure);
|
||||
if (outgoing.replyTo) {
|
||||
const pair = `${outgoing.replyTo}|${outgoing.replySequence}`;
|
||||
if (used.has(pair)) throw new Error("QQ send failed: HTTP 400 code=400304018 duplicate msg_id+msg_seq");
|
||||
used.add(pair);
|
||||
}
|
||||
sent.push(outgoing);
|
||||
}
|
||||
};
|
||||
}
|
||||
|
||||
function createGateway(runtime = new FakeRuntime()) {
|
||||
return { runtime, gateway: new Gateway(policy, runtime) };
|
||||
}
|
||||
@@ -159,7 +181,7 @@ test("runtime errors remain runtime errors and are delivered through the same st
|
||||
assert.deepEqual(sent.map((outgoing) => [outgoing.text, outgoing.replySequence]), [["Agent error: failed", 1]]);
|
||||
});
|
||||
|
||||
test("sendEvent pushes a fresh message to the latest chat target without replyTo", async () => {
|
||||
test("sendEvent uses a fresh passive window with replyTo and an incrementing sequence", async () => {
|
||||
const { runtime, gateway } = createGateway();
|
||||
const sent: OutgoingMessage[] = [];
|
||||
const turn = gateway.receive(message("work", "event-source"), recordingAdapter(sent));
|
||||
@@ -168,11 +190,14 @@ test("sendEvent pushes a fresh message to the latest chat target without replyTo
|
||||
await turn;
|
||||
|
||||
await gateway.sendEvent("qq:chat", "任务完成了");
|
||||
const event = sent.at(-1)!;
|
||||
assert.equal(event.text, "任务完成了");
|
||||
assert.equal(event.replyTo, undefined);
|
||||
assert.equal(event.replySequence, undefined);
|
||||
assert.deepEqual(event.target, { platform: "qq", chatId: "chat", userId: "user", raw: undefined });
|
||||
await gateway.sendEvent("qq:chat", "又完成了一步");
|
||||
// The normal reply already consumed msg_seq 1 for "event-source"; events share the same counter.
|
||||
const events = sent.filter((outgoing) => outgoing.text !== "done");
|
||||
assert.deepEqual(events.map((outgoing) => [outgoing.text, outgoing.replyTo, outgoing.replySequence]), [
|
||||
["任务完成了", "event-source", 2],
|
||||
["又完成了一步", "event-source", 3]
|
||||
]);
|
||||
assert.deepEqual(events[0]!.target, { platform: "qq", chatId: "chat", userId: "user", raw: undefined });
|
||||
});
|
||||
|
||||
test("sendEvent uses the most recent target and adapter for the chat", async () => {
|
||||
@@ -190,7 +215,8 @@ test("sendEvent uses the most recent target and adapter for the chat", async ()
|
||||
await gateway.sendEvent("qq:chat", "event");
|
||||
assert.equal(firstSent.filter((outgoing) => outgoing.text === "event").length, 0);
|
||||
assert.equal(secondSent.at(-1)?.text, "event");
|
||||
assert.equal(secondSent.at(-1)?.replyTo, undefined);
|
||||
assert.equal(secondSent.at(-1)?.replyTo, "two");
|
||||
assert.equal(secondSent.at(-1)?.replySequence, 2); // "two done" consumed seq 1 for message "two"
|
||||
});
|
||||
|
||||
test("sendEvent drops safely when the chat has no known target", async () => {
|
||||
@@ -198,13 +224,187 @@ test("sendEvent drops safely when the chat has no known target", async () => {
|
||||
await gateway.sendEvent("qq:unknown", "event");
|
||||
});
|
||||
|
||||
test("expired group passive window skips the event and reminds on the next inbound", async () => {
|
||||
let now = 1_000_000;
|
||||
const runtime = new FakeRuntime();
|
||||
const gateway = new Gateway(policy, runtime, { now: () => now });
|
||||
const sent: OutgoingMessage[] = [];
|
||||
const adapter = recordingAdapter(sent);
|
||||
|
||||
const turn = gateway.receive(groupMessage("work", "m1"), adapter);
|
||||
await waitFor(() => runtime.turns.length === 1);
|
||||
runtime.turns[0].resolve("done");
|
||||
await turn;
|
||||
|
||||
now += 4 * 60_000; // within the 4.5 minute group window
|
||||
await gateway.sendEvent("qq:group:g1", "fresh event");
|
||||
assert.deepEqual(sent.at(-1), {
|
||||
target: { platform: "qq", chatId: "group:g1", userId: "user", raw: undefined },
|
||||
text: "fresh event", replyTo: "m1", replySequence: 2 // "done" consumed seq 1 for m1
|
||||
});
|
||||
|
||||
now += 2 * 60_000; // past the group window
|
||||
await gateway.sendEvent("qq:group:g1", "stale event");
|
||||
assert.equal(sent.filter((outgoing) => outgoing.text === "stale event").length, 0);
|
||||
|
||||
const status = await gateway.receive(groupMessage("/status", "m2"), adapter, { synchronous: true });
|
||||
assert.match(status.reply || "", /lastEventDelivery=skipped/);
|
||||
assert.match(status.reply || "", /lastEventError=no fresh passive window/);
|
||||
assert.match(status.reply || "", /lastEventAt=\d{4}-\d{2}-\d{2}T\d{2}:\d{2}:\d{2}/);
|
||||
|
||||
const reminder = gateway.receive(groupMessage("hello", "m3"), adapter);
|
||||
await waitFor(() => runtime.turns.length === 2);
|
||||
runtime.turns[1].resolve("ok");
|
||||
await reminder;
|
||||
const reminded = sent.find((outgoing) => outgoing.text === "stale event");
|
||||
assert.equal(reminded?.replyTo, "m3");
|
||||
assert.equal(reminded?.replySequence, 1);
|
||||
|
||||
// The stale event was already flushed once; a later inbound must not repeat it.
|
||||
const again = gateway.receive(groupMessage("hello again", "m4"), adapter);
|
||||
await waitFor(() => runtime.turns.length === 3);
|
||||
runtime.turns[2].resolve("ok");
|
||||
await again;
|
||||
assert.equal(sent.filter((outgoing) => outgoing.text === "stale event").length, 1);
|
||||
});
|
||||
|
||||
test("status reports the last successful event delivery for the current chat", async () => {
|
||||
const { runtime, gateway } = createGateway();
|
||||
const adapter = recordingAdapter();
|
||||
const turn = gateway.receive(message("work", "m1"), adapter);
|
||||
await waitFor(() => runtime.turns.length === 1);
|
||||
runtime.turns[0].resolve("done");
|
||||
await turn;
|
||||
|
||||
await gateway.sendEvent("qq:chat", "event");
|
||||
const status = await gateway.receive(message("/status", "m2"), adapter, { synchronous: true });
|
||||
assert.match(status.reply || "", /lastEventDelivery=success/);
|
||||
assert.match(status.reply || "", /lastEventError=none/);
|
||||
assert.match(status.reply || "", /lastEventAt=\d{4}-\d{2}-\d{2}T/);
|
||||
});
|
||||
|
||||
test("normal replies and fresh events share one msg_seq counter per inbound message", async () => {
|
||||
const runtime = new FakeRuntime();
|
||||
const gateway = new Gateway(policy, runtime);
|
||||
const sent: OutgoingMessage[] = [];
|
||||
const turn = gateway.receive(groupMessage("work", "m1"), qqDedupAdapter(sent));
|
||||
await waitFor(() => runtime.turns.length === 1);
|
||||
runtime.turns[0].resolve("done");
|
||||
await turn;
|
||||
|
||||
// With independent counters these events would reuse (m1, seq 1) and QQ would reject them.
|
||||
await gateway.sendEvent("qq:group:g1", "event one");
|
||||
await gateway.sendEvent("qq:group:g1", "event two");
|
||||
assert.deepEqual(sent.map((outgoing) => [outgoing.text, outgoing.replyTo, outgoing.replySequence]), [
|
||||
["done", "m1", 1],
|
||||
["event one", "m1", 2],
|
||||
["event two", "m1", 3]
|
||||
]);
|
||||
});
|
||||
|
||||
test("a redelivered inbound message id continues its msg_seq counter instead of restarting", async () => {
|
||||
const runtime = new FakeRuntime();
|
||||
const gateway = new Gateway(policy, runtime);
|
||||
const sent: OutgoingMessage[] = [];
|
||||
const adapter = qqDedupAdapter(sent);
|
||||
|
||||
const first = gateway.receive(groupMessage("work", "m1"), adapter);
|
||||
await waitFor(() => runtime.turns.length === 1);
|
||||
runtime.turns[0].resolve("first done");
|
||||
await first;
|
||||
|
||||
// QQ pushed the same msg_id again: the reply must not reuse (m1, seq 1).
|
||||
const second = gateway.receive(groupMessage("work", "m1"), adapter);
|
||||
await waitFor(() => runtime.turns.length === 2);
|
||||
runtime.turns[1].resolve("second done");
|
||||
await second;
|
||||
|
||||
assert.deepEqual(sent.map((outgoing) => [outgoing.text, outgoing.replyTo, outgoing.replySequence]), [
|
||||
["first done", "m1", 1],
|
||||
["second done", "m1", 2]
|
||||
]);
|
||||
});
|
||||
|
||||
test("flushed pending events and the current reply share the msg_seq counter", async () => {
|
||||
let now = 1_000_000;
|
||||
const runtime = new FakeRuntime();
|
||||
const gateway = new Gateway(policy, runtime, { now: () => now });
|
||||
const sent: OutgoingMessage[] = [];
|
||||
const adapter = qqDedupAdapter(sent);
|
||||
|
||||
const first = gateway.receive(groupMessage("work", "m1"), adapter);
|
||||
await waitFor(() => runtime.turns.length === 1);
|
||||
runtime.turns[0].resolve("done");
|
||||
await first;
|
||||
|
||||
now += 6 * 60_000; // past the group window: the event is queued instead of sent
|
||||
await gateway.sendEvent("qq:group:g1", "stale event");
|
||||
assert.equal(sent.filter((outgoing) => outgoing.text === "stale event").length, 0);
|
||||
|
||||
const second = gateway.receive(groupMessage("hello", "m2"), adapter);
|
||||
await waitFor(() => runtime.turns.length === 2);
|
||||
runtime.turns[1].resolve("ok");
|
||||
await second;
|
||||
|
||||
// The redelivery takes (m2, seq 1) and the current turn's reply takes (m2, seq 2); no pair collides.
|
||||
assert.deepEqual(sent.map((outgoing) => [outgoing.text, outgoing.replyTo, outgoing.replySequence]), [
|
||||
["done", "m1", 1],
|
||||
["stale event", "m2", 1],
|
||||
["ok", "m2", 2]
|
||||
]);
|
||||
});
|
||||
|
||||
test("failed pending event redelivery stays queued, records failed, and retries on the next inbound", async () => {
|
||||
let now = 1_000_000;
|
||||
const runtime = new FakeRuntime();
|
||||
const gateway = new Gateway(policy, runtime, { now: () => now });
|
||||
const sent: OutgoingMessage[] = [];
|
||||
let failEvent = true;
|
||||
const adapter = qqDedupAdapter(sent, (outgoing) => failEvent && outgoing.text === "stale event"
|
||||
? "QQ send failed: HTTP 400 code=40034105 proactive message not allowed"
|
||||
: undefined);
|
||||
|
||||
const first = gateway.receive(groupMessage("work", "m1"), adapter);
|
||||
await waitFor(() => runtime.turns.length === 1);
|
||||
runtime.turns[0].resolve("done");
|
||||
await first;
|
||||
|
||||
now += 6 * 60_000;
|
||||
await gateway.sendEvent("qq:group:g1", "stale event"); // skipped, queued
|
||||
|
||||
const second = gateway.receive(groupMessage("hello", "m2"), adapter);
|
||||
await waitFor(() => runtime.turns.length === 2);
|
||||
runtime.turns[1].resolve("ok");
|
||||
await second;
|
||||
// The redelivery failed: nothing was delivered, but the normal reply still went out.
|
||||
assert.equal(sent.filter((outgoing) => outgoing.text === "stale event").length, 0);
|
||||
assert.deepEqual(sent.map((outgoing) => [outgoing.text, outgoing.replyTo, outgoing.replySequence]), [
|
||||
["done", "m1", 1],
|
||||
["ok", "m2", 2]
|
||||
]);
|
||||
|
||||
const status = await gateway.receive(groupMessage("/status", "ms"), adapter, { synchronous: true });
|
||||
assert.match(status.reply || "", /lastEventDelivery=failed/);
|
||||
assert.match(status.reply || "", /lastEventError=HTTP 400 QQ code 40034105/);
|
||||
|
||||
failEvent = false;
|
||||
const third = gateway.receive(groupMessage("hello again", "m3"), adapter);
|
||||
await waitFor(() => runtime.turns.length === 3);
|
||||
runtime.turns[2].resolve("ok again");
|
||||
await third;
|
||||
const redelivered = sent.filter((outgoing) => outgoing.text === "stale event");
|
||||
assert.equal(redelivered.length, 1);
|
||||
assert.equal(redelivered[0]!.replyTo, "m3");
|
||||
assert.equal(redelivered[0]!.replySequence, 1);
|
||||
});
|
||||
|
||||
test("sendEvent delivery failures are logged safely without message content", async () => {
|
||||
const { runtime, gateway } = createGateway();
|
||||
const errors: string[] = [];
|
||||
const adapter: PlatformAdapter = {
|
||||
name: "test", async handleWebhook() { return {}; },
|
||||
async sendMessage(outgoing) {
|
||||
if (!outgoing.replyTo) throw new Error("sensitive provider response");
|
||||
if (outgoing.text.includes("event text")) throw new Error("sensitive provider response");
|
||||
}
|
||||
};
|
||||
const turn = gateway.receive(message("work", "event-failure"), adapter);
|
||||
|
||||
@@ -51,3 +51,40 @@ test("sendMessage uses C2C endpoint and forwards reply sequences as msg_seq", as
|
||||
]);
|
||||
} finally { globalThis.fetch = original; }
|
||||
});
|
||||
|
||||
test("sendMessage without replyTo omits msg_id and msg_seq from the body", async () => {
|
||||
const original = globalThis.fetch; const bodies: Array<Record<string, unknown>> = [];
|
||||
globalThis.fetch = (async (input: string | URL | Request, init?: RequestInit) => {
|
||||
if (String(input).includes("getAppAccessToken")) return new Response(JSON.stringify({ access_token: "token", expires_in: 7200 }), { status: 200 });
|
||||
bodies.push(JSON.parse(String(init?.body)) as Record<string, unknown>);
|
||||
return new Response(JSON.stringify({ id: "sent" }), { status: 200 });
|
||||
}) as typeof fetch;
|
||||
try {
|
||||
const adapter = new QqAdapter(config, { receive: async () => ({ ok: true }) } as never);
|
||||
const target = { platform: "qq", chatId: "group:g1", raw: { group_openid: "g1" } };
|
||||
await adapter.sendMessage({ target, text: "proactive" });
|
||||
assert.deepEqual(bodies, [{ content: "proactive" }]);
|
||||
} finally { globalThis.fetch = original; }
|
||||
});
|
||||
|
||||
test("sendMessage surfaces safe QQ error details (HTTP status and code) without the request body", async () => {
|
||||
const original = globalThis.fetch;
|
||||
globalThis.fetch = (async (input: string | URL | Request) => {
|
||||
if (String(input).includes("getAppAccessToken")) return new Response(JSON.stringify({ access_token: "token", expires_in: 7200 }), { status: 200 });
|
||||
return new Response(JSON.stringify({ code: 40034105, message: "proactive message not allowed" }), { status: 400 });
|
||||
}) as typeof fetch;
|
||||
try {
|
||||
const adapter = new QqAdapter(config, { receive: async () => ({ ok: true }) } as never);
|
||||
const target = { platform: "qq", chatId: "group:g1", raw: { group_openid: "g1" } };
|
||||
await assert.rejects(
|
||||
adapter.sendMessage({ target, text: "secret message content" }),
|
||||
(error: Error) => {
|
||||
assert.match(error.message, /HTTP 400/);
|
||||
assert.match(error.message, /code=40034105/);
|
||||
assert.match(error.message, /proactive message not allowed/);
|
||||
assert.doesNotMatch(error.message, /secret message content/);
|
||||
return true;
|
||||
}
|
||||
);
|
||||
} finally { globalThis.fetch = original; }
|
||||
});
|
||||
|
||||
Reference in New Issue
Block a user