diff --git a/backend/docs/diagrams/product-domains.json b/backend/docs/diagrams/product-domains.json new file mode 100644 index 0000000..01ee6ae --- /dev/null +++ b/backend/docs/diagrams/product-domains.json @@ -0,0 +1,270 @@ +{ + "schema_version": 1, + "mode": "architecture", + "template_type": "architecture", + "style": 1, + "quality_profile": "standard", + "width": 1140, + "height": 760, + "title": "Windup Product Domains", + "subtitle": "最终产品业务能力与核心依赖", + "containers": [ + { + "id": "foundation-domains", + "x": 40, + "y": 120, + "width": 1120, + "height": 120, + "label": "Foundation Domains", + "stroke": "#b9c6d8", + "fill": "#f4f8fc" + }, + { + "id": "workflow-domains", + "x": 40, + "y": 290, + "width": 680, + "height": 160, + "label": "Workflow Domains", + "stroke": "#c9b8d9", + "fill": "#faf5fc" + }, + { + "id": "ai-engine-domain", + "x": 900, + "y": 290, + "width": 200, + "height": 160, + "label": "AI Engine Domain", + "stroke": "#89c4a8", + "fill": "#f0f9f4" + }, + { + "id": "result-domains", + "x": 40, + "y": 490, + "width": 1120, + "height": 120, + "label": "Result Domains", + "stroke": "#b8d0c2", + "fill": "#f5fbf7" + } + ], + "nodes": [ + { + "id": "user", + "kind": "rect", + "x": 80, + "y": 155, + "width": 150, + "height": 55, + "label": "user", + "sublabel": "身份与会话", + "fill": "#e8f0fb", + "stroke": "#6685b2" + }, + { + "id": "quota", + "kind": "rect", + "x": 270, + "y": 155, + "width": 150, + "height": 55, + "label": "quota", + "sublabel": "订阅与积分", + "fill": "#fff2d9", + "stroke": "#c28a36" + }, + { + "id": "project", + "kind": "double_rect", + "x": 460, + "y": 155, + "width": 170, + "height": 55, + "label": "project", + "sublabel": "全局约束容器", + "fill": "#e8f0fb", + "stroke": "#6685b2" + }, + { + "id": "media", + "kind": "rect", + "x": 670, + "y": 155, + "width": 150, + "height": 55, + "label": "media", + "sublabel": "输入参考素材", + "fill": "#e8f0fb", + "stroke": "#6685b2" + }, + { + "id": "character", + "kind": "rect", + "x": 860, + "y": 155, + "width": 180, + "height": 55, + "label": "character / asset", + "sublabel": "角色与产物", + "fill": "#e8f0fb", + "stroke": "#6685b2" + }, + { + "id": "agent", + "kind": "rect", + "x": 70, + "y": 335, + "width": 160, + "height": 65, + "label": "agent", + "sublabel": "懒人智能体", + "fill": "#f2e9f8", + "stroke": "#89659e" + }, + { + "id": "workflow", + "kind": "double_rect", + "x": 270, + "y": 335, + "width": 160, + "height": 65, + "label": "workflow", + "sublabel": "工作流画布", + "fill": "#f2e9f8", + "stroke": "#89659e" + }, + { + "id": "orchestrator", + "kind": "rect", + "x": 510, + "y": 335, + "width": 180, + "height": 65, + "label": "orchestrator", + "sublabel": "生成任务调度", + "fill": "#fff2d9", + "stroke": "#c28a36" + }, + { + "id": "ai-engine", + "kind": "rect", + "x": 930, + "y": 335, + "width": 140, + "height": 65, + "label": "ai_engine", + "sublabel": "模型能力适配", + "fill": "#e4f2ea", + "stroke": "#5c9974" + }, + { + "id": "review", + "kind": "rect", + "x": 230, + "y": 525, + "width": 180, + "height": 55, + "label": "review", + "sublabel": "质检 + 人工审核", + "fill": "#e4f2ea", + "stroke": "#5c9974" + }, + { + "id": "playtest", + "kind": "rect", + "x": 510, + "y": 525, + "width": 180, + "height": 55, + "label": "playtest", + "sublabel": "预览与试玩", + "fill": "#e4f2ea", + "stroke": "#5c9974" + }, + { + "id": "export", + "kind": "rect", + "x": 790, + "y": 525, + "width": 180, + "height": 55, + "label": "export", + "sublabel": "格式转换与下载", + "fill": "#e4f2ea", + "stroke": "#5c9974" + } + ], + "arrows": [ + { + "id": "agent-workflow", + "source": "agent", + "target": "workflow", + "source_port": "right", + "target_port": "left", + "flow": "control", + "label": "", + "label_style": "offset" + }, + { + "id": "workflow-orchestrator", + "source": "workflow", + "target": "orchestrator", + "source_port": "right", + "target_port": "left", + "flow": "control", + "label": "", + "label_style": "offset" + }, + { + "id": "orchestrator-ai", + "source": "orchestrator", + "target": "ai-engine", + "source_port": "right", + "target_port": "left", + "flow": "control", + "label": "", + "label_style": "offset" + }, + { + "id": "review-playtest", + "source": "review", + "target": "playtest", + "source_port": "right", + "target_port": "left", + "flow": "read", + "label": "", + "label_style": "offset" + }, + { + "id": "playtest-export", + "source": "playtest", + "target": "export", + "source_port": "right", + "target_port": "left", + "flow": "read", + "label": "", + "label_style": "offset" + } + ], + "legend_orientation": "horizontal", + "legend_x": 70, + "legend_y": 680, + "legend_locked": false, + "legend": [ + { + "flow": "control", + "label": "业务编排" + }, + { + "flow": "read", + "label": "读取 / 消费" + }, + { + "flow": "write", + "label": "产物 / 输入流" + } + ], + "footer": "Windup · Final Product Domain Map" +} diff --git a/backend/docs/diagrams/product-domains.svg b/backend/docs/diagrams/product-domains.svg new file mode 100644 index 0000000..1444299 --- /dev/null +++ b/backend/docs/diagrams/product-domains.svg @@ -0,0 +1,148 @@ + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + Windup Product Domains + 最终产品业务能力与核心依赖 + + + FOUNDATION DOMAINS + + + + + WORKFLOW DOMAINS + + + + + AI ENGINE DOMAIN + + + + + RESULT DOMAINS + + + + + + + + + + user + 身份与会话 + + + + quota + 订阅与积分 + + + + + project + 全局约束容器 + + + + media + 输入参考素材 + + + + character / asset + 角色与产物 + + + + agent + 懒人智能体 + + + + + workflow + 工作流画布 + + + + orchestrator + 生成任务调度 + + + + ai_engine + 模型能力适配 + + + + review + 质检 + 人工审核 + + + + playtest + 预览与试玩 + + + + export + 格式转换与下载 + + + + + 业务编排 + + 读取 / 消费 + + 产物 / 输入流 + + + Windup · Final Product Domain Map + + \ No newline at end of file diff --git a/backend/docs/module-split.md b/backend/docs/module-split.md new file mode 100644 index 0000000..a741196 --- /dev/null +++ b/backend/docs/module-split.md @@ -0,0 +1,270 @@ +# 后端模块拆分 + +> 当前阶段:MVP 已实现 project / media / character / orchestrator / ai_engine。 +> 按业务域拆分,每个域是一个完整的业务能力单元。 + +--- + +## 项目层级 + +``` +backend/ +├── packages/ +│ ├── common/ # 共享:Response、BizException、BizCode 枚举 +│ ├── framework/ # 基础设施:KodoStorage、ChatProvider、DB 配置 +│ ├── ai_engine/ # 生成管线(独立包) +│ └── app/ # 业务应用 +│ └── src/windup_app/ +│ ├── web/api/ # FastAPI 路由 +│ │ ├── generation.py # 生成 API(前端契约) +│ │ ├── media.py # 媒体 API +│ │ ├── project.py # 项目 API +│ │ └── character.py # 角色 API +│ ├── bootstrap/app.py # 组装入口(composition root) +│ └── server/ # 业务域(按域分组) +│ ├── media/ # [foundation] 用户素材上传 +│ ├── project/ # [foundation] 项目约束配置 +│ ├── character/ # [foundation] 角色资产数据 +│ ├── orchestrator/ # [workflow] 生成任务调度 +│ ├── workflow/ # [workflow] 工作流画布(功能型卡片) +│ └── agent/ # [workflow] Agent 智能体(懒人模式) +``` + +> 注:`[foundation]` / `[workflow]` 标注所属业务域。 +> `ai_engine` 是独立包,不属于 server/ 目录。 + +--- + +## 业务域全景 + +![Windup 最终产品业务域图](diagrams/product-domains.svg) + +--- + +## 依赖方向 + +``` + common + ▲ + framework + ▲ + ┌──────┴──────┐ +ai_engine foundation + ▲ ▲ ▲ + │ ┌────┘ │ + │ workflow result + │ │ │ │ │ + │ │ │ │ │ + └───┘──┘──┘─────┘ + │ │ + orchestrator + │ + agent +``` + +**模块关系**: +- **workflow**(工作流画布):卡片编排,调用 orchestrator 触动生成 +- **orchestrator**(生成任务调度):管理生成任务生命周期,调用 ai_engine +- **agent**(懒人智能体):一句话入口,内部调用 workflow API +- **ai_engine**(生成管线):实际 AI 生成,被 orchestrator 调用 + +**规则**: +- foundation → framework, common +- workflow → orchestrator, foundation, framework, common +- orchestrator → ai_engine, foundation, framework, common +- agent → workflow, framework, common +- result → foundation, framework, common +- ai_engine → framework.providers(接口), common +- **禁止**:foundation → workflow / orchestrator / agent / result / ai_engine +- **禁止**:ai_engine → foundation / workflow / orchestrator / agent / result +- **禁止**:result → workflow / orchestrator / agent + +--- + +## 域边界划分 + +### 为什么按这四个域拆分? + +每个域回答一个业务问题: + +| 域 | 回答的问题 | 包含什么 | +|---|---|---| +| foundation | 数据从哪来、存到哪 | 项目配置、角色数据、用户素材 | +| workflow | 何时生、为谁生、生完怎么办 | 任务调度、工作流编排 | +| pipeline | 怎么生 | 提示词 → 出图 → 抠图 → 截帧 | +| result | 生出来的东西怎么用 | 审核、预览、导出 | + +### 边界职责 + +**foundation(基础业务域)** +- 职责:管理基础数据,被其他域消费 +- 边界:各模块独立 CRUD,不包含生成、审核、导出逻辑 +- 禁止:依赖 workflow / result / ai_engine + +**workflow(工作流域)** +- 职责:编排调度,调用 foundation + ai_engine +- 边界:任务管理、约束加载、积分扣减、结果上传 +- 禁止:依赖 result 域;直接 import ai_engine 内部实现 + +**pipeline(生成管线域)** +- 职责:实际生成过程(提示词 → 出图 → 抠图 → 截帧) +- 边界:AI 调用、帧处理、像素化 +- 禁止:知道 project 约束、character 数据结构、图片存储位置 + +**result(结果域)** +- 职责:处理生成产物,提供预览和导出 +- 边界:读取正式资产,独立演进 +- 禁止:依赖 workflow 域;调用 ai_engine + +--- + +## 基础业务域(Foundation) + +### media — 用户素材上传 + +**对应表:** `windup_media` **接口:** `MediaService` + +| 方法 | 说明 | +|---|---| +| `upload(data, metadata)` | 上传文件到对象存储,返回 URL | + +**文件分类 `MediaCategory`**:`reference-image` / `outfit-preview` / `action-frame` / `general` + +### project — 项目约束配置 + +**对应表:** `windup_project` **接口:** `ProjectService` + +| 字段 | 用途 | +|---|---| +| `character_perspective` | 视角 → 生成朝向 | +| `directional_movement` | 方向数 → 生成方向变体 | +| `sprite_width` / `sprite_height` | 尺寸 → 输出帧大小 | +| `game_style` | 画风 → 提示词风格 | +| `sprite_sample_url` | 风格参考图 → 图生图模式 | + +| 方法 | 说明 | +|---|---| +| `create_project(project)` | 创建项目 | +| `get_project(id)` | 按 ID 查询 | +| `list_projects(page, page_size, user_id)` | 分页查询 | +| `delete_project(id)` | 删除 | + +### character — 角色资产数据 + +**对应表:** `windup_character` **接口:** `CharacterService` + +`character_data` JSONB 三层嵌套:outfit → action → frame + +| 方法 | 说明 | +|---|---| +| `create_character(session, **fields)` | 创建角色 | +| `get_character(session, character_id)` | 按 ID 查询 | +| `list_characters(session, *, project_id, page, page_size)` | 分页查询 | +| `update_character(session, character_id, **fields)` | 更新角色 | +| `delete_character(session, character_id)` | 删除角色(含媒体清理) | + +--- + +## 工作流域(Workflow) + +### orchestrator — 生成任务调度 + +**对应表:** `windup_generation_task` **接口:** `GenerationService` + +管理生成任务生命周期:创建任务 → 加载项目约束 → 调 ai_engine → 上传结果 → 回写状态。 + +| 方法 | 说明 | +|---|---| +| `generate_character_image(input)` | 提交角色图片生成任务 | +| `generate_character_action(input)` | 提交角色动作生成任务 | +| `get_task(project_id, task_id)` | 查询任务状态与结果 | + +### workflow — 工作流画布 + +**对应表:** `windup_workflow` / `windup_canvas_card` / `windup_generation_attempt` +**接口:** `WorkflowService` / `CardService` + +功能型卡片体系:CHARACTER(角色实体)→ CANDIDATE(母版候选)/ ACTION(角色动作)/ EXPORT(资产导出)。 + +**卡片类型 `CardType`:** character / candidate / action / export + +**WorkflowService 接口:** + +| 方法 | 说明 | +|---|---| +| `create_workflow(user_id, project_id, name)` | 创建工作流(自动创建 CHARACTER 根卡片) | +| `get_workflow(workflow_id)` | 获取工作流详情(含全部 active 卡片) | +| `delete_workflow(workflow_id)` | 删除工作流(级联软删除) | + +**CardService 接口:** + +| 方法 | 说明 | +|---|---| +| `create_card(workflow_id, card_type, parent_card_id, ...)` | 创建子卡片(ACTION/EXPORT) | +| `confirm_card(card_id, user_input, spec_overrides?)` | 确认卡片,触发生成 | +| `regenerate_card(card_id, user_input?)` | 重新生成(新 attempt) | +| `delete_card(card_id)` | 删除卡片(级联软删除) | + +**SSE 衔接:** 确认卡片返回 `GenerationAttempt`(含 `task_id`),前端订阅 Task SSE 获取进度。 + +### agent — Agent 智能体 + +**对应表:** `agent_session` / `agent_message` / `agent_tool_call` +**接口:** `AgentService` + +懒人智能体:一句话搞定角色资产生成。Agent 内部调用 Workflow API 完成实际操作。 + +**工具定义:** create_project / create_character / generate_character_image / generate_character_action + +**AgentService 接口:** + +| 方法 | 说明 | +|---|---| +| `create_session(user_id, context?)` | 创建 Agent 会话 | +| `send_message(session_id, content)` | 发送用户消息 | +| `send_choice(session_id, message_id, value)` | 发送用户选择 | +| `get_messages(session_id, limit?, before?)` | 获取会话历史 | + +**SSE 事件:** message / tool_call / tool_result / state_change / error + +--- + +## 生成管线域(Pipeline) + +### ai_engine — 生成管线 + +**接口:** `CharacterGeneratorPort` + +生成管线:提示词 → 出图 → 抠图 → 截帧 → 返回产物。 + +| 组件 | 职责 | +|---|---| +| `strategy/` | 策略分发(VIDEO_I2V / PER_FRAME / PROC_IDLE) | +| `prompt/` | 提示词构建(walk / jump / attack) | +| `slicing/` | 帧提取(imageio/pyav) | +| `postprocess/` | 像素化、脚线对齐、sprite sheet 打包 | +| `master_prep.py` | 母版预处理 | + +--- + +## 结果域(Result) + +| 模块 | 状态 | 职责 | +|---|---|---| +| review | 🟡 前端页面体现 | 质检 + 人工审核 | +| preview | 🟡 前端页面体现 | 预览台:组装可播放数据(帧 + 帧率 + 循环) | +| export | ⬜ 待实现 | GIF / 精灵图 / 引擎格式转换 | + +--- + +## MVP 已实现模块 + +| 模块 | 域 | 数据表 | API | +|---|---|---|---| +| media | foundation | windup_media | POST /media/upload | +| project | foundation | windup_project | POST/GET/DELETE /projects | +| character | foundation | windup_character | POST/GET/PATCH/DELETE /characters | +| orchestrator | workflow | windup_generation_task | POST /generation/image, POST /generation/action, GET /generation/tasks/{id} | +| workflow | workflow | windup_workflow, windup_canvas_card, windup_generation_attempt | POST/GET/DELETE /workflow, POST/PATCH/DELETE /workflow/{id}/cards, POST /cards/{id}/confirm, POST /cards/{id}/regenerate | +| agent | workflow | agent_session, agent_message, agent_tool_call | POST /agent/sessions, POST /agent/sessions/{id}/messages, POST /agent/sessions/{id}/choices, GET /agent/sessions/{id}/messages | +| ai_engine | pipeline | (无独立表) | (内部调用,不暴露 API) | diff --git a/backend/packages/app/src/windup_app/bootstrap/app.py b/backend/packages/app/src/windup_app/bootstrap/app.py index 89f7b43..8a8a33d 100644 --- a/backend/packages/app/src/windup_app/bootstrap/app.py +++ b/backend/packages/app/src/windup_app/bootstrap/app.py @@ -6,10 +6,14 @@ from fastapi import FastAPI +from windup_app.web.api.agent import router as agent_router from windup_app.web.api.media import router as media_router +from windup_app.web.api.workflow import router as workflow_router def create_app() -> FastAPI: app = FastAPI(title="windup", version="0.1.0") app.include_router(media_router) + app.include_router(workflow_router) + app.include_router(agent_router) return app diff --git a/backend/packages/app/src/windup_app/server/agent/__init__.py b/backend/packages/app/src/windup_app/server/agent/__init__.py new file mode 100644 index 0000000..dc087b7 --- /dev/null +++ b/backend/packages/app/src/windup_app/server/agent/__init__.py @@ -0,0 +1 @@ +"""Agent 智能体领域。""" diff --git a/backend/packages/app/src/windup_app/server/agent/interface.py b/backend/packages/app/src/windup_app/server/agent/interface.py new file mode 100644 index 0000000..7f829fa --- /dev/null +++ b/backend/packages/app/src/windup_app/server/agent/interface.py @@ -0,0 +1,91 @@ +"""Agent 智能体领域服务接口。 + +懒人智能体:一句话搞定角色资产生成。 +用户通过自然语言与 Agent 对话,Agent 自动调用工具完成项目创建、角色生成、动作生成等操作。 + +调用流程 +-------- +1. 前端调用 ``POST /agent/sessions`` 创建会话,拿到 ``session_id``。 +2. 前端订阅 ``GET /agent/sessions/{session_id}/stream`` 获取 Agent 事件流。 +3. 前端调用 ``POST /agent/sessions/{session_id}/messages`` 发送用户消息。 +4. Agent 通过 SSE 推送处理结果(message/tool_call/tool_result 等事件)。 +5. 若 tool_call 含 task_id,前端订阅 Task SSE 获取生成进度。 + +工具定义 +-------- +Agent 可调用的工具(通过 SSE tool_call 事件推送): + +- ``create_project``: 创建项目(返回 project_id) +- ``create_character``: 创建角色(返回 character_id) +- ``generate_character_image``: 生成角色图片(返回 task_id,可订阅 Task SSE) +- ``generate_character_action``: 生成角色动作(返回 task_id,可订阅 Task SSE) +""" + +from __future__ import annotations + +from abc import ABC, abstractmethod + +from windup_app.server.agent.model import ( + AgentMessage, + AgentSession, + ToolCall, +) + + +class AgentService(ABC): + """Agent 用例的抽象边界。""" + + # -- 会话管理 ---------------------------------------------------------- + + @abstractmethod + def create_session(self, *, user_id: int, context: dict | None = None) -> AgentSession: + """创建 Agent 会话。 + + 返回的 session_id 用于 SSE 订阅和消息发送。 + """ + + @abstractmethod + def get_session(self, session_id: str) -> AgentSession | None: + """获取会话信息。""" + + @abstractmethod + def close_session(self, session_id: str) -> None: + """关闭会话。""" + + # -- 消息交互 ---------------------------------------------------------- + + @abstractmethod + def send_message(self, session_id: str, *, content: str, message_id: str | None = None) -> AgentMessage: + """发送用户消息。 + + Agent 收到消息后异步处理,通过 SSE 推送结果。 + 返回用户消息记录。 + """ + + @abstractmethod + def send_choice(self, session_id: str, *, message_id: str, value: str) -> None: + """发送用户选择(按钮点击)。 + + 对应 Agent 消息中的 buttons 块。 + """ + + @abstractmethod + def get_messages( + self, + session_id: str, + *, + limit: int = 50, + before: str | None = None, + ) -> list[AgentMessage]: + """获取会话历史消息。 + + 参数: + - limit: 返回条数,默认50 + - before: 分页,此 message_id 之前的消息 + """ + + # -- 工具调用记录 ------------------------------------------------------ + + @abstractmethod + def get_tool_calls(self, session_id: str) -> list[ToolCall]: + """获取会话的全部工具调用记录。""" diff --git a/backend/packages/app/src/windup_app/server/agent/model.py b/backend/packages/app/src/windup_app/server/agent/model.py new file mode 100644 index 0000000..baa182f --- /dev/null +++ b/backend/packages/app/src/windup_app/server/agent/model.py @@ -0,0 +1,121 @@ +"""Agent 智能体领域模型。 + +懒人智能体:一句话搞定角色资产生成。 +用户通过自然语言与 Agent 对话,Agent 自动调用工具完成项目创建、角色生成、动作生成等操作。 +""" + +from __future__ import annotations + +from dataclasses import dataclass, field +from datetime import datetime, timezone +from enum import StrEnum + + +# -- 枚举 ---------------------------------------------------------------- + + +class SessionStatus(StrEnum): + """会话状态。""" + + ACTIVE = "active" + CLOSED = "closed" + + +class MessageRole(StrEnum): + """消息角色。""" + + USER = "user" + ASSISTANT = "assistant" + + +class AgentState(StrEnum): + """Agent 状态(仅 assistant 消息)。""" + + PROCESSING = "processing" # Agent 正在处理 + WAITING_INPUT = "waiting_input" # 等待用户输入 + DONE = "done" # 处理完成 + + +class ToolName(StrEnum): + """已定义的工具名称。""" + + CREATE_PROJECT = "create_project" + CREATE_CHARACTER = "create_character" + GENERATE_CHARACTER_IMAGE = "generate_character_image" + GENERATE_CHARACTER_ACTION = "generate_character_action" + + +# -- 会话 --------------------------------------------------------------- + + +@dataclass +class AgentSession: + """Agent 会话。""" + + id: int | None = None + session_id: str = "" + user_id: int = 0 + project_id: int | None = None + status: SessionStatus = SessionStatus.ACTIVE + context: dict = field(default_factory=dict) + create_at: datetime = field(default_factory=lambda: datetime.now(timezone.utc)) + update_at: datetime = field(default_factory=lambda: datetime.now(timezone.utc)) + + +# -- 消息 --------------------------------------------------------------- + + +@dataclass +class ContentBlock: + """内容块——Agent 消息的基本单元。""" + + type: str = "" # text / image / buttons / confirm / progress / divider / markdown / code + # text / markdown / code + content: str = "" + # image + url: str = "" + caption: str = "" + # buttons / confirm + prompt: str = "" + options: list[dict] = field(default_factory=list) # [{value, label, icon?}] + confirm_label: str = "确认" + reject_label: str = "取消" + # progress + stage: str = "" + current: int = 0 + total: int = 0 + text: str = "" + # code + language: str = "" + + +@dataclass +class AgentMessage: + """Agent 消息。""" + + id: int | None = None + session_id: str = "" + message_id: str = "" + role: MessageRole = MessageRole.USER + blocks: list[ContentBlock] = field(default_factory=list) + state: AgentState | None = None + create_at: datetime = field(default_factory=lambda: datetime.now(timezone.utc)) + + +# -- 工具调用 ----------------------------------------------------------- + + +@dataclass +class ToolCall: + """Agent 工具调用。""" + + id: int | None = None + session_id: str = "" + message_id: str = "" + call_id: str = "" + tool_name: ToolName = ToolName.CREATE_PROJECT + arguments: dict = field(default_factory=dict) + result: dict | None = None + error: str | None = None + task_id: int | None = None + create_at: datetime = field(default_factory=lambda: datetime.now(timezone.utc)) diff --git a/backend/packages/app/src/windup_app/server/agent/schema.py b/backend/packages/app/src/windup_app/server/agent/schema.py new file mode 100644 index 0000000..f8da68c --- /dev/null +++ b/backend/packages/app/src/windup_app/server/agent/schema.py @@ -0,0 +1,173 @@ +"""Agent 智能体 API Schema。 + +定义前端请求/响应的 Pydantic 模型,与 server 层解耦。 +前端团队参考此文件了解接口契约。 +""" + +from __future__ import annotations + +from pydantic import BaseModel, ConfigDict, Field + + +# ══════════════════════════════════════════════════════════════════════════════ +# 请求模型 +# ══════════════════════════════════════════════════════════════════════════════ + + +class SessionCreateRequest(BaseModel): + """创建 Agent 会话。""" + + user_id: int = Field(description="用户 ID") + context: dict | None = Field( + default=None, + description="初始上下文,如 {project_id: 123}。可选,后续通过对话补充。", + ) + + +class MessageSendRequest(BaseModel): + """发送用户消息。""" + + content: str = Field(min_length=1, max_length=2000, description="用户消息内容") + message_id: str | None = Field( + default=None, + description="客户端消息 ID(用于去重),省略则由后端生成", + ) + + +class ChoiceSendRequest(BaseModel): + """发送用户选择(按钮点击)。""" + + message_id: str = Field(description="对应的 Agent 消息 ID(buttons 块所在的 message)") + value: str = Field(description="选择的值(ButtonOption.value)") + + +# ══════════════════════════════════════════════════════════════════════════════ +# 响应模型 +# ══════════════════════════════════════════════════════════════════════════════ + + +class SessionOut(BaseModel): + """会话响应。""" + + model_config = ConfigDict(from_attributes=True) + + session_id: str = Field(description="会话 ID,用于后续 SSE 订阅和消息发送") + created_at: str = Field(description="创建时间(ISO 8601)") + + +class ContentBlockOut(BaseModel): + """内容块响应。""" + + type: str = Field(description="块类型:text/image/buttons/confirm/progress/divider/markdown/code") + # text / markdown / code + content: str | None = None + # image + url: str | None = None + caption: str | None = None + # buttons + prompt: str | None = None + options: list[dict] | None = None # [{value, label, icon?}] + # confirm + confirm_label: str | None = None + reject_label: str | None = None + # progress + stage: str | None = None + current: int | None = None + total: int | None = None + text: str | None = None + # code + language: str | None = None + + +class MessageOut(BaseModel): + """消息响应。""" + + model_config = ConfigDict(from_attributes=True) + + message_id: str = Field(description="消息 ID") + session_id: str = Field(description="会话 ID") + role: str = Field(description="角色:user / assistant") + blocks: list[ContentBlockOut] = Field( + default_factory=list, + description="内容块列表", + ) + state: str | None = Field( + default=None, + description="Agent 状态:processing / waiting_input / done(仅 assistant 消息)", + ) + timestamp: float = Field(description="时间戳(秒)") + + +class ToolCallOut(BaseModel): + """工具调用响应。""" + + model_config = ConfigDict(from_attributes=True) + + call_id: str = Field(description="调用 ID") + tool: str = Field(description="工具名称") + args: dict = Field(default_factory=dict, description="工具参数") + task_id: int | None = Field(default=None, description="任务 ID(生成任务时有值,可订阅 Task SSE)") + message: str | None = Field(default=None, description="说明文字") + + +class ToolResultOut(BaseModel): + """工具结果响应。""" + + model_config = ConfigDict(from_attributes=True) + + call_id: str = Field(description="关联的 tool_call ID") + tool: str = Field(description="工具名称") + result: dict | None = Field(default=None, description="工具返回结果") + error: str | None = Field(default=None, description="错误信息(失败时)") + + +# ══════════════════════════════════════════════════════════════════════════════ +# SSE 事件结构(前端参考) +# ══════════════════════════════════════════════════════════════════════════════ +# +# ── Agent SSE 事件类型 ── +# +# event: message +# data: { +# "message_id": "msg_001", +# "session_id": "session_abc123", +# "role": "assistant", +# "blocks": [ +# {"type": "text", "content": "好的,我需要了解几个细节:"}, +# {"type": "buttons", "prompt": "游戏类型", "options": [ +# {"value": "side_scroller", "label": "横版游戏"}, +# {"value": "top_down", "label": "俯视角"} +# ]} +# ], +# "state": "waiting_input", +# "timestamp": 1234567890.123 +# } +# +# event: tool_call +# data: { +# "call_id": "c1", +# "tool": "create_project", +# "args": {"name": "甲壳虫", "perspective": 1}, +# "task_id": null, +# "message": "正在创建项目..." +# } +# +# event: tool_result +# data: { +# "call_id": "c1", +# "tool": "create_project", +# "result": {"project_id": 123}, +# "error": null +# } +# +# event: state_change +# data: { +# "state": "waiting_input", +# "previous_state": "processing" +# } +# +# event: error +# data: { +# "error": "会话已过期", +# "code": "SESSION_EXPIRED" +# } diff --git a/backend/packages/app/src/windup_app/server/generation/__init__.py b/backend/packages/app/src/windup_app/server/orchestrator/__init__.py similarity index 89% rename from backend/packages/app/src/windup_app/server/generation/__init__.py rename to backend/packages/app/src/windup_app/server/orchestrator/__init__.py index a6701de..c21a7b8 100644 --- a/backend/packages/app/src/windup_app/server/generation/__init__.py +++ b/backend/packages/app/src/windup_app/server/orchestrator/__init__.py @@ -1,6 +1,6 @@ """生成任务领域。""" -from windup_app.server.generation.model import ( +from windup_app.server.orchestrator.model import ( ActionType, CharacterActionFrame, CharacterActionInput, diff --git a/backend/packages/app/src/windup_app/server/generation/interface.py b/backend/packages/app/src/windup_app/server/orchestrator/interface.py similarity index 97% rename from backend/packages/app/src/windup_app/server/generation/interface.py rename to backend/packages/app/src/windup_app/server/orchestrator/interface.py index b43bace..83e38aa 100644 --- a/backend/packages/app/src/windup_app/server/generation/interface.py +++ b/backend/packages/app/src/windup_app/server/orchestrator/interface.py @@ -22,7 +22,7 @@ from abc import ABC, abstractmethod -from windup_app.server.generation.model import ( +from windup_app.server.orchestrator.model import ( CharacterActionInput, CharacterImageInput, GenerationTask, diff --git a/backend/packages/app/src/windup_app/server/generation/model.py b/backend/packages/app/src/windup_app/server/orchestrator/model.py similarity index 100% rename from backend/packages/app/src/windup_app/server/generation/model.py rename to backend/packages/app/src/windup_app/server/orchestrator/model.py diff --git a/backend/packages/app/src/windup_app/server/workflow/__init__.py b/backend/packages/app/src/windup_app/server/workflow/__init__.py new file mode 100644 index 0000000..cc21c20 --- /dev/null +++ b/backend/packages/app/src/windup_app/server/workflow/__init__.py @@ -0,0 +1 @@ +"""工作流画布领域。""" diff --git a/backend/packages/app/src/windup_app/server/workflow/interface.py b/backend/packages/app/src/windup_app/server/workflow/interface.py new file mode 100644 index 0000000..f343406 --- /dev/null +++ b/backend/packages/app/src/windup_app/server/workflow/interface.py @@ -0,0 +1,134 @@ +"""工作流画布领域服务接口。 + +API 层只依赖本模块定义的抽象。具体实现在应用装配层继承后通过依赖注入提供。 + +调用流程 +-------- +1. 前端调用 ``POST /workflow`` 创建工作流,拿到 ``workflow_id``。 +2. 前端调用 ``POST /workflow/{id}/cards`` 创建子卡片(ACTION / EXPORT)。 +3. 前端调用 ``POST /cards/{id}/confirm`` 确认卡片,触发生成。 + - 返回的 ``GenerationAttempt`` 包含 ``task_id``,前端可订阅 Task SSE 获取进度。 +4. 前端通过 ``GET /workflow/{id}`` 获取画布最新状态。 +5. 生成完成后,前端从 ``card.latest_result`` 取出结果。 + +回退流程 +-------- +前端调用 ``POST /cards/{id}/regenerate`` 触发重新生成。 +- CHARACTER:旧 CANDIDATE 全部 INACTIVE,重新生成候选。ACTION / EXPORT 不受影响。 +- ACTION / EXPORT:创建新 attempt,重新执行。 +""" + +from __future__ import annotations + +from abc import ABC, abstractmethod + +from windup_app.server.workflow.model import ( + CanvasCard, + GenerationAttempt, + Workflow, +) + + +class WorkflowService(ABC): + """工作流用例的抽象边界。""" + + # -- 工作流 CRUD -------------------------------------------------------- + + @abstractmethod + def create_workflow(self, *, user_id: int, project_id: int, name: str) -> Workflow: + """创建工作流,自动创建 CHARACTER 根卡片。 + + 返回的工作流已包含一张 DRAFT 状态的 CHARACTER 卡片。 + """ + + @abstractmethod + def get_workflow(self, workflow_id: int) -> Workflow | None: + """获取工作流详情(含全部 active 卡片)。 + + 返回的 Workflow.cards 已按 parent_card_id 组装为树结构。 + """ + + @abstractmethod + def delete_workflow(self, workflow_id: int) -> None: + """删除工作流(级联软删除所有卡片和尝试记录)。""" + + +class CardService(ABC): + """卡片用例的抽象边界。""" + + # -- 卡片 CRUD ---------------------------------------------------------- + + @abstractmethod + def create_card( + self, + workflow_id: int, + *, + card_type: str, + parent_card_id: int, + direction: str | None = None, + user_input: dict | None = None, + spec_overrides: dict | None = None, + ) -> CanvasCard: + """创建子卡片(ACTION / EXPORT)。 + + ACTION 卡片创建时自动从选定的 CANDIDATE 复制母版图到 user_input.master_image_url。 + """ + + @abstractmethod + def update_card( + self, + card_id: int, + *, + user_input: dict | None = None, + position_x: float | None = None, + position_y: float | None = None, + ) -> CanvasCard: + """更新卡片用户输入或位置(不触发生成)。""" + + @abstractmethod + def confirm_card( + self, + card_id: int, + *, + user_input: dict, + spec_overrides: dict | None = None, + ) -> GenerationAttempt: + """确认卡片,触发生成。 + + - CHARACTER:生成候选图 → 创建 CANDIDATE 卡片。 + - ACTION:生成动画帧。 + - EXPORT:打包导出。 + + 返回的 GenerationAttempt 包含 task_id,前端可订阅 Task SSE。 + """ + + @abstractmethod + def regenerate_card( + self, + card_id: int, + *, + user_input: dict | None = None, + ) -> GenerationAttempt: + """重新生成(创建新的 GenerationAttempt)。 + + - CHARACTER:旧 CANDIDATE 全部 INACTION,重新生成候选。ACTION/EXPORT 不受影响。 + - ACTION / EXPORT:创建新 attempt,重新执行。 + """ + + @abstractmethod + def delete_card(self, card_id: int) -> None: + """删除卡片(级联软删除所有子卡片)。""" + + @abstractmethod + def get_card(self, card_id: int) -> CanvasCard | None: + """获取单张卡片。""" + + @abstractmethod + def list_cards(self, workflow_id: int) -> list[CanvasCard]: + """获取工作流下全部 active 卡片。""" + + # -- 生成尝试 ----------------------------------------------------------- + + @abstractmethod + def get_attempts(self, card_id: int) -> list[GenerationAttempt]: + """获取某张卡片的全部生成尝试记录。""" diff --git a/backend/packages/app/src/windup_app/server/workflow/model.py b/backend/packages/app/src/windup_app/server/workflow/model.py new file mode 100644 index 0000000..dde386c --- /dev/null +++ b/backend/packages/app/src/windup_app/server/workflow/model.py @@ -0,0 +1,131 @@ +"""工作流画布领域模型。 + +功能型卡片体系:CHARACTER(角色实体)→ CANDIDATE(母版候选)/ ACTION(角色动作)/ EXPORT(资产导出)。 +前端通过卡片 API 统一操作,后端负责持久化画布结构和生成任务状态。 +""" + +from __future__ import annotations + +from dataclasses import dataclass, field +from datetime import datetime, timezone +from enum import StrEnum + + +# -- 枚举 ---------------------------------------------------------------- + + +class WorkflowStatus(StrEnum): + """工作流状态。""" + + ACTIVE = "active" + ARCHIVED = "archived" + + +class CardType(StrEnum): + """卡片类型——功能型,按职责划分。""" + + CHARACTER = "character" # 角色实体根节点 + CANDIDATE = "candidate" # 母版候选 + ACTION = "action" # 角色动作 + EXPORT = "export" # 资产导出 + + +class CardStatus(StrEnum): + """卡片状态。""" + + DRAFT = "draft" # 已创建,待用户填写 + GENERATING = "generating" # 生成中 + COMPLETED = "completed" # 完成 + FAILED = "failed" # 失败 + INACTIVE = "inactive" # 已失效(软删除) + + +class AttemptStatus(StrEnum): + """生成尝试状态。""" + + PENDING = "pending" + RUNNING = "running" + COMPLETED = "completed" + FAILED = "failed" + + +class Direction(StrEnum): + """角色朝向——多方向扩展时使用。""" + + FRONT = "front" + SIDE = "side" + BACK = "back" + LEFT = "left" + + +# -- 工作流 --------------------------------------------------------------- + + +@dataclass +class Workflow: + """工作流——一个画布实例。""" + + id: int | None = None + user_id: int = 0 + project_id: int | None = None + name: str = "未命名工作流" + status: WorkflowStatus = WorkflowStatus.ACTIVE + project_context: dict = field(default_factory=dict) + schema_version: int = 1 + version: int = 1 + create_at: datetime = field(default_factory=lambda: datetime.now(timezone.utc)) + update_at: datetime = field(default_factory=lambda: datetime.now(timezone.utc)) + + +# -- 画布卡片 ------------------------------------------------------------- + + +@dataclass +class CanvasCard: + """画布卡片——工作流节点。 + + user_input / latest_result 结构按 card_type 不同,见 schema.py 中的注释。 + """ + + id: int | None = None + workflow_id: int = 0 + card_type: CardType = CardType.CHARACTER + status: CardStatus = CardStatus.DRAFT + parent_card_id: int | None = None + direction: Direction | None = None + position_x: float = 0.0 + position_y: float = 0.0 + user_input: dict = field(default_factory=dict) + latest_result: dict | None = None + spec_overrides: dict = field(default_factory=dict) + is_active: bool = True + version: int = 1 + create_at: datetime = field(default_factory=lambda: datetime.now(timezone.utc)) + update_at: datetime = field(default_factory=lambda: datetime.now(timezone.utc)) + + @property + def is_terminal(self) -> bool: + return self.status in (CardStatus.COMPLETED, CardStatus.FAILED, CardStatus.INACTIVE) + + +# -- 生成尝试 ------------------------------------------------------------- + + +@dataclass +class GenerationAttempt: + """生成尝试——每次触发生成创建一条记录。""" + + id: int | None = None + card_id: int = 0 + task_id: int | None = None + attempt_no: int = 1 + status: AttemptStatus = AttemptStatus.PENDING + input_payload: dict = field(default_factory=dict) + result: dict | None = None + error_message: str | None = None + create_at: datetime = field(default_factory=lambda: datetime.now(timezone.utc)) + update_at: datetime = field(default_factory=lambda: datetime.now(timezone.utc)) + + @property + def is_terminal(self) -> bool: + return self.status in (AttemptStatus.COMPLETED, AttemptStatus.FAILED) diff --git a/backend/packages/app/src/windup_app/server/workflow/schema.py b/backend/packages/app/src/windup_app/server/workflow/schema.py new file mode 100644 index 0000000..1cbffe1 --- /dev/null +++ b/backend/packages/app/src/windup_app/server/workflow/schema.py @@ -0,0 +1,306 @@ +"""工作流画布 API Schema。 + +定义前端请求/响应的 Pydantic 模型,与 server 层解耦。 +前端团队参考此文件了解接口契约。 + +卡片输入参数按 card_type 区分,每种类型有明确的 user_input / spec_overrides 结构。 +""" + +from __future__ import annotations + +from enum import StrEnum +from typing import Literal + +from pydantic import BaseModel, ConfigDict, Field + +from windup_app.server.workflow.model import ( + Direction, +) + + +# ══════════════════════════════════════════════════════════════════════════════ +# 工作流 +# ══════════════════════════════════════════════════════════════════════════════ + + +class WorkflowCreateRequest(BaseModel): + """创建工作流。 + + 自动创建一张 CHARACTER 根卡片作为画布起点。 + """ + + project_id: int = Field(description="关联项目 ID,项目约束从这里读取") + name: str = Field(default="未命名工作流", max_length=100, description="工作流名称") + + +class WorkflowOut(BaseModel): + """工作流详情响应(含全部 active 卡片)。""" + + model_config = ConfigDict(from_attributes=True) + + id: int + project_id: int | None = None + name: str + status: str = Field(description="active / archived") + project_context: dict = Field( + default_factory=dict, + description="项目约束快照:{perspective, sprite_width, sprite_height, game_style, ...}", + ) + version: int = Field(description="乐观锁版本号,更新时需携带") + cards: list[CanvasCardOut] = Field(default_factory=list, description="全部有效卡片") + + +# ══════════════════════════════════════════════════════════════════════════════ +# 卡片输入参数(按 card_type 显式定义) +# ══════════════════════════════════════════════════════════════════════════════ + + +# ── CHARACTER 卡片 ──────────────────────────────────────────────────────────── + + +class CharacterConfirmInput(BaseModel): + """CHARACTER 卡片确认时的用户输入。 + + 用户填写角色描述后提交,后端生成 N 张候选图。 + """ + + description: str = Field(min_length=1, max_length=500, description="角色描述,如'甲壳虫战士,手持长剑'") + num_images: int = Field(default=6, ge=1, le=12, description="生成候选图数量,默认6张") + + +class CharacterUpdateInput(BaseModel): + """CHARACTER 卡片更新时的用户输入(选定候选后回填)。""" + + selected_candidate_id: int | None = Field(default=None, description="选定的 CANDIDATE 卡片 ID") + + +# ── ACTION 卡片 ─────────────────────────────────────────────────────────────── + + +class PresetActionType(StrEnum): + """预设动作类型——规格与提示词已预先优化。""" + + IDLE = "idle" + WALK = "walk" + RUN = "run" + JUMP = "jump" + + +class ActionCreateInput(BaseModel): + """ACTION 卡片创建时的用户输入。""" + + action_type: PresetActionType | str = Field( + description="动作类型:预设(idle/walk/run/jump) 或自定义(任意字符串)" + ) + action_name: str | None = Field( + default=None, max_length=50, + description="自定义动作名称(action_type 为自定义时必填)", + ) + description: str = Field( + min_length=1, max_length=300, + description="动作描述,如'大步流星地走'、'向上跳跃'", + ) + reference_image_url: str | None = Field( + default=None, + description="姿势参考图 URL(自定义动作可选,预设动作忽略)", + ) + + +class ActionSpecOverrides(BaseModel): + """ACTION 卡片高级项覆盖。 + + 这些参数有默认值,用户可在"高级项"中覆盖。 + 覆盖后卡片上出现记号标明该处已偏离默认。 + """ + + num_frames: int = Field(default=36, ge=1, le=120, description="生成帧数,默认36帧") + fps: int = Field(default=10, ge=1, le=60, description="帧率,默认10 FPS") + loop: Literal["none", "linear", "pingpong"] = Field( + default="linear", description="循环模式:none=不循环, linear=线性, pingpong=乒乓" + ) + + +# ── EXPORT 卡片 ────────────────────────────────────────────────────────────── + + +class ExportFormat(StrEnum): + """导出格式。""" + + PNG_SEQUENCE = "png_sequence" # PNG 序列帧 + SPRITE_SHEET = "sprite_sheet" # 精灵图集 + PLIST = "plist" # Cocos SpriteFrames + + +class ExportCreateInput(BaseModel): + """EXPORT 卡片创建时的用户输入。""" + + formats: list[ExportFormat] = Field( + min_length=1, + description="导出格式列表,至少选一种", + ) + fps: int = Field(default=10, ge=1, le=60, description="导出帧率") + + +# ── latest_result 结构(按 card_type)───────────────────────────────────────── + + +class CharacterLatestResult(BaseModel): + """CHARACTER 生成结果。""" + + candidate_ids: list[int] = Field(description="生成的 CANDIDATE 卡片 ID 列表") + + +class CandidateLatestResult(BaseModel): + """CANDIDATE 生成结果(与 user_input 相同)。""" + + image_url: str = Field(description="候选图 URL") + + +class ActionFrame(BaseModel): + """动画帧。""" + + index: int = Field(description="帧序号,从0开始") + image_url: str = Field(description="帧图片 URL") + duration_ms: int = Field(default=125, description="帧持续时间(ms)") + + +class ActionLatestResult(BaseModel): + """ACTION 生成结果。""" + + first_frame_url: str | None = Field(default=None, description="首帧图 URL") + frames: list[ActionFrame] = Field(default_factory=list, description="完整动画帧序列") + + +class ExportLatestResult(BaseModel): + """EXPORT 生成结果。""" + + package_url: str = Field(description="导出包下载 URL") + + +# ══════════════════════════════════════════════════════════════════════════════ +# 画布卡片请求/响应 +# ══════════════════════════════════════════════════════════════════════════════ + + +class CardCreateRequest(BaseModel): + """创建子卡片(ACTION / EXPORT)。 + + 由前端"+"菜单触发,parent_card_id 必须指向一张 CHARACTER 卡片。 + ACTION 卡片创建时自动从选定的 CANDIDATE 复制母版图。 + """ + + card_type: Literal["action", "export"] = Field(description="卡片类型") + parent_card_id: int = Field(description="父卡片 ID,必须是 CHARACTER 类型") + direction: Direction | None = Field(default=None, description="方向:front/side/back/left,单方向时省略") + user_input: ActionCreateInput | ExportCreateInput = Field( + description="用户输入:ACTION 用 ActionCreateInput,EXPORT 用 ExportCreateInput", + ) + spec_overrides: ActionSpecOverrides | None = Field( + default=None, + description="ACTION 高级项覆盖(帧数/FPS/循环),EXPORT 省略", + ) + + +class CardUpdateRequest(BaseModel): + """更新卡片用户输入(不触发生成)。 + + CHARACTER:更新 selected_candidate_id(选定候选) + ACTION:更新描述等草稿 + """ + + user_input: CharacterUpdateInput | ActionCreateInput | ExportCreateInput = Field( + description="更新后的用户输入,结构按 card_type" + ) + position_x: float | None = Field(default=None, description="画布 X 坐标(拖动后保存)") + position_y: float | None = Field(default=None, description="画布 Y 坐标(拖动后保存)") + + +class CharacterConfirmRequest(BaseModel): + """CHARACTER 卡片确认请求。""" + + card_type: Literal["character"] = "character" + user_input: CharacterConfirmInput = Field(description="角色描述和候选图数量") + spec_overrides: None = None # CHARACTER 无高级项 + + +class ActionConfirmRequest(BaseModel): + """ACTION 卡片确认请求。""" + + card_type: Literal["action"] = "action" + user_input: ActionCreateInput = Field(description="动作输入(可修改描述后确认)") + spec_overrides: ActionSpecOverrides | None = Field(default=None, description="高级项覆盖") + + +class ExportConfirmRequest(BaseModel): + """EXPORT 卡片确认请求。""" + + card_type: Literal["export"] = "export" + user_input: ExportCreateInput = Field(description="导出格式和帧率") + spec_overrides: None = None # EXPORT 无高级项 + + +# 前端按 card_type 选择对应的 ConfirmRequest 发送 +CardConfirmRequest = CharacterConfirmRequest | ActionConfirmRequest | ExportConfirmRequest + + +class CardRegenerateRequest(BaseModel): + """重新生成(创建新的 GenerationAttempt)。 + + CHARACTER:旧 CANDIDATE 全部 INACTIVE,重新生成候选。ACTION/EXPORT 不受影响。 + ACTION / EXPORT:创建新 attempt,重新执行。 + """ + + user_input: CharacterConfirmInput | ActionCreateInput | ExportCreateInput | None = Field( + default=None, + description="可选:修改输入后重新生成,省略则沿用上次输入", + ) + + +class CanvasCardOut(BaseModel): + """卡片响应。 + + user_input / latest_result 的结构按 card_type 不同,参考对应的 Pydantic 模型。 + """ + + model_config = ConfigDict(from_attributes=True) + + id: int + card_type: str = Field(description="character / candidate / action / export") + status: str = Field(description="draft / generating / completed / failed / inactive") + parent_card_id: int | None = None + direction: str | None = Field(default=None, description="front / side / back / left / null") + position_x: float = 0.0 + position_y: float = 0.0 + user_input: dict = Field( + default_factory=dict, + description="用户输入,结构按 card_type 见对应的 Input 模型", + ) + latest_result: dict | None = Field( + default=None, + description="生成结果,结构按 card_type 见对应的 Result 模型", + ) + spec_overrides: dict = Field( + default_factory=dict, + description="高级项覆盖,ACTION 用 ActionSpecOverrides 结构", + ) + version: int = 1 + + +# ══════════════════════════════════════════════════════════════════════════════ +# 生成尝试 +# ══════════════════════════════════════════════════════════════════════════════ + + +class GenerationAttemptOut(BaseModel): + """生成尝试响应。""" + + model_config = ConfigDict(from_attributes=True) + + id: int + card_id: int + task_id: int | None = Field(default=None, description="关联的 GenerationTask ID,可用于订阅 Task SSE") + attempt_no: int = Field(description="第几次尝试,从 1 开始") + status: str = Field(description="pending / running / completed / failed") + input_payload: dict = Field(default_factory=dict, description="生成输入快照") + result: dict | None = Field(default=None, description="生成结果") + error_message: str | None = None diff --git a/backend/packages/app/src/windup_app/web/api/agent.py b/backend/packages/app/src/windup_app/web/api/agent.py new file mode 100644 index 0000000..051ca96 --- /dev/null +++ b/backend/packages/app/src/windup_app/web/api/agent.py @@ -0,0 +1,166 @@ +"""Agent 智能体 API。 + +懒人智能体:一句话搞定角色资产生成。 +用户通过自然语言与 Agent 对话,Agent 自动调用工具完成项目创建、角色生成、动作生成等操作。 + +端点一览 +-------- +POST /agent/sessions 创建 Agent 会话 +GET /agent/sessions/{session_id}/stream Agent SSE 事件流 +POST /agent/sessions/{session_id}/messages 发送用户消息 +POST /agent/sessions/{session_id}/choices 发送用户选择(按钮点击) +GET /agent/sessions/{session_id}/messages 获取会话历史 + +SSE 事件类型 +----------- +- message: Agent 回复(富内容:text/image/buttons/confirm/progress) +- tool_call: Agent 调用工具(含 task_id,可订阅 Task SSE 获取进度) +- tool_result: 工具返回结果 +- state_change: Agent 状态变更(processing/waiting_input/done) +- error: 错误信息 +""" + +from __future__ import annotations + +import logging + +from fastapi import APIRouter, Depends +from pydantic import BaseModel, ConfigDict, Field +from sqlalchemy.orm import Session + +from windup_common.result import Response, ListResponse +from windup_framework.db import get_session + +logger = logging.getLogger("windup.agent.api") + +router = APIRouter(prefix="/agent", tags=["agent"]) + + +# ══════════════════════════════════════════════════════════════════════════════ +# 请求模型 +# ══════════════════════════════════════════════════════════════════════════════ + + +class SessionCreateRequest(BaseModel): + """创建 Agent 会话。""" + + user_id: int = Field(description="用户 ID") + context: dict | None = Field( + default=None, + description="初始上下文,如 {project_id: 123}。可选,后续通过对话补充。", + ) + + +class MessageSendRequest(BaseModel): + """发送用户消息。""" + + content: str = Field(min_length=1, max_length=2000, description="用户消息内容") + message_id: str | None = Field( + default=None, + description="客户端消息 ID(用于去重),省略则由后端生成", + ) + + +class ChoiceSendRequest(BaseModel): + """发送用户选择(按钮点击)。""" + + message_id: str = Field(description="对应的 Agent 消息 ID(buttons 块所在的 message)") + value: str = Field(description="选择的值(ButtonOption.value)") + + +# ══════════════════════════════════════════════════════════════════════════════ +# 响应模型 +# ══════════════════════════════════════════════════════════════════════════════ + + +class SessionOut(BaseModel): + """会话响应。""" + + model_config = ConfigDict(from_attributes=True) + + session_id: str = Field(description="会话 ID,用于后续 SSE 订阅和消息发送") + created_at: str = Field(description="创建时间(ISO 8601)") + + +class MessageOut(BaseModel): + """消息响应。""" + + model_config = ConfigDict(from_attributes=True) + + message_id: str = Field(description="消息 ID") + role: str = Field(description="角色:user / assistant") + blocks: list[dict] = Field( + default_factory=list, + description="内容块列表,结构见 ContentBlock 定义", + ) + state: str | None = Field( + default=None, + description="Agent 状态:processing / waiting_input / done(仅 assistant 消息)", + ) + timestamp: float = Field(description="时间戳(秒)") + + +# ══════════════════════════════════════════════════════════════════════════════ +# 端点 +# ══════════════════════════════════════════════════════════════════════════════ + + +@router.post("/sessions", response_model=Response[SessionOut]) +def create_session( + body: SessionCreateRequest, + session: Session = Depends(get_session), +) -> Response[SessionOut]: + """创建 Agent 会话。 + + 返回 session_id,前端用于: + 1. 订阅 Agent SSE:GET /agent/sessions/{session_id}/stream + 2. 发送消息:POST /agent/sessions/{session_id}/messages + 3. 发送选择:POST /agent/sessions/{session_id}/choices + """ + # TODO: agent_service.create_session + raise NotImplementedError + + +@router.post("/sessions/{session_id}/messages", response_model=Response[MessageOut]) +def send_message( + session_id: str, + body: MessageSendRequest, + session: Session = Depends(get_session), +) -> Response[MessageOut]: + """发送用户消息。 + + Agent 收到消息后通过 SSE 推送处理结果(message/tool_call/tool_result 等事件)。 + """ + # TODO: agent_service.send_message + raise NotImplementedError + + +@router.post("/sessions/{session_id}/choices", response_model=Response[None]) +def send_choice( + session_id: str, + body: ChoiceSendRequest, + session: Session = Depends(get_session), +) -> Response[None]: + """发送用户选择(按钮点击)。 + + 对应 Agent 消息中的 buttons 块。Agent 收到后继续处理。 + """ + # TODO: agent_service.send_choice + raise NotImplementedError + + +@router.get("/sessions/{session_id}/messages", response_model=ListResponse[MessageOut]) +def get_messages( + session_id: str, + limit: int = 50, + before: str | None = None, + session: Session = Depends(get_session), +) -> ListResponse[MessageOut]: + """获取会话历史消息。 + + 参数: + - limit: 返回条数,默认50 + - before: 分页,此 message_id 之前的消息 + """ + # TODO: agent_service.get_messages + raise NotImplementedError diff --git a/backend/packages/app/src/windup_app/web/api/generation.py b/backend/packages/app/src/windup_app/web/api/generation.py index be3a5e9..872c1c0 100644 --- a/backend/packages/app/src/windup_app/web/api/generation.py +++ b/backend/packages/app/src/windup_app/web/api/generation.py @@ -16,7 +16,7 @@ from windup_common.result import Response from windup_framework.db import get_session -from windup_app.server.generation.model import ( +from windup_app.server.orchestrator.model import ( ActionType, GenerationTask, ) diff --git a/backend/packages/app/src/windup_app/web/api/workflow.py b/backend/packages/app/src/windup_app/web/api/workflow.py new file mode 100644 index 0000000..39929ab --- /dev/null +++ b/backend/packages/app/src/windup_app/web/api/workflow.py @@ -0,0 +1,178 @@ +"""工作流画布 API。 + +契约层:定义端点和请求/响应模型,与 server 层解耦。 +实际逻辑由 server 层实现,本文件只做参数校验和格式转换。 + +端点一览 +-------- +POST /workflow 创建工作流 +GET /workflow/{id} 获取工作流详情 +DELETE /workflow/{id} 删除工作流 +POST /workflow/{wf_id}/cards 创建子卡片 +PATCH /workflow/{wf_id}/cards/{card_id} 更新卡片 +POST /workflow/{wf_id}/cards/{card_id}/confirm 确认卡片(触发生成) +POST /workflow/{wf_id}/cards/{card_id}/regenerate 重新生成 +DELETE /workflow/{wf_id}/cards/{card_id} 删除卡片 +GET /workflow/{wf_id}/cards/{card_id}/attempts 获取生成尝试历史 +""" + +from __future__ import annotations + +import logging + +from fastapi import APIRouter, Depends +from sqlalchemy.orm import Session + +from windup_common.result import Response, ListResponse +from windup_framework.db import get_session + +from windup_app.server.workflow.schema import ( + CanvasCardOut, + CardConfirmRequest, + CardCreateRequest, + CardRegenerateRequest, + CardUpdateRequest, + GenerationAttemptOut, + WorkflowCreateRequest, + WorkflowOut, +) + +logger = logging.getLogger("windup.workflow.api") + +router = APIRouter(prefix="/workflow", tags=["workflow"]) + + +# ── 工作流 CRUD ───────────────────────────────────────────────────────────── + + +@router.post("", response_model=Response[WorkflowOut]) +def create_workflow( + body: WorkflowCreateRequest, + session: Session = Depends(get_session), +) -> Response[WorkflowOut]: + """创建工作流。 + + 自动创建一张 CHARACTER 根卡片作为画布起点。 + 返回的工作流已包含该卡片。 + """ + # TODO: service.create_workflow + raise NotImplementedError + + +@router.get("/{workflow_id}", response_model=Response[WorkflowOut]) +def get_workflow( + workflow_id: int, + session: Session = Depends(get_session), +) -> Response[WorkflowOut]: + """获取工作流详情(含全部 active 卡片)。 + + 前端进入画布时调用,获取完整的卡片树。 + """ + # TODO: service.get_workflow + raise NotImplementedError + + +@router.delete("/{workflow_id}", response_model=Response[None]) +def delete_workflow( + workflow_id: int, + session: Session = Depends(get_session), +) -> Response[None]: + """删除工作流(级联软删除所有卡片)。""" + # TODO: service.delete_workflow + raise NotImplementedError + + +# ── 卡片操作 ──────────────────────────────────────────────────────────────── + + +@router.post("/{workflow_id}/cards", response_model=Response[CanvasCardOut]) +def create_card( + workflow_id: int, + body: CardCreateRequest, + session: Session = Depends(get_session), +) -> Response[CanvasCardOut]: + """创建子卡片(ACTION / EXPORT)。 + + 由前端"+"菜单触发。ACTION 卡片创建时自动复制母版图。 + """ + # TODO: card_service.create_card + raise NotImplementedError + + +@router.patch("/{workflow_id}/cards/{card_id}", response_model=Response[CanvasCardOut]) +def update_card( + workflow_id: int, + card_id: int, + body: CardUpdateRequest, + session: Session = Depends(get_session), +) -> Response[CanvasCardOut]: + """更新卡片用户输入或位置(不触发生成)。 + + 用于保存草稿、拖动画布位置等场景。 + """ + # TODO: card_service.update_card + raise NotImplementedError + + +@router.post("/{workflow_id}/cards/{card_id}/confirm", response_model=Response[GenerationAttemptOut]) +def confirm_card( + workflow_id: int, + card_id: int, + body: CardConfirmRequest, + session: Session = Depends(get_session), +) -> Response[GenerationAttemptOut]: + """确认卡片,触发生成。 + + - CHARACTER:生成候选图 → 自动创建 CANDIDATE 卡片。 + - ACTION:生成动画帧。 + - EXPORT:打包导出。 + + 返回的 GenerationAttempt 包含 task_id,前端应订阅 Task SSE 获取进度: + ``GET /generation/tasks/{task_id}/stream`` + """ + # TODO: card_service.confirm_card + raise NotImplementedError + + +@router.post("/{workflow_id}/cards/{card_id}/regenerate", response_model=Response[GenerationAttemptOut]) +def regenerate_card( + workflow_id: int, + card_id: int, + body: CardRegenerateRequest, + session: Session = Depends(get_session), +) -> Response[GenerationAttemptOut]: + """重新生成(创建新的 GenerationAttempt)。 + + - CHARACTER:旧 CANDIDATE 全部 INACTIVE,重新生成候选。ACTION/EXPORT 不受影响。 + - ACTION / EXPORT:创建新 attempt,重新执行。 + """ + # TODO: card_service.regenerate_card + raise NotImplementedError + + +@router.delete("/{workflow_id}/cards/{card_id}", response_model=Response[None]) +def delete_card( + workflow_id: int, + card_id: int, + session: Session = Depends(get_session), +) -> Response[None]: + """删除卡片(级联软删除所有子卡片)。 + + 删除 CHARACTER 会级联删除其下的 CANDIDATE、ACTION、EXPORT。 + """ + # TODO: card_service.delete_card + raise NotImplementedError + + +@router.get("/{workflow_id}/cards/{card_id}/attempts", response_model=ListResponse[GenerationAttemptOut]) +def get_attempts( + workflow_id: int, + card_id: int, + session: Session = Depends(get_session), +) -> ListResponse[GenerationAttemptOut]: + """获取某张卡片的全部生成尝试记录。 + + 用于展示历史生成结果、调试等。 + """ + # TODO: card_service.get_attempts + raise NotImplementedError diff --git a/backend/packages/app/src/windup_app/web/sse/.gitkeep b/backend/packages/app/src/windup_app/web/sse/.gitkeep deleted file mode 100644 index e69de29..0000000 diff --git a/backend/packages/app/src/windup_app/web/sse/__init__.py b/backend/packages/app/src/windup_app/web/sse/__init__.py new file mode 100644 index 0000000..1984392 --- /dev/null +++ b/backend/packages/app/src/windup_app/web/sse/__init__.py @@ -0,0 +1,10 @@ +"""SSE 模块:Server-Sent Events 实时推送(双模式)。 + +- ``sse_router``: Task mode —— 生成任务进度推送 +- ``session_router``: Session mode —— Agent 多轮对话事件流 +""" + +from windup_app.web.sse.session import router as session_router +from windup_app.web.sse.stream import router as sse_router + +__all__ = ["sse_router", "session_router"] diff --git a/backend/packages/app/src/windup_app/web/sse/event_bus.py b/backend/packages/app/src/windup_app/web/sse/event_bus.py new file mode 100644 index 0000000..21006f2 --- /dev/null +++ b/backend/packages/app/src/windup_app/web/sse/event_bus.py @@ -0,0 +1,83 @@ +"""SSE 内存发布-订阅中心(双模式)。 + +支持两种订阅模式: + +- **Task mode**: ``task_id → events → 终态关闭``(用于生成任务进度推送) +- **Session mode**: ``session_id → events → 持续连接``(用于 Agent 多轮对话) + +线程安全:``publish`` 可从任意线程调用,``subscribe``/``unsubscribe`` 须在 +event loop 中调用(AsyncAPI 路由天然满足)。 +""" + +from __future__ import annotations + +import asyncio +from collections import defaultdict +from dataclasses import dataclass + + +@dataclass(frozen=True) +class SSEEvent: + """统一 SSE 事件。""" + + id: str # 任务 ID(str) 或会话 ID + id_type: str # "task" | "session" + event: str # 事件类型 + data: dict # 事件数据 + + +class EventBus: + """内存发布-订阅:后台线程 publish,异步 SSE generator subscribe。""" + + def __init__(self) -> None: + # Task mode: task_id(str) → subscriber queues + self._task_queues: dict[str, list[asyncio.Queue[SSEEvent]]] = defaultdict(list) + # Session mode: session_id → subscriber queues + self._session_queues: dict[str, list[asyncio.Queue[SSEEvent]]] = defaultdict(list) + + # ── Task mode ───────────────────────────────────────────────────────── + + async def subscribe_task(self, task_id: int) -> asyncio.Queue[SSEEvent]: + """订阅任务事件流。""" + key = str(task_id) + queue: asyncio.Queue[SSEEvent] = asyncio.Queue() + self._task_queues[key].append(queue) + return queue + + async def unsubscribe_task(self, task_id: int, queue: asyncio.Queue[SSEEvent]) -> None: + """取消任务订阅。""" + key = str(task_id) + subscribers = self._task_queues.get(key) + if subscribers and queue in subscribers: + subscribers.remove(queue) + if not subscribers: + del self._task_queues[key] + + def publish_task(self, task_id: int, event: str, data: dict) -> None: + """发布任务事件。线程安全(put_nowait 无阻塞)。""" + key = str(task_id) + sse_event = SSEEvent(id=key, id_type="task", event=event, data=data) + for queue in self._task_queues.get(key, []): + queue.put_nowait(sse_event) + + # ── Session mode ────────────────────────────────────────────────────── + + async def subscribe_session(self, session_id: str) -> asyncio.Queue[SSEEvent]: + """订阅会话事件流(Agent 多轮对话)。""" + queue: asyncio.Queue[SSEEvent] = asyncio.Queue() + self._session_queues[session_id].append(queue) + return queue + + async def unsubscribe_session(self, session_id: str, queue: asyncio.Queue[SSEEvent]) -> None: + """取消会话订阅。""" + subscribers = self._session_queues.get(session_id) + if subscribers and queue in subscribers: + subscribers.remove(queue) + if not subscribers: + del self._session_queues[session_id] + + def publish_session(self, session_id: str, event: str, data: dict) -> None: + """发布会话事件。线程安全(put_nowait 无阻塞)。""" + sse_event = SSEEvent(id=session_id, id_type="session", event=event, data=data) + for queue in self._session_queues.get(session_id, []): + queue.put_nowait(sse_event) diff --git a/backend/packages/app/src/windup_app/web/sse/progress.py b/backend/packages/app/src/windup_app/web/sse/progress.py new file mode 100644 index 0000000..1bcb742 --- /dev/null +++ b/backend/packages/app/src/windup_app/web/sse/progress.py @@ -0,0 +1,37 @@ +"""SSE-aware ProgressPort 实现。 + +实现 ``windup_ai_engine.ports.ProgressPort`` 协议,将生成管线的每一步 +进度通过 ``EventBus`` 广播给 SSE 订阅者。 + +后台线程调用 ``step`` 时,通过 ``loop.call_soon_threadsafe`` 安全地 +将事件发布到 event loop。 +""" + +from __future__ import annotations + +import asyncio + +from windup_app.web.sse.event_bus import EventBus + + +class SSEProgressPort: + """实现 ProgressPort:每步都 publish 到 EventBus。""" + + def __init__( + self, + event_bus: EventBus, + task_id: int, + loop: asyncio.AbstractEventLoop, + ) -> None: + self._bus = event_bus + self._task_id = task_id + self._loop = loop + + def step(self, stage: str, i: int, total: int, note: str = "") -> None: + """进度回调:从后台线程安全地发布到 event loop。""" + self._loop.call_soon_threadsafe( + self._bus.publish_task, + self._task_id, + "progress", + {"stage": stage, "current": i, "total": total, "note": note}, + ) diff --git a/backend/packages/app/src/windup_app/web/sse/session.py b/backend/packages/app/src/windup_app/web/sse/session.py new file mode 100644 index 0000000..de7856b --- /dev/null +++ b/backend/packages/app/src/windup_app/web/sse/session.py @@ -0,0 +1,78 @@ +"""SSE 流端点(Session mode)—— Agent 多轮对话。 + +提供 ``GET /agent/sessions/{session_id}/stream`` 端点, +客户端通过 EventSource 订阅 Agent 事件流。 + +Session mode 特点: +- 持续连接,不自动关闭(除非客户端断开) +- 支持多种事件类型:thinking / message / tool_call / tool_result / status / error +- 双向通信:客户端通过 HTTP POST 发送用户消息 +""" + +from __future__ import annotations + +import asyncio +import json +import logging + +from fastapi import APIRouter, Request +from fastapi.responses import StreamingResponse + +from windup_app.web.sse.event_bus import EventBus, SSEEvent + +logger = logging.getLogger("windup.sse.session") + +router = APIRouter(prefix="/agent", tags=["sse"]) + +# SSE 心跳间隔(秒) +_HEARTBEAT_TIMEOUT = 30.0 + + +@router.get("/sessions/{session_id}/stream") +async def stream_session( + session_id: str, + request: Request, +) -> StreamingResponse: + """SSE:Agent 多轮对话事件流。 + + 事件类型: + - ``thinking``: Agent 思考过程 + - ``message``: Agent 回复文本 + - ``tool_call``: Agent 调用工具 + - ``tool_result``: 工具返回结果 + - ``status``: 会话状态(waiting_input / processing / done) + - ``error``: 错误信息 + + Session mode 不会自动关闭连接,持续监听直到客户端断开。 + """ + event_bus: EventBus = request.app.state.event_bus + queue = await event_bus.subscribe_session(session_id) + logger.debug("SSE Session 订阅: session_id=%s", session_id) + + async def _event_generator(): + try: + while True: + if await request.is_disconnected(): + logger.debug("SSE Session 客户端断开: session_id=%s", session_id) + break + try: + event: SSEEvent = await asyncio.wait_for( + queue.get(), timeout=_HEARTBEAT_TIMEOUT, + ) + payload = json.dumps(event.data, ensure_ascii=False) + yield f"event: {event.event}\ndata: {payload}\n\n" + except asyncio.TimeoutError: + yield ": heartbeat\n\n" + finally: + await event_bus.unsubscribe_session(session_id, queue) + logger.debug("SSE Session 取消订阅: session_id=%s", session_id) + + return StreamingResponse( + _event_generator(), + media_type="text/event-stream", + headers={ + "Cache-Control": "no-cache", + "Connection": "keep-alive", + "X-Accel-Buffering": "no", + }, + ) diff --git a/backend/packages/app/src/windup_app/web/sse/stream.py b/backend/packages/app/src/windup_app/web/sse/stream.py new file mode 100644 index 0000000..f3d9874 --- /dev/null +++ b/backend/packages/app/src/windup_app/web/sse/stream.py @@ -0,0 +1,111 @@ +"""SSE 流端点(Task mode)。 + +提供 ``GET /generation/tasks/{task_id}/stream`` 端点, +客户端通过 EventSource 订阅任务进度推送。 +""" + +from __future__ import annotations + +import asyncio +import dataclasses +import json +import logging + +from fastapi import APIRouter, Depends, Query, Request +from fastapi.responses import StreamingResponse +from sqlalchemy.orm import Session + +from windup_app.server.orchestrator.model import TaskStatus +from windup_app.web.sse.event_bus import EventBus, SSEEvent +from windup_framework.db import get_session + +logger = logging.getLogger("windup.sse.stream") + +router = APIRouter(prefix="/generation", tags=["sse"]) + +# SSE 心跳间隔(秒):超过此时间无事件则发送注释保活 +_HEARTBEAT_TIMEOUT = 30.0 + + +def _task_to_sse_data(task) -> dict | None: + """任务领域对象 → SSE 事件 data(仅终态)。""" + if task.status == TaskStatus.COMPLETED: + result_dict = None + if task.result is not None: + result_dict = dataclasses.asdict(task.result) + return result_dict or {} + if task.status == TaskStatus.FAILED: + return {"error": task.error_message or "未知错误"} + return None + + +@router.get("/tasks/{task_id}/stream") +async def stream_task( + task_id: int, + request: Request, + project_id: int = Query(..., gt=0), + session: Session = Depends(get_session), +) -> StreamingResponse: + """SSE:实时推送任务进度与最终结果。 + + 事件类型: + - ``status``: 任务状态变更(pending/running/completed/failed) + - ``progress``: 生成管线进度(stage/current/total) + - ``completed``: 任务完成,携带最终结果 + - ``failed``: 任务失败,携带错误信息 + + 若客户端订阅时任务已处于终态,立即推送终态事件并关闭连接。 + """ + from windup_app.server.orchestrator.service import service as generation_service + + event_bus: EventBus = request.app.state.event_bus + + # 检查任务初始状态:若已终态,立即推送并关闭 + task = generation_service.get_task(session, project_id, task_id) + if task is not None and task.status in (TaskStatus.COMPLETED, TaskStatus.FAILED): + event = "completed" if task.status == TaskStatus.COMPLETED else "failed" + data = _task_to_sse_data(task) or {} + + async def _immediate(): + yield f"event: {event}\ndata: {json.dumps(data, ensure_ascii=False)}\n\n" + + return StreamingResponse( + _immediate(), + media_type="text/event-stream", + headers={"Cache-Control": "no-cache", "Connection": "keep-alive"}, + ) + + # 订阅事件流 + queue = await event_bus.subscribe_task(task_id) + logger.debug("SSE 订阅: task_id=%d", task_id) + + async def _event_generator(): + try: + while True: + if await request.is_disconnected(): + logger.debug("SSE 客户端断开: task_id=%d", task_id) + break + try: + event: SSEEvent = await asyncio.wait_for( + queue.get(), timeout=_HEARTBEAT_TIMEOUT, + ) + payload = json.dumps(event.data, ensure_ascii=False) + yield f"event: {event.event}\ndata: {payload}\n\n" + if event.event in ("completed", "failed"): + logger.debug("SSE 终态: task_id=%d event=%s", task_id, event.event) + break + except asyncio.TimeoutError: + yield ": heartbeat\n\n" + finally: + await event_bus.unsubscribe_task(task_id, queue) + logger.debug("SSE 取消订阅: task_id=%d", task_id) + + return StreamingResponse( + _event_generator(), + media_type="text/event-stream", + headers={ + "Cache-Control": "no-cache", + "Connection": "keep-alive", + "X-Accel-Buffering": "no", # nginx 透传 + }, + )