Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
9 changes: 9 additions & 0 deletions .changeset/amsg-instant-chunk-probe-auth.md
Original file line number Diff line number Diff line change
@@ -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 字节,没有每请求重编。
22 changes: 22 additions & 0 deletions .changeset/amsg-server-inflight-and-resume.md
Original file line number Diff line number Diff line change
@@ -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 的明文上限、被推送服务拒收。现在按真实取值里最长的那个量,报出来的最大值对所有类型都成立。
9 changes: 9 additions & 0 deletions .changeset/amsg-sw-multipart-give-up-reason.md
Original file line number Diff line number Diff line change
@@ -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`)丢了——宿主拿这个原因做诊断的话,压力下会把发送端发冲突分片这类问题统一误读成存储故障。

现在这四条路径的收尾失败按住不外抛(结论此前已经记进内存兜底表,后续分片和推送服务的重投照样进不来),仍然用本来那条原因广播一次事件,收尾失败另留一条日志。页面收到的事件次数不变。
7 changes: 7 additions & 0 deletions .changeset/blob-store-id-charset-single-source.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,7 @@
---
"@rei-standard/blob-store": patch
---

令牌 id 的字符集判定收到一处

`gc` / `content-scan` / `store` 三处此前各写了一份 `[A-Za-z0-9_]` 的判定,注释互相提醒「要和 `extractRefs` 保持一致」。GC 判断一个 Blob 能不能回收,靠的就是这几处规则完全一致。现在字符集只在 `token.js` 定义一次,四处判定都用它。行为不变,公共 API 没有变化。
64 changes: 30 additions & 34 deletions packages/rei-standard-amsg/instant/src/index.js
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand All @@ -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) {
Expand Down
18 changes: 13 additions & 5 deletions packages/rei-standard-amsg/instant/src/multipart.js
Original file line number Diff line number Diff line change
Expand Up @@ -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 的上限校验 ──────────────────────────────────────────
//
Expand All @@ -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)后的字符数。
*
Expand All @@ -54,17 +62,17 @@ 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
*/
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 位)。
Expand Down
80 changes: 80 additions & 0 deletions packages/rei-standard-amsg/instant/test/handler.test.mjs
Original file line number Diff line number Diff line change
@@ -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,
Expand Down Expand Up @@ -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 ──────────────────────────────
Expand Down Expand Up @@ -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) };
Expand Down
17 changes: 16 additions & 1 deletion packages/rei-standard-amsg/server/src/server/adapters/d1.js
Original file line number Diff line number Diff line change
Expand Up @@ -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 = ?'];
Expand All @@ -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;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -189,8 +189,10 @@
* 它把「为什么失败」透给已失败的行;不实现时退回 getTaskStatus(409 里就
* 没有 lastError)。
* @property {(params: InsertTaskParams, supersedesUuid: string) => Promise<TaskRow & { superseded: boolean }>} [createTaskSuperseding]
* (可选)建新任务的同一事务里取消旧的那条(POST /schedule-message 的
* supersedesUuid)。不实现时 handler 退回「先删再建」两步(失去原子性)。
* (可选)建新任务的同时取消旧的那条,两件事一起成败(POST /schedule-message
* 的 supersedesUuid)。内置的 D1 / pg / neon 都实现了。不实现时 handler 退回
* 两步:先建新、后删旧——删旧失败最坏是两条都留着,客户端重试一次就能收拾,
* 而反过来(先删后建)建新失败就把旧任务白删了,找不回来。
* @property {(userId: string, rows: Array<Object>) => Promise<number>} [appendOutboxMessages]
* (可选;单用户/D1)push 发送前把整批落进 message_outbox(密文 payload),
* (user_id, message_id) 冲突时更新未 ack 的行、不动已 ack 的。
Expand Down
Loading
Loading