From 0e644ed3b7ec77eea46100b8cb439992d13f16e7 Mon Sep 17 00:00:00 2001 From: Tosd0 <65720409+Tosd0@users.noreply.github.com> Date: Thu, 1 Oct 2026 00:10:11 +0900 Subject: [PATCH] feat(amsg-server): configure generation retries per task --- .../amsg-server-generation-retry-policy.md | 5 + packages/rei-standard-amsg/server/README.md | 20 ++ .../server/cloudflare/single-user-worker.js | 1 + .../src/server/handlers/capabilities.js | 2 + .../server/src/server/index.js | 2 + .../server/src/server/lib/agentic-fire.js | 5 + .../src/server/lib/message-processor.js | 32 ++- .../server/src/server/lib/retry-policy.js | 40 ++++ .../server/src/server/lib/run-tick.js | 18 +- .../server/src/server/single-user.js | 2 + .../server/test/capabilities.test.mjs | 1 + .../test/generation-retry-policy.test.mjs | 203 ++++++++++++++++++ 12 files changed, 310 insertions(+), 21 deletions(-) create mode 100644 .changeset/amsg-server-generation-retry-policy.md create mode 100644 packages/rei-standard-amsg/server/src/server/lib/retry-policy.js create mode 100644 packages/rei-standard-amsg/server/test/generation-retry-policy.test.mjs diff --git a/.changeset/amsg-server-generation-retry-policy.md b/.changeset/amsg-server-generation-retry-policy.md new file mode 100644 index 0000000..a637815 --- /dev/null +++ b/.changeset/amsg-server-generation-retry-policy.md @@ -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. diff --git a/packages/rei-standard-amsg/server/README.md b/packages/rei-standard-amsg/server/README.md index 2aaa699..58feb80 100644 --- a/packages/rei-standard-amsg/server/README.md +++ b/packages/rei-standard-amsg/server/README.md @@ -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 里维护的「我刚才说过什么」台账通常是读进内存 → 改 → 整份写回,两条各改各的再写回,后写的必然盖掉前面那条。 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 6d95e1d..56629f0 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 @@ -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)。不配 = 全并发,与以前一致。 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 53fb618..0542649 100644 --- a/packages/rei-standard-amsg/server/src/server/handlers/capabilities.js +++ b/packages/rei-standard-amsg/server/src/server/handlers/capabilities.js @@ -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) { diff --git a/packages/rei-standard-amsg/server/src/server/index.js b/packages/rei-standard-amsg/server/src/server/index.js index 1f77949..d1200e1 100644 --- a/packages/rei-standard-amsg/server/src/server/index.js +++ b/packages/rei-standard-amsg/server/src/server/index.js @@ -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'` 的请求内重试不受它影响。 @@ -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 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 356f2ff..28d3e99 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 @@ -113,6 +113,7 @@ * semantics apply. */ +import { failureRetryDecision } from './retry-policy.js'; import { assertValidDecision, buildSessionContext, @@ -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, @@ -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'); 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 ffaa960..8bb80b5 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 @@ -21,6 +21,7 @@ * 分片发,sw 收齐后还原。 */ +import { resolveMaxDeliveryRetries, resolveMaxGenerationRetries, failureRetryDecision } from './retry-policy.js'; import { randomUUID } from './webcrypto-utils.js'; import { buildContentPush, @@ -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 { @@ -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 @@ -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) { @@ -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 @@ -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 里(见 @@ -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) }; } } @@ -724,7 +742,7 @@ export async function processMessagesByUuid(uuid, ctx, maxRetries = 2, userId, p }; } - while (retryCount <= maxRetries) { + while (true) { let task; try { task = userId @@ -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) { // 确定性失败不进重试:再跑两轮也是同一个错,白让调用方多等、白烧一整轮 @@ -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; diff --git a/packages/rei-standard-amsg/server/src/server/lib/retry-policy.js b/packages/rei-standard-amsg/server/src/server/lib/retry-policy.js new file mode 100644 index 0000000..ad60c32 --- /dev/null +++ b/packages/rei-standard-amsg/server/src/server/lib/retry-policy.js @@ -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, + }; +} 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 2ef0904..de040e7 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 @@ -59,6 +59,8 @@ * @returns {Promise} 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, @@ -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 任务的三轮重试用同一份口径。 @@ -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 @@ -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); @@ -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; } 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 2a0f26a..06547a3 100644 --- a/packages/rei-standard-amsg/server/src/server/single-user.js +++ b/packages/rei-standard-amsg/server/src/server/single-user.js @@ -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: @@ -101,6 +102,7 @@ export function createSingleUserServer(config) { maxScheduledTasksPerFire: config.maxScheduledTasksPerFire, // 定时任务投递失败后的重试次数上限(默认 3)。handlers 用不到,宿主拿这个 // ctx 去调 runScheduledTick 时它跟着走。 + maxGenerationRetries: config.maxGenerationRetries, maxDeliveryRetries: config.maxDeliveryRetries }; diff --git a/packages/rei-standard-amsg/server/test/capabilities.test.mjs b/packages/rei-standard-amsg/server/test/capabilities.test.mjs index 8ced915..d844c6c 100644 --- a/packages/rei-standard-amsg/server/test/capabilities.test.mjs +++ b/packages/rei-standard-amsg/server/test/capabilities.test.mjs @@ -65,6 +65,7 @@ const EXPECTED_FEATURES = [ 'redeliver-committed-batch', 'hook-usage-total', 'max-delivery-retries', + 'max-generation-retries', ]; function makeWorker(extra = {}) { diff --git a/packages/rei-standard-amsg/server/test/generation-retry-policy.test.mjs b/packages/rei-standard-amsg/server/test/generation-retry-policy.test.mjs new file mode 100644 index 0000000..26f8b23 --- /dev/null +++ b/packages/rei-standard-amsg/server/test/generation-retry-policy.test.mjs @@ -0,0 +1,203 @@ +import { test } from 'node:test'; +import assert from 'node:assert/strict'; +import { runTask, runScheduledTick } from '../src/server/lib/run-tick.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); + +async function fixture(options = {}) { + const db = createD1Adapter(createTestD1()); + await db.initSchema(); + await seedPushSubscription(db, USER, MASTER_KEY); + const payload = { + messageType: 'auto', contactName: 'Rei', recurrenceType: 'none', + completePrompt: 'hello', apiUrl: 'https://api.example.com/v1/chat/completions', + apiKey: 'secret-key', primaryModel: 'test', metadata: { interactive: true }, + ...options.payload, + }; + const key = await deriveUserEncryptionKey(USER, MASTER_KEY); + await db.createTask({ + uuid: 'generation-policy', user_id: USER, message_type: payload.messageType, + next_send_at: new Date(Date.now() - 1000).toISOString(), + encrypted_payload: await encryptForStorage(JSON.stringify(payload), key), + }); + const settled = []; + const ctx = { + db, masterKey: MASTER_KEY, + vapid: { email: 'mailto:test@example.com', publicKey: 'pub', privateKey: 'priv' }, + webpush: { async sendNotification() {} }, + hooks: { + onBeforeFire: async () => [{ role: 'user', content: 'hello' }], + onLLMOutput: async () => ({ decision: 'finish', pushPayloads: [{ messageKind: 'content', message: 'hello' }] }), + }, + onFireSettled: async info => settled.push(info), + maxGenerationRetries: task => task.metadata?.interactive ? 0 : undefined, + ...options.ctx, + }; + return { db, ctx, settled, async row() { return (await db.listTasks(USER, { status: 'all', limit: 50 })).tasks[0]; } }; +} + +const failures = [ + ['HTTP 200 error body', () => Response.json({ error: { message: 'no balance', code: 'provider_private_code' } })], + ['HTTP 429', () => Response.json({ error: { message: 'no balance' } }, { status: 429 })], + ['HTTP 503', () => new Response('upstream unavailable', { status: 503 })], + ['network failure', () => { throw new TypeError('fetch failed'); }], +]; + +for (const [label, response] of failures) { + test(`no generation retries: ${label} is terminal on first run`, async t => { + let requests = 0; + t.mock.method(globalThis, 'fetch', async () => { requests++; return response(); }); + const { ctx, row, settled } = await fixture(); + await runTask(ctx, 'generation-policy'); + assert.equal((await row()).status, 'failed'); + assert.equal((await row()).retry_count, 0); + assert.equal(settled[0].willRetry, false); + assert.equal(settled[0].failureStage, 'generation'); + assert.equal(settled[0].error.permanent, undefined, 'policy must not mutate provider errors'); + await runTask(ctx, 'generation-policy'); + assert.equal(requests, 1); + }); +} + +test('pre-generation hook failure follows task policy and passes credential-free task', async () => { + let selected; + const { ctx, row, settled } = await fixture({ ctx: { + maxGenerationRetries: task => { selected = task; return 0; }, + hooks: { onBeforeFire: async () => { throw new Error('pack read failed'); }, onLLMOutput: async () => ({ decision: 'skip-push' }) }, + } }); + await runScheduledTick(ctx); + assert.equal((await row()).status, 'failed'); + assert.equal(selected.metadata.interactive, true); + assert.equal(selected.apiKey, undefined); + assert.equal(selected.encrypted_payload, undefined); + assert.equal(settled[0].willRetry, false); +}); + +test('default task still enters delivery retry backoff', async t => { + t.mock.method(globalThis, 'fetch', async () => { throw new TypeError('fetch failed'); }); + const { ctx, row, settled } = await fixture({ payload: { metadata: { interactive: false } } }); + await runTask(ctx, 'generation-policy'); + assert.equal((await row()).status, 'pending'); + assert.equal((await row()).retry_count, 1); + assert.ok((await ctx.db.getTaskByUuidOnly('generation-policy')).retry_after); + assert.equal(settled[0].willRetry, true); +}); + +test('committed batch keeps push retries and never regenerates with generation limit zero', async t => { + let requests = 0; + t.mock.method(globalThis, 'fetch', async () => { requests++; return Response.json({ choices: [{ message: { content: 'hello' } }] }); }); + let pushes = 0; + const { db, ctx, row, settled } = await fixture({ ctx: { + webpush: { async sendNotification() { pushes++; if (pushes === 1) throw new Error('push network failed'); } }, + } }); + await runTask(ctx, 'generation-policy'); + const retry = await row(); + assert.equal(retry.status, 'pending'); + assert.equal(retry.retry_count, 1); + assert.equal(settled[0].outboxed, true); + assert.equal(settled[0].failureStage, 'delivery'); + assert.equal(settled[0].willRetry, true); + await db.updateTaskById(retry.id, { retry_after: new Date(Date.now() - 1000).toISOString() }); + await runTask(ctx, 'generation-policy'); + assert.equal(await row(), undefined); + assert.equal(requests, 1); + assert.equal(pushes, 2); + assert.equal(settled.length, 1); +}); + +test('frozen prompt path also honors generation limit zero', async t => { + t.mock.method(globalThis, 'fetch', async () => { throw new TypeError('fetch failed'); }); + const { ctx, row } = await fixture({ ctx: { hooks: null, maxGenerationRetries: 0 } }); + await runTask(ctx, 'generation-policy'); + assert.equal((await row()).status, 'failed'); +}); + +test('nonzero generation limit is independent from delivery retry limit', async t => { + t.mock.method(globalThis, 'fetch', async () => { throw new TypeError('fetch failed'); }); + const { db, ctx, row, settled } = await fixture({ ctx: { maxGenerationRetries: 1 } }); + await runTask(ctx, 'generation-policy'); + const first = await row(); + assert.equal(first.retry_count, 1); + assert.equal(settled[0].willRetry, true); + await db.updateTaskById(first.id, { retry_after: new Date(Date.now() - 1000).toISOString() }); + await runTask(ctx, 'generation-policy'); + assert.equal((await row()).status, 'failed'); + assert.equal(settled[1].willRetry, false); +}); + +for (const invalid of [-1, 1.5, null, () => { throw new Error('broken policy'); }]) { + test(`invalid policy ${String(invalid)} fails before generation using default retry handling`, async t => { + let requests = 0; + t.mock.method(globalThis, 'fetch', async () => { requests++; throw new Error('must not call'); }); + const { ctx, row, settled } = await fixture({ ctx: { maxGenerationRetries: invalid } }); + await runTask(ctx, 'generation-policy'); + const stored = await row(); + assert.equal(stored.status, 'pending'); + assert.equal(stored.retry_count, 1); + assert.equal(JSON.parse(stored.last_error).errorCode, 'GENERATION_RETRY_POLICY_INVALID'); + assert.equal(requests, 0); + assert.equal(settled.length, 0, 'policy resolution failed before onBeforeFire was entered'); + }); +} + +test('successful and skipped fire receipts have no failure decision', async t => { + t.mock.method(globalThis, 'fetch', async () => Response.json({ choices: [{ message: { content: 'hello' } }] })); + const success = await fixture(); + await runTask(success.ctx, 'generation-policy'); + assert.equal(success.settled[0].willRetry, null); + assert.equal(success.settled[0].failureStage, null); + const skipped = await fixture({ ctx: { hooks: { onBeforeFire: async () => ({ skip: true }), onLLMOutput: async () => ({ decision: 'skip-push' }) } } }); + await runTask(skipped.ctx, 'generation-policy'); + assert.equal(skipped.settled[0].willRetry, null); +}); + +test('delivery retry exhaustion and permanent failures agree with the receipt', async t => { + t.mock.method(globalThis, 'fetch', async () => Response.json({ choices: [{ message: { content: 'hello' } }] })); + const delivery = await fixture({ ctx: { maxDeliveryRetries: 0, webpush: { async sendNotification() { throw new Error('push failed'); } } } }); + await runTask(delivery.ctx, 'generation-policy'); + assert.equal((await delivery.row()).status, 'failed'); + assert.equal(delivery.settled[0].willRetry, false); + assert.equal(delivery.settled[0].failureStage, 'delivery'); +}); + +test('legacy UUID entry point honors the same task policy', async t => { + const { processMessagesByUuid } = await import('../src/server/lib/message-processor.js'); + let requests = 0; + t.mock.method(globalThis, 'fetch', async () => { requests++; throw new TypeError('fetch failed'); }); + const { ctx, row, settled } = await fixture(); + const result = await processMessagesByUuid('generation-policy', ctx, 2, USER); + assert.equal(result.success, false); + assert.equal(result.error.retriesAttempted, 0); + assert.equal((await row()).status, 'failed'); + assert.equal(settled[0].willRetry, false); + assert.equal(requests, 1); +}); + +test('cancelled attempt cannot advertise another retry after an unrelated in-flight error', async () => { + const { processSingleMessage } = await import('../src/server/lib/message-processor.js'); + let cancelled = false; + const { ctx, settled } = await fixture({ ctx: { + maxGenerationRetries: 3, + isTaskCancelled: () => cancelled, + hooks: { onBeforeFire: async () => { cancelled = true; throw new Error('in-flight request aborted'); }, onLLMOutput: async () => ({ decision: 'skip-push' }) }, + } }); + const task = await ctx.db.getTaskByUuidOnly('generation-policy'); + const result = await processSingleMessage(task, ctx); + assert.equal(result.willRetry, false); + assert.equal(settled[0].willRetry, false); +}); + +test('Cloudflare worker factory forwards the generation policy to its queued task entrypoint', async t => { + const { createSingleUserCloudflareWorker } = await import('../src/server/cloudflare/single-user-worker.js'); + t.mock.method(globalThis, 'fetch', async () => { throw new TypeError('fetch failed'); }); + const { ctx, row, settled } = await fixture(); + const worker = createSingleUserCloudflareWorker(() => ctx); + await worker.runTask('generation-policy', {}); + assert.equal((await row()).status, 'failed'); + assert.equal(settled[0].willRetry, false); +});