diff --git a/.agents/skills/clickhouse-batchwriter/SKILL.md b/.agents/skills/clickhouse-batchwriter/SKILL.md index e9f29c80..942739e0 100644 --- a/.agents/skills/clickhouse-batchwriter/SKILL.md +++ b/.agents/skills/clickhouse-batchwriter/SKILL.md @@ -7,7 +7,7 @@ description: "Wavelet 项目专用:当新增或修改 ClickHouse 批量写入 开始前阅读根目录 `AGENTS.md`。ClickHouse 是辅助 OLAP 存储,**厌恶高频单条写入**(过多小 part);写入路径必须优先批量或异步聚合。 -DDL 与表结构变更见 `database-migration` 技能;本技能只覆盖**运行时写入架构**。 +DDL 与表结构变更见 `database-migration` 技能。日志/分析用途表的判定、三库回落与切换见 `logstore` 技能。本技能只覆盖**运行时写入架构**。 ## 分层职责 @@ -17,7 +17,7 @@ DDL 与表结构变更见 `database-migration` 技能;本技能只覆盖**运 | 批量框架 | `internal/infra/persistence/batchwriter/` | 泛型队列 + 按条数/时间 flush + 非阻塞入队 + 优雅停机;**各业务域独立实例** | | Model | `internal/model/analytics/` | 列定义、`TableName()`、`BatchInsertSQL()`(及可选 `InsertColumns()`) | | Repository | `internal/repository/analytics/` | `BatchInsert*` / `BatchInsertNodeAccessLogs` 等;`PrepareBatch` + 多行 `Append` + 一次 `Send` | -| Apps | `internal/apps//` | 采集、入队、背压;`FlushFunc` 只调 repository,不写 SQL、不 `PrepareBatch` | +| Apps | `internal/apps//` | 采集、入队、背压;`FlushFunc` 只调 logstore / repository,不写 SQL、不 `PrepareBatch` | | 装配 | `internal/platform/bootstrap/bootstrap.go` | 进程启动时调用 `Writer.Start`;初始化时需调用 `lifecycle.OnShutdown` 挂载停机钩子 | | 生命周期 | `internal/platform/lifecycle/lifecycle.go` | 统一协调全局并发优雅停机,业务包无需在 `bootstrap.go` 中硬编码 `Stop` 逻辑 | @@ -37,6 +37,7 @@ writer.Stop(stopCtx) // close 队列 + drain + 最终 flush - `QueueSize`: 10_000 - `MaxBatchSize`: 1_000 +- `MinBatchSize`: 50(未达阈值则跳过按时间 flush,除非设了 `MaxFlushWait`) - `FlushInterval`: 1s 各域可独立覆盖;可观测低频指标可用更小 `MaxBatchSize`(如 100)与更长 `FlushInterval`(如 2–5s),但**不要**退化为逐条 `Send`。 @@ -49,7 +50,8 @@ writer.Stop(stopCtx) // close 队列 + drain + 最终 flush ### FlushFunc 规范 - 签名:`func(ctx context.Context, items []T) error` -- 内部调用 `internal/repository/analytics` 的 `BatchInsert*`(传入 `[]analyticsmodel.X`) +- **日志/分析用途表**:`logstore.Active(ctx)` 再调对应 `BatchInsert*`。禁止 apps 直连 `analyticsrepo` 或 `db.ChConn`。 +- 仅 CH、无需主库回落的分析表:才直接调 `repository/analytics` 的 `BatchInsert*`。 - 在 flush 边界记录一次错误日志,不要把 DB 驱动错误直接暴露给 HTTP 客户端 - `Start` 使用 `context.WithoutCancel(parent)`,避免请求 ctx 取消中断后台 flush @@ -57,13 +59,11 @@ writer.Stop(stopCtx) // close 队列 + drain + 最终 flush 每个业务域拥有自己的 `Writer`、配置与 `FlushFunc`: -| 域 | 表 | 现状 | 目标形态 | -| :--- | :--- | :--- | :--- | -| 管理端审计 | `w_user_access_logs` | `risk_control` → `batchwriter` + `analyticsrepo.BatchInsert` | 已接入 | -| 边缘访问日志 | `of_node_access_logs` | `openflare/chwriter` 异步 flush | 已接入 | -| 可观测时序 | `of_node_metric_snapshots` 等 5 表 | `openflare/chwriter` 五表独立 writer + 进程内短 TTL 去重 | 已接入 | +| 域 | 表 | 写入路径 | +| :--- | :--- | :--- | +| 管理端审计 | `w_user_access_logs` | `risk_control` → `batchwriter` → `logstore.Active` | -**不要**把 audit、access log、observability 并入同一 channel。 +**不要**把不同日志域并入同一 channel。新日志表先按 `logstore` skill 判定,再为本域建独立 writer。 ## 新增 ClickHouse 写入工作流 @@ -73,8 +73,9 @@ writer.Stop(stopCtx) // close 队列 + drain + 最终 flush - `len(items)==0` 直接返回 - `db.ChConn == nil` 返回明确错误 - 一次 `PrepareBatch` → 循环 `Append` → 一次 `Send` -4. **Writer 胶水**(`internal/apps//` 或 `internal/repository/analytics/_writer.go`): +4. **Writer 胶水**(`internal/apps//`): - `New` + `Start`,并在初始化逻辑内通过 `lifecycle.OnShutdown("your_writer_name", Stop)` 注册停机回调 + - 日志表的 `FlushFunc` 调 `logstore.Active`(见 `logstore` skill) - 业务路径 `TryEnqueue`;HTTP 背压用 `IsFull()` 5. **测试**: - repository:mock `ChConn` 验证 `BatchInsertSQL` 与 append 列数 @@ -86,8 +87,7 @@ writer.Stop(stopCtx) // close 队列 + drain + 最终 flush | 场景 | 推荐策略 | | :--- | :--- | | 管理端 API 审计 | 队列满 → `IsFull()` 触发 429(见 `risk_control` middleware) | -| Agent 心跳指标 | 队列满 → `WithDropHandler` 记 warn;不阻塞心跳响应 | -| 边缘 access log | 优先扩大队列与 batch;必要时丢弃最旧或采样 | +| 可丢弃的高频日志 | 队列满 → `WithDropHandler` 记 warn;不阻塞请求 | ## 禁止写法 @@ -123,14 +123,9 @@ var globalChan chan any ```go // internal/platform/bootstrap/bootstrap.go(示意) -var userAccessLogWriter *batchwriter.Writer[*analytics.UserAccessLog] - func RegisterAPI(ctx context.Context) { - // ... - if config.Config.ClickHouse.Enabled { - initUserAccessLogWriter(ctx) // Start writer - risk_control.BindWriter(userAccessLogWriter) // 或逐步替换 InitLogWriter - } + // 日志 writer 不依赖 clickhouse.enabled:flush 时由 logstore 选库 + risk_control.InitLogWriter(ctx) } ``` @@ -149,15 +144,14 @@ make code-check - flush 按 `MaxBatchSize` 与 `FlushInterval` 触发 - `Stop` 能 drain 队列内剩余项 - repository 层无 goroutine、无 channel -- `clickhouse.enabled: false` 时不 `Start` writer、不入队 +- 日志表:`clickhouse.enabled: false` 时 writer 仍 `Start`,flush 走主库 logstore +- 仅 CH 的分析表:未启用 CH 时不要 `Start`、不要入队 ## 相关文件速查 - 框架:`internal/infra/persistence/batchwriter/{config,writer,errs}.go` - 连接:`internal/infra/persistence/clickhouse.go` - 审计写入:`internal/apps/risk_control/logics.go` -- OpenFlare 写入胶水:`internal/apps/openflare/chwriter/writer.go` -- 节点访问日志 repository:`internal/repository/analytics/node_access_log_writer.go` -- 可观测 repository:`internal/repository/analytics/node_observability_writer.go` +- 日志抽象:`internal/repository/logstore` - 生命周期管理器:`internal/platform/lifecycle/lifecycle.go` - Bootstrap:`internal/platform/bootstrap/bootstrap.go` \ No newline at end of file diff --git a/.agents/skills/database-migration/SKILL.md b/.agents/skills/database-migration/SKILL.md index 59f416f4..5a528b7f 100644 --- a/.agents/skills/database-migration/SKILL.md +++ b/.agents/skills/database-migration/SKILL.md @@ -81,7 +81,7 @@ make code-check ClickHouse 是**辅助 OLAP 存储**,与 PostgreSQL/SQLite 主库**完全独立**的迁移与访问管线: - 主库(PG/SQLite):业务事务数据、`goose_db_version`、双方言 SQL。 -- 分析库(ClickHouse):访问日志、统计聚合等分析型数据、`goose_clickhouse_version`、单方言 SQL。 +- 分析库(ClickHouse):分析型数据、`goose_clickhouse_version`、单方言 SQL。日志用途表还必须在主库建回落并走 `logstore`(见该 skill);CH 目录仍只放 CH DDL。 **不要**把 ClickHouse 表结构混入 PG/SQLite 迁移目录,也**不要**在 `support-files/`、`internal/apps/` 或 `internal/repository/` 中手写 DDL。 @@ -118,7 +118,7 @@ ClickHouse 是**辅助 OLAP 存储**,与 PostgreSQL/SQLite 主库**完全独 1. **Model**:在 `internal/model/analytics/` 定义 struct,`gorm:"column:..."` 与 DDL 列名一一对应;实现 `TableName()`,批量写入表可提供 `InsertColumns()` / `BatchInsertSQL()`。 2. **Goose SQL**:在 `internal/infra/persistence/migrator/goose/clickhouse/` 新增递增版本文件(格式同主库,如 `YYYYMMDDNNNN_create_xxx.sql`),编写 `-- +goose Up` / `-- +goose Down`。 3. **Repository**:在 `internal/repository/analytics/` 实现 `BatchInsert*`(`db.ChConn` 一次 `PrepareBatch` + 多行 `Append` + 一次 `Send`)与查询(`db.ChDB`);连接未初始化时返回明确错误,**不要**在 handler 写 SQL,**不要**在 repository 内维护 channel/goroutine。 -4. **Apps**:在 `internal/apps//` 编排采集与入队;高频写入通过 `internal/infra/persistence/batchwriter` 各域独立实例异步 flush(详见 `clickhouse-batchwriter` 技能),`FlushFunc` 只调 repository `BatchInsert*`;管理端统计 API 只读 repository,不触达 DDL。 +4. **Apps**:在 `internal/apps//` 编排采集与入队;高频写入通过 `internal/infra/persistence/batchwriter` 各域独立实例异步 flush(详见 `clickhouse-batchwriter` 技能)。**日志/分析用途表**还要同时建 PG/SQLite 回落并接入 `logstore`(见 `logstore` 技能),`FlushFunc` 调 `logstore.Active` 而不是 `analyticsrepo`;普通业务分析表仍只读 repository。 ### ClickHouse 验证 diff --git a/.agents/skills/logstore/SKILL.md b/.agents/skills/logstore/SKILL.md new file mode 100644 index 00000000..64335519 --- /dev/null +++ b/.agents/skills/logstore/SKILL.md @@ -0,0 +1,92 @@ +--- +name: "logstore" +description: "Wavelet 项目专用:当新增或修改日志/分析用途表(访问日志、审计流水、可观测时序)、接入 internal/repository/logstore、切换日志主库、实现 PG/SQLite 回落,或判断一张表该走业务主库还是日志库时必须使用。" +--- + +# 日志用途表开发 + +开始前阅读根目录 `AGENTS.md`。DDL 用 `database-migration`;高频写入队列用 `clickhouse-batchwriter`;切换任务用 `new-async-task`。本技能只回答:**这张表是不是日志表,以及如何接入可切换的日志主库。** + +分层与切换协议见 [日志用途表](../../../docs/LOGSTORE.md)。 + +## 先判定 + +日志表同时满足: + +- 追加写入、几乎不更新单行 +- 按时间查询/聚合,允许按保留天数删除 +- 关闭 ClickHouse 后仍要能写、能查 +- 不参与用户/配置/任务等事务一致性 + +**不要**做成日志表:用户、配置、任务执行、上传元数据、需要事务或强一致的业务实体。这些走主库 `repository`,不要进 `logstore`。 + +当前框架已接入的日志表:`w_user_access_logs`(管理端 API 访问审计)。 + +## 分层 + +| 层级 | 路径 | 职责 | +| :--- | :--- | :--- | +| 抽象 | `internal/repository/logstore` | 接口 + `Active`/`BuildForMigration`;apps **只**面向这里 | +| CH 实现 | `logstore` 委托 `internal/repository/analytics` | 原生 `PrepareBatch` / `ChDB` 查询 | +| 主库实现 | `logstore` GORM | PG(按月分区)与 SQLite(普通表) | +| Model | `internal/model/analytics` | 实体、`TableName`、`InsertColumns`、`BatchInsertSQL`,无 IO | +| 入队 | `internal/apps/` + `batchwriter` | `FlushFunc` 调 `logstore.Active().….BatchInsert` | +| 切换 | `internal/apps/admin/logs` 的 `logs:db_switch` | 冻结写入 → 排空 → 复制 → 翻转 `log_database` | +| 清理 | `logstore.CleanupExpired`,由 `system:cleanup` 调用 | 按库读取保留天数后 `DeleteBefore` | + +`log_database` ∈ {`postgres`,`sqlite`,`clickhouse`},且只能是「随主库」或 ClickHouse:主库为 PG 时日志不能是 SQLite,反之亦然。`log_database` / `log_db_migration` 受保护,禁止管理端手动改。 + +## 新增一张日志表 + +按顺序做,列名三库必须一致。 + +1. **Model** + 在 `internal/model/analytics/` 定义 struct;实现 `TableName()`;批量写再提供 `InsertColumns()` / `BatchInsertSQL()`。 + +2. **三套 DDL**(`database-migration`) + - ClickHouse:`goose/clickhouse/`,`MergeTree`,`PARTITION BY toYYYYMM(时间列)`。 + - PostgreSQL:`goose/postgres/`,高频表用 `PARTITION BY RANGE (时间列)`,复合主键必须包含分区键。 + - SQLite:`goose/sqlite/`,普通表 + 时间/过滤列索引。 + 不要在 PG/SQLite 上复制 CH 物化视图;聚合在查询时实时算。 + +3. **logstore 接口** + 在对应 Store(现有 `UserAccessLogStore`,或新域自建接口并挂到 `Store`)补齐至少: + - 写入:`BatchInsert`(flush 目标;内调 `ensureWritable`) + - 查询:业务需要的 List/Count/聚合 + - 迁移:`ListForMigration(afterID, limit)`、`MigrationRange`、`DeleteAll`、`EnsurePartitions`(PG 按月预建,CH/SQLite no-op) + - 清理:`DeleteBefore(cutoff)` + +4. **双实现** + - CH:委托 `analyticsrepo`,零额外查询路径。 + - GORM:PG/SQLite 共用一套;方言 SQL 只放小函数(如按日 `to_char` / `strftime`)。零值 `id` 落库前用 `idgen.NextUint64ID()`。 + +5. **`buildStore`** + 在 `provider.go` 的 CH / GORM 分支同时挂上新域。 + +6. **写入** + apps 用独立 `batchwriter` 实例;`FlushFunc` → `logstore.Active(ctx)` → `BatchInsert`。禁止 `analyticsrepo.BatchInsert`、禁止 `db.ChConn`。迁移任务调用域的 `Drain`(等队列空一个 flush 周期,不要 `Stop` writer)。 + +7. **切换任务** + 在 `copy*` 流程增加该表:`DeleteAll` 目标 → `MigrationRange` + `EnsurePartitions` → 按 id 分页复制。不要改切换协议(仍冻结写入、源数据不删、成功才翻转)。 + +8. **清理** + `CleanupExpired` 对该表 `DeleteBefore`;保留天数用已有 `log_retention_days_*`,不要为单表再发明一套 key,除非产品明确要求独立 TTL。 + +## 禁止 + +- apps 直接 `import` `internal/repository/analytics` 或 `db.ChConn` / `db.ChDB` 做日志读写 +- 只建 CH 表、不建 PG/SQLite 回落 +- 在 Handler 里逐条 `PrepareBatch` + `Send` +- 把业务表「顺便」放进 logstore 以便关 CH +- 管理端 API 改 `log_database` / `log_db_migration` + +## 验证 + +```bash +go test ./internal/repository/logstore ./internal/repository/analytics +go test ./internal/apps/admin/logs ./internal/apps/risk_control ./internal/platform/bootstrap +make swagger # 若改了状态/查询 API +make code-check +``` + +对照:`w_user_access_logs` 的 model、三库 goose、`logstore` GORM/CH、`risk_control.InitLogWriter`、`logs.LogDBSwitchHandler`、`system:cleanup`。 diff --git a/AGENTS.md b/AGENTS.md index fd1060e5..79ee2d84 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -76,6 +76,7 @@ Strong success criteria let you loop independently. Weak criteria ("make it work | `new-async-task` | 添加或修改 Asynq 任务、定时任务、TaskHandler、任务元数据 | | `new-setting` | 添加或修改系统/业务/公开设置、`/admin/system` 参数或 `/admin/settings` 图形化设置 | | `database-migration` | 数据库表结构变更、goose SQL 迁移(PG/SQLite/ClickHouse)、seed 数据 | +| `logstore` | 日志/分析用途表、`internal/repository/logstore`、切换日志主库、PG/SQLite 回落 | | `clickhouse-batchwriter` | ClickHouse 批量写入、`internal/infra/persistence/batchwriter` 接入、分析表异步 flush、背压与写入路径改造 | | `file-upload` | 业务上传文件、Worker 程序化摄取、`upload.Ingest` 策略选型、文件访问与 `w_uploads` / 统计排查 | | `cache-framework` | 新增或修改业务缓存(RAM/Redis/DB 三层读路径)、缓存失效、多节点 pub/sub 同步、评估高频读是否应接入缓存 | @@ -97,11 +98,12 @@ Strong success criteria let you loop independently. Weak criteria ("make it work - **分层**:`apps → repository → model`,`repository → infra/persistence`;禁止 `model → repository`。 - `model`:实体、表名、配置 key、查询 DTO、无 IO 规则。禁止 `db.DB` / Redis / CH;禁止 `import repository`。GORM hook 仅可 mutate 自身字段,禁止在 hook 内再查 DB/缓存。 - `repository`:唯一持久化入口。apps/logics 禁止为业务 CRUD 直调 `db.DB`(管理端 SQL 控制台、infra 内部等例外保留)。禁止新增 `model.Get/List/Create/...` 类数据访问 API。 +- 日志/分析表(访问日志、审计流水、可观测时序)走 `internal/repository/logstore`,禁止 apps 直连 `repository/analytics` 或 `db.ChConn`/`db.ChDB`。判定与接入步骤见 `logstore` skill。 ## 技术栈与项目目录结构 ### 技术栈 -- **后端**:Go 1.25+、Gin、GORM、PostgreSQL、ClickHouse、Redis、Asynq、Cobra、Viper、Swaggo、OpenTelemetry、Zap、AWS SDK v2。 +- **后端**:Go 1.25+、Gin、GORM、PostgreSQL、可选 ClickHouse、Redis、Asynq、Cobra、Viper、Swaggo、OpenTelemetry、Zap、AWS SDK v2。 - **前端**:Next.js (App Router)、TypeScript、Tailwind CSS、pnpm、shadcn/ui。 ## 后端开发规范 diff --git a/docs/DEPLOYMENT.md b/docs/DEPLOYMENT.md index 6d1fd483..ff2ee328 100644 --- a/docs/DEPLOYMENT.md +++ b/docs/DEPLOYMENT.md @@ -16,7 +16,7 @@ | **前端服务 (Node.js)** | `pnpm start` | 提供 React/Next.js 页面服务(在分离部署时使用) | 分离模式必选 | | **PostgreSQL** | 关系型主数据库 | 存储用户、系统配置、认证源、任务执行记录等核心数据 | **必选** | | **Redis** | 缓存与消息队列中间件 | 存储 Session 会话、临时缓存以及 Asynq 异步任务队列数据 | **必选** | -| **ClickHouse** | 分析型数据库 | 存储历史数据同步或进行高性能分析 | 可选 | +| **ClickHouse** | 分析型数据库 | 可选的日志主库;关闭时访问审计由 PostgreSQL/SQLite 承接 | 可选 | | **对象存储 (S3)** | 兼容 S3 的云存储/私有云 | 存放用户上传的静态文件、图片等 | 可选 | --- @@ -306,7 +306,7 @@ s3: - **Scheduler 独占**:**【注意】** 为避免重复触发定时 Cron 任务,`wavelet scheduler` 定时调度器进程**同一时间应仅运行单个活跃实例**(主备高可用可以通过容器平台的单实例保障或 K8s Job 机制来限制实例数为 1)。 #### 5. ClickHouse 高并发同步 -在大数据量、高频支付结算场景下,开启 ClickHouse 以接收系统的历史数据同步,通过定时器把 PostgreSQL 的压力转移到 ClickHouse 列式存储中。 +访问审计等日志表默认写在当前业务主库。数据量大、需要列式扫描时,开启 ClickHouse,再在任务管理运行「切换日志数据库」迁到 ClickHouse(迁移期间冻结写入,源数据不删)。开发约定见 [日志用途表](./LOGSTORE.md)。 ```yaml clickhouse: enabled: true diff --git a/docs/LOGSTORE.md b/docs/LOGSTORE.md new file mode 100644 index 00000000..012a7333 --- /dev/null +++ b/docs/LOGSTORE.md @@ -0,0 +1,45 @@ +# 日志用途表 + +Wavelet 的访问审计等日志表不绑死 ClickHouse。`internal/repository/logstore` 按 `log_database` 在 PostgreSQL / SQLite / ClickHouse 之间切换;关闭 ClickHouse 时由当前业务主库承接写入、查询与清理。 + +逐步落地步骤见 `.agents/skills/logstore/SKILL.md`。本文只约定判定、分层与切换协议。 + +## 什么算日志表 + +同时满足才进 logstore: + +- 追加写入,几乎不更新单行 +- 按时间查询或聚合,允许按保留天数删除 +- 关闭 ClickHouse 后仍要能写、能查 +- 不参与用户 / 配置 / 任务等事务一致性 + +用户、系统配置、任务执行、上传元数据走业务主库 `repository`,不要塞进 logstore。 + +当前已接入:`w_user_access_logs`(管理端 API 访问审计),接口 `UserAccessLogStore`。 + +## 分层 + +| 层级 | 路径 | 职责 | +| :--- | :--- | :--- | +| 抽象 | `internal/repository/logstore` | 接口 + `Active` / `BuildForMigration`;apps 只面向这里 | +| CH 实现 | `logstore` 委托 `repository/analytics` | 原生批量与现有查询 | +| 主库实现 | `logstore` GORM | PostgreSQL 按月分区;SQLite 普通表 | +| 入队 | `risk_control` + `batchwriter` | `FlushFunc` → `logstore.Active` | +| 切换 | `logs:db_switch` | 冻结 → 排空 → 复制 → 翻转 | +| 清理 | `logstore.CleanupExpired` | `system:cleanup` 按库读 `log_retention_days_*` 后 `DeleteBefore` | + +`log_database` 只能是「随业务主库」或 `clickhouse`。`log_database` / `log_db_migration` 受保护,管理端不可改。 + +## 切换协议 + +1. 校验 `target` 合法且不等于当前库。 +2. 写 `log_db_migration=migrating`,`Drain` 在途队列(不要 `Stop` writer);写入返回明确错误,不排队。 +3. 清空目标表后按 id 分页复制;PostgreSQL 目标先 `EnsurePartitions`。 +4. 全部成功才翻转 `log_database`;失败清标记,写入继续走源库。 +5. 源数据不删。 + +不要另起切换协议,也不要在任务或 Handler 里直连 `analyticsrepo` / `db.ChConn`。 + +## 新增一张日志表 + +必须同时提供 ClickHouse / PostgreSQL / SQLite 三套 goose,列名一致。接口至少包含 `BatchInsert`、业务查询、`ListForMigration` / `MigrationRange` / `DeleteAll` / `EnsurePartitions`、`DeleteBefore`。`FlushFunc` 调 `logstore.Active`。细节与禁止项见 `logstore` skill。