certflow.controllers.print_queue_controller 源代码

"""打印队列控制器

协调打印队列的异步执行与 UI 交互。封装 QueueWorker 线程与信号机制,
通过 Qt 信号槽实现非阻塞的队列处理,支持入队、暂停、继续、取消、置顶、上移。

架构位置:View → PrintQueueController → PrintQueueService(DB)/ CertificatePrintService(打印)
队列持久化于 print_jobs 表,断电重启后可恢复。

信号约定(小写_下划线 + 语义后缀):
    queue_changed:  队列内容变化(增删/状态变更/重排)时发出
    job_started:    某任务开始打印 (job_id)
    job_finished:   某任务结束 (job_id, success, message)
    status_changed: 控制器状态文本变化 (message)
    progress_updated: 当前任务进度 (current, total)
"""

from __future__ import annotations

from collections.abc import Callable
from typing import TYPE_CHECKING

from PySide6.QtCore import QObject, QThread, Signal

from certflow.services.printer.print_queue import PrintJobStatus, PrintQueueService
from certflow.utils.logger import logger

if TYPE_CHECKING:
    from sqlalchemy.orm import Session

    from certflow.utils.database import DatabaseManager

# print_func 契约:(certificate_id, printer_name, session) -> (success: bool, message: str)
#   或 (success, message, os_job_id: int | None);os_job_id 为 spooler 作业 ID,
#   用于 OS 级取消(BLUEPRINT §2.3 偏差 5 作业 ID 捕获链路)。
PrintFunc = Callable[
    [int | None, str, "Session"],
    "tuple[bool, str] | tuple[bool, str, int | None]",
]


def _default_print_func(
    certificate_id: int | None, printer_name: str, session: Session
) -> tuple[bool, str]:
    """默认打印动作:调用 CertificatePrintService.batch_print 打印整批

    Args:
        certificate_id: Certificate.id(为 None 时直接视为失败)
        printer_name: 打印机名称
        session: 数据库会话(与队列共用,避免跨线程共享)

    Returns:
        tuple[bool, str]: (成功标志, 消息)
    """
    from certflow.services.certificate_print_service import CertificatePrintService

    if certificate_id is None:
        return False, "任务缺少 certificate_id,无法打印", None
    service = CertificatePrintService(session, printer_name=printer_name)
    result = service.batch_print(
        certificate_id=certificate_id,
        printer_name=printer_name or None,
        skip_supplement=True,
    )
    success = result.get("success", 0)
    failed = result.get("failed", 0)
    os_job_id = service.last_os_job_id
    if failed > 0 and success == 0:
        return False, result.get("message", "打印失败"), os_job_id
    return True, result.get("message", "打印完成"), os_job_id


