# src/certflow/controllers/query_controller.py
"""查询控制器
负责查询视图的业务逻辑处理,将 UI 层与数据层解耦。
职责:
- 删除记录(含级联删除、事务回滚)
- 批量更新字段
- 导出选中记录到 Excel
- 自动编号并创建合格证
- 查询数据
- 获取字段可选值
架构:
View (QueryView) → Controller (QueryController) → Service (QueryService/SalePlanService)
→ Handler (ShippedExportHandler)
Usage:
controller = QueryController(session)
success, count, msg = controller.delete_records([{"id": 1}, {"id": 2}])
"""
from __future__ import annotations
from typing import Any
from PySide6.QtCore import QObject, QThread, Signal
from sqlalchemy.orm import Session
from certflow.services import QueryService
from certflow.services.certificate_number_service import CertificateNumberService
from certflow.services.certificate_print_service import CertificatePrintService
from certflow.services.sale_plan_service import SalePlanService
from certflow.utils.logger import logger
[文档]
class FieldValuesWorker(QThread):
"""异步加载字段可选值的工作线程
避免同步 SQL 查询阻塞 UI 线程,通过信号将结果传回主线程。
"""
finished = Signal(str, list)
"""完成信号: (field_name, values)"""
error = Signal(str)
"""错误信号: error_message"""
def __init__(self, service: QueryService, field_name: str, order_desc: bool = True) -> None:
"""初始化工作线程
Args:
service: QueryService 实例,用于执行字段值查询
field_name: 要加载可选值的字段名
order_desc: 是否按数量降序排列,默认为 True
"""
super().__init__()
self.service = service
self.field_name = field_name
self.order_desc = order_desc
[文档]
def run(self) -> None:
"""线程主执行方法(由 QThread.start() 触发)
通过 service 查询字段可选值,成功后发出 finished 信号,
失败则发出 error 信号。
Returns:
None
"""
try:
values = self.service.get_field_values(self.field_name, order_desc=self.order_desc)
self.finished.emit(self.field_name, values)
except Exception as e:
logger.error(f"加载字段值失败 [{self.field_name}]: {e}")
self.error.emit(str(e))
[文档]
class QueryController(QObject):
"""查询控制器 — 业务逻辑层"""
# ============================================================
# Signals
# ============================================================
data_updated = Signal()
"""数据变更信号(删除/更新后触发,视图应刷新)"""
operation_completed = Signal(str)
"""操作完成信号,携带状态消息"""
operation_failed = Signal(str)
"""操作失败信号,携带错误信息"""
field_values_loaded = Signal(str, list)
"""字段可选值加载完成信号: (field_name, values)"""
field_values_load_failed = Signal(str, str)
"""字段可选值加载失败信号: (field_name, error_message)"""
# ============================================================
# Initialization
# ============================================================
def __init__(
self,
session: Session,
db_manager: Any = None,
parent: QObject | None = None,
) -> None:
"""初始化查询控制器
Args:
session: SQLAlchemy 数据库会话,供 Service 层使用
db_manager: 数据库管理器(可选),预留用于共享资源管理
parent: Qt 父对象(可选)
"""
super().__init__(parent)
self.session = session
self.db_manager = db_manager
self.query_service = QueryService(session)
# ============================================================
# 查询
# ============================================================
[文档]
def query(
self,
conditions: dict[str, Any],
page: int,
page_size: int,
order_by: str,
order_desc: bool,
or_filters: Any = None,
) -> dict[str, Any]:
"""执行分页查询
Args:
conditions: 筛选条件字典
page: 页码
page_size: 每页条数
order_by: 排序字段
order_desc: 是否降序
Returns:
查询结果字典,包含 records、page、total_pages、total
Examples:
>>> controller = QueryController(session)
>>> result = controller.query(
... conditions={"product_model": "DN100"},
... page=1,
... page_size=20,
... order_by="id",
... order_desc=True,
... )
>>> print(result["total"], result["records"])
"""
return self.query_service.query(
conditions=conditions,
page=page,
page_size=page_size,
order_by=order_by,
order_desc=order_desc,
or_filters=or_filters,
)
[文档]
def get_field_values(self, field_name: str, order_desc: bool = True) -> list[Any]:
"""获取字段的所有可选值(同步,会阻塞 UI 线程)
推荐使用 load_field_values_async() 异步加载。
Args:
field_name: 字段名
order_desc: 是否按数量降序排列
Returns:
字段值列表
"""
return self.query_service.get_field_values(field_name, order_desc=order_desc)
[文档]
def get_distinct_combos(
self, fields: list[str], conditions: dict[str, Any] | None = None
) -> list[str]:
"""获取字段组合的去重候选项(全量,跨所有分页),用于快查框。
Args:
fields: 字段名列表
conditions: 当前筛选条件(可选)
Returns:
去重组合字符串列表
"""
return self.query_service.get_distinct_field_combos(fields, conditions)
[文档]
def load_field_values_async(
self, field_name: str, order_desc: bool = True
) -> FieldValuesWorker:
"""异步加载字段可选值(推荐)
通过后台线程加载,完成后通过 field_values_loaded 信号通知。
Args:
field_name: 字段名
order_desc: 是否按数量降序排列
Returns:
FieldValuesWorker 实例(调用方需保持引用避免 GC)
"""
worker = FieldValuesWorker(self.query_service, field_name, order_desc)
worker.finished.connect(lambda fn, vals: self.field_values_loaded.emit(fn, vals))
worker.error.connect(lambda err: self.field_values_load_failed.emit(field_name, err))
worker.start()
return worker
[文档]
def get_records_by_ids(self, ids: list[int]) -> list[Any]:
"""根据 ID 列表获取记录
Args:
ids: 记录 ID 列表
Returns:
记录列表
"""
return self.query_service.get_records_by_ids(ids)
[文档]
def get_statistics(self, conditions: dict[str, Any]) -> dict[str, Any]:
"""获取统计信息
Args:
conditions: 筛选条件
Returns:
统计信息字典
"""
return self.query_service.get_statistics(conditions)
# ============================================================
# 删除记录
# ============================================================
[文档]
def delete_records(self, records: list[dict]) -> tuple[bool, int, str]:
"""删除记录(事务回滚)
Args:
records: 要删除的记录列表,每项包含 id
Returns:
tuple[bool, int, str]: (是否成功, 删除数量, 消息)
"""
if not records:
return False, 0, "没有要删除的记录"
ids = [r["id"] for r in records]
success, count, msg = self.query_service.delete_records(ids)
if success:
self.data_updated.emit()
self.operation_completed.emit(f"已删除 {count} 条记录")
else:
self.operation_failed.emit(msg)
return success, count, msg
[文档]
def delete_with_certificates(self, records: list[dict]) -> tuple[bool, int, str]:
"""级联删除记录及关联的合格证
Args:
records: 要删除的记录列表,每项包含 id
Returns:
tuple[bool, int, str]: (是否成功, 删除销售计划数量, 消息)
"""
if not records:
return False, 0, "没有要删除的记录"
ids = [r["id"] for r in records]
success, count, cert_count, msg = self.query_service.delete_with_certificates(ids)
if success:
self.data_updated.emit()
self.operation_completed.emit(msg)
else:
self.operation_failed.emit(msg)
return success, count, msg
[文档]
def check_certificate_exists(self, ids: list[int]) -> int:
"""检查记录是否有关联合格证
Args:
ids: 记录 ID 列表
Returns:
关联的合格证数量
"""
return self.query_service.get_certificate_count(ids)
# ============================================================
# 批量更新
# ============================================================
[文档]
def batch_update_field(self, ids: list[int], field_name: str, value: Any) -> int:
"""批量更新字段(事务回滚)
Args:
ids: 记录 ID 列表
field_name: 字段名
value: 要设置的值
Returns:
更新记录数
"""
if not ids:
return 0
updated = self.query_service.batch_update_field(ids, field_name, value)
self.data_updated.emit()
self.operation_completed.emit(f"已更新 {updated} 条")
return updated
[文档]
def batch_update_production_status(
self, ids: list[int], status: str, update_relations: list[dict] | None = None
) -> int:
"""批量更新生产状态(支持关联字段更新)
Args:
ids: 记录 ID 列表
status: 目标状态
update_relations: 关联更新规则列表
Returns:
更新记录数
"""
if not ids:
return 0
updated = self.query_service.batch_update_production_status(ids, status, update_relations)
self.data_updated.emit()
self.operation_completed.emit(f"已更新 {updated} 条 → {status}")
return updated
# ============================================================
# 导出
# ============================================================
[文档]
def export_selected(self, ids: list[int], columns: list[dict], file_path: str) -> int:
"""导出选中记录到 Excel
Args:
ids: 记录 ID 列表
columns: 列定义列表
file_path: 导出文件路径
Returns:
导出记录数
"""
if not ids:
return 0
records = self.get_records_by_ids(ids)
service = SalePlanService(self.session)
count = service.export_query_result(records, columns, file_path)
self.operation_completed.emit(f"已导出 {count} 条")
return count
[文档]
def archive_to_shipped(self, ids: list[int]) -> int:
"""将选中记录归档到已发货表(Layer 1 → Layer 2)
Args:
ids: 记录 ID 列表
Returns:
归档记录数
"""
service = SalePlanService(self.session)
count = service.archive_to_shipped(ids)
self.data_updated.emit()
self.operation_completed.emit(f"已归档 {count} 条到已发货表")
return count
[文档]
def export_shipped_and_delete(self, ids: list[int], output_dir: str) -> dict[str, Any]:
"""导出完整字段 Excel 并删除(Layer 1 → Layer 3)
Args:
ids: 记录 ID 列表
output_dir: 导出目录
Returns:
{"count": N, "filepath": "..."}
"""
service = SalePlanService(self.session)
result = service.export_and_delete(ids, output_dir)
self.data_updated.emit()
self.operation_completed.emit(f"已导出 {result['count']} 条到 {result['filepath']}")
return result
# ============================================================
# 打印相关
# ============================================================
[文档]
def generate_print_certificates(
self,
ids: list[int],
prefix_override: str | None = None,
print_status: str = "待打印",
) -> tuple[int, int]:
"""统一「打印合格证」入口:去重安全地编号并生成合格证。
合并原 ``auto_number_and_create_certs``(查询页「编号/自动编号」)与
``import_cert_list_from_sales_plan``(「导入合格证清单」)为单一入口。
业务编排在此层,实际生成委托 ``CertificatePrintService``:
跳过已生成合格证的计划(避免重复生成 / 唯一约束冲突),
自动编号并写入 Certificate(待打印),返回 (编号数, 合格证数)。
分层调用链:
View(_goto_print) → Controller(本方法) →
CertificatePrintService(编排) → CertNumberingService(数据访问) → Certificate(Model)。
Args:
ids: 选中的销售计划 ID 列表
prefix_override: 编号前缀覆盖(None 则按规则推导)
print_status: 生成的合格证打印状态,默认 "待打印"
Returns:
tuple[int, int]: (编号数, 合格证生成数)
"""
if not ids:
return 0, 0
print_service = CertificatePrintService(self.session)
numbered, cert_count = print_service.generate_print_certificates(
ids, prefix_override=prefix_override, print_status=print_status
)
self.data_updated.emit()
self.operation_completed.emit(f"已生成 {cert_count} 份合格证(编号 {numbered} 条)")
return numbered, cert_count
[文档]
def manual_number(
self,
ids: list[int],
prefix: str,
start: int,
quantity: int,
suffix: str = "",
) -> int:
"""手动编号并创建合格证
Args:
ids: 记录 ID 列表
prefix: 编号前缀
start: 起始序号
quantity: 每件数量
suffix: 后缀
Returns:
更新记录数
"""
if not ids:
return 0
print_service = CertificatePrintService(self.session)
updated = print_service.manual_number(ids, prefix, start, quantity, suffix)
logger.info(f"手动编号 {updated} 条(前缀={prefix})")
self.operation_completed.emit(f"手动编号 {updated} 条")
return updated
[文档]
def auto_number(self, ids: list[int], prefix_override: str | None = None) -> int:
"""自动编号(正则解析生产令号/供货类型,委托 CertificatePrintService.auto_number)。
为选中且 ``product_code`` 为空/占位符、``needs_numbering!=False`` 的记录生成产品编号,
写入 SalePlan 并即时刷新。对应工具栏「🔢 自动编号」按钮(修复「点击打印提示需先编号
却没有自动编号按钮」的体验缺口)。
Args:
ids: 记录 ID 列表
prefix_override: 手动指定的年月前缀(YYMM),None 取当前年月
Returns:
成功编号的记录数
"""
if not ids:
return 0
print_service = CertificatePrintService(self.session)
updated = print_service.auto_number(ids, prefix_override)
logger.info(f"自动编号 {updated} 条")
self.operation_completed.emit(f"自动编号 {updated} 条")
return updated
[文档]
def generate_batch(
self,
sale_plan_id: int,
batch_qty: int,
ym: str | None = None,
start_override: int | None = None,
) -> dict:
"""按「本批台数」为单个销售计划生成一批产品编号(G4 入口)。
委托 ``CertificateNumberService.generate_batch``:合同总数从 ``SalePlan.quantity``
读取,用 ``shipped_quantity`` 判定剩余与满发,超额直接抛 ``ValueError``(调用方捕获
并转用户提示)。对应工具栏「🔢 自动编号」弹窗输入本批台数的路径。
Args:
sale_plan_id: 销售计划 ID
batch_qty: 本批发货台数(> 0 且 <= 剩余可发)
ym: 目标编号年月(YYMM),None 取当前月
start_override: 起始流水手填覆盖,None 表示自动续号
Returns:
dict: 含 product_code / seq_start / seq_end / remaining / fully_shipped 等
"""
service = CertificateNumberService()
result = service.generate_batch(
self.session, sale_plan_id, batch_qty, ym=ym, start_override=start_override
)
logger.info(
f"批次编号完成 | id={sale_plan_id} | 本批={batch_qty} | "
f"{result.get('product_code')} | 剩余={result.get('remaining')}"
)
self.operation_completed.emit(
f"批次编号完成 | 本批 {batch_qty} 台 | 剩余 {result.get('remaining')} 台"
)
return result
[文档]
def create_certificates(self, ids: list[int]) -> int:
"""生成合格证记录
Args:
ids: 记录 ID 列表
Returns:
生成的合格证数量
"""
if not ids:
return 0
print_service = CertificatePrintService(self.session)
count = print_service.create_certificates(ids, printer_name="默认打印机")
logger.info(f"生成合格证 {count} 条")
self.operation_completed.emit(f"已生成 {count} 份合格证")
return count
# ============================================================
# 复制
# ============================================================
[文档]
def copy_cert_info_to_plan(self, ids: list[int]) -> int:
"""复制合同信息到合格证字段
Args:
ids: 记录 ID 列表
Returns:
更新记录数
"""
if not ids:
return 0
updated = self.query_service.copy_cert_info_to_plan(ids)
self.data_updated.emit()
self.operation_completed.emit(f"已复制 {updated} 条")
return updated
[文档]
def update_cert_fields(self, plan_id: int, fields: dict[str, Any]) -> bool:
"""订正单条记录的合格证打印字段(仅 cert_* + certificate_remarks)
委托 ``SalePlanService.update_cert_fields``,仅更新合格证字段,
不动合同字段。成功后在查询视图即时刷新。
Args:
plan_id: 销售计划记录主键
fields: 待写入字段字典(非合格证字段被服务层忽略)
Returns:
bool: 是否更新成功
"""
service = SalePlanService(self.session)
ok = service.update_cert_fields(plan_id, fields)
if ok:
self.data_updated.emit()
self.operation_completed.emit(f"已订正合格证信息(id={plan_id})")
return ok
[文档]
def update_certificate_fields(self, cert_id: int, fields: dict[str, Any]) -> bool:
"""订正单条合格证记录字段(白名单,对称于 update_cert_fields)
委托 ``CertificateService.update_fields``,仅更新合格证打印字段,
不动编号/状态等系统字段。成功后触发查询视图刷新。
Args:
cert_id: 合格证记录主键
fields: 待写入字段字典(非白名单字段被服务层忽略)
Returns:
bool: 是否更新成功
"""
from certflow.services.certificate_service import CertificateService
service = CertificateService(self.session)
ok = service.update_fields(cert_id, fields)
if ok:
self.data_updated.emit()
self.operation_completed.emit(f"已订正合格证记录(cert_id={cert_id})")
return ok
[文档]
def copy_number(self, ids: list[int], source_number: str) -> int:
"""复制编号到选中行
Args:
ids: 记录 ID 列表
source_number: 源编号
Returns:
更新记录数
"""
if not ids:
return 0
print_service = CertificatePrintService(self.session)
updated = print_service.copy_number(ids, source_number)
logger.info(f"复制编号 {updated} 条: {source_number}")
self.operation_completed.emit(f"复制编号 {updated} 条")
return updated
# ============================================================
# 填充
# ============================================================
[文档]
def fill_plan_no(self, ids: list[int], sequence: str) -> int:
"""填充计划单号
Args:
ids: 记录 ID 列表
sequence: 序号
Returns:
更新记录数
"""
if not ids:
return 0
updated = self.query_service.fill_plan_no(ids, sequence)
self.data_updated.emit()
self.operation_completed.emit(f"已填充 {updated} 条计划单号")
return updated
# ============================================================
# 隐藏/取消隐藏
# ============================================================
[文档]
def hide_records(self, ids: list[int]) -> int:
"""隐藏记录
Args:
ids: 记录 ID 列表
Returns:
更新记录数
"""
return self.batch_update_field(ids, "is_hidden", True)
[文档]
def unhide_records(self, ids: list[int]) -> int:
"""取消隐藏记录
Args:
ids: 记录 ID 列表
Returns:
更新记录数
"""
return self.batch_update_field(ids, "is_hidden", False)
# ============================================================
# 进度标记
# ============================================================
[文档]
def mark_document_done(self, ids: list[int], doc_type: str) -> dict[str, Any]:
"""标记选中行的文档完成
Args:
ids: 记录 ID 列表
doc_type: 文档类型 (certificate/nameplate/test_report/warranty/scan)
Returns:
操作结果字典,包含 success 和 msg
"""
from certflow.services import WorkflowService
service = WorkflowService(self.session)
result = service.batch_mark_done(ids, doc_type)
logger.info(f"标记文档完成: {doc_type}, 成功 {result.get('success', 0)} 条")
self.data_updated.emit()
self.operation_completed.emit(f"已标记 {result.get('success', 0)} 条")
return result
# ============================================================
# 导出 Access
# ============================================================
[文档]
def export_to_access(self, ids: list[int]) -> int:
"""导出到 Access 数据库
Args:
ids: 记录 ID 列表
Returns:
导出记录数
"""
print_service = CertificatePrintService(self.session)
count = print_service.export_to_access(ids)
logger.info(f"导出到 Access: {count} 条")
self.operation_completed.emit(f"已导出 {count} 条到 Access")
return count
[文档]
def get_number_prefix(self, prefix_override: str | None = None) -> str:
"""获取编号前缀
Args:
prefix_override: 前缀覆盖值
Returns:
编号前缀字符串
"""
service = CertificateNumberService(prefix_override=prefix_override)
return service.get_prefix()
[文档]
def get_certificates_by_plan_ids(self, plan_ids: list[int]) -> list[Any]:
"""根据销售计划 ID 查询关联的合格证记录
Args:
plan_ids: 销售计划 ID 列表
Returns:
合格证记录列表
"""
return self.query_service.get_certificates_by_plan_ids(plan_ids)
[文档]
def get_latest_certificate(self, plan_ids: list[int]) -> Any:
"""获取最近创建的一条合格证记录(用于参数补充)
Args:
plan_ids: 销售计划 ID 列表
Returns:
最近创建的合格证记录,或 None
"""
return self.query_service.get_latest_certificate(plan_ids)
[文档]
def supplement_certificate_params(
self,
cert: Any,
parent: Any = None,
supplement_callback: Any = None,
) -> dict[str, Any]:
"""补充合格证参数(委托 Service 层)
Args:
cert: 合格证记录
parent: 父窗口(用于对话框)
supplement_callback: 参数补充回调(UI 对话框)
Returns:
补充结果字典
"""
service = CertificatePrintService(self.session)
return service.check_and_supplement_params(
cert,
parent=parent,
supplement_callback=supplement_callback,
)