aps-agent/server/aps_domain/project_analyze.py

1217 lines
52 KiB
Python
Raw Permalink 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.

# ============================================================
# 项目深度分析(moduleId: domain-project-analyze, 可重生 ✅)
# Pi 已结构化的工程目录 Excel/SQL → 写入当前项目世界 → 明细表 + 齐备度 + 知识库补洞 plan
# ============================================================
from __future__ import annotations
import copy
import json
import os
import re
from typing import Any
from server.aps_domain.readiness import check_readiness
World = dict[str, Any]
_TABLE_EXT_RE = re.compile(r"\.(xlsx|xlsm|csv|txt)$", re.I)
_INPUT_EXT_RE = re.compile(r"\.(xlsx|xlsm|csv|txt|sql)$", re.I)
def _file_sha256(path: str) -> str:
"""资料文件指纹:采用卡冻结来源文件,批准时校验同一份文件未变化。"""
from server.agent_core.fallback_highrisk import file_sha256
return str(file_sha256(path) or "")
def _project_ctx(session_id: str | None) -> tuple[dict | None, str]:
try:
from server.state.projects import get_project_store
snap = get_project_store().snapshot(include_messages=False)
except Exception:
return None, ""
sessions = {s["id"]: s for s in (snap.get("sessions") or []) if isinstance(s, dict)}
projects = {p["id"]: p for p in (snap.get("projects") or []) if isinstance(p, dict)}
sess = sessions.get(session_id or "") or {}
proj = projects.get(sess.get("projectId") or "") if sess.get("projectId") else None
work = (proj.get("workDir") or "").strip() if proj else ""
return proj, work
def _project_ctx_readonly(session_id: str | None) -> tuple[dict | None, str]:
"""Resolve session -> project without ProjectStore seeding writes.
解析顺序:SQLite 只读直查(真实部署主路径)→ 只读 JSON 快照 →
注入/内存态 ProjectStore(单元测试与桌面内存态)。前两条不写任何东西;
最后一条只用于直查拿不到项目的场景,避免只读分析在内存态下直接失效。
"""
sid = str(session_id or "").strip()
if not sid:
return None, ""
try:
from server.auth.context import get_identity
identity = get_identity()
except Exception:
return None, ""
try:
import sqlite3
from sqlalchemy.engine import make_url
from server.db.database import _database_url
url = make_url(_database_url())
if url.get_backend_name() == "sqlite" and url.database:
db_path = os.path.abspath(str(url.database))
if os.path.isfile(db_path):
uri = "file:" + db_path.replace("\\", "/") + "?mode=ro"
connection = sqlite3.connect(uri, uri=True)
connection.row_factory = sqlite3.Row
try:
row = connection.execute(
"SELECT p.id, p.name, p.work_dir "
"FROM aps_chat_sessions s "
"JOIN aps_workspace_projects p "
"ON p.id = s.project_id AND p.tenant_uuid = s.tenant_uuid "
"JOIN aps_project_members m "
"ON m.project_id = p.id AND m.tenant_uuid = p.tenant_uuid "
"WHERE s.id = ? AND s.tenant_uuid = ? AND s.scope = 'project' "
"AND s.deleted = 0 AND p.deleted = 0 "
"AND m.user_id = ? AND m.status = 'active' AND m.deleted = 0 "
"LIMIT 1",
(sid, identity.tenant_uuid, identity.user_id),
).fetchone()
finally:
connection.close()
if row is not None:
project = {
"id": row["id"],
"name": row["name"],
"workDir": row["work_dir"],
}
return project, str(row["work_dir"] or "").strip()
except Exception:
pass
legacy_path = os.environ.get("APS_PROJECTS_PATH") or ""
try:
if legacy_path and os.path.isfile(legacy_path):
with open(legacy_path, "r", encoding="utf-8") as source:
legacy = json.load(source)
sessions = {
str(row.get("id")): row
for row in (legacy.get("sessions") or [])
if isinstance(row, dict)
}
projects = {
str(row.get("id")): row
for row in (legacy.get("projects") or [])
if isinstance(row, dict)
}
session = sessions.get(sid) or {}
project = projects.get(str(session.get("projectId") or "")) or None
if project:
return project, str(project.get("workDir") or "").strip()
except Exception:
pass
try:
from server.state.projects import get_project_store
snap = get_project_store().snapshot(include_messages=False)
sessions = {str(row.get("id")): row for row in (snap.get("sessions") or [])
if isinstance(row, dict)}
projects = {str(row.get("id")): row for row in (snap.get("projects") or [])
if isinstance(row, dict)}
session = sessions.get(sid) or {}
project = projects.get(str(session.get("projectId") or "")) or None
if project:
return project, str(project.get("workDir") or "").strip()
except Exception:
pass
return None, ""
def _analysis_error(message: str, *, read_only: bool = False) -> dict[str, Any]:
return {
"ok": False,
"error": message,
"projectName": None,
"workDir": "",
"sources": [],
"sourcePaths": [],
"applied": {},
"owner": None,
"sqlPreview": None,
"folder": None,
"orders": [],
"materials": [],
"routings": [],
"equipment": [],
"kbActions": [],
"plan": [],
"readiness": {"summary": {}},
"summary": {
"orders": 0, "materials": 0, "routings": 0,
"equipment": 0, "bom": 0, "ready": 0, "blocked": 0,
},
"canSchedule": False,
"readOnly": read_only,
"pendingAdoption": {
"mode": "none", "filename": "", "sourcePaths": [],
"batches": [], "counts": {},
},
"markdown": message,
"knowledgeIngest": {},
}
def _resolve_input_path(input_path: str | None) -> tuple[str, str | None]:
"""确定性校验 Pi 显式传入的文件/目录路径,不解析任何原始问句。"""
raw = str(input_path or "").strip()
if not raw:
return "", None
candidate = os.path.abspath(raw)
if not os.path.exists(candidate):
return "", f"指定的资料路径不存在:{raw}。请核对路径后重新提供 input_path。"
if os.path.isfile(candidate):
if not _INPUT_EXT_RE.search(candidate):
return "", "指定文件类型不支持:仅接受 xlsx、xlsm、csv、txt 或 sql。"
return candidate, None
if os.path.isdir(candidate):
return candidate, None
return "", "指定路径既不是文件也不是目录,请重新提供有效的 input_path。"
def _direct_table_report(
world: World,
paths: list[str],
*,
project_name: str,
work_dir: str,
) -> dict[str, Any]:
"""直接路径旁路:复用 preview_file,但不依赖工作区里的 workDir。"""
from server.aps_domain.folder_pack import KIND_CN
from server.aps_domain.importers import preview_file
required = ("orders", "materials", "routing", "equipment")
files: list[dict[str, Any]] = []
batches: list[dict[str, Any]] = []
kind_counts: dict[str, int] = {}
for path in paths:
name = os.path.basename(path)
try:
with open(path, "rb") as source:
preview = preview_file(name, source.read(), world, soft=True)
file_batches = preview.get("batches") or []
ok_count = int(preview.get("totalOk") or 0)
error_count = int(preview.get("totalErrors") or 0)
primary_kind = next(
(str(batch.get("kind")) for batch in file_batches if batch.get("okCount")),
str((file_batches[0] if file_batches else {}).get("kind") or ""),
)
field_map = next(
(list(batch.get("fieldMap") or []) for batch in file_batches if batch.get("fieldMap")),
[],
)
errors: list[str] = []
warnings: list[str] = []
headers: list[str] = []
samples: list[str] = []
sheets: list[dict[str, Any]] = []
for batch in file_batches:
kind = str(batch.get("kind") or primary_kind or "materials")
count = int(batch.get("okCount") or 0)
if count:
kind_counts[kind] = kind_counts.get(kind, 0) + count
errors.extend(str(item) for item in (batch.get("errors") or [])[:2])
warnings.extend(str(item) for item in (batch.get("warnings") or [])[:2])
if not headers:
headers = [str(item) for item in (batch.get("headersRaw") or []) if item][:8]
for row in (batch.get("okRows") or [])[:1]:
if isinstance(row, dict):
samples.append(";".join(f"{k}={v}" for k, v in list(row.items())[:6]))
sheets.append({
"sheet": batch.get("sheet"), "kind": kind,
"kindCn": KIND_CN.get(kind, kind), "okCount": count,
"errorCount": int(batch.get("errorCount") or 0),
})
if count:
batches.append({
"kind": kind, "sheet": batch.get("sheet") or name,
"okRows": batch.get("okRows") or [], "okCount": count,
"errorCount": int(batch.get("errorCount") or 0),
"errors": (batch.get("errors") or [])[:8],
"warnings": (batch.get("warnings") or [])[:5],
"fieldMap": batch.get("fieldMap") or [], "sourceFile": name,
})
fmap_text = "、".join(
f"{item.get('source')}→{item.get('target')}" for item in field_map
)
issue = (errors or warnings or ([f"字段:{fmap_text}"] if fmap_text else []))
files.append({
"name": name, "size": os.path.getsize(path),
"status": "ok" if ok_count and not error_count else ("warn" if ok_count else "fail"),
"kind": primary_kind, "kindCn": KIND_CN.get(primary_kind, primary_kind or "未识别"),
"okCount": ok_count, "errorCount": error_count,
"errors": errors, "warnings": warnings, "headers": headers,
"samples": samples, "sheets": sheets, "fieldMap": field_map,
"fieldMapText": fmap_text, "issue": issue[0] if issue else "",
})
except Exception as exc: # noqa: BLE001
files.append({
"name": name, "size": os.path.getsize(path) if os.path.isfile(path) else 0,
"status": "fail", "kind": "", "kindCn": "读取失败",
"okCount": 0, "errorCount": 1, "errors": [str(exc)],
"warnings": [], "headers": [], "samples": [], "sheets": [],
"fieldMap": [], "fieldMapText": "", "issue": str(exc),
})
coverage = {kind: kind_counts.get(kind, 0) > 0 for kind in required}
coverage_detail = [{
"key": kind, "label": KIND_CN[kind], "ok": coverage[kind],
"count": kind_counts.get(kind, 0), "fieldMap": [],
"note": (f"读到 {kind_counts[kind]} 行" if coverage[kind] else "还没有有效数据"),
} for kind in required]
missing = [KIND_CN[kind] for kind in required if not coverage[kind]]
return {
"ok": True, "projectName": project_name, "workDir": work_dir,
"files": files, "coverage": coverage, "coverageDetail": coverage_detail,
"kindCounts": kind_counts, "missing": missing, "softMissing": [],
"canSchedule": not missing, "batches": batches,
"totalOk": sum(int(item.get("okCount") or 0) for item in files),
"totalErrors": sum(int(item.get("errorCount") or 0) for item in files),
}
def _preview_apply_template(
world: World,
template: dict[str, Any],
product_code: str,
product_name: str,
) -> None:
"""Apply a routing template to a detached preview world without persistence."""
world["flexRoutings"] = [
row for row in (world.get("flexRoutings") or [])
if row.get("productCode") != product_code
]
operations = {
row.get("code"): row for row in (world.get("flexOperations") or [])
if isinstance(row, dict)
}
equipment = world.get("flexEquipment") or []
for step in template.get("steps") or []:
code = step.get("operationCode")
if not code:
continue
if code not in operations:
world.setdefault("flexOperations", []).append({
"code": code,
"name": step.get("operationName") or code,
"isBottleneck": False,
"changeoverMin": step.get("changeoverMin") or 0,
})
operations[code] = world["flexOperations"][-1]
if not any(code in (row.get("capabilities") or []) for row in equipment):
for row in equipment:
if row.get("status") == "RUNNING":
row.setdefault("capabilities", []).append(code)
row.setdefault("opStdTime", {})[code] = step.get("stdMinDefault") or 1.0
world.setdefault("flexRoutings", []).append({
"productCode": product_code,
"productName": product_name or product_code,
"seq": step.get("seq"),
"operationCode": code,
"requireMold": False,
"stdTimePerUnit": step.get("stdMinDefault") or 1.0,
"stdTimeSource": "模板",
})
template_ops = [row.get("operationCode") for row in (template.get("steps") or [])]
teams = world.get("flexTeams") or []
generic = next((team for team in teams if team.get("code") == "T-TPL"), None)
if generic is None and teams is not None:
generic = {
"code": "T-TPL", "name": "模板工序班组(推断)",
"memberCount": 4, "supportOps": [], "skillLevel": "L2",
}
world.setdefault("flexTeams", []).append(generic)
if generic is not None:
for code in template_ops:
if code and code not in generic.setdefault("supportOps", []):
generic["supportOps"].append(code)
if not any(row.get("code") == product_code for row in (world.get("flexMaterials") or [])):
world.setdefault("flexMaterials", []).append({
"code": product_code, "name": product_name or product_code,
"type": "FINISHED_PRODUCT", "unit": "件",
"stock": 0, "inTransit": 0, "safetyStock": 0, "procurementLeadTime": 0,
})
def _template_catalog_readonly() -> list[dict[str, Any]]:
"""Load the routing-template catalog without initializing or seeding storage."""
try:
from server.knowledge.routing_templates import _BUILTIN
except Exception:
return []
catalog: list[dict[str, Any]] = []
for item in _BUILTIN:
steps = []
for seq, code, name, equipment_type, low, high, setup, changeover in (item.get("steps") or []):
steps.append({
"seq": seq,
"operationCode": code,
"operationName": name,
"equipmentType": equipment_type,
"stdMinLow": low,
"stdMinHigh": high,
"stdMinDefault": round((float(low) + float(high)) / 2, 1),
"setupMin": setup,
"changeoverMin": changeover,
})
catalog.append({
"code": item.get("code"),
"name": item.get("name"),
"industry": "machining",
"category": item.get("category"),
"keywords": str(item.get("keywords") or "").split(","),
"description": item.get("description"),
"assetId": "",
"source": "builtin",
"steps": steps,
})
# Overlay tenant/platform templates when a SQLite database already exists.
# Opening it in mode=ro avoids creating a database or running schema seeding.
try:
import sqlite3
from sqlalchemy.engine import make_url
from server.auth.context import get_identity
from server.db.database import _database_url
url = make_url(_database_url())
if url.get_backend_name() != "sqlite" or not url.database:
return catalog
db_path = os.path.abspath(str(url.database))
if not os.path.isfile(db_path):
return catalog
tenants = ("platform", get_identity().tenant_uuid or "platform")
placeholders = ",".join("?" for _ in tenants)
uri = "file:" + db_path.replace("\\", "/") + "?mode=ro"
connection = sqlite3.connect(uri, uri=True)
connection.row_factory = sqlite3.Row
try:
rows = connection.execute(
"SELECT id, tenant_uuid, code, name, industry, category, keywords, "
"description, asset_id, source FROM routing_templates "
f"WHERE deleted = 0 AND tenant_uuid IN ({placeholders}) ORDER BY id",
tenants,
).fetchall()
if not rows:
return catalog
template_ids = [int(row["id"]) for row in rows]
id_placeholders = ",".join("?" for _ in template_ids)
step_rows = connection.execute(
"SELECT template_id, seq, operation_code, operation_name, equipment_type, "
"std_min_low, std_min_high, setup_min, changeover_min "
"FROM routing_template_steps "
f"WHERE deleted = 0 AND template_id IN ({id_placeholders}) ORDER BY seq",
template_ids,
).fetchall()
steps_by_template: dict[int, list[dict[str, Any]]] = {}
for step in step_rows:
low = float(step["std_min_low"] or 0)
high = float(step["std_min_high"] or 0)
steps_by_template.setdefault(int(step["template_id"]), []).append({
"seq": step["seq"],
"operationCode": step["operation_code"],
"operationName": step["operation_name"],
"equipmentType": step["equipment_type"],
"stdMinLow": low,
"stdMinHigh": high,
"stdMinDefault": round((low + high) / 2, 1),
"setupMin": step["setup_min"],
"changeoverMin": step["changeover_min"],
})
merged = {str(item.get("code")): item for item in catalog}
order = [str(item.get("code")) for item in catalog]
for row in rows:
code = str(row["code"] or "")
if not code:
continue
merged[code] = {
"code": code,
"name": row["name"],
"industry": row["industry"],
"category": row["category"],
"keywords": str(row["keywords"] or "").split(","),
"description": row["description"],
"assetId": row["asset_id"],
"source": row["source"],
"steps": steps_by_template.get(int(row["id"]), []),
}
if code not in order:
order.append(code)
return [merged[code] for code in order]
finally:
connection.close()
except Exception:
return catalog
def _recommend_templates_readonly(text: str, top_k: int = 3) -> list[dict[str, Any]]:
"""Match templates against an immutable preview catalog without seeding."""
text = (text or "").strip()
scored: list[tuple[float, dict[str, Any]]] = []
for template in _template_catalog_readonly():
score = 0.0
for keyword in template.get("keywords") or []:
if keyword and keyword in text:
score += 2.0
if template.get("category") and template["category"] in text:
score += 1.5
for step in template.get("steps") or []:
if step.get("operationName") and step["operationName"] in text:
score += 0.5
if score > 0:
scored.append((score, template))
scored.sort(key=lambda item: -item[0])
return [dict(template, matchScore=score) for score, template in scored[:top_k]]
def _fill_from_knowledge(
world: World,
product_codes: list[str],
*,
apply: bool = True,
) -> list[dict[str, Any]]:
"""缺工艺时用行业模板/知识库补齐;只读模式仅预演模板效果。"""
actions: list[dict[str, Any]] = []
try:
from server.knowledge.routing_templates import recommend_templates, apply_template_to_product
except Exception:
return actions
routes = {r.get("productCode") for r in (world.get("flexRoutings") or [])}
mats = {m.get("code") for m in (world.get("flexMaterials") or [])}
for pc in product_codes:
if not pc:
continue
if pc not in mats:
world.setdefault("flexMaterials", []).append({
"id": len(world.get("flexMaterials") or []) + 1,
"code": pc, "name": pc, "type": "FINISHED_PRODUCT", "unit": "件",
"stock": 0, "inTransit": 0, "safetyStock": 0, "procurementLeadTime": 0,
})
actions.append({"type": "material_stub", "productCode": pc,
"detail": f"物料主数据缺失,已建产品档案占位 {pc}"})
if pc in routes:
continue
name = next((m.get("name") for m in (world.get("flexMaterials") or [])
if m.get("code") == pc), pc)
try:
if apply:
recs = recommend_templates(f"{pc} {name}") or []
else:
recs = _recommend_templates_readonly(f"{pc} {name}") or []
except Exception:
recs = []
if not recs:
actions.append({"type": "need_routing", "productCode": pc,
"detail": f"产品 {pc} 无工艺,知识库也无匹配模板,需客户提供工艺路线"})
continue
tpl = recs[0]
code = tpl.get("code") or tpl.get("templateCode")
try:
if apply:
apply_template_to_product(world, code, pc)
else:
_preview_apply_template(world, tpl, pc, name)
actions.append({
"type": "template_applied", "productCode": pc,
"template": code, "templateName": tpl.get("name"),
"detail": (
f"产品 {pc} 缺工艺 → 已用知识库模板「{tpl.get('name') or code}」生成"
if apply else
f"产品 {pc} 缺工艺 → 只读预演知识库模板「{tpl.get('name') or code}」"
),
})
routes.add(pc)
except Exception as exc:
actions.append({"type": "template_failed", "productCode": pc,
"detail": f"模板 {code} 应用失败:{exc}"})
return actions
def _preview_replace_flex(
world: World,
values: dict[str, Any],
*,
ensure_calendar: bool = False,
) -> dict[str, int]:
"""Load an already parsed flex snapshot into a detached preview world."""
fields = (
("flexOrders", ("flexOrders", "orders")),
("flexMaterials", ("flexMaterials", "materials")),
("flexEquipment", ("flexEquipment", "equipment")),
("flexRoutings", ("flexRoutings", "routing")),
("flexBom", ("flexBom", "bom")),
("flexOperations", ("flexOperations", "operations")),
("flexZones", ("flexZones", "zones")),
("flexTeams", ("flexTeams", "teams")),
("flexPersonnel", ("flexPersonnel", "personnel")),
("flexFactoryResources", ("flexFactoryResources", "factoryResources")),
("flexPartners", ("flexPartners", "partners")),
("flexWip", ("flexWip", "wip")),
)
summary: dict[str, int] = {}
for target, sources in fields:
source = next((name for name in sources if name in values), None)
if source is None:
continue
rows = copy.deepcopy(values.get(source) or [])
world[target] = rows
if rows:
summary[target] = len(rows)
if ensure_calendar and not world.get("flexCalendar"):
world["flexCalendar"] = [{
"shiftCode": "D", "startTime": "08:00", "endTime": "17:00",
"breaks": [{"start": "12:00", "end": "13:00"}],
"workdays": [1, 2, 3, 4, 5],
}]
summary["calendar"] = 1
return summary
def _preview_apply_batches(world: World, batches: list[dict[str, Any]]) -> dict[str, Any]:
"""Apply reviewed import rows to a detached preview world for read-only analysis."""
profile_batches = [batch for batch in batches if batch.get("sourceProfile")]
if profile_batches:
if len(profile_batches) != len(batches):
raise ValueError("完整工作簿与其他导入批次不能混合确认,请分别导入")
from server.importers.planning_workbook import apply_planning_batches
return apply_planning_batches(world, profile_batches)
summary: dict[str, int] = {}
def _next_id(table: str) -> int:
rows = world.get(table) or []
return max(
(row.get("id", 0) for row in rows if isinstance(row.get("id"), int)),
default=0,
) + 1
def _upsert(table: str, row: dict[str, Any], keys: tuple[str, ...]) -> None:
rows = world.setdefault(table, [])
existing = next(
(
item for item in rows
if all(item.get(key) == row.get(key) for key in keys)
),
None,
)
if existing is not None:
existing.update(row)
else:
rows.append({"id": _next_id(table), **row})
for batch in batches:
kind = batch.get("kind")
for row in batch.get("okRows") or []:
if kind == "orders":
_upsert("flexOrders", {
"orderNo": row.get("orderNo"),
"productCode": row.get("productCode"),
"quantity": row.get("quantity"),
"dueDate": row.get("deliveryDate"),
"priority": row.get("priority", 5),
"customerName": row.get("customerName") or "现场客户",
"customerLevel": row.get("customerLevel") or "B",
"isRush": bool(row.get("isRush")),
"status": "RELEASED",
}, ("orderNo",))
summary["flexOrders"] = summary.get("flexOrders", 0) + 1
elif kind == "materials":
_upsert("flexMaterials", dict(row), ("code",))
summary["flexMaterials"] = summary.get("flexMaterials", 0) + 1
elif kind == "calendar":
_upsert("flexCalendar", dict(row), ("shiftCode",))
summary["flexCalendar"] = summary.get("flexCalendar", 0) + 1
elif kind == "equipment":
_upsert("flexEquipment", dict(row), ("code",))
summary["flexEquipment"] = summary.get("flexEquipment", 0) + 1
elif kind == "molds":
_upsert("flexMolds", dict(row), ("code",))
summary["flexMolds"] = summary.get("flexMolds", 0) + 1
elif kind == "operations":
_upsert("flexOperations", dict(row), ("code",))
summary["flexOperations"] = summary.get("flexOperations", 0) + 1
elif kind == "zones":
_upsert("flexZones", dict(row), ("code",))
summary["flexZones"] = summary.get("flexZones", 0) + 1
elif kind == "routing":
_upsert("flexRoutings", dict(row), ("productCode", "seq"))
summary["flexRoutings"] = summary.get("flexRoutings", 0) + 1
elif kind == "bom":
_upsert("flexBom", dict(row), ("productCode", "materialCode"))
summary["flexBom"] = summary.get("flexBom", 0) + 1
return {"summary": summary, "total": sum(summary.values())}
def analyze_project_deep(
world: World,
session_id: str | None = None,
*,
apply_sql: bool = True,
query: str = "",
input_path: str | None = None,
next_id=None,
apply: bool = True,
) -> dict[str, Any]:
"""深度分析显式资料路径或当前项目目录;query 仅为兼容旧调用的弃用参数。"""
del query # 原始自然语言不参与路径、实体或数据语义推断。
try:
from server.agent_core.progress import emit_thinking
except Exception:
def emit_thinking(*_a, **_k): # type: ignore
return None
explicit_input, input_error = _resolve_input_path(input_path)
if input_error:
return _analysis_error(input_error, read_only=not apply)
proj, work = _project_ctx(session_id) if apply else _project_ctx_readonly(session_id)
if explicit_input:
work = explicit_input if os.path.isdir(explicit_input) else os.path.dirname(explicit_input)
if proj is None:
proj = {
"name": os.path.basename(work.rstrip("\\/")) or "当前项目",
"workDir": work,
}
if not work or not os.path.isdir(work):
return _analysis_error(
"请明确要分析的资料路径:提供 input_path,或先选择带有工程目录的项目。"
"不能从原始问句推断路径。",
read_only=not apply,
)
analysis_world = world if apply else copy.deepcopy(world)
sources: list[str] = []
source_paths: list[str] = []
sql_preview: dict[str, Any] | None = None
applied: dict[str, int] = {}
kb_actions: list[dict[str, Any]] = []
mom_applied = False
mom_path = ""
def apply_folder_batches(report: dict[str, Any] | None) -> None:
batches = (report or {}).get("batches") or []
if not batches or next_id is None:
return
if apply:
from server.aps_domain.importers import apply_import_commit
result = apply_import_commit(analysis_world, next_id, batches)
else:
result = _preview_apply_batches(analysis_world, batches)
# apply_import_commit returns {summary: {...}, total: N}; keep accepting
# the historical flat counter shape for callers that still provide it.
counts = (
result.get("summary")
if isinstance(result, dict) and isinstance(result.get("summary"), dict)
else result
) or {}
for key, value in counts.items():
if isinstance(value, (int, float)) and not isinstance(value, bool):
applied[key] = applied.get(key, 0) + int(value)
emit_thinking("定位工程目录", (proj or {}).get("name") or "当前项目", pct=2)
if work and os.path.isdir(work):
names = ([os.path.basename(explicit_input)] if os.path.isfile(explicit_input)
else sorted(os.listdir(work)))
sqls = [n for n in names if n.lower().endswith(".sql") and not n.startswith("~$")]
xls = [n for n in names if _TABLE_EXT_RE.search(n) and not n.startswith("~$")]
source_paths = [
os.path.join(work, name) for name in (*sqls, *xls)
if os.path.isfile(os.path.join(work, name))
]
folder = None
emit_thinking(
"清点目录文件",
f"SQL {len(sqls)} · Excel/CSV {len(xls)}",
pct=5,
)
if sqls:
from server.importers.sql_pack import preview_sql_file, apply_sql_pack_to_world
# 优先最大的 sql(完整库)
sqls_sorted = sorted(sqls, key=lambda n: os.path.getsize(os.path.join(work, n)), reverse=True)
path = os.path.join(work, sqls_sorted[0])
sources.append(f"SQL:{sqls_sorted[0]}")
try:
sql_preview = preview_sql_file(path)
# Read-only analysis must still materialize SQL into the detached
# preview world when callers disable both write switches.
if sql_preview.get("flex") and (apply_sql or not apply):
emit_thinking(
"写入项目世界" if apply else "预演项目世界",
"订单 / 物料 / 工艺 / 设备",
pct=93,
)
# SQL 包是工程目录真相源,始终整包覆盖,避免旧错误工艺残留
if apply:
applied = apply_sql_pack_to_world(
analysis_world, sql_preview["flex"], replace=True)
else:
applied = _preview_replace_flex(
analysis_world,
sql_preview["flex"],
ensure_calendar=True,
)
emit_thinking(
"入库完成",
"、".join(f"{k}{v}" for k, v in applied.items()) or "无新增",
status="done", pct=95,
)
except Exception as exc:
sql_preview = {"error": str(exc), "filename": sqls_sorted[0]}
emit_thinking("SQL 解析失败", str(exc), status="warn")
elif xls:
# 无 SQL 时:识别 MOM 主数据收集表并写入世界(否则深度分析永远是 0)
try:
from server.importers.mom_pack import (
apply_mom_pack_to_world,
is_mom_workbook,
mom_pack_is_adoptable,
mom_pack_row_count,
parse_mom_workbook,
)
mom_files = [
n for n in xls
if is_mom_workbook(os.path.join(work, n))
]
# 文件名含 MOM/主数据 优先;否则对全部 xlsx 试探
candidates = mom_files or xls
for n in candidates:
path = os.path.join(work, n)
if not is_mom_workbook(path):
continue
pack = parse_mom_workbook(path)
if not mom_pack_is_adoptable(pack):
# 结构证据(≥2 张模型表)或主干行(物料/BOM/订单/工艺路线)
# 不足:采用即整表替换,这里不作为 MOM 来源,避免清空既有主数据。
emit_thinking("MOM 主数据表不足以采用", n, status="warn")
continue
sources.append(f"MOM:{n}")
emit_thinking("识别 MOM 主数据表", n, pct=90)
if apply:
applied = apply_mom_pack_to_world(analysis_world, pack, replace=True)
else:
applied = _preview_replace_flex(analysis_world, pack)
result = {"stats": pack.get("stats") or {}}
mom_applied = True
mom_path = path
emit_thinking(
"MOM 入库完成",
"、".join(f"{k}{v}" for k, v in applied.items()) or "无新增",
status="done", pct=95,
)
stats = result.get("stats") or {}
if stats.get("note"):
sources.append(str(stats["note"]))
break
except Exception as exc:
emit_thinking("MOM Excel 解析失败", str(exc), status="warn")
if not mom_applied:
sources.append(f"Excel×{len(xls)}")
try:
emit_thinking("解析并写入 Excel/CSV", f"{len(xls)} 个表格", pct=96)
if explicit_input:
folder = _direct_table_report(
analysis_world,
[os.path.join(work, name) for name in xls],
project_name=(proj or {}).get("name") or "当前项目",
work_dir=work,
)
else:
from server.aps_domain.folder_pack import analyze_work_dir
folder = analyze_work_dir(analysis_world, session_id, apply_sql=apply)
apply_folder_batches(folder)
except Exception as exc:
if folder is None:
folder = {
"ok": False,
"projectName": (proj or {}).get("name"),
"workDir": work,
"files": [], "coverage": {}, "missing": [],
"canSchedule": False, "batches": [],
"totalOk": 0, "totalErrors": 1,
"error": str(exc),
}
else:
folder["error"] = str(exc)
else:
try:
if explicit_input:
folder = _direct_table_report(
analysis_world,
[os.path.join(work, name) for name in xls],
project_name=(proj or {}).get("name") or "当前项目",
work_dir=work,
)
else:
from server.aps_domain.folder_pack import analyze_work_dir
folder = analyze_work_dir(analysis_world, session_id, apply_sql=apply)
except Exception as exc:
if folder is None:
folder = {
"ok": False,
"projectName": (proj or {}).get("name"),
"workDir": work,
"files": [], "coverage": {}, "missing": [],
"canSchedule": False, "batches": [],
"totalOk": 0, "totalErrors": 1,
"error": str(exc),
}
else:
folder["error"] = str(exc)
else:
folder = None
if sqls and xls:
sources.append(f"Excel×{len(xls)}")
try:
emit_thinking("核对 Excel/CSV", f"{len(xls)} 个表格", pct=96)
if explicit_input:
folder = _direct_table_report(
analysis_world,
[os.path.join(work, name) for name in xls],
project_name=(proj or {}).get("name") or "当前项目",
work_dir=work,
)
else:
from server.aps_domain.folder_pack import analyze_work_dir
folder = analyze_work_dir(analysis_world, session_id, apply_sql=apply)
apply_folder_batches(folder)
except Exception as exc:
if folder is None:
folder = {
"ok": False,
"projectName": (proj or {}).get("name"),
"workDir": work,
"files": [], "coverage": {}, "missing": [],
"canSchedule": False, "batches": [],
"totalOk": 0, "totalErrors": 1,
"error": str(exc),
}
else:
folder["error"] = str(exc)
else:
folder = None
emit_thinking("未挂工程目录", "仅分析当前项目世界里的数据", status="warn", pct=10)
# 知识库补洞:缺工艺的产品
emit_thinking("检查齐备度与知识库", "缺工艺时自动匹配行业模板", pct=97)
active = [o for o in (analysis_world.get("flexOrders") or [])
if (o.get("status") or "") not in ("DONE", "CANCELLED")]
pcs = list({str(o.get("productCode") or "") for o in active})
kb_actions = _fill_from_knowledge(analysis_world, pcs, apply=apply)
if kb_actions:
emit_thinking(
"知识库补洞",
f"{len(kb_actions)} 项(模板/占位/待客户补)",
status="done", pct=98,
)
ready = check_readiness(analysis_world)
emit_thinking("整理分析报告", f"订单 {len(active)} 张", status="done", pct=99)
mats = analysis_world.get("flexMaterials") or []
routes = analysis_world.get("flexRoutings") or []
equip = [e for e in (analysis_world.get("flexEquipment") or []) if e.get("status") == "RUNNING"]
bom = analysis_world.get("flexBom") or []
order_rows = [{
"orderNo": o.get("orderNo"),
"productCode": o.get("productCode"),
"productName": o.get("productName") or next(
(m.get("name") for m in mats if m.get("code") == o.get("productCode")), ""),
"quantity": o.get("quantity"),
"dueDate": o.get("dueDate"),
"customerName": o.get("customerName") or "",
"ready": next((r["ready"] for r in ready["orders"] if r["orderNo"] == o.get("orderNo")), False),
} for o in active[:80]]
# 主数据抽样:订单涉及的产品优先
need_codes = {r["productCode"] for r in order_rows if r.get("productCode")}
mat_rows = []
for m in mats:
if need_codes and m.get("code") not in need_codes and len(mat_rows) >= 40:
continue
if m.get("code") in need_codes or len(mat_rows) < 40:
mat_rows.append({
"code": m.get("code"), "name": m.get("name"),
"type": m.get("type"), "unit": m.get("unit"), "stock": m.get("stock"),
})
if len(mat_rows) >= 60:
break
route_rows = []
for r in routes:
if need_codes and r.get("productCode") not in need_codes and len(route_rows) >= 40:
continue
route_rows.append({
"productCode": r.get("productCode"),
"seq": r.get("seq"),
"operationCode": r.get("operationCode"),
"operationName": r.get("operationName") or r.get("operationCode"),
"stdMin": r.get("stdTimePerUnit"),
})
if len(route_rows) >= 80:
break
equip_rows = [{
"code": e.get("code"), "name": e.get("name"),
"capabilities": "、".join((e.get("capabilities") or [])[:6]),
} for e in equip[:40]]
# 缺口 plan:先库后客户(同类缺口合并,避免刷屏)
plan: list[str] = []
need_rt = [a for a in kb_actions if a.get("type") == "need_routing"]
applied_tpl = [a for a in kb_actions if a.get("type") == "template_applied"]
other_kb = [a for a in kb_actions if a.get("type") not in ("need_routing", "template_applied")]
if applied_tpl:
plan.append(f"✓ 已用知识库模板补工艺 {len(applied_tpl)} 个产品")
for a in applied_tpl[:5]:
plan.append(f" · {a.get('productCode')} ← {a.get('templateName') or a.get('template')}")
if len(applied_tpl) > 5:
plan.append(f" · …另有 {len(applied_tpl) - 5} 个")
if need_rt:
samples = "、".join(str(a.get("productCode") or "") for a in need_rt[:8])
more = f" 等共 {len(need_rt)} 个" if len(need_rt) > 8 else f"({len(need_rt)} 个)"
plan.append(
f"△ SQL/知识库仍未挂上工艺的产品:{samples}{more}。"
f"若现场库本有工艺,请重新「分析项目」(工艺按订单 craftl 关联,不是中间件料号)"
)
for a in other_kb[:6]:
plan.append(f"· {a['detail']}")
s = ready["summary"]
if s.get("blocked"):
plan.append(f"△ 仍有 {s['blocked']} 张单齐备度未过,见订单明细「是否可排」列")
if not active:
plan.append("△ 没有订单:请在工程目录放订单表/SQL,或说「新建订单」")
if not equip:
plan.append("△ 没有可用设备:请提供设备台账,或从 SQL 的 md_equipment 导入")
if not plan:
plan.append("✓ 主数据与知识库已齐,可直接说「给我排产」")
proj_name = (proj or {}).get("name") or "当前项目"
md_lines = [
f"**项目深度分析 · {proj_name}**",
"",
f"数据来源:{'、'.join(sources) if sources else '当前项目世界(未挂工程目录文件)'}",
]
# 绑定当前登录用户
try:
from server.knowledge.ingest import bind_world_to_identity, ingest_project_analyze_report
owner_meta = bind_world_to_identity(analysis_world, sources=sources)
md_lines.append(
f"归属:{owner_meta.get('ownerFullname') or owner_meta.get('ownerUsername')} "
f"(用户 {owner_meta.get('ownerUserId')} · 租户 {str(owner_meta.get('tenantUuid') or '')[:8]}…)"
)
except Exception:
owner_meta = {}
ingest_project_analyze_report = None # type: ignore
if apply and applied:
md_lines.append(
"已入库:"
+ "、".join(f"{k} {v}" for k, v in applied.items())
)
if sql_preview and sql_preview.get("fieldMapText"):
md_lines.append(f"SQL 包识别:{sql_preview['fieldMapText']}")
md_lines += [
"",
f"订单 **{len(active)}** · 物料 **{len(mats)}** · 工艺步 **{len(routes)}** · "
f"设备 **{len(equip)}** · BOM **{len(bom)}**",
"",
"### 订单明细(本项目)",
"",
"| 订单号 | 产品 | 名称 | 数量 | 交期 | 客户 | 可排 |",
"| --- | --- | --- | ---: | --- | --- | --- |",
]
for r in order_rows[:30]:
md_lines.append(
f"| {r['orderNo']} | {r['productCode']} | {str(r.get('productName') or '')[:16]} "
f"| {r['quantity']} | {r['dueDate'] or '—'} | {str(r.get('customerName') or '')[:12]} "
f"| {'✓' if r['ready'] else '×'} |"
)
if len(order_rows) > 30:
md_lines.append(f"| … | 另有 {len(order_rows) - 30} 张 | | | | | |")
md_lines += [
"",
"### 主数据(产品/物料抽样)",
"",
"| 编码 | 名称 | 类型 | 单位 | 库存 |",
"| --- | --- | --- | --- | ---: |",
]
for m in mat_rows[:25]:
md_lines.append(
f"| {m['code']} | {str(m.get('name') or '')[:20]} | {m.get('type')} "
f"| {m.get('unit')} | {m.get('stock') or 0} |"
)
md_lines += ["", "### 补数 Plan(先知识库/数据库,再找客户)", ""]
for i, p in enumerate(plan[:12], 1):
md_lines.append(f"{i}. {p}")
def _adoption_counts(values: dict[str, Any] | None) -> dict[str, int]:
return {
str(key): int(value)
for key, value in (values or {}).items()
if isinstance(value, (int, float))
and not isinstance(value, bool)
and int(value) > 0
}
folder_batches = [
copy.deepcopy(batch)
for batch in ((folder or {}).get("batches") or [])
if batch.get("okRows")
]
sql_flex = (sql_preview or {}).get("flex") or {}
sql_available = bool(sql_flex.get("stats"))
if mom_applied and mom_path:
pending_adoption = {
"mode": "mom",
"filename": os.path.basename(mom_path),
"sourcePaths": [mom_path],
"batches": [],
"counts": _adoption_counts(applied),
# 出卡/执行两处都按这个口径复核,避免零业务行的 MOM 表升级成整表替换。
"businessRows": mom_pack_row_count(pack),
"adoptable": mom_pack_is_adoptable(pack),
}
if not apply:
# 只读分析不出卡不入库:把来源路径与指纹冻结给采用卡,批准后按同一份文件采用。
pending_adoption["applyPath"] = mom_path
pending_adoption["applySha256"] = _file_sha256(mom_path)
elif folder_batches:
batch_files = list(dict.fromkeys(
str(batch.get("sourceFile") or "")
for batch in folder_batches
if batch.get("sourceFile")
))
filename = batch_files[0] if len(batch_files) == 1 else ""
batch_file_set = set(batch_files)
selected_paths = [
path for path in source_paths
if os.path.basename(path) in batch_file_set
]
pending_adoption = {
"mode": "folder",
"filename": filename,
"sourcePaths": selected_paths or list(source_paths),
"batches": folder_batches,
"counts": _adoption_counts(applied) or {
kind: sum(
len(batch.get("okRows") or [])
for batch in folder_batches
if batch.get("kind") == kind
)
for kind in {str(batch.get("kind") or "") for batch in folder_batches}
if kind
},
}
elif sql_available:
sql_name = str((sql_preview or {}).get("filename") or "")
pending_adoption = {
"mode": "sql",
"filename": sql_name,
"sourcePaths": [
path for path in source_paths
if sql_name and os.path.basename(path) == sql_name
] or list(source_paths),
"batches": [],
"counts": _adoption_counts(applied),
}
if not apply:
# SQL 没有 okRows 批次:冻结已解析载荷,采用时整包写入而不是重读文件。
payload = copy.deepcopy(sql_flex)
if payload.get("stats"):
pending_adoption["sqlPayload"] = payload
else:
pending_adoption = {
"mode": "none", "filename": "", "sourcePaths": [],
"batches": [], "counts": {},
}
if not apply and pending_adoption["mode"] != "none":
md_lines.append("只读分析:资料尚未采用,确认后再写入项目数据。")
result = {
"ok": True,
"projectName": proj_name,
"workDir": work,
"sources": sources,
"sourcePaths": source_paths,
"applied": applied if apply else {},
"owner": {
"userId": owner_meta.get("ownerUserId"),
"username": owner_meta.get("ownerUsername"),
"fullname": owner_meta.get("ownerFullname"),
"tenantUuid": owner_meta.get("tenantUuid"),
} if owner_meta else None,
"sqlPreview": {
"filename": (sql_preview or {}).get("filename"),
"fieldMapText": (sql_preview or {}).get("fieldMapText"),
"tables": (sql_preview or {}).get("tables") or [],
"error": (sql_preview or {}).get("error"),
} if sql_preview else None,
"folder": folder,
"orders": order_rows,
"materials": mat_rows,
"routings": route_rows,
"equipment": equip_rows,
"kbActions": kb_actions,
"plan": plan,
"readiness": ready,
"summary": {
"orders": len(active), "materials": len(mats),
"routings": len(routes), "equipment": len(equip), "bom": len(bom),
"ready": s.get("ready", 0), "blocked": s.get("blocked", 0),
},
"canSchedule": s.get("ready", 0) > 0 and not ready.get("globalIssues"),
"readOnly": not apply,
"pendingAdoption": pending_adoption,
"markdown": "\n".join(md_lines),
}
# 分析结果写入当前租户知识库(平台知识模块,按登录租户隔离)
knowledge_ingest: dict[str, Any] = {}
if apply:
try:
has_data = bool(
result["summary"]["materials"]
or result["summary"]["orders"]
or result["summary"]["equipment"]
or result["summary"]["bom"]
or result.get("sourcePaths")
)
if ingest_project_analyze_report is not None and has_data:
knowledge_ingest = ingest_project_analyze_report(result, world=analysis_world)
if knowledge_ingest.get("assetId"):
md_lines.append("")
md_lines.append(
f"已写入知识库:《{knowledge_ingest.get('title')}》"
f"({knowledge_ingest.get('chunkCount')} 段,{knowledge_ingest.get('version')})"
)
result["markdown"] = "\n".join(md_lines)
except Exception as exc:
knowledge_ingest = {"error": str(exc)}
result["knowledgeIngest"] = knowledge_ingest
return result