# certflow/services/correction_queue_service.py
"""#30 P1 校正队列服务(黄标校正 + 隔离释放 + upsert 自学习)
替代 VBA D 阶段的「InputBox 人工兜底 + Stop 硬断 + CanShuSheet 回写」:
- 导入后黄标行(flag=True:口径/型号无标准号待复核)与隔离行(quarantine)
统一在校正队列 UI 列出;
- 人工校正后 upsert 进字典表(CaliberMapping / ModelParamMapping),
下次同值导入即自动命中(等价 VBA 自学习回写,但非阻塞、可追溯、团队共享);
- 隔离行可释放重建为 SalePlan(补全要货单号后回归正常导入)。
架构层次:View (CorrectionQueueView) → Controller (CorrectionQueueController)
→ Service (CorrectionQueueService) → Model (SalePlan/QuarantineSalePlan/字典表)
"""
from __future__ import annotations
from datetime import datetime
from typing import Any
from sqlalchemy.orm import Session
from certflow.models.caliber_mapping import CaliberMapping
from certflow.models.model_param_mapping import ModelParamMapping
from certflow.models.param_auto_learn_log import ParamAutoLearnLog
from certflow.models.quarantine_sale_plan import QuarantineSalePlan
from certflow.models.sale_plan import SalePlan
from certflow.utils.logger import logger
[文档]
class CorrectionQueueService:
"""校正队列服务:黄标校正 + 隔离释放 + 字典自学习。"""
def __init__(self, session: Session) -> None:
self.session = session
# ============================================================
# 查询:列出待校正
# ============================================================
[文档]
def list_flagged(self, limit: int = 500) -> list[SalePlan]:
"""列出黄标(flag=True)的销售计划记录(待人工复核口径/标准号)。"""
return (
self.session.query(SalePlan)
.filter(SalePlan.flag == True) # noqa: E712
.order_by(SalePlan.id)
.limit(limit)
.all()
)
[文档]
def list_quarantined(self, limit: int = 500) -> list[QuarantineSalePlan]:
"""列出未释放(resolved=False)的隔离记录。"""
return (
self.session.query(QuarantineSalePlan)
.filter(QuarantineSalePlan.resolved == False) # noqa: E712
.order_by(QuarantineSalePlan.id)
.limit(limit)
.all()
)
[文档]
def get_quarantine(self, quarantine_id: int) -> QuarantineSalePlan | None:
return self.session.get(QuarantineSalePlan, quarantine_id)
[文档]
def get_sale_plan(self, sale_plan_id: int) -> SalePlan | None:
return self.session.get(SalePlan, sale_plan_id)
# ============================================================
# 校正黄标行 + 字典自学习
# ============================================================
[文档]
def correct_sale_plan(
self,
sale_plan_id: int,
corrections: dict[str, Any],
learn: bool = True,
recheck: bool = True,
) -> tuple[SalePlan | None, list[str]]:
"""校正一条黄标销售计划记录,并可选 upsert 进字典表自学习。
补录后做**完整性核对**:仅当该行全部必填空白(口径 + 公称压力)都已
解决才清除黄标;仍有缺口则保留黄标,等后续续补。避免「补了一个口径就
默认整行补齐、误清黄标」(标准号导入期留空、打印回填,不计入黄标)。
**合同口径保护**:`product_spec`(合同口径)与 `product_model`(型号)
参与导入去重的唯一键计算,校正**一律不修改**这两列(空白必须保持空白,
否则重导入时 key 变化无法与原始记录对应)。从型号等推得的口径改写入
**合格证字段** `cert_product_spec`(口径串)与 `spec_norm`(规范化数字,
不参与重导入 key,可安全写)。
Args:
sale_plan_id: 销售计划主键
corrections: 校正字段。可写列:`spec_norm` / `pressure_value` /
`test_standard` / `cert_product_spec`。`product_spec` /
`product_model` 即使传入也被忽略(保护唯一键,仅作学习键读取)。
learn: 是否把校正结果 upsert 进字典表(CaliberMapping/ModelParamMapping),
默认 True。仅校正字段非空且对应原始键存在时才学习。
recheck: 是否做完整性核对决定黄标去留(默认 True)。为 False 时
强制清除黄标(对应 UI「仅标记已复核」)。
Returns:
(更新后的 SalePlan, 仍未解决的字段中文名列表);id 不存在返回 (None, [])。
列表非空表示黄标已保留、仍有 N 处待补录。
"""
sp = self.session.get(SalePlan, sale_plan_id)
if sp is None:
return None, []
old_spec = sp.product_spec or ""
old_model = sp.product_model or ""
applied = False
# 应用校正到可写字段。
# 注意:product_spec / product_model 不在可写列——它们参与导入唯一键,
# 校正不得修改(合同口径空白须保持空白);推得的口径写 cert_product_spec。
for field in (
"spec_norm",
"pressure_value",
"test_standard",
"cert_product_spec",
):
if field in corrections and corrections[field] is not None:
setattr(sp, field, corrections[field])
applied = True
# 自学习:upsert 进字典表(非阻塞,失败仅记日志)
if applied and learn:
self._learn_dicts(sp, old_spec, old_model, corrections)
# 完整性核对:口径与公称压力均解决才清黄标,否则保留待续补
remaining = self._remaining_manual_fields(sp)
if recheck:
sp.flag = bool(remaining)
else:
sp.flag = False # 人工确认已复核,强制清除
remaining = []
sp.updated_at = datetime.now()
self.session.commit()
return sp, remaining
def _remaining_manual_fields(self, sp: SalePlan) -> list[str]:
"""核对该行仍未解决的黄标字段(口径与 IDGenerator 黄标口径一致)。
- 口径:合同口径(product_spec)或合格证口径(cert_product_spec)有来源,
但 spec_norm 非合法数字口径 → 未解决;合同口径空白且未补合格证口径
时不强制(与 IDGenerator「空白口径不黄标」一致);
- 公称压力:product_model 非空但 pressure_value 为空 → 未解决;
- 标准号(test_standard)不计入(导入期留空、由打印回填)。
"""
from certflow.handlers.id_generator import IDGenerator
remaining: list[str] = []
raw_spec = str(sp.product_spec or "").strip()
cert_spec = str(sp.cert_product_spec or "").strip()
spec_norm = str(sp.spec_norm or "").strip()
if (raw_spec or cert_spec) and not IDGenerator._is_numeric_mm(spec_norm):
remaining.append("口径")
model = str(sp.product_model or "").strip()
pressure = str(sp.pressure_value or "").strip()
if model and not pressure:
remaining.append("公称压力")
return remaining
def _learn_dicts(
self,
sp: SalePlan,
old_spec: str,
old_model: str,
corrections: dict[str, Any],
) -> None:
"""把校正结果 upsert 进字典表 + 写自学习日志。
- 口径:product_spec(原始串) → spec_norm(标准化值) → CaliberMapping;
合同口径空白时改用合格证口径串(cert_product_spec)作为学习键。
- 型号/压力/标准号:product_model → pressure_value/test_standard → ModelParamMapping
"""
try:
spec_raw = str(corrections.get("product_spec") or old_spec or "").strip()
if not spec_raw:
# 合同口径空白:用合格证补录的口径串作为口径字典学习键
spec_raw = str(
corrections.get("cert_product_spec") or sp.cert_product_spec or ""
).strip()
spec_val = str(corrections.get("spec_norm") or sp.spec_norm or "").strip()
if spec_raw and spec_val:
self._upsert_caliber(spec_raw, spec_val)
model = str(corrections.get("product_model") or old_model or "").strip()
pn = str(corrections.get("pressure_value") or sp.pressure_value or "").strip()
std = str(corrections.get("test_standard") or sp.test_standard or "").strip()
if model and (pn or std):
self._upsert_model_param(model, pn, std)
except Exception as e: # 自学习失败不应阻断校正
logger.warning(f"校正队列字典自学习失败(忽略,校正已生效): {e}")
def _upsert_caliber(self, raw_text: str, caliber_value: str) -> None:
"""upsert 口径字典(CaliberMapping)+ 自学习日志。"""
existing = (
self.session.query(CaliberMapping).filter(CaliberMapping.raw_text == raw_text).first()
)
if existing:
existing.caliber_value = caliber_value
existing.source = "auto_learn"
existing.updated_at = datetime.now()
else:
self.session.add(
CaliberMapping(
raw_text=raw_text,
caliber_value=caliber_value,
is_imperial=False,
source="auto_learn",
usage_count=0,
)
)
self.session.add(
ParamAutoLearnLog(
table_name="caliber_mappings",
field_name="caliber_value",
lookup_key=raw_text,
user_value=caliber_value,
)
)
def _upsert_model_param(
self, product_model: str, pressure_value: str, test_standard: str
) -> None:
"""upsert 型号参数字典(ModelParamMapping)+ 自学习日志。"""
existing = (
self.session.query(ModelParamMapping)
.filter(ModelParamMapping.product_model == product_model)
.first()
)
if existing:
if pressure_value:
existing.pressure_value = pressure_value
if test_standard:
existing.test_standard = test_standard
existing.source = "auto_learn"
existing.updated_at = datetime.now()
else:
self.session.add(
ModelParamMapping(
product_model=product_model,
pressure_value=pressure_value or "",
test_standard=test_standard or "",
source="auto_learn",
usage_count=0,
)
)
if pressure_value:
self.session.add(
ParamAutoLearnLog(
table_name="model_param_mappings",
field_name="pressure_value",
lookup_key=product_model,
user_value=pressure_value,
)
)
if test_standard:
self.session.add(
ParamAutoLearnLog(
table_name="model_param_mappings",
field_name="test_standard",
lookup_key=product_model,
user_value=test_standard,
)
)
# ============================================================
# 隔离行释放:重建为 SalePlan
# ============================================================
[文档]
def release_quarantine(
self,
quarantine_id: int,
corrections: dict[str, Any] | None = None,
) -> SalePlan | None:
"""释放一条隔离记录:用原始快照 + 人工校正重建 SalePlan,并标记隔离已处理。
Args:
quarantine_id: 隔离记录主键
corrections: 人工补全字段(通常含 sales_order_no/production_order_no/plan_no
等要货单号,使其不再被门控)
Returns:
新建的 SalePlan;id 不存在或已释放返回 None。
"""
q = self.session.get(QuarantineSalePlan, quarantine_id)
if q is None or q.resolved:
return None
corrections = corrections or {}
raw = q.to_dict()
# 原始记录快照优先,校正覆盖
base: dict[str, Any] = dict(raw.get("raw_record") or {})
# 隔离表列字段作为回退(快照可能不含映射后的列名)
for k in (
"contract_no",
"sales_order_no",
"production_order_no",
"plan_no",
"plan_date",
"customer",
"project_name",
"product_name",
"product_model",
"product_spec",
"spec_norm",
"quantity",
"pressure_value",
"test_standard",
"outsource_type",
"shipping_status",
):
if not base.get(k) and raw.get(k):
base[k] = raw[k]
# 人工校正覆盖
base.update({k: v for k, v in corrections.items() if v is not None})
# 重建 SalePlan(走 _create_sale_plan 的轻量等价:直接构造并提交)
sp = self._build_sale_plan(base, q)
self.session.add(sp)
self.session.flush()
# 标记隔离已释放
q.resolved = True
q.resolved_at = datetime.now()
self.session.commit()
return sp
def _build_sale_plan(self, rec: dict[str, Any], q: QuarantineSalePlan) -> SalePlan:
"""从隔离记录 + 校正构造 SalePlan(最小必填字段)。"""
from certflow.handlers.id_generator import IDGenerator
sp = SalePlan(
product_model=str(rec.get("product_model") or "").strip() or "(未填型号)",
contract_no=str(rec.get("contract_no") or ""),
sales_order_no=str(rec.get("sales_order_no") or ""),
production_order_no=str(rec.get("production_order_no") or ""),
plan_no=str(rec.get("plan_no") or ""),
plan_date=str(rec.get("plan_date") or ""),
customer=str(rec.get("customer") or ""),
project_name=str(rec.get("project_name") or ""),
product_name=str(rec.get("product_name") or ""),
product_spec=str(rec.get("product_spec") or ""),
spec_norm=str(rec.get("spec_norm") or ""),
quantity=rec.get("quantity") or 1,
pressure_value=str(rec.get("pressure_value") or ""),
test_standard=str(rec.get("test_standard") or ""),
outsource_type=str(rec.get("outsource_type") or ""),
shipping_status=str(rec.get("shipping_status") or ""),
source_file=q.source_file,
source_sheet=q.source_sheet,
import_batch_id=q.import_batch_id,
import_date=datetime.now(),
flag=False, # 释放即视为已人工复核
)
# 生成唯一键(复用分层去重策略)
try:
sp.unique_key = IDGenerator.generate_unique_key(rec)
except Exception as e:
logger.warning(f"释放隔离行生成 unique_key 失败(用默认): {e}")
return sp
# ============================================================
# 隔离行清理
# ============================================================
[文档]
def delete_quarantine(self, quarantine_id: int) -> bool:
"""直接删除一条隔离记录(确认无效、不释放)。"""
q = self.session.get(QuarantineSalePlan, quarantine_id)
if q is None:
return False
self.session.delete(q)
self.session.commit()
return True
[文档]
def quarantine_stats(self) -> dict[str, int]:
"""隔离统计:未释放 / 已释放 / 总数。"""
total = self.session.query(QuarantineSalePlan).count()
resolved = (
self.session.query(QuarantineSalePlan)
.filter(QuarantineSalePlan.resolved == True) # noqa: E712
.count()
)
return {"total": total, "resolved": resolved, "pending": total - resolved}