aps-agent/server/gateway/app.py

334 lines
20 KiB
Python
Raw Normal View History

2026-07-21 11:05:57 +08:00
# ============================================================
# FastAPI 应用与路由(moduleId: gateway-app, 可重生 ✅)
# 端点清单:
# POST /api/chat 对话(SSE 流:meta/intent/token/command/block/done)
# GET /api/world/summary KPI 摘要
# GET /api/world/gantt 甘特视图数据
# GET /api/world/load 负荷热力数据
# GET /api/world/due 交期承诺看板数据
# GET /api/orders 订单管理页数据(订单 + 成品选项)
# POST /api/orders/stage 订单写入动作暂存确认卡(order.upsert/cancel/complete)
# GET /api/master 主数据管理页数据(资源树/物料BOM/日历维保)
# POST /api/master/stage 主数据写入动作暂存确认卡(master.*.upsert)
# GET /api/mrp MRP 建议单(采购/委外,订单分解产物)
# POST /api/mrp/decompose 订单分解(P1:产出 DRAFT 建议)
# GET /api/timeline 时间线导轨(检查点 + 版本,§4.4)
# POST /api/actions/confirm 确认卡回传(批准/驳回 P2 动作)
# POST /api/actions/scenario/apply 采用某个沙盒方案(P1 正式排产)
# GET /api/gov/pending 门禁管理台·待审批队列(§6.10.2)
# GET /api/gov/audit 门禁管理台·审计链 + 完整性校验(§3.6)
# GET /api/gov/modules 重生中心·可重生模块注册表(§6.10.3)
# GET /api/gov/policy 门禁管理台·权力矩阵投影(M2.5,docs/architecture/harness.md 的 API 化)
# GET /api/gov/tests 重生中心·黄金测试看板(M2.5,读缓存 / ?run=true 重跑)
# GET /api/settings/llm 设置中心·LLM Provider 状态(M2.5)
# GET /api/knowledge/assets 知识资产清单(M3 §8.1)
# GET /api/knowledge/assets/{id} 单个知识资产全文(M3)
# GET /api/reports/{type} 报告导出(Markdown 下载,M3 §9.10)
# ⚠ 文档同步铁律(plan.md §12.9):新增端点必同步 docs/architecture/harness.md 只读端点备案
# ============================================================
from __future__ import annotations # 前向类型引用
import asyncio # 流式分片的微延时
import json # SSE 载荷序列化
import uuid # 会话 ID 生成
from typing import Any, AsyncIterator # 类型标注
from fastapi import FastAPI # Web 框架
from fastapi.middleware.cors import CORSMiddleware # 跨域(开发期前端直连用)
from fastapi.responses import StreamingResponse # SSE 响应
from pydantic import BaseModel, Field # 请求体契约
from server.agent_core import harness # 门禁(待审批队列投影)
from server.agent_core.intent import recognize # 意图两级管线
from server.agent_core.registry import scan_modules, verify_audit_chain # 注册表与审计校验
from server.aps_domain.views import due_view, gantt_view, load_view, world_summary # 世界视图
from server.aps_domain.workflow import execute_confirmed, handle_intent # 工作流编排
from server.aps_domain.orders import confirmation_for_order_action, list_orders, list_products
from server.aps_domain.masterdata import MASTER_ACTIONS, confirmation_for_master_action, master_overview
from server.aps_domain.mrp import decompose_orders, list_mrp, summarize_decomposition
from server.contracts import IntentResult # 方案采用时构造意图
from server.state.checkpoints import get_checkpoints # 成对快照仓(时间线)
from server.state.store import get_store # 世界状态单例
# ---------------- 请求体契约 ----------------
class ChatRequest(BaseModel):
"""对话请求:会话 ID(可空=新会话)+ 用户文本。"""
sessionId: str | None = None # 会话 ID(M1 仅作标识与审计维度)
text: str = Field(min_length=1, max_length=500) # 用户口令(限长防滥用)
class ConfirmRequest(BaseModel):
"""确认卡回传:确认 ID + 批准与否。"""
sessionId: str | None = None # 会话 ID(审计维度)
confirmId: str # 待确认动作的一次性令牌
approve: bool # True=批准执行 False=驳回
class ScenarioApplyRequest(BaseModel):
"""采用沙盒方案:策略 + 引擎(引擎确定性 → 重跑即得卡片同款结果)。"""
sessionId: str | None = None # 会话 ID(审计维度)
strategy: str # 选中方案的策略模板
engine: str = "RULE" # 引擎类型
class OrderStageRequest(BaseModel):
"""订单管理写入动作:只暂存确认卡,不直接改世界。"""
sessionId: str | None = None # 发起会话
action: str # order.upsert / order.cancel / order.complete
payload: dict[str, Any] = Field(default_factory=dict) # 订单载荷
class MasterStageRequest(BaseModel):
"""主数据维护写入动作:只暂存确认卡,不直接改世界(MD-01/02/03)。"""
sessionId: str | None = None # 发起会话
action: str # master.line/material/maintenance/bom/routing.upsert
payload: dict[str, Any] = Field(default_factory=dict) # 主数据载荷
class MrpDecomposeRequest(BaseModel):
"""订单分解请求(P1:产出 DRAFT 采购/委外建议)。"""
sessionId: str | None = None # 发起会话
orderNo: str | None = None # 指定订单号;空=全部可排产订单
# ---------------- SSE 工具 ----------------
def _sse(payload: dict[str, Any]) -> str:
"""把一个事件对象编码为 SSE 帧(data: {json}\\n\\n)。"""
return "data: " + json.dumps(payload, ensure_ascii=False) + "\n\n" # 标准 SSE 数据帧
def create_app() -> FastAPI:
"""应用工厂:注册中间件与全部路由。"""
app = FastAPI(title="APS Planning Agent", version="0.1.0") # 应用实例
app.add_middleware( # 开发期放开跨域(生产由网关收紧)
CORSMiddleware, allow_origins=["*"], allow_methods=["*"], allow_headers=["*"])
# ---------------- 对话:SSE 会话流 ----------------
@app.post("/api/chat")
async def chat(req: ChatRequest) -> StreamingResponse:
"""一轮对话(权力等级:随意图而定;本端点只编排不执行 P2)。"""
store = get_store() # 世界状态
session_id = req.sessionId or uuid.uuid4().hex[:12] # 会话 ID(缺省新建)
async def stream() -> AsyncIterator[str]: # SSE 事件生成器
yield _sse({"type": "meta", "sessionId": session_id}) # ① 会话元信息
try:
intent = await recognize(req.text, store.data) # ② 意图识别(两级管线)
yield _sse({"type": "intent", "intent": intent.model_dump()}) # ③ 意图作为协议证据下发
reply = await handle_intent(store, session_id, intent) # ④ 工作流执行(P0/P1 直通,P2 出卡)
for i in range(0, len(reply.text), 24): # ⑤ 正文分片(打字机流式感)
yield _sse({"type": "token", "text": reply.text[i:i + 24]}) # 24 字符一片
await asyncio.sleep(0.02) # 微延时(不阻塞事件循环)
for cmd in reply.commands: # ⑥ 视口命令逐条下发
yield _sse({"type": "command", "command": cmd.model_dump()})
for block in reply.blocks: # ⑦ UI 块(确认卡等)逐条下发
yield _sse({"type": "block", "block": block.model_dump()})
except Exception as exc: # 兜底:任何异常转为 error 事件(连接不中断)
yield _sse({"type": "error", "message": f"处理失败:{exc}"})
yield _sse({"type": "done"}) # ⑧ 结束标记
return StreamingResponse(stream(), media_type="text/event-stream", # SSE 响应
headers={"Cache-Control": "no-cache", "X-Accel-Buffering": "no"})
# ---------------- 世界视图(P0 只读) ----------------
@app.get("/api/world/summary")
async def summary() -> dict:
"""KPI 摘要(视口工具栏)。"""
return world_summary(get_store().data) # 直接投影
@app.get("/api/world/gantt")
async def gantt() -> dict:
"""甘特视图数据(最新版本工单 + 骨架 + 标记)。"""
return gantt_view(get_store().data) # 直接投影
@app.get("/api/world/load")
async def load(days: int = 14) -> dict:
"""负荷热力数据(产线 × 未来 N 天)。"""
return load_view(get_store().data, days=days) # 直接投影
@app.get("/api/world/due")
async def due() -> dict:
"""交期承诺看板数据。"""
return {"rows": due_view(get_store().data)} # 包一层便于扩展
# ---------------- 订单管理(P0 读;P2 写经确认卡) ----------------
@app.get("/api/orders")
async def orders() -> dict:
"""订单管理页数据:销售订单列表 + 成品下拉选项。"""
world = get_store().data
return {"orders": list_orders(world), "products": list_products(world),
"statuses": ["SUBMITTED", "CONFIRMED", "CANCELLED", "COMPLETED"],
"levels": ["VIP", "A", "B", "C"]}
@app.post("/api/orders/stage")
async def order_stage(req: OrderStageRequest) -> dict:
"""订单写入动作暂存为确认卡(P2 唯一路径仍是 /api/actions/confirm)。"""
if req.action not in ("order.upsert", "order.cancel", "order.complete"):
return {"error": "不支持的订单动作"}
store = get_store()
title, lines = confirmation_for_order_action(store.data, req.action, req.payload)
block = harness.stage_confirmation(req.sessionId or "web", req.action, req.payload,
title=title, summary_lines=lines)
write_action = req.action + ".stage"
from server.agent_core.audit import write_audit
write_audit(store.data, store.next_id, actor=req.sessionId or "web", category="GATE",
action=write_action, target={"type": "SALES_ORDER", "id": req.payload.get("id")},
power="P2", rationale={"confirmId": block.props["confirmId"]})
store.save()
return {"message": f"{title} 已进入 P2 确认队列。", "block": block.model_dump()}
# ---------------- 主数据维护(MD-01/02/03:P0 读;P2 写经确认卡) ----------------
@app.get("/api/master")
async def master() -> dict:
"""主数据管理页数据:资源树 + 物料/BOM/路线 + 日历/维保。"""
return master_overview(get_store().data)
@app.post("/api/master/stage")
async def master_stage(req: MasterStageRequest) -> dict:
"""主数据写入动作暂存为确认卡(P2 唯一执行路径仍是 /api/actions/confirm)。"""
if req.action not in MASTER_ACTIONS:
return {"error": "不支持的主数据动作"}
store = get_store()
try:
title, lines = confirmation_for_master_action(store.data, req.action, req.payload)
except ValueError as exc: # 校验失败:不出卡,直接回错误
return {"error": str(exc)}
block = harness.stage_confirmation(req.sessionId or "web", req.action, req.payload,
title=title, summary_lines=lines)
from server.agent_core.audit import write_audit
write_audit(store.data, store.next_id, actor=req.sessionId or "web", category="GATE",
action=req.action + ".stage", target={"type": "MASTER_DATA", "id": req.payload.get("id")},
power="P2", rationale={"confirmId": block.props["confirmId"]})
store.save()
return {"message": f"{title} 已进入 P2 确认队列。", "block": block.model_dump()}
# ---------------- MRP 订单分解(P0 读;P1 分解直通) ----------------
@app.get("/api/mrp")
async def mrp() -> dict:
"""MRP 建议单投影:采购建议 + 委外建议。"""
return list_mrp(get_store().data)
@app.post("/api/mrp/decompose")
async def mrp_decompose(req: MrpDecomposeRequest) -> dict:
"""订单分解(P1 直通:只写 DRAFT 建议表,经 harness.guard 校验权级)。"""
store = get_store()
try:
result = harness.guard("order.decompose", {"orderNo": req.orderNo},
lambda: decompose_orders(store.data, store.next_id, req.orderNo))
except ValueError as exc:
return {"error": str(exc)}
from server.agent_core.audit import write_audit
write_audit(store.data, store.next_id, actor=req.sessionId or "web", category="ALGO_RUN",
action="order.decompose", target={"type": "MRP", "id": req.orderNo or "ALL"},
power="P1", rationale={"purchase": len(result["purchase"]),
"outsource": len(result["outsource"])})
store.save()
return {"message": summarize_decomposition(result), "result": result}
# ---------------- 时间线导轨(§4.4 中观层数据源) ----------------
@app.get("/api/timeline")
async def timeline() -> dict:
"""检查点 + 排产版本的时间线(前端导轨渲染;P0 只读)。"""
world = get_store().data # 世界状态
versions = [{ # 版本投影(轻量字段)
"id": v["id"], "versionNo": v["versionNo"], "status": v["status"],
"createdAt": v["createdAt"], "conflictCount": v["conflictCount"],
"engineType": v["engineType"],
} for v in world["scheduleVersions"]]
return {"checkpoints": get_checkpoints().list_meta(), "versions": versions} # 两类锚点
# ---------------- 采用沙盒方案(P1:正式排产) ----------------
@app.post("/api/actions/scenario/apply")
async def scenario_apply(req: ScenarioApplyRequest) -> dict:
"""把沙盒对比中选中的策略落为正式排产(引擎确定性保证与卡片一致)。"""
store = get_store() # 世界状态
intent = IntentResult(intent="schedule.run", # 构造等价的排产意图(复用同一执行路径)
params={"strategy": req.strategy, "engine": req.engine},
confidence=1.0, source="RULE_FAST")
reply = await handle_intent(store, req.sessionId or "web", intent) # 走标准工作流(含审计落盘)
# 偏好信号:采用比试排权重更高(明确的方案选择 §8.3)
from server.knowledge import get_preferences # 局部导入(避免循环)
get_preferences().record(req.strategy, source="scenario.apply", actor=req.sessionId or "web")
return {"message": reply.text, "refresh": True} # 前端刷新世界与时间线
# ---------------- 确认卡回传(P2 唯一执行通道) ----------------
@app.post("/api/actions/confirm")
async def confirm(req: ConfirmRequest) -> dict:
"""批准/驳回一条 P2 动作(门禁放行后的执行入口,§3.3)。"""
store = get_store() # 世界状态
message = execute_confirmed(store, req.confirmId, req.approve, # 执行(内部写审计+落盘)
actor=req.sessionId or "planner")
return {"message": message, "refresh": req.approve} # 批准后前端应刷新世界数据
# ---------------- 治理可视化(§6.10 门禁管理台 / 重生中心 数据源) ----------------
@app.get("/api/gov/pending")
async def gov_pending() -> dict:
"""门禁管理台·待审批队列(聊天确认卡与此为同一数据的两个视图)。"""
return {"pending": harness.list_pending()} # 待确认动作投影
@app.get("/api/gov/audit")
async def gov_audit(limit: int = 50) -> dict:
"""门禁管理台·审计链(最近 N 条 + 全链完整性校验结果,§3.6)。"""
events = get_store().data.get("auditEvents", []) # 全部审计事件
return {"events": events[-limit:], # 最近 N 条(新在后)
"chain": verify_audit_chain(events)} # 哈希链完整性校验
@app.get("/api/gov/modules")
async def gov_modules() -> dict:
"""重生中心·可重生模块注册表(扫描源码 moduleId 声明,§9.4/§6.10.3)。"""
return {"modules": scan_modules()} # 注册表清单
@app.get("/api/gov/policy")
async def gov_policy() -> dict:
"""门禁管理台·权力矩阵(_POWER_MAP 只读投影;P0)。"""
return {"policy": harness.list_policy(), # 矩阵条目
"defaultPower": "P3"} # 未登记动作的默认等级(白名单原则)
@app.get("/api/gov/tests")
async def gov_tests(run: bool = False) -> dict:
"""重生中心·黄金测试看板(P0 只读缓存;?run=true 在线程池重跑,不阻塞事件循环)。"""
import asyncio # 线程池调度
from server.gateway.golden import golden_status # 局部导入(避免启动即依赖 pytest)
# 跑测是同步子进程(~10s):必须扔线程池,否则会饿死其他请求(世界视图全超时)
return await asyncio.to_thread(golden_status, run)
# ---------------- 知识库与报告(M3 §8/§9.10,P0 只读) ----------------
@app.get("/api/knowledge/assets")
async def knowledge_assets() -> dict:
"""知识资产元信息清单(知识面板/命令面板数据源)。"""
from server.knowledge import get_knowledge # 局部导入
return {"assets": get_knowledge().list_meta()} # 元信息(不含正文)
@app.get("/api/knowledge/assets/{asset_id}")
async def knowledge_asset(asset_id: str) -> dict:
"""单个知识资产全文(出处点开查看)。"""
from server.knowledge import get_knowledge # 局部导入
asset = get_knowledge().get(asset_id) # 按 ID 取
return asset or {"error": "资产不存在"} # 缺失给错误体(前端兜底)
@app.get("/api/reports/{report_type}")
async def report_export(report_type: str) -> Any:
"""报告导出:即时生成并以 Markdown 附件返回(daily / version-diff)。"""
from urllib.parse import quote # 中文文件名 RFC5987 编码
from fastapi.responses import PlainTextResponse # 文本响应(附件头)
from server.aps_domain.reports import build_report # 报告工厂
report = build_report(get_store().data, report_type) # 生成(数字取自冻结快照)
filename = quote(f"{report['title'].replace(' ', '_')}.md") # 编码后的下载文件名
return PlainTextResponse(report["markdown"], media_type="text/markdown; charset=utf-8",
headers={"Content-Disposition": f"attachment; filename*=UTF-8''{filename}"})
@app.get("/api/settings/llm")
async def settings_llm() -> dict:
"""设置中心·LLM Provider 状态(P0 只读;不暴露密钥)。"""
from server.agent_core.providers import get_provider # 单例
p = get_provider() # 当前配置
return {
"provider": p.provider or None, # deepseek / kimi / None
"model": p.model or None, # 模型名
"baseUrl": p.base_url or None, # 端点(无敏感信息)
"enabled": p.enabled, # False = 降级为规则解析
"fallback": "规则快路解析器(离线可用)", # 降级去向(固定文案)
}
return app # 返回配置完成的应用