certflow.services.correction_queue_service 源代码

# 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}