8.9 KiB
8.9 KiB
网关(Gateway)项目架构流程图
基于代码现状绘制的架构图(backend + frontend + agent 联动)。 端口:后端
:8000、前端:5173、Agent:8001、Redis 远程实例45.207.192.237:56987/0。
一、系统总体架构图(分层)
flowchart TB
subgraph OUT["外部调用方"]
CLI["CLI 客户端 / 模型对话<br/>提交任务(带 auth)"]
end
subgraph FE["前端运维后台(Vue3 + Vite + Element Plus):5173"]
TV["TaskView 任务管理<br/>列表/筛选/详情/取消/重置"]
AV["AgentView Agent 管理<br/>卡片/标签/负载/下线"]
LV["LogView 日志审计<br/>按 request_id/agent_id 筛选"]
CV["ControlView 总览<br/>任务/Agent 统计"]
API["api/index.ts<br/>fetch 封装(/api/admin/*)"]
end
subgraph GW["网关后端(FastAPI + Uvicorn):8000"]
direction TB
subgraph API_L["API 层(app/api/)"]
A1["/api/cli/*<br/>任务提交"]
A2["/api/agent/*<br/>注册/心跳/结果回传"]
A3["/api/admin/*<br/>运维后台查询与管控"]
end
subgraph SVC["服务层(app/services/)"]
S1["TaskService 任务池<br/>受理/取消/重置/超时"]
S2["AgentService Agent 池<br/>注册/心跳/下线/剔除"]
S3["Scheduler 规则调度<br/>选人→绑定→推送→回退"]
S4["RelayService 通信中转<br/>正向推送/反向收结果"]
S5["HeartbeatService<br/>心跳扫描(60s)"]
S6["TimeoutService<br/>任务超时检查(5s)"]
end
subgraph REPO["存储层(app/repository/)"]
R1["TaskRepo"]
R2["AgentRepo"]
R3["LogRepo"]
end
subgraph LOOP["后台循环(scheduler_loop.py)"]
L1["heartbeat_loop<br/>Lock:lock:heartbeat"]
L2["dispatch_loop(2s)<br/>Lock:lock:scheduler"]
L3["timeout_loop<br/>Lock:lock:timeout"]
end
end
subgraph RD["Redis 存储(DB0)"]
K1["task:info:{id} Hash<br/>task:pending / task:running Set"]
K2["agent:info:{id} Hash<br/>agent:all / agent:tag:{tag} Set"]
K3["agent:heartbeat ZSet"]
K4["lock:* 分布式锁"]
K5["log:audit List"]
end
subgraph AG["Agent 服务(ADK ApiServer):8001"]
direction TB
AG0["api_server.py<br/>REST /run /run_sse<br/>SQLite 会话 / InMemory 记忆"]
AG1["task_receiver.py<br/>POST /tasks/{id}<br/>校验 auth → 202 → 后台线程"]
AG2["gateway_client.py<br/>注册 / 心跳 / 回传结果"]
AG3["ADK InMemoryRunner<br/>执行 dev_app(my_agent)"]
end
CLI -->|"POST /api/cli/tasks"| A1
FE --> API
API -->|"/api/admin/*"| A3
A1 --> S1
A2 --> S2
A3 --> S1 & S2
S1 & S2 & S3 & S4 & S5 & S6 --> R1 & R2 & R3
L1 -->|"调用"| S5
L2 -->|"调用"| S3
L3 -->|"调用"| S6
S5 -->|"Lock"| K4
S3 -->|"Lock"| K4
R1 & R2 & R3 --> K1 & K2 & K3 & K5
S3 -->|"Scheduler 选中 Agent"| S4
S4 -->|"POST {endpoint}/tasks/{request_id}<br/>期望 202"| AG1
AG1 -->|"校验通过后后台线程"| AG3
AG3 -->|"执行完成"| AG2
AG2 -->|"POST /api/agent/result"| A2
AG2 -->|"心跳 /api/agent/heartbeat(10s)"| A2
二、任务全生命周期时序图(完整闭环)
sequenceDiagram
autonumber
participant CLI as CLI / 对话模型
participant GW as 网关后端(FastAPI :8000)
participant RD as Redis
participant LOOP as SchedulerLoop 后台循环
participant AG as Agent 服务(:8001)
participant ADK as ADK InMemoryRunner
Note over AG,ADK: Agent 启动:register → 每 10s heartbeat
CLI->>GW: POST /api/cli/tasks {auth, task_type, tags, payload}
GW->>GW: TaskService.submit(幂等:RequestID 复用)
GW->>RD: hset task:info:{id} + sadd task:pending
GW-->>CLI: 200 {request_id}
loop dispatch_loop(每 2s,持 lock:scheduler)
LOOP->>RD: 取 task:pending 成员
LOOP->>GW: Scheduler.dispatch_pending(batch)
GW->>RD: AgentRepo.by_tags(tags) 标签匹配
GW->>GW: 负载过滤(load < max_concurrent)
GW->>GW: 最低负载 + 在线最久者优先
GW->>RD: 任务置 running + 绑定 agent_id<br/>task:pending→task:running,agent 负载 +1
GW->>AG: RelayService.dispatch_command<br/>POST {endpoint}/tasks/{id} {auth, payload}
alt 推送成功(202)
AG->>AG: 校验 auth → 后台线程接收任务
AG-->>GW: 202 {ok, request_id, status:"accepted"}
AG->>ADK: 运行 agent(InMemoryRunner.run_async)
ADK-->>AG: 最终输出文本
AG->>GW: POST /api/agent/result {success, progress:100, result}
GW->>RD: 更新任务 success + 释放 agent 负载<br/>task:running 移除
GW-->>CLI: (前端可查)
else 推送失败(网络错误 / 非 202)
AG-->>GW: 4xx/5xx 或无响应
GW->>RD: 回退:任务置 pending + 解绑 + 释放负载<br/>task:running→task:pending
end
end
Note over LOOP,RD: heartbeat_loop:60s 扫 agent:heartbeat ZSet,<br/>超 120s 未心跳 → 标记 offline<br/>timeout_loop:5s 扫 running,超时 → failed
三、规则调度决策流程(Scheduler)
flowchart LR
START(["dispatch_pending<br/>取 task:pending 批量"]) --> T1{"任务存在且为 pending?"}
T1 -- 否 --> END(["跳过"])
T1 -- 是 --> T2{"按 task_tags 标签匹配<br/>AgentRepo.by_tags(tags)"}
T2 -- 无候选 --> END2(["保持 pending<br/>等待下次调度"])
T2 -- 有候选 --> T3["负载过滤<br/>current_load < max_concurrent"]
T3 -- 空 --> END2
T3 -- 非空 --> T4["排序:最低负载优先<br/>负载相同时在线最久优先"]
T4 --> T5["绑定 Agent<br/>task:info 置 running + agent_id<br/>task:pending→task:running"]
T5 --> T6["agent 负载 +1<br/>adjust_load(agent, +1)"]
T6 --> T7["RelayService.dispatch_command<br/>POST {endpoint}/tasks/{id}"]
T7 --> T8{"HTTP 202?"}
T8 -- 是 --> OK(["下发成功<br/>记录 dispatch 日志"])
T8 -- 否 --> T9["回退:task 置 pending + 解绑<br/>task:running→task:pending<br/>agent 负载 -1"]
T9 --> END2
四、Agent 端任务处理流程
flowchart TB
RECV["POST {endpoint}/tasks/{request_id}"] --> AUTH{"auth == GATEWAY_AUTH?"}
AUTH -- 否 --> 401["401 invalid auth"]
AUTH -- 是 --> THREAD["启动 daemon 后台线程<br/>(立即返回 202 accepted)"]
THREAD --> PROMPT["payload → prompt<br/>(prompt / cmd / 兜底序列化)"]
PROMPT --> RUN["InMemoryRunner.run_async<br/>user_id=gateway, session=task-{id}"]
RUN --> OUT["收集最终文本输出"]
OUT --> REP["gateway_client.report_result<br/>POST /api/agent/result"]
RUN --> ERR["异常捕获"]
ERR --> REPF["report_result(status=failed, error_info)"]
五、Redis 数据结构一览
| 前缀 / 键 | 类型 | 用途 |
|---|---|---|
task:info:{request_id} |
Hash | 任务全量信息(状态/绑定 Agent/进度/结果),TTL 24h |
task:pending |
Set | 待调度任务索引 |
task:running |
Set | 执行中任务索引 |
agent:info:{agent_id} |
Hash | Agent 信息(endpoint/标签/并发上限/负载/心跳时间) |
agent:all |
Set | 全部 Agent 索引 |
agent:tag:{tag} |
Set | 能力标签 → Agent 索引(调度匹配用) |
agent:heartbeat |
ZSet | score=最后心跳时间戳(剔除离线用) |
lock:heartbeat / lock:scheduler / lock:timeout |
String | 分布式锁,防多实例重复执行后台循环 |
log:audit |
List | 全链路审计日志(dispatch/progress/result) |
六、技术栈与端口汇总
flowchart LR
subgraph 后端
B["Python + FastAPI + Uvicorn<br/>redis-py(asyncio) + httpx<br/>端口 8000"]
end
subgraph 前端
F["Vue 3 + Vite + Element Plus<br/>端口 5173(/api 代理到 8000)"]
end
subgraph Agent
A["Python + Google ADK v2.5<br/>FastAPI ApiServer + InMemoryRunner<br/>端口 8001"]
end
subgraph 存储
R["Redis(远程 45.207.192.237:56987/0)<br/>任务池 + Agent 池 + 索引 + 日志 + 锁"]
end
B <--> R
B <-->|"HTTP 推送 / 结果回传"| A
F -->|"/api/admin/*"| B
七、任务状态机
stateDiagram-v2
[*] --> pending: submit 受理(幂等)
pending --> running: 调度器选 Agent 并推送成功
pending --> pending: 无可用 Agent / 推送失败回退
running --> success: Agent 回传 result(success)
running --> failed: Agent 回传 failed
running --> failed: 超时 / 用户取消
pending --> failed: 用户取消
success --> pending: 手动重置(重新调度)
failed --> pending: 手动重置(重新调度)
success --> [*]: TTL 归档
failed --> [*]: TTL 归档