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
5 changes: 5 additions & 0 deletions .changeset/amsg-server-generation-retry-policy.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,5 @@
---
"@rei-standard/amsg-server": minor
---

Add per-task `maxGenerationRetries` configuration and `onFireSettled.willRetry` / `failureStage` receipts. Interactive requests can stop after their first generation failure without provider-specific error codes, while committed outbox batches retain delivery-only retries. Existing tasks and default retry behavior remain compatible.
20 changes: 20 additions & 0 deletions packages/rei-standard-amsg/server/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -703,6 +703,26 @@ LLM 这一条只认上游**答复了、并且拒了**的那几种状态码:Key

`GET /capabilities` 的 features 里对应 `llm-permanent-errors`、`redeliver-committed-batch`、`hook-usage-total`、`max-delivery-retries`。

## 按任务限制生成重试(`maxGenerationRetries`)

用户主动触发的回复可以第一次生成失败就结束,定时任务则保留默认退避。工厂配置和 `runScheduledTick` / `runTask` 支持 `maxGenerationRetries`,值可以是非负整数,也可以是同步函数:

```js
createSingleUserCloudflareWorker((env) => ({
// 这里的 interactive 是宿主自己约定的 metadata,库不认识业务类型。
maxGenerationRetries: (task) => task.metadata?.interactive ? 0 : undefined,
// ...其余配置
}));
```

回调收到与 `serializeBy` 相同的安全任务视图(含 metadata,不含 API Key 等凭据),每次尝试解析一次。返回 `undefined` 继承 `maxDeliveryRetries`(默认 3),返回 0 禁止本轮失败后重新生成;显式数值单独控制生成重试上限。不配置时行为不变,不需要修改任务 payload 或迁移数据库,已有任务也生效。所有服务端工厂和旧的 UUID 即时入口都透传这个配置。

这里的“生成阶段”以**最终回复整批进入 outbox**为边界:前置 hook、模型、工具循环、组装和落库失败都使用生成上限;没有 outbox 的适配器发生推送失败时仍可能需要重新生成,因此也使用生成上限。完整批次已入 outbox 后的推送失败改用 `maxDeliveryRetries`,补投不会再次调用模型。hook 的独立 `emitResult` 不等于最终回复提交。

`onFireSettled` 新增 `willRetry` 与 `failureStage`:失败时前者为本次的重试决策,后者为 `'generation'` 或 `'delivery'`;成功、跳过、未接管时两者均为 `null`。宿主可以用 `info.status === 'failed' && info.willRetry === false && !info.outboxed` 即时显示生成失败,无需修改错误对象或猜测重试次数。此回执发生在调度器写任务状态之前,`willRetry` 描述同源策略决策,不保证后续数据库写入成功;取消导致的失败不会重试。已有永久性错误判定仍然生效。

策略回调抛错或返回非法值时,不调用生成 hook / 模型,以 `GENERATION_RETRY_POLICY_INVALID` 记录配置错误并沿默认投递上限退避,避免一次坏部署立即作废所有任务;此时尚未进入 fire,因此不触发 `onFireSettled`。

## 同一分组的任务不并发(`serializeBy`)

