diff --git a/hzh/Gen2D/00-index.md b/hzh/Gen2D/00-index.md new file mode 100644 index 0000000..d349dbb --- /dev/null +++ b/hzh/Gen2D/00-index.md @@ -0,0 +1,62 @@ +# Gen2D 架构讲解 — 答辩文档索引 + +> AI 驱动的 2D 游戏素材生成工具 + +## 文档导航 + +| # | 文件 | 主题 | 核心要点 | +|---|------|------|----------| +| 1 | [01-system-overview](01-system-overview.md) | 系统总览 | Gin → Handler → Service → 基础设施 → 外部依赖 | +| 2 | [02-worker-pool](02-worker-pool.md) | 协程池 | 有界并发、per-user 限流、背压、优雅关闭 | +| 3 | [03-task-queue](03-task-queue.md) | 任务队列 | 可插拔接口 + Memory / RabbitMQ 双实现 | +| 4 | [04-rabbitmq](04-rabbitmq.md) | RabbitMQ 集成 | AMQP 连接、持久化消息、重试/死信流程 | +| 5 | [05-generation-pipeline](05-generation-pipeline.md) | 生成管线 | Eino 4 阶段图 + 质量回退 + 降级路径 | +| 6 | [06-sprite-processing](06-sprite-processing.md) | 精灵图处理 | 背景移除 → 投影切割 → 后处理 → GIF 预览 | +| 7 | [07-observability](07-observability.md) | 可观测性 | 35 Prometheus 指标 + 3 Grafana 仪表盘 + 10 告警 | +| 8 | [08-sse-push](08-sse-push.md) | SSE 实时推送 | EventBus → SSEHandler → 浏览器 EventSource | +| 9 | [09-rate-limiting](09-rate-limiting.md) | 限流 | Redis Lua 令牌桶 + 双层限流 + Fail-Open | +| 10 | [10-middleware-chain](10-middleware-chain.md) | 中间件链 | Logger → Recovery → Metrics → Auth → RateLimit | +| 11 | [11-consumer-producer](11-consumer-producer.md) | Consumer-Producer 桥接 | TaskQueue → Consumer → WorkerPool 解耦 | +| 12 | [12-three-tier-fallback](12-three-tier-fallback.md) | 三级降级策略 | Queue → Pool → Legacy 降级链 | +| 13 | [13-prompt-engineering](13-prompt-engineering.md) | 标签驱动提示词 | 40+ 标签映射 + 风格一致性 | +| 14 | [14-deployment](14-deployment.md) | 部署架构 | Docker Compose + Prometheus + Grafana | +| 15 | [15-config-cascade](15-config-cascade.md) | 配置级联 | Viper 三层配置:YAML → ENV → Default | + +## 架构总览图 + +``` +┌─────────────────────────────────────────────────────────────────┐ +│ 浏览器 (React) │ +│ EventSource ← SSE ← EventBus ← Pipeline Progress │ +└────────────────────────────┬────────────────────────────────────┘ + │ HTTP +┌────────────────────────────▼────────────────────────────────────┐ +│ Gin HTTP Server │ +│ Logger → Recovery → Metrics → Auth → RateLimit → Handler │ +└────────────────────────────┬────────────────────────────────────┘ + │ +┌────────────────────────────▼────────────────────────────────────┐ +│ Handler 层 │ +│ GenerateHandler / SSEHandler / PromptHandler │ +└──────┬─────────────────────┬─────────────────────┬──────────────┘ + │ │ │ +┌──────▼──────┐ ┌─────────▼─────────┐ ┌─────▼──────┐ +│ TaskQueue │ │ WorkerPool │ │ EventBus │ +│ Memory/RMQ │───→│ NumCPU*4 并发 │ │ Pub/Sub │ +└─────────────┘ │ Per-user 限流 │ └────────────┘ + └─────────┬─────────┘ + │ + ┌─────────▼─────────┐ + │ Eino Pipeline │ + │ 4 阶段生成管线 │ + └─────────┬─────────┘ + │ + ┌──────────────┼──────────────┐ + ▼ ▼ ▼ + Image API Quality Check SplitSprite + (DALL-E 3) (GPT-4o) + GIF Maker +``` + +## 配套图表 + +每份文档开头使用 Mermaid 流程图,可直接在 Markdown 渲染器中查看。 diff --git a/hzh/Gen2D/01-system-overview.md b/hzh/Gen2D/01-system-overview.md new file mode 100644 index 0000000..9bfaffb --- /dev/null +++ b/hzh/Gen2D/01-system-overview.md @@ -0,0 +1,178 @@ +# 01 - 系统总览 + +> **一句话概括**:Gen2D 采用经典分层架构,Gin HTTP Server -> Handler -> Service -> 基础设施 -> 外部依赖,各层职责清晰、可独立替换。 + +## 架构全景 + +```mermaid +graph TB + subgraph Browser["浏览器 (React)"] + UI["三栏工作台"] + end + + subgraph Gin["Gin HTTP Server"] + MW["中间件链
Logger / Recovery / Metrics / Auth / RateLimit"] + end + + subgraph Handler["Handler 层"] + GH["Generate"] + SH["SSE"] + PH["Prompt"] + OH["其他 Handler"] + end + + subgraph Infra["基础设施层"] + WP["WorkerPool
协程池"] + TQ["TaskQueue
Memory / RabbitMQ"] + EB["EventBus
Pub/Sub"] + RL["RateLimiter
Redis 令牌桶"] + end + + subgraph Pipeline["Service 层 (Eino Pipeline)"] + PO["PromptOptimizer"] + AG["AssetGenerator"] + QS["QualitySupervisor"] + FA["FormatAdapter"] + end + + subgraph External["外部依赖"] + DB["MySQL"] + CDN["七牛云 Kodo"] + REDIS["Redis"] + LLM["LLM / Image API"] + end + + UI -->|"HTTP / SSE"| MW + MW --> Handler + GH --> TQ + GH --> WP + SH --> EB + TQ -->|"Consumer"| WP + WP --> Pipeline + PO --> AG --> QS --> FA + AG --> LLM + FA --> CDN + Handler --> DB + RL --> REDIS + Pipeline -->|"Progress"| EB +``` + +## 分层详解 + +| 层级 | 目录 | 核心职责 | 代表组件 | +|------|------|----------|----------| +| **Entry / DI** | `cmd/main.go` | 启动入口、依赖注入、信号处理 | `main()` | +| **HTTP 层** | `internal/handler/` | 请求绑定、参数校验、响应封装 | `Generate`, `SSEHandler`, `PromptOptimize` | +| **中间件** | `internal/mildware/` | 日志、恢复、指标采集、认证、限流 | `Logger`, `Recovery`, `Metrics`, `AuthMiddleware` | +| **业务层** | `internal/service/` | 生成管线、提示词优化、质检、推理调用 | `RunPipeline`, `CheckQuality`, `GenerateImages` | +| **领域模型** | `internal/model/` | 数据库实体、响应体定义 | `User`, `Project`, `Task`, `Asset` | +| **基础设施** | `internal/pkg/` | 协程池、任务队列、事件总线、限流器 | `workerpool.Pool`, `taskqueue.TaskQueue`, `eventbus.Broker` | +| **精灵处理** | `pkg/splitsprite/`, `pkg/gifmaker/` | 精灵表切割、GIF 预览生成 | `splitsprite.Process`, `gifmaker.Encode` | +| **配置** | `internal/config/` | YAML 加载、环境变量绑定、默认值 | `config.Load()` | + +## 依赖注入模式 + +Gen2D 采用轻量级的 **Set\*/Init\* 函数注入** 模式,避免引入 DI 框架。 + +``` +main.go 中的注入链路: + +config.Load() → 加载全局配置 +handler.InitAuthService() → 注入 JWT 配置 +service.InitLLMConfig() → 注入 LLM 配置 +service.InitImageGenConfig() → 注入文生图配置 +handler.InitStorageService() → 注入存储服务 +handler.SetWorkerPool() → 注入协程池 +handler.SetTaskQueue() → 注入任务队列 +eventbus.Init() → 初始化事件总线 +``` + +**设计要点**: + +- Handler 层不直接 import `workerpool`、`taskqueue` 等基础设施包,仅通过注入的接口交互 +- `service.InitImageGenConfig()` 将配置缓存为包级变量,避免在函数签名中传递大量参数 +- 每个 `Set*` 函数对应一个包级全局变量,简单但足够清晰 + +> :bulb: **为什么不用 Wire / Fx?** 项目规模可控,`cmd/main.go` 约 220 行即可完成全部注入,框架级 DI 的复杂度收益比不高。 + +## 配置级联 + +Gen2D 使用 Viper 实现三层配置覆盖,优先级从高到低: + +```mermaid +graph LR + ENV["环境变量
GEN2D_*"] -->|"最高优先级"| VIPER["Viper"] + YAML["YAML 配置文件
config.yaml"] -->|"中等优先级"| VIPER + DEFAULT["代码默认值
SetDefault()"] -->|"最低优先级"| VIPER + VIPER --> CONFIG["Config 结构体"] +``` + +**覆盖规则**:`环境变量 > YAML 文件 > 默认值` + +| 配置源 | 示例 | 说明 | +|--------|------|------| +| 环境变量 | `GEN2D_PORT=9090` | 部署时覆盖,适合容器化场景 | +| YAML 文件 | `server.port: 8080` | 开发时配置,集中管理 | +| 默认值 | `v.SetDefault("server.port", 8080)` | 代码内置,零配置即可启动 | + +配置结构体涵盖 10 个子模块: + +| 配置块 | 对应环境变量前缀 | 关键字段 | +|--------|------------------|----------| +| `server` | `GEN2D_*` | port, mode, max_file_size | +| `database` | `GEN2D_DSN` | dsn | +| `jwt` | `GEN2D_JWT_*` | secret, expire | +| `log` | `GEN2D_LOG_*` | level, format | +| `redis` | `GEN2D_REDIS_*` | addr, password, db | +| `llm` | `GEN2D_LLM_*` | base_url, api_key, model | +| `image_gen` | `GEN2D_IMAGE_*` | base_url, api_key, model, timeout, max_retries | +| `qiniu` | `GEN2D_QINIU_*` | access_key, secret_key, bucket, cdn_host | +| `workerpool` | `GEN2D_WORKERPOOL_*` | workers, queue_size, max_per_user | +| `taskqueue` | `GEN2D_TASKQUEUE_*` | driver, memory.buffer_size, rabbitmq.* | + +## 请求全链路 + +一个素材生成请求的完整生命周期: + +``` +1. 浏览器 → POST /api/v1/generate +2. Gin 中间件链 → Logger → Recovery → Metrics → Auth → RateLimit +3. Handler.Generate() + ├── 参数绑定 & 校验 + ├── 保存任务到 MySQL(status=pending) + ├── 构建 TaskMessage + └── 提交到 TaskQueue(优先)或 WorkerPool(fallback) +4. 返回 { taskId } 给浏览器 +5. Consumer 从 TaskQueue 消费消息 +6. Consumer → WorkerPool.Submit() +7. WorkerPool 的 Worker 执行任务: + ├── 注入 ProgressReporter 到 context + ├── RunPipeline() — Eino Graph 执行 + │ ├── PromptOptimizer(合并风格 + 技术参数) + │ ├── AssetGenerator(调用 AI 出图 API) + │ ├── QualitySupervisor(质检,最多重试 3 次) + │ └── FormatAdapter(精灵表切割 + GIF 预览) + ├── 上传素材到七牛云 + └── 更新 MySQL 任务状态 +8. EventBus.Publish() → SSE 推送给浏览器 +9. 浏览器通过 EventSource 实时接收进度 +``` + +## 关键设计决策 + +| 决策 | 选择 | 理由 | +|------|------|------| +| 分层架构 | Handler-Service-Infra 三层 | 解耦各层职责,便于独立测试和替换 | +| 依赖注入 | Set\*/Init\* 函数 | 轻量、零依赖,项目规模可控 | +| 配置管理 | Viper 三层级联 | 容器化友好,零配置可启动 | +| 异步任务 | 提交-队列-消费-执行 | API 快速返回,长任务不阻塞请求 | +| 事件推送 | EventBus + SSE | 比 WebSocket 轻量,HTTP 原生支持 | +| 限流 | Redis 令牌桶 | 分布式一致,Fail-Open 保可用性 | + +## 关联文档 + +- [索引](00-index.md) — 文档导航与架构总览图 +- [协程池](02-worker-pool.md) — 有界并发与 per-user 限流 +- [任务队列](03-task-queue.md) — 可插拔队列接口与双实现 +- [RabbitMQ 集成](04-rabbitmq.md) — 持久化消息与重试机制 +- [生成管线](05-generation-pipeline.md) — Eino 4 阶段管线与质量回退 diff --git a/hzh/Gen2D/02-worker-pool.md b/hzh/Gen2D/02-worker-pool.md new file mode 100644 index 0000000..9d6af57 --- /dev/null +++ b/hzh/Gen2D/02-worker-pool.md @@ -0,0 +1,249 @@ +# 02 - 协程池 (Worker Pool) + +> **一句话概括**:有界并发协程池,`NumCPU*4` workers,per-user 限流,背压保护,优雅关闭。 + +## 工作流 + +```mermaid +graph TB + subgraph Submit["Submit() 入口"] + CHECK_CLOSED{"池已关闭?"} + CHECK_USER{"per-user 限流
active >= max?"} + CHECK_FULL{"channel 满?"} + end + + subgraph Channel["有界任务队列"] + TASK_CHAN["chan Task
capacity = queueSize"] + end + + subgraph Workers["Worker 协程"] + W1["Worker 0"] + W2["Worker 1"] + W3["Worker ..."] + WN["Worker N-1"] + end + + subgraph Execute["任务执行"] + TIMEOUT["context.WithTimeout
10 min"] + FN["task.Fn(ctx)"] + RELEASE["释放用户槽位"] + end + + CHECK_CLOSED -->|"ErrPoolClosed"| REJECT["返回错误"] + CHECK_CLOSED -->|否| CHECK_USER + CHECK_USER -->|"ErrUserLimitReached"| REJECT + CHECK_USER -->|通过| CHECK_FULL + CHECK_FULL -->|"ErrPoolFull (HTTP 503)"| REJECT + CHECK_FULL -->|通过| TASK_CHAN + TASK_CHAN --> W1 & W2 & W3 & WN + W1 & W2 & W3 & WN --> TIMEOUT --> FN --> RELEASE +``` + +## Pool 结构体 + +```go +type Pool struct { + workers int // worker 数量(默认 NumCPU*4) + queueSize int // 任务队列容量(默认 100) + maxPerUser int // 单用户最大并发数(默认 2) + taskQueue chan Task // 有界任务队列 + wg sync.WaitGroup + ctx context.Context + cancel context.CancelFunc + metrics *Metrics + + mu sync.Mutex + userActive map[string]int // userID -> 当前并发数 + + closed atomic.Bool + closeCh chan struct{} +} +``` + +**核心字段说明**: + +| 字段 | 类型 | 默认值 | 说明 | +|------|------|--------|------| +| `workers` | `int` | `runtime.NumCPU() * 4` | IO 密集型调参 | +| `queueSize` | `int` | `100` | 有界缓冲,触发背压 | +| `maxPerUser` | `int` | `2` | 防止单用户占满池 | +| `taskQueue` | `chan Task` | buffered channel | 有界队列,固定容量 | +| `userActive` | `map[string]int` | — | 记录每用户活跃任务数 | + +## Functional Options 模式 + +协程池采用 Functional Options 模式进行配置,开箱即用、可选覆盖: + +```go +pool := workerpool.New( + workerpool.WithWorkers(16), // 覆盖默认的 NumCPU*4 + workerpool.WithQueueSize(200), // 覆盖默认的 100 + workerpool.WithMaxPerUser(5), // 覆盖默认的 2 +) +``` + +| Option | 默认值 | 约束 | 说明 | +|--------|--------|------|------| +| `WithWorkers(n)` | `NumCPU * 4` | `n < 1` 时强制为 1 | IO 密集型场景推荐 4 倍 CPU 核数 | +| `WithQueueSize(n)` | `100` | `n < 1` 时强制为 1 | 有界缓冲,满时触发背压 | +| `WithMaxPerUser(n)` | `2` | `n < 1` 时置 0(不限制) | 防止单用户占满池 | + +> :bulb: **为什么选择 Functional Options?** 可选参数天然为零值时保持默认,新增配置项无需修改 `New()` 签名,调用方按需指定。 + +## 有界并发 + +协程池的核心是 `make(chan Task, queueSize)` 创建的 **有界缓冲 channel**。 + +**工作原理**: + +1. `Submit()` 尝试将任务写入 channel +2. `N` 个 Worker 协程从 channel 中取任务执行 +3. channel 满时进入背压逻辑(见下文) + +```go +// Submit 核心逻辑 +select { +case p.taskQueue <- task: // 成功入队 + return nil +case <-p.ctx.Done(): // 池已关闭 + return ErrPoolClosed +default: // 队列满,背压 + return ErrPoolFull +} +``` + +## Per-user 限流 + +每个用户同时执行的任务数受到 `maxPerUser` 限制,防止单用户占满整个池。 + +**实现机制**: + +``` +Submit() 调用流程: +1. 加锁检查 userActive[userID] +2. 若 active >= maxPerUser → 返回 ErrUserLimitReached +3. 否则 userActive[userID]++(预留槽位) +4. 任务执行完毕后 defer releaseUserSlot(userID) +``` + +**槽位生命周期**: + +| 阶段 | 操作 | 说明 | +|------|------|------| +| Submit | `userActive[userID]++` | 预留槽位,不等到 Worker 取出 | +| Execute | — | Worker 从 channel 取出后开始执行 | +| Complete | `defer userActive[userID]--` | 任务完成或失败时释放 | + +> :warning: **槽位预留时机**:在 `Submit()` 而非 `worker()` 中预留,确保 channel 满时不会出现"槽位已分配但任务未入队"的不一致状态。 + +## 背压保护 + +当任务队列已满时,协程池通过 `select + default` 实现非阻塞拒绝: + +``` +队列满(channel 已达 capacity) + │ + ▼ +select 进入 default 分支 + │ + ├── 回滚用户槽位(如果有) + ├── 原子递增 RejectedTasks + └── 返回 ErrPoolFull + │ + ▼ + Handler 层映射为 HTTP 503 Service Unavailable + 响应体:"系统繁忙,请稍后重试" +``` + +**背压 vs 阻塞**: + +| 策略 | 行为 | 适用场景 | +|------|------|----------| +| 非阻塞拒绝(Gen2D) | 立即返回错误 | 用户交互型 API,快速失败 | +| 阻塞等待 | 阻塞直到有空位 | 批处理系统,不能丢任务 | + +Gen2D 选择非阻塞拒绝,因为用户期望快速得到反馈,而非无限等待。 + +## 优雅关闭 + +协程池支持优雅关闭,确保正在执行的任务有时间完成: + +``` +收到 SIGINT / SIGTERM + │ + ▼ +consumeCancel() ← 停止消费者,不再接收新任务 + │ + ▼ +pool.Shutdown(ctx) ← 30 秒超时 + │ + ├── closed.Swap(true) ← 停止接收新任务 + ├── cancel() ← 通知 Worker 停止取任务 + ├── wg.Wait() ← 等待所有 Worker 退出 + │ + ├── 成功 → 日志 "workerpool shutdown gracefully" + └── 超时 → 日志 "workerpool shutdown timeout" +``` + +**Shutdown 返回值**: + +| 返回值 | 含义 | +|--------|------| +| `true` | 所有任务正常完成 | +| `false` | 超时,部分任务可能丢失 | + +## Metrics 指标 + +协程池内置原子计数器,支持运行时观测: + +| 指标 | 类型 | 说明 | +|------|------|------| +| `ActiveWorkers` | `atomic.Int32` | 当前正在执行任务的 Worker 数 | +| `QueuedTasks` | `atomic.Int32` | 当前排队等待的任务数 | +| `CompletedTasks` | `atomic.Int64` | 已完成任务总数 | +| `FailedTasks` | `atomic.Int64` | 失败任务总数 | +| `SubmittedTasks` | `atomic.Int64` | 提交任务总数 | +| `RejectedTasks` | `atomic.Int64` | 被拒绝任务总数(队列满或用户限流) | + +所有指标通过 `Metrics()` 方法返回只读快照,同时上报 Prometheus。 + +## Task 结构体 + +```go +type Task struct { + ID string // 任务唯一标识 + UserID string // 所属用户 ID + Fn func(ctx context.Context) error // 任务执行函数 + Priority int // 优先级(预留) +} +``` + +每个任务携带超时 context(默认 10 分钟),从池的根 context 派生,确保优雅关闭时能取消正在执行的任务。 + +## 与 TaskQueue 的协作 + +``` +TaskQueue(全局排队) + │ + ▼ +Consumer(消费消息) + │ + ▼ +WorkerPool.Submit()(单机并发控制) + │ + ▼ +Worker 执行 Pipeline +``` + +| 组件 | 职责 | 范围 | +|------|------|------| +| TaskQueue | 持久化排队,解耦提交与执行 | 全局(可跨机器) | +| WorkerPool | 单机并发控制,per-user 限流 | 单机 | +| Consumer | 桥接 TaskQueue 与 WorkerPool | 单机 | + +## 关联文档 + +- [索引](00-index.md) — 文档导航与架构总览图 +- [系统总览](01-system-overview.md) — 分层架构与依赖注入 +- [任务队列](03-task-queue.md) — 可插拔队列接口 +- [Consumer-Producer 桥接](00-index.md) — TaskQueue 到 WorkerPool 的解耦 diff --git a/hzh/Gen2D/03-task-queue.md b/hzh/Gen2D/03-task-queue.md new file mode 100644 index 0000000..2a706be --- /dev/null +++ b/hzh/Gen2D/03-task-queue.md @@ -0,0 +1,247 @@ +# 03 - 任务队列 (Task Queue) + +> **一句话概括**:可插拔任务队列接口,支持 Memory 和 RabbitMQ 双实现,通过工厂模式一行切换。 + +## 架构设计 + +```mermaid +graph TB + subgraph Producer["生产者"] + HANDLER["Handler.Generate()"] + end + + subgraph Interface["TaskQueue 接口"] + SUBMIT["Submit(ctx, msg)"] + CONSUME["Consume(ctx, handler)"] + CLOSE["Close()"] + end + + subgraph Memory["MemoryQueue"] + MEM_CH["chan TaskMessage
buffered"] + MEM_SUBMIT["阻塞写入"] + MEM_CONSUME["持续消费"] + end + + subgraph RabbitMQ["RabbitMQQueue"] + RMQ_PUB["channel.Publish
Persistent"] + RMQ_CONSUME["channel.Consume
Manual ACK"] + end + + HANDLER --> SUBMIT + SUBMIT --> Memory + SUBMIT --> RabbitMQ + CONSUME --> Memory + CONSUME --> RabbitMQ + Memory --- MEM_CH + RabbitMQ --- RMQ_PUB +``` + +## TaskQueue 接口 + +```go +type TaskQueue interface { + Submit(ctx context.Context, msg TaskMessage) error + Consume(ctx context.Context, handler func(TaskMessage) error) error + Close() error +} +``` + +| 方法 | 语义 | 错误处理 | +|------|------|----------| +| `Submit` | 提交任务到队列 | 内存队列 buffer 满时阻塞;RabbitMQ 发送失败时返回错误 | +| `Consume` | 持续消费任务,直到 ctx 取消 | handler 返回错误时,内存队列丢弃,RabbitMQ NACK 重试 | +| `Close` | 关闭连接,释放资源 | 返回 `errors.Join` 聚合错误 | + +## 工厂模式 + +通过配置驱动,一行切换队列实现: + +```go +tq, err := taskqueue.New(cfg.TaskQueue) +``` + +```go +func New(cfg config.TaskQueueConfig) (TaskQueue, error) { + switch cfg.Driver { + case "rabbitmq": + return NewRabbitMQQueue(cfg.RabbitMQ) + default: + return NewMemoryQueue(cfg.Memory.BufferSize), nil + } +} +``` + +| `cfg.Driver` | 实现 | 适用场景 | +|--------------|------|----------| +| `"memory"` (默认) | `MemoryQueue` | 单机开发、演示环境 | +| `"rabbitmq"` | `RabbitMQQueue` | 多机生产部署 | + +## TaskMessage 消息结构 + +```go +type TaskMessage struct { + TaskID string `json:"task_id"` + UserID string `json:"user_id"` + ProjectID string `json:"project_id"` + Input service.PipelineInput `json:"input"` + CreatedAt time.Time `json:"created_at"` + RetryCount int `json:"retry_count"` +} +``` + +**关键设计**: + +- `Input` 直接嵌入 `PipelineInput`,Consumer 反序列化后即可传入管线 +- `RetryCount` 供 RabbitMQ 实现判断是否超过最大重试次数 +- JSON 序列化,兼容内存队列和 RabbitMQ 两种传输 + +## MemoryQueue 实现 + +基于 Go channel 的内存队列,零外部依赖。 + +**核心逻辑**: + +```go +type MemoryQueue struct { + ch chan TaskMessage // 有界缓冲 + closed atomic.Bool + closeCh chan struct{} +} +``` + +### Submit + +```go +func (q *MemoryQueue) Submit(ctx context.Context, msg TaskMessage) error { + select { + case q.ch <- msg: // 成功入队 + return nil + case <-ctx.Done(): // 上下文取消 + return ctx.Err() + case <-q.closeCh: // 队列已关闭 + return ErrQueueClosed + } +} +``` + +> :bulb: **阻塞语义**:MemoryQueue 的 Submit 是阻塞的——当 buffer 满时,调用方会阻塞直到有空位或 ctx 取消。这与 WorkerPool 的非阻塞拒绝形成对比。 + +### Consume + +```go +func (q *MemoryQueue) Consume(ctx context.Context, handler func(TaskMessage) error) error { + for { + select { + case msg := <-q.ch: // 取出消息 + if err := handler(msg); err != nil { + // handler 失败,记录日志并丢弃(不重试) + slog.Error("handler failed, dropping message", ...) + } + case <-ctx.Done(): // 上下文取消 + return ctx.Err() + case <-q.closeCh: // 队列已关闭 + return ErrQueueClosed + } + } +} +``` + +**MemoryQueue 特性**: + +| 特性 | 行为 | +|------|------| +| 持久化 | 无(进程重启后任务丢失) | +| 重试 | 无(handler 错误直接丢弃) | +| 背压 | 阻塞直到有空位 | +| 依赖 | 零外部依赖 | + +> :warning: **可接受的任务丢失**:AI 生成任务可以重新提交,进程重启丢失排队中的任务是可接受的权衡。 + +## RabbitMQQueue 实现 + +基于 AMQP 的持久化消息队列,支持手动 ACK 和重试。详见 [04-rabbitmq](04-rabbitmq.md)。 + +**RabbitMQQueue 特性**: + +| 特性 | 行为 | +|------|------| +| 持久化 | 消息 `DeliveryMode=Persistent`,队列 `durable=true` | +| 重试 | 失败 + retry < max → NACK+requeue;超过 → ACK 丢弃 | +| 背压 | 发布失败时返回错误(不阻塞) | +| 依赖 | 需要 RabbitMQ 服务 | + +## 双实现对比 + +| 维度 | MemoryQueue | RabbitMQQueue | +|------|-------------|---------------| +| **依赖** | 无 | RabbitMQ 服务 | +| **持久化** | 无 | 有(磁盘持久化) | +| **重试** | 无 | NACK + requeue | +| **背压** | 阻塞等待 | 返回错误 | +| **部署** | 单机 | 多机 | +| **适用** | 开发/演示 | 生产环境 | +| **消息丢失** | 进程重启丢失 | 服务重启不丢失 | + +## Consumer 桥接 + +Consumer 从 TaskQueue 消费消息,提交到 WorkerPool 执行,实现队列与并发控制的解耦: + +```go +func (c *Consumer) Start(ctx context.Context) error { + return c.queue.Consume(ctx, func(msg TaskMessage) error { + return c.pool.Submit(workerpool.Task{ + ID: msg.TaskID, + UserID: msg.UserID, + Fn: func(taskCtx context.Context) error { + c.handler(taskCtx, msg) // 执行 Pipeline + return nil + }, + }) + }) +} +``` + +``` +TaskQueue → Consumer → WorkerPool → Pipeline + 全局排队 桥接 单机并发 业务逻辑 +``` + +## 三级降级策略 + +Gen2D 在 `cmd/main.go` 中实现了三级降级链: + +| 优先级 | 组件 | 条件 | 行为 | +|--------|------|------|------| +| 1 | TaskQueue + Consumer | 初始化成功 | 队列排队 → Consumer 消费 → WorkerPool 执行 | +| 2 | WorkerPool(直接) | TaskQueue 初始化失败 | 直接提交到协程池 | +| 3 | Legacy FIFO Queue | 均不可用 | 串行队列兜底 | + +```go +if taskQueue != nil { + // 优先:通过队列提交 + taskQueue.Submit(ctx, msg) +} else if workerPool != nil { + // 回退:直接提交到协程池 + workerPool.Submit(task) +} else { + // 兜底:旧串行队列 + generateQueue.Enqueue(job) +} +``` + +## Metrics 指标 + +| Prometheus 指标 | 类型 | Label | 说明 | +|-----------------|------|-------|------| +| `gen2d_queue_depth` | Gauge | `driver` | 当前队列积压深度 | +| `gen2d_queue_submitted_total` | Counter | `driver` | 入队总量 | +| `gen2d_queue_consumed_total` | Counter | `driver` | 出队总量 | +| `gen2d_queue_submit_duration_seconds` | Histogram | `driver` | 入队耗时 | +| `gen2d_queue_errors_total` | Counter | `driver`, `error_type` | 队列错误总量 | + +## 关联文档 + +- [索引](00-index.md) — 文档导航与架构总览图 +- [系统总览](01-system-overview.md) — 分层架构与依赖注入 +- [协程池](02-worker-pool.md) — 有界并发与 per-user 限流 +- [RabbitMQ 集成](04-rabbitmq.md) — 持久化消息与重试机制 diff --git a/hzh/Gen2D/04-rabbitmq.md b/hzh/Gen2D/04-rabbitmq.md new file mode 100644 index 0000000..cd621d6 --- /dev/null +++ b/hzh/Gen2D/04-rabbitmq.md @@ -0,0 +1,209 @@ +# 04 - RabbitMQ 集成 + +> **一句话概括**:基于 AMQP 的持久化消息队列,支持手动 ACK、失败重试和死信丢弃。 + +## 消息流 + +```mermaid +graph TB + subgraph Producer["生产者"] + HANDLER["Handler.Generate()"] + end + + subgraph RabbitMQ["RabbitMQ"] + EXCHANGE["Default Exchange
(Direct)"] + QUEUE["gen2d:tasks
durable=true"] + end + + subgraph Consumer["消费者"] + CONSUME["channel.Consume
autoAck=false"] + end + + subgraph Decision["ACK/NACK 决策树"] + SUCCESS{"handler 成功?"} + RETRY{"retry < maxRetry?"} + ACK_OK["ACK
确认消费"] + NACK["NACK + requeue
重新入队"] + ACK_DISCARD["ACK (discard)
丢弃死信"] + end + + HANDLER -->|"Publish
Persistent"| EXCHANGE + EXCHANGE --> QUEUE + QUEUE --> CONSUME + CONSUME --> SUCCESS + SUCCESS -->|是| ACK_OK + SUCCESS -->|否| RETRY + RETRY -->|是| NACK + RETRY -->|否| ACK_DISCARD + NACK -.->|"重新投递"| QUEUE +``` + +## 连接流程 + +RabbitMQQueue 在初始化时完成连接、声明队列、设置 QoS: + +``` +amqp.Dial(cfg.URL) + │ + ▼ +conn.Channel() + │ + ▼ +ch.QueueDeclare(name, durable=true, autoDelete=false, exclusive=false) + │ + ▼ +ch.Qos(prefetch=1, prefetchSize=0, global=false) + │ + ▼ +RabbitMQQueue{conn, channel, queue, maxRetry, prefetch} +``` + +**参数说明**: + +| 参数 | 默认值 | 说明 | +|------|--------|------| +| `URL` | `amqp://guest:guest@localhost:5672/` | AMQP 连接地址 | +| `Queue` | `gen2d:tasks` | 队列名称 | +| `Prefetch` | `1` | 每次预取消息数,1 保证公平调度 | +| `MaxRetry` | `3` | 失败最大重试次数 | + +> :bulb: **Prefetch=1 的含义**:每个 Consumer 同时只处理 1 条消息,处理完(ACK)后才接收下一条。这避免了消息堆积在 Consumer 端,配合协程池的并发控制实现精确的任务调度。 + +## 消息发布 (Submit) + +```go +func (q *RabbitMQQueue) Submit(ctx context.Context, msg TaskMessage) error { + body, _ := json.Marshal(msg) + return q.channel.PublishWithContext(ctx, + "", // exchange(默认直连) + q.queue, // routing key = queue name + false, // mandatory + false, // immediate + amqp.Publishing{ + ContentType: "application/json", + DeliveryMode: amqp.Persistent, // 持久化消息 + Body: body, + Timestamp: time.Now(), + Headers: amqp.Table{ + "x-retry-count": msg.RetryCount, + }, + }, + ) +} +``` + +**关键配置**: + +| 属性 | 值 | 说明 | +|------|-----|------| +| `DeliveryMode` | `Persistent (2)` | 消息写磁盘,RabbitMQ 重启不丢失 | +| `ContentType` | `application/json` | JSON 序列化 | +| `x-retry-count` | `int` (header) | 当前重试次数,供消费端判断 | + +## 消息消费 (Consume) + +```go +func (q *RabbitMQQueue) Consume(ctx context.Context, handler func(TaskMessage) error) error { + deliveries, _ := q.channel.Consume( + q.queue, // queue + "", // consumer name(自动生成) + false, // autoAck = false(手动 ACK) + false, // exclusive + false, // noLocal + false, // noWait + nil, // args + ) + // 持续消费循环... +} +``` + +**手动 ACK 模式**: + +手动 ACK 给予消费者完全的控制权——只有当消息被成功处理后才确认,否则可以选择重试或丢弃。 + +## ACK/NACK 决策树 + +```mermaid +graph TB + MSG["收到消息"] --> PARSE{"JSON 解析成功?"} + PARSE -->|否| ACK_DISCARD1["ACK (discard)
格式错误无法恢复"] + PARSE -->|是| HANDLER{"handler(msg)
执行成功?"} + HANDLER -->|是| ACK_OK["ACK
确认消费"] + HANDLER -->|否| RETRY_CHECK{"msg.RetryCount
< maxRetry?"} + RETRY_CHECK -->|是| NACK["NACK(requeue=true)
重新入队等待重试"] + RETRY_CHECK -->|否| ACK_DISCARD2["ACK (discard)
超过最大重试,记录死信日志"] +``` + +| 场景 | 操作 | 说明 | +|------|------|------| +| handler 成功 | `d.Ack(false)` | 确认消费,消息从队列移除 | +| handler 失败 + retry < max | `d.Nack(false, true)` | 拒绝并重新入队,retry count 递增 | +| handler 失败 + retry >= max | `d.Ack(false)` + 日志 | 超过最大重试,丢弃(可扩展为死信队列) | +| JSON 解析失败 | `d.Ack(false)` | 格式错误无法恢复,直接丢弃 | + +**重试计数传递**: + +```go +// 发布时:写入 header +Headers: amqp.Table{ + "x-retry-count": msg.RetryCount, +} + +// 消费时:从 header 读取 +if retry, ok := d.Headers["x-retry-count"].(int32); ok { + msg.RetryCount = int(retry) +} +``` + +> :warning: **NACK requeue 的行为**:`Nack(false, true)` 会将消息重新放回队列头部。如果消费者立即再次消费,可能导致"毒消息"反复重试。Gen2D 通过 `maxRetry=3` 限制重试次数,并在超过后 ACK 丢弃来规避此问题。 + +## Metrics 指标 + +| Prometheus 指标 | 类型 | Label | 说明 | +|-----------------|------|-------|------| +| `gen2d_rabbitmq_connection_status` | Gauge | — | 连接状态(1=connected, 0=disconnected) | +| `gen2d_queue_submitted_total` | Counter | `driver=rabbitmq` | 发布消息总量 | +| `gen2d_queue_consumed_total` | Counter | `driver=rabbitmq` | 成功消费总量 | +| `gen2d_queue_errors_total` | Counter | `driver=rabbitmq`, `error_type` | 错误总量 | +| `gen2d_queue_submit_duration_seconds` | Histogram | `driver=rabbitmq` | 发布耗时 | + +## 关闭流程 + +```go +func (q *RabbitMQQueue) Close() error { + metrics.RabbitMQConnectionStatus.Set(0) // 标记断开 + var errs []error + errs = append(errs, q.channel.Close()) // 先关 channel + errs = append(errs, q.conn.Close()) // 再关连接 + return errors.Join(errs...) +} +``` + +**关闭顺序**:Channel 先于 Connection 关闭,确保所有未确认的消息被释放回队列。 + +## 配置参考 + +```yaml +# config.yaml +taskqueue: + driver: "rabbitmq" + rabbitmq: + url: "amqp://guest:guest@localhost:5672/" + queue: "gen2d:tasks" + prefetch: 1 + max_retry: 3 +``` + +| 配置项 | 环境变量 | 默认值 | 说明 | +|--------|----------|--------|------| +| `url` | `GEN2D_TASKQUEUE_RABBITMQ_URL` | `amqp://guest:guest@localhost:5672/` | AMQP 连接地址 | +| `queue` | `GEN2D_TASKQUEUE_RABBITMQ_QUEUE` | `gen2d:tasks` | 队列名称 | +| `prefetch` | `GEN2D_TASKQUEUE_RABBITMQ_PREFETCH` | `1` | 预取消息数 | +| `max_retry` | `GEN2D_TASKQUEUE_RABBITMQ_MAX_RETRY` | `3` | 最大重试次数 | + +## 关联文档 + +- [索引](00-index.md) — 文档导航与架构总览图 +- [任务队列](03-task-queue.md) — 可插拔接口与 MemoryQueue 实现 +- [协程池](02-worker-pool.md) — 单机并发控制 +- [系统总览](01-system-overview.md) — 分层架构与配置级联 diff --git a/hzh/Gen2D/05-generation-pipeline.md b/hzh/Gen2D/05-generation-pipeline.md new file mode 100644 index 0000000..b4e1de9 --- /dev/null +++ b/hzh/Gen2D/05-generation-pipeline.md @@ -0,0 +1,324 @@ +# 05 - 生成管线 (Generation Pipeline) + +> **一句话概括**:基于 CloudWeGo Eino 的 4 阶段生成管线,带质量回退和降级,确保任务不阻塞。 + +## 管线拓扑 + +```mermaid +graph TB + START(["START"]) --> PO["PromptOptimizer
提示词优化"] + PO --> AG["AssetGenerator
素材生成"] + AG --> QS["QualitySupervisor
质量检查"] + + QS -->|pass| FA["FormatAdapter
格式适配"] + QS -->|fail + retry < 3| PO + QS -->|fail + retry >= 3| FA + + FA --> END(["END"]) + + style QS fill:#fff3cd,stroke:#ffc107 + style PO fill:#d1ecf1,stroke:#17a2b8 + style AG fill:#d4edda,stroke:#28a745 + style FA fill:#d4edda,stroke:#28a745 +``` + +**四阶段详解**: + +| 阶段 | 节点名 | 职责 | 耗时特征 | +|------|--------|------|----------| +| 1 | `PromptOptimizer` | 合并风格、注入重试原因、追加技术参数 | < 1ms(纯 CPU) | +| 2 | `AssetGenerator` | 调用 AI 推理 API 出图 | 10-120s(IO 密集) | +| 3 | `QualitySupervisor` | 质检,决定路由目标 | 1-5s(可含 LLM 调用) | +| 4 | `FormatAdapter` | 精灵表切割 + GIF 预览 | 1-3s(CPU 密集) | + +## PipelineState 全局状态 + +Eino `compose.Graph` 通过 `WithGenLocalState` 注入全局状态,各节点通过 `StatePreHandler` / `StatePostHandler` 读写: + +```go +type PipelineState struct { + Input PipelineInput // 原始输入(首次运行时保存) + FinalPrompt string // PromptOptimizer 输出 + RawImages []GeneratedImage // AssetGenerator 输出 + PassQuality bool // QualitySupervisor 质检结果 + RejectReason string // 质检不通过原因 + RetryCount int // 重试次数(最多 3 次) + NextNode string // QualitySupervisor 设置的路由目标 +} +``` + +**状态流转**: + +``` +RetryCount=0, Input=原始输入 + │ + ▼ PromptOptimizer +FinalPrompt = "风格约束:...。用户描述:...。技术参数:..." + │ + ▼ AssetGenerator +RawImages = [GeneratedImage, ...] + │ + ▼ QualitySupervisor + ├── pass=true → NextNode="format_adapter" + ├── pass=false, RetryCount<3 → RetryCount++, NextNode="prompt_optimizer" + └── pass=false, RetryCount>=3 → NextNode="format_adapter"(降级) +``` + +## 节点详解 + +### 1. PromptOptimizer — 提示词优化 + +**职责**:合并多源提示词,输出最终提示词(纯字符串拼接,无 LLM 调用)。 + +``` +输入:PipelineInput +输出:string(最终提示词) + +拼接顺序: +1. 工程全局风格提示词 (GlobalStylePrompt) +2. 用户原始提示词 (Prompt) +3. 风格描述 (ProjectStyle + TaskStyle → "风格约束:色调:暖色;线条:粗线条") +4. 重试拒绝原因 (RejectReason → "注意修正以下问题:风格不一致") +5. 技术参数段 ("素材类型: sprite;分辨率: 64;输出格式: spritesheet") +``` + +**StatePreHandler**: + +| 场景 | 行为 | +|------|------| +| 首次运行 (`RetryCount=0`) | 保存原始输入到 `state.Input` | +| 重试运行 | 注入 `RejectReason` 到输入 | + +**StatePostHandler**:将最终提示词写入 `state.FinalPrompt`。 + +### 2. AssetGenerator — 素材生成 + +**职责**:调用 AI 推理 API 生成图片,支持三级图片来源回退。 + +``` +输入:string(提示词) +输出:[]GeneratedImage(原始图片列表) + +图片来源优先级: +1. 任务级参考图 (ReferenceImageData) → 图生图 (/images/edits) +2. 工程级参考图 (ProjectReferenceImage URL) → 下载后图生图 +3. 无参考图 → 纯文生图 (/images/generations) +``` + +**三级回退逻辑**: + +```go +if len(refData) > 0 { + // 1. 任务级参考图:直接使用 base64 解码后的数据 + return GenerateImagesFromRef(ctx, prompt, refData, params) +} +if projectRefURL != "" { + // 2. 工程级参考图:下载后使用 + refData, err = downloadImage(ctx, projectRefURL) + if err == nil { + return GenerateImagesFromRef(ctx, prompt, refData, params) + } + // 下载失败,回退到纯文生图 +} +// 3. 纯文生图 +return GenerateImages(ctx, prompt, params) +``` + +**外部 API 弹性**: + +| HTTP 状态 | 行为 | 说明 | +|-----------|------|------| +| `200 OK` | 解析响应 | 成功 | +| `5xx` | 自动重试(默认 2 次,间隔 5s) | 服务端临时故障 | +| `4xx` | 立即失败,不重试 | 客户端错误(参数错误等) | +| 超时 | 10 分钟(由 WorkerPool context 控制) | 长任务保护 | + +**Mock 降级**: + +当 `APIKey` 为空时,生成彩色占位图(纯色 PNG),便于开发和演示: + +```go +if imgCfg.APIKey == "" { + l.Warn("image API key not configured, using mock") + return generateMockImages(size, count) +} +``` + +**StatePostHandler**:将原始图片写入 `state.RawImages`。 + +### 3. QualitySupervisor — 质量检查 + +**职责**:检查生成图片的质量和风格一致性,决定路由目标。 + +``` +输入:[]GeneratedImage +输出:PipelineInput(用于路由分支读取) + +决策逻辑: +1. 调用 CheckQuality(ctx, images, style) +2. pass=true → NextNode = "format_adapter" +3. pass=false + RetryCount < 3 → RetryCount++, NextNode = "prompt_optimizer" +4. pass=false + RetryCount >= 3 → NextNode = "format_adapter"(降级) +``` + +**路由分支**(Eino `AddBranch`): + +```go +g.AddBranch(nodeQualitySupervisor, compose.NewGraphBranch( + func(ctx context.Context, _ PipelineInput) (string, error) { + var next string + compose.ProcessState[*PipelineState](ctx, func(_ context.Context, state *PipelineState) error { + next = state.NextNode + return nil + }) + return next, nil + }, + map[string]bool{nodePromptOptimizer: true, nodeFormatAdapter: true}, +)) +``` + +**质量检查器**: + +| 检查器 | 行为 | 用途 | +|--------|------|------| +| `defaultCheckQuality` | 始终返回 `true` | 生产环境(当前默认) | +| `NewCountedQualityChecker(n)` | 第 n 次调用后通过 | 测试重试逻辑 | +| `AlwaysFailQualityChecker` | 始终返回 `false` | 测试降级路径 | + +> :bulb: **可扩展性**:`QualityChecker` 是一个可替换的函数变量,未来可接入 LLM 视觉模型进行真正的质量评估。 + +### 4. FormatAdapter — 格式适配 + +**职责**:将原始图片转换为游戏引擎友好的格式。 + +``` +输入:PipelineInput +输出:PipelineOutput + +处理模式: +├── 精灵表模式 (Format="spritesheet", 单张图) +│ ├── PNG 解码 +│ ├── splitsprite.Process() 切割为独立帧 +│ ├── gifmaker.Encode() 生成 GIF 预览 +│ └── 输出:原始精灵表 + 帧列表 + GIF 预览 +│ +└── 普通模式(多图或非精灵表) + └── 原样透传 +``` + +**精灵表处理流程**: + +``` +单张精灵表 PNG + │ + ▼ splitsprite.Process() +独立帧列表 [frame_0, frame_1, ..., frame_N] + │ + ├── 保留原始精灵表 + ├── 编码每帧为独立 PNG + └── gifmaker.Encode() → 预览 GIF + │ + ▼ PipelineOutput{ + Assets: [spritesheet.png, frame_000.png, ..., preview.gif] + Metadata: {FrameWidth, FrameHeight, FrameCount, Directions, GIFURL} + } +``` + +**降级策略**: + +| 失败点 | 降级行为 | +|--------|----------| +| `splitsprite.Process()` 失败 | 整张图作为单帧输出 | +| `gifmaker.Encode()` 失败 | 跳过 GIF 预览,不阻塞管线 | +| 单帧编码失败 | 返回错误(无法降级) | + +## 进度上报 + +管线通过 context 注入 `ProgressReporter` 回调,实时上报进度: + +```go +type ProgressReporter func(stage string, progress int) +``` + +**进度点**: + +| 阶段 | 进度值 | 触发时机 | +|------|--------|----------| +| `prompt_builder` | 10 | 首次进入 PromptOptimizer | +| `prompt_builder` | 30 + retry*10 | 重试进入 PromptOptimizer | +| `asset_generator` | 35 | PromptOptimizer 完成 | +| `quality_supervisor` | 60 | AssetGenerator 完成 | +| `quality_supervisor` | 50 + retry*10 | 质检不通过(重试) | +| `format_adapter` | 85 | 质检通过或降级 | + +进度通过 `EventBus.Publish()` 推送给 SSE 客户端。 + +## 阶段耗时监控 + +`WithStageTimer` 注入阶段计时器到 context,`reportProgress` 在每次进度上报时计算上一阶段的耗时并上报 Prometheus: + +```go +metrics.PipelineStageDuration.WithLabelValues(prevStage).Observe(elapsed.Seconds()) +``` + +**监控指标**: + +| Prometheus 指标 | 类型 | Label | 说明 | +|-----------------|------|-------|------| +| `gen2d_pipeline_total` | Counter | `status` | 管线执行总量 (success/fail/timeout) | +| `gen2d_pipeline_duration_seconds` | Histogram | `status` | 端到端耗时 | +| `gen2d_pipeline_stage_duration_seconds` | Histogram | `stage` | 各阶段耗时 | +| `gen2d_pipeline_retries_total` | Counter | `stage` | 重试次数 | +| `gen2d_pipeline_tasks_active` | Gauge | — | 当前活跃管线数 | + +## Eino Graph 编译与执行 + +```go +func RunPipeline(ctx context.Context, in PipelineInput) (*PipelineOutput, error) { + // 1. 创建 Graph + g, err := NewGenerateGraph() + + // 2. 编译(maxRunSteps=20 防止无限循环) + r, err := g.Compile(ctx, compose.WithMaxRunSteps(20)) + + // 3. 执行 + output, err := r.Invoke(WithStageTimer(ctx), in) + + return &output, nil +} +``` + +**安全机制**: + +| 机制 | 说明 | +|------|------| +| `WithMaxRunSteps(20)` | 最大执行步数,防止无限循环(3 次重试 * 4 节点 + 余量) | +| `context.WithTimeout(10min)` | WorkerPool 层面的任务超时 | +| `context.Canceled` | 优雅关闭时取消正在执行的管线 | + +## PipelineInput 完整结构 + +```go +type PipelineInput struct { + ProjectID string // 工程 ID + TaskID string // 任务 ID + Prompt string // 用户原始文本 + AssetType string // sprite / background / ui / animation + ProjectStyle map[string]string // 工程风格键值对 + TaskStyle map[string]string // 任务风格覆盖 + Params AssetParams // 技术参数 + RejectReason string // 重试时由 state 注入 + Tags []string // 用户选择的标签 + UserNote string // 用户额外描述 + ReferenceImageData []byte // 任务级参考图(base64 解码后) + ProjectReferenceImage string // 工程级参考图 CDN URL + GlobalStylePrompt string // 工程全局风格提示词 +} +``` + +## 关联文档 + +- [索引](00-index.md) — 文档导航与架构总览图 +- [系统总览](01-system-overview.md) — 分层架构与 Service 层定位 +- [协程池](02-worker-pool.md) — 管线执行的并发控制 +- [任务队列](03-task-queue.md) — 管线任务的排队机制 diff --git a/hzh/Gen2D/06-sprite-processing.md b/hzh/Gen2D/06-sprite-processing.md new file mode 100644 index 0000000..b79176f --- /dev/null +++ b/hzh/Gen2D/06-sprite-processing.md @@ -0,0 +1,214 @@ +# 06 — 精灵图处理管线 + +> **一句话概括**:自动精灵图处理管线 — 背景移除 → 投影检测 → 切割 → 对齐 → GIF 预览,将一张 AI 生成的 Sprite Sheet 无缝转化为可用的逐帧动画资源。 + +--- + +```mermaid +flowchart LR + A["🖼️ Input PNG
Sprite Sheet"] --> B["🪄 Background
Removal"] + B --> C["📊 Gap
Detection"] + C --> D["✂️ Tile
Extract"] + D --> E["🔍 Filter
MinFill"] + E --> F["📐 Trim
Alpha"] + F --> G["🎯 Align
padToLargest"] + G --> H["🎬 GIF
Preview"] + + style A fill:#e3f2fd,stroke:#1976d2 + style B fill:#fff3e0,stroke:#f57c00 + style C fill:#e8f5e9,stroke:#388e3c + style D fill:#fce4ec,stroke:#c62828 + style E fill:#f3e5f5,stroke:#7b1fa2 + style F fill:#e0f7fa,stroke:#00838f + style G fill:#fff8e1,stroke:#f9a825 + style H fill:#e8eaf6,stroke:#303f9f +``` + +--- + +## 📋 处理管线总览 + +AI 生成的 Sprite Sheet 通常包含白色/绿色背景、不均匀间距和尺寸差异。Gen2D 的精灵图处理管线自动完成从"原始 PNG"到"可用动画帧"的全部转换工作,无需用户手动操作。 + +**管线核心入口**:`splitsprite.Process(img, opts)` + +| 步骤 | 函数 | 作用 | +|:---:|------|------| +| 1 | `removeWhiteBg` / `removeGreenScreen` | 移除背景色,生成透明通道 | +| 2 | `projectionSplit` / `fixedGridSplit` | 检测帧边界,确定切割位置 | +| 3 | `MinFillRatio` 过滤 | 丢弃空白/低填充率的无效帧 | +| 4 | `trimAlpha` | 去除每帧的透明边框,减少冗余像素 | +| 5 | `padToLargest` | 统一画布尺寸,底部居中对齐 | + +--- + +## 🪄 步骤 1:背景移除 + +AI 生成的图片通常带有纯色背景,管线支持两种模式: + +### 白底模式(WhiteBg) + +```go +// 阈值距离 #FFFFFF + alpha 渐变 +dist := max(255-R, max(255-G, 255-B)) +if dist < threshold/2 → 全透明 +if dist < threshold → alpha 线性渐变(抗锯齿) +``` + +- **WhiteThreshold** 默认 40,距离纯白 40 以内的像素被处理 +- 采用渐变 alpha 实现平滑过渡,避免硬边缘 + +### 绿幕模式(GreenScreen) + +```go +// 绿色主导检测 +gDominance := G - (R+B)/2 +if gDominance > tolerance*255 → 透明化 +``` + +- **GreenTolerance** 默认 0.2,控制绿色检测灵敏度 +- 适用于绿色背景的 AI 生成图 + +> 💡 **设计选择**:两种模式互斥,`Process()` 根据 `opts.WhiteBg` / `opts.GreenScreen` 自动选择。 + +--- + +## 📊 步骤 2:切割策略 + +管线提供两种切割方式,根据配置自动切换: + +### 投影检测(Projection Detection) + +这是 Gen2D 的核心创新点。传统全局投影在检测列边界时,会因为武器、尾巴、翅膀等突出物被"稀释"而导致切割不准。 + +**两阶段投影算法**: + +1. **全局行投影**:统计每行非透明像素占比,检测水平间隙 + - 行间隙通常是干净且全宽的,全局投影效果好 + +2. **逐行列投影**:在每个行段内独立计算列密度 + - 防止突出物(武器/尾巴)被其他行稀释 + - 每行段获得独立的列边界,互不干扰 + +```mermaid +flowchart TB + A["输入图像"] --> B["全局行投影
检测行间隙"] + B --> C["行段 1"] + B --> D["行段 2"] + B --> E["行段 N"] + C --> F["逐行列投影"] + D --> G["逐行列投影"] + E --> H["逐行列投影"] + F --> I["tiles[]"] + G --> I + H --> I + + style B fill:#e8f5e9,stroke:#388e3c + style F fill:#fff3e0,stroke:#f57c00 + style G fill:#fff3e0,stroke:#f57c00 + style H fill:#fff3e0,stroke:#f57c00 +``` + +**峰值检测算法**(`findCuts`): +- 对密度曲线做滑动平均平滑(kernel = minGap) +- 以均值为参考检测峰值段(start: mean*1.05, continue: mean*0.95) +- 在相邻峰之间的谷底确定切割位置 +- 回退机制:峰值不足时降级到阈值法(`findCutsByGap`) + +### 固定网格(Fixed Grid) + +当 `GridRows > 0 && GridCols > 0` 时启用,将图像等分为 Rows×Cols 个单元格。 + +- **GridPadding**:单元格间间距(默认 2px) +- 适用于已知行列数的标准 Sprite Sheet + +--- + +## 🔍 步骤 3:过滤 — MinFillRatio + +切割后的每个 tile 都计算填充率: + +``` +fillRatio = 非透明像素数 / 总像素数 +``` + +- **MinFillRatio** 默认 0.14(14%) +- 低于阈值的 tile 被判定为空白帧,自动丢弃 +- 有效过滤因间隙检测误差产生的空白切片 + +--- + +## 📐 步骤 4:裁剪 — trimAlpha + +对每个 tile 执行透明边框裁剪: + +1. 扫描四边界,找到非透明像素的最小包围矩形 +2. 裁切到该矩形,去除四周透明区域 +3. 减少冗余像素,为后续对齐做准备 + +--- + +## 🎯 步骤 5:对齐 — padToLargest + +动画播放时,如果每帧尺寸不同且内容未对齐,会导致角色"抖动"。 + +**底部居中锚定**(Bottom-Center Anchor): + +``` +canvasW = maxW * 110% // 最大帧宽度 + 10% padding +canvasH = maxH * 110% // 最大帧高度 + 10% padding + +每帧偏移: + ox = (canvasW - frameW) / 2 // 水平居中 + oy = canvasH - frameH // 底部对齐(脚踏同一水平线) +``` + +- 所有帧共享统一画布尺寸 +- 水平居中保证角色在同一屏幕位置 +- 底部对齐保证角色"脚踏实地",防止上下漂移 + +--- + +## 🎬 GIF Maker + +`gifmaker.Encode()` 将处理后的帧序列编码为动画 GIF: + +| 特性 | 实现 | +|------|------| +| 统一画布 | 所有帧归一化到 maxW × maxH | +| 透明色 | 调色板索引 0 = 完全透明 | +| 防鬼影 | `DisposalBackground` 每帧清除前一帧 | +| 调色板 | 采样像素构建(每 3px 取样),最多 256 色 | +| 循环 | `LoopCount = 0`(无限循环) | + +```go +anim.Disposal = append(anim.Disposal, gif.DisposalBackground) +anim.BackgroundIndex = 0 // 透明色 +``` + +> ⚠️ **DisposalBackground 的重要性**:如果不设置此选项,GIF 播放器会在前一帧基础上叠加新帧,产生"残影"效果。 + +--- + +## 📦 Options 配置速查 + +| 参数 | 类型 | 默认值 | 说明 | +|------|------|--------|------| +| `WhiteBg` | bool | true | 白底移除模式 | +| `WhiteThreshold` | uint8 | 40 | 白色阈值(0-255) | +| `GreenScreen` | bool | false | 绿幕移除模式 | +| `GreenTolerance` | float64 | 0.2 | 绿色容差(0-1) | +| `GridRows` / `GridCols` | int | 0 | 固定网格行列数 | +| `GapThreshold` | float64 | 0.03 | 间隙判定阈值 | +| `MinGapWidth` | int | 2 | 最小间隙宽度(px) | +| `MinFillRatio` | float64 | 0.14 | 最小填充率 | +| `Trim` | bool | true | 透明边框裁剪 | +| `CenterAlign` | bool | true | 底部居中对齐 | + +--- + +## 🔗 关联文档 + +- [← 返回索引](00-index.md) +- [05 — 生成管线](05-generation-pipeline.md) — 管线中 SplitSprite 节点的调用方 +- [08 — SSE 实时推送](08-sse-push.md) — 处理进度的实时推送 diff --git a/hzh/Gen2D/07-observability.md b/hzh/Gen2D/07-observability.md new file mode 100644 index 0000000..b0bb690 --- /dev/null +++ b/hzh/Gen2D/07-observability.md @@ -0,0 +1,202 @@ +# 07 — 可观测性 + +> **一句话概括**:35 个 Prometheus 指标 + 3 个 Grafana 仪表盘 + 10 条告警规则,覆盖全栈,让系统运行状态一目了然。 + +--- + +```mermaid +flowchart LR + A["🌐 Request"] --> B["📊 Metrics
Middleware"] + B --> C["⚙️ Application
Logic"] + C --> D["📈 Prometheus
Scrape"] + D --> E["📉 Grafana
Dashboard"] + D --> F["🚨 AlertManager
Notify"] + + style A fill:#e3f2fd,stroke:#1976d2 + style B fill:#fff3e0,stroke:#f57c00 + style C fill:#e8f5e9,stroke:#388e3c + style D fill:#fce4ec,stroke:#c62828 + style E fill:#e8eaf6,stroke:#303f9f + style F fill:#ffebee,stroke:#b71c1c +``` + +--- + +## 📐 指标体系总览 + +Gen2D 遵循 Prometheus 命名最佳实践,所有指标使用 `gen2d_` 前缀,共 **35 个指标**,分为 **6 大组**: + +| 组 | 指标数 | 采集方式 | 侵入性 | +|:--:|:------:|---------|:------:| +| HTTP 层 | 5 | Gin 中间件自动采集 | 零侵入 | +| 限流层 | 2 | 限流中间件自动采集 | 零侵入 | +| 队列层 | 5 | 队列实现内部埋点 | 低 | +| 协程池层 | 6 | 池内部埋点 | 低 | +| Pipeline 层 | 5 | 管线节点回调 | 低 | +| 基础设施层 | 5 | 连接状态监控 | 低 | + +--- + +## 📊 五大指标组详解 + +### 1️⃣ HTTP 层指标 + +Gin 中间件自动采集,**零业务代码侵入**。 + +| 指标名 | 类型 | 标签 | 说明 | +|--------|------|------|------| +| `gen2d_http_requests_total` | Counter | method, path, status | 请求总量 | +| `gen2d_http_request_duration_seconds` | Histogram | method, path | 请求延迟分布 | +| `gen2d_http_request_size_bytes` | Histogram | method, path | 请求体大小 | +| `gen2d_http_response_size_bytes` | Histogram | method, path | 响应体大小 | +| `gen2d_http_requests_in_flight` | Gauge | — | 当前并发请求数 | + +> 💡 **FullPath() 的关键作用**:使用路由模板 `/api/v1/tasks/:taskId` 而非实际路径 `/api/v1/tasks/abc123`,避免高基数标签导致 Prometheus 内存爆炸。 + +### 2️⃣ 限流层指标 + +| 指标名 | 类型 | 标签 | 说明 | +|--------|------|------|------| +| `gen2d_ratelimit_requests_total` | Counter | scope, endpoint, result | 限流决策总量 | +| `gen2d_ratelimit_remaining_tokens` | Gauge | scope, endpoint | 剩余令牌数 | + +- `scope`:`user` / `global` +- `result`:`allowed` / `denied` + +### 3️⃣ 任务队列层指标 + +| 指标名 | 类型 | 标签 | 说明 | +|--------|------|------|------| +| `gen2d_queue_depth` | Gauge | driver | 当前队列积压深度 | +| `gen2d_queue_submitted_total` | Counter | driver | 入队总量 | +| `gen2d_queue_consumed_total` | Counter | driver | 出队总量 | +| `gen2d_queue_submit_duration_seconds` | Histogram | driver | 入队耗时 | +| `gen2d_queue_errors_total` | Counter | driver, error_type | 队列错误总量 | + +- `driver`:`memory` / `rabbitmq` + +### 4️⃣ 协程池层指标 + +| 指标名 | 类型 | 标签 | 说明 | +|--------|------|------|------| +| `gen2d_pool_active_workers` | Gauge | — | 活跃 worker 数 | +| `gen2d_pool_queued_tasks` | Gauge | — | 池内排队任务数 | +| `gen2d_pool_submitted_total` | Counter | — | 提交到池的任务总量 | +| `gen2d_pool_completed_total` | CounterVec | result | 完成的任务总量 | +| `gen2d_pool_rejected_total` | Counter | — | 被拒绝的任务 | +| `gen2d_pool_task_duration_seconds` | Histogram | — | 任务执行耗时 | + +### 5️⃣ Pipeline 业务层指标 + +| 指标名 | 类型 | 标签 | 说明 | +|--------|------|------|------| +| `gen2d_pipeline_total` | CounterVec | status | Pipeline 执行总量 | +| `gen2d_pipeline_duration_seconds` | HistogramVec | status | 端到端耗时 | +| `gen2d_pipeline_stage_duration_seconds` | HistogramVec | stage | 各阶段耗时 | +| `gen2d_pipeline_retries_total` | CounterVec | stage | 各阶段重试次数 | +| `gen2d_pipeline_tasks_active` | Gauge | — | 当前执行中的 Pipeline 数 | + +> 🔍 **stage_duration 定位瓶颈**:通过 `stage` 标签(如 `asset_generator`、`quality_check`)可以精确定位哪个阶段是性能瓶颈。 + +### 6️⃣ 基础设施层指标 + +| 指标名 | 类型 | 标签 | 说明 | +|--------|------|------|------| +| `gen2d_redis_operations_total` | CounterVec | op, result | Redis 操作总量 | +| `gen2d_redis_operation_duration_seconds` | HistogramVec | op | Redis 操作延迟 | +| `gen2d_redis_connection_pool_size` | Gauge | — | Redis 连接池大小 | +| `gen2d_rabbitmq_connection_status` | Gauge | — | RabbitMQ 连接状态 | +| `gen2d_rabbitmq_reconnect_total` | Counter | — | RabbitMQ 重连次数 | + +--- + +## 🔧 中间件集成 + +### Metrics 中间件工作流程 + +```go +func Metrics() gin.HandlerFunc { + return func(c *gin.Context) { + path := c.FullPath() // 路由模板,避免高基数 + HTTPRequestsInFlight.Inc() // 进入时 +1 + defer HTTPRequestsInFlight.Dec() // 离开时 -1 + + reqSize := c.Request.ContentLength // 请求体大小 before + c.Next() // 执行后续链 + elapsed := time.Since(start) // 耗时 after + HTTPRequestsTotal.WithLabelValues(...).Inc() + } +} +``` + +**关键设计**: +- `in_flight` 使用 `defer` 保证异常退出也能正确递减 +- 请求体大小在 `c.Next()` 之前采集(此时 Content-Length 已知) +- 响应体大小在 `c.Next()` 之后采集(此时 Writer 已写入) + +--- + +## 📉 Grafana 仪表盘 + +| 仪表盘 | 用途 | 关键面板 | +|--------|------|---------| +| **Overview** | 全局概览 | 请求量、错误率、延迟 P50/P95/P99、活跃连接 | +| **Pipeline** | 管线监控 | 成功率、各阶段耗时、重试率、活跃任务数 | +| **Infrastructure** | 基础设施 | Redis/RabbitMQ 状态、队列深度、池饱和度 | + +--- + +## 🚨 告警规则 + +共 **10 条告警规则**,覆盖限流、队列、协程池、管线和基础设施: + +| 告警名 | 级别 | 条件 | 说明 | +|--------|:----:|------|------| +| `RateLimitHighDenialRate` | ⚠️ | 限流拒绝率 > 5%(持续 5m) | 可能遭受攻击或配置过严 | +| `QueueBacklog` | ⚠️ | 队列深度 > 50(持续 2m) | 消费能力不足 | +| `QueueErrors` | 🔴 | 5 分钟内错误 > 5 次 | 队列服务异常 | +| `PoolSaturation` | ⚠️ | 活跃 worker 占比 > 90%(持续 5m) | 考虑扩容 | +| `PoolTaskRejected` | ⚠️ | 5 分钟内拒绝 > 5 个 | 池容量不足 | +| `PipelineSuccessRateLow` | 🔴 | 成功率 < 90%(持续 10m) | 生成服务异常 | +| `RedisDown` | 🔴 | 连接池大小 = 0(持续 1m) | Redis 不可用 | +| `RabbitMQDisconnected` | 🔴 | 连接状态 = 0(持续 1m) | RabbitMQ 断连 | +| `HighErrorRate` | 🔴 | 5xx 错误率 > 5%(持续 5m) | 服务异常 | +| `HighLatency` | ⚠️ | P95 延迟 > 5s(持续 5m) | 影响用户体验 | + +> 🛡️ **告警级别说明**:🔴 Critical 表示需要立即处理,⚠️ Warning 表示需要关注但不紧急。 + +--- + +## 📦 基础设施指标 + +除业务指标外,Gen2D 还监控外部依赖的健康状态: + +```mermaid +flowchart LR + subgraph "基础设施监控" + R["Redis"] -->|"operations_total"| M["Prometheus"] + R -->|"connection_pool_size"| M + Q["RabbitMQ"] -->|"connection_status"| M + Q -->|"reconnect_total"| M + end + M --> G["Grafana"] + M --> A["AlertManager"] + + style R fill:#fce4ec,stroke:#c62828 + style Q fill:#fff3e0,stroke:#f57c00 + style M fill:#e8f5e9,stroke:#388e3c + style G fill:#e8eaf6,stroke:#303f9f + style A fill:#ffebee,stroke:#b71c1c +``` + +- **Redis**:操作延迟、成功率、连接池大小 +- **RabbitMQ**:连接状态(1=connected, 0=disconnected)、重连次数 + +--- + +## 🔗 关联文档 + +- [← 返回索引](00-index.md) +- [10 — 中间件链](10-middleware-chain.md) — Metrics 中间件的挂载位置 +- [09 — 限流](09-rate-limiting.md) — 限流指标的采集方式 +- [14 — 部署架构](14-deployment.md) — Prometheus + Grafana 的部署配置 diff --git a/hzh/Gen2D/08-sse-push.md b/hzh/Gen2D/08-sse-push.md new file mode 100644 index 0000000..973d53f --- /dev/null +++ b/hzh/Gen2D/08-sse-push.md @@ -0,0 +1,219 @@ +# 08 — SSE 实时推送 + +> **一句话概括**:内存 EventBus 发布/订阅,SSE 推送管线进度到浏览器,让用户实时看到生成过程。 + +--- + +```mermaid +flowchart LR + P["⚙️ Pipeline
Callback"] -->|"Publish"| EB["📡 EventBus
Broker"] + EB -->|"Subscribe
taskID"| H1["🌐 SSE Handler
/tasks/:id/stream"] + EB -->|"SubscribeAll
global"| H2["🌐 SSE Handler
/projects/:id/stream"] + H1 -->|"text/event-stream"| B1["🖥️ Browser
EventSource"] + H2 -->|"text/event-stream"| B2["🖥️ Browser
EventSource"] + + style P fill:#e8f5e9,stroke:#388e3c + style EB fill:#fff3e0,stroke:#f57c00 + style H1 fill:#e3f2fd,stroke:#1976d2 + style H2 fill:#e3f2fd,stroke:#1976d2 + style B1 fill:#f3e5f5,stroke:#7b1fa2 + style B2 fill:#f3e5f5,stroke:#7b1fa2 +``` + +--- + +## 📡 EventBus 架构 + +EventBus 是 Gen2D 的内存事件总线,负责在 Pipeline 执行过程中发布进度事件,并由 SSE Handler 订阅推送给客户端。 + +### 核心设计 + +```go +type Broker struct { + mu sync.RWMutex + subs map[string][]chan TaskEvent // per-task 订阅 + allSubs []chan TaskEvent // 全局订阅 +} +``` + +| 组件 | 作用 | +|------|------| +| `subs` | 按 taskID 索引的订阅者列表 | +| `allSubs` | 全局订阅者(接收所有事件) | +| `sync.RWMutex` | 读写锁保护并发访问 | + +### TaskEvent 数据结构 + +```go +type TaskEvent struct { + TaskID string `json:"task_id"` + ProjectID string `json:"project_id,omitempty"` + Status string `json:"status"` // pending|running|saving|completed|failed + Stage string `json:"stage,omitempty"` // prompt_builder|asset_generator|... + Progress int `json:"progress"` // 0-100 + Error string `json:"error,omitempty"` +} +``` + +--- + +## 🔔 订阅模式 + +### Subscribe(taskID) — 任务级订阅 + +```go +func (b *Broker) Subscribe(taskID string) <-chan TaskEvent { + ch := make(chan TaskEvent, 16) // 有界缓冲,容量 16 + b.mu.Lock() + b.subs[taskID] = append(b.subs[taskID], ch) + b.mu.Unlock() + return ch +} +``` + +- 用于 `GET /api/v1/tasks/:taskId/stream` +- 仅接收指定任务的状态变更 +- Buffer 容量 **16**,足够应对正常进度更新频率 + +### SubscribeAll() — 全局订阅 + +```go +func (b *Broker) SubscribeAll() <-chan TaskEvent { + ch := make(chan TaskEvent, 64) // 有界缓冲,容量 64 + b.mu.Lock() + b.allSubs = append(b.allSubs, ch) + b.mu.Unlock() + return ch +} +``` + +- 用于 `GET /api/v1/projects/:projectId/stream` +- 接收所有任务的事件,在 Handler 层按 projectID 过滤 +- Buffer 容量 **64**,因为全局事件量更大 + +--- + +## 📤 Publish — 扇出分发 + +```go +func (b *Broker) Publish(taskID string, event TaskEvent) { + // 1. 发送到任务级订阅者 + for _, ch := range subs { + select { + case ch <- event: + default: // 满则丢弃,非阻塞 + slog.Warn("subscriber buffer full, dropping event") + } + } + // 2. 发送到全局订阅者 + for _, ch := range allSubs { + select { ... } + } +} +``` + +**关键特性**: +- **Fan-out**:同时发送到 task 级和 global 级订阅者 +- **Non-blocking send**:使用 `select default` 防止慢消费者阻塞发布方 +- **慢消费者丢弃**:缓冲满时静默丢弃,保证 Pipeline 不被 SSE 拖慢 + +--- + +## 🌐 SSE Handler + +### Stream — 任务级流 + +``` +GET /api/v1/tasks/:taskId/stream +``` + +```go +func (h *SSEHandler) Stream(c *gin.Context) { + // 1. 设置 SSE 响应头 + c.Header("Content-Type", "text/event-stream") + c.Header("Cache-Control", "no-cache") + c.Header("X-Accel-Buffering", "no") // 禁用 nginx 缓冲 + + // 2. 订阅事件 + ch := h.broker.Subscribe(taskID) + defer h.broker.Unsubscribe(taskID, ch) + + // 3. 事件循环 + for { + select { + case event := <-ch: + c.SSEvent("status", event) + c.Writer.Flush() + // 终态自动关闭 + if event.Status == "completed" || event.Status == "failed" { + return + } + case <-c.Request.Context().Done(): + return // 客户端断开 + } + } +} +``` + +**SSE 响应头**: + +| Header | 值 | 作用 | +|--------|---|------| +| `Content-Type` | `text/event-stream` | 标识 SSE 流 | +| `Cache-Control` | `no-cache` | 禁用缓存 | +| `X-Accel-Buffering` | `no` | 禁用 nginx 代理缓冲 | + +### StreamProject — 工程级流 + +``` +GET /api/v1/projects/:projectId/stream +``` + +- 使用 `SubscribeAll()` 订阅全局事件 +- 在 Handler 层按 `event.ProjectID != projectID` 过滤 +- 工程级流不会因单个任务完成而关闭,持续监听新任务 + +--- + +## 📊 数据流全景 + +```mermaid +sequenceDiagram + participant P as Pipeline + participant DB as Database + participant EB as EventBus + participant SSE as SSE Handler + participant B as Browser + + P->>DB: updateTaskInDB(status, progress) + P->>EB: Publish(taskID, event) + EB->>SSE: ch <- event + SSE->>B: data: {"status":"running","progress":45} + Note over B: EventSource.onmessage() + P->>DB: updateTaskInDB(completed) + P->>EB: Publish(taskID, terminal event) + EB->>SSE: ch <- event + SSE->>B: data: {"status":"completed","progress":100} + Note over SSE: 终态,关闭连接 +``` + +--- + +## 🛡️ 容错设计 + +| 场景 | 处理方式 | +|------|---------| +| 慢消费者 | Buffer 满时丢弃事件,Pipeline 不阻塞 | +| 客户端断开 | `c.Request.Context().Done()` 触发,自动 Unsubscribe | +| 终态到达 | completed/failed 后自动关闭 SSE 连接 | +| 无订阅者 | Publish 静默返回,不报错 | +| Broker 关闭 | Close() 关闭所有 channel,SSE 循环退出 | + +--- + +## 🔗 关联文档 + +- [← 返回索引](00-index.md) +- [05 — 生成管线](05-generation-pipeline.md) — Pipeline 中的进度回调 +- [07 — 可观测性](07-observability.md) — SSE 连接的监控 +- [10 — 中间件链](10-middleware-chain.md) — SSE 端点的中间件配置 diff --git a/hzh/Gen2D/09-rate-limiting.md b/hzh/Gen2D/09-rate-limiting.md new file mode 100644 index 0000000..68d80fc --- /dev/null +++ b/hzh/Gen2D/09-rate-limiting.md @@ -0,0 +1,207 @@ +# 09 — 限流 + +> **一句话概括**:Redis Lua 原子令牌桶 + 双层限流 + Fail-Open 降级,保护系统免受过载。 + +--- + +```mermaid +flowchart LR + A["🌐 Request"] --> B["🌍 Global
Limiter"] + B -->|pass| C["👤 User
Limiter"] + B -->|deny| F["❌ 429"] + B -->|redis-fail| C + C -->|pass| D["✅ Handler"] + C -->|deny| F + C -->|redis-fail| D + + style A fill:#e3f2fd,stroke:#1976d2 + style B fill:#fff3e0,stroke:#f57c00 + style C fill:#e8f5e9,stroke:#388e3c + style D fill:#c8e6c9,stroke:#2e7d32 + style F fill:#ffcdd2,stroke:#c62828 +``` + +--- + +## ⚙️ 令牌桶算法 + +Gen2D 使用 **Redis + Lua 脚本** 实现分布式令牌桶限流,保证原子性和一致性。 + +### Lua 脚本核心逻辑 + +```lua +-- KEYS[1] = 限流 key +-- ARGV[1] = rate(每秒令牌数) +-- ARGV[2] = burst(桶容量) +-- ARGV[3] = now(当前时间戳,毫秒) +-- ARGV[4] = expiration(key 过期时间) + +-- 1. 获取当前桶状态 +local data = redis.call('HMGET', key, 'tokens', 'ts') + +-- 2. 首次访问,初始化满桶 +if tokens == nil then + tokens = burst + last_ts = now +end + +-- 3. 计算时间差,补充令牌 +local delta = now - last_ts +if delta > 0 and rate > 0 then + local refill = (delta * rate) / 1000 + tokens = math.min(burst, tokens + refill) +end + +-- 4. 判断是否允许 +if tokens >= 1 then + tokens = tokens - 1 + allowed = 1 +else + retry_after = math.ceil((deficit * 1000) / rate) +end + +-- 5. 更新 Redis +redis.call('HSET', key, 'tokens', tokens, 'ts', last_ts) +redis.call('PEXPIRE', key, expiration) + +return {allowed, tokens, retry_after} +``` + +### 固定窗口变体 + +当 `Rate = 0` 时,令牌桶退化为固定窗口模式: + +- 桶初始化为满(Burst 个令牌) +- 用完后不补充(`rate = 0` 时跳过 refill) +- 等待 key 过期后重置(`Expiration` 控制窗口大小) + +> 💡 **适用场景**:24 小时维度的配额控制,如"每天 30 次提示词优化"。 + +--- + +## 🔀 双层限流配置 + +Gen2D 对核心接口实施**全局限流 + 用户限流**双重保护: + +| 接口 | 全局限流 | 用户限流 | 窗口 | +|------|---------|---------|------| +| `/api/v1/prompt/optimize` | 1000 次/24h | 30 次/24h | 25h 过期 | +| `/api/v1/generate` | 500 次/24h | 15 次/24h | 25h 过期 | + +### 执行顺序 + +```mermaid +flowchart TB + A["请求进入"] --> B["全局限流检查"] + B -->|通过| C["用户限流检查"] + B -->|拒绝| D["429 Too Many Requests"] + C -->|通过| E["执行 Handler"] + C -->|拒绝| D + + style B fill:#fff3e0,stroke:#f57c00 + style C fill:#e8f5e9,stroke:#388e3c + style D fill:#ffcdd2,stroke:#c62828 + style E fill:#c8e6c9,stroke:#2e7d32 +``` + +**全局限流在前**:先检查系统总体配额,避免单个用户耗尽全局配额。 + +### 配置参数 + +```go +type Config struct { + Rate int // 每秒令牌数;0 = 固定窗口 + Burst int // 桶容量(窗口内总量上限) + KeyPrefix string // Redis key 前缀 + Expiration time.Duration // key 过期时间 +} +``` + +| 配置项 | Prompt User | Prompt Global | Generate User | Generate Global | +|--------|:-----------:|:-------------:|:-------------:|:---------------:| +| Rate | 0 | 0 | 0 | 0 | +| Burst | 30 | 1000 | 15 | 500 | +| KeyPrefix | `ratelimit:prompt:user:` | `ratelimit:prompt:global:` | `ratelimit:generate:user:` | `ratelimit:generate:global:` | +| Expiration | 25h | 25h | 25h | 25h | + +--- + +## 🛡️ Fail-Open 降级 + +当 Redis 不可用时限流器自动降级为 **Fail-Open** 模式: + +```go +func (l *TokenBucketLimiter) Allow(ctx context.Context, key string) (bool, int, time.Duration) { + result, err := l.script.Run(ctx, l.client, ...).Int64Slice() + if err != nil { + // Redis 不可用时 fail-open,放行请求 + return true, l.config.Burst, 0 + } + // ... +} +``` + +**设计权衡**: + +| 策略 | 优点 | 缺点 | +|------|------|------| +| **Fail-Open** ✅ | 保证可用性,用户体验不受影响 | 可能短暂失去限流保护 | +| Fail-Close | 严格限流保护 | Redis 故障导致全站不可用 | + +> 🛡️ **选择 Fail-Open**:在"偶尔超限"和"完全不可用"之间,优先保证服务可用性。 + +--- + +## 📡 中间件响应 + +限流中间件返回标准化的 HTTP 响应: + +### 允许通过 + +```http +HTTP/1.1 200 OK +X-RateLimit-Remaining: 12 +``` + +### 被限流 + +```http +HTTP/1.1 429 Too Many Requests +Retry-After: 3600 +X-RateLimit-Remaining: 0 + +{ + "code": 429, + "message": "请求过于频繁,请稍后再试" +} +``` + +| Header | 说明 | +|--------|------| +| `X-RateLimit-Remaining` | 剩余令牌数 | +| `Retry-After` | 建议重试等待秒数 | + +--- + +## 📊 指标采集 + +限流中间件自动采集 Prometheus 指标: + +```go +metrics.RateLimitRequestsTotal.WithLabelValues(scope, endpoint, result).Inc() +metrics.RateLimitRemainingTokens.WithLabelValues(scope, endpoint).Set(float64(remaining)) +``` + +- `scope`:`user` / `global` +- `endpoint`:`prompt` / `generate` +- `result`:`allowed` / `denied` + +配合告警规则 `RateLimitHighDenialRate`(拒绝率 > 5%),及时发现异常流量。 + +--- + +## 🔗 关联文档 + +- [← 返回索引](00-index.md) +- [10 — 中间件链](10-middleware-chain.md) — 限流中间件在链中的位置 +- [07 — 可观测性](07-observability.md) — 限流指标和告警规则 diff --git a/hzh/Gen2D/10-middleware-chain.md b/hzh/Gen2D/10-middleware-chain.md new file mode 100644 index 0000000..28baae3 --- /dev/null +++ b/hzh/Gen2D/10-middleware-chain.md @@ -0,0 +1,305 @@ +# 10 — 中间件链 + +> **一句话概括**:Logger → Recovery → Metrics → Auth → RateLimit → Handler,洋葱模型,层层守护请求处理。 + +--- + +```mermaid +flowchart LR + A["🌐 Request"] --> B["📝 Logger"] + B --> C["🛡️ Recovery"] + C --> D["📊 Metrics"] + D --> E["🔐 Auth"] + E --> F["🚦 RateLimit"] + F --> G["⚙️ Handler"] + G --> F + F --> E + E --> D + D --> C + C --> B + B --> H["📡 Response"] + + style A fill:#e3f2fd,stroke:#1976d2 + style B fill:#e8f5e9,stroke:#388e3c + style C fill:#fff3e0,stroke:#f57c00 + style D fill:#fce4ec,stroke:#c62828 + style E fill:#f3e5f5,stroke:#7b1fa2 + style F fill:#e0f7fa,stroke:#00838f + style G fill:#fff8e1,stroke:#f9a825 + style H fill:#e8eaf6,stroke:#303f9f +``` + +--- + +## 🧅 洋葱模型 + +Gin 的中间件采用**洋葱模型**:请求从外到内穿过各中间件,响应从内到外返回。每个中间件可以在 `c.Next()` 前后执行逻辑。 + +``` +请求 → Logger.enter → Recovery.enter → Metrics.enter → Auth.enter → RateLimit.enter → Handler +响应 ← Logger.leave ← Recovery.leave ← Metrics.leave ← Auth.leave ← RateLimit.leave ← Handler +``` + +--- + +## 📝 Logger — 请求日志 + +**职责**:为每个请求生成唯一 ID,记录请求详情。 + +### 核心逻辑 + +```go +func Logger() gin.HandlerFunc { + return func(c *gin.Context) { + requestID := generateRequestID() // 时间戳 + 8位随机hex + c.Set("request_id", requestID) + c.Header("X-Request-ID", requestID) + c.Next() + // 根据状态码选择日志级别 + if status >= 500 → Error + if status >= 400 → Warn + else → Info + } +} +``` + +### RequestID 生成 + +```go +func generateRequestID() string { + b := make([]byte, 8) + rand.Read(b) + return fmt.Sprintf("%d-%x", time.Now().UnixMilli(), b) +} +``` + +格式:`1717286400000-a1b2c3d4e5f67890` + +| 组件 | 作用 | +|------|------| +| 时间戳(毫秒) | 保证时间有序性 | +| 8 字节随机 hex | 保证唯一性 | + +### 日志级别映射 + +| HTTP 状态码 | 日志级别 | 含义 | +|:-----------:|:-------:|------| +| 5xx | ERROR | 服务器错误 | +| 4xx | WARN | 客户端错误 | +| 2xx/3xx | INFO | 正常请求 | + +--- + +## 🛡️ Recovery — Panic 恢复 + +**职责**:捕获未处理的 panic,防止服务崩溃。 + +```go +func Recovery() gin.HandlerFunc { + return func(c *gin.Context) { + defer func() { + if r := recover(); r != nil { + // 记录完整堆栈 + l.Error("panic recovered", + "error", fmt.Sprintf("%v", r), + "stack", string(debug.Stack()), + ) + // 返回结构化 500 + c.AbortWithStatusJSON(500, gin.H{ + "code": 500, + "message": "服务器内部错误", + }) + } + }() + c.Next() + } +} +``` + +**关键特性**: +- 捕获所有未处理的 panic +- 记录完整的 request 上下文和堆栈信息 +- 返回统一格式的 500 错误,避免泄露内部信息 +- 使用 `defer` 保证即使 panic 也能执行清理逻辑 + +--- + +## 📊 Metrics — 指标采集 + +**职责**:自动采集 HTTP 请求的性能指标。 + +```go +func Metrics() gin.HandlerFunc { + return func(c *gin.Context) { + path := c.FullPath() // 路由模板,避免高基数 + HTTPRequestsInFlight.Inc() + defer HTTPRequestsInFlight.Dec() + + reqSize := c.Request.ContentLength // 请求体大小(before) + c.Next() + elapsed := time.Since(start) // 耗时(after) + HTTPRequestsTotal.WithLabelValues(...).Inc() + } +} +``` + +**采集时机**: + +| 指标 | 采集时机 | 原因 | +|------|---------|------| +| InFlight | 进入时 +1,离开时 -1 | 使用 defer 保证异常时也能递减 | +| RequestSize | `c.Next()` 之前 | Content-Length 此时已知 | +| ResponseSize | `c.Next()` 之后 | Writer 此时已写入 | +| Duration | `c.Next()` 之后 | 需要计算总耗时 | + +> 💡 **FullPath() 的关键作用**:使用路由模板 `/api/v1/tasks/:taskId` 而非实际路径,避免高基数标签导致 Prometheus 内存爆炸。 + +--- + +## 🔐 Auth — JWT 认证 + +**职责**:验证 Bearer token,提取用户身份。 + +```go +func AuthMiddleware(jwtSecret string) gin.HandlerFunc { + return func(c *gin.Context) { + // 1. 提取 Authorization header + authHeader := c.GetHeader("Authorization") + tokenString := strings.TrimPrefix(authHeader, "Bearer ") + + // 2. 解析和验证 JWT + token, _ := jwt.Parse(tokenString, func(token *jwt.Token) (interface{}, []byte) { + return []byte(jwtSecret), nil + }) + + // 3. 提取 sub claim → userID + claims := token.Claims.(jwt.MapClaims) + c.Set("userID", claims["sub"]) + + c.Next() + } +} +``` + +**作用范围**:仅应用于 `v1Auth` 路由组,公开接口(如 `/health`)不经过认证。 + +**错误处理**: + +| 场景 | HTTP 状态码 | 消息 | +|------|:-----------:|------| +| 无 token | 401 | 未提供认证令牌 | +| 格式错误 | 401 | 认证格式错误,需为 Bearer \ | +| token 无效 | 401 | 令牌无效或已过期 | +| 解析失败 | 401 | 令牌解析失败 | + +--- + +## 🚦 RateLimit — 限流 + +**职责**:按路由配置执行双层限流(全局 + 用户)。 + +```go +// 每个路由独立配置 +v1Auth.POST("/generate", + ratelimit.RateLimit(globalLimiter, ..., "global", "generate"), + ratelimit.RateLimit(userLimiter, ..., "user", "generate"), + handler.Generate, +) +``` + +详见 [09 — 限流](09-rate-limiting.md)。 + +--- + +## 🔗 中间件挂载 + +### 全局链(所有请求) + +```go +r := gin.New() +r.Use(mildware.Logger()) // 1. 请求日志 +r.Use(mildware.Recovery()) // 2. Panic 恢复 +r.Use(mildware.Metrics()) // 3. 指标采集 +``` + +### 路由组级链(需认证) + +```go +v1Auth := r.Group("/api/v1") +v1Auth.Use(mildware.AuthMiddleware(cfg.JWT.Secret)) // 4. JWT 认证 +``` + +### 端点级链(需限流) + +```go +v1Auth.POST("/generate", + ratelimit.RateLimit(globalLimiter, ...), // 5. 全局限流 + ratelimit.RateLimit(userLimiter, ...), // 6. 用户限流 + handler.Generate, // 7. 业务处理 +) +``` + +--- + +## 📋 端点链示例:/api/v1/generate + +```mermaid +flowchart TB + A["POST /api/v1/generate"] --> B["Logger
生成 request_id"] + B --> C["Recovery
注册 panic 恢复"] + C --> D["Metrics
InFlight++, 记录 reqSize"] + D --> E["Auth
验证 JWT, 提取 userID"] + E --> F["RateLimit Global
检查全局限额"] + F --> G["RateLimit User
检查用户限额"] + G --> H["Handler.Generate
执行生成逻辑"] + H --> I["Metrics
记录 duration, respSize"] + I --> J["Recovery
检查是否 panic"] + J --> K["Logger
记录请求日志"] + K --> L["📡 Response"] + + style A fill:#e3f2fd,stroke:#1976d2 + style B fill:#e8f5e9,stroke:#388e3c + style C fill:#fff3e0,stroke:#f57c00 + style D fill:#fce4ec,stroke:#c62828 + style E fill:#f3e5f5,stroke:#7b1fa2 + style F fill:#e0f7fa,stroke:#00838f + style G fill:#e0f7fa,stroke:#00838f + style H fill:#fff8e1,stroke:#f9a825 + style L fill:#e8eaf6,stroke:#303f9f +``` + +### 完整请求生命周期 + +| 阶段 | 中间件 | 动作 | +|:----:|--------|------| +| 1 | Logger | 生成 request_id,注入 context 和 header | +| 2 | Recovery | 注册 defer panic 恢复 | +| 3 | Metrics | InFlight +1,记录请求体大小 | +| 4 | Auth | 验证 JWT,提取 userID | +| 5 | RateLimit | 检查全局限额 | +| 6 | RateLimit | 检查用户限额 | +| 7 | Handler | 执行业务逻辑 | +| 8 | Metrics | 记录耗时、响应体大小,InFlight -1 | +| 9 | Recovery | 检查是否发生 panic | +| 10 | Logger | 记录请求日志(含状态码、耗时) | + +--- + +## 📊 中间件职责矩阵 + +| 中间件 | 请求进入 | 请求离开 | 异常处理 | 作用范围 | +|--------|---------|---------|---------|---------| +| Logger | 生成 request_id | 记录日志 | — | 全局 | +| Recovery | 注册 defer | — | panic → 500 | 全局 | +| Metrics | InFlight++, reqSize | duration, respSize | — | 全局 | +| Auth | 验证 JWT | — | 无效 → 401 | v1Auth | +| RateLimit | 检查配额 | — | 超限 → 429 | 端点级 | + +--- + +## 🔗 关联文档 + +- [← 返回索引](00-index.md) +- [07 — 可观测性](07-observability.md) — Metrics 中间件采集的指标 +- [09 — 限流](09-rate-limiting.md) — RateLimit 中间件的详细实现 +- [08 — SSE 实时推送](08-sse-push.md) — SSE 端点的中间件配置 diff --git a/hzh/Gen2D/11-consumer-producer.md b/hzh/Gen2D/11-consumer-producer.md new file mode 100644 index 0000000..7f82b64 --- /dev/null +++ b/hzh/Gen2D/11-consumer-producer.md @@ -0,0 +1,185 @@ +# 11. Consumer-Producer 桥接模式 + +> **一句话概括**:`Consumer` 结构体桥接 `TaskQueue` 和 `WorkerPool`,实现生产者与消费者的彻底解耦。 + +--- + +## 架构总览 + +```mermaid +flowchart LR + subgraph Producer["生产者"] + A[Generate Handler] + end + + subgraph Queue["TaskQueue 接口"] + B((Memory\nQueue)) + C((RabbitMQ\nQueue)) + end + + subgraph Bridge["Consumer 桥接层"] + D{{"Consumer\n(bridge)"}} + end + + subgraph Pool["WorkerPool"] + E[Worker 1] + F[Worker 2] + G[Worker N] + end + + subgraph Pipeline["业务逻辑"] + H[Eino Pipeline] + end + + A -->|"Submit(msg)"| B + A -->|"Submit(msg)"| C + B -->|"Consume()"| D + C -->|"Consume()"| D + D -->|"Submit(task)"| E + D -->|"Submit(task)"| F + D -->|"Submit(task)"| G + E --> H + F --> H + G --> H + + style D fill:#f9a825,stroke:#333,color:#000 +``` + +--- + +## 核心结构体 + +`Consumer` 是整个任务调度体系的**桥梁**,它只做一件事:从队列取消息,提交到协程池。 + +```go +// worker/consumer.go +type Consumer struct { + queue taskqueue.TaskQueue // 可插拔队列接口 + pool *workerpool.Pool // 有界协程池 + handler TaskHandler // 业务逻辑注入点 + logger *slog.Logger +} +``` + +| 字段 | 类型 | 职责 | +|------|------|------| +| `queue` | `TaskQueue` 接口 | 消息来源,支持 Memory / RabbitMQ 替换 | +| `pool` | `*workerpool.Pool` | 并发执行引擎,控制单机并行度 | +| `handler` | `TaskHandler` | 回调函数,由 handler 层注入实际业务逻辑 | + +--- + +## 工作流程 + +### 启动消费循环 + +```mermaid +sequenceDiagram + participant Main as main.go + participant Consumer + participant Queue as TaskQueue + participant Pool as WorkerPool + participant Handler as RunFromTaskMessage + + Main->>Consumer: Start(ctx) + loop 持续消费 + Consumer->>Queue: Consume(ctx, callback) + Queue-->>Consumer: TaskMessage + Consumer->>Consumer: processMessage() + Consumer->>Pool: Submit(workerpool.Task) + Pool-->>Consumer: accepted + Pool->>Handler: Fn(ctx) + Handler->>Handler: runPipelineBg() + end +``` + +四步循环: + +1. **消费** — `Consumer` 调用 `queue.Consume(ctx, handler)`,阻塞等待消息 +2. **转换** — 将 `taskqueue.TaskMessage` 包装为 `workerpool.Task` +3. **提交** — 调用 `pool.Submit(task)` 送入协程池执行 +4. **执行** — `Task.Fn` 回调实际的 `RunFromTaskMessage`,驱动 Eino 管线 + +--- + +## 解耦的三层设计 + +```mermaid +flowchart TB + subgraph "第 1 层:消息源" + Q["TaskQueue 接口\n(Memory / RabbitMQ)"] + end + subgraph "第 2 层:桥接" + C["Consumer\n(只关心 消费→提交)"] + end + subgraph "第 3 层:执行引擎" + P["WorkerPool\n(只关心 并发控制)"] + end + subgraph "第 4 层:业务逻辑" + H["TaskHandler 回调\n(RunFromTaskMessage)"] + end + + Q --> C --> P --> H +``` + +| 组件 | 知道什么 | 不知道什么 | +|------|----------|------------| +| **TaskQueue** | 消息的存储与投递 | WorkerPool 的存在 | +| **WorkerPool** | 任务的并发执行 | 消息来自哪个队列 | +| **Consumer** | 如何桥接两者 | 具体的业务逻辑 | +| **TaskHandler** | 生成管线的执行 | 消息来自内存还是 RabbitMQ | + +> **设计哲学**:每个组件只关心自己的职责边界,可独立替换、测试、扩展。 + +--- + +## Handler 层注入 + +`Consumer` 不硬编码业务逻辑,而是通过 `TaskHandler` 函数签名由外部注入: + +```go +// worker/consumer.go — 定义 +type TaskHandler func(ctx context.Context, msg taskqueue.TaskMessage) + +// main.go — 注入 +consumer := worker.NewConsumer(tq, pool, handler.RunFromTaskMessage) +``` + +`RunFromTaskMessage` 负责将队列消息还原为 `GenerateRequest`,再调用 `runPipelineBg` 驱动完整的 Eino 管线。 + +--- + +## 信号处理与优雅关闭 + +```mermaid +sequenceDiagram + participant OS as 操作系统 + participant Main as main.go + participant Consumer + participant Pool as WorkerPool + + OS->>Main: SIGINT / SIGTERM + Main->>Consumer: cancel() 停止消费 + Note over Consumer: 不再接收新消息 + Main->>Pool: Shutdown(30s timeout) + Note over Pool: 等待正在执行的任务完成 + Pool-->>Main: true (正常) / false (超时) + Main->>Main: 进程退出 +``` + +关闭顺序至关重要: + +1. **先停 Consumer** — 不再从队列拉取新消息 +2. **再关 WorkerPool** — 等待已提交的任务执行完毕(最多 30 秒) +3. **最后关闭队列连接** — 释放 RabbitMQ / 内存资源 + +--- + +## 关联文档 + +| 文档 | 关系 | +|------|------| +| [03 - 任务队列](03-task-queue.md) | Consumer 的消息来源 | +| [02 - 协程池](02-worker-pool.md) | Consumer 的执行引擎 | +| [12 - 三级降级策略](12-three-tier-fallback.md) | Consumer 不参与降级,降级在 Handler 层 | +| [00 - 索引](00-index.md) | 返回文档总览 | diff --git a/hzh/Gen2D/12-three-tier-fallback.md b/hzh/Gen2D/12-three-tier-fallback.md new file mode 100644 index 0000000..ff21bc0 --- /dev/null +++ b/hzh/Gen2D/12-three-tier-fallback.md @@ -0,0 +1,183 @@ +# 12. 三级降级策略 + +> **一句话概括**:Generate handler 实现三级降级 — TaskQueue -> WorkerPool -> Legacy FIFO,宁可降级也不拒绝服务。 + +--- + +## 降级链总览 + +```mermaid +flowchart TD + R["HTTP Request\nPOST /api/v1/generate"] --> A{"TaskQueue\n可用?"} + + A -->|"Submit 成功"| S1["200 OK\ntaskId 返回"] + A -->|"Submit 失败"| B{"WorkerPool\n可用?"} + + B -->|"Submit 成功"| S2["200 OK\ntaskId 返回"] + B -->|"ErrPoolFull"| E2["503 Service\nUnavailable"] + B -->|"ErrUserLimit"| E3["429 Too Many\nRequests"] + B -->|"Pool 不可用"| C{"Legacy FIFO\nQueue 可用?"} + + C -->|"Enqueue 成功"| S3["200 OK\ntaskId 返回"] + C -->|"Queue 不可用"| E4["500 Internal\nServer Error"] + + style A fill:#4caf50,stroke:#333,color:#fff + style B fill:#ff9800,stroke:#333,color:#fff + style C fill:#f44336,stroke:#333,color:#fff + style S1 fill:#8bc34a,stroke:#333,color:#fff + style S2 fill:#8bc34a,stroke:#333,color:#fff + style S3 fill:#8bc34a,stroke:#333,color:#fff +``` + +--- + +## 三级详解 + +### 第一级:TaskQueue(优先路径) + +| 属性 | 说明 | +|------|------| +| **组件** | `taskqueue.TaskQueue` 接口(Memory / RabbitMQ) | +| **特点** | 支持持久化、分布式、消息确认 | +| **提交** | `taskQueue.Submit(ctx, msg)` | +| **失败时** | 返回 HTTP 500,标记任务为 `failed` | + +```go +// handler/generate.go — 第一级 +if taskQueue != nil { + msg := taskqueue.TaskMessage{...} + if err := taskQueue.Submit(ctx, msg); err != nil { + c.JSON(500, "提交任务失败") + return + } + c.JSON(200, taskId) + return // 成功,不再降级 +} +``` + +### 第二级:WorkerPool(有界并发) + +| 属性 | 说明 | +|------|------| +| **组件** | `workerpool.Pool` | +| **特点** | 有界队列 + per-user 并发限制 | +| **提交** | `workerPool.Submit(task)` | +| **失败类型** | `ErrPoolFull` -> 503 / `ErrUserLimitReached` -> 429 | + +```go +// handler/generate.go — 第二级 +if workerPool != nil { + err := workerPool.Submit(workerpool.Task{...}) + switch { + case errors.Is(err, workerpool.ErrPoolFull): + c.JSON(503, "系统繁忙,请稍后重试") + case errors.Is(err, workerpool.ErrUserLimitReached): + c.JSON(429, "您的生成任务已达上限") + default: + c.JSON(500, "提交任务失败") + } + return +} +``` + +### 第三级:Legacy FIFO(最后保底) + +| 属性 | 说明 | +|------|------| +| **组件** | `service.TaskQueue`(串行 FIFO 队列) | +| **特点** | 零配置、串行执行、无并发控制 | +| **用途** | 兜底,确保系统在任何配置下都能运行 | + +```go +// handler/generate.go — 第三级 +if generateQueue != nil { + generateQueue.Enqueue(&service.TaskJob{...}) +} +c.JSON(200, taskId) +``` + +--- + +## HTTP 状态码映射 + +```mermaid +flowchart LR + subgraph "降级路径" + TQ["TaskQueue"] + WP["WorkerPool"] + LF["Legacy FIFO"] + end + + subgraph "HTTP 响应" + E200["200 OK\n任务已接受"] + E429["429 Too Many Requests\n用户限流"] + E500["500 Internal Server Error\n系统错误"] + E503["503 Service Unavailable\n系统繁忙"] + end + + TQ -->|"成功"| E200 + TQ -->|"失败"| E500 + WP -->|"成功"| E200 + WP -->|"队列满"| E503 + WP -->|"用户限流"| E429 + LF -->|"成功"| E200 + LF -->|"不可用"| E500 +``` + +| 状态码 | 含义 | 触发条件 | +|--------|------|----------| +| `200` | 任务已接受 | 任意一级提交成功 | +| `429` | 用户限流 | WorkerPool 用户并发达到上限 | +| `500` | 系统内部错误 | TaskQueue 提交失败 / 所有级别不可用 | +| `503` | 服务暂不可用 | WorkerPool 队列已满 | + +--- + +## 为什么需要三级? + +```mermaid +flowchart TB + subgraph "第一级:分布式能力" + TQ["TaskQueue\n持久化 + 消息确认\n支持 RabbitMQ 横向扩展"] + end + subgraph "第二级:并发控制" + WP["WorkerPool\n有界队列背压\nper-user 限流保护"] + end + subgraph "第三级:可用性兜底" + LF["Legacy FIFO\n零依赖、零配置\n确保始终可运行"] + end + + TQ -->|"不可用时降级到"| WP + WP -->|"不可用时降级到"| LF +``` + +| 级别 | 核心价值 | 典型场景 | +|------|----------|----------| +| TaskQueue | 分布式 + 持久化 | 生产环境,多实例部署 | +| WorkerPool | 并发控制 + 背压 | 单机部署,需要限制资源 | +| Legacy FIFO | 可用性兜底 | 开发/测试环境,或队列组件故障 | + +--- + +## 设计哲学 + +> **宁可降级,也不能拒绝服务。** + +三级降级的核心思想是**渐进式降级**: + +1. **功能完整** — TaskQueue 提供持久化、重试、死信等高级特性 +2. **性能可控** — WorkerPool 提供有界并发和用户隔离 +3. **始终可用** — Legacy FIFO 确保在任何配置下都能接收任务 + +每一级降级都意味着功能的减少,但服务的可用性始终得到保障。前端收到 `200 OK` 后,通过 SSE 或轮询获取任务进度,对用户而言体验一致。 + +--- + +## 关联文档 + +| 文档 | 关系 | +|------|------| +| [11 - Consumer-Producer 桥接](11-consumer-producer.md) | Consumer 连接 TaskQueue 和 WorkerPool | +| [02 - 协程池](02-worker-pool.md) | WorkerPool 的背压和限流机制 | +| [03 - 任务队列](03-task-queue.md) | TaskQueue 接口的可插拔设计 | +| [00 - 索引](00-index.md) | 返回文档总览 | diff --git a/hzh/Gen2D/13-prompt-engineering.md b/hzh/Gen2D/13-prompt-engineering.md new file mode 100644 index 0000000..3d87157 --- /dev/null +++ b/hzh/Gen2D/13-prompt-engineering.md @@ -0,0 +1,216 @@ +# 13. 标签驱动提示词工程 + +> **一句话概括**:40+ 预定义标签映射到精确的图像生成指令,保障风格一致性与管线友好性。 + +--- + +## 处理流程 + +```mermaid +flowchart LR + A["用户选择标签\n+ 输入描述"] --> B["Tag Mapper\n标签→指令映射"] + B --> C["Prompt Template\n三段式结构组装"] + C --> D{"LLM 可用?"} + D -->|"是"| E["PromptOptimizer\nLLM 精炼"] + D -->|"否"| F["Fallback Template\n模板生成"] + E --> G["Final Prompt\n最终提示词"] + F --> G + G --> H["AssetGenerator\n文生图 API"] + + style B fill:#7c4dff,stroke:#333,color:#fff + style D fill:#ff9800,stroke:#333,color:#fff + style G fill:#4caf50,stroke:#333,color:#fff +``` + +--- + +## 标签分类体系 + +系统内置 40+ 预定义标签,分为 **7 大类别**,每条标签精确映射到一条图像生成指令。 + +### 内容类型 — 决定布局与格式 + +| 标签 | 映射指令要点 | +|------|-------------| +| `瓦片/图块` | 瓦片集,可无缝拼接,网格排列,白色背景 | +| `背景/远景` | 三层结构(远景/中景/前景),视差滚动适配 | +| `前景/装饰` | 独立遮挡物,白色背景,可叠加 | +| `图集` | 多元素打包,网格排列,白色背景 | +| `纹理` | 可无缝平铺,白色背景 | +| `序列帧` | 连续动画帧,网格排列,标注方向和帧数 | +| `纸娃娃部件` | 可组合散件,统一比例和锚点 | + +### 美术风格 — 决定渲染技术 + +| 标签 | 映射指令要点 | +|------|-------------| +| `像素` | 严格像素网格,无抗锯齿,有限色盘 | +| `卡通` | 粗轮廓线,明亮饱和色彩,夸张比例 | +| `手绘` | 自然笔触纹理,不规则线条 | +| `矢量` | 干净几何形状,平滑曲线 | +| `扁平` | 无阴影或极少阴影,纯色块面 | + +### 色调配色 + +| 标签 | 映射指令要点 | +|------|-------------| +| `暖色` | 红橙黄为主,温暖活力氛围 | +| `冷色` | 蓝青紫为主,冷静神秘氛围 | +| `鲜艳` | 高饱和度,强烈对比 | +| `柔和` | 低饱和度,温和内敛 | +| `单色` | 单一色相,明暗层次 | + +### 线条粗细 + +| 标签 | 映射指令要点 | +|------|-------------| +| `无` | 无线条轮廓,纯色块面 | +| `细线` | 0.5-1px,细腻精致 | +| `中等` | 1-2px,清晰明确 | +| `粗线` | 2-4px,粗犷有力 | + +### 场景氛围 + +| 标签 | 映射指令要点 | +|------|-------------| +| `森林` | 树木、藤蔓、蘑菇、落叶 | +| `地牢` | 石砖墙壁、火把、蛛网 | +| `城市` | 建筑、街道、路灯 | +| `太空` | 星云、行星、飞船 | +| `水下` | 珊瑚、水草、气泡 | +| `沙漠` | 沙丘、仙人掌、绿洲 | + +### 光照效果 + +| 标签 | 映射指令要点 | +|------|-------------| +| `明亮` | 充足自然光,无深阴影 | +| `昏暗` | 低光照,柔和阴影 | +| `戏剧` | 强烈明暗对比,聚光灯效果 | +| `霓虹` | 高饱和彩色光源,赛博朋克辉光 | + +### 情绪基调 + +| 标签 | 映射指令要点 | +|------|-------------| +| `欢快` | 明亮色彩,圆润造型 | +| `黑暗` | 深色调,尖锐造型 | +| `神秘` | 朦胧效果,未知元素 | +| `史诗` | 宏大场面,壮观远景 | +| `平静` | 柔和色彩,开阔空间 | + +--- + +## 精灵图特殊处理 + +```mermaid +flowchart TD + T["用户标签"] --> G{"包含网格标签?\n瓦片/图集/序列帧/纸娃娃"} + G -->|"是"| S["注入网格布局指令\n白色背景 + 8-16px 间隙\n投影法可检测"] + G -->|"否"| N["标准布局指令"] + S --> P["Prompt 构建"] + N --> P +``` + +当标签命中 `gridLayoutTags` 集合时,自动注入网格布局专用指令: + +- **纯白色背景** (`#FFFFFF`) +- **帧间 8-16px 纯白间隙** +- **间隙内不得有任何像素** +- **确保投影法能可靠检测间隙** + +这些指令对下游的 `SplitSprite` 切割算法至关重要。 + +--- + +## 未知标签回退 + +当用户输入系统未预定义的标签时,不会报错,而是降级为通用风格描述: + +```go +func tagToInstruction(tag string) string { + if inst, ok := tagInstruction[tag]; ok { + return inst + } + // 未知标签:作为风格修饰词处理 + return fmt.Sprintf("风格特征: %s(用户自定义标签,按其字面含义应用到画面中)", tag) +} +``` + +这保证了系统的**向前兼容** — 新增标签无需修改核心逻辑。 + +--- + +## PromptOptimizer 节点 + +PromptOptimizer 是 Eino 管线的第一个节点,负责将标签指令组装为最终提示词。 + +```mermaid +flowchart LR + subgraph "输入" + A["全局风格提示词\nprojectStyle"] + B["任务描述\nprompt + userNote"] + C["标签指令\ntagInstructions"] + D["拒绝原因\nqualityFeedback\n(重试时)"] + end + + subgraph "PromptOptimizer" + E["buildMetaPrompt\n组装元提示词"] + F{"LLM 可用?"} + G["callLLMRefine\nLLM 精炼"] + H["fallbackRefine\n模板生成"] + end + + subgraph "输出" + I["Final Prompt\n三段式结构"] + end + + A --> E + B --> E + C --> E + D --> E + E --> F + F -->|"API Key 已配置"| G + F -->|"未配置 / 调用失败"| H + G --> I + H --> I +``` + +三段式输出结构: + +| 段落 | 内容 | 示例 | +|------|------|------| +| 【主题】 | 画面主体与场景 | "一个融合像素风格的游戏角色精灵图..." | +| 【风格】 | 艺术风格与视觉特征 | "美术风格: 像素画;色调: 暖色系..." | +| 【技术】 | 分辨率和格式参数 | "输出格式: spritesheet;纯白色背景..." | + +--- + +## LLM 回退机制 + +当 LLM 不可用时(API Key 未配置 / 网络故障),系统自动降级到模板生成: + +```mermaid +flowchart TD + A["callLLMRefine"] --> B{"API Key\n已配置?"} + B -->|"否"| F["fallbackRefine\n模板回退"] + B -->|"是"| C["调用 Chat API"] + C --> D{"调用成功?"} + D -->|"是"| E["返回 LLM 优化结果"] + D -->|"否"| F + F --> G["解析标签\n组装三段式提示词"] + G --> H["返回模板生成结果"] +``` + +模板回退同样遵循标签驱动逻辑,保证即使没有 LLM 参与,生成的提示词也具备结构化和一致性。 + +--- + +## 关联文档 + +| 文档 | 关系 | +|------|------| +| [05 - 生成管线](05-generation-pipeline.md) | PromptOptimizer 是管线第一阶段 | +| [06 - 精灵图处理](06-sprite-processing.md) | 网格布局指令影响切割算法 | +| [01 - 系统总览](01-system-overview.md) | 提示词工程在整体架构中的位置 | +| [00 - 索引](00-index.md) | 返回文档总览 | diff --git a/hzh/Gen2D/14-deployment.md b/hzh/Gen2D/14-deployment.md new file mode 100644 index 0000000..01d0f1b --- /dev/null +++ b/hzh/Gen2D/14-deployment.md @@ -0,0 +1,200 @@ +# 14. 部署架构 + +> **一句话概括**:Docker Compose 编排 + Prometheus 监控 + Grafana 可视化 + 自动化部署脚本。 + +--- + +## 容器拓扑 + +```mermaid +flowchart TB + subgraph "Docker Compose" + subgraph "应用层" + FE["frontend-v2\nReact SPA\nNginx :80"] + BE["backend-v2\nGo HTTP Server\n:8080"] + end + + subgraph "数据层" + DB["MySQL\n持久化存储"] + RD["Redis\n缓存 + 限流"] + RMQ["RabbitMQ\n消息队列"] + end + + subgraph "监控层" + PR["Prometheus\n指标采集 :9090"] + GF["Grafana\n可视化 :3000"] + end + + subgraph "外部服务" + AI["LLM API\n(OpenAI 兼容)"] + IG["Image Gen API\n文生图"] + QN["七牛云 CDN\n对象存储"] + end + end + + FE -->|"HTTP API"| BE + BE -->|"SQL"| DB + BE -->|"Redis 协议"| RD + BE -->|"AMQP"| RMQ + BE -->|"HTTP"| AI + BE -->|"HTTP"| IG + BE -->|"Upload"| QN + PR -->|"/metrics\n每 10s"| BE + GF -->|"PromQL"| PR + FE -->|"CDN URL"| QN + + style FE fill:#61dafb,stroke:#333,color:#000 + style BE fill:#00add8,stroke:#333,color:#fff + style PR fill:#e6522c,stroke:#333,color:#fff + style GF fill:#f46800,stroke:#333,color:#fff +``` + +--- + +## 服务组成 + +| 服务 | 镜像 | 端口 | 职责 | +|------|------|------|------| +| `backend-v2` | 自建 Go 镜像 | `9001:8080` | HTTP API + 管线执行 | +| `frontend-v2` | 自建 Nginx 镜像 | `9000:80` | React SPA 静态资源 | +| `redis` | `redis:7-alpine` | `6379` | 令牌桶限流 + 缓存 | +| `rabbitmq` | `rabbitmq:3-management` | `5672/15672` | 可选消息队列 | +| `prometheus` | `prom/prometheus` | `9090` | 指标采集与告警 | +| `grafana` | `grafana/grafana` | `3000` | 仪表盘可视化 | + +--- + +## 监控栈 + +### Prometheus 配置 + +```yaml +# deploy/prometheus/prometheus.yml +scrape_configs: + - job_name: "gen2d-backend" + static_configs: + - targets: ["backend-v2:8080"] + metrics_path: "/metrics" + scrape_interval: 10s +``` + +Prometheus 每 **10 秒**抓取一次后端的 `/metrics` 端点,采集全部 35+ 指标。 + +### Grafana 自动化 + +```mermaid +flowchart LR + subgraph "Provisioning" + DS["datasources/\n自动配置数据源"] + DB["dashboards/\n自动导入仪表盘"] + end + + subgraph "Grafana" + GF["Grafana Server\n:3000"] + end + + DS -->|"启动时加载"| GF + DB -->|"启动时加载"| GF +``` + +Grafana 通过 provisioning 机制自动加载: +- **数据源配置** — 指向 Prometheus 实例 +- **仪表盘 JSON** — 预定义的 3 个仪表盘 + +### 告警规则 + +10 条告警规则覆盖全栈关键指标: + +| 告警名称 | 条件 | 严重程度 | +|----------|------|----------| +| 后端实例宕机 | `up == 0` 持续 1 分钟 | Critical | +| HTTP 5xx 错误率 | `> 5%` 持续 2 分钟 | Warning | +| 请求延迟过高 | `P99 > 5s` 持续 3 分钟 | Warning | +| 管线任务失败率 | `> 10%` 持续 5 分钟 | Warning | +| 协程池队列满 | `pool_queued_tasks > 90` | Warning | +| 协程池拒绝率 | `> 20%` 持续 2 分钟 | Critical | +| 用户限流触发 | `rate_limit_rejected > 50/min` | Info | +| Redis 连接失败 | `redis_up == 0` | Warning | +| RabbitMQ 连接失败 | `rabbitmq_up == 0` | Warning | +| 磁盘空间不足 | `< 10%` 可用 | Critical | + +--- + +## 部署脚本 + +```mermaid +flowchart TD + A["deploy.sh"] --> B["拉取最新代码"] + B --> C["构建后端镜像\ndocker build"] + C --> D["构建前端镜像\ndocker build"] + D --> E["docker compose up -d"] + E --> F["健康检查\n等待服务就绪"] + F --> G["部署完成"] + + style A fill:#4caf50,stroke:#333,color:#fff +``` + +部署脚本 `deploy.sh` 负责: +1. 拉取最新代码 +2. 构建 Docker 镜像 +3. 使用 Docker Compose 启动所有服务 +4. 执行健康检查验证部署 + +--- + +## 配置管理 + +```mermaid +flowchart LR + subgraph "配置源" + Y["config.yaml\n文件配置"] + E[".env\n环境变量"] + D["默认值\n代码内置"] + end + + subgraph "加载优先级" + P["ENV > YAML > Default"] + end + + subgraph "运行时" + C["Config Struct\n全局配置对象"] + end + + D --> P + Y --> P + E --> P + P --> C +``` + +| 环境 | 配置方式 | 说明 | +|------|----------|------| +| 开发 | `config.yaml` 文件 | 本地开发,配置直观 | +| 测试 | 环境变量覆盖 | CI/CD 管线注入 | +| 生产 | K8s ConfigMap / Secret | 容器编排平台管理 | + +--- + +## 网络与存储 + +### Docker 网络 + +所有服务加入 `gen2d-v2-net` 桥接网络,容器间通过服务名互相访问。 + +### 数据卷 + +| 卷名 | 挂载点 | 用途 | +|------|--------|------| +| `backend-v2-data` | `/data` | 后端持久化数据 | +| MySQL data | 默认 | 数据库持久化 | +| Redis data | 默认 | 缓存持久化 | + +--- + +## 关联文档 + +| 文档 | 关系 | +|------|------| +| [07 - 可观测性](07-observability.md) | Prometheus 指标与 Grafana 仪表盘详情 | +| [15 - 配置级联](15-config-cascade.md) | YAML / ENV / Default 三层配置机制 | +| [01 - 系统总览](01-system-overview.md) | 部署架构在整体系统中的位置 | +| [00 - 索引](00-index.md) | 返回文档总览 | diff --git a/hzh/Gen2D/15-config-cascade.md b/hzh/Gen2D/15-config-cascade.md new file mode 100644 index 0000000..f8d92ce --- /dev/null +++ b/hzh/Gen2D/15-config-cascade.md @@ -0,0 +1,197 @@ +# 15. 配置级联机制 + +> **一句话概括**:Viper 三层配置级联 — YAML 文件 -> 环境变量 -> 默认值,一处配置随处运行。 + +--- + +## 配置加载流程 + +```mermaid +flowchart LR + subgraph "配置源(优先级从低到高)" + D["默认值\nsetDefaults()"] + Y["YAML 文件\nconfig.yaml"] + E["环境变量\n.env / system env"] + end + + subgraph "Viper 引擎" + V["Viper\n合并 + 覆盖"] + end + + subgraph "输出" + C["Config Struct\n全局配置对象"] + S["各组件\nServer / DB / Redis / ..."] + end + + D -->|"1. 设置默认值"| V + Y -->|"2. 读取 YAML"| V + E -->|"3. 绑定 ENV"| V + V -->|"Unmarshal"| C + C -->|"注入"| S + + style D fill:#9e9e9e,stroke:#333,color:#fff + style Y fill:#2196f3,stroke:#333,color:#fff + style E fill:#f44336,stroke:#333,color:#fff + style V fill:#ff9800,stroke:#333,color:#fff +``` + +**优先级规则**:`ENV > YAML > Default` + +环境变量始终拥有最高优先级,可以覆盖任何 YAML 配置;YAML 配置覆盖默认值。 + +--- + +## 配置结构体 + +```go +type Config struct { + Server ServerConfig // HTTP 服务 + Database DatabaseConfig // MySQL 数据库 + JWT JWTConfig // JWT 签名 + Log LogConfig // 日志 + Redis RedisConfig // Redis 缓存 + LLM LLMConfig // 大语言模型 + ImageGen ImageGenConfig // 文生图 API + Qiniu QiniuConfig // 七牛云存储 + WorkerPool WorkerPoolConfig // 协程池 + TaskQueue TaskQueueConfig // 任务队列 +} +``` + +| 配置块 | 关键字段 | 默认值 | +|--------|----------|--------| +| `Server` | `port`, `mode` | `8080`, `debug` | +| `Database` | `dsn` | `gen2d:password@tcp(127.0.0.1:3306)/gen2d` | +| `JWT` | `secret`, `expire` | `gen2d-dev-secret`, `7200s` | +| `Log` | `level`, `format` | `info`, `text` | +| `Redis` | `addr`, `password`, `db` | `localhost:6379`, `""`, `0` | +| `LLM` | `base_url`, `api_key`, `model` | `api.openai.com/v1`, `""`, `gpt-4o` | +| `ImageGen` | `base_url`, `model`, `timeout` | `api.suchuang.vip/v1`, `gpt-image-2-token`, `120s` | +| `Qiniu` | `bucket`, `cdn_host`, `url_expire` | `""`, `""`, `3600s` | +| `WorkerPool` | `workers`, `queue_size`, `max_per_user` | `NumCPU*4`, `100`, `2` | +| `TaskQueue` | `driver`, `memory.buffer_size` | `memory`, `100` | + +--- + +## 环境变量绑定 + +每个配置字段都有对应的环境变量绑定,命名规则为 `GEN2D_` 前缀 + 大写下划线格式: + +```go +func bindEnvVars(v *viper.Viper) { + v.BindEnv("server.port", "GEN2D_PORT") + v.BindEnv("server.mode", "GEN2D_MODE") + v.BindEnv("database.dsn", "GEN2D_DSN") + v.BindEnv("jwt.secret", "GEN2D_JWT_SECRET") + v.BindEnv("llm.api_key", "GEN2D_LLM_API_KEY") + v.BindEnv("redis.addr", "GEN2D_REDIS_ADDR") + v.BindEnv("qiniu.access_key", "GEN2D_QINIU_ACCESS_KEY") + // ... 共 25+ 个绑定 +} +``` + +| 配置路径 | 环境变量 | 说明 | +|----------|----------|------| +| `server.port` | `GEN2D_PORT` | HTTP 端口 | +| `server.mode` | `GEN2D_MODE` | `debug` / `release` | +| `database.dsn` | `GEN2D_DSN` | MySQL 连接串 | +| `jwt.secret` | `GEN2D_JWT_SECRET` | JWT 签名密钥 | +| `llm.base_url` | `GEN2D_LLM_BASE_URL` | LLM API 地址 | +| `llm.api_key` | `GEN2D_LLM_API_KEY` | LLM API 密钥 | +| `image_gen.base_url` | `GEN2D_IMAGE_BASE_URL` | 文生图 API 地址 | +| `redis.addr` | `GEN2D_REDIS_ADDR` | Redis 地址 | +| `taskqueue.driver` | `GEN2D_TASKQUEUE_DRIVER` | `memory` / `rabbitmq` | + +--- + +## 使用场景 + +```mermaid +flowchart TD + subgraph "开发环境" + D1["config.yaml\n本地配置文件"] + D2[".env\n环境变量覆盖敏感值"] + D3["默认值\n开箱即用"] + end + + subgraph "测试环境" + T1["环境变量\nCI/CD 管线注入"] + T2["默认值\n兜底"] + end + + subgraph "生产环境" + P1["K8s ConfigMap\n非敏感配置"] + P2["K8s Secret\n敏感配置(密钥、DSN)"] + P3["默认值\n兜底"] + end + + D1 --> D3 + D2 --> D1 + T1 --> T2 + P1 --> P3 + P2 --> P1 +``` + +| 环境 | 主配置源 | 敏感配置 | 说明 | +|------|----------|----------|------| +| 开发 | `config.yaml` | `.env` 文件 | 直观、可版本控制 | +| 测试 | 环境变量 | 环境变量 | CI/CD 管线注入 | +| 生产 | ConfigMap | Secret | K8s 原生管理 | + +--- + +## 加载过程详解 + +```mermaid +sequenceDiagram + participant Main as main.go + participant Viper + participant YAML as config.yaml + participant ENV as 环境变量 + participant Cfg as Config Struct + + Main->>Viper: Load() + Main->>Viper: setDefaults(v) + Note over Viper: 设置 25+ 默认值 + Viper->>YAML: ReadInConfig() + YAML-->>Viper: 配置内容 + Viper->>Viper: bindEnvVars(v) + Note over Viper: 绑定 25+ 环境变量 + Viper->>ENV: 读取环境变量 + ENV-->>Viper: 覆盖对应字段 + Viper->>Cfg: Unmarshal(&cfg) + Cfg-->>Main: 返回完整配置 +``` + +关键步骤: + +1. **加载 `.env` 文件** — `godotenv.Load()` 从项目根目录读取 `.env` +2. **设置默认值** — `setDefaults(v)` 为所有字段提供合理的默认值 +3. **读取 YAML** — `v.ReadInConfig()` 从 `internal/config/config.yaml` 加载 +4. **绑定环境变量** — `bindEnvVars(v)` 将每个字段映射到 `GEN2D_*` 环境变量 +5. **反序列化** — `v.Unmarshal(&cfg)` 将合并后的配置映射到 Go 结构体 + +> **容错设计**:YAML 文件不存在时不报错,仅使用环境变量 + 默认值。这确保了零配置即可启动。 + +--- + +## 与部署的关系 + +配置级联机制与部署架构紧密配合: + +| 部署方式 | 配置策略 | 示例 | +|----------|----------|------| +| `go run` 本地开发 | YAML + `.env` | `config.yaml` 配置 DB,`.env` 配置 API Key | +| `docker compose` | `.env` 文件注入 | `env_file: ./backend/.env` | +| Kubernetes | ConfigMap + Secret | `GEN2D_DSN` 从 Secret 注入 | + +--- + +## 关联文档 + +| 文档 | 关系 | +|------|------| +| [14 - 部署架构](14-deployment.md) | 配置管理在部署中的应用 | +| [01 - 系统总览](01-system-overview.md) | 配置在启动流程中的位置 | +| [02 - 协程池](02-worker-pool.md) | WorkerPoolConfig 控制并发参数 | +| [00 - 索引](00-index.md) | 返回文档总览 |