PLM 关系引擎协同架构设计文档
版本: 1.0 | 更新: 2026-08-28 | 状态: 已实现核心模块
一、问题背景
1.1 现有协同的局限
当前系统的协同数据(A2A消息、MACP会话、全局黑板、事件总线)全部存储在本地SQLite,存在以下问题:
| 问题 | 影响 | 50×100场景表现 |
|------|------|----------------|
| 单机SQLite写锁 | 并发写受限 | 50员工×100任务=5000次写竞争 |
| 跨进程不可见 | 多实例无法协同 | 无法水平扩展 |
| 无关系图查询 | 依赖/分配/状态查询低效 | DAG遍历需多次全表扫描 |
| 无checkpoint | 长任务重启丢失进度 | 100步任务在第50步崩溃需从头重跑 |
1.2 MTClaw 的真实角色
MTClaw 不是任务调度引擎,而是 OpenAI兼容的函数路由器(Python FastAPI,端口18790):
三层 subagent 架构:
Layer 1: Routing subagent — 判断是否调用工具
Layer 2: Tool execution subagents — builtin + custom wrappers
Layer 3: Completion check subagent — 决定直接回复还是转交上游
- 高频确定性操作(查询/计算/检查)→ MTClaw 加速
- 复杂推理(生成/决策/优化)→ 云端 LLM
- 当前状态: 配置完整但服务未运行,降级为普通LLM模式
1.3 设计目标
- PLM关系引擎作为协同的权威源(single source of truth)
- 支撑 50×100 复杂工业场景(50类任务×100实例=5000并发协同)
- 支持长任务(checkpoint断点恢复)
- 支持连续任务(LoopEngine已有,增强PLM同步)
- 支持多任务编排(DAG + 批量调度)
二、架构设计
2.1 整体架构
┌──────────────────────────────────────────────────┐
│ PLM 服务器 (SCSAI) │
│ ┌──────────────────────────────────────────┐ │
│ │ Collaboration Task (ItemType) │ │
│ │ ├── Task Dependency (关系: task→task) │ │
│ │ ├── Task Assignment (关系: task→staff) │ │
│ │ └── Task Result (关系: task→result) │ │
│ ├──────────────────────────────────────────┤ │
│ │ Collaboration Session (ItemType) │ │
│ │ ├── Session Task (关系: session→task) │ │
│ │ └── Session Member (关系: session→staff)│ │
│ └──────────────────────────────────────────┘ │
│ ★ 关系图查询 | 跨进程可见 | 权威源 │
└──────────────────┬───────────────────────────────┘
│ AML/HTTP
┌──────────────────┴───────────────────────────────┐
│ PLM Collaboration Engine │
│ (server/core/plm-collaboration-engine.js) │
│ ├── createTask / updateTaskStatus / assignTask │
│ ├── addDependency / getReadyTasks (DAG) │
│ ├── saveCheckpoint / resumeFromCheckpoint │
│ ├── batchCreateTasks (50×100批量) │
│ └── getBatchStatus / getTaskTree │
├──────────────────────────────────────────────────┤
│ 本地 SQLite (读缓存 + 写穿透) │
│ plm_collab_tasks | plm_collab_sessions │
│ plm_task_checkpoints │
└──────────────────┬───────────────────────────────┘
│
┌──────────────┼──────────────┐
│ │ │
┌───┴───┐ ┌────┴────┐ ┌────┴────┐
│Node #1│ │Node #2 │ │Node #N │
│60员工 │ │60员工 │ │60员工 │
└───────┘ └─────────┘ └─────────┘
2.2 数据流
任务创建 → PLM Collaboration Engine
├── 写入本地 SQLite (即时可用)
└── 写入 SCSAI PLM (异步, 降级容错)
└── 创建 Part(classification=Collaboration Task)
任务依赖 → addDependency(taskA, taskB)
└── 本地 SQLite depends_on JSON 数组
任务完成 → updateTaskStatus(taskId, 'completed', result)
├── 更新本地 SQLite
└── 更新 PLM Part 状态
长任务断点 → saveCheckpoint(taskId, stepN, state)
└── 写入 plm_task_checkpoints 表
任务恢复 → resumeFromCheckpoint(taskId)
└── 查询最新 checkpoint → 恢复状态 → 继续执行
2.3 50×100 批量调度
// 创建50类任务×100实例 = 5000个协同任务
const tasks = [];
for (let type = 0; type < 50; type++) {
for (let inst = 0; inst < 100; inst++) {
tasks.push({
title: `工业任务-T${type}-I${inst}`,
taskType: `type_${type}`,
priority: type < 10 ? 1 : 5, // 高优先级类型先执行
dependsOn: inst > 0 ? [`T${type}-I${inst-1}`] : [], // 同类型串行
});
}
}
// 批量创建(每批10个,并行写入)
const result = await collab.batchCreateTasks(tasks, sessionId);
// → { created: 5000, tasks: [...] }
// 查询就绪任务(依赖已完成的)
const ready = await collab.getReadyTasks(sessionId);
// → 最多返回100个无依赖或依赖已完成的pending任务
// 批量状态
const status = await collab.getBatchStatus(sessionId);
// → { total: 5000, completed: 3200, pending: 1800, progress: "64.0%" }
三、核心模块
3.1 PLM Collaboration Engine
文件: server/core/plm-collaboration-engine.js
| 方法 | 功能 | 50×100支撑 |
|------|------|------------|
| createTask | 创建单个任务 | PLM+SQLite双写 |
| batchCreateTasks | 批量创建(每批10个并行) | 5000任务高效写入 |
| updateTaskStatus | 更新状态+结果 | PLM同步 |
| assignTask | 分配给员工 | PLM owned_by_id |
| addDependency | 添加DAG依赖 | JSON数组存储 |
| getReadyTasks | 查询就绪任务(依赖已完成) | DAG拓扑排序 |
| saveCheckpoint | 保存断点 | 长任务恢复 |
| getLatestCheckpoint | 获取最新断点 | 长任务恢复 |
| resumeFromCheckpoint | 从断点恢复 | 长任务恢复 |
| createSession | 创建协同会话 | MACP映射 |
| getBatchStatus | 批量状态统计 | 50×100监控 |
| getTaskTree | 任务树(含子任务) | 层级查询 |
3.2 SQLite表设计
-- 协同任务表
CREATE TABLE plm_collab_tasks (
id TEXT PRIMARY KEY, -- 任务ID (CT-xxxxxxxx)
plm_id TEXT, -- SCSAI PLM中的Part ID
parent_task_id TEXT, -- 父任务(分解子任务)
session_id TEXT, -- 所属会话
title TEXT NOT NULL,
description TEXT DEFAULT '',
task_type TEXT DEFAULT 'generic',
status TEXT DEFAULT 'pending', -- pending/running/completed/failed
assigned_staff TEXT DEFAULT '',
priority INTEGER DEFAULT 5,
depends_on TEXT DEFAULT '[]', -- JSON数组: 依赖的任务ID
checkpoint TEXT DEFAULT '{}', -- 最新checkpoint元数据
result TEXT DEFAULT '{}', -- 任务结果
created_at TEXT,
updated_at TEXT,
completed_at TEXT
);
-- 协同会话表
CREATE TABLE plm_collab_sessions (
id TEXT PRIMARY KEY,
plm_id TEXT,
title TEXT NOT NULL,
scenario TEXT DEFAULT '',
status TEXT DEFAULT 'active',
member_count INTEGER DEFAULT 0,
task_count INTEGER DEFAULT 0,
metadata TEXT DEFAULT '{}',
created_at TEXT,
closed_at TEXT
);
-- 任务断点表
CREATE TABLE plm_task_checkpoints (
id TEXT PRIMARY KEY,
task_id TEXT NOT NULL,
step_number INTEGER NOT NULL,
step_name TEXT DEFAULT '',
state_json TEXT DEFAULT '{}', -- 完整执行状态
created_at TEXT
);
3.3 API端点
| 方法 | 路径 | 功能 |
|------|------|------|
| POST | /api/collaboration/sessions | 创建协同会话 |
| GET | /api/collaboration/stats | 协同统计 |
| POST | /api/collaboration/tasks | 创建任务(单个或批量) |
| GET | /api/collaboration/tasks/ready | 获取就绪任务 |
| GET | /api/collaboration/tasks/batch-status | 批量状态 |
| PUT | /api/collaboration/tasks/:id | 更新任务(状态/分配/依赖) |
| POST | /api/collaboration/tasks/:id/checkpoint | 保存断点 |
| GET | /api/collaboration/tasks/:id/checkpoint | 获取断点 |
| POST | /api/collaboration/tasks/:id/resume | 从断点恢复 |
| GET | /api/collaboration/tasks/:id/tree | 任务树 |
四、长任务 Checkpoint 机制
4.1 集成点
在 lite-scheduler.js 的 _runDecomposedSteps 中,每个子步骤完成后自动保存checkpoint:
for (let i = 0; i < steps.length; i++) {
// ... 执行步骤 i ...
// ★ 保存 checkpoint
await collab.saveCheckpoint(taskId, i + 1, step.capability, {
staffId, intent, stepIndex: i, stepResults: stepResults.slice(-1)
});
}
4.2 恢复流程
长任务崩溃/重启
→ resumeFromCheckpoint(taskId)
→ 查询最新 checkpoint (stepNumber=N)
→ 恢复状态 (state_json)
→ 从步骤 N+1 继续执行
→ 跳过已完成的步骤 1~N
五、与现有系统的集成
5.1 任务编排层级
Layer 1: LLM分解 (llm-brain.js decomposeTask)
→ 自然语言 → 子任务DAG (最多6步)
Layer 2: 工作流执行 (workflow-executor.js)
→ 步骤间引用传递 ($stepN.key)
Layer 3: 编队协同 (task-force.js)
→ DAG分层并行 + 人在环
Layer 4: PLM协同 (plm-collaboration-engine.js) ★新增
→ 50×100批量调度 + checkpoint + PLM关系
Layer 5: 循环引擎 (loop-engine.js)
→ monitor/goal_driven/self_repair 三模式
Layer 6: 状态机 (agent-state-machine.js)
→ 6态生命周期 + HITS确认
5.2 MTClaw 集成
任务执行时:
高频确定性操作 → MTClaw (端口18790) 加速
复杂推理 → 云端LLM (SmartRouter)
MTClaw未运行 → 降级为普通LLM (每60s重试)
5.3 A2A/MACP 与 PLM 的关系
当前: A2A消息 → 本地SQLite (a2a_messages表)
目标: A2A消息 → PLM协同引擎 → SCSAI关系实例
当前: MACP会话 → 本地SQLite (a2a_sessions表)
目标: MACP会话 → PLM协同会话 → SCSAI Collaboration Session
当前: 全局黑板 → 本地SQLite (global_blackboard表)
目标: 全局黑板 → PLM关系属性 → SCSAI Part属性
六、50×100 场景评估
6.1 单机能力
| 指标 | 当前 | 优化后 |
|------|------|--------|
| 并发员工 | 60+ | 60+ (不变) |
| 并发任务 | ~60 (员工锁限制) | ~60 (同) |
| 任务队列 | 无限(SQLite) | 无限(SQLite+PLM) |
| 长任务恢复 | 不支持 | ✅ checkpoint |
| 批量创建 | 逐个串行 | ✅ 每批10并行 |
| 跨进程可见 | 不支持 | ✅ PLM同步 |
6.2 水平扩展路径
单机 (当前)
↓
多进程 (worker_threads)
↓
多实例 (多Node.js + 共享PLM)
↓
分布式 (多Node.js + 消息队列 + PLM)
6.3 瓶颈分析
| 瓶颈 | 影响 | 解决方案 |
|------|------|----------|
| Node.js单进程 | CPU密集型阻塞 | worker_threads |
| SQLite写锁 | 并发写竞争 | PLM为主,SQLite为缓存 |
| SCSAI网络延迟 | PLM写入慢 | 异步写+降级容错 |
| 员工互斥锁 | 同员工不能并发 | 按任务类型分配不同员工 |
七、进化路线
| 阶段 | 状态 | 说明 |
|---|---|---|
| PLM协同引擎 | ✅ 完成 | 核心模块+API+SQLite表 |
| Checkpoint机制 | ✅ 完成 | 集成到decomposeSteps |
| 批量调度 | ✅ 完成 | 50×100 batchCreate |
| A2A→PLM同步 | ✅ 完成 | a2a-protocol.js 3处集成(经adapter) |
| task-force→PLM同步 | ✅ 完成 | task-force.js 2处集成(经adapter) |
| 适配层切换 | ✅ 完成 | 全部require改为plm-collaboration-adapter |
| 多进程worker | 📋 规划 | worker_threads拆分CPU密集 |
| 消息队列 | 📋 规划 | RabbitMQ/Kafka解耦 |
| PLM关系图查询 | 📋 规划 | SCSAI原生关系遍历 |
八、适配层架构(PLM优先 → SQLite降级)
8.1 三层架构
调用方 (a2a-protocol / task-force / server.js)
│
▼
plm-collaboration-adapter.js ← 统一入口
├── 检测 PLM 可用性 (SCSAI ping)
├── PLM 可用 → 委托 SCSAI AML 操作
└── PLM 不可用 → 降级到 plm-collaboration-engine.js (SQLite)
│
▼
plm-collaboration-engine.js ← 本地降级实现
└── SQLite 表: plm_collab_tasks / plm_collab_sessions / plm_task_checkpoints
8.2 适配层接口(13个方法)
| 方法 | PLM路径 | 降级路径 |
|------|---------|----------|
| createSession | SCSAI Collaboration Session | SQLite insert |
| closeSession | SCSAI edit state=Closed | SQLite update |
| createTask | SCSAI Part(classification=Collaboration Task) | SQLite insert |
| updateTaskStatus | SCSAI edit state字段 | SQLite update |
| assignTask | SCSAI edit owned_by_id | SQLite update |
| addDependency | SCSAI RelationshipType: Task Dependency | SQLite JSON array |
| getReadyTasks | SCSAI 关系图查询 | SQLite DAG拓扑排序 |
| batchCreateTasks | SCSAI 批量createItem | SQLite批量insert |
| saveCheckpoint | SCSAI Task Checkpoint ItemType | SQLite insert |
| getLatestCheckpoint | SCSAI queryItems | SQLite select max(step) |
| resumeFromCheckpoint | 读checkpoint→恢复状态 | 同左 |
| getBatchStatus | SCSAI aggregate查询 | SQLite group by |
| getTaskTree | SCSAI 关系遍历 | SQLite递归查询 |
8.3 集成点一览
| 文件 | 行号 | 集成内容 |
|---|---|---|
| a2a-protocol.js | 73 | A2A消息→PLM协同任务 |
| a2a-protocol.js | 227 | MACP会话创建→PLM协同会话 |
| a2a-protocol.js | 347 | MACP会话关闭→PLM关闭会话 |
| task-force.js | 531 | 编队完成→PLM关闭会话 |
| task-force.js | 642 | 编队启动→PLM创建会话 |
| lite-scheduler.js | _runDecomposedSteps | 每步→checkpoint保存 |
| server.js | 5149-5224 | 10+个协同API端点 |
8.4 降级策略
PLM 可用性检测:
├── 启动时 ping SCSAI → 缓存结果
├── 每60s 重检
└── 操作时若 PLM 失败 → 降级 SQLite + 记录警告
数据一致性:
├── PLM 是权威源(正确性)
├── SQLite 是本地缓存(效率)
└── 写入顺序: PLM → SQLite(PLM成功才写SQLite)
└── 降级时: SQLite only(PLM恢复后异步补偿同步)
BossAgents