同一个角色可能有好几条定时任务。撞在一起并发跑的话,用户一口气收到两条互不知情的消息;宿主在 hook 里维护的「我刚才说过什么」台账通常是读进内存 → 改 → 整份写回,两条各改各的再写回,后写的必然盖掉前面那条。
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -326,6 +326,7 @@ export function createSingleUserCloudflareWorker(buildConfig, options = {}) {
onFireSettled: cfg.onFireSettled,
// 一次触发投递失败后最多再重试几次(默认 3,0 = 第一次失败就终审;见
// lib/run-tick.js 的 DEFAULT_MAX_DELIVERY_RETRIES)。
maxGenerationRetries: cfg.maxGenerationRetries,
maxDeliveryRetries: cfg.maxDeliveryRetries,
// 分组串行:(task) => 分组标识 | null。同一分组的任务同时只跑一条,
// 跨跳也算(见 lib/run-tick.js)。不配 = 全并发,与以前一致。
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -125,6 +125,8 @@ export const SERVER_FEATURES = Object.freeze([
'hook-usage-total',
// 工厂配置认 maxDeliveryRetries(投递失败的重试次数上限,默认 3)。
'max-delivery-retries',
// Per-task pre-commit retry limit and fire receipt retry decision.
'max-generation-retries',
]);

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 @@ -62,6 +62,7 @@ import { normalizeVapidSubject } from '@rei-standard/amsg-shared';
* 即可)。一条 push 装不下的思考过程要切片发,切多大、最多几片、重组窗口多长由
* 接收端说了算——发送端不知道这份配置的话,切出来的分片到了那边会被逐片拒收,
* 或者整批没能在重组窗口内发完,一条也拼不回来。不配 = 两边都用默认值。
* @property {number | ((task: Object) => number | undefined)} [maxGenerationRetries] - Retry limit before a complete batch is committed; defaults to maxDeliveryRetries.
* @property {number} [maxDeliveryRetries] - 定时任务一次触发投递失败后最多再重试
* 几次(默认 3,0 = 第一次失败就终审)。只管 `/send-notifications` 的退避阶梯;
* `messageType: 'instant'` 的请求内重试不受它影响。
Expand Down Expand Up @@ -144,6 +145,7 @@ export async function createReiServer(config) {
multipart: config.multipart || null,
// 定时任务投递失败后的重试次数上限(默认 3),send-notifications 展开 ctx
// 时带进 runScheduledTick。
maxGenerationRetries: config.maxGenerationRetries,
maxDeliveryRetries: config.maxDeliveryRetries,
tenant: {
initSecret
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -113,6 +113,7 @@
* semantics apply.
*/

import { failureRetryDecision } from './retry-policy.js';
import {
assertValidDecision,
buildSessionContext,
Expand Down Expand Up @@ -791,6 +792,9 @@ export async function runAgenticFire({ task, decryptedPayload, userKey, ctx }) {
await notifyFireSettled(ctx, {
task,
status: settledStatus,
...(settledStatus === 'failed' && ctx._deliveryState
? failureRetryDecision(ctx._deliveryState, settledError)
: { willRetry: null, failureStage: null }),
skipReason: settledStatus === 'skipped' ? progress.skipReason : null,
sentCount: progress.sentCount,
pushedCount: progress.pushedCount,
Expand Down Expand Up @@ -1259,6 +1263,7 @@ async function sendHookPushPayloads({
}
outboxed = await appendPushesToOutbox({ db: ctx.db, userId: task.user_id, userKey, pushes: finalized });
progress.outboxed = outboxed;
if (ctx._deliveryState) ctx._deliveryState.outboxed = outboxed;

if (!ctx.vapid || !ctx.vapid.email || !ctx.vapid.publicKey || !ctx.vapid.privateKey) {
throw new Error('VAPID configuration missing - push notifications cannot be sent');
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,7 @@
* 分片发,sw 收齐后还原。
*/

import { resolveMaxDeliveryRetries, resolveMaxGenerationRetries, failureRetryDecision } from './retry-policy.js';
import { randomUUID } from './webcrypto-utils.js';
import {
buildContentPush,
Expand All @@ -38,7 +39,7 @@ import { MAX_PUSH_PAYLOAD_BYTES, measurePushPayload } from './webpush-webcrypto.

import { decryptFromStorage, deriveUserEncryptionKey } from './encryption.js';
import { callLlm } from './llm.js';
import { runAgenticFire, taskNeedsLlm, occurrenceSuffix, occurrenceMsOf, stampTaskIdentity } from './agentic-fire.js';
import { buildHookTask, runAgenticFire, taskNeedsLlm, occurrenceSuffix, occurrenceMsOf, stampTaskIdentity } from './agentic-fire.js';
import { resolvePushSubscription } from './push-subscription-store.js';
import { hasChatCredRef, resolveFireCredentials } from './llm-credentials-store.js';
import {
Expand Down Expand Up @@ -344,6 +345,8 @@ function positiveIntegerOr(value, fallback) {

/**
* @typedef {Object} ProcessorContext
* @property {number} [maxDeliveryRetries]
* @property {number | ((task: Object) => number | undefined)} [maxGenerationRetries]
* @property {Object} webpush - The web-push module instance (already VAPID-configured).
* @property {Object} vapid - { email, publicKey, privateKey }
* @property {import('../adapters/interface.js').DbAdapter} db
Expand Down Expand Up @@ -429,10 +432,19 @@ 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 }>}
* @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 }>}
* 失败时 `pushStatusCode` 是推送服务回的 HTTP 状态码(不是推送阶段炸的 → null)。
*/
export async function processSingleMessage(task, ctx, providedMasterKey, predecrypted = null) {
const deliveryState = {
outboxed: false,
deliveryLimit: resolveMaxDeliveryRetries(ctx),
generationLimit: resolveMaxDeliveryRetries(ctx),
retryCount: task.retry_count ?? 0,
isCancelled: ctx.isTaskCancelled,
};
// One attempt owns this state; concurrent tasks never share their progress.
ctx = { ...ctx, _deliveryState: deliveryState };
try {
const masterKey = providedMasterKey || ctx.masterKey;
if (!masterKey) {
Expand All @@ -459,7 +471,11 @@ export async function processSingleMessage(task, ctx, providedMasterKey, predecr
taskUuid: task.uuid,
occurrenceMs: occurrenceMsOf(task),
});
if (committed) return await redeliverCommittedBatch(task, ctx, userKey, decryptedPayload, committed);
if (committed) {
deliveryState.outboxed = true;
return await redeliverCommittedBatch(task, ctx, userKey, decryptedPayload, committed);
}
deliveryState.generationLimit = resolveMaxGenerationRetries(ctx, buildHookTask(task, decryptedPayload));

// Fire-time hooks: when the host configured onBeforeFire and the task
// needs the LLM, offer the agentic path first. onBeforeFire → null
Expand Down Expand Up @@ -599,6 +615,7 @@ export async function processSingleMessage(task, ctx, providedMasterKey, predecr
// 之前的整条——补收走的是 HTTP,没有单条体积上限。
const pushesToSend = reasoningPush ? [reasoningPush, ...contentPushes] : contentPushes;
const outboxed = await appendPushesToOutbox({ db: ctx.db, userId: task.user_id, userKey, pushes: pushesToSend });
deliveryState.outboxed = outboxed;

// VAPID 与订阅排在落行之后:内容已经生成好了,之后哪一步失败都只该重试投递,
// 而「只重试投递」的前提是这一批已经落定在 outbox 里(见
Expand Down Expand Up @@ -695,7 +712,8 @@ export async function processSingleMessage(task, ctx, providedMasterKey, predecr
error: error.message,
errorCode: error.code || null,
pushStatusCode: readPushStatusCode(error),
permanent: isNonRetryableError(error)
permanent: isNonRetryableError(error),
...failureRetryDecision(deliveryState, error)
};
}
}
Expand Down Expand Up @@ -724,7 +742,7 @@ export async function processMessagesByUuid(uuid, ctx, maxRetries = 2, userId, p
};
}

while (retryCount <= maxRetries) {
while (true) {
let task;
try {
task = userId
Expand All @@ -749,7 +767,7 @@ export async function processMessagesByUuid(uuid, ctx, maxRetries = 2, userId, p

// 上一轮生成成功、只是推送失败的话,processSingleMessage 会认出这次触发已
// 经落定的批次,这一轮只补推送、不再把 LLM 跑一遍。
const result = await processSingleMessage(task, ctx, masterKey, null);
const result = await processSingleMessage({ ...task, retry_count: retryCount }, { ...ctx, maxDeliveryRetries: maxRetries }, masterKey, null);

if (!result.success) {
// 确定性失败不进重试:再跑两轮也是同一个错,白让调用方多等、白烧一整轮
Expand All @@ -762,7 +780,7 @@ export async function processMessagesByUuid(uuid, ctx, maxRetries = 2, userId, p
pushStatus: result.pushStatusCode
});

if (!permanent && retryCount < maxRetries) {
if (!permanent && result.willRetry !== false && retryCount < (result.retryLimit ?? maxRetries)) {
retryCount++;
await new Promise(resolve => setTimeout(resolve, 1000 * retryCount));
continue;
Expand Down
40 changes: 40 additions & 0 deletions packages/rei-standard-amsg/server/src/server/lib/retry-policy.js
Original file line number Diff line number Diff line change
@@ -0,0 +1,40 @@
import { DeploymentConfigError, isNonRetryableError, isPermanentDeliveryFailure, isTaskCancelledError, readPushStatusCode } from './errors.js';

/** Number of retries after the first attempt when no policy is configured. */
export const DEFAULT_MAX_DELIVERY_RETRIES = 3;

export function resolveMaxDeliveryRetries(ctx) {
const value = ctx.maxDeliveryRetries;
return Number.isInteger(value) && value >= 0 ? value : DEFAULT_MAX_DELIVERY_RETRIES;
}

/** Resolve once per attempt, before any generation hook or model request. */
export function resolveMaxGenerationRetries(ctx, safeTask) {
let value;
try {
value = typeof ctx.maxGenerationRetries === 'function'
? ctx.maxGenerationRetries(safeTask)
: ctx.maxGenerationRetries;
} catch (cause) {
throw new DeploymentConfigError('maxGenerationRetries callback failed', { code: 'GENERATION_RETRY_POLICY_INVALID', cause });
}
if (value === undefined) return resolveMaxDeliveryRetries(ctx);
if (!Number.isInteger(value) || value < 0) {
throw new DeploymentConfigError('maxGenerationRetries must return a non-negative integer or undefined', { code: 'GENERATION_RETRY_POLICY_INVALID' });
}
return value;
}

/** Shared by fire receipts and both delivery entry points; never alters errors. */
export function failureRetryDecision(state, error) {
const retryLimit = state.outboxed ? state.deliveryLimit : state.generationLimit;
const permanent = isPermanentDeliveryFailure({
permanent: isNonRetryableError(error), errorCode: error?.code,
pushStatus: readPushStatusCode(error),
});
return {
failureStage: state.outboxed ? 'delivery' : 'generation',
retryLimit,
willRetry: !state.isCancelled?.() && !isTaskCancelledError(error) && !permanent && state.retryCount < retryLimit,
};
}
18 changes: 4 additions & 14 deletions packages/rei-standard-amsg/server/src/server/lib/run-tick.js
Original file line number Diff line number Diff line change
Expand Up @@ -59,6 +59,8 @@
* @returns {Promise<Object>} summary { totalTasks, successCount, failedCount, processedAt, executionTime, details }
*/

import { resolveMaxDeliveryRetries } from './retry-policy.js';
export { DEFAULT_MAX_DELIVERY_RETRIES } from './retry-policy.js';
import { hmacSha256, bytesToBase64Url, utf8 } from './webcrypto-utils.js';
import {
buildErrorExtra,
Expand Down Expand Up @@ -101,12 +103,6 @@ export const DEFAULT_HEARTBEAT_LEASE_TTL_MS = 90 * 1000;
// 宿主可用 ctx.staleAfterMs 覆盖(与 claimLeaseMs 同一模式)。
export const STALE_AFTER_MS = 60 * 60 * 1000;

// 一次触发投递失败后最多再重试几次(退避 2 / 4 / 6 分钟)。宿主可用
// ctx.maxDeliveryRetries 调低(0 = 不重试)。重试不一定重新生成:内容已经落进
// outbox 的,重试只补推送(见 lib/message-processor.js 的 redeliverCommittedBatch);
// 生成本身失败的才会把整条生成再跑一遍,每跑一遍都花钱。
export const DEFAULT_MAX_DELIVERY_RETRIES = 3;

// 「重试也好不了」的判定(永久性错误码、终态推送状态码、payload 超限)住在
// lib/errors.js —— 定时任务的退避阶梯和 instant 任务的三轮重试用同一份口径。

Expand Down Expand Up @@ -241,12 +237,6 @@ function resolveStaleAfterMs(ctx) {
return positiveNumber(ctx.staleAfterMs) || STALE_AFTER_MS;
}

/** ctx.maxDeliveryRetries 取非负整数,别的值(没配、负数、小数)一律用默认值。 */
function resolveMaxDeliveryRetries(ctx) {
const value = ctx.maxDeliveryRetries;
return Number.isInteger(value) && value >= 0 ? value : DEFAULT_MAX_DELIVERY_RETRIES;
}

/**
* @typedef {Object} StaleSkipInfo
* @property {'stale'} reason
Expand Down Expand Up @@ -832,7 +822,7 @@ async function deliverTasks(ctx, tasks) {
// 不写出去的话下游只知道「失败了」,不知道该让用户重建订阅还是裁短内容。
const errorExtra = buildErrorExtra(errorCode, pushStatus);
try {
if (permanent || task.retry_count >= maxDeliveryRetries) {
if (permanent || failure.willRetry === false || task.retry_count >= (failure.retryLimit ?? maxDeliveryRetries)) {
const encrypted = await encryptPayloadWithLastError(task, decryptedPayload, userKey, reason, errorExtra);
if (isRecurringType(recurrenceType)) {
const nextSendAt = nextFutureOccurrence(Date.parse(task.next_send_at), recurrenceType, Date.now(), tzId);
Expand Down Expand Up @@ -1113,7 +1103,7 @@ async function deliverTasks(ctx, tasks) {
}
await handleDeliveryFailure(
task, sendResult.error || '消息发送失败', recurrenceType, decryptedPayload, userKey,
{ errorCode: sendResult.errorCode || null, permanent: sendResult.permanent === true, pushStatus: sendResult.pushStatusCode }
{ errorCode: sendResult.errorCode || null, permanent: sendResult.permanent === true, pushStatus: sendResult.pushStatusCode, willRetry: sendResult.willRetry, retryLimit: sendResult.retryLimit }
);
return;
}
Expand Down
2 changes: 2 additions & 0 deletions packages/rei-standard-amsg/server/src/server/single-user.js
Original file line number Diff line number Diff line change
Expand Up @@ -25,6 +25,7 @@
* @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)。
* @param {number} [config.maxScheduledTasksPerFire] - 一次 fire 里 hook 用 ctx.scheduleTask() 最多能建几条后续任务(默认 2,0 表示不许自排)。
* @param {number | ((task: Object) => number | undefined)} [config.maxGenerationRetries] - Retry limit before a complete batch is committed; defaults to maxDeliveryRetries.
* @param {number} [config.maxDeliveryRetries] - 定时任务一次触发投递失败后最多再重试几次(默认 3,0 = 第一次失败就终审)。
* 只管 runScheduledTick 的退避阶梯:返回的 ctx 交给 runScheduledTick 时带上它;instant 的请求内重试不受影响。
* @param {function} [config.onAfterSend] - 推送发出(或发挂)之后的可选 hook:
Expand Down Expand Up @@ -101,6 +102,7 @@ export function createSingleUserServer(config) {
maxScheduledTasksPerFire: config.maxScheduledTasksPerFire,
// 定时任务投递失败后的重试次数上限(默认 3)。handlers 用不到,宿主拿这个
// ctx 去调 runScheduledTick 时它跟着走。
maxGenerationRetries: config.maxGenerationRetries,
maxDeliveryRetries: config.maxDeliveryRetries
};

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -65,6 +65,7 @@ const EXPECTED_FEATURES = [
'redeliver-committed-batch',
'hook-usage-total',
'max-delivery-retries',
'max-generation-retries',
];

function makeWorker(extra = {}) {
Expand Down
Loading
Loading