aps-agent/server/engines/external_engine.py

333 lines
13 KiB
Python
Raw Normal View History

# ============================================================
# 外部算法引擎适配器(moduleId: engines-external, 可重生 ✅)
# ISchedulingEngine:组中性 DTO → HTTP/local skill → 写回 world
# ============================================================
from __future__ import annotations
from copy import deepcopy
from datetime import datetime
from typing import Any
from zoneinfo import ZoneInfo
import httpx
from server.aps_domain.scheduling_dto import (
SchedulingSolution, apply_flex_solution, world_to_flex_problem,
)
from server.contracts import ScheduleResult
from server.engines.base import EngineParams, ISchedulingEngine
class ExternalEngine(ISchedulingEngine):
"""外部算法排产:通过 SkillRegistry 选定 skill,HTTP 或 local://stub。"""
name = "EXTERNAL"
supports_anytime = False
def __init__(self, skill_id: str | None = None) -> None:
self.skill_id = skill_id
def solve(self, world: dict[str, Any], params: EngineParams, next_id) -> ScheduleResult:
from server.agent_core.skills import get_skills
skill_id = self.skill_id or (
params.meta_skill_id if hasattr(params, "meta_skill_id") else None
)
skill_id = skill_id or world.pop("_externalSkillId", None)
registry = get_skills()
skill = (
registry.get_enabled(skill_id, track="flex")
or registry.get_enabled(skill_id)
or registry.get_enabled()
)
if not skill:
raise ValueError("No enabled external scheduling skill is available")
order_ids = list(params.orderIds or []) or None
start_time = (
params.startDate + " 08:00"
if params.startDate and " " not in params.startDate
else params.startDate
)
if _is_explicit_demo_world(world):
problem = world_to_flex_problem(
world,
order_ids=order_ids,
start_time=start_time,
)
solution = self._call_skill(skill, problem.model_dump())
summary = apply_flex_solution(
world,
solution,
next_id,
skill_id=skill["skill_id"],
)
summary["executionMode"] = "EXPLICIT_DEMO_LEGACY"
else:
summary, solution, problem = _run_closed_loop_external_candidate(
world,
next_id,
skill,
order_ids=order_ids,
start_time=start_time,
)
evidence_refs = [f"skill:{skill['skill_id']}", f"run:{solution.runId}"]
if summary.get("planningProblemId"):
evidence_refs.append(f"planning-problem:{summary['planningProblemId']}")
return ScheduleResult(
versionId=summary["versionId"],
versionNo=summary["versionNo"],
engineType="EXTERNAL",
strategy=f"EXTERNAL:{skill['skill_id']}",
status="DRAFT",
orderCount=len(problem.orders),
poCount=summary["vlCount"],
woCount=summary["woCount"],
conflictCount=summary["conflictCount"],
totalTardiness=float((solution.kpi or {}).get("totalTardiness") or 0),
avgUtilization=float(summary.get("avgUtilization") or 0),
evidenceRefs=evidence_refs,
solveStatus=str(summary.get("solveStatus") or solution.status),
solveTimeSec=float((solution.solverMeta or {}).get("timeSec") or 0) or None,
)
def _call_skill(self, skill: dict[str, Any], problem: dict) -> SchedulingSolution:
ep = skill.get("endpoint") or ""
if ep.startswith("local://"):
from server.integrations.algo_skill_stub import solve_problem
return solve_problem(problem)
url = ep.rstrip("/") + "/schedule"
headers = {"Content-Type": "application/json"}
if skill.get("auth"):
headers["Authorization"] = f"Bearer {skill['auth']}"
timeout = float(skill.get("timeout_sec") or 60)
try:
with httpx.Client(timeout=timeout) as client:
resp = client.post(url, json=problem, headers=headers)
resp.raise_for_status()
return SchedulingSolution(**resp.json())
except Exception as exc:
raise ValueError(f"外部算法 skill「{skill.get('name')}」调用失败:{exc}") from exc
def _is_explicit_demo_world(world: dict[str, Any]) -> bool:
return any(
str(factory.get("name") or "") == "\u6f14\u793a\u5de5\u5382"
for factory in world.get("factories") or []
)
def _run_closed_loop_external_candidate(
world: dict[str, Any],
next_id,
skill: dict[str, Any],
*,
order_ids: list[int] | None,
start_time: str | None,
) -> tuple[dict[str, Any], SchedulingSolution, Any]:
"""Validate a legacy Skill result against V2 before materializing it."""
from server.aps_domain.closed_loop_problem import build_closed_loop_problem
from server.aps_domain.closed_loop_runtime import (
closed_loop_to_problem_v2,
flex_version_to_solution_v2,
persist_closed_loop_projection,
project_admitted_demands_to_flex_orders,
)
from server.aps_domain.scheduling_problem_v2 import SolveStatus
from server.aps_domain.scheduling_validator import validate_solution
selected_order_nos = None
if order_ids:
selected = set(order_ids)
selected_order_nos = sorted({
str(row.get("salesOrderNo") or row.get("orderNo"))
for row in world.get("flexOrders") or []
if row.get("id") in selected and (row.get("salesOrderNo") or row.get("orderNo"))
})
business_date = str(
world.get("businessDate")
or datetime.now(ZoneInfo("Asia/Shanghai")).date().isoformat()
)[:10]
closed_loop = build_closed_loop_problem(
world,
business_date=business_date,
order_nos=selected_order_nos,
strict=True,
)
problem_v2 = closed_loop_to_problem_v2(world, closed_loop)
persist_closed_loop_projection(world, closed_loop, problem_v2)
admitted = [
row
for row in closed_loop.manufacturing_demands
if row.release_status in {"READY", "READY_FOR_SCHEDULING"}
and row.routing_status == "READY"
and row.resource_status == "READY"
]
if not admitted:
blocker_counts = closed_loop.stats.get("blockerCounts") or {}
details = ", ".join(f"{code}={count}" for code, count in sorted(blocker_counts.items())[:5])
raise ValueError(f"External Skill scheduling admission blocked: {details or 'no admitted MAKE demand'}")
candidate = deepcopy(world)
projected = project_admitted_demands_to_flex_orders(candidate, closed_loop)
legacy_problem = world_to_flex_problem(
candidate,
order_ids=projected["orderIds"],
start_time=start_time,
)
engine = ExternalEngine(skill_id=skill["skill_id"])
legacy_solution = engine._call_skill(skill, legacy_problem.model_dump())
summary = apply_flex_solution(
candidate,
legacy_solution,
next_id,
skill_id=skill["skill_id"],
)
version_id = int(summary["versionId"])
version = next(
row for row in candidate["flexScheduleVersions"] if row.get("id") == version_id
)
business_version_no = f"EXT{business_date.replace('-', '')}-{version_id:03d}"
version.update({
"versionNo": business_version_no,
"createdAt": f"{business_date} 00:00",
"orderCount": len(projected["orderIds"]),
"demandCount": len(closed_loop.manufacturing_demands),
"admittedDemandCount": len(admitted),
"unscheduledDemandCount": max(0, len(admitted) - int(summary.get("vlCount") or 0)),
"planningProblemId": closed_loop.problem_id,
"planningSourceHash": closed_loop.source_revision,
})
solution_v2 = flex_version_to_solution_v2(
candidate,
closed_loop,
problem_v2,
version_id,
)
report = validate_solution(problem_v2, solution_v2, world=candidate)
if not report.valid or solution_v2.solveStatus != SolveStatus.FEASIBLE:
details = "; ".join(
f"{issue.code}:{issue.message}"
for issue in report.hardViolations[:5]
)
raise ValueError(
f"External Skill solution rejected by SchedulingProblemV2 validator: {details or solution_v2.solveStatus.value}"
)
version["schedulingSolutionV2"] = solution_v2.model_dump(mode="json")
version["validationReport"] = report.model_dump(mode="json")
version["solveStatus"] = solution_v2.solveStatus.value
version["engineType"] = "EXTERNAL"
for key in (
"flexOrders",
"flexRoutings",
"flexScheduleVersions",
"flexVirtualLines",
"flexWorkOrders",
"flexConflicts",
):
world[key] = candidate[key]
summary.update({
"versionNo": business_version_no,
"executionMode": "CLOSED_LOOP_V2_VALIDATED",
"solveStatus": solution_v2.solveStatus.value,
"planningProblemId": closed_loop.problem_id,
"planningSourceHash": closed_loop.source_revision,
"validation": version["validationReport"],
"projectedFlexOrders": projected,
})
return summary, legacy_solution, legacy_problem
def run_external_flex(store, *, skill_id: str | None = None,
order_ids: list[int] | None = None, actor: str = "planner") -> dict[str, Any]:
"""Run external flexible scheduling through the closed-loop V2 validator."""
from datetime import timedelta
from server.agent_core.audit import write_audit
from server.agent_core.skills import get_skills
from server.aps_domain.mrp import decompose_orders
from server.aps_domain.orders import sync_flex_orders_to_sales
from server.aps_domain.sourcing import annotate_world_sourcing
from server.timeutil import add_minutes, fmt_dt, today0
registry = get_skills()
skill = (
registry.get_enabled(skill_id, track="flex")
or registry.get_enabled(skill_id)
or registry.get_enabled()
)
if not skill:
raise ValueError("No enabled external scheduling skill is available")
if _is_explicit_demo_world(store.data):
problem = world_to_flex_problem(
store.data,
order_ids=order_ids,
start_time=fmt_dt(add_minutes(today0(), 24 * 60)),
)
engine = ExternalEngine(skill_id=skill["skill_id"])
solution = engine._call_skill(skill, problem.model_dump())
summary = apply_flex_solution(
store.data,
solution,
store.next_id,
skill_id=skill["skill_id"],
)
summary["executionMode"] = "EXPLICIT_DEMO_LEGACY"
else:
synced = sync_flex_orders_to_sales(store.data)
annotate_world_sourcing(store.data)
decomposition = decompose_orders(store.data, store.next_id)
business_date = str(
store.data.get("businessDate")
or datetime.now(ZoneInfo("Asia/Shanghai")).date().isoformat()
)[:10]
start_day = datetime.fromisoformat(business_date).date() + timedelta(days=1)
summary, solution, problem = _run_closed_loop_external_candidate(
store.data,
store.next_id,
skill,
order_ids=order_ids,
start_time=f"{start_day.isoformat()} 08:00",
)
summary["salesOrdersSynced"] = synced
summary["decompose"] = {
"orders": len(decomposition.get("orders") or []),
"make": len(decomposition.get("make") or []),
"purchase": len(decomposition.get("purchase") or []),
"outsource": len(decomposition.get("outsource") or []),
}
write_audit(
store.data,
store.next_id,
actor=actor,
category="ALGO_RUN",
action="external.flex.schedule",
target={"type": "FLEX_VERSION", "id": summary["versionId"]},
power="P1",
rationale={
"skillId": skill["skill_id"],
"runId": solution.runId,
"executionMode": summary.get("executionMode"),
"solveStatus": summary.get("solveStatus") or solution.status,
"planningProblemId": summary.get("planningProblemId"),
"planningSourceHash": summary.get("planningSourceHash"),
"validationValid": (summary.get("validation") or {}).get("valid"),
},
evidence_refs=[
f"skill:{skill['skill_id']}",
f"run:{solution.runId}",
*(
[f"planning-problem:{summary['planningProblemId']}"]
if summary.get("planningProblemId")
else []
),
],
)
store.save()
return summary