PLM 关系引擎协同架构设计文档

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.js73A2A消息→PLM协同任务
a2a-protocol.js227MACP会话创建→PLM协同会话
a2a-protocol.js347MACP会话关闭→PLM关闭会话
task-force.js531编队完成→PLM关闭会话
task-force.js642编队启动→PLM创建会话
lite-scheduler.js_runDecomposedSteps每步→checkpoint保存
server.js5149-522410+个协同API端点

8.4 降级策略

PLM 可用性检测:
  ├── 启动时 ping SCSAI → 缓存结果
  ├── 每60s 重检
  └── 操作时若 PLM 失败 → 降级 SQLite + 记录警告

数据一致性:
  ├── PLM 是权威源(正确性)
  ├── SQLite 是本地缓存(效率)
  └── 写入顺序: PLM → SQLite(PLM成功才写SQLite)
  └── 降级时: SQLite only(PLM恢复后异步补偿同步)
← 返回案例列表
分享:
🤖 Try Now →
🤖
🎁