refine long-running task notices
This commit is contained in:
@@ -89,7 +89,7 @@ gori-agent list
|
|||||||
- `src/core/gateway.ts`
|
- `src/core/gateway.ts`
|
||||||
- 入站 allowlist、群 mention、命令、per-chat lock、队列状态和回复。
|
- 入站 allowlist、群 mention、命令、per-chat lock、队列状态和回复。
|
||||||
- policy 后 `/status`、`/cancel`、`/help` 必须绕过 chat lock,`/new` 必须保持串行。
|
- policy 后 `/status`、`/cancel`、`/help` 必须绕过 chat lock,`/new` 必须保持串行。
|
||||||
- 异步平台普通消息轮到后立即 prompt,不等待提示发送;15 秒内完成只发结果,否则按 queued/running 提示,queued 转 running 时补开始提示,并在入站后 60 秒、180 秒及之后每 300 秒结合 ACP `idleSeconds` 低频提醒;同步 webhook 不提示。
|
- 异步平台普通消息轮到后立即 prompt,不等待提示发送;15 秒内完成只发结果,否则按 queued/running 提示,queued 转 running 时补开始提示;主动提示从入站绝对时间按 15 秒、60 秒、180 秒、480 秒触发,之后每 600 秒一次,只使用用户可感知文案并隐藏 ACP、`phase`、`idleSeconds` 等内部指标,技术状态只保留在 `/status`;同步 webhook 不提示。
|
||||||
- 每条异步入站使用独立串行 ReplyStream,`replySequence` 从 1 动态递增;发送失败只写安全日志且不阻止任务或后续发送,runtime/delivery 错误分离,最终消息入流后释放 chat lock。
|
- 每条异步入站使用独立串行 ReplyStream,`replySequence` 从 1 动态递增;发送失败只写安全日志且不阻止任务或后续发送,runtime/delivery 错误分离,最终消息入流后释放 chat lock。
|
||||||
- 每条入站只维护一个基于绝对 deadline 的可取消 timer;完成先标记 finished 并清 timer,shutdown 清理全部 timer,生产 timer 必须 `unref`,Gateway 可注入 fake clock/schedule 供测试。
|
- 每条入站只维护一个基于绝对 deadline 的可取消 timer;完成先标记 finished 并清 timer,shutdown 清理全部 timer,生产 timer 必须 `unref`,Gateway 可注入 fake clock/schedule 供测试。
|
||||||
- `src/core/command-router.ts`
|
- `src/core/command-router.ts`
|
||||||
|
|||||||
@@ -212,7 +212,7 @@ Bot fingerprint 包含 Bot ID、workspace、persona、agent、permissions、skil
|
|||||||
- `/new` 保持串行,清除当前 chat binding;下一条普通消息创建新 native session。
|
- `/new` 保持串行,清除当前 chat binding;下一条普通消息创建新 native session。
|
||||||
- `/roles`、`/role`、`/agents`、`/agent` 会返回固定 retired 提示,不会转发给 ACP。
|
- `/roles`、`/role`、`/agents`、`/agent` 会返回固定 retired 提示,不会转发给 ACP。
|
||||||
|
|
||||||
同一 chat 的普通消息和 `/new` 串行执行,不同 chat 可并发。异步平台的普通消息轮到后立即开始 ACP prompt,不等待状态提示发送:15 秒内完成只发送真实结果;15 秒仍未完成时,running 发送“这条还在处理,完成后我会直接回复。”,queued 发送“前面还有任务,这条还在排队,轮到后我马上处理。”;发送过 queued 提示的消息真正开始时再发送“轮到这条了,我开始处理。”。入站后 60 秒、180 秒及之后每 300 秒按 queued/running 和 ACP `idleSeconds` 发送低频口语化提醒;延迟执行的 timer 不补发历史提醒。同步 webhook 不发送这些提示。
|
同一 chat 的普通消息和 `/new` 串行执行,不同 chat 可并发。异步平台的普通消息轮到后立即开始 ACP prompt,不等待状态提示发送:15 秒内完成只发送真实结果;15 秒仍未完成时按 running/queued 发送口语化提示,发送过 queued 提示的消息真正开始时再补充开始提示。主动提示从入站绝对时间按 15 秒、60 秒、180 秒、480 秒触发,480 秒后每 600 秒一次(18、28、38 分钟……);它们只描述用户可感知的处理或等待状态,不暴露 ACP、`phase`、`idleSeconds` 等内部技术指标,技术状态只保留在 `/status`。延迟执行的 timer 不补发历史提醒;同步 webhook 不发送这些提示。
|
||||||
|
|
||||||
每条异步入站有独立、串行的回复流,`replySequence` 从 1 动态递增;因此短任务最终回复使用 `msg_seq=1`,已发送提示时后续消息使用下一序号。平台发送失败只记录不含消息正文或 provider 错误详情的安全日志,不影响 ACP 任务或回复流中的后续发送,也不会把成功的 ACP 结果误报为 `Agent error`。最终结果加入回复流后即释放 per-chat lock,平台发送延迟不会阻塞下一项任务;完成与 shutdown 都会清理状态提醒 timer。
|
每条异步入站有独立、串行的回复流,`replySequence` 从 1 动态递增;因此短任务最终回复使用 `msg_seq=1`,已发送提示时后续消息使用下一序号。平台发送失败只记录不含消息正文或 provider 错误详情的安全日志,不影响 ACP 任务或回复流中的后续发送,也不会把成功的 ACP 结果误报为 `Agent error`。最终结果加入回复流后即释放 per-chat lock,平台发送延迟不会阻塞下一项任务;完成与 shutdown 都会清理状态提醒 timer。
|
||||||
|
|
||||||
|
|||||||
+19
-35
@@ -104,11 +104,11 @@ export class Gateway {
|
|||||||
stream: new ReplyStream(message, adapter)
|
stream: new ReplyStream(message, adapter)
|
||||||
};
|
};
|
||||||
this.asyncInbounds.add(inbound);
|
this.asyncInbounds.add(inbound);
|
||||||
this.armTimer(inbound, inbound.receivedAt + 15_000, message);
|
this.armTimer(inbound, inbound.receivedAt + 15_000);
|
||||||
|
|
||||||
const locked = await this.withChatLock(chatKey, async () => {
|
const locked = await this.withChatLock(chatKey, async () => {
|
||||||
inbound.phase = "running";
|
inbound.phase = "running";
|
||||||
if (inbound.queuedNoticeSent) void inbound.stream.enqueue("轮到这条了,我开始处理。");
|
if (inbound.queuedNoticeSent) void inbound.stream.enqueue("轮到这件事了,我现在开始处理。");
|
||||||
try {
|
try {
|
||||||
const reply = (await this.runtime.prompt({
|
const reply = (await this.runtime.prompt({
|
||||||
platform: message.platform,
|
platform: message.platform,
|
||||||
@@ -200,7 +200,7 @@ export class Gateway {
|
|||||||
return await this.runtime.cancel(message.platform, message.chatId) ? "Cancellation requested." : "No active turn to cancel.";
|
return await this.runtime.cancel(message.platform, message.chatId) ? "Cancellation requested." : "No active turn to cancel.";
|
||||||
}
|
}
|
||||||
|
|
||||||
private armTimer(inbound: AsyncInbound, deadline: number, message: IncomingMessage): void {
|
private armTimer(inbound: AsyncInbound, deadline: number): void {
|
||||||
if (this.closed || inbound.phase === "finished") return;
|
if (this.closed || inbound.phase === "finished") return;
|
||||||
inbound.timer?.cancel();
|
inbound.timer?.cancel();
|
||||||
inbound.timer = this.schedule(() => {
|
inbound.timer = this.schedule(() => {
|
||||||
@@ -208,41 +208,19 @@ export class Gateway {
|
|||||||
if (this.closed || inbound.phase === "finished") return;
|
if (this.closed || inbound.phase === "finished") return;
|
||||||
const now = this.now();
|
const now = this.now();
|
||||||
if (now < deadline) {
|
if (now < deadline) {
|
||||||
this.armTimer(inbound, deadline, message);
|
this.armTimer(inbound, deadline);
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
|
const elapsed = now - inbound.receivedAt;
|
||||||
if (!inbound.initialNoticeSent) {
|
if (!inbound.initialNoticeSent) {
|
||||||
inbound.initialNoticeSent = true;
|
inbound.initialNoticeSent = true;
|
||||||
inbound.queuedNoticeSent = inbound.phase === "queued";
|
inbound.queuedNoticeSent = inbound.phase === "queued";
|
||||||
void inbound.stream.enqueue(inbound.phase === "queued"
|
|
||||||
? "前面还有任务,这条还在排队,轮到后我马上处理。"
|
|
||||||
: "这条还在处理,完成后我会直接回复。");
|
|
||||||
} else {
|
|
||||||
void inbound.stream.enqueue(this.reminderText(inbound, message));
|
|
||||||
}
|
}
|
||||||
this.armTimer(inbound, nextReminderDeadline(inbound.receivedAt, now), message);
|
void inbound.stream.enqueue(reminderText(inbound.phase, elapsed));
|
||||||
|
this.armTimer(inbound, nextReminderDeadline(inbound.receivedAt, now));
|
||||||
}, Math.max(0, deadline - this.now()));
|
}, Math.max(0, deadline - this.now()));
|
||||||
}
|
}
|
||||||
|
|
||||||
private reminderText(inbound: AsyncInbound, message: IncomingMessage): string {
|
|
||||||
let idleSeconds: number | undefined;
|
|
||||||
try {
|
|
||||||
const value = this.runtime.status(message.platform, message.chatId).idleSeconds;
|
|
||||||
if (typeof value === "number" && Number.isFinite(value) && value >= 0) idleSeconds = Math.floor(value);
|
|
||||||
} catch {
|
|
||||||
// A status sampling failure must not affect the active turn.
|
|
||||||
}
|
|
||||||
const activity = activityText(idleSeconds);
|
|
||||||
if (inbound.phase === "queued") {
|
|
||||||
return activity
|
|
||||||
? `前面的任务还在处理,${activity};这条仍在排队,轮到后我马上处理。`
|
|
||||||
: "前面的任务还在处理,这条仍在排队,轮到后我马上处理。";
|
|
||||||
}
|
|
||||||
return activity
|
|
||||||
? `这条还在处理中,${activity};完成后我会直接回复。`
|
|
||||||
: "这条还在处理中,完成后我会直接回复。";
|
|
||||||
}
|
|
||||||
|
|
||||||
private finishInbound(inbound: AsyncInbound): void {
|
private finishInbound(inbound: AsyncInbound): void {
|
||||||
if (inbound.phase === "finished") return;
|
if (inbound.phase === "finished") return;
|
||||||
inbound.phase = "finished";
|
inbound.phase = "finished";
|
||||||
@@ -295,12 +273,18 @@ function nextReminderDeadline(receivedAt: number, now: number): number {
|
|||||||
const elapsed = now - receivedAt;
|
const elapsed = now - receivedAt;
|
||||||
if (elapsed < 60_000) return receivedAt + 60_000;
|
if (elapsed < 60_000) return receivedAt + 60_000;
|
||||||
if (elapsed < 180_000) return receivedAt + 180_000;
|
if (elapsed < 180_000) return receivedAt + 180_000;
|
||||||
return receivedAt + 180_000 + (Math.floor((elapsed - 180_000) / 300_000) + 1) * 300_000;
|
if (elapsed < 480_000) return receivedAt + 480_000;
|
||||||
|
return receivedAt + 480_000 + (Math.floor((elapsed - 480_000) / 600_000) + 1) * 600_000;
|
||||||
}
|
}
|
||||||
|
|
||||||
function activityText(idleSeconds: number | undefined): string | undefined {
|
function reminderText(phase: AsyncInbound["phase"], elapsed: number): string {
|
||||||
if (idleSeconds === undefined) return undefined;
|
if (phase === "queued") {
|
||||||
if (idleSeconds < 30) return "刚刚还有进展";
|
return elapsed < 60_000
|
||||||
if (idleSeconds < 120) return `最近 ${idleSeconds} 秒没新进展`;
|
? "前面还有一件事没处理完,这条我记着,轮到后马上处理。"
|
||||||
return `最近 ${Math.floor(idleSeconds / 60)} 分钟没新进展`;
|
: "前面的事情还没处理完,这条还在等。我没有漏掉,轮到后会接着做。";
|
||||||
|
}
|
||||||
|
if (elapsed < 60_000) return "收到,我还在处理,稍等我一下。";
|
||||||
|
if (elapsed < 180_000) return "还在弄,暂时还没出结果,弄好我马上回你。";
|
||||||
|
if (elapsed < 480_000) return "这件事比预想中多花了一点时间,我还在继续处理。你不用一直盯着,弄好我会直接回你。";
|
||||||
|
return "我还在处理这件事,确实花了些时间。想看看现在的状态可以发 /status,不想继续等也可以发 /cancel。";
|
||||||
}
|
}
|
||||||
|
|||||||
+23
-19
@@ -148,7 +148,7 @@ test("completion and timer callbacks are safe in both orders at 15 seconds", asy
|
|||||||
runtime.turns[0].resolve("done");
|
runtime.turns[0].resolve("done");
|
||||||
await turn;
|
await turn;
|
||||||
assert.deepEqual(sent.map((outgoing) => [outgoing.text, outgoing.replySequence]), [
|
assert.deepEqual(sent.map((outgoing) => [outgoing.text, outgoing.replySequence]), [
|
||||||
["这条还在处理,完成后我会直接回复。", 1],
|
["收到,我还在处理,稍等我一下。", 1],
|
||||||
["done", 2]
|
["done", 2]
|
||||||
]);
|
]);
|
||||||
}
|
}
|
||||||
@@ -162,14 +162,14 @@ test("running and queued notices are independent and queued work announces start
|
|||||||
const second = gateway.receive(message("two", "second"), adapter);
|
const second = gateway.receive(message("two", "second"), adapter);
|
||||||
await waitFor(() => runtime.turns.length === 1);
|
await waitFor(() => runtime.turns.length === 1);
|
||||||
await clock.advanceTo(15_000);
|
await clock.advanceTo(15_000);
|
||||||
assert.deepEqual(messagesFor(sent, "first").map((outgoing) => outgoing.text), ["这条还在处理,完成后我会直接回复。"]);
|
assert.deepEqual(messagesFor(sent, "first").map((outgoing) => outgoing.text), ["收到,我还在处理,稍等我一下。"]);
|
||||||
assert.deepEqual(messagesFor(sent, "second").map((outgoing) => outgoing.text), ["前面还有任务,这条还在排队,轮到后我马上处理。"]);
|
assert.deepEqual(messagesFor(sent, "second").map((outgoing) => outgoing.text), ["前面还有一件事没处理完,这条我记着,轮到后马上处理。"]);
|
||||||
|
|
||||||
runtime.turns[0].resolve("first done");
|
runtime.turns[0].resolve("first done");
|
||||||
await waitFor(() => runtime.turns.length === 2);
|
await waitFor(() => runtime.turns.length === 2);
|
||||||
assert.deepEqual(messagesFor(sent, "second").map((outgoing) => [outgoing.text, outgoing.replySequence]), [
|
assert.deepEqual(messagesFor(sent, "second").map((outgoing) => [outgoing.text, outgoing.replySequence]), [
|
||||||
["前面还有任务,这条还在排队,轮到后我马上处理。", 1],
|
["前面还有一件事没处理完,这条我记着,轮到后马上处理。", 1],
|
||||||
["轮到这条了,我开始处理。", 2]
|
["轮到这件事了,我现在开始处理。", 2]
|
||||||
]);
|
]);
|
||||||
runtime.turns[1].resolve("second done");
|
runtime.turns[1].resolve("second done");
|
||||||
await Promise.all([first, second]);
|
await Promise.all([first, second]);
|
||||||
@@ -177,29 +177,32 @@ test("running and queued notices are independent and queued work announces start
|
|||||||
assert.deepEqual(messagesFor(sent, "second").map((outgoing) => outgoing.replySequence), [1, 2, 3]);
|
assert.deepEqual(messagesFor(sent, "second").map((outgoing) => outgoing.replySequence), [1, 2, 3]);
|
||||||
});
|
});
|
||||||
|
|
||||||
test("running reminders fire at 60, 180, and 480 seconds with dynamic activity text", async () => {
|
test("running reminders use staged wording and wait from 8 to 18 minutes", async () => {
|
||||||
const { runtime, clock, gateway } = createGateway();
|
const { runtime, clock, gateway } = createGateway();
|
||||||
const sent: OutgoingMessage[] = [];
|
const sent: OutgoingMessage[] = [];
|
||||||
const turn = gateway.receive(message("work", "reminders"), recordingAdapter(sent));
|
const turn = gateway.receive(message("work", "reminders"), recordingAdapter(sent));
|
||||||
await waitFor(() => runtime.turns.length === 1);
|
await waitFor(() => runtime.turns.length === 1);
|
||||||
await clock.advanceTo(15_000);
|
await clock.advanceTo(15_000);
|
||||||
runtime.idleSeconds = 10;
|
|
||||||
await clock.advanceTo(60_000);
|
await clock.advanceTo(60_000);
|
||||||
runtime.idleSeconds = 60;
|
|
||||||
await clock.advanceTo(180_000);
|
await clock.advanceTo(180_000);
|
||||||
runtime.idleSeconds = 180;
|
|
||||||
await clock.advanceTo(480_000);
|
await clock.advanceTo(480_000);
|
||||||
assert.deepEqual(sent.map((outgoing) => outgoing.text), [
|
assert.deepEqual(sent.map((outgoing) => outgoing.text), [
|
||||||
"这条还在处理,完成后我会直接回复。",
|
"收到,我还在处理,稍等我一下。",
|
||||||
"这条还在处理中,刚刚还有进展;完成后我会直接回复。",
|
"还在弄,暂时还没出结果,弄好我马上回你。",
|
||||||
"这条还在处理中,最近 60 秒没新进展;完成后我会直接回复。",
|
"这件事比预想中多花了一点时间,我还在继续处理。你不用一直盯着,弄好我会直接回你。",
|
||||||
"这条还在处理中,最近 3 分钟没新进展;完成后我会直接回复。"
|
"我还在处理这件事,确实花了些时间。想看看现在的状态可以发 /status,不想继续等也可以发 /cancel。"
|
||||||
]);
|
]);
|
||||||
|
await clock.advanceTo(780_000);
|
||||||
|
assert.equal(sent.length, 4);
|
||||||
|
await clock.advanceTo(1_080_000);
|
||||||
|
assert.equal(sent.at(-1)?.text,
|
||||||
|
"我还在处理这件事,确实花了些时间。想看看现在的状态可以发 /status,不想继续等也可以发 /cancel。");
|
||||||
|
assert.equal(sent.length, 5);
|
||||||
runtime.turns[0].resolve("done");
|
runtime.turns[0].resolve("done");
|
||||||
await turn;
|
await turn;
|
||||||
});
|
});
|
||||||
|
|
||||||
test("queued reminders use queue wording and sample the active ACP status", async () => {
|
test("queued reminders use staged queue wording without ACP activity details", async () => {
|
||||||
const { runtime, clock, gateway } = createGateway();
|
const { runtime, clock, gateway } = createGateway();
|
||||||
const sent: OutgoingMessage[] = [];
|
const sent: OutgoingMessage[] = [];
|
||||||
const adapter = recordingAdapter(sent);
|
const adapter = recordingAdapter(sent);
|
||||||
@@ -207,10 +210,11 @@ test("queued reminders use queue wording and sample the active ACP status", asyn
|
|||||||
const second = gateway.receive(message("two", "queued"), adapter);
|
const second = gateway.receive(message("two", "queued"), adapter);
|
||||||
await waitFor(() => runtime.turns.length === 1);
|
await waitFor(() => runtime.turns.length === 1);
|
||||||
await clock.advanceTo(15_000);
|
await clock.advanceTo(15_000);
|
||||||
runtime.idleSeconds = 45;
|
|
||||||
await clock.advanceTo(60_000);
|
await clock.advanceTo(60_000);
|
||||||
assert.equal(messagesFor(sent, "queued").at(-1)?.text,
|
assert.deepEqual(messagesFor(sent, "queued").map((outgoing) => outgoing.text), [
|
||||||
"前面的任务还在处理,最近 45 秒没新进展;这条仍在排队,轮到后我马上处理。");
|
"前面还有一件事没处理完,这条我记着,轮到后马上处理。",
|
||||||
|
"前面的事情还没处理完,这条还在等。我没有漏掉,轮到后会接着做。"
|
||||||
|
]);
|
||||||
runtime.turns[0].resolve("one done");
|
runtime.turns[0].resolve("one done");
|
||||||
await waitFor(() => runtime.turns.length === 2);
|
await waitFor(() => runtime.turns.length === 2);
|
||||||
runtime.turns[1].resolve("two done");
|
runtime.turns[1].resolve("two done");
|
||||||
@@ -275,7 +279,7 @@ test("prompt does not wait for a notice send and final delivery delay releases t
|
|||||||
runtime.turns[1].resolve("second done");
|
runtime.turns[1].resolve("second done");
|
||||||
await Promise.all([first, second]);
|
await Promise.all([first, second]);
|
||||||
assert.deepEqual(messagesFor(sent, "first").map((outgoing) => outgoing.text), [
|
assert.deepEqual(messagesFor(sent, "first").map((outgoing) => outgoing.text), [
|
||||||
"这条还在处理,完成后我会直接回复。", "first done"
|
"收到,我还在处理,稍等我一下。", "first done"
|
||||||
]);
|
]);
|
||||||
});
|
});
|
||||||
|
|
||||||
@@ -302,7 +306,7 @@ test("delivery failures are logged safely and do not block later sends", async (
|
|||||||
} finally {
|
} finally {
|
||||||
console.error = originalError;
|
console.error = originalError;
|
||||||
}
|
}
|
||||||
assert.deepEqual(sent.map((outgoing) => outgoing.text), ["这条还在处理,完成后我会直接回复。", "done"]);
|
assert.deepEqual(sent.map((outgoing) => outgoing.text), ["收到,我还在处理,稍等我一下。", "done"]);
|
||||||
assert.deepEqual(errors, ["Gateway delivery send failed (platform=qq)"]);
|
assert.deepEqual(errors, ["Gateway delivery send failed (platform=qq)"]);
|
||||||
});
|
});
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user