Skip to content
Open
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
8 changes: 8 additions & 0 deletions .dev.vars.example
Original file line number Diff line number Diff line change
Expand Up @@ -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
8 changes: 8 additions & 0 deletions .env.example
Original file line number Diff line number Diff line change
Expand Up @@ -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
25 changes: 25 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -77,6 +77,7 @@ OpenList-Worker 是官方 [OpenListTeam/OpenList](https://github.com/OpenListTea
- **离线任务**:后台任务队列,支持批量操作与异步处理。
- **外部接口**:将聚合存储以 WebDAV 或 S3 兼容协议对外暴露,便于挂载到第三方工具。
- **MCP 服务**:提供 Model Context Protocol 端点,可被 AI 助手等客户端集成调用。
- **并发去重(singleflight)**:同一时刻相同的目录列表 / 元信息请求只向上游网盘发起一次,其余请求复用同一结果;一个部署存在多个 isolate,因此除进程内合并外,还会用数据库表 `x_singleflight` 做跨实例协调(详见[并发去重](#并发去重-singleflight))。

### 权限管理

Expand Down Expand Up @@ -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)
Expand Down
6 changes: 5 additions & 1 deletion package.json
Original file line number Diff line number Diff line change
Expand Up @@ -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."
}
}
},
Expand Down Expand Up @@ -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",
Expand Down
16 changes: 13 additions & 3 deletions src/backend/durable-objects/OpenListDB.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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)
}
}
Expand Down
14 changes: 11 additions & 3 deletions src/backend/internal/model/store/driver/d1.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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 绑定接口形态。
Expand Down Expand Up @@ -46,8 +50,12 @@ const d1Inited = new WeakMap<object, boolean>()

async function ensureSchema(db: any, env?: any): Promise<void> {
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)
Expand Down
10 changes: 7 additions & 3 deletions src/backend/internal/model/store/driver/mysql.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down Expand Up @@ -52,8 +52,12 @@ let _schemaInitedPrefix: string | null = null
async function ensureSchema(pool: any, env?: any): Promise<void> {
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
Expand Down
78 changes: 78 additions & 0 deletions src/backend/internal/model/store/schema.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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
}
48 changes: 48 additions & 0 deletions src/backend/internal/op/storage.ts
Original file line number Diff line number Diff line change
@@ -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"
Expand Down Expand Up @@ -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<ReturnType<typeof resolvePath>>,
virtualPath: string,
requestContext?: StorageRequestContext,
): Promise<{ content: FileItem[]; provider: string; storage?: any }> {
let items: FileItem[] = []
let driverName = "Virtual"

Expand Down Expand Up @@ -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<ReturnType<typeof resolvePath>>,
virtualPath: string,
requestContext?: StorageRequestContext,
): Promise<{ item: FileItem; provider: string; rawUrl: string }> {
if (resolved.isVirtual) {
const name = resolved.cleanPath.split("/").filter(Boolean).pop() || "root"
return {
Expand Down
Loading
Loading