From 168cf7ca3eab5416eac1dd81df8d98c7958b7f2f Mon Sep 17 00:00:00 2001 From: Wudarensheng Date: Wed, 23 Sep 2026 18:30:24 +0800 Subject: [PATCH] =?UTF-8?q?feat(singleflight):=20=E5=B9=B6=E5=8F=91?= =?UTF-8?q?=E7=9B=B8=E5=90=8C=E8=B0=83=E7=94=A8=E5=8E=BB=E9=87=8D=EF=BC=88?= =?UTF-8?q?=E8=BF=9B=E7=A8=8B=E5=86=85=E5=90=88=E5=B9=B6=20+=20=E8=B7=A8?= =?UTF-8?q?=E5=AE=9E=E4=BE=8B=E8=A1=A8=E5=8D=8F=E8=B0=83=EF=BC=89?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 一个部署会同时存在多个 isolate,进程内 Map 无法跨实例合并重复调用,导致 同一目录列表 / 元信息的并发请求各自打一次上游网盘 API —— 这正是限流、封号、 CPU 超时的主要来源。新增 singleflight 做两级去重: - L1 进程内:Map 合并,所有模式生效,零成本 - L2 跨实例:数据库表 x_singleflight 抢锁 + 结果共享(仅 db 模式) · 抢锁靠 INSERT OR IGNORE + 回读 owner,不依赖事务 · 执行者按 lockTtl/3 心跳续租,崩溃后锁到期可被接管 · 等待方轮询读取结果;超时则自行执行,不会把请求挂死 模式由 SINGLEFLIGHT 控制(auto/db/memory/off),默认 auto:有 SQL 库 (d1 / mysql / do)时用表协调,否则退回进程内合并。 关键设计: - 不做结果缓存。只有「到达时观察到执行者处于 running」的调用才复用结果, 新到达者会回收旧记录并重新执行,因此不会出现「删除文件后刷新仍看到被删 文件」这类脏读(handoffMs 只是等待方读取结果的交接窗口,不是缓存)。 - 业务异常原样抛出且不重试;协调层异常(表缺失 / DB 超时 / 结果不可序列化 / 执行者崩溃)一律降级为「各自执行一次」,业务不会失败。 - x_singleflight 属运行时数据,刻意不加入 TABLE_NAMES —— 否则 sqlFormat.save() 的整表 DELETE 会在每次保存配置时清空在途记录。 接入:internal/op/storage.ts 的 listItems / getItem,覆盖 /fs/list、/fs/get、 /fs/dir、搜索、WebDAV、MCP。运行统计见 /api/debug/info 的 singleflight 字段。 配置:SINGLEFLIGHT / _HANDOFF_MS / _LOCK_TTL_MS / _POLL_MS / _WAIT_MS, 已同步到 .env.example、.dev.vars.example、wrangler.jsonc 与 README。 新增单测 17 例(含 mock SQL 驱动),并纳入 test:all 的 test:pkg。 --- .dev.vars.example | 8 + .env.example | 8 + README.md | 25 + package.json | 6 +- src/backend/durable-objects/OpenListDB.ts | 16 +- src/backend/internal/model/store/driver/d1.ts | 14 +- .../internal/model/store/driver/mysql.ts | 10 +- src/backend/internal/model/store/schema.ts | 78 ++ src/backend/internal/op/storage.ts | 48 ++ src/backend/pkg/singleflight.test.ts | 540 ++++++++++++ src/backend/pkg/singleflight.ts | 776 ++++++++++++++++++ src/backend/server/debug.ts | 4 + wrangler.jsonc | 23 + 13 files changed, 1546 insertions(+), 10 deletions(-) create mode 100644 src/backend/pkg/singleflight.test.ts create mode 100644 src/backend/pkg/singleflight.ts diff --git a/.dev.vars.example b/.dev.vars.example index 276f2cbe..82344e67 100644 --- a/.dev.vars.example +++ b/.dev.vars.example @@ -48,3 +48,11 @@ CF_ACCOUNT= CF_KV_UUID= # 具有 KV 读写权限的 API Token CF_API_KEY= + +# singleflight(并发相同调用去重)。一个部署会同时存在多个 isolate,进程内去重 +# 无法跨实例生效,因此存在 SQL 库(d1 / mysql / do)时用表 x_singleflight 做协调。 +# auto(默认)- 有 SQL 库用表,否则退回进程内(memory) +# db - 强制用表协调;协调层出错时自动降级,不影响业务 +# memory - 仅进程内合并(等价 Go 的进程内 singleflight) +# off - 关闭去重,每次调用都直接执行 +SINGLEFLIGHT=auto diff --git a/.env.example b/.env.example index 54ab7470..ae026b70 100644 --- a/.env.example +++ b/.env.example @@ -46,3 +46,11 @@ CF_ACCOUNT= CF_KV_UUID= # 具有 KV 读写权限的 API Token CF_API_KEY= + +# singleflight(并发相同调用去重)。一个部署会同时存在多个 isolate,进程内去重 +# 无法跨实例生效,因此存在 SQL 库(d1 / mysql / do)时用表 x_singleflight 做协调。 +# auto(默认)- 有 SQL 库用表,否则退回进程内(memory) +# db - 强制用表协调;协调层出错时自动降级,不影响业务 +# memory - 仅进程内合并(等价 Go 的进程内 singleflight) +# off - 关闭去重,每次调用都直接执行 +SINGLEFLIGHT=auto diff --git a/README.md b/README.md index 83fdca45..723f1122 100644 --- a/README.md +++ b/README.md @@ -77,6 +77,7 @@ OpenList-Worker 是官方 [OpenListTeam/OpenList](https://github.com/OpenListTea - **离线任务**:后台任务队列,支持批量操作与异步处理。 - **外部接口**:将聚合存储以 WebDAV 或 S3 兼容协议对外暴露,便于挂载到第三方工具。 - **MCP 服务**:提供 Model Context Protocol 端点,可被 AI 助手等客户端集成调用。 +- **并发去重(singleflight)**:同一时刻相同的目录列表 / 元信息请求只向上游网盘发起一次,其余请求复用同一结果;一个部署存在多个 isolate,因此除进程内合并外,还会用数据库表 `x_singleflight` 做跨实例协调(详见[并发去重](#并发去重-singleflight))。 ### 权限管理 @@ -284,6 +285,30 @@ CF_API_KEY=your_api_token `des-cbc-hmac` / `3des-cbc-hmac`(后两者仅兼容用途,不安全) - `ADMIN_PASS`:初始管理员密码(可选,设置后跳过安装向导自动初始化 admin) +#### 并发去重(singleflight) + +同一时刻对同一存储、同一路径的重复读请求(目录列表、元信息)只向上游发起一次调用, +其余并发请求等待并复用同一结果,用于规避网盘限流与无谓的 CPU 消耗。 + +| 变量 | 默认值 | 说明 | +| --- | --- | --- | +| `SINGLEFLIGHT` | `auto` | `auto`:有 SQL 库(d1/mysql/do)时用表协调,否则退回进程内;`db`:强制用表;`memory`:仅进程内;`off`:关闭 | +| `SINGLEFLIGHT_HANDOFF_MS` | `1000` | 执行完成后结果在表中保留多久(ms),供**正在等待的**并发实例读取 | +| `SINGLEFLIGHT_LOCK_TTL_MS` | `30000` | 执行者持锁上限(ms),期间按 1/3 周期续租;执行者崩溃后到期可被接管 | +| `SINGLEFLIGHT_POLL_MS` | `50` | 等待方轮询间隔(ms) | +| `SINGLEFLIGHT_WAIT_MS` | `15000` | 等待方最长等待(ms),超时后自行执行,避免被慢请求拖死 | + +行为要点: + +- **两级去重**:进程内合并(所有模式都生效)+ 数据库表协调(`db` 模式,跨 isolate 生效)。 +- **不做结果缓存**:只有「到达时观察到执行者正在运行」的请求才会复用结果;新到达的请求 + 会回收旧记录并重新执行,因此不存在「删除文件后刷新仍看到旧文件」的脏读。 +- **可用性优先**:表不存在、数据库超时、结果无法序列化、执行者崩溃等**协调层**异常一律 + 降级为「各自执行一次」,业务不会失败;而**业务函数自身的异常会原样抛出且不重试**。 +- **表结构**:`x_singleflight` 由 SQL 驱动在初始化时自动创建,属运行时数据,不参与配置读写。 +- 运行统计(`executed` / `coalesced` / `shared` / `fallback`)可在 `/api/debug/info` 的 + `singleflight` 字段查看,用于确认跨实例协调是否真的生效。 + #### 其他配置 详细配置说明请参考 [官方文档](https://doc.oplist.org/guide/configuration) diff --git a/package.json b/package.json index 8662a0cd..8618bc84 100644 --- a/package.json +++ b/package.json @@ -22,6 +22,9 @@ }, "DB_CIPHER": { "description": "At-rest cipher for sensitive fields (drive credentials, 2FA secrets, password hashes): `none` (default, no encryption), `aes-256-gcm` (HKDF-derived AES-256-GCM key — the `enc:v2:` envelope written by existing encrypted deployments, recommended), `aes-256-gcm-pbkdf2` (legacy `enc:v1:` envelope, PBKDF2 per field, slow), `aes-256-cbc-hmac` (`enc:v3:`), `chacha20-poly1305` (`enc:v4:`, RFC 8439, pure-JS implementation), `des-cbc-hmac` / `3des-cbc-hmac` (`enc:v5:`/`enc:v6:`, compatibility only — single DES is 56-bit and 3DES is deprecated, do not use for real data). Ciphertexts carry an `enc:vN:` prefix and are decrypted by prefix on read, so switching back to `none` (or changing algorithms) never makes existing data unreadable — it only changes how newly written data is stored (legacy ciphertext is migrated to plaintext on the next save). Key derivation is cached per isolate, and unchanged fields are not re-encrypted on save. This setting does NOT affect the shared secret: with `none`, a secret is still generated and persisted during setup when JWT_SECRET is not provided, because JWT signing needs it." + }, + "SINGLEFLIGHT": { + "description": "Dedup of concurrent identical calls (singleflight). `auto` (default) uses the `x_singleflight` table when a SQL backend (d1/mysql/do) is available, otherwise falls back to in-process merging. `db` forces table coordination, `memory` keeps it per-isolate, `off` disables dedup." } } }, @@ -55,10 +58,11 @@ "test:server": "tsx --test \"src/backend/server/*.test.ts\"", "test:store": "tsx --test src/backend/internal/model/store/store.test.ts", "test:model": "tsx --test \"src/backend/internal/model/*.test.ts\"", + "test:pkg": "tsx --test \"src/backend/pkg/*.test.ts\"", "test:regress": "tsx scripts/_regress.mjs", "test:deploy": "node scripts/test-deploy.js", "env:check": "tsx scripts/env-check.mjs", - "test:all": "npm run test:189 && npm run test:drivers && npm run test:server && npm run test:store && npm run test:model && npm run test:regress", + "test:all": "npm run test:189 && npm run test:drivers && npm run test:server && npm run test:store && npm run test:model && npm run test:pkg && npm run test:regress", "sls:deploy": "serverless deploy", "sls:package": "serverless package", "sls:info": "serverless info", diff --git a/src/backend/durable-objects/OpenListDB.ts b/src/backend/durable-objects/OpenListDB.ts index 007cd077..eeb0f182 100644 --- a/src/backend/durable-objects/OpenListDB.ts +++ b/src/backend/durable-objects/OpenListDB.ts @@ -13,7 +13,11 @@ * 通过 RPC 调用(stub.kvGet / sqlQuery 等),每个实例用 `idFromName` 定位到 * 固定实例,保证数据持久在同一 DO 实例。 */ -import { D1_SCHEMA, KV_SCHEMA_SQLITE } from "../internal/model/store/schema" +import { + D1_SCHEMA, + KV_SCHEMA_SQLITE, + SINGLEFLIGHT_SCHEMA_SQLITE, +} from "../internal/model/store/schema" export class OpenListDB { private state: any @@ -27,9 +31,15 @@ export class OpenListDB { return (this.state.storage as any).sql } - /** 建表(幂等):KV 表(map/key 格式)+ 列式表(sql 格式)。 */ + /** + * 建表(幂等):KV 表(map/key 格式)+ 列式表(sql 格式)+ singleflight 协调表。 + */ private ensureSchema(): void { - for (const ddl of [...KV_SCHEMA_SQLITE, ...D1_SCHEMA]) { + for (const ddl of [ + ...KV_SCHEMA_SQLITE, + ...D1_SCHEMA, + ...SINGLEFLIGHT_SCHEMA_SQLITE, + ]) { this.sql.exec(ddl) } } diff --git a/src/backend/internal/model/store/driver/d1.ts b/src/backend/internal/model/store/driver/d1.ts index 5fb7070c..1c9a5400 100644 --- a/src/backend/internal/model/store/driver/d1.ts +++ b/src/backend/internal/model/store/driver/d1.ts @@ -6,7 +6,11 @@ * - OPENLIST_DB (别名) */ import type { Driver } from "../types" -import { buildDdl, KV_SCHEMA_SQLITE } from "../schema" +import { + buildDdl, + buildSingleFlightDdl, + KV_SCHEMA_SQLITE, +} from "../schema" /** * 判断对象是否具备 D1 绑定接口形态。 @@ -46,8 +50,12 @@ const d1Inited = new WeakMap() async function ensureSchema(db: any, env?: any): Promise { if (d1Inited.get(db)) return - // KV 表(map/key 格式)+ 列式表(sql 格式)一并创建 - for (const ddl of [...KV_SCHEMA_SQLITE, ...buildDdl("sqlite", env)]) { + // KV 表(map/key 格式)+ 列式表(sql 格式)+ singleflight 协调表一并创建 + for (const ddl of [ + ...KV_SCHEMA_SQLITE, + ...buildDdl("sqlite", env), + ...buildSingleFlightDdl("sqlite"), + ]) { await db.prepare(ddl).run() } d1Inited.set(db, true) diff --git a/src/backend/internal/model/store/driver/mysql.ts b/src/backend/internal/model/store/driver/mysql.ts index 73604457..9df33f94 100644 --- a/src/backend/internal/model/store/driver/mysql.ts +++ b/src/backend/internal/model/store/driver/mysql.ts @@ -8,7 +8,7 @@ * - MYSQL_HOST, MYSQL_PORT, MYSQL_USER, MYSQL_PASS, MYSQL_NAME */ import type { Driver } from "../types" -import { buildDdl, getTablePrefix, KV_SCHEMA_MYSQL } from "../schema" +import { buildDdl, buildSingleFlightDdl, getTablePrefix, KV_SCHEMA_MYSQL } from "../schema" function isNode(): boolean { return typeof process !== "undefined" && process.release?.name === "node" @@ -52,8 +52,12 @@ let _schemaInitedPrefix: string | null = null async function ensureSchema(pool: any, env?: any): Promise { const prefix = getTablePrefix(env) if (_schemaInitedPrefix === prefix) return - // KV 表(map/key 格式)+ 列式表(sql 格式)一并创建 - for (const ddl of [...KV_SCHEMA_MYSQL, ...buildDdl("mysql", env)]) { + // KV 表(map/key 格式)+ 列式表(sql 格式)+ singleflight 协调表一并创建 + for (const ddl of [ + ...KV_SCHEMA_MYSQL, + ...buildDdl("mysql", env), + ...buildSingleFlightDdl("mysql"), + ]) { await pool.query(ddl) } _schemaInitedPrefix = prefix diff --git a/src/backend/internal/model/store/schema.ts b/src/backend/internal/model/store/schema.ts index b5ce1a0d..11a4a76d 100644 --- a/src/backend/internal/model/store/schema.ts +++ b/src/backend/internal/model/store/schema.ts @@ -487,3 +487,81 @@ export const KV_SCHEMA_SQLITE: string[] = [ export const KV_SCHEMA_MYSQL: string[] = [ "CREATE TABLE IF NOT EXISTS `kv` (`key` VARCHAR(512) PRIMARY KEY, `value` LONGTEXT NOT NULL)", ] + +/** + * singleflight 协调表的逻辑名(不含 `x_` 前缀)。 + * + * 与 Go 后端**无对应表**:Go 是单进程模型,`golang.org/x/sync/singleflight` + * 只需进程内 Map;TS 侧一个部署会同时跑多个 isolate/实例,进程内 Map 无法 + * 互相看见,因此需要一张共享表来做跨实例的单飞协调(见 pkg/singleflight.ts)。 + */ +export const SINGLEFLIGHT_TABLE = "singleflight" + +/** singleflight 表的完整 SQL 表名(含 `x_` 前缀)。 */ +export function singleflightTableName(env?: any): string { + return getTablePrefix(env) + SINGLEFLIGHT_TABLE +} + +/** + * singleflight 协调表(SQLite / Cloudflare D1 / Durable Object 方言)。 + * + * 列语义: + * - `key` 去重键(调用方给定,如 `fs.list:1:/movies`) + * - `owner` 执行者实例标识;只有 owner 能写回结果(防止被接管后串写) + * - `state` running | done | error + * - `started_at` 开始时间(epoch ms,诊断用) + * - `expires_at` 过期时间(epoch ms):running 超时表示执行者疑似崩溃, + * 可被其他实例接管;done/error 则是结果交接窗口的截止时间 + * - `result` 成功结果的 JSON 序列化(仅 done 有值) + * - `error` 失败原因文本(仅 error 有值) + * + * 注意:该表存的是**运行时在途记录**,绝不参与 db.ts 的配置往返, + * 因此刻意不放进 TABLE_NAMES —— 否则 sqlFormat.save() 的整表 DELETE + * 会在每次保存配置时把正在执行中的单飞记录清空。 + */ +export const SINGLEFLIGHT_SCHEMA_SQLITE: string[] = [ + `CREATE TABLE IF NOT EXISTS ${quote(singleflightTableName())} (` + + [ + `${quote("key")} TEXT PRIMARY KEY`, + `${quote("owner")} TEXT NOT NULL`, + `${quote("state")} TEXT NOT NULL`, + `${quote("started_at")} INTEGER NOT NULL`, + `${quote("expires_at")} INTEGER NOT NULL`, + `${quote("result")} TEXT`, + `${quote("error")} TEXT`, + ].join(", ") + + `)`, + // 回收过期行(崩溃残留 / 交接窗口结束)按 expires_at 过滤,走索引避免全表扫。 + `CREATE INDEX IF NOT EXISTS ${quote("idx_singleflight_expires")} ` + + `ON ${quote(singleflightTableName())} (${quote("expires_at")})`, +] + +export const SINGLEFLIGHT_SCHEMA_MYSQL: string[] = [ + `CREATE TABLE IF NOT EXISTS ${quote(singleflightTableName())} (` + + [ + `${quote("key")} VARCHAR(512) PRIMARY KEY`, + `${quote("owner")} VARCHAR(128) NOT NULL`, + `${quote("state")} VARCHAR(16) NOT NULL`, + `${quote("started_at")} BIGINT NOT NULL`, + `${quote("expires_at")} BIGINT NOT NULL`, + `${quote("result")} LONGTEXT`, + `${quote("error")} LONGTEXT`, + ].join(", ") + + `)`, + // MySQL 不支持 `CREATE INDEX IF NOT EXISTS`,而索引缺失只影响过期行的回收效率 + // (建表失败会中断整个 init),因此这里刻意不建索引。 +] + +/** + * 生成 singleflight 表的建表语句。 + * + * 与 KV_SCHEMA / D1_SCHEMA 并列,由各 SQL 驱动在 ensureSchema 时一并执行。 + * 非 SQL 驱动(KV / Blob / 内存)没有该表,此时 singleflight 自动退回内存级。 + */ +export function buildSingleFlightDdl( + dialect: "sqlite" | "mysql", +): string[] { + return dialect === "mysql" + ? SINGLEFLIGHT_SCHEMA_MYSQL + : SINGLEFLIGHT_SCHEMA_SQLITE +} diff --git a/src/backend/internal/op/storage.ts b/src/backend/internal/op/storage.ts index d2900a96..20abc7f1 100644 --- a/src/backend/internal/op/storage.ts +++ b/src/backend/internal/op/storage.ts @@ -1,5 +1,6 @@ import { resolvePath, getDb, getSettings, saveDb } from "../model/db" import { encodeDownloadPath } from "../../pkg/path" +import { singleflight } from "../../pkg/singleflight" import { canUseProxyEndpoint, normalizeExtList } from "../driver/proxy" import { FileItem, StorageDriver, calcFileType } from "../driver/base" import { Onedrive } from "../../drivers/onedrive/driver" @@ -1330,11 +1331,45 @@ export async function flushPendingDriverState( await scheduleStoragePersistence(requestContext?.waitUntil, persistence) } +/** + * fs 读操作的 singleflight 去重键。 + * + * 组成:操作名 + 存储 id + 存储修订号 + 虚拟路径。 + * - 带存储 id/修订号:不同存储、以及存储配置被编辑后,不会复用旧结果; + * - 带虚拟路径(调用方已做过 base_path 处理):不会跨用户串味。 + * + * 该键同时用于进程内(L1)与数据库表(L2)两级去重,因此必须稳定可复现: + * 绝不能包含时间戳、随机数、用户对象等每次调用都不同的内容。 + */ +function fsSingleflightKey( + op: string, + resolved: any, + virtualPath: string, +): string { + const storageId = resolved?.storage?.id ?? "virtual" + const revision = resolved?.storage?.modified ?? "" + return `${op}:${storageId}:${revision}:${virtualPath}` +} + export async function listItems( virtualPath: string, requestContext?: StorageRequestContext, ): Promise<{ content: FileItem[]; provider: string; storage?: any }> { const resolved = await resolvePath(virtualPath, requestContext?.env) + // 并发相同的「同一存储 + 同一目录」只向上游发起一次 list。 + // 键必须在 resolvePath 之后才能确定(需要存储 id),因此去重发生在路径解析之后。 + return singleflight( + fsSingleflightKey("fs.list", resolved, virtualPath), + () => listItemsResolved(resolved, virtualPath, requestContext), + { env: requestContext?.env }, + ) +} + +async function listItemsResolved( + resolved: Awaited>, + virtualPath: string, + requestContext?: StorageRequestContext, +): Promise<{ content: FileItem[]; provider: string; storage?: any }> { let items: FileItem[] = [] let driverName = "Virtual" @@ -1469,6 +1504,19 @@ export async function getItem( requestContext?: StorageRequestContext, ): Promise<{ item: FileItem; provider: string; rawUrl: string }> { const resolved = await resolvePath(virtualPath, requestContext?.env) + // 与 listItems 同理:并发相同的「同一存储 + 同一路径」只向上游取一次元信息。 + return singleflight( + fsSingleflightKey("fs.get", resolved, virtualPath), + () => getItemResolved(resolved, virtualPath, requestContext), + { env: requestContext?.env }, + ) +} + +async function getItemResolved( + resolved: Awaited>, + virtualPath: string, + requestContext?: StorageRequestContext, +): Promise<{ item: FileItem; provider: string; rawUrl: string }> { if (resolved.isVirtual) { const name = resolved.cleanPath.split("/").filter(Boolean).pop() || "root" return { diff --git a/src/backend/pkg/singleflight.test.ts b/src/backend/pkg/singleflight.test.ts new file mode 100644 index 00000000..51b8b580 --- /dev/null +++ b/src/backend/pkg/singleflight.test.ts @@ -0,0 +1,540 @@ +import assert from "node:assert/strict" +import { test } from "node:test" +import { + singleflight, + getSingleFlightStats, + __resetSingleFlightForTest, + __singleFlightInsertIgnoreForTest, +} from "./singleflight" +import { + buildSingleFlightDdl, + singleflightTableName, +} from "../internal/model/store/schema" + +/** + * singleflight 单飞去重回归测试。 + * + * 背景:本仓库跑在 Workers 上,一个部署有多个 isolate,进程内 Map 无法跨实例 + * 合并重复调用,因此需要数据库表做二级协调。本文件锁定以下行为: + * + * 1. L1(进程内)必须把并发的同 key 调用合并为一次执行; + * 2. L2(数据库表)必须能让等待方读到执行者的结果,而不是各自再执行一次; + * 3. **绝不脏读**:到达时若该 key 上已有「已结束」的记录,必须重新执行, + * 不能把上一轮的结果当缓存返回(否则会出现「删除后刷新仍看到被删文件」); + * 4. 业务异常原样抛出且不重试;协调层异常则降级为直接执行(可用性优先); + * 5. off / memory / db 三种模式与 SQL 方言分支正确。 + */ + +const delay = (ms: number) => new Promise((r) => setTimeout(r, ms)) + +interface MockRow { + key: string + owner: string + state: string + started_at: number + expires_at: number + result: string | null + error: string | null +} + +/** + * 基于局部 Map 的简易 SQL 驱动。 + * + * 只实现 pkg/singleflight.ts 实际会发出的那几条语句(与 store.test.ts 的 + * mock 驱动思路一致),并保留 rows 便于用例直接构造「别的实例写入的行」。 + */ +function createMockSqlDriver(dialect: "sqlite" | "mysql" = "sqlite") { + const rows = new Map() + let pendingFailure: Error | null = null + + const maybeFail = () => { + if (pendingFailure) { + const err = pendingFailure + pendingFailure = null + throw err + } + } + + const driver = { + name: dialect === "mysql" ? "mysql" : "d1", + + async query(sql: string, params: any[]): Promise { + maybeFail() + const s = sql.trim() + if (/^SELECT `owner` FROM/i.test(s)) { + const r = rows.get(String(params[0])) + return r ? [{ owner: r.owner }] : [] + } + if (/^SELECT `state`/i.test(s)) { + const r = rows.get(String(params[0])) + return r + ? [ + { + state: r.state, + result: r.result, + error: r.error, + expires_at: r.expires_at, + }, + ] + : [] + } + throw new Error(`mock-sql unsupported query: ${sql}`) + }, + + async execute(sql: string, params: any[]): Promise { + maybeFail() + const s = sql.trim() + + // 回收:过期行 或 已结束(state <> running)的行 + if ( + /^DELETE FROM .* WHERE `key` = \? AND \(`expires_at` <= \? OR `state` <> \?\)/i.test( + s, + ) + ) { + const [key, now, running] = params + const r = rows.get(String(key)) + if (r && (r.expires_at <= Number(now) || r.state !== String(running))) { + rows.delete(String(key)) + } + return + } + + // 执行者主动释放 + if (/^DELETE FROM .* WHERE `key` = \? AND `owner` = \?/i.test(s)) { + const [key, owner] = params + const r = rows.get(String(key)) + if (r && r.owner === String(owner)) rows.delete(String(key)) + return + } + + // 抢占:INSERT [OR] IGNORE,主键冲突即忽略(模拟数据库的原子性) + const ins = s.match( + /^INSERT (?:OR )?IGNORE INTO .* \(([^)]+)\) VALUES \(([^)]+)\)/i, + ) + if (ins) { + const cols = ins[1].split(",").map((c) => c.trim().replace(/`/g, "")) + const key = String(params[0]) + if (rows.has(key)) return + const row: any = {} + cols.forEach((c, i) => { + row[c] = params[i] + }) + rows.set(key, row as MockRow) + return + } + + // 写回结果 + if ( + /^UPDATE .* SET `state` = \?, `result` = \?, `error` = \?, `expires_at` = \? WHERE `key` = \? AND `owner` = \?/i.test( + s, + ) + ) { + const [state, result, error, expires, key, owner] = params + const r = rows.get(String(key)) + if (r && r.owner === String(owner)) { + r.state = String(state) + r.result = result ?? null + r.error = error ?? null + r.expires_at = Number(expires) + } + return + } + + // 续租 + if ( + /^UPDATE .* SET `expires_at` = \? WHERE `key` = \? AND `owner` = \? AND `state` = \?/i.test( + s, + ) + ) { + const [expires, key, owner, state] = params + const r = rows.get(String(key)) + if (r && r.owner === String(owner) && r.state === String(state)) { + r.expires_at = Number(expires) + } + return + } + + throw new Error(`mock-sql unsupported execute: ${sql}`) + }, + + isAvailable: async () => true, + init: async () => {}, + get: async () => null, + put: async () => {}, + delete: async () => {}, + list: async () => [], + health: async () => ({ connected: true }), + + /** 测试钩子:直接读写表内容,用于构造「别的实例」的行。 */ + rows, + failNext(err: Error) { + pendingFailure = err + }, + } + + return driver +} + +// ── 建表语句 ──────────────────────────────────────────────────────────────── + +test("schema: singleflight 表名带 x_ 前缀,且 DDL 覆盖两种方言", () => { + assert.equal(singleflightTableName(), "x_singleflight") + + const sqlite = buildSingleFlightDdl("sqlite").join("\n") + assert.match(sqlite, /CREATE TABLE IF NOT EXISTS `x_singleflight`/) + assert.match(sqlite, /`expires_at` INTEGER NOT NULL/) + // 回收过期行需要按 expires_at 过滤,SQLite 侧建索引 + assert.match(sqlite, /CREATE INDEX IF NOT EXISTS `idx_singleflight_expires`/) + + const mysql = buildSingleFlightDdl("mysql").join("\n") + assert.match(mysql, /CREATE TABLE IF NOT EXISTS `x_singleflight`/) + assert.match(mysql, /`key` VARCHAR\(512\) PRIMARY KEY/) + assert.match(mysql, /`expires_at` BIGINT NOT NULL/) + // MySQL 不支持 CREATE INDEX IF NOT EXISTS,因此不建索引 + assert.doesNotMatch(mysql, /CREATE INDEX/) +}) + +test("schema: singleflight 表不得进入配置往返(否则保存配置会清空在途记录)", async () => { + const schema = await import("../internal/model/store/schema") + assert.ok( + !(schema.TABLE_NAMES as readonly string[]).includes("singleflight"), + "singleflight 是运行时数据,不能放进 TABLE_NAMES", + ) + assert.ok( + !(schema.DDL_TABLE_NAMES as readonly string[]).includes("singleflight"), + "建表由驱动单独执行(buildSingleFlightDdl),不混入列式表 DDL", + ) +}) + +test("SQL 方言:mysql 用 INSERT IGNORE,d1/do(SQLite)用 INSERT OR IGNORE", () => { + assert.equal(__singleFlightInsertIgnoreForTest("mysql"), "INSERT IGNORE") + assert.equal(__singleFlightInsertIgnoreForTest("d1"), "INSERT OR IGNORE") + assert.equal(__singleFlightInsertIgnoreForTest("do"), "INSERT OR IGNORE") +}) + +// ── L1:进程内合并 ────────────────────────────────────────────────────────── + +test("memory: 并发相同 key 只执行一次,所有调用方拿到同一结果", async () => { + __resetSingleFlightForTest() + let executed = 0 + + const fn = async () => { + executed++ + await delay(30) + return { value: "shared" } + } + + const results = await Promise.all( + Array.from({ length: 10 }, () => + singleflight("mem-a", fn, { mode: "memory", env: {} }), + ), + ) + + assert.equal(executed, 1, "10 个并发调用只应触发 1 次执行") + for (const r of results) assert.deepEqual(r, { value: "shared" }) + + const stats = getSingleFlightStats() + assert.equal(stats.executed, 1) + assert.equal(stats.coalesced, 9, "其余 9 次应被 L1 合并") +}) + +test("memory: 执行结束后不缓存结果,下一次调用会重新执行", async () => { + __resetSingleFlightForTest() + let executed = 0 + const fn = async () => { + executed++ + return executed + } + + assert.equal(await singleflight("mem-b", fn, { mode: "memory", env: {} }), 1) + assert.equal(await singleflight("mem-b", fn, { mode: "memory", env: {} }), 2) + assert.equal(executed, 2, "singleflight 只合并并发,不做结果缓存") +}) + +test("memory: 失败传播给所有等待方,且不残留导致后续调用被误合并", async () => { + __resetSingleFlightForTest() + let executed = 0 + const failing = async () => { + executed++ + await delay(10) + throw new Error("upstream boom") + } + + const settled = await Promise.allSettled([ + singleflight("mem-c", failing, { mode: "memory", env: {} }), + singleflight("mem-c", failing, { mode: "memory", env: {} }), + ]) + assert.equal(executed, 1) + assert.ok(settled.every((s) => s.status === "rejected")) + assert.match(String((settled[0] as any).reason?.message), /upstream boom/) + + // 在途表必须已清理:下一次调用应重新执行而不是复用已失败的 Promise + assert.equal( + await singleflight("mem-c", async () => "ok", { mode: "memory", env: {} }), + "ok", + ) + assert.equal(getSingleFlightStats().active, 0) +}) + +test("off: 完全关闭去重,每次调用都直接执行", async () => { + __resetSingleFlightForTest() + let executed = 0 + const fn = async () => { + executed++ + await delay(20) + return executed + } + + await Promise.all( + Array.from({ length: 3 }, () => + singleflight("off-a", fn, { mode: "off", env: {} }), + ), + ) + assert.equal(executed, 3) + assert.equal(getSingleFlightStats().mode, "off") +}) + +// ── L2:数据库表协调 ──────────────────────────────────────────────────────── + +test("db: 到达时发现执行者正在运行 → 等待并共享其结果(不重复执行)", async () => { + __resetSingleFlightForTest() + const driver = createMockSqlDriver() + const now = Date.now() + + // 模拟「另一个 isolate 已在执行」:running 行由 other 持有 + driver.rows.set("db-share", { + key: "db-share", + owner: "other-instance", + state: "running", + started_at: now, + expires_at: now + 5000, + result: null, + error: null, + }) + // 30ms 后对方写回结果 + setTimeout(() => { + const row = driver.rows.get("db-share")! + row.state = "done" + row.result = JSON.stringify({ from: "leader" }) + row.expires_at = Date.now() + 1000 + }, 30) + + let executed = 0 + const out = await singleflight( + "db-share", + async () => { + executed++ + return { from: "self" } + }, + { mode: "db", driver, env: {}, pollMs: 5, waitTimeoutMs: 3000 }, + ) + + assert.deepEqual(out, { from: "leader" }) + assert.equal(executed, 0, "等待方不应自己执行") + assert.equal(getSingleFlightStats().shared, 1) +}) + +test("db: 到达时记录已结束 → 必须重新执行,绝不返回旧结果(防脏读)", async () => { + __resetSingleFlightForTest() + const driver = createMockSqlDriver() + const now = Date.now() + + // 上一轮调用留下的交接记录:内容已经过期 + driver.rows.set("db-stale", { + key: "db-stale", + owner: "other-instance", + state: "done", + started_at: now - 100, + expires_at: now + 10_000, + result: JSON.stringify({ value: "stale" }), + error: null, + }) + + let executed = 0 + const out = await singleflight( + "db-stale", + async () => { + executed++ + return { value: "fresh" } + }, + { mode: "db", driver, env: {}, pollMs: 5, waitTimeoutMs: 500 }, + ) + + assert.deepEqual(out, { value: "fresh" }, "新到达者必须拿到新鲜结果") + assert.equal(executed, 1) + assert.equal(getSingleFlightStats().shared, 0, "不应命中交接记录") +}) + +test("db: 抢到锁后写回结果,行状态为 done 且带交接窗口", async () => { + __resetSingleFlightForTest() + const driver = createMockSqlDriver() + + const out = await singleflight("db-lead", async () => ({ n: 7 }), { + mode: "db", + driver, + env: {}, + handoffMs: 1234, + }) + + assert.deepEqual(out, { n: 7 }) + const row = driver.rows.get("db-lead") + assert.ok(row, "执行者应留下交接记录供等待方读取") + assert.equal(row!.state, "done") + assert.equal(row!.result, JSON.stringify({ n: 7 })) + assert.equal(row!.owner.length > 0, true) + assert.ok( + row!.expires_at > Date.now() + 1000 && row!.expires_at <= Date.now() + 1400, + "expires_at 应约为 now + handoffMs", + ) + assert.equal(getSingleFlightStats().executed, 1) +}) + +test("db: handoffMs=0 时执行完立即释放,不留交接记录", async () => { + __resetSingleFlightForTest() + const driver = createMockSqlDriver() + + await singleflight("db-nohand", async () => "v", { + mode: "db", + driver, + env: {}, + handoffMs: 0, + }) + + assert.equal(driver.rows.size, 0, "交接窗口为 0 应立即删除行") +}) + +test("db: 业务异常原样抛出,且不会被降级重试(只执行一次)", async () => { + __resetSingleFlightForTest() + const driver = createMockSqlDriver() + let executed = 0 + + await assert.rejects( + singleflight( + "db-err", + async () => { + executed++ + throw new Error("driver exploded") + }, + { mode: "db", driver, env: {} }, + ), + /driver exploded/, + ) + + assert.equal(executed, 1, "业务异常不得触发降级重试") + assert.equal(getSingleFlightStats().fallback, 0) + // 失败也应写回,让等待方共享失败而不是各自重试上游 + assert.equal(driver.rows.get("db-err")!.state, "error") +}) + +test("db: 协调层异常(DB 不可用)→ 降级为直接执行,业务不受影响", async () => { + __resetSingleFlightForTest() + const driver = createMockSqlDriver() + driver.failNext(new Error("D1_ERROR: no such table: x_singleflight")) + + let executed = 0 + const out = await singleflight( + "db-degrade", + async () => { + executed++ + return "business-ok" + }, + { mode: "db", driver, env: {} }, + ) + + assert.equal(out, "business-ok") + assert.equal(executed, 1) + assert.equal(getSingleFlightStats().fallback, 1) +}) + +test("db: 等待超时后自行执行,不会无限挂起", async () => { + __resetSingleFlightForTest() + const driver = createMockSqlDriver() + const now = Date.now() + + // 执行者一直 running 且迟迟不写回 + driver.rows.set("db-timeout", { + key: "db-timeout", + owner: "other-instance", + state: "running", + started_at: now, + expires_at: now + 60_000, + result: null, + error: null, + }) + + let executed = 0 + const out = await singleflight( + "db-timeout", + async () => { + executed++ + return "self-result" + }, + { mode: "db", driver, env: {}, pollMs: 5, waitTimeoutMs: 40 }, + ) + + assert.equal(out, "self-result") + assert.equal(executed, 1, "等待超时后应自行执行,保证请求不被拖死") +}) + +// ── auto 模式 ─────────────────────────────────────────────────────────────── + +test("auto: 探测不到 SQL 驱动 → memory(不会因为无数据库而报错)", async () => { + __resetSingleFlightForTest() + let executed = 0 + const out = await singleflight( + "auto-mem", + async () => { + executed++ + return "ok" + }, + { mode: "auto", driver: null, env: {} }, + ) + + assert.equal(out, "ok") + assert.equal(executed, 1) + assert.equal(getSingleFlightStats().mode, "memory") +}) + +test("auto: 探测到 SQL 驱动 → db", async () => { + __resetSingleFlightForTest() + const driver = createMockSqlDriver() + const out = await singleflight("auto-db", async () => "ok", { + mode: "auto", + driver, + env: {}, + }) + + assert.equal(out, "ok") + assert.equal(getSingleFlightStats().mode, "db") + assert.equal(driver.rows.get("auto-db")!.state, "done") +}) + +test("环境变量 SINGLEFLIGHT 可强制模式,非法值回退 auto", async () => { + __resetSingleFlightForTest() + let executed = 0 + const fn = async () => { + executed++ + await delay(15) + return executed + } + + // off:3 次调用全部执行 + await Promise.all( + Array.from({ length: 3 }, () => + singleflight("env-off", fn, { env: { SINGLEFLIGHT: "off" } }), + ), + ) + assert.equal(executed, 3) + assert.equal(getSingleFlightStats().mode, "off") + + __resetSingleFlightForTest() + executed = 0 + // 非法值 → auto → 无 SQL 驱动 → memory(仍然合并) + await Promise.all( + Array.from({ length: 3 }, () => + singleflight("env-bad", fn, { env: { SINGLEFLIGHT: "nonsense" } }), + ), + ) + assert.equal(executed, 1) + assert.equal(getSingleFlightStats().mode, "memory") +}) diff --git a/src/backend/pkg/singleflight.ts b/src/backend/pkg/singleflight.ts new file mode 100644 index 00000000..18af772f --- /dev/null +++ b/src/backend/pkg/singleflight.ts @@ -0,0 +1,776 @@ +/** + * singleflight(单飞)—— 把「同一时刻的同一个调用」合并成一次执行。 + * + * ## 为什么 TS 侧需要它(Go 侧不需要) + * + * Go 后端是**单进程**模型,`golang.org/x/sync/singleflight` 用一个进程内 Map + * 就够了:同一时刻只有一个 goroutine 真正执行,其余等待者共享结果。 + * + * 本仓库跑在 Cloudflare Workers / EdgeOne Edge Functions 上,**一个部署会同时 + * 存在多个 isolate / 实例**,且请求会被随机分发到不同 isolate。进程内 Map 只能 + * 合并「落在同一个 isolate 且时间重叠」的调用,跨 isolate 的重复请求依然会各自 + * 打一次上游网盘 API —— 这正是限流、封号、CPU 超时的主要来源。 + * + * ## 两级去重 + * + * L1 进程内(永远开启):`Map`,零成本,同一 isolate 内直接复用 + * 同一个 Promise。任何模式下都生效。 + * L2 跨实例(仅 db 模式):用数据库表 `x_singleflight` 做协调 —— 抢到行的一方 + * 执行,其余实例轮询该行,执行完成后**直接读取结果**, + * 因此上游只被调用一次。 + * + * ## 模式(环境变量 SINGLEFLIGHT) + * + * auto(默认)有 SQL 库(d1 / mysql / do)→ db;否则 → memory + * db 强制走数据库表协调;协调层出错时自动降级为直接执行 + * memory 仅进程内合并(等价于 Go 的进程内 singleflight) + * off 关闭去重,每次调用都直接执行 + * + * ## 可用性优先 + * + * 单飞是**优化**而非**依赖**:任何协调层面的异常(表不存在、DB 超时、结果无法 + * 序列化、执行者崩溃、等待超时)都不会让业务失败,最坏情况退化为「各自执行一次」。 + * 但**业务函数自身的异常会原样抛出**,且不会被误判为协调失败而重试 —— 见 + * `FlightOutcome` 的三态设计。 + */ + +import type { Driver } from "../internal/model/store/types" +import { getStorageBackend } from "../internal/model/store/backend" +import { singleflightTableName } from "../internal/model/store/schema" + +/** 单飞模式。 */ +export type SingleFlightMode = "off" | "memory" | "db" + +/** 模式 + auto(由运行时探测决定实际模式)。 */ +export type SingleFlightModeInput = SingleFlightMode | "auto" + +/** 支持 SQL 的驱动子集(d1 / mysql / do 具备 query + execute)。 */ +type SqlDriverLike = Pick & + Partial> + +/** 单次调用的可选项(多数场景无需传,用默认值 + 环境变量即可)。 */ +export interface SingleFlightOptions { + /** 覆盖模式;默认取环境变量 SINGLEFLIGHT(缺省 auto)。 */ + mode?: SingleFlightModeInput + /** 环境上下文(用于读取环境变量、探测存储驱动)。 */ + env?: any + /** + * 显式指定用于协调的 SQL 驱动。 + * 一般不需要:留空时自动取当前存储后端(见 getStorageBackend)。 + */ + driver?: SqlDriverLike | null + /** + * 结果交接窗口(ms,默认 1000)。 + * + * 执行完成后,结果/错误在表中保留这么久,供**正在等待的**并发实例读取。 + * 它**不是**结果缓存:新到达的调用若看到已结束的记录,会回收该行并重新执行 + * (见 dbFlight 注释),因此不存在「窗口内读到旧数据」的脏读。 + * 设 0 表示执行完立即释放,此时等待方基本拿不到结果、会各自执行一次。 + */ + handoffMs?: number + /** + * 执行者持有锁的最长时间(ms,默认 30000)。 + * 期间执行者会按 1/3 周期续租;若执行者崩溃,锁到期后其他实例可接管。 + */ + lockTtlMs?: number + /** 等待方的轮询间隔(ms,默认 50)。 */ + pollMs?: number + /** + * 等待方最长等待时间(ms,默认 15000)。 + * 超时后不再等待,自行执行(保证请求不会因为别人慢而被拖死)。 + */ + waitTimeoutMs?: number +} + +/** 运行期统计,用于 /debug/info 与排障。 */ +export interface SingleFlightStats { + /** 最近一次生效的模式(off/memory/db)。 */ + mode: SingleFlightMode + /** 真正执行了业务函数的次数。 */ + executed: number + /** 被进程内合并掉的次数(L1 命中)。 */ + coalesced: number + /** 从数据库表读到他人结果的次数(L2 命中)。 */ + shared: number + /** DB 协调异常而降级的次数。 */ + fallback: number + /** 业务函数抛错的次数。 */ + failed: number + /** 当前在途键数量。 */ + active: number +} + +// ── 默认值与环境变量 ──────────────────────────────────────────────────────── + +const DEFAULT_HANDOFF_MS = 1000 +const DEFAULT_LOCK_TTL_MS = 30_000 +const DEFAULT_POLL_MS = 50 +const DEFAULT_WAIT_TIMEOUT_MS = 15_000 + +/** 读取配置:env 绑定优先,其次 process.env(Node / 本地开发)。 */ +function readEnvValue(env: any, name: string): string | undefined { + try { + const fromEnv = env && typeof env === "object" ? env[name] : undefined + if (fromEnv !== undefined && fromEnv !== null && String(fromEnv) !== "") { + return String(fromEnv) + } + } catch { + // env 可能是 Proxy 等异常对象,忽略 + } + try { + const proc = typeof process !== "undefined" ? (process as any) : undefined + const v = proc?.env?.[name] + if (v !== undefined && v !== null && String(v) !== "") return String(v) + } catch { + // 忽略 + } + return undefined +} + +function readEnvInt(env: any, name: string, fallback: number): number { + const raw = readEnvValue(env, name) + if (raw === undefined) return fallback + const n = parseInt(raw, 10) + return Number.isFinite(n) && n >= 0 ? n : fallback +} + +/** 解析 SINGLEFLIGHT 环境变量(非法值按 auto 处理并告警一次)。 */ +let warnedBadMode = false +function readModeInput(env: any): SingleFlightModeInput { + const raw = readEnvValue(env, "SINGLEFLIGHT")?.trim().toLowerCase() + if (!raw) return "auto" + if (raw === "auto") return "auto" + if (raw === "off" || raw === "false" || raw === "0" || raw === "none") { + return "off" + } + if (raw === "memory" || raw === "mem" || raw === "local" || raw === "proc") { + return "memory" + } + if (raw === "db" || raw === "database" || raw === "sql" || raw === "table") { + return "db" + } + if (!warnedBadMode) { + warnedBadMode = true + console.warn( + `[singleflight] Unknown SINGLEFLIGHT="${raw}", falling back to auto ` + + `(expected one of: auto, off, memory, db)`, + ) + } + return "auto" +} + +// ── 实例标识 ──────────────────────────────────────────────────────────────── + +/** + * 本实例(isolate / 进程)的随机标识。 + * + * 用途:只有写入该值的执行者才能回写结果,避免「锁被接管后,原执行者迟到写回」 + * 造成结果错乱。 + */ +const INSTANCE_ID = (() => { + try { + const c: any = typeof crypto !== "undefined" ? crypto : undefined + if (c && typeof c.randomUUID === "function") return c.randomUUID() + } catch { + // 忽略 + } + return `inst-${Math.random().toString(36).slice(2)}${Date.now().toString(36)}` +})() + +// ── 统计 ──────────────────────────────────────────────────────────────────── + +const stats: SingleFlightStats = { + mode: "memory", + executed: 0, + coalesced: 0, + shared: 0, + fallback: 0, + failed: 0, + active: 0, +} + +/** 当前统计快照(含当前在途键数量)。 */ +export function getSingleFlightStats(): SingleFlightStats { + return { ...stats, active: inflight.size } +} + +// ── L1:进程内合并 ────────────────────────────────────────────────────────── + +/** key → 在途 Promise。同一 isolate 内所有模式共用。 */ +const inflight = new Map>() + +/** + * 进程内合并:命中则复用已有 Promise;未命中则执行 factory 并登记。 + * + * 清理用 `then(onOk, onErr)` 而非 `finally()`:后者会产生一个「派生 Promise」, + * 若业务失败而调用方只 await 了原 Promise,派生 Promise 的 rejection 无人处理, + * 在 Node 下会触发 unhandledRejection 告警。 + */ +function coalesce(key: string, factory: () => Promise): Promise { + const existing = inflight.get(key) + if (existing) { + stats.coalesced++ + return existing as Promise + } + + const pending = factory() + inflight.set(key, pending) + const release = () => { + // 仅清理自己登记的那一个,避免误删后来者的在途项 + if (inflight.get(key) === pending) inflight.delete(key) + } + pending.then(release, release) + return pending +} + +// ── L2:数据库表协调 ──────────────────────────────────────────────────────── + +/** + * 协调层的结果三态。 + * + * 关键点:**业务异常与协调异常必须分开**。 + * - 业务异常:原样抛给调用方,绝不重试(否则上游会被打两次); + * - 协调异常:降级为「自己直接执行一次」,保证业务可用。 + * 用三态返回值而非异常标记,可避免给用户错误对象打标记(可能被冻结)的坑。 + */ +type FlightOutcome = + | { kind: "value"; value: T } + /** 执行失败;error 为业务函数抛出的原始异常,或等待方收到的远端失败原因 */ + | { kind: "business-error"; error: unknown } + /** 抢锁 / 轮询 / 写回等协调环节出错 —— 调用方应降级为直接执行 */ + | { kind: "coordination-error"; error: unknown } + +/** 反引号包裹标识符(SQLite 与 MySQL 均支持)。 */ +function q(name: string): string { + return "`" + name + "`" +} + +function supportsSql(driver: any): driver is SqlDriverLike { + return ( + !!driver && + typeof driver.query === "function" && + typeof driver.execute === "function" + ) +} + +/** SQL 方言:mysql 驱动用 `INSERT IGNORE`,其余(d1 / do,均 SQLite)用 `INSERT OR IGNORE`。 */ +function insertIgnoreClause(driver: SqlDriverLike): string { + return String(driver.name || "").toLowerCase() === "mysql" + ? "INSERT IGNORE" + : "INSERT OR IGNORE" +} + +/** 解析用于协调的 SQL 驱动;不可用时返回 null(调用方降级为内存级)。 */ +async function resolveSqlDriver( + env: any, + explicit: SqlDriverLike | null | undefined, +): Promise { + if (explicit !== undefined) return explicit + try { + const { driver } = await getStorageBackend(env) + return supportsSql(driver) ? driver : null + } catch { + // 无可用存储后端(如 serverless 未绑定)→ 内存级 + return null + } +} + +interface ResolvedOptions { + mode: SingleFlightMode + /** 解析出的协调驱动(仅 db 模式有值)。 */ + driver: SqlDriverLike | null + handoffMs: number + lockTtlMs: number + pollMs: number + waitTimeoutMs: number +} + +async function resolveOptions( + options: SingleFlightOptions, +): Promise { + const env = options.env + const input = options.mode ?? readModeInput(env) + + // 只在需要时才探测存储驱动:off / memory 模式完全不需要碰存储层。 + const needsDriver = input === "auto" || input === "db" + const driver = needsDriver + ? await resolveSqlDriver(env, options.driver) + : null + + const mode: SingleFlightMode = + input === "auto" ? (driver ? "db" : "memory") : input + + return { + mode, + driver, + handoffMs: + options.handoffMs ?? + readEnvInt(env, "SINGLEFLIGHT_HANDOFF_MS", DEFAULT_HANDOFF_MS), + lockTtlMs: + options.lockTtlMs ?? + readEnvInt(env, "SINGLEFLIGHT_LOCK_TTL_MS", DEFAULT_LOCK_TTL_MS), + pollMs: + options.pollMs ?? + readEnvInt(env, "SINGLEFLIGHT_POLL_MS", DEFAULT_POLL_MS), + waitTimeoutMs: + options.waitTimeoutMs ?? + readEnvInt(env, "SINGLEFLIGHT_WAIT_MS", DEFAULT_WAIT_TIMEOUT_MS), + } +} + +function sleep(ms: number): Promise { + return new Promise((resolve) => setTimeout(resolve, ms)) +} + +/** + * 抢占 key。 + * + * 三步走,全部依赖数据库自身的原子性(不依赖「读后写」,因此无需事务): + * 1. 回收可回收的行 —— 崩溃残留的过期 running 行,以及**已结束**的 + * done/error 行(它们只是留给等待方读取的交接记录,谁先来谁回收); + * 2. `INSERT OR IGNORE` —— 主键冲突即表示已有他人在执行,我们成为等待方; + * 3. 回读 owner 确认归属 —— 避免「插入被忽略但仍以为自己是执行者」。 + * + * 注意第 1 步的删除条件:**未过期的 running 行绝不能被删**,否则正在执行的 + * 调用会被并发接管,单飞形同虚设。 + */ +async function dbAcquire( + driver: SqlDriverLike, + table: string, + key: string, + owner: string, + opts: ResolvedOptions, + env: any, +): Promise { + const t = q(table) + const now = Date.now() + + await driver.execute!( + `DELETE FROM ${t} WHERE ${q("key")} = ? AND (` + + `${q("expires_at")} <= ? OR ${q("state")} <> ?)`, + [key, now, "running"], + env, + ) + + await driver.execute!( + `${insertIgnoreClause(driver)} INTO ${t} (` + + `${q("key")},${q("owner")},${q("state")},${q("started_at")},${q("expires_at")},` + + `${q("result")},${q("error")}) VALUES (?,?,?,?,?,NULL,NULL)`, + [key, owner, "running", now, now + opts.lockTtlMs], + env, + ) + + const rows = await driver.query!( + `SELECT ${q("owner")} FROM ${t} WHERE ${q("key")} = ?`, + [key], + env, + ) + return rows.length > 0 && String((rows[0] as any).owner) === owner +} + +/** + * 执行者续租:把 expires_at 往后推。 + * + * 为什么需要:目录很大时上游 list 可能超过 lockTtl,若不续租,等待方会误判 + * 「执行者已崩溃」而接管,导致上游被调用两次 —— 恰好是单飞要消除的行为。 + */ +function startHeartbeat( + driver: SqlDriverLike, + table: string, + key: string, + owner: string, + opts: ResolvedOptions, + env: any, +): any { + try { + const interval = Math.max(1000, Math.floor(opts.lockTtlMs / 3)) + const timer = setInterval(() => { + Promise.resolve( + driver.execute!( + `UPDATE ${q(table)} SET ${q("expires_at")} = ? ` + + `WHERE ${q("key")} = ? AND ${q("owner")} = ? AND ${q("state")} = ?`, + [Date.now() + opts.lockTtlMs, key, owner, "running"], + env, + ), + ).catch(() => { + // 续租失败不影响业务:最坏情况是被接管后多执行一次 + }) + }, interval) + // Node 下不要因为心跳定时器而阻止进程退出(Workers 无此方法,可选调用) + ;(timer as any)?.unref?.() + return timer + } catch { + // 运行时不支持定时器(极端环境):放弃续租,退化为「锁到期可被接管」 + return null + } +} + +/** 停止心跳;定时器不可用时为 no-op。 */ +function stopHeartbeat(timer: any): void { + if (timer === null || timer === undefined) return + try { + clearInterval(timer) + } catch { + // 忽略 + } +} + +/** + * 执行者写回结果。 + * + * `WHERE key = ? AND owner = ?` 是必需的:若锁已被接管(原执行者卡顿超过 + * lockTtl 被他人接管),迟到写回必须作废,否则会覆盖接管者的正确结果。 + * + * 写回失败只告警不抛出:协调是尽力而为,业务结果已拿到,不该因为写不进去而失败。 + */ +async function dbPublish( + driver: SqlDriverLike, + table: string, + key: string, + owner: string, + state: "done" | "error", + result: string | null, + error: string | null, + opts: ResolvedOptions, + env: any, +): Promise { + const t = q(table) + try { + if (opts.handoffMs <= 0) { + // 交接窗口为 0:直接释放,让等待方立刻可以重新抢锁 + await driver.execute!( + `DELETE FROM ${t} WHERE ${q("key")} = ? AND ${q("owner")} = ?`, + [key, owner], + env, + ) + return + } + await driver.execute!( + `UPDATE ${t} SET ${q("state")} = ?, ${q("result")} = ?, ${q("error")} = ?, ` + + `${q("expires_at")} = ? WHERE ${q("key")} = ? AND ${q("owner")} = ?`, + [state, result, error, Date.now() + opts.handoffMs, key, owner], + env, + ) + } catch (err) { + console.warn( + `[singleflight] Failed to publish result for "${key}" (ignored):`, + err instanceof Error ? err.message : err, + ) + } +} + +/** 等待方的等待结果。 */ +type WaitOutcome = + | { status: "ok"; value: any } + | { status: "error"; error: string } + /** 行消失 / 执行者疑似崩溃 / 等待超时 —— 交给调用方自行处理 */ + | { status: "give-up" } + +/** 表中某个 key 的当前行(null 表示不存在)。 */ +interface ProbeRow { + state: string + result: any + error: any + expires_at: number +} + +async function dbProbe( + driver: SqlDriverLike, + table: string, + key: string, + env: any, +): Promise { + const rows = await driver.query!( + `SELECT ${q("state")},${q("result")},${q("error")},${q("expires_at")} ` + + `FROM ${q(table)} WHERE ${q("key")} = ?`, + [key], + env, + ) + if (rows.length === 0) return null + const row: any = rows[0] + return { + state: String(row.state ?? ""), + result: row.result, + error: row.error, + expires_at: Number(row.expires_at), + } +} + +/** + * 轮询等待执行者写回结果。 + * + * 仅在「确认执行者处于 running 状态」后调用,因此这里读到 done/error 就是 + * 本次并发去重的目标结果。 + */ +async function dbWait( + driver: SqlDriverLike, + table: string, + key: string, + opts: ResolvedOptions, + env: any, +): Promise { + const deadline = Date.now() + opts.waitTimeoutMs + + for (;;) { + const now = Date.now() + if (now >= deadline) return { status: "give-up" } + await sleep(Math.max(1, Math.min(opts.pollMs, deadline - now))) + + const row = await dbProbe(driver, table, key, env) + // 行已被回收(交接记录被新到达者回收,或执行者主动释放)→ 不再等待 + if (!row) return { status: "give-up" } + + if (row.state === "done") { + if (row.result === null || row.result === undefined) { + return { status: "give-up" } + } + try { + return { status: "ok", value: JSON.parse(String(row.result)) } + } catch { + // 结果不可反序列化(理论上不该发生)→ 自行执行 + return { status: "give-up" } + } + } + if (row.state === "error") { + return { status: "error", error: String(row.error ?? "") } + } + // running 且已过期 → 执行者疑似崩溃 + if (row.expires_at <= Date.now()) return { status: "give-up" } + } +} + +/** 执行者路径:真正调用业务函数,并把结果/异常写回表。 */ +async function dbLead( + driver: SqlDriverLike, + table: string, + key: string, + owner: string, + fn: () => Promise, + opts: ResolvedOptions, + env: any, +): Promise> { + const heartbeat = startHeartbeat(driver, table, key, owner, opts, env) + + let value: T + try { + value = await fn() + } catch (err) { + stopHeartbeat(heartbeat) + await dbPublish( + driver, + table, + key, + owner, + "error", + null, + err instanceof Error ? err.message : String(err), + opts, + env, + ) + return { kind: "business-error", error: err } + } + stopHeartbeat(heartbeat) + + // 结果必须可 JSON 往返才能共享给其他实例;不可序列化时直接释放锁, + // 让等待方自行执行(本次调用方仍然拿到正确的返回值)。 + let payload: string | null = null + try { + payload = JSON.stringify(value === undefined ? null : value) + } catch { + payload = null + } + if (payload === null) { + await dbPublish( + driver, + table, + key, + owner, + "done", + null, + null, + { + ...opts, + handoffMs: 0, + }, + env, + ) + } else { + await dbPublish(driver, table, key, owner, "done", payload, null, opts, env) + } + return { kind: "value", value } +} + +/** + * 数据库表协调的完整流程。 + * + * ## 关键判定:什么情况下才接受「别人的结果」 + * + * 只有在**本调用到达时执行者正处于 running 状态**时,才等待并接受它的结果 —— + * 这代表两者确实并发重叠,是真正的「单飞去重」。 + * + * 反过来,如果到达时该 key 上已经是一行**已结束**的 done/error(上一轮调用留下的 + * 交接记录),本次调用**绝不接受**它,而是回收该行、自己成为新的执行者。 + * + * 这一点是「不引入脏读」的核心:交接记录因此只是给**正在等待的**并发者读取的, + * 不会变成一个有 TTL 的结果缓存。否则会出现「删除文件后 1 秒内刷新列表仍看到 + * 被删文件」这类问题(handoffMs 窗口内命中旧结果)。 + * + * @returns 三态结果,由调用方决定是否降级 + */ +async function dbFlight( + driver: SqlDriverLike, + key: string, + fn: () => Promise, + opts: ResolvedOptions, + env: any, +): Promise> { + const table = singleflightTableName(env) + const owner = INSTANCE_ID + + try { + // 最多三轮:足够覆盖「探测到 running → 等待」「等待落空 → 接管」 + // 「抢占失败(别人刚抢到)→ 再探测」这几种交错。 + for (let attempt = 0; attempt < 3; attempt++) { + const row = await dbProbe(driver, table, key, env) + + if (row && row.state === "running" && row.expires_at > Date.now()) { + // 到达时正在执行 → 并发去重,等待并接受其结果 + const wait = await dbWait(driver, table, key, opts, env) + if (wait.status === "ok") { + stats.shared++ + return { kind: "value", value: wait.value as T } + } + if (wait.status === "error") { + // 共享他人的失败:与 Go singleflight 语义一致,等待方同样收到错误, + // 而不是各自再打一次上游(那正是要消除的行为)。 + stats.shared++ + return { + kind: "business-error", + error: new Error(wait.error || "singleflight: shared call failed"), + } + } + // 执行者崩溃 / 等待超时 → 下一轮尝试接管 + continue + } + + // 无行,或行已结束(done/error 交接记录)→ 回收并抢占,成为新执行者 + if (await dbAcquire(driver, table, key, owner, opts, env)) { + stats.executed++ + const outcome = await dbLead(driver, table, key, owner, fn, opts, env) + if (outcome.kind === "business-error") stats.failed++ + return outcome + } + // 抢占失败:别人刚成为执行者 → 下一轮会探测到 running 并等待 + } + return { kind: "coordination-error", error: null } + } catch (err) { + return { kind: "coordination-error", error: err } + } +} + +// ── 公开 API ──────────────────────────────────────────────────────────────── + +/** + * 以 `key` 为粒度执行 `fn`,同一时刻相同 key 只执行一次。 + * + * @param key 去重键。**必须能唯一标识一次调用的语义**, + * 例如 `fs.list:12:/movies`(存储 id + 物理路径)。 + * 不同 env / 不同存储务必使用不同 key。 + * @param fn 业务函数。注意其返回值需要可 JSON 往返(db 模式下要共享给 + * 其他实例),因此不要返回函数、Symbol、循环引用等。 + * @param options 可选覆盖项,见 SingleFlightOptions。 + * + * @example + * const items = await singleflight(`fs.list:${storage.id}:${path}`, () => + * driver.list(path), + * ) + */ +export async function singleflight( + key: string, + fn: () => Promise, + options: SingleFlightOptions = {}, +): Promise { + const opts = await resolveOptions(options) + stats.mode = opts.mode + + if (opts.mode === "off") { + stats.executed++ + try { + return await fn() + } catch (err) { + stats.failed++ + throw err + } + } + + // L1:进程内合并。db 模式同样先走这一步 —— 同一 isolate 内没必要为同一个 + // key 反复往返数据库。 + return coalesce(key, async () => { + if (opts.mode === "memory") { + stats.executed++ + try { + return await fn() + } catch (err) { + stats.failed++ + throw err + } + } + + const driver = opts.driver + if (!driver) { + // 探测不到 SQL 驱动 → 降级为内存级(L1 仍然生效) + stats.fallback++ + stats.executed++ + try { + return await fn() + } catch (err) { + stats.failed++ + throw err + } + } + + const outcome = await dbFlight(driver, key, fn, opts, options.env) + switch (outcome.kind) { + case "value": + return outcome.value + case "business-error": + throw outcome.error + default: + // 协调失败 → 降级为直接执行,业务照常可用 + stats.fallback++ + stats.executed++ + try { + return await fn() + } catch (err) { + stats.failed++ + throw err + } + } + }) +} + +/** + * 仅测试用:清空进程内在途表与统计,避免用例间互相污染。 + */ +export function __resetSingleFlightForTest(): void { + inflight.clear() + stats.mode = "memory" + stats.executed = 0 + stats.coalesced = 0 + stats.shared = 0 + stats.fallback = 0 + stats.failed = 0 + warnedBadMode = false +} + +/** + * 仅测试用:暴露 SQL 方言分支(mysql 与 sqlite 的 INSERT 语法差异)。 + */ +export function __singleFlightInsertIgnoreForTest(driverName: string): string { + return insertIgnoreClause({ name: driverName }) +} + +/** 仅测试用:暴露实例标识,便于断言 owner 写入。 */ +export function __singleFlightInstanceIdForTest(): string { + return INSTANCE_ID +} diff --git a/src/backend/server/debug.ts b/src/backend/server/debug.ts index ee32a805..80e0f1c2 100644 --- a/src/backend/server/debug.ts +++ b/src/backend/server/debug.ts @@ -1,6 +1,7 @@ import { Hono } from "hono" import { getDb, getStoreStatus } from "../internal/model/db" import { checkAdminAuth } from "../pkg/utils" +import { getSingleFlightStats } from "../pkg/singleflight" export const debugRouter = new Hono() @@ -13,6 +14,9 @@ debugRouter.get("/info", async (c) => { timestamp: new Date().toISOString(), // 后端驱动信息非敏感,未登录也返回,便于确认 D1/KV/MySQL 是否生效 store: await getStoreStatus(c.env), + // singleflight 去重统计:用于确认跨实例协调是否真的生效 + // (coalesced/shared 持续为 0 说明没有并发重复调用,或协调未生效) + singleflight: getSingleFlightStats(), } if (isAdmin) { diff --git a/wrangler.jsonc b/wrangler.jsonc index a7a08f3a..47797d66 100644 --- a/wrangler.jsonc +++ b/wrangler.jsonc @@ -75,6 +75,29 @@ // 自动生成并持久化共享密钥(JWT 签名需要)。 "DB_CIPHER": "none", + // ── singleflight(并发相同调用去重)────────────────────────────────── + // + // 一个部署会同时存在多个 isolate,进程内去重无法跨实例生效,因此当存在 + // SQL 库(d1 / mysql / do)时用表 `x_singleflight` 做跨实例协调。 + // auto(默认)- 有 SQL 库用表,否则退回进程内(memory) + // db - 强制用表协调;协调层出错时自动降级,不影响业务 + // memory - 仅进程内合并(等价 Go 的进程内 singleflight) + // off - 关闭去重,每次调用都直接执行 + "SINGLEFLIGHT": "auto", + + // 执行完成后结果在表中保留多久(ms),供**正在等待的**并发实例读取。 + // 它不是结果缓存:新到达的调用会回收旧记录并重新执行,不会读到过期数据。 + // "SINGLEFLIGHT_HANDOFF_MS": "1000", + // + // 执行者持锁上限(ms);期间按 1/3 周期续租,执行者崩溃后到期可被接管。 + // "SINGLEFLIGHT_LOCK_TTL_MS": "30000", + // + // 等待方轮询间隔(ms) + // "SINGLEFLIGHT_POLL_MS": "50", + // + // 等待方最长等待(ms);超时后自行执行,避免被慢请求拖死 + // "SINGLEFLIGHT_WAIT_MS": "15000", + // ── 安全鉴权(推荐以 Secret 类型配置,见 .dev.vars.example)────────── // // JWT 签名密钥(>=16 字符,兼作可选的字段加密密钥与定时任务鉴权)