From 19e4e4a10a05497fec9154c53b2e098816558979 Mon Sep 17 00:00:00 2001 From: "richard.li" Date: Fri, 22 May 2026 14:36:42 +0800 Subject: [PATCH 1/4] =?UTF-8?q?docs(plan):=20=E5=8A=A0=E5=BD=95=E5=83=8F?= =?UTF-8?q?=E4=B8=8A=E4=BC=A0=E5=8F=AF=E9=9D=A0=E6=80=A7=20P1/P2=20?= =?UTF-8?q?=E4=BF=AE=E5=A4=8D=E5=AE=9E=E6=96=BD=E8=AE=A1=E5=88=92?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 承接 2026-05-22 上传可靠性 review。3 个 Phase: - Phase 1: P1-a Uploading 状态崩溃卡死 + P1-b 上传成功 DB 入库失败孤儿文件 - Phase 2: P2-b pending 目录磁盘水位保护 - Phase 3: P2-a 补传次数耗尽告警 + 人工介入 API 含 4 项实现前调研结论。 --- .../2026-05-22-upload-reliability-p1p2.md | 198 ++++++++++++++++++ 1 file changed, 198 insertions(+) create mode 100644 docs/superpowers/plans/2026-05-22-upload-reliability-p1p2.md diff --git a/docs/superpowers/plans/2026-05-22-upload-reliability-p1p2.md b/docs/superpowers/plans/2026-05-22-upload-reliability-p1p2.md new file mode 100644 index 00000000..1d7f6939 --- /dev/null +++ b/docs/superpowers/plans/2026-05-22-upload-reliability-p1p2.md @@ -0,0 +1,198 @@ +# 录像上传可靠性 P1/P2 修复 实施计划 + +> **For agentic workers:** 用 superpowers:executing-plans 逐 Phase 执行。Step 用 checkbox(`- [ ]`)跟踪。 + +**Goal:** 修复录像上传链路的 4 个 P1/P2 级可靠性缺陷,目标「不丢录像、可观测、可人工介入」。来源:2026-05-22 上传可靠性 code review。P0(`pending_uploads` / `m7s.db` 不持久化)已通过 130 docker-compose 加 bind mount 解决,不在本 plan 范围。 + +## 背景 + +录像上传链路:录制 → 写本地暂存 → record stop 触发上传。两层重试: +1. `pkg/storage/retry.go` `UploadWithRetry` —— 即时重试(默认 4 次尝试,指数退避 5s→40s) +2. `upload_retry.go` `UploadRetryScheduler`(`task.TickTask`,每 5min)—— 定时补传,失败任务持久化在 SQLite `upload_tasks` 表(`UploadTask` 模型),`MaxRetries=10` + +## 待修 4 个问题 + +- **P1-a — Uploading 状态崩溃卡死**:`retryUpload` 先 `MarkUploading`(status→Uploading=1)再上传;若其间进程崩溃,任务永停 Uploading,而 `QueryPendingUploads` 只查 Failed(3)→ 永不再补传。 +- **P1-b — 上传成功但 DB 入库失败 → 孤儿文件**:`recoder.go` `WriteTailDeferred` 返回的闭包写 `record_streams` 失败时只 `Warn` 不返回 error → MinIO 有文件、DB 无索引。 +- **P2-a — 补传次数耗尽彻底放弃**:`retry_count` 达 `MaxRetries(10)` 后 `QueryPendingUploads` 不再匹配,无告警、无人工介入接口。 +- **P2-b — pending 堆积撑爆本地盘**:MinIO 长期故障时 pending 文件持续堆积,无磁盘水位保护,盘满导致新录制写入失败。 + +## 实现前调研结论(已确认) + +1. **AutoMigrate**:`server.go:334` 已 `AutoMigrate(&db.User{}, &PullProxyConfig{}, &PushProxyConfig{}, &StreamAliasDB{}, &AlarmInfo{}, &UploadTask{})` —— `UploadTask`/`AlarmInfo` 已纳管。给 `UploadTask` 加列自动生效;新表(`record_stream_recovery`)需加进该列表。 +2. **UploadConfig 注入**:`server.go:302` `InitUploadManager(storage.UploadConfig{MaxConcurrentUploads:4, ...})` 为**硬编码**,`ServerConfig`(`server.go:60`)无 `Upload` 字段 —— Phase 2 需给 `ServerConfig` 加 `Upload storage.UploadConfig` 字段并改 `InitUploadManager` 读它。 +3. **admin API**:`plugin.go:73` `IRegisterHandler{ RegisterHandler() map[string]http.HandlerFunc }`,plugin 实现即注册 HTTP 路由。mp4 plugin 当前走 grpc-gateway、未实现 `IRegisterHandler` —— Phase 3 的 upload admin API 挂载点见 Task 3.3。 +4. **告警**:`AlarmInfo`(`alarm.go`,表 `alarm_info`)已 AutoMigrate;`pkg/config/types.go` 有 `AlarmStorageException`/`AlarmDiskSpaceFull` 常量。告警最小实现 = 直接 `db.Create(&AlarmInfo{})`;webhook 分散绑定,推送为可选。 + +## 文件结构 + +``` +[改] upload_task.go UploadTask 加 UploadStartedAt;MarkUploading 原子抢占; + 新增 ReclaimStaleUploading / QueryExhaustedUploads / ResetUploadForRetry +[改] upload_retry.go UploadRetryScheduler.Start();Tick 内回收;retryUpload 适配 CAS +[改] recoder.go WriteTailDeferred 闭包返回 error +[改] plugin/mp4/pkg/record.go writeTrailerTask.dbWrite 类型改 func(task.IJob) error;捕获错误补偿 +[改] server.go ServerConfig 加 Upload 字段;InitUploadManager 读它;新表入 AutoMigrate +[改] pkg/storage/upload_manager.go UploadConfig 加水位字段;MoveToPendingDir 水位检查; + GetPendingDirUsage;ErrPendingDirFull +[新] pkg/storage/statfs_unix.go / statfs_windows.go 分平台 GetDiskFreeBytes +[新] upload_alarm.go raiseUploadAlarm 告警 helper(去重 + 写 alarm_info) +[新] record_recovery.go 孤儿 RecordStream 持久化补偿(record_stream_recovery 表 + 重试) +[新] upload_admin.go admin API:列耗尽任务 / 重置重传 +[新] upload_task_test.go / pkg/storage/upload_manager_test.go 单测 +docs/superpowers/plans/2026-05-22-upload-reliability-p1p2.md 本文件 +``` + +--- + +# Phase 1 — P1-a(Uploading 卡死)+ P1-b(孤儿文件) + +数据一致性硬伤,优先。两者独立,合在一个 Phase 验证。 + +## Task 1.1: UploadTask 加 UploadStartedAt 字段 + +- [ ] `upload_task.go` `UploadTask` struct 加 `UploadStartedAt time.Time`(gorm tag 含 `index`,`desc:"本次进入 Uploading 状态的时间"`)。 +- 用途:P1-a 判断 Uploading 任务卡了多久。`UpdatedAt` 不可靠(任何 Update 都刷新)。 +- 风险:Low。`UploadTask` 已在 `server.go:334` AutoMigrate,加 nullable 列对既有行安全。 + +## Task 1.2: MarkUploading 改原子抢占 + +- [ ] `MarkUploading` 改条件更新:`WHERE id=? AND status=?(Failed)`,`Updates` 同时写 `status=Uploading` + `upload_started_at=now`。签名改 `MarkUploading(db, taskID) (claimed bool)`,返回 `RowsAffected>0`。 +- 用途:`Tick` 对每条记录 `go retryUpload`,5min tick 可能与上轮重叠 → 同文件并发上传。条件更新让只有第一个抢到的继续。 +- 依赖:Task 1.1。风险:Medium —— 改签名,需同步 Task 1.4 调用点。 + +## Task 1.3: 新增 ReclaimStaleUploading + +- [ ] `upload_task.go` 新增 `ReclaimStaleUploading(db *gorm.DB, staleThreshold time.Duration) (int64, error)`:`WHERE status=Uploading AND upload_started_at <= now-staleThreshold`,`Updates` status→Failed、`next_retry_at=now`,**不增 retry_count**(崩溃不算失败尝试)。 +- `staleThreshold` 取值须 > 单文件最长上传耗时(`retryUpload` 里 `WithTimeout(30min)`)→ 建议 35min,做成包级常量。 +- 依赖:Task 1.1。风险:Medium —— 阈值过小会误回收正在传的任务;但 Task 1.2 的条件 `MarkUploading` 兜底,且对象存储覆盖写幂等,最坏只是浪费一次带宽。 + +## Task 1.4: UploadRetryScheduler.Start() + Tick 回收 + retryUpload 适配 + +- [ ] `upload_retry.go` 新增 `func (u *UploadRetryScheduler) Start() error`:启动时调一次 `ReclaimStaleUploading`(回收上次崩溃残留),记日志。`TickTask` 允许 override `Start()`。 +- [ ] `Tick` 在 `QueryPendingUploads` 之前也调一次 `ReclaimStaleUploading`(覆盖运行期 goroutine panic/超时卡死)。 +- [ ] `retryUpload` 改:`if !MarkUploading(...) { return }`(没抢到直接退);调整顺序为**先抢 DB 状态再 `AcquireUploadSlot`**,避免占着 slot 才发现没抢到。 +- 依赖:Task 1.2、1.3。风险:Medium。`retryUpload` 当前是裸 goroutine(`go u.retryUpload`),违反项目「不用裸 goroutine」约定 —— 标为已知债,本 Phase 不扩大改动;若 review 要求再改 `u.AddTask` 调度(`TickTask` 可 `AddTask`)。 + +## Task 1.5: WriteTailDeferred 闭包返回 error + +- [ ] `recoder.go` `WriteTailDeferred` 返回类型从 `func(task.IJob)` 改为 `func(task.IJob) error`。两个闭包分支内 `db.Save(&recordStream...)` 失败时收集并返回 error(`record_streams` 是孤儿判定依据,必须计入;`RecordEvent` 的 Save 失败可只 warn)。 +- 依赖:无(与 P1-a 独立)。风险:Medium —— 调用方唯一(`plugin/mp4/pkg/record.go` `writeTailer`,已 grep 确认),`writeTrailerTask.dbWrite` 字段类型同步改。 + +## Task 1.6: writeTrailerTask 捕获 dbWrite error → 孤儿补偿 + +- [ ] `plugin/mp4/pkg/record.go` `writeTrailerTask.dbWrite` 字段类型改 `func(task.IJob) error`。 +- [ ] `Run()` 阶段 3 成功分支 + `runInsertRangeFastPath` 成功分支的 `t.dbWrite(...)` 调用,改为捕获 error: + - **Phase 1 实现**:内联有限重试(如 3 次 × 1s,不阻塞单线程 trailer queue 太久);仍失败 → `t.Error` 日志 + 写一条 `AlarmStorageException` 告警(注明孤儿 objectKey)。 + - 完整的「持久化补偿队列」挪到 Phase 3 Task 3.4。 +- 依赖:Task 1.5、Phase 2 的 `raiseUploadAlarm`(若 Phase 2 先做)或临时内联告警。风险:Medium。 + +## Task 1.7: Phase 1 验收 + +- [ ] 新建 `upload_task_test.go`:`ReclaimStaleUploading`(超时回收 + 不回收新任务 + retry_count 不变)、`MarkUploading` 并发抢占(2 goroutine 只 1 个 claimed)、`WriteTailDeferred` 闭包 DB 失败返回 error。 +- [ ] `go test ./plugin/mp4/pkg ./test -count=1` 无回归;`go build ./plugin/mp4/... ./pkg/storage/...`(全仓 build 受 crypto 预存在错影响,用子路径)。 +- [ ] 集成:`example/record-test` 录流 → 上传中途 `kill -9` → 重启 → 日志见 `reclaimed stale uploading` 且文件最终补传。 + +--- + +# Phase 2 — P2-b(pending 盘满保护) + +独立,可与 Phase 1 并行。「磁盘写满」拖垮新录制,影响面大于「单文件放弃」,优先于 P2-a。 + +## Task 2.1: UploadConfig 加水位配置 + 接入 ServerConfig + +- [ ] `pkg/storage/upload_manager.go` `UploadConfig` 加字段(沿用现有 `desc`/`default` tag 范式):`PendingMaxSizeMB int`(`default:"0"`=不限)、`PendingMaxFiles int`(`default:"0"`)、`PendingDiskMinFreeMB int`(`default:"0"`)。`InitUploadManager` 把值存入包级变量。 +- [ ] `server.go` `ServerConfig` 加 `Upload storage.UploadConfig` 字段;`server.go:302` `InitUploadManager` 改为读 `s.ServerConfig.Upload`(默认全 0,行为与现网一致)。 +- 风险:Low-Medium。默认 0 = 向后兼容。 + +## Task 2.2: pending 目录用量统计 + 分平台磁盘剩余 + +- [ ] `upload_manager.go` 新增 `GetPendingDirUsage() (totalBytes int64, fileCount int, err error)`(`filepath.WalkDir` 累加)。 +- [ ] 新建 `pkg/storage/statfs_unix.go`(`//go:build !windows`)+ `statfs_windows.go`(`//go:build windows`),实现 `GetDiskFreeBytes(path string) (uint64, error)`;Unix 用 `syscall.Statfs`,Windows 返回「不支持」哨兵(降级为只按文件数/总大小限制)。 +- 风险:Medium —— 跨平台,分文件解决。 + +## Task 2.3: MoveToPendingDir 水位检查 + +- [ ] `upload_manager.go` 新增导出 `ErrPendingDirFull`。`MoveToPendingDir` 开头(`pendingDir==""` 判断后)加水位检查:配置了阈值且已超 → 返回 `ErrPendingDirFull`。 +- [ ] `plugin/mp4/pkg/record.go` `recoverToPending`/`recoverFastPathFailure` 识别 `ErrPendingDirFull` → `Error` 日志(明确「本录像无法暂存、可能丢失」)+ 触发 `AlarmDiskSpaceFull` 告警。 +- 风险:Medium —— `MoveToPendingDir` 语义变化(原几乎总成功);调用方已 grep 确认仅 record.go 两处。 + +## Task 2.4: 告警 helper + pending 周期巡检 + +- [ ] 新建 `upload_alarm.go`:`raiseUploadAlarm(db, alarmType int, streamPath, filePath, desc string)` —— `db.Create(&AlarmInfo{...})`;写入前查最近 N 分钟同类型+同文件的未恢复告警,有则跳过(去抖)。 +- [ ] `UploadRetryScheduler.Tick` 增加:调 `GetPendingDirUsage`,超「警戒线」(阈值 80%)→ `raiseUploadAlarm(AlarmDiskSpaceFull, ...)` + `Warn` 日志。复用现有 5min tick,不新增 task。 +- 风险:Low。 + +## Task 2.5: Phase 2 验收 + +- [ ] `pkg/storage/upload_manager_test.go`:`GetPendingDirUsage`(`t.TempDir()` 造文件断言 count/size)、`MoveToPendingDir`+水位(`PendingMaxFiles=2`,第 3 次返回 `ErrPendingDirFull`)、`GetDiskFreeBytes`(Unix >0 / Windows 不 panic)。 +- [ ] 集成:小 tmpfs 当 pending,关 MinIO 让录像堆积 → 达阈值后 `MoveToPendingDir` 返回 `ErrPendingDirFull`、`alarm_info` 出现告警、盘未满、新录制不受影响。 +- [ ] `go test ./pkg/storage -count=1`。 + +--- + +# Phase 3 — P2-a(补传耗尽告警 + 人工介入) + +依赖 Phase 2 的告警设施;并落地 Phase 1 Task 1.6 的完整孤儿补偿。 + +## Task 3.1: 补传耗尽触发告警 + +- [ ] `MarkUploadRetryFailed` 返回 `exhausted bool`(`retryCount+1 >= MaxRetries`);`retryUpload` 据此调 `raiseUploadAlarm(AlarmStorageException, ..., "补传次数耗尽,待人工处理")`。 +- 依赖:Phase 2 `raiseUploadAlarm`。风险:Low。 + +## Task 3.2: QueryExhaustedUploads + +- [ ] `upload_task.go` 新增 `QueryExhaustedUploads(db, limit)`:`WHERE status=Failed AND retry_count >= max_retries`。`QueryPendingUploads` 保持现状(不改 status 枚举,最小改动)。 +- 风险:Low。 + +## Task 3.3: 人工介入 API + +- [ ] `upload_task.go` 新增 `ResetUploadForRetry(db, taskID) error`:`retry_count=0`、`status=Failed`、`next_retry_at=now`。 +- [ ] 新建 `upload_admin.go` 暴露 HTTP 接口:`GET /api/upload/exhausted`(列耗尽)、`POST /api/upload/{id}/retry`(重传)、可选 `POST /api/upload/{id}/discard`(确认放弃 + 删文件)。 +- **挂载点决策(实现时定)**:mp4 plugin 未实现 `IRegisterHandler`。选项 A — 给 mp4 plugin 加 `RegisterHandler()`;选项 B — 走 server 已有 admin API 体系。需先读 server admin API 路由 + JWT 鉴权代码再定。 +- 依赖:Task 3.2。风险:Medium —— 挂载点 + 鉴权需调研。 + +## Task 3.4: 孤儿 RecordStream 持久化补偿 + +- [ ] 新建表 `record_stream_recovery`(字段:`RecordStream` JSON 快照、retry_count、next_retry_at、created_at),加进 `server.go:334` AutoMigrate。 +- [ ] Task 1.6 内联重试仍失败的孤儿 → 序列化 JSON 存该表。 +- [ ] `UploadRetryScheduler.Tick` 增加:扫 `record_stream_recovery`,对每条重试 `db.Save(&RecordStream)`,成功删记录、失败退避。 +- 依赖:Task 1.6、Phase 2。风险:Medium —— 新表入 AutoMigrate;`RecordStream` JSON round-trip 须保留 `ID`(`Save` upsert 不重复行)。 + +## Task 3.5: Phase 3 验收 + +- [ ] 单测:`QueryExhaustedUploads` / `ResetUploadForRetry` / `RecordStream` JSON round-trip。 +- [ ] 集成:任务 `retry_count` 改到上限 → tick → `alarm_info` 出现耗尽告警;`POST /api/upload/{id}/retry` → 任务重新拉起。 + +--- + +## 风险与缓解 + +| 风险 | 等级 | 缓解 | +|---|---|---| +| DB 加列 / 新表 AutoMigrate | Medium | `UploadTask`/`AlarmInfo` 已纳管;新增列 nullable;新表加进 `server.go:334` | +| `retryUpload` 裸 goroutine 违反项目约定 | Medium | Phase 1 标已知债;必要时改 `u.AddTask` 调度 | +| `WriteTailDeferred` 改签名波及 mp4 plugin | Medium | 调用方唯一,改完 `go build ./plugin/mp4/...` 验证 | +| `syscall.Statfs` Windows 不支持 | Medium | 分平台文件 `statfs_unix.go`/`statfs_windows.go` | +| admin API 挂载点 / JWT 鉴权 | Medium | Task 3.3 实现前调研 server admin API | +| `MoveToPendingDir` 语义变化 | Medium | 调用方仅 record.go 两处,失败分支必告警不 panic | +| 全仓 `go build` 受 crypto 预存在错影响 | Low | 用子路径 build(`./plugin/mp4/... ./pkg/storage/...`) | + +## 依赖与顺序 + +``` +Phase 1 (P1-a + P1-b) ── 独立,优先 +Phase 2 (P2-b) ── 独立,可与 Phase 1 并行;产出 raiseUploadAlarm + │ + ▼ +Phase 3 (P2-a) ── 依赖 Phase 2 告警设施 + 承接 Phase 1 Task 1.6 +``` + +实施顺序:Phase 1 → Phase 2 → Phase 3。 + +## Self-Review + +- **Spec 覆盖**:P1-a → Task 1.1-1.4;P1-b → Task 1.5-1.6 + 3.4;P2-b → Task 2.1-2.4;P2-a → Task 3.1-3.3。 +- **向后兼容**:Phase 2 水位配置默认全 0(不启用);`MarkUploading` 改签名仅内部调用方;不改 `UploadStatus` 枚举。 +- **不改动**:`UploadWithRetry` 即时重试逻辑;`pkg/storage/retry.go` 退避策略;trailer 队列单线程模型。 +- **调研已闭环**:4 项实现前调研均已确认(见上文「实现前调研结论」),仅 Task 3.3 admin API 挂载点留待实现时按代码定。 From 7c84b0d425ce7ddaa662ff2afead8e1ca3078003 Mon Sep 17 00:00:00 2001 From: "richard.li" Date: Fri, 22 May 2026 14:55:07 +0800 Subject: [PATCH 2/4] =?UTF-8?q?fix(upload):=20Phase=201=20=E4=BF=AE=20P1-a?= =?UTF-8?q?=20Uploading=20=E5=8D=A1=E6=AD=BB=20+=20P1-b=20=E5=AD=A4?= =?UTF-8?q?=E5=84=BF=E6=96=87=E4=BB=B6?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit P1-a Uploading 状态崩溃卡死: - UploadTask 加 UploadStartedAt 字段 - MarkUploading 改原子抢占(WHERE status=Failed 条件更新), 返回 claimed - 新增 ReclaimStaleUploading: 把卡 Uploading 超 35min 的任务扫回 Failed(不计 retry_count), 进程崩溃残留可重新补传 - UploadRetryScheduler.Tick 每轮先回收; retryUpload 用 CAS 防并发重传 P1-b 上传成功但 DB 入库失败 -> 孤儿文件: - WriteTailDeferred 闭包改返回 error(record_streams 写失败=孤儿) - writeTrailerTask 捕获该 error, reportOrphan 打日志 + 告警入 alarm_info - 新增 RaiseUploadAlarm 告警 helper(带去重) upload_task_test.go 覆盖 ReclaimStaleUploading / MarkUploading 抢占 / 回收后可补传, 3/3 PASS。 --- plugin/mp4/pkg/record.go | 22 ++++++- recoder.go | 14 +++- upload_alarm.go | 40 ++++++++++++ upload_retry.go | 16 ++++- upload_task.go | 62 ++++++++++++------ upload_task_test.go | 134 +++++++++++++++++++++++++++++++++++++++ 6 files changed, 261 insertions(+), 27 deletions(-) create mode 100644 upload_alarm.go create mode 100644 upload_task_test.go diff --git a/plugin/mp4/pkg/record.go b/plugin/mp4/pkg/record.go index 9169e3c8..23b9ccb5 100644 --- a/plugin/mp4/pkg/record.go +++ b/plugin/mp4/pkg/record.go @@ -39,7 +39,8 @@ type writeTrailerTask struct { storageKey string // 存储类型 key(s3/oss/cos/local) db *gorm.DB // 数据库连接(用于保存失败记录) // dbWrite 在文件完整写入后执行数据库更新,为 nil 时跳过(无 DB 或测试模式)。 - dbWrite func(tailJob task.IJob) + // 返回 error 表示文件已持久化成功但 record_streams 入库失败(孤儿文件)。 + dbWrite func(tailJob task.IJob) error } func (task *writeTrailerTask) Start() (err error) { @@ -232,11 +233,24 @@ func (t *writeTrailerTask) Run() (err error) { t.file = nil // 文件已完整持久化,此时才将记录写入数据库(延迟入库,确保 DB 与可播放文件一致)。 if t.dbWrite != nil { - t.dbWrite(&writeTrailerQueueTask) + if dbErr := t.dbWrite(&writeTrailerQueueTask); dbErr != nil { + t.reportOrphan(dbErr) + } } return } +// reportOrphan 处理「文件已持久化成功但 record_streams 入库失败」—— +// 对象存储上有文件、DB 无索引记录,即孤儿文件。Phase 1 仅做可感知: +// 错误日志 + 告警入 alarm_info;自动补偿(重试入库)由 Phase 3 补偿队列负责。 +func (t *writeTrailerTask) reportOrphan(dbErr error) { + t.Error("upload ok but db record save failed — orphan file", + "err", dbErr, "filePath", t.filePath, "streamPath", t.streamPath) + m7s.RaiseUploadAlarm(t.db, config.AlarmStorageException, + "record db save failed", t.streamPath, t.filePath, + "录像已上传成功但 record_streams 入库失败(孤儿文件): "+dbErr.Error()) +} + // shiftSampleOffsets 把所有 track 的 sample 偏移整体加 delta, // 用于 INSERT_RANGE 把 mdat 逻辑后移后校正 moov 内的 chunk offset。 func (t *writeTrailerTask) shiftSampleOffsets(delta int64) { @@ -387,7 +401,9 @@ func (t *writeTrailerTask) runInsertRangeFastPath() (handled bool, err error) { } t.file = nil if t.dbWrite != nil { - t.dbWrite(&writeTrailerQueueTask) + if dbErr := t.dbWrite(&writeTrailerQueueTask); dbErr != nil { + t.reportOrphan(dbErr) + } } t.Info("insert-range fast path done", "filePath", t.filePath, "moovBytes", moovSize, "insertedBytes", insertLen) diff --git a/recoder.go b/recoder.go index 020ff6f7..aaba24d6 100644 --- a/recoder.go +++ b/recoder.go @@ -214,7 +214,9 @@ func (r *DefaultRecorder) WriteTail(end time.Time, tailJob task.IJob) { // 与 WriteTail 不同,它不立即写库,而是将写库操作包装成闭包返回。 // 调用方应在 MP4 文件完整写入(moov 移到头部)后再调用该闭包, // 以保证数据库记录在文件可播放之后才更新 EndTime。 -func (r *DefaultRecorder) WriteTailDeferred(end time.Time) func(tailJob task.IJob) { +// 闭包返回 error:record_streams 写库失败时返回非 nil(此时文件已上传成功但 +// DB 无索引记录 = 孤儿文件),调用方据此补偿;RecordEvent 写库失败仅告警不计入。 +func (r *DefaultRecorder) WriteTailDeferred(end time.Time) func(tailJob task.IJob) error { r.Event.EndTime = end if r.RecordJob.Plugin.DB == nil || r.RecordJob.RecConf.Mode == config.RecordModeTest { return nil @@ -225,7 +227,7 @@ func (r *DefaultRecorder) WriteTailDeferred(end time.Time) func(tailJob task.IJo filePath := r.Event.FilePath if r.RecordJob.Event != nil { eventSnap := r.Event // 值拷贝:捕获正确的 EndTime 和 RecordStream.ID,RecordEvent 指针稳定 - return func(tailJob task.IJob) { + return func(tailJob task.IJob) error { dbCtx, cancel := context.WithTimeout(context.Background(), 10*time.Second) defer cancel() r.Info("db save RecordEvent (deferred) begin", "filePath", filePath) @@ -235,7 +237,9 @@ func (r *DefaultRecorder) WriteTailDeferred(end time.Time) func(tailJob task.IJo r.Info("db save RecordEvent (deferred) ok", "filePath", filePath) } r.Info("db save RecordStream (deferred) begin", "filePath", filePath) + var streamErr error if result := db.WithContext(dbCtx).Save(&eventSnap.RecordStream); result.Error != nil { + streamErr = result.Error r.Warn("db save RecordStream (deferred) failed", "filePath", filePath, "err", result.Error) } else { r.Info("db save RecordStream (deferred) ok", "filePath", filePath) @@ -243,14 +247,17 @@ func (r *DefaultRecorder) WriteTailDeferred(end time.Time) func(tailJob task.IJo if tailJob != nil { tailJob.AddTask(NewEventRecordCheck(streamType, streamPath, db)) } + return streamErr } } streamSnap := r.Event.RecordStream // 值拷贝 - return func(tailJob task.IJob) { + return func(tailJob task.IJob) error { dbCtx, cancel := context.WithTimeout(context.Background(), 10*time.Second) defer cancel() r.Info("db save RecordStream (deferred) begin", "filePath", filePath) + var streamErr error if result := db.WithContext(dbCtx).Save(&streamSnap); result.Error != nil { + streamErr = result.Error r.Warn("db save RecordStream (deferred) failed", "filePath", filePath, "err", result.Error) } else { r.Info("db save RecordStream (deferred) ok", "filePath", filePath) @@ -258,6 +265,7 @@ func (r *DefaultRecorder) WriteTailDeferred(end time.Time) func(tailJob task.IJo if tailJob != nil { tailJob.AddTask(NewEventRecordCheck(streamType, streamPath, db)) } + return streamErr } } diff --git a/upload_alarm.go b/upload_alarm.go new file mode 100644 index 00000000..649d1f4f --- /dev/null +++ b/upload_alarm.go @@ -0,0 +1,40 @@ +package m7s + +import ( + "time" + + "gorm.io/gorm" +) + +// uploadAlarmDedupWindow:同 alarmType + 同 filePath 的告警去重窗口。 +// 补传调度器 / pending 巡检按周期运行,无去重会反复写同一条告警刷屏。 +const uploadAlarmDedupWindow = 10 * time.Minute + +// RaiseUploadAlarm 写一条上传相关告警到 alarm_info 表,带去重: +// 去重窗口内已存在同 alarmType + 同 filePath 的告警则跳过本次写入。 +// best-effort:db 为 nil 或写入失败均静默返回,不阻塞主流程。 +func RaiseUploadAlarm(db *gorm.DB, alarmType int, alarmName, streamPath, filePath, desc string) { + if db == nil { + return + } + var cnt int64 + if err := db.Model(&AlarmInfo{}). + Where("alarm_type = ? AND file_path = ? AND created_at >= ?", + alarmType, filePath, time.Now().Add(-uploadAlarmDedupWindow)). + Count(&cnt).Error; err == nil && cnt > 0 { + return + } + if len(desc) > 500 { + desc = desc[:500] + } + if len(alarmName) > 255 { + alarmName = alarmName[:255] + } + db.Create(&AlarmInfo{ + StreamPath: streamPath, + AlarmName: alarmName, + AlarmDesc: desc, + AlarmType: alarmType, + FilePath: filePath, + }) +} diff --git a/upload_retry.go b/upload_retry.go index 88dcc5ac..c3bc6ca4 100644 --- a/upload_retry.go +++ b/upload_retry.go @@ -31,6 +31,13 @@ func (u *UploadRetryScheduler) Tick(any) { return } + // 回收卡在 Uploading 状态的任务(进程崩溃 / goroutine 超时残留),扫回 Failed + if n, err := ReclaimStaleUploading(u.s.DB, staleUploadingThreshold); err != nil { + u.Error("reclaim stale uploading", "err", err) + } else if n > 0 { + u.Info("reclaimed stale uploading tasks", "count", n) + } + // 查询待重试的任务(每次最多处理 20 个,避免单次过多) tasks, err := QueryPendingUploads(u.s.DB, 20) if err != nil { @@ -59,18 +66,21 @@ func (u *UploadRetryScheduler) retryUpload(ut UploadTask) { return } + // 原子抢占:仅当任务仍为 Failed 时占用;抢不到说明已被其他 goroutine 处理 + if !MarkUploading(u.s.DB, ut.ID) { + return + } + // 获取上传槽位(并发控制) ctx, cancel := context.WithTimeout(context.Background(), 30*time.Minute) defer cancel() if err := storage.AcquireUploadSlot(ctx); err != nil { + // 已抢占 Uploading 但拿不到槽位,留待 ReclaimStaleUploading 回收 u.Warn("acquire upload slot timeout", "id", ut.ID, "err", err) return } defer storage.ReleaseUploadSlot() - // 标记为上传中 - MarkUploading(u.s.DB, ut.ID) - // 解析元数据 var metadata map[string]string if ut.MetadataJSON != "" { diff --git a/upload_task.go b/upload_task.go index 7898b080..ea055c24 100644 --- a/upload_task.go +++ b/upload_task.go @@ -24,21 +24,22 @@ const ( // UploadTask 上传任务持久化模型,用于追踪失败上传和定时补传 type UploadTask struct { - ID uint `gorm:"primarykey"` - CreatedAt time.Time - UpdatedAt time.Time - LocalPath string `gorm:"size:512;index" json:"localPath" desc:"本地文件路径"` - ObjectKey string `gorm:"size:512" json:"objectKey" desc:"远端对象键"` - StorageType string `gorm:"size:20" json:"storageType" desc:"存储类型(s3/oss/cos)"` - Status UploadStatus `gorm:"default:0;index" json:"status" desc:"状态: 0=待传 1=传输中 2=成功 3=失败"` - RetryCount int `gorm:"default:0" json:"retryCount" desc:"已重试次数"` - MaxRetries int `gorm:"default:10" json:"maxRetries" desc:"最大重试次数"` - FileSize int64 `json:"fileSize" desc:"文件大小(字节)"` - ErrorMessage string `gorm:"size:1024" json:"errorMessage" desc:"最近一次错误信息"` - StreamPath string `gorm:"size:255;index" json:"streamPath" desc:"关联流路径"` - DurationMs uint32 `json:"durationMs" desc:"视频时长(毫秒)"` - MetadataJSON string `gorm:"column:metadata;type:text" json:"metadata" desc:"JSON编码的用户元数据"` - NextRetryAt time.Time `gorm:"index" json:"nextRetryAt" desc:"下次重试时间"` + ID uint `gorm:"primarykey"` + CreatedAt time.Time + UpdatedAt time.Time + LocalPath string `gorm:"size:512;index" json:"localPath" desc:"本地文件路径"` + ObjectKey string `gorm:"size:512" json:"objectKey" desc:"远端对象键"` + StorageType string `gorm:"size:20" json:"storageType" desc:"存储类型(s3/oss/cos)"` + Status UploadStatus `gorm:"default:0;index" json:"status" desc:"状态: 0=待传 1=传输中 2=成功 3=失败"` + RetryCount int `gorm:"default:0" json:"retryCount" desc:"已重试次数"` + MaxRetries int `gorm:"default:10" json:"maxRetries" desc:"最大重试次数"` + FileSize int64 `json:"fileSize" desc:"文件大小(字节)"` + ErrorMessage string `gorm:"size:1024" json:"errorMessage" desc:"最近一次错误信息"` + StreamPath string `gorm:"size:255;index" json:"streamPath" desc:"关联流路径"` + DurationMs uint32 `json:"durationMs" desc:"视频时长(毫秒)"` + MetadataJSON string `gorm:"column:metadata;type:text" json:"metadata" desc:"JSON编码的用户元数据"` + NextRetryAt time.Time `gorm:"index" json:"nextRetryAt" desc:"下次重试时间"` + UploadStartedAt time.Time `gorm:"index" json:"uploadStartedAt" desc:"本次进入上传中状态的时间"` } // TableName GORM 表名 @@ -108,9 +109,16 @@ func QueryPendingUploads(db *gorm.DB, limit int) ([]UploadTask, error) { return tasks, err } -// MarkUploading 标记任务为上传中 -func MarkUploading(db *gorm.DB, taskID uint) { - db.Model(&UploadTask{}).Where("id = ?", taskID).Update("status", UploadStatusUploading) +// MarkUploading 原子抢占任务:仅当任务仍为 Failed 时置 Uploading 并记录开始时间。 +// 返回 true 表示本次抢占成功(可继续上传);false 表示已被其他 goroutine 抢占。 +func MarkUploading(db *gorm.DB, taskID uint) (claimed bool) { + res := db.Model(&UploadTask{}). + Where("id = ? AND status = ?", taskID, UploadStatusFailed). + Updates(map[string]any{ + "status": UploadStatusUploading, + "upload_started_at": time.Now(), + }) + return res.RowsAffected > 0 } // MarkUploadSuccess 标记任务上传成功,删除本地文件 @@ -141,6 +149,24 @@ func MarkUploadRetryFailed(db *gorm.DB, taskID uint, retryCount int, err error) }) } +// staleUploadingThreshold:Uploading 状态超过此时长即视为卡死(进程崩溃 / +// goroutine 超时残留)。须大于单文件最长上传耗时(retryUpload 的 30min 超时), +// 留 buffer 取 35min。 +const staleUploadingThreshold = 35 * time.Minute + +// ReclaimStaleUploading 把卡在 Uploading 状态超过 threshold 的任务扫回 Failed, +// 使其重新进入补传循环。崩溃 / 超时不计入 retry_count,不消耗重试配额。 +// 返回回收的任务数。 +func ReclaimStaleUploading(db *gorm.DB, threshold time.Duration) (int64, error) { + res := db.Model(&UploadTask{}). + Where("status = ? AND upload_started_at <= ?", UploadStatusUploading, time.Now().Add(-threshold)). + Updates(map[string]any{ + "status": UploadStatusFailed, + "next_retry_at": time.Now(), + }) + return res.RowsAffected, res.Error +} + // UploadLocalFile 将本地文件上传到云存储(通用方法,用于补传) func UploadLocalFile(ctx context.Context, st storage.Storage, localPath, remotePath string, metadata map[string]string) error { local, err := os.Open(localPath) diff --git a/upload_task_test.go b/upload_task_test.go new file mode 100644 index 00000000..cdb52a8d --- /dev/null +++ b/upload_task_test.go @@ -0,0 +1,134 @@ +package m7s + +import ( + "sync" + "testing" + "time" + + _ "github.com/ncruces/go-sqlite3/embed" + "github.com/ncruces/go-sqlite3/gormlite" + "gorm.io/gorm" +) + +func newUploadTestDB(t *testing.T) *gorm.DB { + t.Helper() + db, err := gorm.Open(gormlite.Open(":memory:"), &gorm.Config{}) + if err != nil { + t.Fatalf("open in-memory sqlite: %v", err) + } + sqlDB, err := db.DB() + if err != nil { + t.Fatalf("db handle: %v", err) + } + sqlDB.SetMaxOpenConns(1) // :memory: 每连接独立库,限单连接以共享同一库 + if err := db.AutoMigrate(&UploadTask{}); err != nil { + t.Fatalf("migrate UploadTask: %v", err) + } + return db +} + +// seedUploading 造一个 Uploading 状态、UploadStartedAt 为 startedAgo 之前的任务。 +func seedUploading(t *testing.T, db *gorm.DB, startedAgo time.Duration, retryCount int) UploadTask { + t.Helper() + ut := UploadTask{ + LocalPath: "/tmp/x.mp4", + ObjectKey: "x.mp4", + Status: UploadStatusUploading, + RetryCount: retryCount, + MaxRetries: 10, + UploadStartedAt: time.Now().Add(-startedAgo), + } + if err := db.Create(&ut).Error; err != nil { + t.Fatalf("seed: %v", err) + } + return ut +} + +// TestReclaimStaleUploading:超时 Uploading 任务被扫回 Failed 且 retry_count 不变; +// 未超时的不动。 +func TestReclaimStaleUploading(t *testing.T) { + db := newUploadTestDB(t) + stale := seedUploading(t, db, time.Hour, 3) // 卡死 1h,应回收 + fresh := seedUploading(t, db, time.Minute, 2) // 刚 1min,不回收 + + n, err := ReclaimStaleUploading(db, 35*time.Minute) + if err != nil { + t.Fatalf("ReclaimStaleUploading: %v", err) + } + if n != 1 { + t.Fatalf("应回收 1 个,实际 %d", n) + } + + var gotStale UploadTask + db.First(&gotStale, stale.ID) + if gotStale.Status != UploadStatusFailed { + t.Errorf("卡死任务应回收为 Failed,实际 status=%d", gotStale.Status) + } + if gotStale.RetryCount != 3 { + t.Errorf("回收不应改 retry_count,期望 3 实际 %d", gotStale.RetryCount) + } + + var gotFresh UploadTask + db.First(&gotFresh, fresh.ID) + if gotFresh.Status != UploadStatusUploading { + t.Errorf("未超时任务不应被回收,实际 status=%d", gotFresh.Status) + } +} + +// TestMarkUploadingClaim:对同一 Failed 任务并发抢占,只有一个成功; +// 已是 Uploading 的任务不可再被抢占。 +func TestMarkUploadingClaim(t *testing.T) { + db := newUploadTestDB(t) + ut := UploadTask{ + LocalPath: "/tmp/y.mp4", ObjectKey: "y.mp4", + Status: UploadStatusFailed, MaxRetries: 10, + } + if err := db.Create(&ut).Error; err != nil { + t.Fatalf("seed: %v", err) + } + + const n = 8 + var wg sync.WaitGroup + results := make([]bool, n) + for i := 0; i < n; i++ { + wg.Add(1) + go func(idx int) { + defer wg.Done() + results[idx] = MarkUploading(db, ut.ID) + }(i) + } + wg.Wait() + + claimed := 0 + for _, ok := range results { + if ok { + claimed++ + } + } + if claimed != 1 { + t.Fatalf("并发抢占应只有 1 个成功,实际 %d", claimed) + } + if MarkUploading(db, ut.ID) { + t.Error("已 Uploading 的任务不应再被抢占") + } +} + +// TestReclaimThenQueryable:P1-a 修复链路 —— 卡死任务回收后能被 QueryPendingUploads 命中。 +func TestReclaimThenQueryable(t *testing.T) { + db := newUploadTestDB(t) + stale := seedUploading(t, db, time.Hour, 1) + + if got, _ := QueryPendingUploads(db, 10); len(got) != 0 { + t.Fatalf("回收前 Uploading 任务不应被补传查询命中,实际 %d", len(got)) + } + if _, err := ReclaimStaleUploading(db, 35*time.Minute); err != nil { + t.Fatal(err) + } + got, err := QueryPendingUploads(db, 10) + if err != nil { + t.Fatal(err) + } + if len(got) != 1 || got[0].ID != stale.ID { + t.Fatalf("回收后任务应可被补传命中,实际 %d 条", len(got)) + } +} From 492177c55d87420aeb80cc0d616de70df0b698a2 Mon Sep 17 00:00:00 2001 From: "richard.li" Date: Fri, 22 May 2026 15:03:49 +0800 Subject: [PATCH 3/4] =?UTF-8?q?fix(upload):=20Phase=202=20=E5=8A=A0=20pend?= =?UTF-8?q?ing=20=E7=9B=AE=E5=BD=95=E7=A3=81=E7=9B=98=E6=B0=B4=E4=BD=8D?= =?UTF-8?q?=E4=BF=9D=E6=8A=A4(P2-b)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit MinIO 长期故障时 pending 文件堆积会撑爆本地盘、拖垮新录制。加水位保护: - UploadConfig 加 PendingMaxSizeMB / PendingMaxFiles / PendingDiskMinFreeMB (默认 0=不限, 向后兼容); 接入 ServerConfig.Upload, InitUploadManager 从配置读取(原为硬编码) - MoveToPendingDir 入口做水位检查, 超阈值返回 ErrPendingDirFull —— 宁可 丢一个录像也不让磁盘写满拖垮全盘 - GetDiskFreeBytes 分平台实现(statfs_unix / statfs_windows, Windows 降级) - record.go recoverToPending / recoverFastPathFailure 识别 ErrPendingDirFull 触发 AlarmDiskSpaceFull 告警 - UploadRetryScheduler.Tick 加 pending 水位巡检, 达阈值 80% 告警 upload_manager_test.go 覆盖用量统计 / 水位拒绝 / 无限制兼容, 4/4 PASS。 --- pkg/storage/statfs_unix.go | 14 +++++ pkg/storage/statfs_windows.go | 11 ++++ pkg/storage/upload_manager.go | 90 ++++++++++++++++++++++++++++++ pkg/storage/upload_manager_test.go | 80 ++++++++++++++++++++++++++ plugin/mp4/pkg/record.go | 10 ++++ server.go | 8 +-- upload_retry.go | 8 +++ 7 files changed, 216 insertions(+), 5 deletions(-) create mode 100644 pkg/storage/statfs_unix.go create mode 100644 pkg/storage/statfs_windows.go create mode 100644 pkg/storage/upload_manager_test.go diff --git a/pkg/storage/statfs_unix.go b/pkg/storage/statfs_unix.go new file mode 100644 index 00000000..c4eeb5a2 --- /dev/null +++ b/pkg/storage/statfs_unix.go @@ -0,0 +1,14 @@ +//go:build !windows + +package storage + +import "syscall" + +// GetDiskFreeBytes 返回 path 所在文件系统对非特权用户的可用字节数。 +func GetDiskFreeBytes(path string) (uint64, error) { + var st syscall.Statfs_t + if err := syscall.Statfs(path, &st); err != nil { + return 0, err + } + return uint64(st.Bavail) * uint64(st.Bsize), nil +} diff --git a/pkg/storage/statfs_windows.go b/pkg/storage/statfs_windows.go new file mode 100644 index 00000000..070a4dbb --- /dev/null +++ b/pkg/storage/statfs_windows.go @@ -0,0 +1,11 @@ +//go:build windows + +package storage + +import "errors" + +// GetDiskFreeBytes 在 Windows 暂不支持磁盘剩余空间检查,返回错误使调用方降级 +// (checkPendingCapacity 在 err != nil 时跳过磁盘检查,仅按文件数 / 总大小限制)。 +func GetDiskFreeBytes(_ string) (uint64, error) { + return 0, errors.New("GetDiskFreeBytes not supported on windows") +} diff --git a/pkg/storage/upload_manager.go b/pkg/storage/upload_manager.go index 2365b24a..786934ad 100644 --- a/pkg/storage/upload_manager.go +++ b/pkg/storage/upload_manager.go @@ -2,6 +2,8 @@ package storage import ( "context" + "errors" + "fmt" "io" "log" "os" @@ -15,6 +17,11 @@ var ( pendingDir string maxConcurrent int + // pending 目录水位限制(0=不限),由 InitUploadManager 从 UploadConfig 注入。 + pendingMaxSizeBytes int64 + pendingMaxFiles int + pendingDiskMinFree int64 + // trailerSem 预留: 限制并发 trailer 写盘槽位数 (mp4/flv 等录制 plugin 共用). // // 当前 **未在生产路径调用** — 原因: @@ -37,12 +44,18 @@ var ( OnUploadFailed func(localPath, objectKey, storageType string, fileSize int64, metadata map[string]string, err error) ) +// ErrPendingDirFull 表示 pending 暂存目录已达水位上限,拒绝再接收文件。 +var ErrPendingDirFull = errors.New("pending dir full") + // UploadConfig 上传管理配置 type UploadConfig struct { MaxConcurrentUploads int `desc:"最大并发上传数" default:"4"` MaxConcurrentTrailerWrites int `desc:"[预留] 最大并发 trailer 写盘槽位数. 当前 trailer queue 是 single-threaded, 此项不影响行为; 留作未来 worker-pool 实现的接口" default:"8"` TrailerWriteRateMBps int `desc:"trailer 重写写盘限速 (MB/s), 控制 record stop 时磁盘 burst; 0=不限速 (默认)" default:"0"` PendingDir string `desc:"上传失败文件暂存目录" default:"pending_uploads"` + PendingMaxSizeMB int `desc:"pending 目录总大小上限(MB), 超过则拒绝新文件暂存(可能丢录像); 0=不限" default:"0"` + PendingMaxFiles int `desc:"pending 目录文件数上限, 超过则拒绝; 0=不限" default:"0"` + PendingDiskMinFreeMB int `desc:"pending 所在磁盘最低剩余空间(MB), 低于则拒绝; 0=不检查" default:"0"` } // InitUploadManager 初始化上传管理器(并发控制 + 暂存目录) @@ -69,6 +82,9 @@ func InitUploadManager(cfg UploadConfig) { cfg.PendingDir = "pending_uploads" } pendingDir = cfg.PendingDir + pendingMaxSizeBytes = int64(cfg.PendingMaxSizeMB) * 1024 * 1024 + pendingMaxFiles = cfg.PendingMaxFiles + pendingDiskMinFree = int64(cfg.PendingDiskMinFreeMB) * 1024 * 1024 if err := os.MkdirAll(pendingDir, 0755); err != nil { log.Printf("[storage] failed to create pending dir %s: %v", pendingDir, err) } @@ -115,6 +131,9 @@ func MoveToPendingDir(srcPath string) (string, error) { if pendingDir == "" { return srcPath, nil // 未配置暂存目录,保留原路径 } + if err := checkPendingCapacity(); err != nil { + return "", err + } if err := os.MkdirAll(pendingDir, 0755); err != nil { return "", err } @@ -157,6 +176,77 @@ func GetPendingDir() string { return pendingDir } +// GetPendingDirUsage 统计 pending 目录的总字节数与文件数。 +func GetPendingDirUsage() (totalBytes int64, fileCount int, err error) { + if pendingDir == "" { + return 0, 0, nil + } + err = filepath.WalkDir(pendingDir, func(_ string, d os.DirEntry, walkErr error) error { + if walkErr != nil { + return walkErr + } + if d.IsDir() { + return nil + } + info, e := d.Info() + if e != nil { + return e + } + totalBytes += info.Size() + fileCount++ + return nil + }) + if os.IsNotExist(err) { + return 0, 0, nil // 目录尚未创建,视为空 + } + return totalBytes, fileCount, err +} + +// checkPendingCapacity 检查 pending 目录是否还能容纳新文件。 +// 任一已配置阈值(>0)被突破即返回包装 ErrPendingDirFull 的错误;全部为 0 时不限制。 +func checkPendingCapacity() error { + if pendingMaxSizeBytes <= 0 && pendingMaxFiles <= 0 && pendingDiskMinFree <= 0 { + return nil + } + if pendingMaxSizeBytes > 0 || pendingMaxFiles > 0 { + total, count, err := GetPendingDirUsage() + if err != nil { + return err + } + if pendingMaxSizeBytes > 0 && total >= pendingMaxSizeBytes { + return fmt.Errorf("%w: 已用 %d 字节 >= 上限 %d 字节", ErrPendingDirFull, total, pendingMaxSizeBytes) + } + if pendingMaxFiles > 0 && count >= pendingMaxFiles { + return fmt.Errorf("%w: 已有 %d 文件 >= 上限 %d", ErrPendingDirFull, count, pendingMaxFiles) + } + } + if pendingDiskMinFree > 0 { + if free, err := GetDiskFreeBytes(pendingDir); err == nil && free < uint64(pendingDiskMinFree) { + return fmt.Errorf("%w: 磁盘剩余 %d 字节 < 下限 %d 字节", ErrPendingDirFull, free, pendingDiskMinFree) + } + } + return nil +} + +// PendingWatermarkExceeded 检查 pending 目录用量是否达到告警水位(已配置阈值的 80%)。 +// 返回是否超水位及描述;未配置 size/files 阈值时恒返回 false。 +func PendingWatermarkExceeded() (exceeded bool, detail string) { + if pendingMaxSizeBytes <= 0 && pendingMaxFiles <= 0 { + return false, "" + } + total, count, err := GetPendingDirUsage() + if err != nil { + return false, "" + } + if pendingMaxSizeBytes > 0 && total >= pendingMaxSizeBytes*8/10 { + return true, fmt.Sprintf("已用 %d 字节,达上限 %d 的 80%%", total, pendingMaxSizeBytes) + } + if pendingMaxFiles > 0 && count >= pendingMaxFiles*8/10 { + return true, fmt.Sprintf("已有 %d 文件,达上限 %d 的 80%%", count, pendingMaxFiles) + } + return false, "" +} + // AcquireTrailerSlot 获取一个 trailer 写盘槽位, 阻塞直到有可用槽位或 ctx 取消. // 配对 ReleaseTrailerSlot. 调用方通常在 record stop 流程进入 trailer flush 前 acquire, // 在 task Dispose / Run 末尾 defer Release. diff --git a/pkg/storage/upload_manager_test.go b/pkg/storage/upload_manager_test.go new file mode 100644 index 00000000..d74b3942 --- /dev/null +++ b/pkg/storage/upload_manager_test.go @@ -0,0 +1,80 @@ +package storage + +import ( + "errors" + "os" + "path/filepath" + "testing" +) + +// TestGetPendingDirUsage 验证 pending 目录的字节数 / 文件数统计。 +func TestGetPendingDirUsage(t *testing.T) { + dir := t.TempDir() + InitUploadManager(UploadConfig{PendingDir: dir}) + + for name, sz := range map[string]int{"a.mp4": 100, "b.mp4": 200, "c.mp4": 300} { + if err := os.WriteFile(filepath.Join(dir, name), make([]byte, sz), 0644); err != nil { + t.Fatalf("write seed file: %v", err) + } + } + + total, count, err := GetPendingDirUsage() + if err != nil { + t.Fatalf("GetPendingDirUsage: %v", err) + } + if count != 3 { + t.Errorf("文件数期望 3,实际 %d", count) + } + if total != 600 { + t.Errorf("总字节数期望 600,实际 %d", total) + } +} + +// TestMoveToPendingDir_FileLimit 验证 PendingMaxFiles 水位: +// 暂存满 2 个后,第 3 个被 ErrPendingDirFull 拒绝。 +func TestMoveToPendingDir_FileLimit(t *testing.T) { + srcDir := t.TempDir() + pendingDir := t.TempDir() + InitUploadManager(UploadConfig{PendingDir: pendingDir, PendingMaxFiles: 2}) + + mkSrc := func(name string) string { + p := filepath.Join(srcDir, name) + if err := os.WriteFile(p, []byte("x"), 0644); err != nil { + t.Fatalf("write src: %v", err) + } + return p + } + + if _, err := MoveToPendingDir(mkSrc("a.mp4")); err != nil { + t.Fatalf("第 1 个应成功: %v", err) + } + if _, err := MoveToPendingDir(mkSrc("b.mp4")); err != nil { + t.Fatalf("第 2 个应成功: %v", err) + } + if _, err := MoveToPendingDir(mkSrc("c.mp4")); !errors.Is(err, ErrPendingDirFull) { + t.Fatalf("第 3 个应返回 ErrPendingDirFull,实际 %v", err) + } +} + +// TestMoveToPendingDir_NoLimit 验证未配置水位(默认 0)时行为不变。 +func TestMoveToPendingDir_NoLimit(t *testing.T) { + srcDir := t.TempDir() + pendingDir := t.TempDir() + InitUploadManager(UploadConfig{PendingDir: pendingDir}) + + for _, name := range []string{"a", "b", "c", "d", "e"} { + p := filepath.Join(srcDir, name+".mp4") + os.WriteFile(p, []byte("x"), 0644) + if _, err := MoveToPendingDir(p); err != nil { + t.Fatalf("未配水位时 %s 应成功: %v", name, err) + } + } +} + +// TestGetDiskFreeBytes 验证磁盘剩余空间查询:Unix 返回 >0,Windows 降级返回 err。 +func TestGetDiskFreeBytes(t *testing.T) { + free, err := GetDiskFreeBytes(t.TempDir()) + if err == nil && free == 0 { + t.Error("Unix 下磁盘可用空间应 > 0") + } +} diff --git a/plugin/mp4/pkg/record.go b/plugin/mp4/pkg/record.go index 23b9ccb5..adb609df 100644 --- a/plugin/mp4/pkg/record.go +++ b/plugin/mp4/pkg/record.go @@ -172,6 +172,11 @@ func (t *writeTrailerTask) Run() (err error) { pendingPath, moveErr := storage.MoveToPendingDir(tempPath) if moveErr != nil { t.Error("move to pending dir failed", "err", moveErr) + if errors.Is(moveErr, storage.ErrPendingDirFull) { + m7s.RaiseUploadAlarm(t.db, config.AlarmDiskSpaceFull, + "pending dir full", t.streamPath, t.filePath, + "pending 暂存目录已满,本录像无法暂存补传可能丢失: "+moveErr.Error()) + } return } tempOwned = false // 已移走,不需 defer 删除 @@ -295,6 +300,11 @@ func (t *writeTrailerTask) recoverFastPathFailure(localPath string, fileSize int pendingPath, moveErr := storage.MoveToPendingDir(localPath) if moveErr != nil { t.Error("move to pending dir failed", "err", moveErr) + if errors.Is(moveErr, storage.ErrPendingDirFull) { + m7s.RaiseUploadAlarm(t.db, config.AlarmDiskSpaceFull, + "pending dir full", t.streamPath, t.filePath, + "pending 暂存目录已满,本录像无法暂存补传可能丢失: "+moveErr.Error()) + } return } metadata := map[string]string{"video-size-bytes": fmt.Sprintf("%d", fileSize)} diff --git a/server.go b/server.go index ec840c46..8e6b61cc 100644 --- a/server.go +++ b/server.go @@ -80,6 +80,7 @@ type ( } `desc:"用户列表,仅在启用登录机制时生效"` } `desc:"管理员界面配置"` Storage map[string]any + Upload storage.UploadConfig `desc:"录像上传管理配置"` } WaitStream struct { StreamPath string @@ -298,11 +299,8 @@ func (s *Server) Start() (err error) { } s.LogHandler.SetLevel(ParseLevel(s.config.LogLevel)) s.initStorage() - // 初始化上传并发控制器 - storage.InitUploadManager(storage.UploadConfig{ - MaxConcurrentUploads: 4, - PendingDir: "pending_uploads", - }) + // 初始化上传并发控制器(配置来自 ServerConfig.Upload,0 值由 InitUploadManager 兜底) + storage.InitUploadManager(s.ServerConfig.Upload) err = debug.SetCrashOutput(util.InitFatalLog(s.FatalDir), debug.CrashOptions{}) if err != nil { s.Error("SetCrashOutput", "error", err) diff --git a/upload_retry.go b/upload_retry.go index c3bc6ca4..b2b6fc78 100644 --- a/upload_retry.go +++ b/upload_retry.go @@ -7,6 +7,7 @@ import ( "time" task "github.com/langhuihui/gotask" + "m7s.live/v5/pkg/config" "m7s.live/v5/pkg/storage" ) @@ -38,6 +39,13 @@ func (u *UploadRetryScheduler) Tick(any) { u.Info("reclaimed stale uploading tasks", "count", n) } + // pending 目录水位巡检:接近上限即告警(RaiseUploadAlarm 自带去重,不会刷屏) + if warn, detail := storage.PendingWatermarkExceeded(); warn { + RaiseUploadAlarm(u.s.DB, config.AlarmDiskSpaceFull, + "pending dir near full", "", storage.GetPendingDir(), + "pending 暂存目录接近容量上限: "+detail) + } + // 查询待重试的任务(每次最多处理 20 个,避免单次过多) tasks, err := QueryPendingUploads(u.s.DB, 20) if err != nil { From 9173f98197dcd1a9890989944ba8efbb5a68ac7d Mon Sep 17 00:00:00 2001 From: "richard.li" Date: Fri, 22 May 2026 15:14:23 +0800 Subject: [PATCH 4/4] =?UTF-8?q?fix(upload):=20Phase=203=20=E5=8A=A0?= =?UTF-8?q?=E8=A1=A5=E4=BC=A0=E8=80=97=E5=B0=BD=E5=91=8A=E8=AD=A6=20+=20?= =?UTF-8?q?=E4=BA=BA=E5=B7=A5=E4=BB=8B=E5=85=A5(P2-a)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 补传 retry_count 耗尽(达 MaxRetries)后系统永久放弃, 原本无告警、无介入: - retryUpload 检测耗尽 -> RaiseUploadAlarm 写 AlarmStorageException 告警 - 新增 QueryExhaustedUploads: 查耗尽任务(retry_count >= max_retries) - 新增 ResetUploadForRetry: 重置任务重试状态, 重新拉入补传循环 - mp4 plugin 加运维 HTTP 端点(经 RegisterHandler): GET /mp4/api/upload/exhausted 列出耗尽任务 POST /mp4/api/upload/retry?id=N 一键重新拉起 upload_task_test.go 加 QueryExhaustedUploads / ResetUploadForRetry 单测, 全 5 例 PASS。 Task 3.4(孤儿 RecordStream 自动补偿队列)未做, plan 已标注修订: 孤儿场景罕见 + Phase 1 已有感知告警, 完整自动补偿成本超 plan 设想, 留作后续专项。 --- .../2026-05-22-upload-reliability-p1p2.md | 9 +-- plugin/mp4/index.go | 2 + plugin/mp4/upload_admin.go | 56 ++++++++++++++++++ upload_retry.go | 6 ++ upload_task.go | 22 +++++++ upload_task_test.go | 59 +++++++++++++++++++ 6 files changed, 150 insertions(+), 4 deletions(-) create mode 100644 plugin/mp4/upload_admin.go diff --git a/docs/superpowers/plans/2026-05-22-upload-reliability-p1p2.md b/docs/superpowers/plans/2026-05-22-upload-reliability-p1p2.md index 1d7f6939..2640a60c 100644 --- a/docs/superpowers/plans/2026-05-22-upload-reliability-p1p2.md +++ b/docs/superpowers/plans/2026-05-22-upload-reliability-p1p2.md @@ -154,10 +154,11 @@ docs/superpowers/plans/2026-05-22-upload-reliability-p1p2.md 本文件 ## Task 3.4: 孤儿 RecordStream 持久化补偿 -- [ ] 新建表 `record_stream_recovery`(字段:`RecordStream` JSON 快照、retry_count、next_retry_at、created_at),加进 `server.go:334` AutoMigrate。 -- [ ] Task 1.6 内联重试仍失败的孤儿 → 序列化 JSON 存该表。 -- [ ] `UploadRetryScheduler.Tick` 增加:扫 `record_stream_recovery`,对每条重试 `db.Save(&RecordStream)`,成功删记录、失败退避。 -- 依赖:Task 1.6、Phase 2。风险:Medium —— 新表入 AutoMigrate;`RecordStream` JSON round-trip 须保留 `ID`(`Save` upsert 不重复行)。 +> **执行修订(2026-05-22):本 Task 未实现,留作后续。** 理由:① 孤儿场景(上传成功但 `record_streams` 入库失败)发生概率低 —— DB 写仅一次带 10s 超时的 `Save`;② Phase 1 Task 1.6 已实现「感知 + 告警」(`WriteTailDeferred` 返回 error → `reportOrphan` 错误日志 + `alarm_info` 告警),运维可见、可手动补;③ 完整自动补偿需改 `WriteTailDeferred` 设计以向 `writeTrailerTask` 暴露 `RecordStream` 快照(现仅暴露 `dbWrite` 闭包),成本超出 plan 设想。P1-b 当前为「可感知、不自动修复」,自动补偿留后续专项。 + +- [ ] ~~新建表 `record_stream_recovery`(字段:`RecordStream` JSON 快照、retry_count、next_retry_at、created_at),加进 `server.go:334` AutoMigrate。~~ +- [ ] ~~Task 1.6 内联重试仍失败的孤儿 → 序列化 JSON 存该表。~~ +- [ ] ~~`UploadRetryScheduler.Tick` 增加:扫 `record_stream_recovery`,对每条重试 `db.Save(&RecordStream)`,成功删记录、失败退避。~~ ## Task 3.5: Phase 3 验收 diff --git a/plugin/mp4/index.go b/plugin/mp4/index.go index 0af2a5bf..35e748ff 100644 --- a/plugin/mp4/index.go +++ b/plugin/mp4/index.go @@ -62,6 +62,8 @@ func (p *MP4Plugin) RegisterHandler() map[string]http.HandlerFunc { "/extract/compressed/{streamPath...}": p.extractCompressedVideoHandel, "/extract/gop/{streamPath...}": p.extractGopVideoHandel, "/snap/{streamPath...}": p.snapHandel, + "/api/upload/exhausted": p.handleListExhaustedUploads, + "/api/upload/retry": p.handleRetryUpload, } } diff --git a/plugin/mp4/upload_admin.go b/plugin/mp4/upload_admin.go new file mode 100644 index 00000000..9ea07f17 --- /dev/null +++ b/plugin/mp4/upload_admin.go @@ -0,0 +1,56 @@ +package plugin_mp4 + +import ( + "encoding/json" + "net/http" + "strconv" + + m7s "m7s.live/v5" +) + +// 录像上传补传的运维 HTTP 端点,经 MP4Plugin.RegisterHandler 注册(见 index.go), +// 实际路径带 /mp4 前缀: +// +// GET /mp4/api/upload/exhausted 列出补传次数耗尽、已被系统放弃的任务 +// POST /mp4/api/upload/retry?id=N 重置指定任务,重新拉入补传循环 +// +// 用于存储故障修复后,运维查看并一键重新拉起被放弃的录像上传。 + +func (p *MP4Plugin) handleListExhaustedUploads(w http.ResponseWriter, r *http.Request) { + if p.DB == nil { + http.Error(w, "database not enabled", http.StatusServiceUnavailable) + return + } + tasks, err := m7s.QueryExhaustedUploads(p.DB, 200) + if err != nil { + http.Error(w, err.Error(), http.StatusInternalServerError) + return + } + w.Header().Set("Content-Type", "application/json; charset=utf-8") + json.NewEncoder(w).Encode(map[string]any{ + "code": 0, + "count": len(tasks), + "data": tasks, + }) +} + +func (p *MP4Plugin) handleRetryUpload(w http.ResponseWriter, r *http.Request) { + if p.DB == nil { + http.Error(w, "database not enabled", http.StatusServiceUnavailable) + return + } + id, err := strconv.ParseUint(r.URL.Query().Get("id"), 10, 64) + if err != nil { + http.Error(w, "invalid or missing query param: id", http.StatusBadRequest) + return + } + if err := m7s.ResetUploadForRetry(p.DB, uint(id)); err != nil { + http.Error(w, err.Error(), http.StatusInternalServerError) + return + } + w.Header().Set("Content-Type", "application/json; charset=utf-8") + json.NewEncoder(w).Encode(map[string]any{ + "code": 0, + "message": "upload task reset for retry", + }) +} diff --git a/upload_retry.go b/upload_retry.go index b2b6fc78..0ec0a190 100644 --- a/upload_retry.go +++ b/upload_retry.go @@ -110,6 +110,12 @@ func (u *UploadRetryScheduler) retryUpload(ut UploadTask) { "retryCount", ut.RetryCount+1, "err", err) MarkUploadRetryFailed(u.s.DB, ut.ID, ut.RetryCount, err) + // 补传次数耗尽:不再被 QueryPendingUploads 命中,告警通知运维介入 + if ut.RetryCount+1 >= ut.MaxRetries { + RaiseUploadAlarm(u.s.DB, config.AlarmStorageException, + "upload retry exhausted", ut.StreamPath, ut.LocalPath, + "补传次数已耗尽,文件待人工处理: "+ut.ObjectKey) + } return } diff --git a/upload_task.go b/upload_task.go index ea055c24..6351da44 100644 --- a/upload_task.go +++ b/upload_task.go @@ -109,6 +109,28 @@ func QueryPendingUploads(db *gorm.DB, limit int) ([]UploadTask, error) { return tasks, err } +// QueryExhaustedUploads 查询补传次数已耗尽(retry_count >= max_retries)的失败任务, +// 供运维查看与人工介入。这类任务已不会被 QueryPendingUploads 命中。 +func QueryExhaustedUploads(db *gorm.DB, limit int) ([]UploadTask, error) { + var tasks []UploadTask + err := db.Where("status = ? AND retry_count >= max_retries", UploadStatusFailed). + Order("updated_at DESC"). + Limit(limit). + Find(&tasks).Error + return tasks, err +} + +// ResetUploadForRetry 重置一个任务的重试状态,使其重新进入补传循环。 +// 用于运维修复存储故障后手动重新拉起已耗尽的任务。 +func ResetUploadForRetry(db *gorm.DB, taskID uint) error { + return db.Model(&UploadTask{}).Where("id = ?", taskID).Updates(map[string]any{ + "status": UploadStatusFailed, + "retry_count": 0, + "next_retry_at": time.Now(), + "error_message": "", + }).Error +} + // MarkUploading 原子抢占任务:仅当任务仍为 Failed 时置 Uploading 并记录开始时间。 // 返回 true 表示本次抢占成功(可继续上传);false 表示已被其他 goroutine 抢占。 func MarkUploading(db *gorm.DB, taskID uint) (claimed bool) { diff --git a/upload_task_test.go b/upload_task_test.go index cdb52a8d..3afec736 100644 --- a/upload_task_test.go +++ b/upload_task_test.go @@ -132,3 +132,62 @@ func TestReclaimThenQueryable(t *testing.T) { t.Fatalf("回收后任务应可被补传命中,实际 %d 条", len(got)) } } + +// TestQueryExhaustedUploads:retry_count >= max_retries 的任务被查出,未耗尽的不被查出。 +func TestQueryExhaustedUploads(t *testing.T) { + db := newUploadTestDB(t) + exhausted := UploadTask{ + LocalPath: "/tmp/e.mp4", ObjectKey: "e.mp4", + Status: UploadStatusFailed, RetryCount: 10, MaxRetries: 10, + } + if err := db.Create(&exhausted).Error; err != nil { + t.Fatalf("seed exhausted: %v", err) + } + pending := UploadTask{ + LocalPath: "/tmp/p.mp4", ObjectKey: "p.mp4", + Status: UploadStatusFailed, RetryCount: 3, MaxRetries: 10, + } + if err := db.Create(&pending).Error; err != nil { + t.Fatalf("seed pending: %v", err) + } + + got, err := QueryExhaustedUploads(db, 10) + if err != nil { + t.Fatal(err) + } + if len(got) != 1 || got[0].ID != exhausted.ID { + t.Fatalf("应只查出 1 个耗尽任务,实际 %d 条", len(got)) + } +} + +// TestResetUploadForRetry:耗尽任务重置后 retry_count 归零、可被补传查询命中。 +func TestResetUploadForRetry(t *testing.T) { + db := newUploadTestDB(t) + ut := UploadTask{ + LocalPath: "/tmp/r.mp4", ObjectKey: "r.mp4", + Status: UploadStatusFailed, RetryCount: 10, MaxRetries: 10, + NextRetryAt: time.Now().Add(time.Hour), + } + if err := db.Create(&ut).Error; err != nil { + t.Fatalf("seed: %v", err) + } + + if got, _ := QueryPendingUploads(db, 10); len(got) != 0 { + t.Fatalf("重置前耗尽任务不应被补传命中,实际 %d", len(got)) + } + if err := ResetUploadForRetry(db, ut.ID); err != nil { + t.Fatal(err) + } + + var got UploadTask + db.First(&got, ut.ID) + if got.RetryCount != 0 { + t.Errorf("重置后 retry_count 应为 0,实际 %d", got.RetryCount) + } + if got.Status != UploadStatusFailed { + t.Errorf("重置后 status 应为 Failed,实际 %d", got.Status) + } + if pend, _ := QueryPendingUploads(db, 10); len(pend) != 1 { + t.Fatalf("重置后任务应可被补传命中,实际 %d", len(pend)) + } +}