[文档] class QueueWorker(QThread): """队列执行工作线程 在独立线程中循环取出 pending 任务并打印,避免阻塞 UI 主线程。 每完成一个任务后检查暂停/取消标志,保证安全中断。 Signals: job_started: 任务开始 (job_id) job_finished: 任务结束 (job_id, success, message) status_changed: 状态文本 (message) queue_changed: 队列内容变化 """ job_started = Signal(int) job_finished = Signal(int, bool, str) status_changed = Signal(str) queue_changed = Signal() def __init__(self, controller: PrintQueueController) -> None: """初始化工作线程 Args: controller: 所属控制器(读取其标志与打印动作) """ super().__init__() self._controller = controller self._stop = False
[文档] def stop(self) -> None: """请求停止工作线程""" self._stop = True
[文档] def run(self) -> None: # noqa: D102 (QThread 覆写) while not self._stop: ctrl = self._controller if ctrl._paused: self.status_changed.emit("队列已暂停") self.msleep(200) continue session = ctrl._session_factory() try: queue = PrintQueueService(session) job = queue.get_next_pending() if job is None: # 补丁38:空闲即退出线程(而非永久 sleep 死守),使「入队」仅排队、 # 不被常驻 worker 立即打印;需要连续批量时由「启动」/「入队启动」唤醒。 self.status_changed.emit("队列空闲,工作线程退出") return # —— 补丁38 修正:队列清空后工作线程自动退出 —— # 旧实现在 job is None 时 msleep(300)+continue 永久死守线程,导致一旦 # 启动过一次(「入队启动」或队列面板「启动」),worker 常驻后台,后续 # 任何「入队」(本应仅排队)都会被立即取走打印,与「入队启动」无差别。 # 现改为空闲即退出线程,仅「启动」/「入队启动」显式唤醒,确保 # 「入队」=仅入队(PENDING),「入队启动」=入队+打印,二者为明确的 # 父子/短路关系。 queue.mark_running(job.id) self.job_started.emit(job.id) self.status_changed.emit(f"正在打印任务 #{job.id}") self.queue_changed.emit() try: result = ctrl._print_func(job.certificate_id, ctrl._printer_name, session) # 兼容 2 元组 (success, message) 与 3 元组 (success, message, os_job_id) if isinstance(result, tuple): success, message, os_job_id = (tuple(result) + (None,))[:3] else: success, message, os_job_id = result, "", None except Exception as e: # noqa: BLE001 - 打印异常需落库,不崩溃 logger.exception(f"打印任务 #{job.id} 执行异常") success, message, os_job_id = False, str(e), None # 回填 OS 打印作业 ID(偏差 5 作业 ID 捕获链路):打印动作产生的 # spooler 作业 ID 持久化到任务,供 OS 级取消使用。 if os_job_id is not None: queue.set_os_job_id(job.id, os_job_id) # 处理运行中的取消标记 if ctrl._cancel_job_id == job.id: queue.mark_cancelled(job.id, "打印中取消") self.job_finished.emit(job.id, False, "已取消") elif success: queue.mark_done(job.id) self.job_finished.emit(job.id, True, message) else: queue.mark_error(job.id, message) self.job_finished.emit(job.id, False, message) self.queue_changed.emit() self.status_changed.emit("任务完成,继续下一任务") finally: session.close()
[文档] class PrintQueueController(QObject): """打印队列控制器 管理队列生命周期:入队、启动、暂停、继续、取消、重排、清理与恢复。 通过 Qt 信号槽与视图解耦。 Attributes: _session_factory: 返回新 Session 的工厂(多线程各自创建,避免共享) _printer_name: 默认打印机名称 _print_func: 打印动作(可注入,便于测试) _worker: 当前队列工作线程 _paused: 暂停标志 _cancel_job_id: 运行中待取消的任务 ID(None 表示无) """ queue_changed = Signal() """队列内容变化信号""" job_started = Signal(int) """任务开始信号 (job_id)""" job_finished = Signal(int, bool, str) """任务结束信号 (job_id, success, message)""" status_changed = Signal(str) """状态文本信号 (message)""" progress_updated = Signal(int, int) """进度信号 (current, total)""" def __init__( self, db_manager: DatabaseManager | None = None, session_factory: Callable[[], Session] | None = None, printer_name: str = "", print_func: PrintFunc | None = None, ) -> None: """初始化打印队列控制器 Args: db_manager: 数据库管理器(提供 get_session);与 session_factory 二选一 session_factory: 自定义会话工厂(优先于 db_manager) printer_name: 默认打印机名称 print_func: 自定义打印动作(测试注入用) """ super().__init__() if session_factory is not None: self._session_factory: Callable[[], Session] = session_factory elif db_manager is not None: self._session_factory = db_manager.get_session else: from certflow.bootstrap import get_default_context self._session_factory = lambda: get_default_context().get_session() self._printer_name = printer_name self._print_func: PrintFunc = print_func or _default_print_func self._worker: QueueWorker | None = None self._paused = False self._cancel_job_id: int | None = None # ============================================================ # 入队 # ============================================================
[文档] def enqueue( self, certificate_id: int | None, sn: str | None = None, priority: int = 0, *, replace_pending: bool = False, ) -> int: """入队一个打印任务 Args: certificate_id: 关联合格证批次 ID sn: 出厂编号(可选) priority: 优先级 replace_pending: 为 True 时入队前先取消该证已有 pending 任务,避免 重复点击/订正后再入队产生同证多份 pending 而重复打印。 Returns: int: 新任务 ID """ session = self._session_factory() try: svc = PrintQueueService(session) job = svc.enqueue( certificate_id, sn=sn, priority=priority, replace_pending=replace_pending ) job_id = job.id finally: session.close() self.queue_changed.emit() return job_id
[文档] def enqueue_certificates( self, certificate_ids: list[int], priority: int = 0, *, replace_pending: bool = False ) -> list[int]: """批量入队多个 Certificate Args: certificate_ids: Certificate.id 列表 priority: 统一优先级 replace_pending: 为 True 时每个证入队前先取消其已有 pending 任务(去重)。 Returns: list[int]: 创建的任务 ID 列表 """ session = self._session_factory() try: svc = PrintQueueService(session) jobs = svc.bulk_enqueue( [(cid, None) for cid in certificate_ids], priority=priority, replace_pending=replace_pending, ) ids = [j.id for j in jobs] finally: session.close() self.queue_changed.emit() return ids
# ============================================================ # 队列控制 # ============================================================
[文档] def start(self) -> None: """启动队列工作线程(若未运行)""" if self._worker is not None and self._worker.isRunning(): return self._worker = QueueWorker(self) self._worker.job_started.connect(self.job_started) self._worker.job_finished.connect(self.job_finished) self._worker.status_changed.connect(self.status_changed) self._worker.queue_changed.connect(self.queue_changed) self._worker.start() self.status_changed.emit("队列已启动")
[文档] def stop(self) -> None: """停止队列工作线程(等待退出)""" if self._worker is not None: self._worker.stop() self._worker.wait(3000) self._worker = None self.status_changed.emit("队列已停止")
[文档] def pause(self) -> None: """暂停队列(不再启动新任务;当前任务完成后停止) 注:中途 ESC/P-K 指令流的安全中断依赖真机验证(P0), 当前版本采用任务边界暂停语义。interrupt_current() 预留 OS 级取消钩子。 """ self._paused = True self.status_changed.emit("队列暂停中") self.queue_changed.emit()
[文档] def resume(self) -> None: """继续队列""" self._paused = False self.status_changed.emit("队列继续") if self._worker is None or not self._worker.isRunning(): self.start() self.queue_changed.emit()
[文档] def is_paused(self) -> bool: """是否处于暂停状态 Returns: bool: True 表示已暂停 """ return self._paused
[文档] def is_running(self) -> bool: """工作线程是否在运行 Returns: bool: True 表示运行中 """ return self._worker is not None and self._worker.isRunning()
# ============================================================ # 任务操作 # ============================================================
[文档] def cancel_job(self, job_id: int) -> None: """取消任务 - pending:立即标记取消(不再打印) - running:设置取消标记,当前任务结束后标记为取消,并尝试 OS 级中断 Args: job_id: 任务 ID """ session = self._session_factory() try: svc = PrintQueueService(session) job = svc.get_job(job_id) if job is None: return if job.status == PrintJobStatus.PENDING: svc.mark_cancelled(job_id) elif job.status == PrintJobStatus.RUNNING: self._cancel_job_id = job_id self._interrupt_current(job_id) else: # 终态任务无需操作 return finally: session.close() self.queue_changed.emit()
def _interrupt_current(self, job_id: int) -> None: """OS 级中断当前打印作业(best-effort,依赖真实打印机) 经捕获的 OS spooler 作业 ID(``PrintJob.os_job_id``,由偏差 5 作业 ID 捕获链路从 ``LQ635KIIPrinter._flush`` / ``CertPrintEngine`` 透传而来) 调用 ``PrinterManager.cancel_job`` 真正下发 ``JOB_CONTROL_CANCEL``。 headless / 无 win32print / 无作业 ID 时安全降级(仅记日志)。 Args: job_id: 任务 ID(用于日志关联与读取 os_job_id) """ try: from certflow.services.printer_manager import PrinterManager session = self._session_factory() try: svc = PrintQueueService(session) job = svc.get_job(job_id) os_job_id = job.os_job_id if job is not None else None finally: session.close() if os_job_id is None: logger.info( f"中断运行中任务 #{job_id}:无 OS 作业 ID(headless 或驱动未捕获)," "跳过 OS 级取消" ) return manager = PrinterManager() ok = manager.cancel_job(os_job_id, self._printer_name or None) if ok: logger.info(f"已下发 OS 级取消:任务 #{job_id} → spooler 作业 {os_job_id}") else: logger.warning( f"OS 级取消未生效:任务 #{job_id} → spooler 作业 {os_job_id}" "(可能已完成或打印机不可达)" ) except Exception as e: # noqa: BLE001 logger.debug(f"OS 级中断不可用: {e}")
[文档] def move_top(self, job_id: int) -> None: """置顶任务""" session = self._session_factory() try: PrintQueueService(session).reorder_top(job_id) finally: session.close() self.queue_changed.emit()
[文档] def move_up(self, job_id: int) -> None: """上移任务""" session = self._session_factory() try: PrintQueueService(session).reorder_up(job_id) finally: session.close() self.queue_changed.emit()
[文档] def clear_finished(self) -> int: """清理已完成任务 Returns: int: 删除的任务数 """ session = self._session_factory() try: deleted = PrintQueueService(session).clear_finished() finally: session.close() self.queue_changed.emit() return deleted
[文档] def recover(self) -> int: """启动时恢复未完成任务 Returns: int: 被重置为 pending 的任务数 """ session = self._session_factory() try: recovered = PrintQueueService(session).recover() finally: session.close() self.queue_changed.emit() return recovered
# ============================================================ # 查询 # ============================================================
[文档] def list_jobs(self, include_terminal: bool = True) -> list[dict]: """列出队列任务(字典形式,供视图渲染) Args: include_terminal: 是否包含终态任务 Returns: list[dict]: 任务字典列表 """ session = self._session_factory() try: jobs = PrintQueueService(session).list_jobs(include_terminal) return [j.to_dict() for j in jobs] finally: session.close()
[文档] def count_pending(self) -> int: """统计 pending 任务数 Returns: int: pending 数量 """ session = self._session_factory() try: return PrintQueueService(session).count_pending() finally: session.close()
[文档] def set_printer(self, printer_name: str) -> None: """设置默认打印机 Args: printer_name: 打印机名称 """ self._printer_name = printer_name
# ============================================================ # 静态门面:收口 QueuePanel 对 model.PrintJobStatus 枚举的直连 # (即使 PrintQueueController 未注入实例,视图仍可经类名调用,向后兼容) # ============================================================ _STATUS_LABELS: dict[str, str] = { PrintJobStatus.PENDING: "待打印", PrintJobStatus.RUNNING: "打印中", PrintJobStatus.PAUSED: "已暂停", PrintJobStatus.DONE: "已完成", PrintJobStatus.CANCELLED: "已取消", } _STATUS_COLORS: dict[str, str] = { PrintJobStatus.PENDING: "#000000", PrintJobStatus.RUNNING: "#0078d4", PrintJobStatus.PAUSED: "#d83b01", PrintJobStatus.DONE: "#107c10", PrintJobStatus.CANCELLED: "#888888", }
[文档] @classmethod def status_label(cls, status: str) -> str: """打印任务状态中文标签(替代视图对 ``PrintJobStatus`` 枚举的直连)。""" return cls._STATUS_LABELS.get(status, status)
[文档] @classmethod def status_color(cls, status: str) -> str: """打印任务状态配色(替代视图对 ``PrintJobStatus`` 枚举的直连)。""" return cls._STATUS_COLORS.get(status, "#000000")