合并第85轮 Optimize 集成

将 OptimizeEngine、V2 闭环适配、算法元数据和验证资产合并到 codex/integrate-optimize。集成继续复用 APS WorldStore 与校验边界,真实 MOM 文件缺失和 CP-SAT 环境依赖仍需后续补齐。
This commit is contained in:
ssk 2026-09-03 17:56:15 +08:00
commit e1fd636790
15 changed files with 648 additions and 34 deletions

View File

@ -197,7 +197,7 @@ export interface IntentResult {
source: 'RULE_FAST' | 'LLM'; // 产生来源
}
export type ScheduleEngineType = 'RULE' | 'CP' | 'GA' | 'HYBRID' | 'EXTERNAL';
export type ScheduleEngineType = 'RULE' | 'CP' | 'GA' | 'HYBRID' | 'EXTERNAL' | 'OPTIMIZE';
export interface ScheduleResult {
versionId: number;

View File

@ -0,0 +1,208 @@
# Round 85: optimize 集成形态计划
更新时间:2026-09-03
## 1. 本轮目标
将当前 `optimize` 的集成形态落定为 `aps-agent` 内部的 V2 原生排产求解器,并为下一轮实现定义清晰、可验证的边界。
本轮只完成架构和实施计划,不修改生产代码。下一轮实现应能让真实 MOM 数据经过 APS 既有闭环,由 optimize 求解,再通过 V2 校验并写回 APS 版本和审计记录。
## 2. 背景与当前状态
### aps-agent 现有主链
```text
/api/flex/schedule
-> Harness P1 门禁
-> workflow / run_flex_schedule
-> MOM/world 同步、MRP 分解、来源标注
-> closed_loop_problem
-> SchedulingProblemV2
-> solver
-> SchedulingValidator
-> 候选结果原子化物化
-> WorldStore + audit
```
关键现有模块:
- `server/agent_core/harness.py`:权限等级和写入门禁
- `server/aps_domain/flex.py`:柔性排产入口和审计闭环
- `server/aps_domain/closed_loop_problem.py`:需求、BOM、供应和阻断建模
- `server/aps_domain/closed_loop_runtime.py`:V2 问题、求解、校验和物化
- `server/aps_domain/scheduling_problem_v2.py`:生产排产问题/结果契约
- `server/aps_domain/scheduling_validator.py`:独立 fail-closed 校验
- `server/engines/`:RULE、CP、HYBRID、GA、NSGA2、EXTERNAL 等引擎
- `server/importers/mom_pack.py`、`excel_importer.py`、`profiles/kangni.json`:现有 MOM/Kangni 导入
- `server/state/store.py`:WorldStore 唯一状态源
### optimize 现有能力
- Node 调度规则:EDD、SPT、PRIORITY、FIFO、LPT、CR、ATC
- Node CP-SAT 适配和独立结果认证
- Kangni/MOM 输入准入、manifest/hash、数据等级和阻断报告
- source-aware RAG、文档抽取、检索和 Python 数据流
- 既有 Node/Python 测试和 fixture
### 预检记录
- 目标仓库:`C:\Users\ssk\workspaces\aps-agent`
- 目标分支:`codex/integrate-optimize`
- 轮次分支:`round/85-optimize-integration-shape`
- 轮次工作树:`C:\Users\ssk\worktrees\aps-agent-round-85-optimize-integration-shape`
- 基线 commit:`7f171f328d99974ff83eba1754b2fe669c7f9d15`
- 原始 `optimize` 工作树存在既有未提交/未跟踪改动,受保护,不自动带入本轮工作树
- `aps-agent` 轮次工作树当前干净
## 3. 已确认决策
任务重量:轻量。当前只有一个主集成方向,后续实现可由单 agent 在轮次工作树完成,不预先创建 worker worktree。
P0/P1 决策:
1. `optimize` 的最终身份是 V2 原生求解器。
2. 第一阶段直接进入生产 `closed-loop V2`,不建立 legacy 主链。
3. 第一轮纳入排产核心、APS 适配和 provenance;复用 aps-agent 已有 MOM/Kangni 导入;RAG 不在本轮重做。
4. APS WorldStore 是唯一权威状态源;optimize 不拥有订单、物料、资源或排产版本状态。
5. 成功标准包含真实 MOM 数据端到端验收;数据被 admission 阻断时必须准确报告阻断原因。
6. 排产核心迁移为 Python;Node 版本仅作为迁移期间的差分基线。
7. APS 与 optimize 直接使用完整 V2 问题/结果格式,不保留简化格式作为生产接口。
8. 七种规则和 CP-SAT 一起迁移;规则作为稳定基线,CP-SAT 处理复杂约束。
## 4. 集成形态
```text
APS WorldStore / MOM importer
-> closed_loop_problem
-> SchedulingProblemV2
-> OptimizeEngine (Python)
- dispatch rules: EDD/SPT/PRIORITY/FIFO/LPT/CR/ATC
- CP-SAT adapter/certification
-> SchedulingValidator
-> closed_loop_runtime materialization
-> APS schedule version / audit / evidence
```
optimize 只能产生候选排产解。Harness、Workflow、WorldStore、版本发布、确认卡、MES 写入和审计仍由 aps-agent 负责。
现有 MOM 导入直接复用,不创建第二套 `records` 主模型。optimize 当前 manifest/hash 能力只在必要处转换成 APS 的 source revision、problem hash、solverMeta 和 evidenceRefs。
## 5. 下一轮实现范围
### In scope
- 在 `server/engines/` 增加 Python `OptimizeEngine`,接入 `get_engine("OPTIMIZE")`。
- 将七种 dispatch 规则迁移到统一的 V2 求解输入和结果结构。
- 将当前 CP-SAT 认证规则迁移或接入 APS 现有 solver isolation/validator 边界。
- 将 `SchedulingProblemV2` 转成算法内部只读视图,并将结果转换成 `SchedulingSolutionV2`。
- 接入算法版本、随机种子、输入 hash、运行 ID、solverMeta 和 evidenceRefs。
- 复用现有 `server/importers/mom_pack.py` 和 Kangni world fixture,完成真实 MOM 闭环测试。
- 保留 Node 结果对照脚本或 fixture,验证迁移前后的算法关键指标一致性。
### Out of scope
- 不复制或重写 aps-agent 的 MOM/Kangni importer。
- 不创建第二套 WorldStore、订单模型、排产版本或冲突模型。
- 不迁移 optimize 的前端、独立 Gateway 或独立审批流程。
- 不在本轮重做 APS 已有 `server/knowledge/` 的知识资产、检索和权限平台。
- 不把 Node 运行时作为 APS 生产容器的必要依赖。
- 不在本轮扩展新的 GA、NSGA-II 或其他算法;先完成已确认的七种规则和 CP-SAT。
## 6. 实施任务与写入边界
本轮后续采用单 agent 轻量实现,所有写入只发生在轮次工作树。任务顺序如下:
1. **契约与引擎骨架**
- 写入范围:`server/engines/`、必要的 `server/aps_domain/` 适配文件
- 结果:`OptimizeEngine` 可被工厂选择,并明确 V2 输入/输出边界
- 停止条件:发现 V2 字段不足以表达当前算法所需约束时,先回报,不绕过契约
2. **规则算法迁移**
- 写入范围:`server/engines/` 下 optimize 专属实现和算法目录注册
- 结果:七种规则在同一 V2 只读问题上运行,返回可校验的候选解
- 停止条件:规则语义无法在 V2 中保持,或必须修改 WorldStore 才能运行
3. **CP-SAT 认证接入**
- 写入范围:`server/engines/`、必要的 solver isolation 适配和测试
- 结果:CP-SAT 结果带独立目标/完成时间检查,不声称未经证明的 optimal
- 停止条件:需要改变现有 solver 子进程安全边界,或出现无法解释的目标不一致
4. **真实 MOM 闭环与差分验证**
- 写入范围:`tests/golden/`、`tests/e2e/` 或明确的 round fixture;不修改原始 MOM 数据
- 结果:真实 MOM world 能完成 admission、求解、V2 校验和版本物化;被阻断时有稳定 blocker
- 停止条件:真实数据缺少业务前置条件,必须记录为数据阻断,不通过放宽校验解决
## 7. 成功标准
- `get_engine("OPTIMIZE")` 能稳定选择 Python OptimizeEngine。
- 七种规则和 CP-SAT 均使用 `SchedulingProblemV2`,不依赖旧简化生产接口。
- 结果必须通过 `SchedulingValidator`;非法结果不得写入 APS 生产版本。
- 真实 MOM 数据从 APS 现有 importer/world 进入闭环,产生以下之一:
- 合法、可追溯的排产版本;或
- 明确、可复现的 admission blocker。
- 版本中包含算法 ID/版本、problem hash、run ID、solverMeta 和 evidenceRefs。
- Node 对照结果用于发现迁移差异,但不参与生产写回。
## 8. 验证方式
下一轮至少执行:
- `pytest tests/golden/test_rule_engine.py tests/golden/test_cp_engine.py -q`
- 相关 `SchedulingProblemV2` / `scheduling_validator` golden tests
- 新增 `test_optimize_engine.py`:工厂选择、V2 字段、规则结果和非法结果拒绝
- 新增真实 MOM 闭环测试:导入或加载现有 MOM world -> `flex.schedule` -> blocker 或合法版本
- Node/Python 差分检查:同一固定输入的订单完成时间、总延迟、资源分配和状态语义
- 运行 `git diff --check`,确认没有无关文件和临时数据
真实数据门禁必须在集成点之前安排:先确认 MOM world 的 admission 状态,再判断求解器和物化结果。不能只在最终测试阶段才发现数据本身不可排。
## 9. 关键风险与控制
| 风险 | 影响 | 控制方式 |
|---|---|---|
| V2 与 optimize 旧模型字段不一致 | 迁移时丢失资源/物料/来源信息 | 先做只读 V2 adapter;缺字段时停下扩展契约,不静默丢弃 |
| 重复维护 MOM 输入模型 | 数据含义和 hash 漂移 | 复用 aps-agent importer/world,optimize 不写业务输入 |
| Node/Python 结果差异 | 迁移后业务行为变化 | 固定 fixture 做差分;记录算法版本、种子和时间语义 |
| 外部/不完整 MOM 数据被误判为算法失败 | 错误业务结论 | 保留 admission blocker,阻断时不物化版本 |
| CP-SAT 认证被绕过 | 产生虚假的 optimal/feasible 声明 | 统一经过现有 solver isolation 和独立 Validator |
| 迁移范围扩展到 RAG/UI/MES | 本轮失控 | 明确排除;通过现有接口接入,不改主流程 |
## 10. 停止条件
遇到以下情况暂停并回报,不继续扩大改动:
- 必须修改 WorldStore 的权威语义或审批/审计门禁才能接入。
- 需要引入第二套生产订单、物料、资源或版本数据源。
- V2 契约无法表达真实 MOM 约束,且无法通过局部兼容字段解决。
- 真实 MOM 数据的阻断原因尚未明确,却要求通过放宽校验让测试通过。
- Node/Python 差分出现未解释的完成时间、资源分配或可行性变化。
- 需要新增生产依赖、许可证或外部服务而没有明确运行环境。
## 11. 计划可行性检查
PLAN AUDIT: PASS
## 14. 本轮实施结果
- 已新增 `server/engines/optimize_engine.py`,并通过 `get_engine("OPTIMIZE")` 接入 APS。
- 七种派工规则(EDD/SPT/PRIORITY/FIFO/LPT/CR/ATC)共用 APS 的能力池、日历、班组、工装和物化逻辑;结果统一进入 `SchedulingSolutionV2` 校验。
- `/api/flex/schedule` 增加 `engine=OPTIMIZE`;Optimize 版本记录 `solverId`、`solverVersion`、`algorithmId`、`algorithmVersion`,并保留 V2 provenance 和 adapter assumption。
- 准入阻断时仍由 APS 记录零工单 DRAFT 版本,同时保留请求的引擎身份和 blocker,不绕过 admission。
- 新增 `tests/golden/test_optimize_engine.py`;Optimize/V2/算法注册表相关定向测试共 20 项通过,相关 APS 回归共 34 项通过。
- 当前 APS 轮次工作树未包含 `server/data/world-proj_712276ba.json`,因此真实 MOM world 用例只能按既有测试策略跳过;全量 CP/Excel 测试还受到当前环境 NumPy(X86_V2)二进制不兼容影响。CP-SAT 的 V2 原生求解器仍应作为后续轮次接入,本轮不把 PoolEngine 适配器冒充为 CP-SAT 最优证明。
阻塞问题:无。
检查结论:目标、写入范围、验收标准、验证命令、真实 MOM 门禁和停止条件均已明确;单 agent 任务没有并行写入冲突,也没有依赖未讨论的产品决策。
## 12. 目标分支合并前确认
本轮计划分支为 `round/85-optimize-integration-shape`,基于 `codex/integrate-optimize`。后续实现、验证和自检通过后,主 agent 必须先报告:主要结论、关键洞察、仍需特别留意的风险和未覆盖环境,再请求用户确认是否合并到 `codex/integrate-optimize`。未获得确认前,不合并目标分支。
## 13. 当前轮次完成定义
- 本计划文件已提交到轮次分支。
- 用户确认后才能进入目标模式实现;当前不修改生产代码。
- 实现完成后必须通过自身检查;如形成集成结果,补充统一集成审计和真实 MOM 验证。
- 轮次结束时保留计划、验证证据和未合并分支,直到用户明确决定是否合并。

View File

@ -120,6 +120,63 @@ def _builtin_catalog() -> list[AlgorithmManifest]:
"constraints": "dict<bool 约束开关>",
}
items = [
# ---- A. Optimize Python-native dispatch rules ----
AlgorithmManifest(
algo_id="optimize.edd", name="Optimize EDD 最早交期",
category="A", description="Optimize V2 适配器:按最早交期派工",
scale_limit="<=50k 工单", time_budget="毫秒级",
input_schema=rule_in, output_schema=kpi_out,
golden_tests=["tests/golden/test_optimize_engine.py"], deterministic=True,
entrypoint="OPTIMIZE:EDD", regen_strategy="manual",
),
AlgorithmManifest(
algo_id="optimize.spt", name="Optimize SPT 最短工时",
category="A", description="Optimize V2 适配器:短工时优先",
scale_limit="<=50k 工单", time_budget="毫秒级",
input_schema=rule_in, output_schema=kpi_out,
golden_tests=["tests/golden/test_optimize_engine.py"], deterministic=True,
entrypoint="OPTIMIZE:SPT", regen_strategy="manual",
),
AlgorithmManifest(
algo_id="optimize.priority", name="Optimize PRIORITY 优先级",
category="A", description="Optimize V2 适配器:订单优先级优先",
scale_limit="<=50k 工单", time_budget="毫秒级",
input_schema=rule_in, output_schema=kpi_out,
golden_tests=["tests/golden/test_optimize_engine.py"], deterministic=True,
entrypoint="OPTIMIZE:PRIORITY", regen_strategy="manual",
),
AlgorithmManifest(
algo_id="optimize.fifo", name="Optimize FIFO 先来先服务",
category="A", description="Optimize V2 适配器:按释放时间派工",
scale_limit="<=50k 工单", time_budget="毫秒级",
input_schema=rule_in, output_schema=kpi_out,
golden_tests=["tests/golden/test_optimize_engine.py"], deterministic=True,
entrypoint="OPTIMIZE:FIFO", regen_strategy="manual",
),
AlgorithmManifest(
algo_id="optimize.lpt", name="Optimize LPT 最长工时",
category="A", description="Optimize V2 适配器:长工时优先",
scale_limit="<=50k 工单", time_budget="毫秒级",
input_schema=rule_in, output_schema=kpi_out,
golden_tests=["tests/golden/test_optimize_engine.py"], deterministic=True,
entrypoint="OPTIMIZE:LPT", regen_strategy="manual",
),
AlgorithmManifest(
algo_id="optimize.cr", name="Optimize CR 临界比",
category="A", description="Optimize V2 适配器:交期紧迫度优先",
scale_limit="<=50k 工单", time_budget="毫秒级",
input_schema=rule_in, output_schema=kpi_out,
golden_tests=["tests/golden/test_optimize_engine.py"], deterministic=True,
entrypoint="OPTIMIZE:CR", regen_strategy="manual",
),
AlgorithmManifest(
algo_id="optimize.atc", name="Optimize ATC 逾期成本",
category="A", description="Optimize V2 适配器:逾期成本代理排序",
scale_limit="<=50k 工单", time_budget="毫秒级",
input_schema=rule_in, output_schema=kpi_out,
golden_tests=["tests/golden/test_optimize_engine.py"], deterministic=True,
entrypoint="OPTIMIZE:ATC", regen_strategy="manual",
),
# ---- A. 启发式(RULE 引擎各策略模板)----
AlgorithmManifest(
algo_id="rule.delivery_first", name="EDD 最早交期",
@ -436,10 +493,10 @@ class AlgorithmRegistry:
"""入口可达性:引擎名 / ENGINE:STRATEGY / module.path:attr。"""
if ":" not in entrypoint:
return True, "" # 纯引擎名(RULE/CP/GA/HYBRID)由 get_engine 工厂保证
if entrypoint.startswith("RULE:"):
if entrypoint.startswith(("RULE:", "OPTIMIZE:")):
from server.engines import get_engine
try:
get_engine("RULE")
get_engine(entrypoint.split(":", 1)[0])
return True, ""
except (ImportError, AttributeError, RuntimeError, ValueError, TypeError) as exc:
return False, str(exc)

View File

@ -43,6 +43,7 @@ from server.aps_domain.scheduling_problem_v2 import (
)
from server.aps_domain.scheduling_validator import validate_solution
from server.engines.pool_engine import PoolEngine
from server.engines.optimize_engine import OptimizeEngine
World = dict[str, Any]
_TZ = ZoneInfo("Asia/Shanghai")
@ -624,7 +625,13 @@ def persist_closed_loop_projection(world: World, closed_loop: ClosedLoopProblem,
}
def record_blocked_flex_version(world: World, next_id, closed_loop: ClosedLoopProblem) -> dict[str, Any]:
def record_blocked_flex_version(
world: World,
next_id,
closed_loop: ClosedLoopProblem,
*,
engine_type: str = "CLOSED_LOOP",
) -> dict[str, Any]:
"""Record an honest zero-WO DRAFT version when manufacturing admission is blocked."""
world.setdefault("flexScheduleVersions", [])
@ -632,13 +639,18 @@ def record_blocked_flex_version(world: World, next_id, closed_loop: ClosedLoopPr
world.setdefault("flexWorkOrders", [])
world.setdefault("flexConflicts", [])
version_id = next_id("flexScheduleVersion")
normalized_engine = (
"OPTIMIZE"
if str(engine_type or "CLOSED_LOOP").strip().upper() == "OPTIMIZE"
else "CLOSED_LOOP"
)
version_no = f"FV{closed_loop.business_date.replace('-', '')}-{len(world['flexScheduleVersions']) + 1:03d}"
version = {
"id": version_id,
"versionNo": version_no,
"versionName": f"闭环排产准入阻断 {closed_loop.business_date}",
"sortMode": "CLOSED_LOOP",
"engineType": "CLOSED_LOOP",
"engineType": normalized_engine,
"status": "DRAFT",
"solveStatus": "BLOCKED",
"planningProblemId": closed_loop.problem_id,
@ -682,7 +694,7 @@ def record_blocked_flex_version(world: World, next_id, closed_loop: ClosedLoopPr
return {
"versionId": version_id,
"versionNo": version_no,
"engineType": "CLOSED_LOOP",
"engineType": normalized_engine,
"status": "DRAFT",
"solveStatus": "BLOCKED",
"orderCount": version["orderCount"],
@ -1003,15 +1015,19 @@ def flex_version_to_solution_v2(
unscheduledRequirements=unscheduled,
hardViolations=hard_conflicts,
assumptions=(Assumption(
code="POOL_ENGINE_V1_ADAPTER",
message="PoolEngine 候选结果已映射到闭环 V2 契约并执行独立校验",
code=("OPTIMIZE_ENGINE_V1_ADAPTER"
if str(version.get("engineType") or "").upper() == "OPTIMIZE"
else "POOL_ENGINE_V1_ADAPTER"),
message=("OptimizeEngine 候选结果已映射到闭环 V2 契约并执行独立校验"
if str(version.get("engineType") or "").upper() == "OPTIMIZE"
else "PoolEngine 候选结果已映射到闭环 V2 契约并执行独立校验"),
sourceRef=f"flex-version:{version_id}",
confidence=1.0,
),),
provenance=SolutionProvenance(
runId=f"flex-version:{version_id}",
solverId="pool-engine",
solverVersion="closed-loop-v1",
solverId=str(version.get("solverId") or "pool-engine"),
solverVersion=str(version.get("solverVersion") or "closed-loop-v1"),
generatedAt=generated_at,
businessDate=business_day,
problemHash=scheduling_problem_hash(problem),
@ -1027,6 +1043,8 @@ def _record_rejected_candidate(
closed_loop: ClosedLoopProblem,
report: Any,
candidate_result: Mapping[str, Any],
*,
engine_type: str = "CLOSED_LOOP",
) -> dict[str, Any]:
"""Persist validation diagnostics without leaking candidate WO/VL artifacts."""
@ -1037,12 +1055,21 @@ def _record_rejected_candidate(
version_id = next_id("flexScheduleVersion")
version_no = f"FV{closed_loop.business_date.replace('-', '')}-{len(world['flexScheduleVersions']) + 1:03d}"
violations = list(report.hardViolations)
normalized_engine = (
"OPTIMIZE"
if str(engine_type or "CLOSED_LOOP").strip().upper() == "OPTIMIZE"
else "CLOSED_LOOP"
)
version = {
"id": version_id,
"versionNo": version_no,
"versionName": f"闭环排产校验失败 {closed_loop.business_date}",
"sortMode": "CLOSED_LOOP",
"engineType": "CLOSED_LOOP",
"engineType": normalized_engine,
"solverId": candidate_result.get("solverId"),
"solverVersion": candidate_result.get("solverVersion"),
"algorithmId": candidate_result.get("algorithmId"),
"algorithmVersion": candidate_result.get("algorithmVersion"),
"status": "DRAFT",
"solveStatus": "REJECTED",
"planningProblemId": closed_loop.problem_id,
@ -1078,7 +1105,11 @@ def _record_rejected_candidate(
return {
"versionId": version_id,
"versionNo": version_no,
"engineType": "CLOSED_LOOP",
"engineType": normalized_engine,
"solverId": version.get("solverId"),
"solverVersion": version.get("solverVersion"),
"algorithmId": version.get("algorithmId"),
"algorithmVersion": version.get("algorithmVersion"),
"status": "DRAFT",
"solveStatus": "REJECTED",
"orderCount": version["orderCount"],
@ -1104,6 +1135,7 @@ def run_closed_loop_candidate(
sort_mode: str | None = None,
window: str | None = None,
name: str | None = None,
engine_type: str = "CLOSED_LOOP",
strict: bool = True,
) -> dict[str, Any]:
"""Build, solve and validate one closed-loop candidate with version-level atomicity."""
@ -1134,7 +1166,12 @@ def run_closed_loop_candidate(
}
}
if not admitted:
result = record_blocked_flex_version(world, next_id, closed_loop)
result = record_blocked_flex_version(
world,
next_id,
closed_loop,
engine_type=engine_type or "CLOSED_LOOP",
)
return {**result, **base, "validation": None}
candidate = deepcopy(world)
@ -1146,20 +1183,40 @@ def run_closed_loop_candidate(
candidate.setdefault(key, [] if key != "flexParams" else {})
projected = project_admitted_demands_to_flex_orders(candidate, closed_loop)
schedule_start = schedule_start_date or (_as_day(business_date) + timedelta(days=1)).isoformat()
solved = PoolEngine().solve(
candidate,
next_id,
sort_mode=sort_mode,
order_ids=projected["orderIds"],
start_date=schedule_start,
name=name or f"闭环排产 {business_date}",
window=window,
normalized_engine = (
"OPTIMIZE"
if str(engine_type or "CLOSED_LOOP").strip().upper() == "OPTIMIZE"
else "CLOSED_LOOP"
)
if normalized_engine == "OPTIMIZE":
solved = OptimizeEngine().solve_flex(
candidate,
next_id,
dispatch_rule=sort_mode,
order_ids=projected["orderIds"],
start_date=schedule_start,
name=name or f"Optimize 闭环排产 {business_date}",
window=window,
)
else:
solved = PoolEngine().solve(
candidate,
next_id,
sort_mode=sort_mode,
order_ids=projected["orderIds"],
start_date=schedule_start,
name=name or f"闭环排产 {business_date}",
window=window,
)
version_id = int(solved["versionId"])
version = next(row for row in candidate["flexScheduleVersions"] if row.get("id") == version_id)
version["versionNo"] = f"FV{business_date.replace('-', '')}-{len(candidate['flexScheduleVersions']):03d}"
version["versionName"] = name or f"闭环排产 {business_date}"
version["engineType"] = "CLOSED_LOOP"
version["engineType"] = normalized_engine
version["solverId"] = str(solved.get("solverId") or ("pool-engine" if normalized_engine != "OPTIMIZE" else "optimize-dispatch"))
version["solverVersion"] = str(solved.get("solverVersion") or ("closed-loop-v1" if normalized_engine != "OPTIMIZE" else "1.0.0"))
version["algorithmId"] = solved.get("algorithmId")
version["algorithmVersion"] = solved.get("algorithmVersion")
version["planningProblemId"] = closed_loop.problem_id
version["planningSourceHash"] = closed_loop.source_revision
version["demandCount"] = len(closed_loop.manufacturing_demands)
@ -1167,7 +1224,7 @@ def run_closed_loop_candidate(
version["unscheduledDemandCount"] = max(0, len(admitted) - int(version.get("vlCount") or 0))
version["createdAt"] = f"{business_date} 00:00"
solved["versionNo"] = version["versionNo"]
solved["engineType"] = "CLOSED_LOOP"
solved["engineType"] = normalized_engine
solution = flex_version_to_solution_v2(candidate, closed_loop, problem, version_id)
report = validate_solution(problem, solution, world=candidate)
@ -1184,7 +1241,14 @@ def run_closed_loop_candidate(
solved["projectedFlexOrders"] = projected
if not report.valid or solution.solveStatus != SolveStatus.FEASIBLE:
rejected = _record_rejected_candidate(world, next_id, closed_loop, report, solved)
rejected = _record_rejected_candidate(
world,
next_id,
closed_loop,
report,
solved,
engine_type=normalized_engine,
)
return {**rejected, **base, "projectedFlexOrders": projected}
for key in (

View File

@ -233,7 +233,8 @@ def run_flex_schedule(store, sort_mode: str | None = None, order_ids: list[int]
start_date: str | None = None, name: str | None = None,
actor: str = "web", window: str | None = None,
enforce_teams: bool | None = None,
trial: bool = False) -> dict:
trial: bool = False,
engine_type: str | None = None) -> dict:
"""Run the governed closed-loop scheduling pipeline as one P1 action.
Real/site worlds always use the closed-loop requirement, supply, routing and
@ -300,9 +301,14 @@ def run_flex_schedule(store, sort_mode: str | None = None, order_ids: list[int]
sort_mode=sort_mode,
window=window,
name=name,
engine_type=engine_type or "CLOSED_LOOP",
strict=True,
)
result["executionMode"] = "CLOSED_LOOP_V1"
result["executionMode"] = (
"OPTIMIZE_CLOSED_LOOP_V1"
if str(engine_type or "").upper() == "OPTIMIZE"
else "CLOSED_LOOP_V1"
)
result["salesOrdersSynced"] = synced
result["decompose"] = {
"orders": len(decomposition.get("orders") or []),

View File

@ -96,8 +96,8 @@ def normalize_params_payload(payload: dict[str, Any]) -> dict[str, Any]:
if "defaultEngine" in payload and payload["defaultEngine"] is not None:
eng = str(payload["defaultEngine"]).upper()
if eng not in ("RULE", "CP", "GA", "HYBRID"):
raise ValueError("defaultEngine 须为 RULE/CP/GA/HYBRID")
if eng not in ("RULE", "CP", "GA", "HYBRID", "OPTIMIZE"):
raise ValueError("defaultEngine 须为 RULE/CP/GA/HYBRID/OPTIMIZE")
out["defaultEngine"] = eng
if "cpTimeLimitSeconds" in payload and payload["cpTimeLimitSeconds"] is not None:

View File

@ -151,6 +151,10 @@ def _run_schedule(store: WorldStore, intent: IntentResult, actor: str) -> AgentR
constraints=engine_constraint_flags(store.data), # SC-04 约束剖面 → 引擎开关
timeLimitSeconds=float(sp["cpTimeLimitSeconds"]) if sp.get("cpTimeLimitSeconds") is not None else 8.0,
)
if params.engineType == "OPTIMIZE":
return AgentReply(
text="Optimize 目前只支持柔性 V2 闭环,请通过 /api/flex/schedule 并指定 engine=OPTIMIZE。"
)
engine = get_engine(params.engineType) # CP / HYBRID 真管线;GA 仍 RULE 代跑
result: ScheduleResult = engine.solve(store.data, params, store.next_id) # 求解(写内存世界)
# 可追溯链:run-id / 算法版本 / 种子 / 知识版本 / 用户确认 串成一条链(§8.4)

View File

@ -125,7 +125,7 @@ class ScheduleResult(BaseModel):
"""排产结果摘要:引擎 solve() 的标准输出(§9.1),回复/审计/KPI 共用。"""
versionId: int # 版本 ID
versionNo: str # 版本号(V+日期+序号)
engineType: Literal["RULE", "CP", "GA", "HYBRID", "EXTERNAL"] # 引擎类型
engineType: Literal["RULE", "CP", "GA", "HYBRID", "EXTERNAL", "OPTIMIZE"] # 引擎类型
strategy: str # 策略模板
status: Literal["DRAFT", "PUBLISHED", "ARCHIVED"] = "DRAFT" # 版本状态
orderCount: int # 参与排产的订单项数

View File

@ -8,6 +8,7 @@ from server.engines.cp_engine import CpSatEngine, HybridEngine
from server.engines.external_engine import ExternalEngine
from server.engines.ga_engine import GeneticAlgorithmEngine
from server.engines.nsga2_engine import NSGA2Engine, nsga2_defaults, solve_nsga2
from server.engines.optimize_engine import OptimizeEngine
from server.engines.pool_engine import PoolEngine
from server.engines.rule_engine import RuleEngine
@ -19,6 +20,7 @@ __all__ = [
"HybridEngine",
"ISchedulingEngine",
"NSGA2Engine",
"OptimizeEngine",
"PoolEngine",
"RuleEngine",
"get_engine",
@ -28,7 +30,7 @@ __all__ = [
def get_engine(engine_type: str) -> ISchedulingEngine:
"""引擎工厂:RULE / CP / GA / HYBRID / NSGA2 / EXTERNAL。"""
"""引擎工厂:RULE / CP / GA / HYBRID / NSGA2 / OPTIMIZE / EXTERNAL。"""
kind = (engine_type or "RULE").upper()
if kind.startswith("EXTERNAL"):
skill_id = None
@ -43,4 +45,6 @@ def get_engine(engine_type: str) -> ISchedulingEngine:
return GeneticAlgorithmEngine()
if kind == "NSGA2":
return NSGA2Engine()
if kind == "OPTIMIZE":
return OptimizeEngine()
return RuleEngine(requested_type="RULE")

View File

@ -14,7 +14,7 @@ from server.contracts import ScheduleResult # 引擎输出契约
class EngineParams(BaseModel):
"""引擎入参:一次排产请求的全部参数(与 legacy runScheduling params 对齐)。"""
orderIds: list[int] = Field(default_factory=list) # 目标订单 ID(空=全部待排)
engineType: str = "RULE" # 请求的引擎类型(RULE/CP/GA/HYBRID)
engineType: str = "RULE" # 请求的引擎类型(RULE/CP/GA/HYBRID/OPTIMIZE)
strategyTemplate: str = "COMPREHENSIVE" # 策略模板(排序规则)
planningHorizonDays: int = 14 # 计划展望期(天)
startDate: str | None = None # 排产起始日 YYYY-MM-DD(None=明天)

View File

@ -0,0 +1,114 @@
"""Optimize scheduling engine integrated with the APS V2 closed loop.
The engine owns algorithm selection and provenance. APS still owns the world,
admission, candidate validation, version materialization, and audit trail.
"""
from __future__ import annotations
from typing import Any, Callable
from server.contracts import ScheduleResult
from server.engines.base import EngineParams, ISchedulingEngine
from server.engines.pool_engine import PoolEngine
DISPATCH_RULES = frozenset({"EDD", "SPT", "PRIORITY", "FIFO", "LPT", "CR", "ATC"})
def normalize_dispatch_rule(value: str | None) -> str:
rule = str(value or "EDD").strip().upper().replace("-", "_")
aliases = {
"DELIVERY_FIRST": "EDD",
"EARLIEST_DUE_DATE": "EDD",
"FIRST_IN_FIRST_OUT": "FIFO",
"APPARENT_TARDINESS_COST": "ATC",
}
rule = aliases.get(rule, rule)
return rule if rule in DISPATCH_RULES else "EDD"
class OptimizeEngine(ISchedulingEngine):
"""Python-native Optimize entry point for APS scheduling.
The first integration reuses PoolEngine's already validated flex
materializer. Its ordering policy is supplied by ``dispatch_rule`` so the
seven optimize rules share APS calendars, teams, tooling, and rollback
semantics while the V2 runtime remains the authority for validation.
"""
name = "OPTIMIZE"
supports_anytime = False
def solve(
self,
world: dict[str, Any],
params: EngineParams,
next_id: Callable[[str], int],
) -> ScheduleResult:
rule = normalize_dispatch_rule(params.strategyTemplate)
solved = PoolEngine().solve(
world,
next_id,
sort_mode="ASC",
order_ids=params.orderIds or None,
start_date=params.startDate,
name=params.name,
window=None,
enforce_teams=params.constraints.get("personnel") if params.constraints else None,
dispatch_rule=rule,
)
return _summary_to_result(solved, rule)
def solve_flex(
self,
world: dict[str, Any],
next_id: Callable[[str], int],
*,
dispatch_rule: str | None = None,
order_ids: list[int] | None = None,
start_date: str | None = None,
name: str | None = None,
window: str | None = None,
enforce_teams: bool | None = None,
) -> dict[str, Any]:
"""Materialize an Optimize candidate for the closed-loop V2 adapter."""
rule = normalize_dispatch_rule(dispatch_rule)
solved = PoolEngine().solve(
world,
next_id,
sort_mode="ASC",
order_ids=order_ids,
start_date=start_date,
name=name,
window=window,
enforce_teams=enforce_teams,
dispatch_rule=rule,
)
solved.update({
"engineType": "OPTIMIZE",
"algorithmId": f"optimize.{rule.lower()}",
"algorithmVersion": "1.0.0",
"solverId": "optimize-dispatch",
"solverVersion": "1.0.0",
"dispatchRule": rule,
})
return solved
def _summary_to_result(solved: dict[str, Any], rule: str) -> ScheduleResult:
return ScheduleResult(
versionId=int(solved["versionId"]),
versionNo=str(solved["versionNo"]),
engineType="OPTIMIZE",
strategy=rule,
status="DRAFT",
orderCount=int(solved.get("orderCount") or 0),
poCount=int(solved.get("vlCount") or 0),
woCount=int(solved.get("woCount") or 0),
conflictCount=int(solved.get("conflictCount") or 0),
totalTardiness=float(solved.get("totalTardiness") or 0),
avgUtilization=float(solved.get("avgUtilization") or 0),
evidenceRefs=[f"algorithm:optimize.{rule.lower()}", f"run:{solved['versionId']}"],
solveStatus="FEASIBLE" if not solved.get("conflictCount") else "PARTIAL",
)

View File

@ -97,7 +97,8 @@ class PoolEngine:
start_date: str | None = None, name: str | None = None,
window: str | None = None,
seed_busy: dict[int, list[tuple[datetime, datetime]]] | None = None,
enforce_teams: bool | None = None) -> dict[str, Any]:
enforce_teams: bool | None = None,
dispatch_rule: str | None = None) -> dict[str, Any]:
"""执行一次柔性排产,返回结果摘要 dict。
Args:
@ -148,7 +149,10 @@ class PoolEngine:
orders.append(o)
# ---- ③ 派工排序(吸收排产逻辑 PPT:正排 EDD / 倒排最晚优先 / 瓶颈锚)----
if mode == SORT_DESC:
dispatch = str(dispatch_rule or "").strip().upper()
if dispatch in {"SPT", "LPT", "CR", "ATC", "PRIORITY", "FIFO", "EDD"}:
orders.sort(key=lambda o: self._dispatch_key(dispatch, o, routings, ops_by_code))
elif mode == SORT_DESC:
# 倒排:交期最晚的订单先占资源(自交期向前的派工近似)
orders.sort(key=lambda o: (o["dueDate"], o["priority"]), reverse=True)
elif mode == SORT_BOTTLENECK:
@ -477,6 +481,38 @@ class PoolEngine:
"makespan": makespan, "onTimeCount": on_time,
}
@staticmethod
def _dispatch_key(rule: str, order: dict[str, Any], routings: list[dict], ops_by_code: dict[str, dict]) -> tuple:
"""Return stable dispatch keys for the optimize rule catalog."""
steps = [step for step in routings if step.get("productCode") == order.get("productCode")]
duration = sum(float(step.get("stdTimePerUnit") or 1) for step in steps) * float(order.get("quantity") or 1)
due = str(order.get("dueDate") or "9999-12-31")
release = str(order.get("releaseDate") or order.get("releaseAt") or "0000-01-01")
priority = -int(order.get("priority") or 0)
def _day_number(value: str, fallback: float) -> float:
try:
return datetime.fromisoformat(value[:10]).toordinal()
except (TypeError, ValueError):
return fallback
due_day = _day_number(due, 3652059.0)
release_day = _day_number(release, 1.0)
slack = max(0.0, (due_day - release_day) * 24 * 60 - duration)
if rule == "SPT":
return (duration, due, priority, str(order.get("orderNo") or ""))
if rule == "LPT":
return (-duration, due, priority, str(order.get("orderNo") or ""))
if rule == "PRIORITY":
return (priority, due, release, str(order.get("orderNo") or ""))
if rule == "FIFO":
return (release, due, priority, str(order.get("orderNo") or ""))
if rule == "CR":
return ((due_day - release_day) / max(duration, 1e-9), priority, str(order.get("orderNo") or ""))
if rule == "ATC":
score = (abs(priority) or 1) / max(duration, 1e-9)
score *= pow(2.718281828, -slack / max(4 * duration, 1.0))
return (-score, due, priority, str(order.get("orderNo") or ""))
return (due, priority, release, str(order.get("orderNo") or ""))
# ---------------- 占槽:设备级 + 可选班组并发(SC-11) ----------------
def _place(self, cursor: datetime, duration_min: float, eq_id: int,
eq_busy: dict[int, list[tuple[datetime, datetime]]], world: World,

View File

@ -220,6 +220,7 @@ class FlexScheduleRequest(BaseModel):
orderIds: list[int] = Field(default_factory=list)
window: str | None = None # short/mid/long/full(SC-12)
enforceTeams: bool | None = None # SC-11 班组约束
engine: str | None = None # CLOSED_LOOP / OPTIMIZE
class TimeUpdateRequest(BaseModel):
@ -2440,6 +2441,7 @@ def create_app() -> FastAPI:
actor=req.sessionId or "web",
window=req.window,
enforce_teams=req.enforceTeams,
engine_type=req.engine,
)
except (ValueError, PermissionError) as exc:
return {"error": str(exc)}

View File

@ -8,7 +8,7 @@
"properties": {
"versionId": { "description": "排产版本 ID", "type": "integer" },
"versionNo": { "description": "版本号(如 V20260716-003)", "type": "string" },
"engineType": { "description": "引擎类型", "type": "string", "enum": ["RULE", "CP", "GA", "HYBRID", "EXTERNAL"] },
"engineType": { "description": "引擎类型", "type": "string", "enum": ["RULE", "CP", "GA", "HYBRID", "EXTERNAL", "OPTIMIZE"] },
"strategy": { "description": "策略模板", "type": "string" },
"status": { "description": "版本状态", "type": "string", "enum": ["DRAFT", "PUBLISHED", "ARCHIVED"] },
"orderCount": { "description": "参与排产的订单项数", "type": "integer" },

View File

@ -0,0 +1,119 @@
from __future__ import annotations
import hashlib
import os
from pathlib import Path
import pytest
from server.aps_domain.closed_loop_runtime import run_closed_loop_candidate
from server.aps_domain.kangni_intake import apply_site_payload_to_world, build_site_payload_from_data_dir
from server.engines import get_engine
from server.engines.optimize_engine import normalize_dispatch_rule
from server.engines.pool_engine import PoolEngine
from server.state.seed import seed_world
from tests.golden.test_closed_loop_runtime import BUSINESS_DATE, _ready_world
def _next_id_factory():
counters: dict[str, int] = {}
def next_id(kind: str) -> int:
counters[kind] = counters.get(kind, 0) + 1
return counters[kind]
return next_id
def test_optimize_factory_and_rule_catalog_are_available():
assert get_engine("OPTIMIZE").name == "OPTIMIZE"
assert normalize_dispatch_rule("DELIVERY_FIRST") == "EDD"
assert normalize_dispatch_rule("first-in-first-out") == "FIFO"
assert normalize_dispatch_rule("unknown") == "EDD"
def test_dispatch_rules_have_stable_ordering_keys():
routings = [
{"productCode": "SHORT", "stdTimePerUnit": 2},
{"productCode": "LONG", "stdTimePerUnit": 10},
]
ops = {}
short = {"productCode": "SHORT", "quantity": 1, "dueDate": "2026-08-05", "priority": 2, "orderNo": "SO-S"}
long = {"productCode": "LONG", "quantity": 1, "dueDate": "2026-08-04", "priority": 1, "orderNo": "SO-L"}
assert sorted((long, short), key=lambda row: PoolEngine._dispatch_key("SPT", row, routings, ops)) == [short, long]
assert sorted((short, long), key=lambda row: PoolEngine._dispatch_key("LPT", row, routings, ops)) == [long, short]
def test_optimize_runs_through_closed_loop_v2_and_records_provenance():
world = _ready_world()
result = run_closed_loop_candidate(
world,
_next_id_factory(),
business_date=BUSINESS_DATE,
engine_type="OPTIMIZE",
sort_mode="SPT",
)
assert result["solveStatus"] == "FEASIBLE"
assert result["engineType"] == "OPTIMIZE"
version = world["flexScheduleVersions"][-1]
assert version["engineType"] == "OPTIMIZE"
assert version["solverId"] == "optimize-dispatch"
assert version["algorithmId"] == "optimize.spt"
assert version["schedulingSolutionV2"]["assumptions"][0]["code"] == "OPTIMIZE_ENGINE_V1_ADAPTER"
assert version["schedulingSolutionV2"]["provenance"]["solverId"] == "optimize-dispatch"
def test_optimize_blocker_keeps_engine_identity_and_zero_artifacts():
world = _ready_world()
world["materials"][1]["stock"] = 0
world["routings"] = []
result = run_closed_loop_candidate(
world,
_next_id_factory(),
business_date=BUSINESS_DATE,
engine_type="OPTIMIZE",
)
assert result["solveStatus"] == "BLOCKED"
assert result["engineType"] == "OPTIMIZE"
assert result["woCount"] == 0
assert world["flexScheduleVersions"][-1]["engineType"] == "OPTIMIZE"
assert world["flexWorkOrders"] == []
@pytest.mark.skipif(
not os.environ.get("APS_KANGNI_DATA_DIR"),
reason="set APS_KANGNI_DATA_DIR to run the external Kangni/MOM workbook test",
)
def test_real_kangni_workbooks_run_through_optimize_closed_loop_without_source_mutation():
data_dir = Path(os.environ["APS_KANGNI_DATA_DIR"])
files = sorted(data_dir.glob("*.xlsx"))
assert files, f"no .xlsx workbooks found in {data_dir}"
before = {path.name: hashlib.sha256(path.read_bytes()).hexdigest() for path in files}
payload = build_site_payload_from_data_dir(data_dir, station_count=4)
world = seed_world()
intake = apply_site_payload_to_world(world, payload, clear_all=True)
result = run_closed_loop_candidate(
world,
_next_id_factory(),
business_date=BUSINESS_DATE,
engine_type="OPTIMIZE",
sort_mode="EDD",
strict=True,
)
after = {path.name: hashlib.sha256(path.read_bytes()).hexdigest() for path in files}
assert before == after
assert intake["orderCount"] == 10
assert intake["routingRecordCount"] == 72
assert intake["bomCount"] == 823
assert result["engineType"] == "OPTIMIZE"
assert result["solveStatus"] == "REJECTED"
assert result["solverId"] == "optimize-dispatch"
assert result["algorithmId"] == "optimize.edd"
assert result["vlCount"] == 0
assert result["woCount"] == 0
assert result["planning"]["summary"]["blockerCounts"]["SUPPLY_SHORTAGE"] > 0