From 0822c473b1699a95adc4c5c2ed917ae8d649917f Mon Sep 17 00:00:00 2001 From: Tosd0 <65720409+Tosd0@users.noreply.github.com> Date: Wed, 23 Sep 2026 11:11:04 +0800 Subject: [PATCH 1/4] =?UTF-8?q?fix(amsg-server):=20=E6=94=B9=E4=BB=BB?= =?UTF-8?q?=E5=8A=A1=E4=B8=8D=E5=86=8D=E8=AE=A9=E8=BF=99=E6=AC=A1=E8=A7=A6?= =?UTF-8?q?=E5=8F=91=E9=87=8D=E6=96=B0=E7=94=9F=E6=88=90=EF=BC=9B=E9=A1=B6?= =?UTF-8?q?=E6=9B=BF=E6=97=A7=E4=BB=BB=E5=8A=A1=E5=9C=A8=20pg=20/=20neon?= =?UTF-8?q?=20=E4=B8=8A=E4=B8=80=E8=B5=B7=E6=88=90=E8=B4=A5=EF=BC=9B?= =?UTF-8?q?=E6=8A=95=E9=80=92=E4=B8=AD=E4=B8=8D=E6=8E=A5=E5=8F=97=E6=94=B9?= =?UTF-8?q?=E6=8E=92=E6=9C=9F?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - 「这次触发的内容有没有落定」不再看 retry_count:改任务会把重试计数清零,于是 用户在重试窗口里随手改一下,重试那一跳就重新调 LLM 生成一整条,旧那批还在 outbox 里等补收,同一时刻冒出两份内容。现在每次投递都去 outbox 认这次触发的 批次,认到只补推送。 - pg / neon 补上 createTaskSuperseding(一条数据修改型 CTE,INSERT 抛错时 DELETE 跟着回滚)。自定义适配器的两步退路改成先建新、后删旧:反过来的话建新失败会把 旧任务白删,客户端还以为它在。 - 三个适配器的 updateTaskByUuid 在改 next_send_at 时加租约门:任务正被投递占着 就不改。PUT /update-message 回 409 TASK_IN_FLIGHT,renewTask 回 { renewed: false, reason: 'in_flight' },不再假装成功。只改正文不受影响。 - 分片上限校验的探针改用最长的 messageKind(tool_request),之前按 reasoning 量 少算 3 字节,配在上限附近的部署里 tool_request 的分片会被整批拒收。 回归守卫:改任务后仍只补推送、顶替失败时旧任务还在、投递中改排期被挡(适配器 / handler / renewTask 三处)、放行的最大分片对每种 messageKind 都成立。 Claude-Session: https://claude.ai/code/session_017fqgW4ShQ4yjJJwqAiJrcT --- .changeset/amsg-server-inflight-and-resume.md | 22 ++++++ .../server/src/server/adapters/d1.js | 17 ++++- .../server/src/server/adapters/interface.js | 6 +- .../server/src/server/adapters/neon.js | 21 +++++- .../server/src/server/adapters/pg-shared.js | 47 +++++++++++++ .../server/src/server/adapters/pg.js | 20 +++++- .../src/server/handlers/schedule-message.js | 15 ++-- .../src/server/handlers/update-message.js | 15 ++++ .../server/src/server/lib/agentic-fire.js | 13 +++- .../src/server/lib/message-processor.js | 65 +++++++++-------- .../server/src/server/lib/run-tick.js | 5 +- .../server/test/d1-adapter.test.mjs | 26 +++++++ .../server/test/pg-neon-adapter.test.mjs | 47 +++++++++++++ .../test/processor-multipart-pacing.test.mjs | 37 +++++++++- .../server/test/spend-guards.test.mjs | 33 +++++++++ .../test/update-message-fields.test.mjs | 40 +++++++++++ .../server/test/upstream-additions.test.mjs | 69 +++++++++++++++++++ 17 files changed, 450 insertions(+), 48 deletions(-) create mode 100644 .changeset/amsg-server-inflight-and-resume.md diff --git a/.changeset/amsg-server-inflight-and-resume.md b/.changeset/amsg-server-inflight-and-resume.md new file mode 100644 index 0000000..b214b6c --- /dev/null +++ b/.changeset/amsg-server-inflight-and-resume.md @@ -0,0 +1,22 @@ +--- +"@rei-standard/amsg-server": minor +--- + +改任务不再让这次触发重新生成一遍;顶替旧任务在 pg / neon 上也一起成败;投递进行中不接受改排期 + +**改任务不再让已经落定的那次触发重新生成。** 一次触发的整批内容落进收件箱之后推送失败、任务进入重试,这期间用户调 `PUT /update-message` 改了任何字段(哪怕只是联系人名字),重试那一跳之前会重新调 LLM 生成一整条:旧那批还在收件箱里等客户端补收,同一个时刻于是冒出两份内容,生成的钱也白花一次。原因是「这次触发有没有落定的内容」之前靠 `retry_count > 0` 判断,而改任务会把重试计数清零(这是它该做的——修好 apiKey 的任务不该背着旧账)。现在每次投递都直接去收件箱里认这次触发的批次,认到就只补推送。代价是每次定时触发多一次收件箱查询(D1 走已有索引;pg / neon 没有收件箱,这一步直接跳过)。 + +一次触发的内容一旦落定,这次就按落定的那份发完;改动从下一次触发开始生效。一次性任务在这个窗口里改内容,改的这一版就用不上了。 + +**`supersedesUuid` 顶替旧任务在 pg / neon 上也是一条语句。** 之前只有 D1 是「删旧 + 建新一起成败」,pg / neon 退回先删旧、再建新两步:建新那一步失败(uuid 撞了、连接一时抖了)时旧任务已经删掉,接口却回失败,客户端以为旧任务还在,其实再也找不回来。现在 `pg` / `neon` 也实现了 `createTaskSuperseding`(数据修改型 CTE,INSERT 抛错时 DELETE 跟着回滚)。没实现这个方法的自定义适配器仍走两步退路,但顺序改成先建新、后删旧——最坏是两条都留着,客户端重试一次就能收拾干净。 + +**任务正在投递时不接受改排期。** 投递收尾会按它领取时看到的排期推进下一次(一次性任务标成已发送),这期间写进去的新时刻随后就被盖掉,接口却已经回了成功。现在三个内置适配器的 `updateTaskByUuid` 在改 `next_send_at` 时多一道租约门: + +| 入口 | 碰上正在投递的任务 | +|---|---| +| `PUT /update-message` 带 `nextSendAt` | 409 `TASK_IN_FLIGHT`(任务不存在仍是 404,行不是 pending 仍是 409 `UPDATE_CONFLICT`) | +| fire hook 的 `ctx.renewTask(uuid, nextSendAt)` | `{ renewed: false, reason: 'in_flight' }` | + +投递一般几秒到几十秒(agentic 链路默认最多 240 秒),worker 中途没了的话租约在 ~90 秒内到期,之后照常能改。只改正文这类字段不受这道门约束:它们只影响以后的触发,收尾本来就不覆盖用户刚保存的内容。等重试的那几分钟任务没被占用,照常能改。 + +**分片大小上限的校验按最长的 `messageKind` 算。** `multipart.maxChunkBytes` 的上限校验用一条探针量信封开销,探针之前固定用 `reasoning`,比最长的 `tool_request` 短 3 字节。把 `maxChunkBytes` 配在上限附近(例如照着校验失败时建议的最大值去配)的部署,`content` / `reasoning` 的分片发得出去,`tool_request` 的分片每一片都超出单条 push 的明文上限、被推送服务拒收。现在按真实取值里最长的那个量,报出来的最大值对所有类型都成立。 diff --git a/packages/rei-standard-amsg/server/src/server/adapters/d1.js b/packages/rei-standard-amsg/server/src/server/adapters/d1.js index f076d9c..b62e170 100644 --- a/packages/rei-standard-amsg/server/src/server/adapters/d1.js +++ b/packages/rei-standard-amsg/server/src/server/adapters/d1.js @@ -384,6 +384,15 @@ export class D1Adapter { return this._db.prepare('SELECT * FROM scheduled_messages WHERE id = ?').bind(taskId).first(); } + /** + * 按 uuid 改一条 pending 任务(PUT /update-message、fire hook 的 renewTask)。 + * + * 改排期(extraFields 里带 next_send_at)时多一道租约门:这条任务正被一次投递 + * 占着的话不改,返回 null。投递收尾会按自己领取时看到的排期推进下一次(一次性 + * 任务干脆标成已发送),这期间写进去的新时刻随后就被盖掉,接口却已经回了成功 + * ——用户以为改期生效了,实际什么都没留下。正文这类字段不受这道门约束:它们 + * 只影响以后的触发,收尾那边本来就不会覆盖(见 run-tick 的收尾守卫)。 + */ async updateTaskByUuid(uuid, userId, encryptedPayload, extraFields) { const now = this._now(); const sets = ['encrypted_payload = ?', 'updated_at = ?']; @@ -399,9 +408,15 @@ export class D1Adapter { } values.push(uuid, userId); + let leaseGate = ''; + if (extraFields && Object.prototype.hasOwnProperty.call(extraFields, 'next_send_at')) { + leaseGate = ' AND (lease_until IS NULL OR lease_until <= ?)'; + values.push(now); + } + const res = await this._db.prepare( `UPDATE scheduled_messages SET ${sets.join(', ')} - WHERE uuid = ? AND user_id = ? AND status = 'pending'` + WHERE uuid = ? AND user_id = ? AND status = 'pending'${leaseGate}` ).bind(...values).run(); if (!res.meta.changes) return null; diff --git a/packages/rei-standard-amsg/server/src/server/adapters/interface.js b/packages/rei-standard-amsg/server/src/server/adapters/interface.js index fd97868..09a62bd 100644 --- a/packages/rei-standard-amsg/server/src/server/adapters/interface.js +++ b/packages/rei-standard-amsg/server/src/server/adapters/interface.js @@ -189,8 +189,10 @@ * 它把「为什么失败」透给已失败的行;不实现时退回 getTaskStatus(409 里就 * 没有 lastError)。 * @property {(params: InsertTaskParams, supersedesUuid: string) => Promise} [createTaskSuperseding] - * (可选)建新任务的同一事务里取消旧的那条(POST /schedule-message 的 - * supersedesUuid)。不实现时 handler 退回「先删再建」两步(失去原子性)。 + * (可选)建新任务的同时取消旧的那条,两件事一起成败(POST /schedule-message + * 的 supersedesUuid)。内置的 D1 / pg / neon 都实现了。不实现时 handler 退回 + * 两步:先建新、后删旧——删旧失败最坏是两条都留着,客户端重试一次就能收拾, + * 而反过来(先删后建)建新失败就把旧任务白删了,找不回来。 * @property {(userId: string, rows: Array) => Promise} [appendOutboxMessages] * (可选;单用户/D1)push 发送前把整批落进 message_outbox(密文 payload), * (user_id, message_id) 冲突时更新未 ack 的行、不动已 ack 的。 diff --git a/packages/rei-standard-amsg/server/src/server/adapters/neon.js b/packages/rei-standard-amsg/server/src/server/adapters/neon.js index 0563df3..a879274 100644 --- a/packages/rei-standard-amsg/server/src/server/adapters/neon.js +++ b/packages/rei-standard-amsg/server/src/server/adapters/neon.js @@ -116,6 +116,12 @@ export class NeonAdapter { return rows[0] || null; } + // 建新任务 + 删掉被顶替的旧任务,一条语句一起成败(SQL 与语义见 pg-shared.js)。 + async createTaskSuperseding(params, supersedesUuid) { + const sql = this._getSql(); + return pgShared.createTaskSuperseding((text, args) => sql.query(text, args), params, supersedesUuid); + } + async getTaskByUuid(uuid, userId) { const sql = this._getSql(); const rows = await sql.query( @@ -181,6 +187,15 @@ export class NeonAdapter { return rows[0] || null; } + /** + * 按 uuid 改一条 pending 任务(PUT /update-message、fire hook 的 renewTask)。 + * + * 改排期(extraFields 里带 next_send_at)时多一道租约门:这条任务正被一次投递 + * 占着的话不改,返回 null。投递收尾会按自己领取时看到的排期推进下一次(一次性 + * 任务干脆标成已发送),这期间写进去的新时刻随后就被盖掉,接口却已经回了成功 + * ——用户以为改期生效了,实际什么都没留下。正文这类字段不受这道门约束:它们 + * 只影响以后的触发,收尾那边本来就不会覆盖(见 run-tick 的收尾守卫)。 + */ async updateTaskByUuid(uuid, userId, encryptedPayload, extraFields) { const sql = this._getSql(); const sets = ['encrypted_payload = $1', 'updated_at = NOW()']; @@ -199,9 +214,13 @@ export class NeonAdapter { } values.push(uuid, userId); + // 改排期时多一道租约门(语义见方法头注释)。 + const leaseGate = extraFields && Object.prototype.hasOwnProperty.call(extraFields, 'next_send_at') + ? ' AND (lease_until IS NULL OR lease_until <= NOW())' + : ''; const rows = await sql.query( `UPDATE scheduled_messages SET ${sets.join(', ')} - WHERE uuid = $${idx} AND user_id = $${idx + 1} AND status = 'pending' + WHERE uuid = $${idx} AND user_id = $${idx + 1} AND status = 'pending'${leaseGate} RETURNING uuid, updated_at`, values ); diff --git a/packages/rei-standard-amsg/server/src/server/adapters/pg-shared.js b/packages/rei-standard-amsg/server/src/server/adapters/pg-shared.js index 009b74f..8e19deb 100644 --- a/packages/rei-standard-amsg/server/src/server/adapters/pg-shared.js +++ b/packages/rei-standard-amsg/server/src/server/adapters/pg-shared.js @@ -113,6 +113,53 @@ export async function renewTaskLease(query, taskId, leaseUntil) { return rows.length > 0; } +/** + * 建新任务的同时把被顶替的旧任务删掉(POST /schedule-message 带 + * supersedesUuid 时走这条)。 + * + * 删旧和建新必须一起成败:分成两次查询的话,建新那步失败(uuid 撞了、连接 + * 一时抖了)时旧任务已经删掉了,接口却回失败——客户端以为旧任务还在,实际 + * 上它已经没了,没有任何办法找回来。 + * + * 一条语句解决,不用显式事务:数据修改型 CTE 里的 DELETE 一定会执行到底, + * 而 INSERT 抛错时整条语句一起回滚(pg 与 neon 的 HTTP 驱动都是单语句一个 + * 隐式事务,用法完全一致)。superseded 由 DELETE 到底删没删到行推出来。 + * + * @param {PgQuery} query + * @param {{ user_id: string, uuid: string, encrypted_payload: string, + * next_send_at: string|Date, message_type: string }} params - 新任务 + * @param {string} supersedesUuid - 要顶替掉的旧任务 uuid(同一个用户名下) + * @returns {Promise<{ id: number, uuid: string, next_send_at: any, status: string, + * created_at: any, superseded: boolean }|null>} + */ +export async function createTaskSuperseding(query, params, supersedesUuid) { + const rows = await query( + `WITH deleted AS ( + DELETE FROM scheduled_messages + WHERE uuid = $6 AND user_id = $1 + RETURNING id + ), inserted AS ( + INSERT INTO scheduled_messages + (user_id, uuid, encrypted_payload, next_send_at, message_type, status, retry_count, created_at, updated_at) + VALUES ($1, $2, $3, $4, $5, 'pending', 0, NOW(), NOW()) + RETURNING id, uuid, next_send_at, status, created_at + ) + SELECT inserted.*, (SELECT count(*) FROM deleted) > 0 AS superseded + FROM inserted`, + [ + params.user_id, + params.uuid, + params.encrypted_payload, + params.next_send_at, + params.message_type, + supersedesUuid, + ] + ); + const row = rows[0]; + if (!row) return null; + return { ...row, superseded: row.superseded === true }; +} + /** * 状态 + 失败摘要(GET /message 用它把「为什么失败」透给已失败的行)。 * diff --git a/packages/rei-standard-amsg/server/src/server/adapters/pg.js b/packages/rei-standard-amsg/server/src/server/adapters/pg.js index 3961407..002d77f 100644 --- a/packages/rei-standard-amsg/server/src/server/adapters/pg.js +++ b/packages/rei-standard-amsg/server/src/server/adapters/pg.js @@ -130,6 +130,11 @@ export class PgAdapter { return rows[0] || null; } + // 建新任务 + 删掉被顶替的旧任务,一条语句一起成败(SQL 与语义见 pg-shared.js)。 + async createTaskSuperseding(params, supersedesUuid) { + return pgShared.createTaskSuperseding((text, args) => this._query(text, args), params, supersedesUuid); + } + async getTaskByUuid(uuid, userId) { const rows = await this._query( `SELECT ${TASK_DETAIL_COLUMNS} @@ -191,6 +196,15 @@ export class PgAdapter { return rows[0] || null; } + /** + * 按 uuid 改一条 pending 任务(PUT /update-message、fire hook 的 renewTask)。 + * + * 改排期(extraFields 里带 next_send_at)时多一道租约门:这条任务正被一次投递 + * 占着的话不改,返回 null。投递收尾会按自己领取时看到的排期推进下一次(一次性 + * 任务干脆标成已发送),这期间写进去的新时刻随后就被盖掉,接口却已经回了成功 + * ——用户以为改期生效了,实际什么都没留下。正文这类字段不受这道门约束:它们 + * 只影响以后的触发,收尾那边本来就不会覆盖(见 run-tick 的收尾守卫)。 + */ async updateTaskByUuid(uuid, userId, encryptedPayload, extraFields) { const sets = ['encrypted_payload = $1', 'updated_at = NOW()']; const values = [encryptedPayload]; @@ -208,9 +222,13 @@ export class PgAdapter { } values.push(uuid, userId); + // 改排期时多一道租约门(语义见方法头注释)。 + const leaseGate = extraFields && Object.prototype.hasOwnProperty.call(extraFields, 'next_send_at') + ? ' AND (lease_until IS NULL OR lease_until <= NOW())' + : ''; const rows = await this._query( `UPDATE scheduled_messages SET ${sets.join(', ')} - WHERE uuid = $${idx} AND user_id = $${idx + 1} AND status = 'pending' + WHERE uuid = $${idx} AND user_id = $${idx + 1} AND status = 'pending'${leaseGate} RETURNING uuid, updated_at`, values ); diff --git a/packages/rei-standard-amsg/server/src/server/handlers/schedule-message.js b/packages/rei-standard-amsg/server/src/server/handlers/schedule-message.js index 00f28a5..10bcf3e 100644 --- a/packages/rei-standard-amsg/server/src/server/handlers/schedule-message.js +++ b/packages/rei-standard-amsg/server/src/server/handlers/schedule-message.js @@ -235,9 +235,14 @@ export function createScheduleMessageHandler(ctx) { next_send_at: effectiveSendTime, message_type: payload.messageType }; - // supersede:建这条的同时取消旧的那条。适配器支持原子形态(删旧 + 建新落 - // 在同一事务,见 D1 的 createTaskSuperseding)就走它;不支持的退回「先删 - // 再建」两步——语义相同,只是失去原子性和那次省下的往返。 + // supersede:建这条的同时取消旧的那条。适配器支持原子形态(删旧 + 建新一起 + // 成败,见 D1 和 pg-shared 的 createTaskSuperseding)就走它,内置的三个适配 + // 器都支持。 + // + // 自定义适配器没实现的退回两步,顺序是先建新、后删旧:反过来的话建新失败 + // (uuid 撞了、连接一时抖了)时旧任务已经删掉,接口却回失败,客户端以为旧 + // 任务还在,其实再也找不回来了。这个顺序最坏是删旧失败、两条都留着,客户端 + // 重试一次就能收拾干净。 const supersedesUuid = payload.supersedesUuid || null; let superseded = false; let dbResult; @@ -246,10 +251,10 @@ export function createScheduleMessageHandler(ctx) { dbResult = await db.createTaskSuperseding(createParams, supersedesUuid); superseded = !!(dbResult && dbResult.superseded); } else { - if (supersedesUuid) { + dbResult = await db.createTask(createParams); + if (supersedesUuid && dbResult) { superseded = await db.deleteTaskByUuid(supersedesUuid, userId); } - dbResult = await db.createTask(createParams); } } catch (error) { if (isUniqueViolation(error)) { diff --git a/packages/rei-standard-amsg/server/src/server/handlers/update-message.js b/packages/rei-standard-amsg/server/src/server/handlers/update-message.js index f677860..93f2efc 100644 --- a/packages/rei-standard-amsg/server/src/server/handlers/update-message.js +++ b/packages/rei-standard-amsg/server/src/server/handlers/update-message.js @@ -325,6 +325,21 @@ export function createUpdateMessageHandler(ctx) { const result = await db.updateTaskByUuid(taskUuid, userId, encryptedPayload, extraFields); if (!result) { + // 改排期时行还在、还是 pending,那就是被一次正在跑的投递占着(见适配器 + // updateTaskByUuid 的租约门)。这跟「任务没了」得分开说:客户端过几十秒 + // 重试就能成,不该让用户以为任务被删了。 + const inFlight = updates.nextSendAt + && typeof db.getTaskByUuid === 'function' + && !!(await db.getTaskByUuid(taskUuid, userId)); + if (inFlight) { + return { + status: 409, + body: { + success: false, + error: { code: 'TASK_IN_FLIGHT', message: '这条任务正在投递,等这次发完再改排期' } + } + }; + } return { status: 409, body: { success: false, error: { code: 'UPDATE_CONFLICT', message: '任务更新失败,任务可能已被修改或删除' } } }; } 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 a2de346..356f2ff 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 @@ -660,8 +660,11 @@ export async function runAgenticFire({ task, decryptedPayload, userKey, ctx }) { * @param {string} uuid * @param {string} nextSendAt - ISO 8601 * @returns {Promise<{ renewed: true, uuid: string, nextSendAt: string } - * | { renewed: false, reason: 'not_found' }>} not_found = 行不存在或已不是 - * pending(宿主自己决定要不要转头 scheduleTask 一条新的)。 + * | { renewed: false, reason: 'not_found' | 'in_flight' }>} + * not_found = 行不存在或已不是 pending(宿主自己决定要不要转头 scheduleTask + * 一条新的);in_flight = 那条任务此刻正被一次投递占着,排期没改动——它正要 + * 发出去,收尾会按自己那份排期推进,硬改进去也留不住。宿主想推迟的那次已经 + * 在路上了,通常没什么好补救的;真要顺延下一次,等这次发完再调一遍。 */ const renewTask = async (uuid, nextSendAt) => { if (typeof uuid !== 'string' || !uuid.trim()) { @@ -702,7 +705,11 @@ export async function runAgenticFire({ task, decryptedPayload, userKey, ctx }) { retry_count: 0, ...(typeof ctx.db.claimTask === 'function' ? { retry_after: null } : {}), }); - if (!updated) return { renewed: false, reason: 'not_found' }; + if (!updated) { + // 行还在、还是 pending 的话,改不动的原因是它正被一次投递占着。 + const stillPending = await ctx.db.getTaskByUuid(uuid, task.user_id); + return { renewed: false, reason: stillPending ? 'in_flight' : 'not_found' }; + } return { renewed: true, uuid, nextSendAt: nextSendAtIso }; }; 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 3945cd0..ffaa960 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 @@ -28,6 +28,7 @@ import { buildReasoningPush, readReasoningContent, stripReasoningTags, + MESSAGE_KIND, DEFAULT_MULTIPART_CHUNK_BYTES, DEFAULT_MULTIPART_MAX_CHUNKS, DEFAULT_MULTIPART_MAX_TOTAL_BYTES, @@ -286,7 +287,7 @@ function assertChunkBytesFitPushLimit({ maxChunkBytes, maxChunks, ttlMs }) { // 最小探针:3 字节原文 → base64url 后恰好 4 字符,信封开销 = 总长 - 4。 const PROBE_CHUNK_BYTES = 3; const [probe] = buildMultipartPushPayloads( - { messageKind: 'reasoning' }, + { messageKind: longestMessageKind() }, { serializedPayload: 'x'.repeat(PROBE_CHUNK_BYTES), maxChunkBytes: PROBE_CHUNK_BYTES, ttlMs } ); // index / total 在真实批次里最多到 maxChunks(探针里各只有 1 位)。 @@ -319,6 +320,17 @@ function base64UrlLength(n) { return Math.ceil(n * 4 / 3); } +/** + * 取值最长的那个 messageKind。分片信封里原样带着原消息的 kind,探针得按最长的 + * 那个量,否则 kind 短的类型量出来的开销偏小,把 maxChunkBytes 配在上限边上时 + * 只有长 kind 的那几类会被推送服务整批拒收。 + * + * @returns {string} + */ +function longestMessageKind() { + return Object.values(MESSAGE_KIND).reduce((a, b) => (b.length > a.length ? b : a)); +} + /** * @param {unknown} value * @param {number} fallback @@ -409,23 +421,18 @@ async function redeliverCommittedBatch(task, ctx, userKey, decryptedPayload, bat /** * Process a single database task row: decrypt → generate content → push. * - * 重试同一次触发时(`options.resumeCommittedBatch`),先看这次触发的整批是不是 - * 已经落进了 outbox:落进去了就只补推送、不再生成(见 redeliverCommittedBatch), - * 返回值带 `redelivered: true`。 + * 每次投递都先看这次触发的整批是不是已经落进了 outbox:落进去了就只补推送、 + * 不再生成(见 redeliverCommittedBatch),返回值带 `redelivered: true`。 * * @param {import('../adapters/interface.js').TaskRow} task * @param {ProcessorContext} ctx * @param {string} [providedMasterKey] * @param {{ userKey: string, payload: Object } | null} [predecrypted] - 调用方 * (run-tick 的预扫描)已经解好的 payload;传了就不再解第二遍。 - * @param {{ resumeCommittedBatch?: boolean }} [options] - resumeCommittedBatch:这 - * 一跳是不是同一次触发的重试(是才去 outbox 里找落定的批次,首次触发不多花这 - * 次查询)。不传时按 `task.retry_count > 0` 判断——定时任务的重试计数就记在 - * 这一列上;processMessagesByUuid 的请求内重试不改这一列,由它显式传。 * @returns {Promise<{ success: boolean, messagesSent: number, redelivered?: boolean, pushedCount?: number, error?: string, errorCode?: string|null, pushStatusCode?: number|null, permanent?: boolean }>} * 失败时 `pushStatusCode` 是推送服务回的 HTTP 状态码(不是推送阶段炸的 → null)。 */ -export async function processSingleMessage(task, ctx, providedMasterKey, predecrypted = null, options = {}) { +export async function processSingleMessage(task, ctx, providedMasterKey, predecrypted = null) { try { const masterKey = providedMasterKey || ctx.masterKey; if (!masterKey) { @@ -437,21 +444,22 @@ export async function processSingleMessage(task, ctx, providedMasterKey, predecr const decryptedPayload = (predecrypted && predecrypted.payload) || JSON.parse(await decryptFromStorage(task.encrypted_payload, userKey)); - // 同一次触发的重试:内容已经落定的话只补推送。放在所有生成路径(agentic 与 - // 冻结 prompt)之前,两条路落进 outbox 的批次都认。 - const resumeCommittedBatch = options && typeof options.resumeCommittedBatch === 'boolean' - ? options.resumeCommittedBatch - : (task.retry_count || 0) > 0; - if (resumeCommittedBatch) { - const committed = await findCommittedBatch({ - db: ctx.db, - userId: task.user_id, - userKey, - taskUuid: task.uuid, - occurrenceMs: occurrenceMsOf(task), - }); - if (committed) return await redeliverCommittedBatch(task, ctx, userKey, decryptedPayload, committed); - } + // 内容已经落定的话只补推送。放在所有生成路径(agentic 与冻结 prompt)之前, + // 两条路落进 outbox 的批次都认。 + // + // 每次投递都查,不按「这是不是重试」预判:重试计数会被 PUT /update-message + // 和 renewTask 清零(那是它们该做的——修好 apiKey 的任务不该背着旧账),一 + // 旦拿它当「这次触发已经有内容了」的标记,用户在重试窗口里改一下任务,这次 + // 触发就会重新生成一整条,旧那批还躺在 outbox 里等客户端收,同一时刻冒出两 + // 份内容,LLM 的钱也白花一次。首次触发这一查是空的,多一次索引查询而已。 + const committed = await findCommittedBatch({ + db: ctx.db, + userId: task.user_id, + userKey, + taskUuid: task.uuid, + occurrenceMs: occurrenceMsOf(task), + }); + if (committed) return await redeliverCommittedBatch(task, ctx, userKey, decryptedPayload, committed); // Fire-time hooks: when the host configured onBeforeFire and the task // needs the LLM, offer the agentic path first. onBeforeFire → null @@ -739,12 +747,9 @@ export async function processMessagesByUuid(uuid, ctx, maxRetries = 2, userId, p return { success: false, error: { code: 'TASK_NOT_FOUND', message: '任务不存在或已处理' } }; } - // 请求内的重试不动任务行的 retry_count,「这是不是同一次触发的重试」得由这 - // 里显式告诉 processSingleMessage:上一轮生成成功、只是推送失败的话,这一轮 - // 只补推送,不再把 LLM 跑一遍。 - const result = await processSingleMessage(task, ctx, masterKey, null, { - resumeCommittedBatch: retryCount > 0 || (task.retry_count || 0) > 0, - }); + // 上一轮生成成功、只是推送失败的话,processSingleMessage 会认出这次触发已 + // 经落定的批次,这一轮只补推送、不再把 LLM 跑一遍。 + const result = await processSingleMessage(task, ctx, masterKey, null); if (!result.success) { // 确定性失败不进重试:再跑两轮也是同一个错,白让调用方多等、白烧一整轮 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 631275e..2ef0904 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 @@ -1085,10 +1085,7 @@ async function deliverTasks(ctx, tasks) { isTaskCancelled: () => lease.lost, }, masterKey, - { userKey, payload: decryptedPayload }, - // 重试计数就记在这一列上:大于 0 说明这是同一次触发的重试,内容已经落 - // 进 outbox 的话只补推送、不再生成(见 redeliverCommittedBatch)。 - { resumeCommittedBatch: (task.retry_count || 0) > 0 } + { userKey, payload: decryptedPayload } ); } catch (error) { if (lease.lost) { diff --git a/packages/rei-standard-amsg/server/test/d1-adapter.test.mjs b/packages/rei-standard-amsg/server/test/d1-adapter.test.mjs index 264c7b4..091690b 100644 --- a/packages/rei-standard-amsg/server/test/d1-adapter.test.mjs +++ b/packages/rei-standard-amsg/server/test/d1-adapter.test.mjs @@ -74,6 +74,32 @@ test('updateTaskByUuid updates only pending rows and returns {uuid, updated_at}' assert.equal(await adapter.updateTaskByUuid('missing', USER, 'enc2'), null); }); +test('updateTaskByUuid:任务正被投递占着时改不动排期,正文照改', async () => { + // 投递收尾会按自己领取时看到的排期推进下一次,这期间写进去的新时刻随后就被 + // 盖掉——接口回了成功,用户的改期却什么也没留下。所以改排期这一下直接不做, + // 让调用方知道。正文只影响以后的触发,收尾不覆盖它,照常放行。 + const { adapter, db } = await freshAdapter(); + const row = await adapter.createTask(baseTask({ uuid: 'busy', next_send_at: '2026-01-01T00:00:00.000Z' })); + await adapter.updateTaskById(row.id, { lease_until: new Date(Date.now() + 90_000).toISOString() }); + + assert.equal( + await adapter.updateTaskByUuid('busy', USER, 'enc-new', { next_send_at: '2027-01-01T00:00:00.000Z' }), + null, + '正在投递时不接受改排期' + ); + assert.equal(readRow(db, 'busy').next_send_at, '2026-01-01T00:00:00.000Z'); + + assert.ok( + await adapter.updateTaskByUuid('busy', USER, 'enc-content-only'), + '只改正文不受租约影响' + ); + + // 租约过期(投递那边没了)之后照常能改。 + await adapter.updateTaskById(row.id, { lease_until: new Date(Date.now() - 1000).toISOString() }); + assert.ok(await adapter.updateTaskByUuid('busy', USER, 'enc3', { next_send_at: '2027-01-01T00:00:00.000Z' })); + assert.equal(readRow(db, 'busy').next_send_at, '2027-01-01T00:00:00.000Z'); +}); + // lease_until 是占位用的内部列,适配器的取任务方法不返回它,测试直接读行。 function readRow(db, uuid) { return db._raw.prepare('SELECT next_send_at, lease_until FROM scheduled_messages WHERE uuid = ?').get(uuid); diff --git a/packages/rei-standard-amsg/server/test/pg-neon-adapter.test.mjs b/packages/rei-standard-amsg/server/test/pg-neon-adapter.test.mjs index 2ddda42..b1942a2 100644 --- a/packages/rei-standard-amsg/server/test/pg-neon-adapter.test.mjs +++ b/packages/rei-standard-amsg/server/test/pg-neon-adapter.test.mjs @@ -222,6 +222,53 @@ async function d1TaskSelectColumns() { return taskSelectColumns(adapter, () => calls.at(-1).sql); } +for (const backend of BACKENDS) { + test(`${backend.name}: updateTaskByUuid 改排期时带租约门,只改正文时不带`, async () => { + // 投递收尾会按它领取时看到的排期推进下一次,投递期间写进去的新时刻随后就被 + // 盖掉。所以改排期这一下要在 SQL 里挡住(回 null 让调用方知道),正文这类 + // 字段不挡——收尾本来就不覆盖它们。 + const { adapter, calls } = backend.make(() => [{ uuid: 'u', updated_at: 'now' }]); + + await adapter.updateTaskByUuid('u', 'user-1', 'cipher', { next_send_at: '2027-01-01T00:00:00.000Z' }); + assert.match( + flat(calls.at(-1).text), + /lease_until IS NULL OR lease_until <= NOW\(\)/i, + '改排期必须带租约门' + ); + + await adapter.updateTaskByUuid('u', 'user-1', 'cipher', { retry_count: 0 }); + assert.doesNotMatch(flat(calls.at(-1).text), /lease_until/i, '不改排期的写入不该被租约挡住'); + }); +} + +for (const backend of BACKENDS) { + test(`${backend.name}: createTaskSuperseding 删旧 + 建新落在同一条语句里`, async () => { + // 分成两次查询的话,建新失败时旧任务已经删掉,接口却回失败——客户端以为旧 + // 任务还在,其实已经没了。一条语句里的数据修改型 CTE 是一个隐式事务, + // INSERT 抛错时 DELETE 跟着回滚。 + const { adapter, calls } = backend.make(() => [{ + id: 5, uuid: 'new-uuid', next_send_at: '2026-01-01T00:00:00.000Z', + status: 'pending', created_at: '2026-01-01T00:00:00.000Z', superseded: true + }]); + + const row = await adapter.createTaskSuperseding({ + user_id: 'u-1', + uuid: 'new-uuid', + encrypted_payload: 'cipher', + next_send_at: '2026-01-01T00:00:00.000Z', + message_type: 'fixed' + }, 'old-uuid'); + + assert.equal(calls.length, 1, '删旧和建新必须在同一条语句里,不能各发一次'); + const { text, params } = calls[0]; + assert.match(flat(text), /DELETE FROM scheduled_messages/i); + assert.match(flat(text), /INSERT INTO scheduled_messages/i); + assert.deepEqual(params, ['u-1', 'new-uuid', 'cipher', '2026-01-01T00:00:00.000Z', 'fixed', 'old-uuid']); + assert.equal(row.superseded, true); + assert.equal(row.id, 5); + }); +} + test('投递链路和读接口各自的列集,三个适配器逐字一致', async () => { const pg = recordingPg(countAwareRows); const neon = recordingNeon(countAwareRows); diff --git a/packages/rei-standard-amsg/server/test/processor-multipart-pacing.test.mjs b/packages/rei-standard-amsg/server/test/processor-multipart-pacing.test.mjs index f7c7440..d86cfa0 100644 --- a/packages/rei-standard-amsg/server/test/processor-multipart-pacing.test.mjs +++ b/packages/rei-standard-amsg/server/test/processor-multipart-pacing.test.mjs @@ -9,7 +9,7 @@ import { describe, it } from 'node:test'; import assert from 'node:assert/strict'; -import { DEFAULT_MULTIPART_TTL_MS } from '@rei-standard/amsg-shared'; +import { buildMultipartPushPayloads, DEFAULT_MULTIPART_TTL_MS, MESSAGE_KIND } from '@rei-standard/amsg-shared'; import { processSingleMessage } from '../src/server/lib/message-processor.js'; import { MAX_PUSH_PAYLOAD_BYTES, measurePushPayload } from '../src/server/lib/webpush-webcrypto.js'; import { deriveUserEncryptionKey, encryptForStorage } from '../src/server/lib/encryption.js'; @@ -180,6 +180,41 @@ describe('maxChunkBytes 的上限校验', () => { assert.ok(withinLimit, `分片信封 ${bytes} 字节,超过单条 push 的 ${MAX_PUSH_PAYLOAD_BYTES} 字节上限`); } }); + + it('错误信息给出的最大值:每一种 messageKind 都装得下', async () => { + // 分片信封里原样带着原消息的 messageKind,名字越长信封越大。校验量开销用 + // 的探针只有一种 kind,要是挑的不是最长的那个,照着错误信息里的最大值去配 + // 的人,content / reasoning 的分片发得出去,tool_request 这类长名字的分片 + // 每一片都超出上限、被推送服务整批拒收。 + const { result } = await deliverWithClock({ + reasoningChars: 20_000, + multipart: { maxChunkBytes: 3000 }, + }); + const maxAllowed = Number(/最大 (\d+)/.exec(result.error || '')?.[1]); + assert.ok(Number.isInteger(maxAllowed) && maxAllowed > 0, `错误信息里没读到最大值:${result.error}`); + + // 切满 100 片,让 index / total 都占到 3 位数——分片默认最多 128 片,位数 + // 正是校验按最坏情况预留的那部分,只切一片量不出真正的上限。 + const CHUNKS = 100; + for (const kind of Object.values(MESSAGE_KIND)) { + const chunks = buildMultipartPushPayloads( + { messageKind: kind }, + { + serializedPayload: 'x'.repeat(maxAllowed * CHUNKS), + maxChunkBytes: maxAllowed, + ttlMs: DEFAULT_MULTIPART_TTL_MS, + } + ); + assert.equal(chunks.length, CHUNKS); + const worst = chunks[chunks.length - 1]; + const { bytes, withinLimit } = measurePushPayload(JSON.stringify(worst)); + assert.ok( + withinLimit, + `messageKind=${kind} 按放行的最大值 ${maxAllowed} 切出来的信封 ${bytes} 字节,` + + `超过单条 push 的 ${MAX_PUSH_PAYLOAD_BYTES} 字节上限` + ); + } + }); }); describe('宿主配的分片限额传得到发送端', () => { diff --git a/packages/rei-standard-amsg/server/test/spend-guards.test.mjs b/packages/rei-standard-amsg/server/test/spend-guards.test.mjs index 837a9db..bcefced 100644 --- a/packages/rei-standard-amsg/server/test/spend-guards.test.mjs +++ b/packages/rei-standard-amsg/server/test/spend-guards.test.mjs @@ -555,6 +555,39 @@ describe('推送 5xx 一次后恢复:重试只补推送,不重新生成', () assert.equal(llm.calls.length, 1); assert.deepEqual(webpush.received.map((p) => p.message), ['甲1。', '乙1。']); }); + + test('用户在重试窗口里改了任务:补推那一跳照样只补推送', async () => { + const adapter = await makeAdapter(); + await seedTask(adapter, { uuid: 'edited', payload: LLM_PAYLOAD }); + const webpush = scriptedWebpush([1]); // 第一条推送 503,之后恢复 + const ctx = tickCtx(adapter, webpush); + const llm = stubLlm(); + try { + const first = await runScheduledTick(ctx); + assert.equal(first.failedCount, 1, '第一跳推送失败,这一批已经落进 outbox'); + + // 用户这时改了个跟投递无关的字段(联系人名字)。PUT /update-message 落到 + // 行上的就是下面这几样:新密文 + 重试计数清零 + 退避放掉(见 + // handlers/update-message.js)。计数清零是它该做的——修好 apiKey 的任务不 + // 该背着旧账——但这次触发的内容已经落定了,不能因此重新生成一遍。 + const userKey = await deriveUserEncryptionKey(USER, MASTER_KEY); + const edited = await encryptForStorage( + JSON.stringify({ recurrenceType: 'none', ...LLM_PAYLOAD, contactName: '改过的名字' }), + userKey + ); + await adapter.updateTaskByUuid('edited', USER, edited, { retry_count: 0, retry_after: null }); + + const second = await runScheduledTick(ctx); + assert.equal(second.successCount, 1); + assert.equal(second.details.redeliveredTasks.length, 1); + } finally { + llm.restore(); + } + assert.equal(llm.calls.length, 1, '改过任务也不该把这次触发重新生成一遍'); + // 设备上只有第一次生成的那一份,不会冒出第二份内容。 + assert.deepEqual(webpush.received.map((p) => p.message), ['甲1。', '乙1。']); + assert.equal((await outboxOf(adapter, 'edited')).length, 2, 'outbox 里只有这一批'); + }); }); describe('没落进 outbox 的批次(没有收件箱 / 落行失败)', () => { diff --git a/packages/rei-standard-amsg/server/test/update-message-fields.test.mjs b/packages/rei-standard-amsg/server/test/update-message-fields.test.mjs index c242a3f..aa9fc8c 100644 --- a/packages/rei-standard-amsg/server/test/update-message-fields.test.mjs +++ b/packages/rei-standard-amsg/server/test/update-message-fields.test.mjs @@ -136,3 +136,43 @@ test('PUT /update-message 的 updatedFields 只报真正落库的字段', async assert.equal(stored.contactName, '改过的名字'); assert.equal('contactname' in stored, false); }); + +// 任务正在投递的那几十秒里用户改排期:写进去也会被这次投递的收尾盖掉(收尾按 +// 它领取时看到的排期推进下一次,一次性任务干脆标成已发送)。所以这一下要如实 +// 回报,别让客户端拿着一个「改成功了」的 200 走。 +test('PUT /update-message:任务正在投递时改排期回 409 TASK_IN_FLIGHT', async () => { + const { server, db } = await makeServer(); + const uuid = await scheduleFixed(server); + const row = await db.getTaskByUuid(uuid, USER); + await db.updateTaskById(row.id, { lease_until: new Date(Date.now() + 90_000).toISOString() }); + + const rescheduled = await server.handlers.updateMessage.PUT( + `/update-message?id=${uuid}`, + HEADERS, + await encBody({ nextSendAt: '2999-06-01T00:00:00.000Z' }) + ); + assert.equal(rescheduled.status, 409); + assert.equal(rescheduled.body.error.code, 'TASK_IN_FLIGHT'); + assert.equal( + (await db.getTaskByUuid(uuid, USER)).next_send_at, + '2999-01-01T00:00:00.000Z', + '排期没动' + ); + + // 只改正文不碰排期的照常成功:收尾那边本来就不会覆盖用户刚保存的正文。 + const edited = await server.handlers.updateMessage.PUT( + `/update-message?id=${uuid}`, + HEADERS, + await encBody({ contactName: '投递中也能改的名字' }) + ); + assert.equal(edited.status, 200, JSON.stringify(edited.body)); + assert.equal((await readStoredPayload(db, uuid)).contactName, '投递中也能改的名字'); + + // 任务本来就不在时,报的还是原来那句「可能已被修改或删除」。 + const missing = await server.handlers.updateMessage.PUT( + '/update-message?id=99999999-8888-4777-8666-555555555551', + HEADERS, + await encBody({ nextSendAt: '2999-06-01T00:00:00.000Z' }) + ); + assert.equal(missing.status, 404); +}); diff --git a/packages/rei-standard-amsg/server/test/upstream-additions.test.mjs b/packages/rei-standard-amsg/server/test/upstream-additions.test.mjs index 5c59622..c5b2462 100644 --- a/packages/rei-standard-amsg/server/test/upstream-additions.test.mjs +++ b/packages/rei-standard-amsg/server/test/upstream-additions.test.mjs @@ -210,6 +210,50 @@ describe('schedule-message immediate / supersede', () => { assert.equal((await res2.json()).data.superseded, false); }); + test('supersedesUuid:适配器没有原子形态时,建新失败也不能把旧任务白删', async () => { + // 自定义适配器(没实现 createTaskSuperseding)走两步退路。建新这一步撞了 + // uuid,接口回 409 让客户端换个 uuid 重来——这时旧任务必须还在,客户端才有 + // 东西可顶替。反过来「先删后建」的话,失败的请求会把旧任务顺手带走。 + const d1 = createTestD1(); + const worker = createSingleUserCloudflareWorker((env) => { + const base = createD1Adapter(env.DB); + return { + db: new Proxy(base, { + get(target, prop) { + if (prop === 'createTaskSuperseding') return undefined; + if (prop === 'createTask') { + return async () => { throw new Error('duplicate key value violates unique constraint'); }; + } + const value = target[prop]; + return typeof value === 'function' ? value.bind(target) : value; + }, + }), + masterKey: MASTER_KEY, + vapid: VAPID, + webpush: { async sendNotification() {} }, + }; + }); + const env = { DB: d1 }; + await worker.fetch(new Request('https://w.dev/init-tenant', { method: 'POST' }), env); + const adapter = createD1Adapter(d1); + const oldUuid = 'aaaaaaaa-bbbb-4ccc-8ddd-eeeeeeeeee77'; + await seed(adapter, { uuid: oldUuid, nextSendAt: new Date(Date.now() + 3600_000).toISOString() }); + + const res = await worker.fetch(new Request('https://w.dev/schedule-message', { + method: 'POST', headers: ENC_HEADERS, + body: await encBody({ + contactName: 'Rei', messageType: 'fixed', userMessage: 'v2', + firstSendTime: new Date(Date.now() + 3600_000).toISOString(), + uuid: '77777777-6666-4555-8444-333333333332', + supersedesUuid: oldUuid + }) + }), env); + + assert.equal(res.status, 409); + assert.equal((await res.json()).error.code, 'TASK_UUID_CONFLICT'); + assert.ok(await adapter.getTaskByUuid(oldUuid, USER), '请求失败了,旧任务得还在'); + }); + test('supersedesUuid 撞新任务自己的 uuid 被拒', async () => { const d1 = createTestD1(); const worker = makeWorker(); @@ -444,6 +488,31 @@ describe('fire ctx cancelTask / renewTask', () => { assert.equal(row.next_send_at, newAt); assert.match(String(tooSoonError), /至少要比现在晚/); }); + + test('renewTask 碰上正在投递的那条任务:报 in_flight,不假装改成功', async () => { + // 两条任务差不多同时到点,其中一条的 hook 去给另一条改期,而那一条此刻正被 + // 一次投递占着(租约还没到期)。硬写进去的新时刻会被那次投递的收尾按它自己 + // 领取时的排期盖掉,所以这里根本不写,直接如实回报。 + const { adapter } = await freshAdapter(); + const busyUuid = '31313131-4242-4535-8626-717171717171'; + const busyAt = new Date(Date.now() + 3600_000).toISOString(); + await seed(adapter, { uuid: busyUuid, nextSendAt: busyAt }); + const busyRow = await adapter.getTaskByUuid(busyUuid, USER); + await adapter.updateTaskById(busyRow.id, { lease_until: new Date(Date.now() + 90_000).toISOString() }); + + const res = await fireCtxOf(adapter, async (ctx) => { + assert.deepEqual( + await ctx.renewTask(busyUuid, new Date(Date.now() + 2 * 3600_000).toISOString()), + { renewed: false, reason: 'in_flight' } + ); + }); + assert.equal(res.successCount, 1); + assert.equal( + (await adapter.getTaskByUuid(busyUuid, USER)).next_send_at, + busyAt, + '排期一个字都不该动' + ); + }); }); // ─── message_outbox ──────────────────────────────────────────────────────── From 14406f1327756b92c36339124fade72fc6e71271 Mon Sep 17 00:00:00 2001 From: Tosd0 <65720409+Tosd0@users.noreply.github.com> Date: Wed, 23 Sep 2026 11:11:04 +0800 Subject: [PATCH 2/4] =?UTF-8?q?fix(amsg-instant):=20=E5=88=86=E7=89=87?= =?UTF-8?q?=E4=B8=8A=E9=99=90=E6=A0=A1=E9=AA=8C=E6=8C=89=E6=9C=80=E9=95=BF?= =?UTF-8?q?=20messageKind=20=E7=AE=97=EF=BC=8CX-Client-Token=20=E6=A0=A1?= =?UTF-8?q?=E9=AA=8C=E5=90=88=E5=B9=B6=E6=88=90=E4=B8=80=E4=BB=BD?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - 探针之前固定用 reasoning,比最长的 tool_request 短 3 字节:把 maxChunkBytes 配 在上限附近(例如照着报错里建议的最大值配)的部署,tool_request 的分片每片都超 出单条 push 的明文上限、被推送服务拒收。现在从 shared 的枚举里现取最长的那个。 - 导出的 validateClientAuth 与 handler 内部的校验此前是两份内容相同的实现,现在 共用同一份 checkClientToken。两个函数的签名和行为不变,handler 仍用启动时编好 的 token 字节。 Claude-Session: https://claude.ai/code/session_017fqgW4ShQ4yjJJwqAiJrcT --- .changeset/amsg-instant-chunk-probe-auth.md | 9 +++ .../rei-standard-amsg/instant/src/index.js | 64 +++++++-------- .../instant/src/multipart.js | 18 +++-- .../instant/test/handler.test.mjs | 80 +++++++++++++++++++ 4 files changed, 132 insertions(+), 39 deletions(-) create mode 100644 .changeset/amsg-instant-chunk-probe-auth.md diff --git a/.changeset/amsg-instant-chunk-probe-auth.md b/.changeset/amsg-instant-chunk-probe-auth.md new file mode 100644 index 0000000..67e9e26 --- /dev/null +++ b/.changeset/amsg-instant-chunk-probe-auth.md @@ -0,0 +1,9 @@ +--- +"@rei-standard/amsg-instant": patch +--- + +分片大小上限的校验按最长的 messageKind 算;X-Client-Token 的校验只留一份实现 + +`multipart.maxChunkBytes` 的上限校验用一条探针量信封开销,探针之前固定用 `reasoning`,比最长的 `tool_request` 短 3 字节。把 `maxChunkBytes` 配在上限附近(例如照着 `createInstantHandler` 抛错时建议的最大值去配)的部署,`content` / `reasoning` 的分片发得出去,`tool_request` 的分片每一片都超出单条 push 的明文上限、被推送服务拒收。现在按真实取值里最长的那个量,报出来的最大值对所有类型都成立。 + +导出的 `validateClientAuth` 与 `createInstantHandler` 内部的校验现在共用同一份实现(存在性检查、常时比较、401 响应体都是同一处),两边不会再各自漂。函数签名和行为不变,handler 仍然用启动时编好的 token 字节,没有每请求重编。 diff --git a/packages/rei-standard-amsg/instant/src/index.js b/packages/rei-standard-amsg/instant/src/index.js index aa25ccc..7c522c6 100644 --- a/packages/rei-standard-amsg/instant/src/index.js +++ b/packages/rei-standard-amsg/instant/src/index.js @@ -935,12 +935,35 @@ function isVapidConfigValid(vapid) { return true; } +/** + * X-Client-Token 校验的唯一实现:presence 检查 + 常时比较。 + * + * handler 内部(verifyClientToken)和导出给宿主的 validateClientAuth 都走这里, + * 两条路径的 401 响应体因此不会各自漂。期望值收的是编码好的字节,handler 在 + * 启动时编一次就够,不用每个请求再编一遍。 + * + * @param {Request} request + * @param {Uint8Array} expectedBytes - 期望的 token,UTF-8 字节 + * @returns {Object | null} 校验不过返回 401 响应体,通过返回 null + */ +function checkClientToken(request, expectedBytes) { + const received = getHeader(request, 'x-client-token'); + if (!received) { + return { success: false, error: { code: 'INVALID_CLIENT_TOKEN', message: '缺少 X-Client-Token' } }; + } + if (!timingSafeEqualBytes(utf8(received), expectedBytes)) { + return { success: false, error: { code: 'INVALID_CLIENT_TOKEN', message: 'X-Client-Token 无效' } }; + } + return null; +} + /** * X-Client-Token 的独立校验口(导出)。 * - * createInstantHandler 内部走的就是这一套(presence 检查 + 常时比较)。导出 - * 是给「在同一个 worker 里挂自己路由」的宿主用的:那些路由的鉴权语义应该和 - * 本 handler 完全一致,宿主此前只能照抄内部实现——抄的那份不会跟着上游修。 + * createInstantHandler 内部走的就是这一套(presence 检查 + 常时比较),两边共用 + * 同一份 checkClientToken。导出是给「在同一个 worker 里挂自己路由」的宿主用的: + * 那些路由的鉴权语义应该和本 handler 完全一致,宿主此前只能照抄内部实现——抄的 + * 那份不会跟着上游修。 * * @param {Request} request * @param {string} expectedToken - 部署配置里的共享密钥(AMSG_CLIENT_TOKEN) @@ -952,40 +975,13 @@ function isVapidConfigValid(vapid) { export function validateClientAuth(request, expectedToken) { const expected = expectedToken ? String(expectedToken) : ''; if (!expected) return { ok: true }; - const received = getHeader(request, 'x-client-token'); - if (!received) { - return { - ok: false, - status: 401, - body: { success: false, error: { code: 'INVALID_CLIENT_TOKEN', message: '缺少 X-Client-Token' } } - }; - } - if (!timingSafeEqualBytes(utf8(received), utf8(expected))) { - return { - ok: false, - status: 401, - body: { success: false, error: { code: 'INVALID_CLIENT_TOKEN', message: 'X-Client-Token 无效' } } - }; - } - return { ok: true }; + const failureBody = checkClientToken(request, utf8(expected)); + return failureBody ? { ok: false, status: 401, body: failureBody } : { ok: true }; } function verifyClientToken(request, expectedBytes, respond) { - const received = getHeader(request, 'x-client-token'); - if (!received) { - return respond(401, { - success: false, - error: { code: 'INVALID_CLIENT_TOKEN', message: '缺少 X-Client-Token' } - }); - } - const receivedBytes = utf8(received); - if (!timingSafeEqualBytes(receivedBytes, expectedBytes)) { - return respond(401, { - success: false, - error: { code: 'INVALID_CLIENT_TOKEN', message: 'X-Client-Token 无效' } - }); - } - return null; + const failureBody = checkClientToken(request, expectedBytes); + return failureBody ? respond(401, failureBody) : null; } async function verifyBearerToken(request, signingKey, respond) { diff --git a/packages/rei-standard-amsg/instant/src/multipart.js b/packages/rei-standard-amsg/instant/src/multipart.js index 0a56781..2ff4f2a 100644 --- a/packages/rei-standard-amsg/instant/src/multipart.js +++ b/packages/rei-standard-amsg/instant/src/multipart.js @@ -11,7 +11,7 @@ export { buildMultipartPushPayloads, } from '@rei-standard/amsg-shared'; -import { buildMultipartPushPayloads } from '@rei-standard/amsg-shared'; +import { buildMultipartPushPayloads, MESSAGE_KIND } from '@rei-standard/amsg-shared'; // ─── maxChunkBytes 的上限校验 ────────────────────────────────────────── // @@ -34,6 +34,14 @@ const MAX_PUSH_PAYLOAD_BYTES = WEB_PUSH_MAX_BODY_BYTES - WEB_PUSH_ENCRYPTION_OVE const PROBE_BYTE_ENCODER = new TextEncoder(); +/** + * 真实取值里最长的 messageKind。分片信封把原消息的 messageKind 原样写进 + * `originalMessageKind`,所以信封开销的最坏情况由最长的那个取值决定;从 shared + * 的枚举里现取,以后加了更长的类型能自动跟上。 + */ +const LONGEST_MESSAGE_KIND = Object.values(MESSAGE_KIND) + .reduce((longest, kind) => (kind.length > longest.length ? kind : longest), ''); + /** * n 字节编成 base64url(不带 padding)后的字符数。 * @@ -54,9 +62,9 @@ function base64UrlLength(n) { * 到每次投递才炸。 * * 上限不写死成常量,而是现算:拿真实的 buildMultipartPushPayloads 造一片最小 - * 探针量出信封开销(id / createdAt 等定宽字段取的就是真实值),再把 index / - * total 的位数按 maxChunks 补足到最坏情况——shared 那边信封格式变了,这里跟着 - * 变,不会留下一个过时的魔数。 + * 探针量出信封开销(id / createdAt 等定宽字段取的就是真实值,messageKind 取真实 + * 取值里最长的那个),再把 index / total 的位数按 maxChunks 补足到最坏情况—— + * shared 那边信封格式变了,这里跟着变,不会留下一个过时的魔数。 * * @param {{ maxChunkBytes: number, maxChunks: number, ttlMs: number }} resolved */ @@ -64,7 +72,7 @@ export function assertChunkBytesFitPushLimit({ maxChunkBytes, maxChunks, ttlMs } // 最小探针:3 字节原文 → base64url 后恰好 4 字符,信封开销 = 总长 - 4。 const PROBE_CHUNK_BYTES = 3; const [probe] = buildMultipartPushPayloads( - { messageKind: 'reasoning' }, + { messageKind: LONGEST_MESSAGE_KIND }, { serializedPayload: 'x'.repeat(PROBE_CHUNK_BYTES), maxChunkBytes: PROBE_CHUNK_BYTES, ttlMs } ); // index / total 在真实批次里最多到 maxChunks(探针里各只有 1 位)。 diff --git a/packages/rei-standard-amsg/instant/test/handler.test.mjs b/packages/rei-standard-amsg/instant/test/handler.test.mjs index cd60d09..ecee4f8 100644 --- a/packages/rei-standard-amsg/instant/test/handler.test.mjs +++ b/packages/rei-standard-amsg/instant/test/handler.test.mjs @@ -1,11 +1,15 @@ import { describe, it, before } from 'node:test'; import assert from 'node:assert/strict'; +import { MESSAGE_KIND } from '@rei-standard/amsg-shared'; + import { createInstantHandler, sendPushWithMaybeBlob, validateInstantPayload, + validateClientAuth, } from '../src/index.js'; +import { assertChunkBytesFitPushLimit, buildMultipartPushPayloads } from '../src/multipart.js'; import { generateTestVapid, generateTestSubscription, @@ -513,6 +517,42 @@ describe('createInstantHandler — clientToken', () => { assert.equal(doneReceived, true); await waitForPushCalls(router, 1); }); + + // 导出 validateClientAuth 是为了让宿主的鉴权跟着上游一起改,所以它和 handler + // 内部那条路径必须给出一模一样的 401 —— 谁单独改了都算回归。 + it('validateClientAuth 与 handler 内部路径给出同一个 401 响应体', async () => { + const clientToken = 'shared-secret-xyz'; + const handler = createInstantHandler({ vapid, clientToken }); + + for (const { name, headers } of [ + { name: '缺头', headers: {} }, + { name: '头不匹配', headers: { 'x-client-token': 'wrong-token' } }, + ]) { + const handlerRes = await handler(makeRequest({ body: makeValidPayload(), headers })); + assert.equal(handlerRes.status, 401, name); + const handlerBody = await handlerRes.json(); + + const exported = validateClientAuth(makeRequest({ body: makeValidPayload(), headers }), clientToken); + assert.equal(exported.ok, false, name); + assert.equal(exported.status, handlerRes.status, name); + assert.deepEqual(exported.body, handlerBody, `${name}:两条路径的 401 响应体必须完全一致`); + } + }); + + it('validateClientAuth 没配共享密钥时一律放行(与 handler 的开放模式一致)', () => { + const request = makeRequest({ body: makeValidPayload() }); + for (const expected of [undefined, null, '']) { + assert.deepEqual(validateClientAuth(request, expected), { ok: true }); + } + // 配了密钥 + 头匹配也放行,且不返回多余字段。 + assert.deepEqual( + validateClientAuth( + makeRequest({ body: makeValidPayload(), headers: { 'x-client-token': 'shared-secret-xyz' } }), + 'shared-secret-xyz' + ), + { ok: true } + ); + }); }); // ─── Handler: happy path & push delivery ────────────────────────────── @@ -680,6 +720,46 @@ describe('createInstantHandler — multipart.maxChunkBytes 上限', () => { createInstantHandler({ vapid, multipart: { maxChunkBytes: 900, maxChunks: 32, ttlMs: 120_000 } }); }); + it('校验放行的最大值:最长 messageKind 切出来的分片信封仍在单条 push 明文上限内', () => { + // 推送服务限的是加密后 body 4096 字节,aes128gcm 固定开销 103 字节,明文 + // 只剩 3993 字节(推导见 src/multipart.js)。校验放行一个值,就等于对部署 + // 承诺「这个值切出来的每一片都发得出去」——对最长的 messageKind 也得成立。 + const MAX_PUSH_PAYLOAD_BYTES = 4096 - 103; + const maxChunks = 10; + const ttlMs = 60_000; + + // 当前配置下的最大值由校验自己说了算:故意配超,从报错里读出来。 + let thrown; + try { + assertChunkBytesFitPushLimit({ maxChunkBytes: MAX_PUSH_PAYLOAD_BYTES, maxChunks, ttlMs }); + } catch (err) { + thrown = err; + } + assert.ok(thrown, '配到明文上限必然超限'); + const maxAllowed = Number(thrown.message.match(/最大 (\d+)/)[1]); + assert.ok(Number.isInteger(maxAllowed) && maxAllowed > 0); + assertChunkBytesFitPushLimit({ maxChunkBytes: maxAllowed, maxChunks, ttlMs }); + + // 最坏的一片:originalMessageKind 取最长的取值,index / total 都满位。 + const longestKind = Object.values(MESSAGE_KIND) + .reduce((longest, kind) => (kind.length > longest.length ? kind : longest), ''); + const parts = buildMultipartPushPayloads( + { messageKind: longestKind }, + { serializedPayload: 'x'.repeat(maxAllowed * maxChunks), maxChunkBytes: maxAllowed, ttlMs } + ); + assert.equal(parts.length, maxChunks, 'index / total 要真的走到 maxChunks 的位数'); + + const encoder = new TextEncoder(); + for (const part of parts) { + const envelopeBytes = encoder.encode(JSON.stringify(part)).byteLength; + assert.ok( + envelopeBytes <= MAX_PUSH_PAYLOAD_BYTES, + `messageKind=${longestKind} 第 ${part.multipart.index} 片信封 ${envelopeBytes} 字节,` + + `超过单条 push 明文上限 ${MAX_PUSH_PAYLOAD_BYTES} 字节` + ); + } + }); + it('直接调 sendPushWithMaybeBlob(自己攒 ctx)同样被拦下,一片都不发出', async () => { const router = llmRouter('unused.'); const oversized = { messageKind: 'reasoning', reasoningContent: 'x'.repeat(5000) }; From 0b13bda754ff0c8caaae358029281618215a904d Mon Sep 17 00:00:00 2001 From: Tosd0 <65720409+Tosd0@users.noreply.github.com> Date: Wed, 23 Sep 2026 11:11:04 +0800 Subject: [PATCH 3/4] =?UTF-8?q?fix(amsg-sw):=20=E6=94=BE=E5=BC=83=E5=88=86?= =?UTF-8?q?=E7=89=87=E6=97=B6=E6=8A=A5=E7=9A=84=E5=8E=9F=E5=9B=A0=E4=B8=8D?= =?UTF-8?q?=E5=86=8D=E8=A2=AB=E5=A2=93=E7=A2=91=E5=86=99=E5=A4=B1=E8=B4=A5?= =?UTF-8?q?=E7=9B=96=E6=88=90=E5=AD=98=E5=82=A8=E6=95=85=E9=9A=9C?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 四条当场放弃的路径(窗口走完、分片说法冲突、超 maxTotalBytes、拼不回来)此前没接 收尾抛出的错误,墓碑写失败时异常冒到外层兜底,页面收到的原因一律变成 storage-failed,本来那条具体原因丢了。现在收尾失败按住不外抛、只留日志,仍按本来 那条原因广播一次事件——结论此前已经记进内存兜底表,后续分片照样收不进来。 Claude-Session: https://claude.ai/code/session_017fqgW4ShQ4yjJJwqAiJrcT --- .../amsg-sw-multipart-give-up-reason.md | 9 +++ packages/rei-standard-amsg/sw/src/index.js | 58 +++++++++++++++---- .../multipart-failure-visibility.test.mjs | 50 ++++++++++++++++ 3 files changed, 107 insertions(+), 10 deletions(-) create mode 100644 .changeset/amsg-sw-multipart-give-up-reason.md diff --git a/.changeset/amsg-sw-multipart-give-up-reason.md b/.changeset/amsg-sw-multipart-give-up-reason.md new file mode 100644 index 0000000..74664e1 --- /dev/null +++ b/.changeset/amsg-sw-multipart-give-up-reason.md @@ -0,0 +1,9 @@ +--- +"@rei-standard/amsg-sw": patch +--- + +放弃一条分片消息时,报的原因不再被墓碑写失败盖成「存储故障」 + +分片重组有四条当场放弃的路径:重组窗口走完、同一个 id 的分片对 total / encoding 各说各话、累计字节超过 `maxTotalBytes`、分片齐了却拼不回原 payload。这几条都会先写「这个 id 到此为止」的墓碑再广播 `MULTIPART_EXPIRED`,而墓碑写失败(存储压力大的时候正是它容易失败)时异常会冒到外层兜底,页面收到的原因于是一律变成 `storage-failed`,本来那条具体原因(`chunk-conflict` / `size-limit-exceeded` / `restore-failed` / `ttl-expired`)丢了——宿主拿这个原因做诊断的话,压力下会把发送端发冲突分片这类问题统一误读成存储故障。 + +现在这四条路径的收尾失败按住不外抛(结论此前已经记进内存兜底表,后续分片和推送服务的重投照样进不来),仍然用本来那条原因广播一次事件,收尾失败另留一条日志。页面收到的事件次数不变。 diff --git a/packages/rei-standard-amsg/sw/src/index.js b/packages/rei-standard-amsg/sw/src/index.js index 51a7c79..2529cff 100644 --- a/packages/rei-standard-amsg/sw/src/index.js +++ b/packages/rei-standard-amsg/sw/src/index.js @@ -1294,8 +1294,9 @@ async function acceptMultipartChunkInternal(sw, normalized, options) { '[rei-standard-amsg-sw] multipart reassembly window elapsed; giving up on this multipart id:', { id: existing.id, total: existing.total, receivedCount: existing.receivedCount } ); - await settleMultipartId(existing, existing.total, options); - await dispatchMultipartExpired(sw, existing); + await giveUpMultipartId( + sw, existing, existing.total, options, MULTIPART_FAILURE_REASON.TTL_EXPIRED + ); return null; } @@ -1321,8 +1322,13 @@ async function acceptMultipartChunkInternal(sw, normalized, options) { incoming: { total: normalized.total, encoding: normalized.encoding }, } ); - await settleMultipartId(base, Math.max(base.total, normalized.total), options); - await dispatchMultipartExpired(sw, base, MULTIPART_FAILURE_REASON.CHUNK_CONFLICT); + await giveUpMultipartId( + sw, + base, + Math.max(base.total, normalized.total), + options, + MULTIPART_FAILURE_REASON.CHUNK_CONFLICT + ); return null; } @@ -1344,8 +1350,9 @@ async function acceptMultipartChunkInternal(sw, normalized, options) { maxTotalBytes: options.maxTotalBytes, } ); - await settleMultipartId(base, base.total, options); - await dispatchMultipartExpired(sw, base, MULTIPART_FAILURE_REASON.SIZE_LIMIT_EXCEEDED); + await giveUpMultipartId( + sw, base, base.total, options, MULTIPART_FAILURE_REASON.SIZE_LIMIT_EXCEEDED + ); return null; } @@ -1370,8 +1377,9 @@ async function acceptMultipartChunkInternal(sw, normalized, options) { '[rei-standard-amsg-sw] multipart restore failed; giving up on this multipart id:', error ); - await settleMultipartId(base, base.total, options); - await dispatchMultipartExpired(sw, base, MULTIPART_FAILURE_REASON.RESTORE_FAILED); + await giveUpMultipartId( + sw, base, base.total, options, MULTIPART_FAILURE_REASON.RESTORE_FAILED + ); return null; } @@ -1388,6 +1396,34 @@ async function acceptMultipartChunkInternal(sw, normalized, options) { return restored; } +/** + * 中途放弃一条 multipart id 的统一出口:先收尾(见 {@link settleMultipartId}), + * 再按 `reason` 告诉页面这条为什么不用再等了。 + * + * 收尾里的持久墓碑那笔写可能失败,这里把它按住不外抛:那时候「这个 id 已有结 + * 论」已经记进了内存兜底表,后面的分片和推送服务的重投照样收不进来,缺的只是 + * 墓碑那笔写。让错误往上冒的话,外层兜底(见 acceptMultipartChunkSafely)会接 + * 手,把页面收到的原因改写成笼统的 STORAGE_FAILED —— 宿主拿这个原因做诊断,就 + * 分不出到底是发送端一直发冲突分片、还是存储真的挂了。 + * + * @param {ServiceWorkerGlobalScope} sw + * @param {{ id: string, total?: number, ttlMs?: number }} record + * @param {number} total - 要清掉的分片数(冲突时取两边的较大值,别漏删) + * @param {{ ttlMs: number }} options + * @param {string} reason - {@link MULTIPART_FAILURE_REASON} 之一,走到这一步的那条路。 + */ +async function giveUpMultipartId(sw, record, total, options, reason) { + try { + await settleMultipartId(record, total, options); + } catch (error) { + console.error( + `[rei-standard-amsg-sw] multipart cleanup while giving up (${reason}) failed:`, + error + ); + } + await dispatchMultipartExpired(sw, record, reason); +} + /** * 这个 multipart id 到此为止:先写 done 墓碑,再清掉 pending 记录和已收的分片。 * 收齐还原了、和中途放弃了(分片对不上、超限、拼不回来),走的是同一套收尾。 @@ -1404,8 +1440,10 @@ async function acceptMultipartChunkInternal(sw, normalized, options) { * 墓碑比重组窗口活得久(两倍),推送服务重投旧分片时也不会再触发一次业务事件。 * * 持久墓碑本身写失败时,同样的结论会先落进内存兜底表(见 - * {@link memoryFallbackMultipartDone})再把错误往上抛:调用方各自的失败处理不变, - * 但「已有结论」这件事在 SW 存活期内不丢。 + * {@link memoryFallbackMultipartDone})再把错误往上抛:「已有结论」这件事在 SW + * 存活期内不丢,剩下的交给调用方决定。现在的两个调用方都是把错误按住、只留一 + * 条日志——还原成功那条照常把消息交付出去,放弃那条照常广播自己的原因(见 + * {@link giveUpMultipartId})。 * * @param {{ id: string, ttlMs?: number }} record * @param {number} total - 要清掉的分片数(冲突时取两边的较大值,别漏删) diff --git a/packages/rei-standard-amsg/sw/test/multipart-failure-visibility.test.mjs b/packages/rei-standard-amsg/sw/test/multipart-failure-visibility.test.mjs index 80d7204..57707a8 100644 --- a/packages/rei-standard-amsg/sw/test/multipart-failure-visibility.test.mjs +++ b/packages/rei-standard-amsg/sw/test/multipart-failure-visibility.test.mjs @@ -736,3 +736,53 @@ test('multipart: 墓碑写失败后,重投的旧分片不能把已交付的消 assert.equal(notifications.length, 1, '通知也只该弹第一次那一条'); assert.deepEqual(expiredEvents(postedMessages), [], '交付过的消息更不能被报成丢了'); }); + +test('multipart: 墓碑写失败也不能把放弃的原因盖成「存储故障」', async () => { + const { sw, notifications, postedMessages, triggerPush } = createSwMock(); + const business = install(sw); + + const original = { messageKind: 'content', messageId: 'msg_mp_conflict_done_failed', message: 'x'.repeat(300) }; + const wide = buildMultipartPayloads(original, { id: 'mp_conflict_done_failed', maxChunkBytes: 120 }); + const narrow = buildMultipartPayloads(original, { id: 'mp_conflict_done_failed', maxChunkBytes: 40 }); + assert.notEqual(wide.length, narrow.length, '同一个 id 要有两份不同的 total 才谈得上冲突'); + + await triggerPush(wide[0]); + assert.equal(postedMessages.length, 0, '第一片正常落库,什么都不该广播'); + + // 第二片自相矛盾,走「放弃这条 id」那条路;同时把收尾的第一步(写 done 墓碑) + // 打断。页面要的是「发送端发了冲突分片」这个诊断,不能因为收尾那笔写顺带挂了 + // 就被改写成笼统的存储故障。 + const conn = fake.lastConnection(QUEUE_DB_NAME); + const restore = breakDoneStoreWrites(conn); + let lines; + try { + ({ lines } = await captureErrors(() => triggerPush(narrow[1]))); + } finally { + restore(); + } + + assert.equal(notifications.length, 0); + assert.equal(business.length, 0); + assert.ok( + lines.some((line) => line.includes('multipart chunks disagree on total/encoding')), + `冲突本身要留下能归因的日志:${JSON.stringify(lines)}`, + ); + // 这条测试要的就是「收尾那笔写挂掉」,没打断到的话下面全是假阳性。 + assert.ok( + lines.some((line) => line.includes('quota exceeded')), + `没打断到 done 墓碑那笔写,这条测试白测了:${JSON.stringify(lines)}`, + ); + + const expired = expiredEvents(postedMessages); + assert.equal(expired.length, 1, '收尾失败不能让页面收到两条 MULTIPART_EXPIRED'); + assert.equal(expired[0].id, 'mp_conflict_done_failed'); + assert.equal( + expired[0].reason, 'chunk-conflict', + '报的得是这条路本来的原因,不是外层兜底的 storage-failed', + ); + + assert.ok( + lines.some((line) => line.includes('multipart cleanup while giving up (chunk-conflict) failed')), + `收尾挂掉这件事本身不能被吞掉,得留下能归因的日志:${JSON.stringify(lines)}`, + ); +}); From 26ba771b1ab237d054a0e5d47c1df2cb09eafff5 Mon Sep 17 00:00:00 2001 From: Tosd0 <65720409+Tosd0@users.noreply.github.com> Date: Wed, 23 Sep 2026 11:11:04 +0800 Subject: [PATCH 4/4] =?UTF-8?q?refactor(blob-store):=20=E4=BB=A4=E7=89=8C?= =?UTF-8?q?=20id=20=E7=9A=84=E5=AD=97=E7=AC=A6=E9=9B=86=E5=88=A4=E5=AE=9A?= =?UTF-8?q?=E6=94=B6=E5=88=B0=20token.js=20=E4=B8=80=E5=A4=84?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit gc / content-scan / store 此前各写一份 [A-Za-z0-9_] 的正则,注释互相提醒要和 extractRefs 保持一致;GC 能不能安全回收就取决于这几处完全一致。现在字符集只在 token.js 定义一次,四处判定都走它。行为不变,公共 API 没动。 Claude-Session: https://claude.ai/code/session_017fqgW4ShQ4yjJJwqAiJrcT --- .../blob-store-id-charset-single-source.md | 7 ++++ .../src/content-scan.js | 7 ++-- packages/rei-standard-blob-store/src/gc.js | 7 ++-- packages/rei-standard-blob-store/src/store.js | 13 +++---- packages/rei-standard-blob-store/src/token.js | 19 ++++++++++- .../rei-standard-blob-store/test/gc.test.mjs | 34 +++++++++++++++++++ 6 files changed, 68 insertions(+), 19 deletions(-) create mode 100644 .changeset/blob-store-id-charset-single-source.md diff --git a/.changeset/blob-store-id-charset-single-source.md b/.changeset/blob-store-id-charset-single-source.md new file mode 100644 index 0000000..cbf0b7e --- /dev/null +++ b/.changeset/blob-store-id-charset-single-source.md @@ -0,0 +1,7 @@ +--- +"@rei-standard/blob-store": patch +--- + +令牌 id 的字符集判定收到一处 + +`gc` / `content-scan` / `store` 三处此前各写了一份 `[A-Za-z0-9_]` 的判定,注释互相提醒「要和 `extractRefs` 保持一致」。GC 判断一个 Blob 能不能回收,靠的就是这几处规则完全一致。现在字符集只在 `token.js` 定义一次,四处判定都用它。行为不变,公共 API 没有变化。 diff --git a/packages/rei-standard-blob-store/src/content-scan.js b/packages/rei-standard-blob-store/src/content-scan.js index 03a0c80..97fd674 100644 --- a/packages/rei-standard-blob-store/src/content-scan.js +++ b/packages/rei-standard-blob-store/src/content-scan.js @@ -9,10 +9,7 @@ // 绝不把「读不到」说成「没有重复」——宿主拿着一份假的「没有重复」会以为清干净了。 // 单条 blob 读失败或算不出哈希只跳过这一条(计入 skipped),整轮照常出结果。 -import { parseIdTimestamp } from './token.js'; - -// 与 extractRefs 的 id 边界字符集保持一致(见 token.js;gc.js / store.js 的同名常量同源) -const ID_CHARSET = /^[A-Za-z0-9_]+$/; +import { isIdCharset, parseIdTimestamp } from './token.js'; /** * @typedef {Object} DuplicateGroup @@ -105,7 +102,7 @@ export async function runContentScan({ adapter, prefix }, opts) { try { // 字符集外的 id 整条跳过:extractRefs 按该字符集划边界,这类 id 在引用面上提不全, // 宿主没法把指向它的引用可靠地改写成 canonical,合并进去就是破图(与 gc.js 安全阀 5 同源) - if (!ID_CHARSET.test(id)) { skipped++; continue; } + if (!isIdCharset(id)) { skipped++; continue; } const blob = await adapter.get(id); // keys() 之后、读到之前被删掉了,或者适配器吐了个不是 Blob 的东西 if (!blob || typeof blob.arrayBuffer !== 'function') { skipped++; continue; } diff --git a/packages/rei-standard-blob-store/src/gc.js b/packages/rei-standard-blob-store/src/gc.js index 1517f9f..f2095ad 100644 --- a/packages/rei-standard-blob-store/src/gc.js +++ b/packages/rei-standard-blob-store/src/gc.js @@ -18,13 +18,10 @@ //(mark 只认本 store 前缀、sweep 扫整张表,多前缀共表会互删活图);GC 一轮进行中 // 引用不得在面间搬家(mark 不是一致性快照,瞬间从所有面消失就会被误判孤儿)。 -import { extractRefs, parseIdTimestamp } from './token.js'; +import { extractRefs, isIdCharset, parseIdTimestamp } from './token.js'; const DEFAULT_MIN_AGE_MS = 72 * 3600 * 1000; -// 与 extractRefs 的 id 边界字符集保持一致(见 token.js) -const ID_CHARSET = /^[A-Za-z0-9_]+$/; - /** * @typedef {Object} GcOptions * @property {(Iterable | AsyncIterable) & object} refSources 全部可能含令牌的持久化面;吐字符串。`& object` 把裸字符串挡在类型层——string 本身满足 Iterable,而那恰是会「逐字符迭代、什么都标记不到」的最危险误用 @@ -98,7 +95,7 @@ export async function runGc({ adapter, prefix }, opts) { // 串行删除是刻意的——GC 是后台活儿,并行只会压满 IDB。 for (const id of ids) { if (used.has(id)) { kept++; continue; } - if (!ID_CHARSET.test(id)) { kept++; continue; } // 安全阀 5:mark 不可能命中的 id,无引用不构成证据 + if (!isIdCharset(id)) { kept++; continue; } // 安全阀 5:mark 不可能命中的 id,无引用不构成证据 const ts = parseIdTimestamp(id, now); if (ts !== null && now - ts < minAgeMs) { kept++; continue; } // 安全阀 6:与某个在用 id 互为前缀 = 令牌边界出过事(复合键拼接 / 分块切开)的痕迹,不删。 diff --git a/packages/rei-standard-blob-store/src/store.js b/packages/rei-standard-blob-store/src/store.js index 90391d7..a83008d 100644 --- a/packages/rei-standard-blob-store/src/store.js +++ b/packages/rei-standard-blob-store/src/store.js @@ -2,16 +2,11 @@ // 本模块完全不碰 IndexedDB。错误哲学:读失败 null、删失败吞、put 失败上抛 //(调用方必须知道图没存进去)、迁移失败回退原串。 -import { DEFAULT_PREFIX, genId } from './token.js'; +import { DEFAULT_PREFIX, genId, isIdCharset } from './token.js'; import { dataUrlToBlob, blobToDataUrl } from './dataurl.js'; import { runGc } from './gc.js'; import { runContentScan } from './content-scan.js'; -// 与 extractRefs 的 id 边界字符集保持一致(见 token.js;gc.js 的 ID_CHARSET 同源)。 -// restore 按它拒收字符集外的 id:这类 id 一旦写入,extractRefs 在引用面上提不全它、 -// GC 只能靠安全阀永久豁免,等于制造永不可回收的存量。 -const ID_CHARSET = /^[A-Za-z0-9_]+$/; - /** * @typedef {Object} StorageAdapter * @property {(id: string) => Promise} get @@ -64,14 +59,16 @@ export function createBlobStore(options) { * @param {Blob} blob * @returns {Promise} * @throws {TypeError} token 不是本 store 的令牌、id 为空或含字符集外字符——这是编程/ - * 数据错误,吵着抛(字符集外的 id 会成为 GC 永不可回收的存量,见 ID_CHARSET 注释); + * 数据错误,吵着抛(字符集外的 id 会成为 GC 永不可回收的存量,见 token.js 的 isIdCharset); * blob 的鸭子判定与 put 相同,拒收非 Blob。 */ async restore(token, blob) { if (!isRef(token)) { throw new TypeError(`restore: 需要本 store 前缀的令牌(形如 ${prefix})`); } - if (!ID_CHARSET.test(idOf(token))) { + // 字符集外的 id 一旦写进去,extractRefs 在引用面上提不全它、GC 只能靠安全阀永久豁免, + // 等于亲手制造永不可回收的存量,所以在入口拦掉。 + if (!isIdCharset(idOf(token))) { throw new TypeError('restore: 令牌 id 须非空且完整落在 [A-Za-z0-9_] 内——字符集外的 id 引用提取不全、GC 永不能回收'); } if (!blob || typeof blob.arrayBuffer !== 'function' || typeof blob.slice !== 'function') { diff --git a/packages/rei-standard-blob-store/src/token.js b/packages/rei-standard-blob-store/src/token.js index cd8af8f..ff012c7 100644 --- a/packages/rei-standard-blob-store/src/token.js +++ b/packages/rei-standard-blob-store/src/token.js @@ -4,6 +4,11 @@ export const DEFAULT_PREFIX = 'blobref:'; +// id 的合法字符集,全包只此一份。extractRefs 逐字符扫 ID_CHAR 来划令牌在文本里的边界, +// 要整串判定的地方(gc / content-scan / store)走下面的 isIdCharset——两种用法同源。 +const ID_CHAR = /[A-Za-z0-9_]/; +const ID_CHARSET = new RegExp(`^${ID_CHAR.source}+$`); + let seq = 0; /** @@ -32,6 +37,18 @@ export function parseIdTimestamp(id, now = Date.now()) { return ts; } +/** + * 整个 id 是否非空且完整落在令牌字符集内。 + * 字符集外的 id(比如存量数据直接拿带 `-` 的 UUID 当 id)在引用面上提不全—— + * extractRefs 扫到越界字符就收尾,提出来的只是半截。所以 GC 只能豁免这类 id、 + * content-scan 只能跳过、restore 干脆拒收,判定都落在这个函数上。 + * @param {string} id + * @returns {boolean} + */ +export function isIdCharset(id) { + return ID_CHARSET.test(id); +} + /** * 从任意字符串提取全部令牌。prefix 之后取最长的 [A-Za-z0-9_] 段作为 id, * 所以 JSON 串里内嵌的令牌(后随引号)也能正确截断。 @@ -50,7 +67,7 @@ export function extractRefs(str, prefix = DEFAULT_PREFIX) { while ((i = str.indexOf(prefix, i)) !== -1) { let j = i + prefix.length; if (j < runEnd) j = runEnd; // i 落在上一条已扫 run 内:[j, runEnd) 都是词字符,直接续用终点 - while (j < str.length && /[A-Za-z0-9_]/.test(str[j])) j++; + while (j < str.length && ID_CHAR.test(str[j])) j++; runEnd = j; if (j > i + prefix.length) refs.push(str.slice(i, j)); // 只跳过 prefix 本身、不跳过整个 id 段:'blobref:blobref:b_x' 里第二个令牌 diff --git a/packages/rei-standard-blob-store/test/gc.test.mjs b/packages/rei-standard-blob-store/test/gc.test.mjs index f02d4d8..056317d 100644 --- a/packages/rei-standard-blob-store/test/gc.test.mjs +++ b/packages/rei-standard-blob-store/test/gc.test.mjs @@ -1,6 +1,7 @@ import test from 'node:test'; import assert from 'node:assert/strict'; import { createBlobStore } from '../src/store.js'; +import { extractRefs } from '../src/token.js'; import { memoryAdapter } from './helpers.mjs'; const blobOf = (s) => new Blob([s], { type: 'text/plain' }); @@ -275,3 +276,36 @@ test('新鲜豁免先于边界歧义豁免:又新鲜又互为前缀的 id 记 }); assert.deepEqual(result, { deleted: 0, kept: 1, keptBoundary: 0, aborted: false }); }); + +test('id 字符集的判定在 extractRefs / restore / gc 三处同源:能被完整提取的 id 才准写入、才会被当孤儿删,提不全的一律拒收并豁免', async () => { + // 三处判定共用 token.js 的一份字符集。任何一处单独放宽或收紧,下面的比对就会不一致: + // 提不全却准写入 = 制造永不可回收的存量;提得全却被豁免 = 真孤儿永远删不掉。 + const ids = [ + 'b_abc_0_deadbe', // SDK 生成的格式 + 'A9_z', // 字符集内的其他形状 + '550e8400-e29b-41d4-a716-446655440000', // UUID,`-` 越界 + 'thumb.png', // `.` 越界 + 'a b', // 空格越界 + ]; + for (const id of ids) { + const token = 'blobref:' + id; + // 基准:把令牌放进一段文本,extractRefs 能不能把整个 id 提回来 + const boundaryOk = extractRefs(`{"pic":"${token}"}`)[0] === token; + + const adapter = memoryAdapter(); + const store = createBlobStore({ adapter }); + let restoreOk = true; + try { + await store.restore(token, blobOf('x')); + } catch (err) { + assert.ok(err instanceof TypeError, `restore(${id}) 应该只因字符集抛 TypeError`); + restoreOk = false; + } + assert.equal(restoreOk, boundaryOk, `restore 与 extractRefs 对 ${id} 的判定不一致`); + + // 无人引用的老 id:字符集内的该删,字符集外的该豁免 + adapter.map.set(id, blobOf('x')); + const { deleted } = await store.gc({ refSources: ['{}'], minAgeMs: 0 }); + assert.equal(deleted === 1, boundaryOk, `gc 与 extractRefs 对 ${id} 的判定不一致`); + } +});