aps-agent/server/aps_domain/mes.py

322 lines
14 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

# ============================================================
# MES 下发与报工(moduleId: domain-mes, 可重生 ✅)
# EX-05 下发 / EX-09 报工回流:经 Mock MES 桩,幂等键防重复。
# ============================================================
from __future__ import annotations
from typing import Any
from server.integrations.mes_stub import get_mes_client
from server.timeutil import fmt_date, today0
World = dict[str, Any]
def mes_connection_status() -> dict:
return get_mes_client().status()
def _latest_flex(world: World) -> dict | None:
vers = world.get("flexScheduleVersions") or []
return vers[-1] if vers else None
def _latest_fixed(world: World, *, published_only: bool = True) -> dict | None:
vers = world.get("scheduleVersions") or []
if published_only:
pubs = [v for v in vers if v.get("status") == "PUBLISHED"]
return pubs[-1] if pubs else None
return vers[-1] if vers else None
def _linked_keys(world: World) -> set[str]:
return {l.get("idemKey") for l in world.get("mesLinks", []) if l.get("kind") == "dispatch"}
def preview_dispatch(world: World, track: str = "flex") -> dict:
"""预览可下发工单(不改世界)。track ∈ flex|fixed。"""
from server.state.seed import ensure_flex_seed
client = get_mes_client()
tr = (track or "flex").lower()
if tr == "flex":
ensure_flex_seed(world)
ver = _latest_flex(world)
if not ver:
return {"track": "flex", "connection": client.status(), "items": [],
"summary": "无柔性排产版本可下发"}
vid = ver["id"]
wos = [w for w in world.get("flexWorkOrders", [])
if w.get("versionId") == vid and not w.get("frozen")]
items = []
linked = _linked_keys(world)
for w in wos[:80]:
idem = f"flex:{vid}:{w['id']}"
items.append({
"woId": w["id"], "orderNo": w.get("flexOrderNo"),
"operation": w.get("operationCode"), "equipment": w.get("equipmentCode"),
"start": w.get("plannedStartTime"), "end": w.get("plannedEndTime"),
"idemKey": idem, "already": idem in linked,
"mesExternalId": w.get("mesExternalId"),
})
new_n = sum(1 for i in items if not i["already"])
return {
"track": "flex", "connection": client.status(),
"versionId": vid, "versionNo": ver.get("versionNo"),
"versionStatus": ver.get("status"),
"items": items, "newCount": new_n,
"summary": (f"柔性版本 {ver.get('versionNo')} 可新下发 {new_n} / "
f"共 {len(items)} 条工序"),
}
ver = _latest_fixed(world, published_only=True)
if not ver:
return {"track": "fixed", "connection": client.status(), "items": [],
"summary": "无已发布固定排产版本(请先「发布版本」)"}
vid = ver["id"]
wos = [w for w in world.get("workOrders", []) if w.get("versionId") == vid]
linked = _linked_keys(world)
items = []
for w in wos[:80]:
idem = f"fixed:{vid}:{w['id']}"
items.append({
"woId": w["id"], "orderNo": w.get("productionOrderNo") or w.get("orderNo"),
"operation": w.get("operationCode") or w.get("operationName"),
"equipment": w.get("workstationName") or w.get("lineName"),
"start": w.get("plannedStartTime"), "end": w.get("plannedEndTime"),
"idemKey": idem, "already": idem in linked,
"mesExternalId": w.get("mesExternalId"),
})
new_n = sum(1 for i in items if not i["already"])
return {
"track": "fixed", "connection": client.status(),
"versionId": vid, "versionNo": ver.get("versionNo"),
"versionStatus": ver.get("status"),
"items": items, "newCount": new_n,
"summary": (f"固定版本 {ver.get('versionNo')} 可新下发 {new_n} / "
f"共 {len(items)} 条工序"),
}
def apply_dispatch(store, track: str = "flex", actor: str = "web") -> dict:
"""执行下发:幂等推送 Mock MES,回写 mesExternalId。"""
from server.agent_core.audit import write_audit
from server.state.seed import ensure_flex_seed
tr = (track or "flex").lower()
if tr == "flex":
ensure_flex_seed(store.data)
world = store.data
preview = preview_dispatch(world, tr)
client = get_mes_client()
created, duped = [], []
wo_by_id = {}
if tr == "flex":
wo_by_id = {w["id"]: w for w in world.get("flexWorkOrders", [])}
else:
wo_by_id = {w["id"]: w for w in world.get("workOrders", [])}
for it in preview.get("items") or []:
if it.get("already"):
duped.append(it["idemKey"])
continue
payload = {
"track": tr, "apsWoId": it["woId"], "orderNo": it.get("orderNo"),
"operation": it.get("operation"), "equipment": it.get("equipment"),
"start": it.get("start"), "end": it.get("end"),
"versionId": preview.get("versionId"), "versionNo": preview.get("versionNo"),
}
r = client.create_work_order(payload, it["idemKey"])
ext = r["externalWo"]
if r["duplicate"]:
duped.append(it["idemKey"])
else:
created.append(ext["id"])
wo = wo_by_id.get(it["woId"])
if wo:
wo["mesExternalId"] = ext["id"]
wo["mesIdemKey"] = it["idemKey"]
if wo.get("status") in (None, "PENDING", "DRAFT"):
wo["status"] = "RELEASED"
wo.setdefault("progressPct", 0)
wo.setdefault("qtyDone", 0)
if not r["duplicate"]:
world.setdefault("mesLinks", []).append({
"kind": "dispatch", "idemKey": it["idemKey"],
"externalWoId": ext["id"], "woId": it["woId"], "track": tr,
"versionId": preview.get("versionId"),
"syncedAt": fmt_date(today0()), "actor": actor,
})
# 版本标记已下发
vid = preview.get("versionId")
if tr == "flex" and vid:
for v in world.get("flexScheduleVersions", []):
if v["id"] == vid:
v["status"] = "DISPATCHED"
v["dispatchedAt"] = fmt_date(today0())
break
elif tr == "fixed" and vid:
for v in world.get("scheduleVersions", []):
if v["id"] == vid:
v["dispatchedAt"] = fmt_date(today0())
v["mesDispatched"] = True
break
journal = {
"id": store.next_id("mesJournal"),
"direction": "dispatch", "track": tr, "at": fmt_date(today0()),
"actor": actor, "created": len(created), "duplicates": len(duped),
"versionNo": preview.get("versionNo"),
"summary": f"下发 {len(created)} 新 / {len(duped)} 跳过",
}
world.setdefault("mesSyncJournal", []).append(journal)
write_audit(world, store.next_id, actor=actor, category="INTEGRATION",
action="mes.dispatch", target={"type": "MES", "id": tr},
power="P3", rationale={"created": len(created), "duped": len(duped)})
store.save()
return {
"journal": journal, "created": created, "duplicates": duped,
"message": (f"MES 下发完成 ✅ 新外部工单 {len(created)},"
f"幂等跳过 {len(duped)}({preview.get('versionNo')})。"),
}
def list_execution(world: World, track: str = "flex") -> dict:
"""已下发工单进度清单(EX-09)。"""
from server.state.seed import ensure_flex_seed
tr = (track or "flex").lower()
if tr == "flex":
ensure_flex_seed(world)
ver = _latest_flex(world)
vid = ver["id"] if ver else None
wos = [w for w in world.get("flexWorkOrders", [])
if (not vid or w.get("versionId") == vid) and w.get("mesExternalId")]
else:
ver = _latest_fixed(world, published_only=False)
vid = ver["id"] if ver else None
wos = [w for w in world.get("workOrders", [])
if (not vid or w.get("versionId") == vid) and w.get("mesExternalId")]
rows = [{
"woId": w["id"], "orderNo": w.get("flexOrderNo") or w.get("productionOrderNo"),
"operation": w.get("operationCode") or w.get("operationName"),
"equipment": w.get("equipmentCode") or w.get("workstationName"),
"status": w.get("status"), "progressPct": w.get("progressPct", 0),
"qtyDone": w.get("qtyDone", 0), "mesExternalId": w.get("mesExternalId"),
"start": w.get("plannedStartTime"), "end": w.get("plannedEndTime"),
} for w in wos]
done = sum(1 for r in rows if r["status"] == "COMPLETED" or (r["progressPct"] or 0) >= 100)
return {
"track": tr, "versionNo": (ver or {}).get("versionNo"),
"total": len(rows), "completed": done,
"rows": rows,
"connection": get_mes_client().status(),
}
def apply_report(store, wo_id: int, *, track: str = "flex",
progress_pct: int | None = None, finish: bool = False,
actor: str = "web") -> dict:
"""报工回流:更新 APS 工单进度并写 Mock MES。"""
from server.agent_core import harness
from server.agent_core.audit import write_audit
from server.state.seed import ensure_flex_seed
def _run():
tr = (track or "flex").lower()
if tr == "flex":
ensure_flex_seed(store.data)
wos = store.data.get("flexWorkOrders", [])
else:
wos = store.data.get("workOrders", [])
wo = next((w for w in wos if w["id"] == wo_id), None)
if not wo:
raise ValueError(f"工单 #{wo_id} 不存在")
ext_id = wo.get("mesExternalId")
if not ext_id:
raise ValueError(f"工单 #{wo_id} 尚未下发 MES,无法报工")
pct = 100 if finish else int(progress_pct if progress_pct is not None else 50)
pct = max(0, min(100, pct))
status = "COMPLETED" if pct >= 100 else "RUNNING"
qty = wo.get("qtyDone") or 0
# 柔性订单数量作参考
if tr == "flex":
fo = next((o for o in store.data.get("flexOrders", [])
if o.get("orderNo") == wo.get("flexOrderNo")), None)
target_qty = int((fo or {}).get("quantity") or 1)
else:
target_qty = 1
qty_done = target_qty if pct >= 100 else max(qty, int(target_qty * pct / 100))
client = get_mes_client()
client.post_report(ext_id, {
"progressPct": pct, "qtyDone": qty_done, "status": status, "actor": actor,
})
wo["progressPct"] = pct
wo["qtyDone"] = qty_done
wo["status"] = status
if status == "COMPLETED":
wo["actualEndTime"] = wo.get("plannedEndTime")
# 同订单工序全完工 → 柔性订单完成
order_done = False
if tr == "flex" and status == "COMPLETED":
order_no = wo.get("flexOrderNo")
vid = wo.get("versionId")
sibs = [w for w in store.data.get("flexWorkOrders", [])
if w.get("flexOrderNo") == order_no and w.get("versionId") == vid
and not w.get("frozen")]
if sibs and all((w.get("status") == "COMPLETED" or (w.get("progressPct") or 0) >= 100)
for w in sibs):
fo = next((o for o in store.data.get("flexOrders", [])
if o.get("orderNo") == order_no), None)
if fo:
fo["status"] = "COMPLETED"
order_done = True
store.data.setdefault("mesLinks", []).append({
"kind": "report", "woId": wo_id, "externalWoId": ext_id,
"track": tr, "progressPct": pct, "status": status,
"syncedAt": fmt_date(today0()), "actor": actor,
})
write_audit(store.data, store.next_id, actor=actor, category="INTEGRATION",
action="mes.report", target={"type": "WORK_ORDER", "id": wo_id},
power="P1", rationale={"pct": pct, "status": status, "orderDone": order_done})
store.save()
msg = f"报工已回写 ✅ WO#{wo_id} → {pct}%({status})"
if order_done:
msg += f";订单 {wo.get('flexOrderNo')} 已完工"
return {"woId": wo_id, "progressPct": pct, "status": status,
"orderDone": order_done, "message": msg}
return harness.guard("mes.report", {"woId": wo_id}, _run)
def confirmation_for_dispatch(world: World, track: str = "flex") -> tuple[str, list[str]]:
p = preview_dispatch(world, track)
tr_cn = "柔性" if (track or "flex").lower() == "flex" else "固定"
return f"MES 下发确认({tr_cn})", [
p["summary"],
f"系统 {p['connection'].get('system')} / 工厂 {p['connection'].get('plant')}(Mock)",
"外部副作用:写入 Mock MES;幂等键防重复工单(P3)",
]
def stage_dispatch(store, track: str = "flex", *, session_id: str, actor: str = "web") -> dict:
from server.agent_core import harness
from server.agent_core.audit import write_audit
tr = (track or "flex").lower()
title, lines = confirmation_for_dispatch(store.data, tr)
preview = preview_dispatch(store.data, tr)
if not preview.get("items"):
return {"staged": False, "message": preview.get("summary") or "无可下发工单", "block": None}
if preview.get("newCount", 0) == 0:
return {"staged": False, "message": "全部工序已下发(幂等),无需重复。", "block": None}
block = harness.stage_confirmation(
session_id, "mes.dispatch", {"track": tr},
title=title, summary_lines=lines)
write_audit(store.data, store.next_id, actor=actor, category="GATE",
action="mes.dispatch.stage", target={"type": "MES", "id": tr},
power="P3", rationale={"confirmId": block.props["confirmId"]})
store.save()
return {"staged": True, "message": f"{title} 属于 P3,需要你确认后执行。", "block": block}