Skip to content
Merged
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
16 changes: 16 additions & 0 deletions .changeset/amsg-server-before-fire-defer.md
Original file line number Diff line number Diff line change
@@ -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`。
51 changes: 49 additions & 2 deletions packages/rei-standard-amsg/server/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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 上的口子 | 干什么 |
Expand All @@ -383,14 +393,49 @@ 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` 里):

| 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 都没跑过,而那正是它要留痕迹的时候。
Expand All @@ -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` 的每种结局(发完、跳过、失败)都带——失败也花了钱:
Expand Down Expand Up @@ -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)。

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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 的
Expand Down Expand Up @@ -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) {
Expand All @@ -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 管 |

Expand Down
2 changes: 2 additions & 0 deletions packages/rei-standard-amsg/server/src/server/cloudflare.js
Original file line number Diff line number Diff line change
Expand Up @@ -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';
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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 秒内就能
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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) {
Expand Down
2 changes: 2 additions & 0 deletions packages/rei-standard-amsg/server/src/server/index.js
Original file line number Diff line number Diff line change
Expand Up @@ -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 这一份,
// 别手抄第二份)。
Expand Down
Loading
Loading