diff --git a/.changeset/amsg-server-before-fire-defer.md b/.changeset/amsg-server-before-fire-defer.md new file mode 100644 index 0000000..7e145a5 --- /dev/null +++ b/.changeset/amsg-server-before-fire-defer.md @@ -0,0 +1,16 @@ +--- +"@rei-standard/amsg-server": minor +--- + +onBeforeFire 新增 `{ defer: { afterMs } }` 出口:现在不合适,过一会儿再来问 + +任务到点时宿主这会儿不方便生成(比如同一段对话里正有另一轮回复在生成),`onBeforeFire` 返回 `{ defer: { afterMs } }` 即可把这次触发往后推。这一次不调 LLM、不推送,也不调 `onLLMOutput` / `onAfterSend`;任务行只把 `retry_after` 写成现在 + `afterMs` 并放掉租约,`retry_count`、`last_error`、`next_send_at`、`status` 都不动,循环任务不推进。到点之后 cron 把它重新捞起来,从 `onBeforeFire` 再走一遍,可以继续推迟。 + +它不是失败:抛错会占一格重试、写 `last_error`、按 2 / 4 / 6 分钟退避,推迟这三样都没有。 + +- `afterMs` 必须是有限正数,最大 24 小时(导出为 `MAX_DEFER_AFTER_MS`);不合法按 `AGENTIC_BAD_BEFORE_FIRE` 一跳终审。 +- `onFireSettled` 的 `status` 多一种 `deferred`,载荷新增 `retryAfter`(ISO 时刻,其余结局为 `null`)。 +- tick 汇总新增 `details.deferredTasks`(`{ taskId, retryAfter }`),不计入成功或失败;`runTask` 在到点之前再调会回 `retry_pending`。 +- 过期线照常生效:名义时刻过去超过 `staleAfterMs` 的任务按过期处理;某次推迟的唤醒时刻已经越过这条线时当场按过期收场,重试链上的任务也一样。 +- 只有定时投递(`runScheduledTick` / `runTask`)且适配器实现了 `claimTask` 时可用。请求内当场投递的 `instant` 任务、没实现 `claimTask` 的自定义适配器返回 `{ defer }` 会得到 `AGENTIC_DEFER_UNSUPPORTED` 配置错误。 +- `GET /capabilities` 的 features 新增 `before-fire-defer`。 diff --git a/packages/rei-standard-amsg/server/README.md b/packages/rei-standard-amsg/server/README.md index 9ced586..b6552d1 100644 --- a/packages/rei-standard-amsg/server/README.md +++ b/packages/rei-standard-amsg/server/README.md @@ -364,6 +364,16 @@ hook 在 `pushPayloads` 里自己写了这几个字段的话会被库覆盖: 配上 `hooks: { onBeforeFire, onLLMOutput, executeToolCalls }` 之后,AI 类任务的 prompt 不再是排程那一刻冻结的文本,而是 cron 触发时现场组装,工具也在服务端就地跑完,全程不需要客户端在线。完整用法见 [`examples/cloudflare-single-user/README.md`](https://github.com/Tosd0/ReiStandard/blob/main/packages/rei-standard-amsg/server/examples/cloudflare-single-user/README.md) 的「Fire 时刻 hooks」。 +`onBeforeFire` 的返回值决定这次触发怎么走: + +| 返回值 | 这次触发怎么走 | +|---|---| +| 消息数组,或 `{ messages, maxToolIterations?, totalTimeoutMs?, tools?, toolChoice? }` | 用这份 prompt 生成 | +| `{ skip: true }` | 这次不发,按零推送的成功收场:一次性任务删掉、循环任务推进到下一次 | +| `{ defer: { afterMs } }` | 现在不合适,`afterMs` 毫秒后再来问。见下面[推迟这次触发](#推迟这次触发defer) | +| `null` | 交还给排程时冻结的 prompt | +| 抛错 | 按投递失败处理:占一次重试、写 `last_error`、2 / 4 / 6 分钟退避 | + 三个 hook 拿到的 ctx 上都有这几个口子: | ctx 上的口子 | 干什么 | @@ -383,6 +393,41 @@ hook 在 `pushPayloads` 里自己写了这几个字段的话会被库覆盖: `sessionId` 是给日志和去重用的不透明字符串,格式随版本变,别拆它拿上面这些值。 +### 推迟这次触发(defer) + +任务到点了,但宿主这会儿不方便生成——比如同一段对话里正有另一轮回复在生成,想等它结束再说。`onBeforeFire` 返回 `{ defer: { afterMs } }`: + +```js +async onBeforeFire(ctx) { + if (await isConversationBusy(ctx)) return { defer: { afterMs: 30_000 } }; + return buildMessages(ctx); +} +``` + +这一次不调 LLM、不推送,也不调 `onLLMOutput` / `onAfterSend`。任务行上只有两处变化:`retry_after` 写成现在 + `afterMs`,租约放掉。到点之后 cron 把它重新捞起来,从 `onBeforeFire` 再走一遍,那时可以继续推迟。 + +跟抛错重试的区别: + +| | 抛错 | `{ defer }` | +|---|---|---| +| `retry_count` | 加一,用完就终审 | 不变 | +| `last_error` | 写上这次的原因 | 不写,也不清已有的 | +| 多久之后再来 | 固定 2 / 4 / 6 分钟 | 宿主给的 `afterMs` | +| `onFireSettled` 的 `status` | `failed` | `deferred`,带 `retryAfter` | + +两边相同的是 `next_send_at` 都不动:它一直是这次触发的名义时刻,`ctx.task.nextSendAt`、`occurrenceMs`、推送的 `messageId` 在推迟前后都是同一套。循环任务被推迟时不会推进到下一次。 + +几条规矩: + +- `afterMs` 必须是有限正数,最大 24 小时(导出为 `MAX_DEFER_AFTER_MS`)。不合法的值按 hook 契约违约处理(`AGENTIC_BAD_BEFORE_FIRE`,一跳终审)。没有下限,但 cron 一分钟一跳,实际是「到点之后的第一跳」才会再问;用 `runTask` 自己触发的话,到点之前调它会得到 `retry_pending`。 +- 推迟不会让任务永远挂着。名义时刻过去超过 `staleAfterMs`(默认 60 分钟)的任务照常按过期处理:一次性任务标 `failed`、循环任务快进到下一次,并调 `onStaleSkip`。某次推迟的唤醒时刻已经越过这条线时,库不再等那一轮,当场按过期收场——这时宿主会先收到 `status: 'deferred'` 的 `onFireSettled`,紧接着收到 `onStaleSkip`。 +- 每次被重新问到都是一次新的 fire:`scratch` 是新的,`onFireSettled` 各调各的。 +- 只有定时投递(`runScheduledTick`、单用户 Worker 的 `scheduled()`、`runTask`)支持推迟,并且适配器要实现 `claimTask`(任务表有 `retry_after` 列,内置适配器都有)。`messageType: 'instant'` 在请求里当场投递的任务、以及没实现 `claimTask` 的自定义适配器返回 `{ defer }` 会得到 `AGENTIC_DEFER_UNSUPPORTED`,按[部署配错了算可重试](#部署配错了算可重试)处理。 +- 内容已经落进收件箱、只是推送没发完的重试不调 `onBeforeFire`(见[生成成功之后推送失败:只补推送](#生成成功之后推送失败只补推送)),那一跳也就不会被推迟。 +- 配了 `serializeBy` 时,被推迟的任务不占着分组:同组里排在它后面的任务可以先跑。 + +tick 汇总的 `details.deferredTasks`(`{ taskId, retryAfter }`)列出这一跳被推迟的任务,它们不计入 `successCount` / `failedCount`。`GET /capabilities` 的 features 里有 `before-fire-defer`。 + ### config 级 hook 这两个挂在 worker 工厂 config 的顶层(不在 `hooks` 里): @@ -390,7 +435,7 @@ hook 在 `pushPayloads` 里自己写了这几个字段的话会被库覆盖: | hook | 什么时候调 | 载荷 | |---|---|---| | `onAfterSend` | fire 的 pushPayloads 逐段发完,或中途发挂 | `{ task, sentCount, pushedCount, total, error, usage, usageTotal, llmCalls, outboxed, scratch, readState, writeState, emitResult }` | -| `onFireSettled` | 一次 fire 收尾——只要 `onBeforeFire` 被调用过,什么结局都调一次 | `{ task, status, skipReason, sentCount, pushedCount, total, iterations, error, metadata, usage, usageTotal, llmCalls, outboxed, scratch, readState, writeState, emitResult }` | +| `onFireSettled` | 一次 fire 收尾——只要 `onBeforeFire` 被调用过,什么结局都调一次 | `{ task, status, skipReason, retryAfter, sentCount, pushedCount, total, iterations, error, metadata, usage, usageTotal, llmCalls, outboxed, scratch, readState, writeState, emitResult }` | | `onStaleSkip` | 任务错过触发时刻超过 60 分钟、这一次(或这几次)不再补发 | `{ reason, action, metadata, recurrenceType, occurrenceMs, skippedCount, skippedOccurrences, skippedTruncated, nextSendAt, readState, writeState, emitResult }` | 三个 hook 都自带 `readState` / `writeState` / `emitResult`,作用于当前用户,语义与 fire 级那套一致。`onStaleSkip` 尤其需要:服务停摆恢复后的第一跳里可能一次 fire 都没跑过,而那正是它要留痕迹的时候。 @@ -405,7 +450,9 @@ hook 在 `pushPayloads` 里自己写了这几个字段的话会被库覆盖: |---|---| | `sent` | pushPayloads 全部发完(`sentCount === total`) | | `skipped` | 这次不发。`skipReason` 区分是 `onBeforeFire` 直接 `{ skip: true }`(`'before-fire'`)还是模型跑完后判定不发(`'skip-push'`) | +| `deferred` | `onBeforeFire` 返回了 `{ defer: { afterMs } }`:这次没生成也没发,任务还在。`retryAfter`(ISO 时刻)之后会从 `onBeforeFire` 重新走一遍;这个字段在其余结局都是 `null` | | `failed` | 链路抛错,`error` 带原始错误。发到第 k 段挂了也是这个:`sentCount = k`、`total` 是原本要发的段数 | +| `cancelled` | 任务在投递期间被取消或顶替,见[取消正在执行的任务](#取消正在执行的生成与工具) | | `not-handled` | `onBeforeFire` 返回 `null`,这条任务交还给排程时冻结的 prompt 老链路。那条链路不归 fire hook 管,它后面发没发出去不体现在这里 | **用量记账**看这三个字段,两个 hook 都带,`onFireSettled` 的每种结局(发完、跳过、失败)都带——失败也花了钱: @@ -507,7 +554,7 @@ const { messageId, pushed } = await ctx.emitResult({ ### hook 契约违约算确定性失败 -宿主 hook 返回了库不认的东西(`onBeforeFire` 的返回形状、`onLLMOutput` 的决策标签),或者建后续任务时 `createTask` 没把行交回来——这些错误带 `permanent: true` 和一个稳定的 `code`(`AGENTIC_BAD_BEFORE_FIRE` / `AGENTIC_BAD_DECISION` / `AGENTIC_SCHEDULE_FAILED` / `TASK_PAYLOAD_TOO_LARGE`),投递侧据此跳过退避阶梯:一次性任务直接标 `failed`,循环任务作废本次 occurrence。重试也是同一个结果,而每重试一轮都要把 `onBeforeFire` 和一整轮 LLM 重跑一遍。 +宿主 hook 返回了库不认的东西(`onBeforeFire` 的返回形状、`{ defer }` 里不合法的 `afterMs`、`onLLMOutput` 的决策标签),或者建后续任务时 `createTask` 没把行交回来——这些错误带 `permanent: true` 和一个稳定的 `code`(`AGENTIC_BAD_BEFORE_FIRE` / `AGENTIC_BAD_DECISION` / `AGENTIC_SCHEDULE_FAILED` / `TASK_PAYLOAD_TOO_LARGE`),投递侧据此跳过退避阶梯:一次性任务直接标 `failed`,循环任务作废本次 occurrence。重试也是同一个结果,而每重试一轮都要把 `onBeforeFire` 和一整轮 LLM 重跑一遍。 分界线是「谁写错了」:契约由宿主代码定死,重掷一次还是同一个形状;而模型这一轮掷出了什么则是每轮都可能不同的。所以「tool-request 决策里没有能解析的 `toolCalls`」(`AGENTIC_EMPTY_TOOL_REQUEST`)和「轮数用尽也没等到 `finish` / `skip-push`」(`AGENTIC_LOOP_EXCEEDED`)带 `code` 但不带 `permanent`,留在退避阶梯上——隔两分钟重掷一次多半就正常收尾了,判终态的话一次性任务第一次掷歪就永久 `failed`,行离开 `pending` 之后连 `PUT /update-message` 都救不回来(回 409)。 diff --git a/packages/rei-standard-amsg/server/examples/cloudflare-single-user/README.md b/packages/rei-standard-amsg/server/examples/cloudflare-single-user/README.md index ca2f4b4..0885934 100644 --- a/packages/rei-standard-amsg/server/examples/cloudflare-single-user/README.md +++ b/packages/rei-standard-amsg/server/examples/cloudflare-single-user/README.md @@ -157,6 +157,9 @@ export default createSingleUserCloudflareWorker((env) => ({ // 或返回 { skip: true }:这次不生成,零推送直接算成功结束(不调 LLM)。 // 一次性任务照删、循环任务照推进到下次。适合排程后对话已有新进展、 // 这条到点已多余的情况。 + // 或返回 { defer: { afterMs: 30_000 } }:现在不合适,30 秒后再来问。 + // 这次不调 LLM、不推送,任务原样留着(不占重试次数、不留报错、触发 + // 时刻不变),到点后从 onBeforeFire 重新走一遍。afterMs 最大 24 小时。 }, // 每轮 LLM 输出后分类。ctx 形状与 @rei-standard/amsg-instant 的 @@ -317,13 +320,13 @@ export default createSingleUserCloudflareWorker((env) => ({ ## 一次 fire 的收尾回执(onFireSettled) -这是啥:**只要 `onBeforeFire` 被调用过**,这次 fire 无论是发完了、跳过了、还是半路抛错,都会调一次 `onFireSettled`。「开始时占点什么、结束时放掉」的写法挂这个。 +这是啥:**只要 `onBeforeFire` 被调用过**,这次 fire 无论是发完了、跳过了、推迟了、还是半路抛错,都会调一次 `onFireSettled`。「开始时占点什么、结束时放掉」的写法挂这个。 ```js export default createSingleUserCloudflareWorker((env) => ({ // ...其余 config async onFireSettled(info) { - // info: { task, status, skipReason, sentCount, pushedCount, total, iterations, + // info: { task, status, skipReason, retryAfter, sentCount, pushedCount, total, iterations, // error, metadata, usage, usageTotal, llmCalls, outboxed, // scratch, readState, writeState } if (info.scratch.scheduledFollowUp) { @@ -336,12 +339,13 @@ export default createSingleUserCloudflareWorker((env) => ({ })); ``` -`status` 四种: +`status` 常见的几种: | status | 什么时候 | |---|---| | `sent` | pushPayloads 全部发完(`sentCount === total`) | | `skipped` | 这次不发。`skipReason` 区分是 `onBeforeFire` 直接 `{ skip: true }`(`'before-fire'`)还是模型跑完后判定不发(`'skip-push'`) | +| `deferred` | `onBeforeFire` 返回了 `{ defer: { afterMs } }`:这次没生成也没发,`retryAfter`(ISO 时刻)之后再从 `onBeforeFire` 走一遍 | | `failed` | 链路抛错,`error` 带原始错误。发到第 k 段挂了也是这个:`sentCount = k`、`total` 是原本要发的段数 | | `not-handled` | `onBeforeFire` 返回 `null`,这条任务交还给排程时冻结的 prompt 老链路。那条链路不归 fire hook 管 | diff --git a/packages/rei-standard-amsg/server/src/server/cloudflare.js b/packages/rei-standard-amsg/server/src/server/cloudflare.js index a952a84..7bf86f9 100644 --- a/packages/rei-standard-amsg/server/src/server/cloudflare.js +++ b/packages/rei-standard-amsg/server/src/server/cloudflare.js @@ -35,6 +35,8 @@ export { // 投递失败重试次数的默认值(config 的 maxDeliveryRetries 不配时用它)。 DEFAULT_MAX_DELIVERY_RETRIES, } from './lib/run-tick.js'; +// onBeforeFire 的 { defer: { afterMs } } 里 afterMs 的上限。 +export { MAX_DEFER_AFTER_MS } from './lib/agentic-fire.js'; // Schema 自查 / 补齐:升级后老部署的表没跟上时,cron 会每分钟静默挂在缺的那 // 一列上。见 lib/schema-version.js。 export { getSchemaVersion, ensureSchema, SCHEMA_VERSION } from './lib/schema-version.js'; diff --git a/packages/rei-standard-amsg/server/src/server/cloudflare/single-user-worker.js b/packages/rei-standard-amsg/server/src/server/cloudflare/single-user-worker.js index 56629f0..b7c7928 100644 --- a/packages/rei-standard-amsg/server/src/server/cloudflare/single-user-worker.js +++ b/packages/rei-standard-amsg/server/src/server/cloudflare/single-user-worker.js @@ -54,6 +54,9 @@ * tasks replay the schedule-time frozen prompt exactly as before. See * lib/agentic-fire.js. * + * `onBeforeFire` 除了返回 prompt,还可以返回 `{ skip: true }`(这次不发)或 + * `{ defer: { afterMs } }`(现在不合适,afterMs 毫秒后再来问;不占重试次数)。 + * * scheduled() 每次触发都会先给任务占位(在行的 lease_until 上写租约),同一 * 条任务不会被相邻两跳重复触发(见 lib/run-tick.js)。投递期间租约按心跳滚动 * 续租(默认 30s 心跳 / 90s 租约):isolate 中途被回收时任务在 ~90 秒内就能 diff --git a/packages/rei-standard-amsg/server/src/server/handlers/capabilities.js b/packages/rei-standard-amsg/server/src/server/handlers/capabilities.js index 0542649..5d9c4d7 100644 --- a/packages/rei-standard-amsg/server/src/server/handlers/capabilities.js +++ b/packages/rei-standard-amsg/server/src/server/handlers/capabilities.js @@ -127,6 +127,9 @@ export const SERVER_FEATURES = Object.freeze([ 'max-delivery-retries', // Per-task pre-commit retry limit and fire receipt retry decision. 'max-generation-retries', + // onBeforeFire 认 { defer: { afterMs } }:这次不生成,过一会儿再来问(不占重试 + // 次数);onFireSettled 的 status 多一种 'deferred',带 retryAfter。 + 'before-fire-defer', ]); export function createCapabilitiesHandler(ctx) { diff --git a/packages/rei-standard-amsg/server/src/server/index.js b/packages/rei-standard-amsg/server/src/server/index.js index d1200e1..452d050 100644 --- a/packages/rei-standard-amsg/server/src/server/index.js +++ b/packages/rei-standard-amsg/server/src/server/index.js @@ -204,6 +204,8 @@ export { DEFAULT_TOTAL_TIMEOUT_MS, DEFAULT_MAX_SCHEDULED_TASKS_PER_FIRE, MIN_SCHEDULE_LEAD_MS, + // onBeforeFire 的 { defer: { afterMs } } 里 afterMs 的上限。 + MAX_DEFER_AFTER_MS, } from './lib/agentic-fire.js'; // 单用户 worker 的 CORS 允许头/方法列表(外层再包路由的宿主 import 这一份, // 别手抄第二份)。 diff --git a/packages/rei-standard-amsg/server/src/server/lib/agentic-fire.js b/packages/rei-standard-amsg/server/src/server/lib/agentic-fire.js index 17eadf6..2c3a408 100644 --- a/packages/rei-standard-amsg/server/src/server/lib/agentic-fire.js +++ b/packages/rei-standard-amsg/server/src/server/lib/agentic-fire.js @@ -5,6 +5,7 @@ * the completePrompt frozen at schedule time. At fire time instead: * * onBeforeFire(fireCtx) → fresh messages (may read client_state) | { skip: true } + * | { defer: { afterMs } }(现在不合适,过一会儿再来问) * → callLlm → onLLMOutput(sessionCtx) → decision * ├─ 'finish' → push decision.pushPayloads, done * ├─ 'skip-push' → record, done (task counts as delivered) @@ -28,11 +29,22 @@ * * 收尾回执:onAfterSend 只走「有 push 要发」这条路。**只要 onBeforeFire 被 * 调用过**,无论结局是发完、跳过(skip / skip-push)还是抛错,可选的 - * ctx.onFireSettled?.({ task, status, skipReason, sentCount, pushedCount, - * total, iterations, error, usage, usageTotal, llmCalls, outboxed, scratch, - * readState, writeState }) 都会被调用一次(见 notifyFireSettled)。「开始时占 + * ctx.onFireSettled?.({ task, status, skipReason, retryAfter, sentCount, + * pushedCount, total, iterations, error, usage, usageTotal, llmCalls, outboxed, + * scratch, readState, writeState }) 都会被调用一次(见 notifyFireSettled)。 + * 推迟(defer)也算一种结局,同样调一次。「开始时占 * 点什么、结束时放掉」的写法挂这个才不会漏。 * + * 推迟(defer):onBeforeFire 返回 `{ defer: { afterMs } }` 表示「现在不合适, + * 过 afterMs 毫秒再来问我」——比如任务到点时宿主那边正有另一件事占着同一段 + * 对话。这次不调 LLM、不推送、不调 onLLMOutput / onAfterSend,任务行只写 + * `retry_after = now + afterMs` 并放掉租约:`retry_count` 不涨、`last_error` 不 + * 动、`next_send_at`(本次触发的名义时刻)不动,循环任务也不推进。到点之后 + * cron 再把它捞起来,从 onBeforeFire 重新走一遍(可以再推迟)。它不是失败: + * 抛错会占一格重试、留下 last_error,推迟两样都不占。收尾回执的 status 是 + * `deferred`,带 `retryAfter`。afterMs 的范围见 MAX_DEFER_AFTER_MS;能不能推迟 + * 取决于投递入口,见 runFireChain 里的 AGENTIC_DEFER_UNSUPPORTED。 + * * 失败之后的重试:LLM 上游明确拒了这次请求(401 / 403 / 400 …)一跳终审, * 错误上带 `permanent: true`(见 lib/errors.js 的 isPermanentLlmFailure)。 * finish 的整批落进 outbox 之后推送才失败的,重试那一跳只补推送、不再调这里 @@ -161,6 +173,12 @@ export const DEFAULT_MAX_SCHEDULED_TASKS_PER_FIRE = 2; // 自排任务的最小提前量:cron 一分钟一跳,比这更近等于让下一跳立刻捡走。 export const MIN_SCHEDULE_LEAD_MS = 60_000; +// onBeforeFire 的 `{ defer: { afterMs } }` 最多能把这次触发往后推多久。推迟是 +// 给「等手头这件事结束」用的短等待;想等得更久,该改的是任务的排期。上限同时 +// 保证 now + afterMs 一定是个合法的时刻。没有下限:cron 一分钟一跳,再小的 +// afterMs 也是「到点之后的第一跳」才会被重新问到。 +export const MAX_DEFER_AFTER_MS = 24 * 60 * 60 * 1000; + // 自排任务允许的类型。instant 归 POST /schedule-message 那条同步路径管。 const SCHEDULABLE_MESSAGE_TYPES = new Set(['auto', 'prompted', 'fixed']); // 与 validateScheduleMessagePayload 同一套;run-tick 只认得这三种,别的值会让 @@ -249,6 +267,52 @@ export function buildHookTask(task, decryptedPayload) { }); } +/** + * onBeforeFire 的返回值。 + * + * @typedef {Array + * | { messages: Array, maxToolIterations?: number, totalTimeoutMs?: number, tools?: Array, toolChoice?: any } + * | { skip: true } + * | { defer: { afterMs: number } } + * | null + * | undefined} BeforeFireResult + * - 消息数组 / `{ messages, ... }`:用这份 prompt 生成; + * - `{ skip: true }`:这次不发,按零推送的成功收场(一次性任务删行、循环任务 + * 推进到下一次); + * - `{ defer: { afterMs } }`:现在不合适,afterMs 毫秒后再来问。afterMs 必须 + * 是有限正数,且不超过 MAX_DEFER_AFTER_MS(24 小时); + * - `null` / `undefined`:交还给排程时冻结的 prompt 老链路。 + */ + +/** + * 收尾回执(onFireSettled)的 status。 + * + * @typedef {'sent'|'skipped'|'deferred'|'failed'|'cancelled'|'not-handled'} FireSettledStatus + */ + +const BAD_BEFORE_FIRE_MESSAGE = + 'AGENTIC_BAD_BEFORE_FIRE: onBeforeFire must return ChatMessage[] | { messages, maxToolIterations?, totalTimeoutMs?, tools?, toolChoice? } | { skip: true } | { defer: { afterMs } } | null'; + +/** + * 校验 `{ defer: { afterMs } }` 并取出 afterMs。形状不对、不是有限正数、或超过 + * MAX_DEFER_AFTER_MS,都按 onBeforeFire 的契约违约处理(确定性失败,不重试)。 + * + * @param {unknown} defer + * @returns {number} afterMs + */ +function readDeferAfterMs(defer) { + const afterMs = defer && typeof defer === 'object' ? /** @type {any} */ (defer).afterMs : undefined; + if (typeof afterMs !== 'number' || !Number.isFinite(afterMs) || afterMs <= 0 || afterMs > MAX_DEFER_AFTER_MS) { + throw markPermanent( + new TypeError( + `${BAD_BEFORE_FIRE_MESSAGE}(defer.afterMs 必须是有限正数,且不超过 ${MAX_DEFER_AFTER_MS} 毫秒)` + ), + 'AGENTIC_BAD_BEFORE_FIRE' + ); + } + return afterMs; +} + function normalizeBeforeFireResult(result) { if (Array.isArray(result)) { return { messages: result }; @@ -265,12 +329,7 @@ function normalizeBeforeFireResult(result) { toolChoice: result.toolChoice, }; } - throw markPermanent( - new TypeError( - 'AGENTIC_BAD_BEFORE_FIRE: onBeforeFire must return ChatMessage[] | { messages, maxToolIterations?, totalTimeoutMs?, tools?, toolChoice? } | { skip: true } | null' - ), - 'AGENTIC_BAD_BEFORE_FIRE' - ); + throw markPermanent(new TypeError(BAD_BEFORE_FIRE_MESSAGE), 'AGENTIC_BAD_BEFORE_FIRE'); } function firstPositiveInt(values, fallback) { @@ -295,12 +354,18 @@ function firstPositiveNumber(values, fallback) { * @param {Object} args.decryptedPayload - decrypted task payload (has credentials; they stop here) * @param {string} args.userKey - per-user storage key (for readState decryption) * @param {Object} args.ctx - processor ctx ({ db, webpush, vapid, hooks, maxToolIterations, totalTimeoutMs, maxScheduledTasksPerFire }) - * @returns {Promise<{ handled: false } | { handled: true, result: { success: true, messagesSent: number, status: 'finished'|'skipped', iterations: number } }>} + * @returns {Promise<{ handled: false } + * | { handled: true, result: { success: true, messagesSent: number, status: 'finished'|'skipped', iterations: number } } + * | { handled: true, result: { success: false, deferred: true, retryAfter: string, messagesSent: 0, status: 'deferred', iterations: 0 } }>} * `handled: false` → caller falls back to the legacy frozen-prompt path * (onBeforeFire returned null). * `onBeforeFire` may also return `{ skip: true }` to complete the fire * before the first LLM call → `status: 'skipped', iterations: 0`, same * success handling as the post-LLM skip-push path. + * `{ defer: { afterMs } }` → `status: 'deferred'`,`retryAfter` 是「什么时候 + * 再来问」的 ISO 时刻。`success` 给 false 而不是 true:这次什么都没发,调用 + * 方得看 `deferred` 走推迟的收尾;万一有调用方不认识这个字段,也只会把它当 + * 成没发成,不会当成发完了把任务删掉。 * Failures (timeout / loop exceeded / config errors) throw — the caller's * existing error handling turns them into task retry/failure. * @@ -798,6 +863,8 @@ export async function runAgenticFire({ task, decryptedPayload, userKey, ctx }) { // 结局默认按 failed 记:下面只要有任何一步抛出去,finally 里发出的就是这个。 let settledStatus = 'failed'; let settledError = null; + // 只有 deferred 这个结局带值:什么时候再来问(ISO 时刻)。 + let settledRetryAfter = null; try { const outcome = await runFireChain({ task, decryptedPayload, userKey, ctx, hooks, nowFn, sleep, @@ -806,9 +873,11 @@ export async function runAgenticFire({ task, decryptedPayload, userKey, ctx }) { sessionId, messageIdBase, occurrenceMs, }); throwIfCancelled(); - settledStatus = !outcome.handled - ? 'not-handled' - : (outcome.result.status === 'skipped' ? 'skipped' : 'sent'); + if (!outcome.handled) settledStatus = 'not-handled'; + else if (outcome.result.status === 'deferred') { + settledStatus = 'deferred'; + settledRetryAfter = outcome.result.retryAfter; + } else settledStatus = outcome.result.status === 'skipped' ? 'skipped' : 'sent'; return outcome; } catch (error) { if (isCancelled() && !isTaskCancelledError(error)) { @@ -828,6 +897,7 @@ export async function runAgenticFire({ task, decryptedPayload, userKey, ctx }) { ? failureRetryDecision(ctx._deliveryState, settledError) : { willRetry: null, failureStage: null }), skipReason: settledStatus === 'skipped' ? progress.skipReason : null, + retryAfter: settledStatus === 'deferred' ? settledRetryAfter : null, sentCount: progress.sentCount, pushedCount: progress.pushedCount, total: progress.total, @@ -883,6 +953,41 @@ async function runFireChain({ return { handled: true, result: { success: true, messagesSent: 0, status: 'skipped', iterations: 0 } }; } + // 推迟:现在不合适,过一会儿再来问。这里只算出「什么时候再来」,写库(把 + // retry_after 推到那个时刻、放掉租约)是 run-tick 的事。 + if (typeof before === 'object' && !Array.isArray(before) && before.defer != null) { + const afterMs = readDeferAfterMs(before.defer); + // 推迟靠 retry_after 这一列落地,只有定时投递(runScheduledTick / runTask) + // 且适配器实现了 claimTask 时才有这一列可写,由 run-tick 用 ctx._deferSupported + // 告诉这里。别的入口都不行: + // - 没实现 claimTask 的自定义适配器只能把时间写进 next_send_at,而那是 + // 这次触发的名义时刻(宿主拿它当触发身份、循环推进拿它当基准),不能改; + // - `messageType: 'instant'` 在请求里当场投递(processMessagesByUuid), + // 没有「过一会儿再捞起来」这回事。 + // 这是部署 / 用法层面的问题,按配置错误报(可重试、不判终态),不悄悄当成 + // 跳过或立刻再问。 + if (ctx._deferSupported !== true) { + throw new DeploymentConfigError( + 'AGENTIC_DEFER_UNSUPPORTED: onBeforeFire 返回了 { defer },但这条投递路径不支持推迟——' + + '只有 runScheduledTick / runTask 且适配器实现了 claimTask(任务表有 retry_after 列)时可用,' + + '请求内当场投递的 instant 任务不可用', + { code: 'AGENTIC_DEFER_UNSUPPORTED' } + ); + } + return { + handled: true, + result: { + success: false, + deferred: true, + // 用真实时钟:这个时刻要落进 retry_after,跟数据库捞取条件里的「现在」比。 + retryAfter: new Date(Date.now() + afterMs).toISOString(), + messagesSent: 0, + status: 'deferred', + iterations: 0, + }, + }; + } + const normalized = normalizeBeforeFireResult(before); const maxToolIterations = firstPositiveInt( [normalized.maxToolIterations, ctx.maxToolIterations], @@ -1144,11 +1249,11 @@ async function notifyAfterSend(ctx, info) { } /** - * ctx.onFireSettled?.({ task, status, skipReason, sentCount, pushedCount, - * total, iterations, error, metadata, usage, usageTotal, llmCalls, outboxed, - * scratch, readState, writeState }) + * ctx.onFireSettled?.({ task, status, skipReason, retryAfter, sentCount, + * pushedCount, total, iterations, error, metadata, usage, usageTotal, llmCalls, + * outboxed, scratch, readState, writeState }) * —— 一次 fire 收尾的可选 hook。**onBeforeFire 被调用过,这个就一定会被调用 - * 一次**,无论这次是发完了、跳过了、还是半路抛错。 + * 一次**,无论这次是发完了、跳过了、推迟了、还是半路抛错。 * * sentCount 是这批走完了几段,pushedCount 是其中真的占用了推送通道的有几条 * (不会弹通知的段默认只落收件箱,见 lib/push-policy.js)。 @@ -1183,12 +1288,20 @@ async function notifyAfterSend(ctx, info) { * 了,但记账的代码挂在发送后,这次没发成就没人记,那条任务从此只活在数据库 * 里;以及 fire 开头拿的锁没有可靠的释放点,一次 skip 就把资源占满整个 TTL。 * - * status 五种: + * status 六种(见 FireSettledStatus): * - `cancelled` —— 任务被取消/顶替;error.code 为 TASK_CANCELLED,不重试 * - `sent` —— pushPayloads 全部发完(sentCount === total) * - `skipped` —— 这次不发。skipReason 区分是 onBeforeFire 直接 * `{ skip: true }`(`'before-fire'`),还是模型跑完之后 * 判定不发(`'skip-push'`) + * - `deferred` —— onBeforeFire 返回了 `{ defer: { afterMs } }`:这次没 + * 生成也没发,任务还在,retryAfter(ISO 时刻)之后会从 + * onBeforeFire 重新走一遍。到那时又是一次新的 fire,有 + * 自己的 scratch 和自己的一次收尾回执。retryAfter 只在这 + * 个结局有值,其余为 null。回执发生在写库之前:之后这次 + * 推迟也可能因为唤醒时刻越过了过期线而直接按过期收场 + * (见 lib/run-tick.js 的 handleDeferred),那时宿主会再 + * 收到一次 onStaleSkip * - `failed` —— 链路抛错,error 带原始错误。部分失败也是这个:发到第 * k 段挂了 → sentCount = k、total 是原本要发的段数 * - `not-handled` —— onBeforeFire 返回 null,这条任务交还给排程时冻结的 diff --git a/packages/rei-standard-amsg/server/src/server/lib/message-processor.js b/packages/rei-standard-amsg/server/src/server/lib/message-processor.js index 065a288..dc0191b 100644 --- a/packages/rei-standard-amsg/server/src/server/lib/message-processor.js +++ b/packages/rei-standard-amsg/server/src/server/lib/message-processor.js @@ -371,7 +371,8 @@ function positiveIntegerOr(value, fallback) { * 样,客户端补收拿到的、推送收到的、首次推到一半已经收到的,全是同一份。 * * 推送再失败就照常抛出去,由调用方走既有的重试 / 终审逻辑——下一跳还是来这里补 - * 推,照样不生成。 + * 推,照样不生成。onBeforeFire 既然不调,这一跳也就不存在「推迟」(defer): + * 内容已经定了,剩下的只是把它送到。 * * @param {Object} task * @param {ProcessorContext} ctx @@ -432,8 +433,11 @@ async function redeliverCommittedBatch(task, ctx, userKey, decryptedPayload, bat * @param {string} [providedMasterKey] * @param {{ userKey: string, payload: Object } | null} [predecrypted] - 调用方 * (run-tick 的预扫描)已经解好的 payload;传了就不再解第二遍。 - * @returns {Promise<{ success: boolean, messagesSent: number, redelivered?: boolean, pushedCount?: number, error?: string, errorCode?: string|null, pushStatusCode?: number|null, permanent?: boolean, willRetry?: boolean, retryLimit?: number, failureStage?: string }>} + * @returns {Promise<{ success: boolean, messagesSent: number, redelivered?: boolean, pushedCount?: number, deferred?: boolean, retryAfter?: string, error?: string, errorCode?: string|null, pushStatusCode?: number|null, permanent?: boolean, willRetry?: boolean, retryLimit?: number, failureStage?: string }>} * 失败时 `pushStatusCode` 是推送服务回的 HTTP 状态码(不是推送阶段炸的 → null)。 + * `deferred: true` 表示 onBeforeFire 把这次触发推迟了(`retryAfter` 是唤醒 + * 时刻):`success` 为 false 但不是失败,调用方要先看这个字段。只有 ctx 带 + * `_deferSupported: true`(run-tick 在适配器有 retry_after 列时给)才会出现。 */ export async function processSingleMessage(task, ctx, providedMasterKey, predecrypted = null) { const deliveryState = { @@ -484,6 +488,8 @@ export async function processSingleMessage(task, ctx, providedMasterKey, predecr // is byte-identical. onBeforeFire → { skip: true } completes the fire // here as a zero-push success (no LLM call, no frozen-prompt fallback): // use it when the host can tell at fire time the message is moot. + // onBeforeFire → { defer: { afterMs } } 也在这里结束:结果带 deferred: true + // 和 retryAfter,没生成也没推送,由 run-tick 把行推迟到那个时刻。 if (ctx.hooks && typeof ctx.hooks.onBeforeFire === 'function' && taskNeedsLlm(decryptedPayload)) { const agentic = await runAgenticFire({ task, decryptedPayload, userKey, ctx }); if (agentic.handled) return agentic.result; @@ -730,6 +736,9 @@ export async function processSingleMessage(task, ctx, providedMasterKey, predecr * `reasoningError` 只在正文都发出去了、思考过程那一条没发成时出现(思考过程是 * 附赠内容,它发不出去不算这条消息失败)。调用方拿它提示用户这次没有思考过程, * 不带这个字段就是整轮都送到了。 + * + * 这条路是请求里当场投递,不支持 onBeforeFire 的 `{ defer }`:返回它会得到 + * AGENTIC_DEFER_UNSUPPORTED 的配置错误,按这里的失败流程处理。 */ export async function processMessagesByUuid(uuid, ctx, maxRetries = 2, userId, providedMasterKey) { let retryCount = 0; diff --git a/packages/rei-standard-amsg/server/src/server/lib/run-tick.js b/packages/rei-standard-amsg/server/src/server/lib/run-tick.js index 793cd6a..2246583 100644 --- a/packages/rei-standard-amsg/server/src/server/lib/run-tick.js +++ b/packages/rei-standard-amsg/server/src/server/lib/run-tick.js @@ -13,6 +13,10 @@ * 任务列表要读它、循环任务推进下一次也要拿它当基准。投递收尾时把租约放掉, * 这条任务立刻可以被下一跳接手。 * + * retry_after 还有第二个用途:onBeforeFire 返回 `{ defer: { afterMs } }` 时, + * 这一列写成「什么时候再来问」,别的字段都不动(见 deliverClaimedTask 里的 + * handleDeferred)。两种用途靠 retry_count 涨没涨分得开——推迟不涨。 + * * 投递失败的退避写在另一列(retry_after)上,不跟租约挤在一起:租约的意思只 * 有「这条正在跑」,而正在等重试的任务其实闲着——两件事共用一列的话,分组串行 * (见下)会把一条闲着的任务当成「这个分组忙着」,同组别的任务白等一轮退避。 @@ -100,6 +104,8 @@ export const DEFAULT_HEARTBEAT_LEASE_TTL_MS = 90 * 1000; // 补发一天。正在重试链上的任务(retry_count > 0)不算过期——它的 next_send_at // 一直是名义时刻,重试拖过一小时不等于用户错过了它——但重试时刻(retry_after) // 本身也被拖过了这个时长的除外:那说明停摆发生在重试窗口里,内容一样旧。 +// 被 onBeforeFire 推迟(defer)的任务同样受它管:推迟不涨 retry_count,一直推 +// 下去的任务最后在这里收场。 // 宿主可用 ctx.staleAfterMs 覆盖(与 claimLeaseMs 同一模式)。 export const STALE_AFTER_MS = 60 * 60 * 1000; @@ -360,8 +366,11 @@ async function cleanupExpiredClientState(ctx) { * - `already_settled`:行还在,但已经是 sent / failed(`status` 带上是哪 * 个)。适配器没实现 `getTaskStatusByUuidOnly` 时这种情况并进 `not_found`。 * - `not_due`:还没到 `next_send_at`(`nextSendAt` 带上是什么时候)。 - * - `retry_pending`:上次投递失败,还在退避窗口里(`retryAfter` 带上什么时 - * 候到点)。 + * - `retry_pending`:上次投递失败还在退避窗口里,或者上次 onBeforeFire 把 + * 它推迟了(`{ defer }`)还没到点(`retryAfter` 带上什么时候到点)。 + * + * 这一次跑的时候 onBeforeFire 返回 `{ defer }` 的话,结果是 `ran: true`,任务 + * 列在 summary.details.deferredTasks 里;到点之前再调都回 `retry_pending`。 * * @param {Object} ctx - 与 runScheduledTick 同一份 ctx * @param {string} uuid - 任务 uuid(pending 行) @@ -435,6 +444,7 @@ async function deliverTasks(ctx, tasks) { deletedOnceOffTasks: 0, updatedRecurringTasks: 0, staleTasks: [], + deferredTasks: [], cancelledTasks: [], reasoningSkippedTasks: [], redeliveredTasks: [], @@ -983,9 +993,9 @@ async function deliverTasks(ctx, tasks) { const retryAfterMs = task.retry_after ? Date.parse(task.retry_after) : NaN; const notOnFreshRetryChain = (task.retry_count || 0) === 0 || (Number.isFinite(retryAfterMs) && Date.now() - retryAfterMs > staleAfterMs); - if (Number.isFinite(occurrenceMs) - && Date.now() - occurrenceMs > staleAfterMs - && notOnFreshRetryChain) { + // 过期收场本身。拆成一个函数是因为除了下面这道守卫,推迟(defer)把唤醒 + // 时刻推过了过期线时也走它,见 handleDeferred。 + const settleAsStale = async () => { try { // hook 的 client_state 读写口:过期跳过往往正是宿主要留一条痕迹的时 // 候(服务停摆恢复后的第一跳,此前这个 tick 里可能一次 fire 都没跑 @@ -1061,9 +1071,60 @@ async function deliverTasks(ctx, tasks) { results.failedCount++; results.failedTasks.push({ taskId: task.id, reason: error.message || '过期任务处理失败', status: 'stale_update_failed' }); } + }; + + if (Number.isFinite(occurrenceMs) + && Date.now() - occurrenceMs > staleAfterMs + && notOnFreshRetryChain) { + await settleAsStale(); return; } + /** + * onBeforeFire 返回了 `{ defer: { afterMs } }`:这次不发,到 retryAfter 再 + * 来问。 + * + * 行上只动两个字段:`retry_after` 写成唤醒时刻(捞取条件会滤掉没到点的 + * 行,到点自然放行),`lease_until` 置空把租约放掉。其余一概不碰—— + * `retry_count` 不涨(这不是失败,不占重试次数)、`last_error` 不写也不清、 + * `next_send_at` 保持名义时刻(它是这次触发的身份,也是循环推进和过期判定 + * 的基准)、`status` 仍是 pending,循环任务不推进。这一跳既不计成功也不计 + * 失败,记在 details.deferredTasks 里。 + * + * 过期线照常管着它:唤醒时刻离名义时刻超过 staleAfterMs 的话,这次推迟 + * 不落库,当场按过期收场(一次性任务标 failed、循环任务快进,调 + * onStaleSkip)。对没失败过的任务,这跟「先写下推迟、下次捞起来再被上面那 + * 道守卫判过期」是同一个结局,只是不白等一轮;对正在重试链上的任务 + * (retry_count > 0)则是唯一的出口——上面那道守卫看的是 retry_after 新不 + * 新,而每推迟一次都会把它刷新,光靠那道守卫的话这条任务可以被无限期推 + * 下去。 + * + * @param {string} retryAfter - 唤醒时刻(ISO),runAgenticFire 算好的 + */ + const handleDeferred = async (retryAfter) => { + const wakeMs = Date.parse(retryAfter); + if (Number.isFinite(occurrenceMs) && wakeMs - occurrenceMs > staleAfterMs) { + await settleAsStale(); + return; + } + try { + const updated = await updateTaskWithLastError(task.id, { + retry_after: retryAfter, + lease_until: null + }); + if (rowVanished(updated)) { + await recordCancelled(task, 'cancelled_mid_delivery'); + return; + } + results.deferredTasks.push({ taskId: task.id, retryAfter }); + } catch (error) { + // 推迟没写进去:行还是 pending,租约到期后下一跳会把它重新捞起来, + // 等于提前再问一次 onBeforeFire。 + results.failedCount++; + results.failedTasks.push({ taskId: task.id, reason: error.message || '推迟写库失败', status: 'defer_update_failed' }); + } + }; + let sendResult; try { // 预扫描解好的 payload 一并递过去,投递侧不再解第二遍。 @@ -1078,6 +1139,10 @@ async function deliverTasks(ctx, tasks) { webpush: guardWebpushWithLease(ctx.webpush, lease), isTaskCancelled: () => lease.lost, signal: lease.signal, + // onBeforeFire 的 { defer } 要靠 retry_after 这一列落地,没实现 + // claimTask 的适配器没有这一列(见 lib/agentic-fire.js 的 + // AGENTIC_DEFER_UNSUPPORTED)。 + _deferSupported: supportsClaim, }, masterKey, { userKey, payload: decryptedPayload } @@ -1099,6 +1164,17 @@ async function deliverTasks(ctx, tasks) { return; } + // 宿主说「现在不合适,过一会儿再来问」。排在失败分支前面:推迟的结果不带 + // success,但它不是失败。 + if (sendResult.deferred) { + if (lease.lost) { + await recordCancelled(task, 'cancelled_mid_delivery'); + return; + } + await handleDeferred(sendResult.retryAfter); + return; + } + if (!sendResult.success) { // 取消是拦下来的,不是发失败——按失败走会给一条已经不存在的行排重试, // 也会把这件事混进 failedTasks 里。 @@ -1232,6 +1308,9 @@ async function deliverTasks(ctx, tasks) { deletedOnceOffTasks: results.deletedOnceOffTasks, updatedRecurringTasks: results.updatedRecurringTasks, staleTasks: results.staleTasks, + // onBeforeFire 返回 { defer } 被推迟的任务({ taskId, retryAfter }):这一 + // 跳没生成也没发,到 retryAfter 之后再问。不计入 successCount / failedCount。 + deferredTasks: results.deferredTasks, // 投递期间行被取消 / 顶替的任务。`cancelled_mid_delivery` = 推送在发出去 // 之前被拦下;`cancelled_after_delivery` = 推送已经发完,收尾写库才发现 // 行没了。两种都不计入 successCount / failedCount。 diff --git a/packages/rei-standard-amsg/server/src/server/single-user.js b/packages/rei-standard-amsg/server/src/server/single-user.js index 06547a3..6b8f02f 100644 --- a/packages/rei-standard-amsg/server/src/server/single-user.js +++ b/packages/rei-standard-amsg/server/src/server/single-user.js @@ -21,6 +21,10 @@ * @param {Object} [config.hooks] - optional fire-time hooks (see lib/agentic-fire.js): * { onBeforeFire, onLLMOutput, executeToolCalls }. When omitted, AI tasks * replay the schedule-time frozen prompt (legacy behavior, unchanged). + * @param {(fireCtx: Object) => import('./lib/agentic-fire.js').BeforeFireResult | Promise} [config.hooks.onBeforeFire] + * 返回消息数组或 `{ messages, ... }` 继续生成;`{ skip: true }` 这次不发; + * `{ defer: { afterMs } }` 现在不合适、afterMs 毫秒后再来问(不占重试次数, + * 只在 runScheduledTick / runTask 投递时可用);`null` 交还给冻结 prompt。 * @param {number} [config.maxToolIterations] - factory default LLM-round cap for the agentic loop (default 5). * @param {number} [config.totalTimeoutMs] - factory default wall-time ceiling for the agentic loop (default 240000). * @param {number} [config.maxStateValueBytes] - client_state 单条 value 的总上限(默认 5MB)。超过 200KB 的值由服务端透明分块存储(见 lib/state-chunks.js)。 @@ -38,11 +42,13 @@ * 上抛之前调用完。hook 自身抛错只记日志,不影响主流程(见 * lib/agentic-fire.js)。 * @param {function} [config.onFireSettled] - 一次 fire 收尾的可选 hook: - * ({ task, status, skipReason, sentCount, pushedCount, total, iterations, - * error, metadata, usage, usageTotal, llmCalls, outboxed, scratch, readState, - * writeState }) => void|Promise。onBeforeFire 被调用过 - * 就一定会调一次,无论这次是发完(status 'sent')、跳过('skipped')、抛错 - * ('failed')还是交还给冻结 prompt 老链路('not-handled')。onAfterSend 只 + * ({ task, status, skipReason, retryAfter, sentCount, pushedCount, total, + * iterations, error, metadata, usage, usageTotal, llmCalls, outboxed, scratch, + * readState, writeState }) => void|Promise。onBeforeFire 被调用过 + * 就一定会调一次,无论这次是发完(status 'sent')、跳过('skipped')、推迟 + * ('deferred',retryAfter 是什么时候再来问)、抛错('failed')、被取消 + * ('cancelled')还是交还给冻结 prompt 老链路('not-handled')。status 的 + * 类型见 lib/agentic-fire.js 的 FireSettledStatus。onAfterSend 只 * 走「有 push 要发」那条路,「开始时占点什么、结束时放掉」的写法挂这个才不 * 会漏(见 lib/agentic-fire.js)。 * @returns {{ handlers: Object, ctx: Object }} diff --git a/packages/rei-standard-amsg/server/test/before-fire-defer.test.mjs b/packages/rei-standard-amsg/server/test/before-fire-defer.test.mjs new file mode 100644 index 0000000..ac6127b --- /dev/null +++ b/packages/rei-standard-amsg/server/test/before-fire-defer.test.mjs @@ -0,0 +1,435 @@ +/** + * onBeforeFire 的 `{ defer: { afterMs } }`:现在不合适,过一会儿再来问。 + * + * 钉住的是任务行在推迟前后的样子(只动 retry_after 和租约)、这一跳不生成不 + * 推送、到点之后从 onBeforeFire 重新走一遍,以及过期线照常管着被推迟的任务。 + */ +import { test } from 'node:test'; +import assert from 'node:assert/strict'; +import { runScheduledTick, runTask } from '../src/server/lib/run-tick.js'; +import { processMessagesByUuid } from '../src/server/lib/message-processor.js'; +import { MAX_DEFER_AFTER_MS } from '../src/server/lib/agentic-fire.js'; +import { createD1Adapter } from '../src/server/adapters/d1.js'; +import { createTestD1 } from './helpers/sqlite-d1.mjs'; +import { deriveUserEncryptionKey, encryptForStorage } from '../src/server/lib/encryption.js'; +import { seedPushSubscription } from './helpers/push-subscription.mjs'; + +const USER = '550e8400-e29b-41d4-a716-446655440000'; +const MASTER_KEY = 'a'.repeat(64); +const VAPID = { email: 'mailto:x@example.com', publicKey: 'pub', privateKey: 'priv' }; + +const SECOND = 1000; +const MINUTE = 60 * SECOND; +const DAY = 24 * 60 * MINUTE; + +function isoFromNow(ms) { + return new Date(Date.now() + ms).toISOString(); +} + +/** + * 一套现成的场景:一条刚到点的 LLM 任务 + 记账用的 hook。 + * + * onBeforeFire 的返回值由 `beforeFire(fireCtx, callIndex)` 决定;LLM(fetch)、 + * onLLMOutput、onAfterSend、推送被调到了几次都记下来,推迟的那一跳它们应当全 + * 是 0。 + */ +async function fixture(t, { uuid = 'defer-task', recurrenceType = 'none', nextSendAt, beforeFire, messageType = 'auto', ctx = {} } = {}) { + const d1 = createTestD1(); + const adapter = createD1Adapter(d1); + await adapter.initSchema(); + await seedPushSubscription(adapter, USER, MASTER_KEY); + const userKey = await deriveUserEncryptionKey(USER, MASTER_KEY); + const dueAt = nextSendAt || isoFromNow(-30 * SECOND); + await adapter.createTask({ + user_id: USER, + uuid, + encrypted_payload: await encryptForStorage(JSON.stringify({ + contactName: 'Rei', + messageType, + recurrenceType, + apiUrl: 'https://example.com/v1/chat/completions', + apiKey: 'key', + primaryModel: 'test', + completePrompt: 'frozen prompt', + metadata: { charId: 'c1' }, + }), userKey), + next_send_at: dueAt, + message_type: messageType, + }); + + const calls = { beforeFire: [], llm: 0, llmOutput: 0, afterSend: 0, settled: [], stale: [] }; + t.mock.method(globalThis, 'fetch', async () => { + calls.llm++; + return Response.json({ choices: [{ message: { role: 'assistant', content: 'hello' } }] }); + }); + const webpush = { sent: [], async sendNotification(_sub, payload) { webpush.sent.push(payload); } }; + + const tickCtx = { + db: adapter, + masterKey: MASTER_KEY, + vapid: VAPID, + webpush, + // 心跳用的是真定时器,这些用例里用不上。 + leaseHeartbeatMs: 0, + hooks: { + onBeforeFire: async (fireCtx) => { + calls.beforeFire.push({ nextSendAt: fireCtx.task.nextSendAt, retryCount: fireCtx.task.retryCount }); + return beforeFire(fireCtx, calls.beforeFire.length - 1); + }, + onLLMOutput: async () => { calls.llmOutput++; return { decision: 'skip-push' }; }, + }, + onAfterSend: async () => { calls.afterSend++; }, + onFireSettled: async (info) => { calls.settled.push(info); }, + onStaleSkip: async (task, info) => { calls.stale.push({ uuid: task.uuid, info }); }, + ...ctx, + }; + + /** 行的全部列(含 lease_until / last_error,投递用的列集里没有这两个)。 */ + const row = () => d1.prepare('SELECT * FROM scheduled_messages WHERE uuid = ?').bind(uuid).first(); + + return { adapter, tickCtx, calls, webpush, row, dueAt }; +} + +/** 把 Date 换成可拨的钟(定时器不动),返回「往后拨 ms」的函数。 */ +function controllableClock(t) { + let nowMs = Date.now(); + t.mock.timers.enable({ apis: ['Date'], now: nowMs }); + return (ms) => { + nowMs += ms; + t.mock.timers.setTime(nowMs); + }; +} + +test('defer:只写 retry_after 并放掉租约,不调 LLM、不推送', async (t) => { + const f = await fixture(t, { beforeFire: () => ({ defer: { afterMs: 45 * SECOND } }) }); + + const before = Date.now(); + const res = await runScheduledTick(f.tickCtx); + const after = Date.now(); + + const row = await f.row(); + assert.equal(row.status, 'pending'); + assert.equal(row.next_send_at, f.dueAt, '名义触发时刻不能动'); + assert.equal(row.retry_count, 0, '推迟不占重试次数'); + assert.equal(row.last_error, null, '推迟不留报错'); + assert.equal(row.lease_until, null, '租约要放掉'); + const retryAfterMs = Date.parse(row.retry_after); + assert.ok(retryAfterMs >= before + 45 * SECOND && retryAfterMs <= after + 45 * SECOND, + `retry_after 应当是 now + afterMs,实际 ${row.retry_after}`); + + assert.equal(f.calls.llm, 0); + assert.equal(f.calls.llmOutput, 0); + assert.equal(f.calls.afterSend, 0); + assert.equal(f.webpush.sent.length, 0); + + assert.equal(res.successCount, 0); + assert.equal(res.failedCount, 0); + assert.deepEqual(res.details.deferredTasks, [{ taskId: row.id, retryAfter: row.retry_after }]); + assert.equal(res.details.deletedOnceOffTasks, 0); + assert.deepEqual(res.details.failedTasks, []); +}); + +test('defer:收尾回执是 deferred,带 retryAfter', async (t) => { + const f = await fixture(t, { beforeFire: () => ({ defer: { afterMs: 45 * SECOND } }) }); + await runScheduledTick(f.tickCtx); + + assert.equal(f.calls.settled.length, 1, 'onBeforeFire 调过一次,收尾也只调一次'); + const info = f.calls.settled[0]; + assert.equal(info.status, 'deferred'); + assert.equal(info.retryAfter, (await f.row()).retry_after); + assert.equal(info.skipReason, null); + assert.equal(info.error, null); + assert.equal(info.willRetry, null); + assert.equal(info.failureStage, null); + assert.equal(info.llmCalls, 0); + assert.equal(info.iterations, 0); + assert.equal(info.sentCount, 0); + assert.equal(info.outboxed, false); + assert.deepEqual(info.metadata, { charId: 'c1' }); +}); + +test('别的结局的收尾回执上 retryAfter 是 null', async (t) => { + const f = await fixture(t, { beforeFire: () => ({ skip: true }) }); + await runScheduledTick(f.tickCtx); + assert.equal(f.calls.settled[0].status, 'skipped'); + assert.equal(f.calls.settled[0].retryAfter, null); +}); + +test('defer:已有的 retry_count 和 last_error 原样留着', async (t) => { + const f = await fixture(t, { beforeFire: () => ({ defer: { afterMs: 45 * SECOND } }) }); + const lastError = JSON.stringify({ at: isoFromNow(-3 * MINUTE), reason: '上一次的失败' }); + await f.adapter.updateTaskById((await f.row()).id, { + retry_count: 2, + retry_after: isoFromNow(-SECOND), + last_error: lastError, + }); + + await runScheduledTick(f.tickCtx); + + const row = await f.row(); + assert.equal(row.retry_count, 2); + assert.equal(row.last_error, lastError); + assert.equal(row.status, 'pending'); + assert.ok(Date.parse(row.retry_after) > Date.now()); + assert.equal(f.calls.beforeFire[0].retryCount, 2); +}); + +test('defer:到 retry_after 之前不被捞,之后从 onBeforeFire 重新走一遍', async (t) => { + const advance = controllableClock(t); + const f = await fixture(t, { + beforeFire: (_ctx, callIndex) => (callIndex < 2 ? { defer: { afterMs: 2 * MINUTE } } : { skip: true }), + }); + + await runScheduledTick(f.tickCtx); + assert.equal(f.calls.beforeFire.length, 1); + + // 还没到点:这一跳捞不到它,hook 不会被调。 + advance(MINUTE); + const idle = await runScheduledTick(f.tickCtx); + assert.equal(idle.totalTasks, 0); + assert.equal(f.calls.beforeFire.length, 1); + + // 到点:重新问一次,这次又推迟。 + advance(MINUTE + SECOND); + const second = await runScheduledTick(f.tickCtx); + assert.equal(f.calls.beforeFire.length, 2); + assert.equal(second.details.deferredTasks.length, 1); + assert.equal((await f.row()).retry_count, 0); + + // 再到点:这次放行(skip 收场),一次性任务删掉。 + advance(2 * MINUTE + SECOND); + const third = await runScheduledTick(f.tickCtx); + assert.equal(f.calls.beforeFire.length, 3); + assert.equal(third.successCount, 1); + assert.equal(await f.row(), null); + + // 三次问到的是同一次触发:名义时刻、重试次数都没变。 + assert.deepEqual(f.calls.beforeFire, Array(3).fill({ nextSendAt: f.dueAt, retryCount: 0 })); + assert.deepEqual(f.calls.settled.map((s) => s.status), ['deferred', 'deferred', 'skipped']); + assert.equal(f.calls.llm, 0); +}); + +test('defer:循环任务不推进到下一次', async (t) => { + const advance = controllableClock(t); + const f = await fixture(t, { + recurrenceType: 'daily', + beforeFire: (_ctx, callIndex) => (callIndex === 0 ? { defer: { afterMs: MINUTE } } : { skip: true }), + }); + + const res = await runScheduledTick(f.tickCtx); + assert.equal(res.details.updatedRecurringTasks, 0); + assert.equal((await f.row()).next_send_at, f.dueAt); + + // 到点后放行,这才推进到下一次,推迟留下的 retry_after 也清掉。 + advance(MINUTE + SECOND); + const next = await runScheduledTick(f.tickCtx); + assert.equal(next.details.updatedRecurringTasks, 1); + const row = await f.row(); + assert.equal(row.next_send_at, new Date(Date.parse(f.dueAt) + DAY).toISOString()); + assert.equal(row.retry_after, null); +}); + +test('defer 之后名义时刻过了 staleAfterMs:下次捞起来直接判过期,不再问 onBeforeFire', async (t) => { + const advance = controllableClock(t); + const f = await fixture(t, { beforeFire: () => ({ defer: { afterMs: 10 * MINUTE } }) }); + + await runScheduledTick(f.tickCtx); + assert.equal(f.calls.beforeFire.length, 1); + + advance(61 * MINUTE); + const res = await runScheduledTick(f.tickCtx); + + assert.equal(f.calls.beforeFire.length, 1, '过期守卫排在 onBeforeFire 前面'); + assert.deepEqual(res.details.staleTasks.map((s) => s.action), ['expired']); + const row = await f.row(); + assert.equal(row.status, 'failed'); + assert.equal(JSON.parse(row.last_error).reason, 'stale'); + assert.equal(f.calls.stale.length, 1); +}); + +test('反复 defer:唤醒时刻越过过期线的那一次当场判过期', async (t) => { + const advance = controllableClock(t); + const f = await fixture(t, { beforeFire: () => ({ defer: { afterMs: 10 * MINUTE } }) }); + + // 名义时刻是 30 秒前、过期线 60 分钟。每 10 分钟被问一次:前五次的唤醒时刻 + // 都在线内,第六次(第 50 分钟问,要推到第 60 分钟)越线。 + for (let i = 0; i < 5; i++) { + const res = await runScheduledTick(f.tickCtx); + assert.equal(res.details.deferredTasks.length, 1, `第 ${i + 1} 次应当照常推迟`); + advance(10 * MINUTE + SECOND); + } + assert.equal((await f.row()).status, 'pending'); + + const last = await runScheduledTick(f.tickCtx); + assert.equal(f.calls.beforeFire.length, 6); + assert.deepEqual(last.details.deferredTasks, []); + assert.deepEqual(last.details.staleTasks.map((s) => s.action), ['expired']); + + const row = await f.row(); + assert.equal(row.status, 'failed'); + assert.equal(row.retry_count, 0); + assert.equal(JSON.parse(row.last_error).reason, 'stale'); + assert.equal(f.calls.stale.length, 1); + assert.equal(f.calls.stale[0].info.action, 'expired'); + assert.equal(f.calls.llm, 0); + + // 之后不会再被捞起来。 + advance(10 * MINUTE); + assert.equal((await runScheduledTick(f.tickCtx)).totalTasks, 0); + assert.equal(f.calls.beforeFire.length, 6); +}); + +// 重试链上的任务(retry_count > 0)不看名义时刻,只看 retry_after 新不新;而 +// 每次推迟都会把 retry_after 刷新。推迟自己不把过期线带上的话,这种任务可以 +// 被无限期推下去。 +test('反复 defer:已经在重试链上的任务一样会过期', async (t) => { + const advance = controllableClock(t); + const f = await fixture(t, { beforeFire: () => ({ defer: { afterMs: 10 * MINUTE } }) }); + await f.adapter.updateTaskById((await f.row()).id, { retry_count: 1, retry_after: isoFromNow(-SECOND) }); + + let ticks = 0; + let res; + do { + res = await runScheduledTick(f.tickCtx); + ticks++; + advance(10 * MINUTE + SECOND); + } while (res.details.deferredTasks.length === 1 && ticks < 20); + + assert.equal(ticks, 6, '第六次的唤醒时刻越线'); + assert.deepEqual(res.details.staleTasks.map((s) => s.action), ['expired']); + assert.equal((await f.row()).status, 'failed'); + assert.equal(f.calls.stale.length, 1); +}); + +test('反复 defer:循环任务越线时快进到下一次', async (t) => { + const advance = controllableClock(t); + const f = await fixture(t, { recurrenceType: 'daily', beforeFire: () => ({ defer: { afterMs: 20 * MINUTE } }) }); + + let res; + for (let i = 0; i < 3; i++) { + res = await runScheduledTick(f.tickCtx); + advance(20 * MINUTE + SECOND); + } + + assert.deepEqual(res.details.staleTasks.map((s) => s.action), ['fast_forwarded']); + const row = await f.row(); + assert.equal(row.status, 'pending'); + assert.equal(row.next_send_at, new Date(Date.parse(f.dueAt) + DAY).toISOString()); + assert.equal(row.retry_after, null); + assert.equal(f.calls.stale[0].info.action, 'fast_forwarded'); +}); + +for (const [label, defer] of [ + ['0', { afterMs: 0 }], + ['负数', { afterMs: -1000 }], + ['NaN', { afterMs: NaN }], + ['Infinity', { afterMs: Infinity }], + ['字符串', { afterMs: '1000' }], + ['缺 afterMs', {}], + ['defer 不是对象', true], + ['超过上限', { afterMs: MAX_DEFER_AFTER_MS + 1 }], +]) { + test(`非法的 defer(${label}):按 AGENTIC_BAD_BEFORE_FIRE 一跳终审`, async (t) => { + const f = await fixture(t, { beforeFire: () => ({ defer }) }); + const res = await runScheduledTick(f.tickCtx); + + assert.equal(res.failedCount, 1); + assert.deepEqual(res.details.deferredTasks, []); + assert.equal(res.details.failedTasks[0].status, 'permanently_failed'); + assert.match(res.details.failedTasks[0].reason, /AGENTIC_BAD_BEFORE_FIRE/); + assert.match(res.details.failedTasks[0].reason, /\{ defer: \{ afterMs \} \}/, '合法返回值列表里要有 defer'); + + const row = await f.row(); + assert.equal(row.status, 'failed'); + assert.equal(JSON.parse(row.last_error).errorCode, 'AGENTIC_BAD_BEFORE_FIRE'); + + assert.equal(f.calls.settled.length, 1); + assert.equal(f.calls.settled[0].status, 'failed'); + assert.equal(f.calls.settled[0].error.permanent, true); + assert.equal(f.calls.settled[0].retryAfter, null); + assert.equal(f.calls.llm, 0); + }); +} + +test('afterMs 正好等于上限:照常推迟', async (t) => { + // 过期线放宽到两天,免得唤醒时刻先撞上它。 + const f = await fixture(t, { + beforeFire: () => ({ defer: { afterMs: MAX_DEFER_AFTER_MS } }), + ctx: { staleAfterMs: 2 * DAY }, + }); + const res = await runScheduledTick(f.tickCtx); + assert.equal(res.details.deferredTasks.length, 1); + assert.equal(MAX_DEFER_AFTER_MS, DAY); +}); + +test('runTask:defer 后回 ran: true,到点之前再调回 retry_pending', async (t) => { + const f = await fixture(t, { beforeFire: () => ({ defer: { afterMs: 45 * SECOND } }) }); + + const first = await runTask(f.tickCtx, 'defer-task'); + assert.equal(first.ran, true); + assert.equal(first.summary.details.deferredTasks.length, 1); + + const row = await f.row(); + const second = await runTask(f.tickCtx, 'defer-task'); + assert.deepEqual(second, { ran: false, reason: 'retry_pending', retryAfter: row.retry_after }); + assert.equal(f.calls.beforeFire.length, 1); +}); + +test('defer 期间任务被取消:不把行写回来,记成取消', async (t) => { + const f = await fixture(t, { + beforeFire: async () => { + assert.equal(await f.adapter.deleteTaskByUuid('defer-task', USER), true); + return { defer: { afterMs: 45 * SECOND } }; + }, + }); + + const res = await runScheduledTick(f.tickCtx); + + assert.equal(await f.row(), null); + assert.deepEqual(res.details.deferredTasks, []); + assert.deepEqual(res.details.cancelledTasks.map((c) => c.status), ['cancelled_mid_delivery']); + assert.equal(res.failedCount, 0); +}); + +// 没有 retry_after 列可写的适配器只能改 next_send_at,而那是这次触发的名义时 +// 刻。所以这里不推迟,明确报配置错误(可重试,不判终态)。 +test('没实现 claimTask 的适配器:defer 报 AGENTIC_DEFER_UNSUPPORTED', async (t) => { + const f = await fixture(t, { beforeFire: () => ({ defer: { afterMs: 45 * SECOND } }) }); + const db = new Proxy(f.adapter, { + get(target, prop) { + if (prop === 'claimTask') return undefined; + const value = target[prop]; + return typeof value === 'function' ? value.bind(target) : value; + }, + }); + + const res = await runScheduledTick({ ...f.tickCtx, db }); + + assert.equal(res.failedCount, 1); + assert.deepEqual(res.details.deferredTasks, []); + assert.match(res.details.failedTasks[0].reason, /AGENTIC_DEFER_UNSUPPORTED/); + assert.ok(res.details.failedTasks[0].nextRetryAt, '配置错误留在退避阶梯上'); + const row = await f.row(); + assert.equal(row.status, 'pending'); + assert.equal(row.retry_count, 1); + assert.equal(f.calls.settled[0].status, 'failed'); + assert.equal(f.calls.settled[0].error.code, 'AGENTIC_DEFER_UNSUPPORTED'); + assert.notEqual(f.calls.settled[0].error.permanent, true); + assert.equal(f.calls.llm, 0); +}); + +test('请求内当场投递的 instant 任务:defer 报 AGENTIC_DEFER_UNSUPPORTED', async (t) => { + const f = await fixture(t, { + messageType: 'instant', + beforeFire: () => ({ defer: { afterMs: 45 * SECOND } }), + }); + + const result = await processMessagesByUuid('defer-task', f.tickCtx, 0, USER, MASTER_KEY); + + assert.equal(result.success, false); + assert.match(result.error.message, /AGENTIC_DEFER_UNSUPPORTED/); + assert.equal((await f.row()).status, 'failed', '不能当成发完了把任务删掉,也不能悄悄留着'); + assert.equal(f.calls.settled[0].status, 'failed'); + assert.equal(f.calls.llm, 0); + assert.equal(f.webpush.sent.length, 0); +}); diff --git a/packages/rei-standard-amsg/server/test/capabilities.test.mjs b/packages/rei-standard-amsg/server/test/capabilities.test.mjs index d844c6c..63f62b8 100644 --- a/packages/rei-standard-amsg/server/test/capabilities.test.mjs +++ b/packages/rei-standard-amsg/server/test/capabilities.test.mjs @@ -66,6 +66,7 @@ const EXPECTED_FEATURES = [ 'hook-usage-total', 'max-delivery-retries', 'max-generation-retries', + 'before-fire-defer', ]; function makeWorker(extra = {}) {