certflow.services.printer.print_queue 源代码

"""打印队列服务(持久化队列管理)

提供基于数据库(print_jobs 表)的打印任务队列管理能力,与 UI/线程解耦。
负责任务的入队、出队、状态机流转、优先级重排与断电恢复。

设计要点:
- 所有状态变更均落库(commit),保证断电可恢复;
- 出队顺序由 (priority DESC, created_at ASC) 决定,支持置顶/上移;
- 本服务为纯逻辑层,不直接调用打印机,打印动作由 QueueWorker 驱动。
"""

from __future__ import annotations

from datetime import datetime
from typing import TYPE_CHECKING

from sqlalchemy.orm import Session

from certflow.models.print_job import PrintJob, PrintJobStatus
from certflow.utils.database import DatabaseManager
from certflow.utils.logger import logger

if TYPE_CHECKING:
    from collections.abc import Sequence


[文档] class PrintQueueService: """打印队列服务 围绕 print_jobs 表的队列管理。每个实例绑定一个 SQLAlchemy Session, 调用方需保证该 Session 在正确的线程中使用(跨线程请各自创建实例)。 Examples: >>> svc = PrintQueueService(session) >>> job = svc.enqueue(certificate_id=1, priority=5) >>> next_job = svc.get_next_pending() >>> svc.mark_running(next_job.id) >>> svc.mark_done(next_job.id) """ def __init__(self, session: Session) -> None: """初始化打印队列服务 Args: session: SQLAlchemy 数据库会话对象 """ self.session = session # ============================================================ # 入队 # ============================================================
[文档] def cancel_pending_by_certificate(self, certificate_id: int) -> int: """取消指定合格证的所有 pending 任务(用于 replace_pending 去重)。 仅取消 **pending**(尚未开始)的任务;running/paused/终态任务不动—— 正在打印的任务中途取消有风险,且终态任务无需处理。 Args: certificate_id: 合格证批次 ID Returns: int: 被取消的 pending 任务数量 """ jobs = ( self.session.query(PrintJob) .filter( PrintJob.certificate_id == certificate_id, PrintJob.status == PrintJobStatus.PENDING, ) .all() ) for job in jobs: job.status = PrintJobStatus.CANCELLED job.error = "被新入队任务替换(replace_pending 去重)" job.finished_at = datetime.now() if jobs: DatabaseManager.commit_with_retry(self.session) logger.info( f"replace_pending: 取消同证旧 pending {len(jobs)} 个, cert_id={certificate_id}" ) return len(jobs)
[文档] def enqueue( self, certificate_id: int | None, sn: str | None = None, priority: int = 0, *, replace_pending: bool = False, ) -> PrintJob: """入队一个打印任务 Args: certificate_id: 关联合格证批次 ID(可为 None 表示纯展示任务) sn: 出厂编号(可选,用于单台任务或展示代表号) priority: 优先级,数值越大越优先出队 replace_pending: 为 True 时,入队前先取消该证已有的 pending 任务, 避免"重复点击入队"或"订正后再入队"产生同一证的多份 pending 导致重复打印。仅在 certificate_id 非空时生效。 Returns: PrintJob: 新创建的任务对象(已落库,含自增 id) """ if replace_pending and certificate_id is not None: self.cancel_pending_by_certificate(certificate_id) job = PrintJob( certificate_id=certificate_id, sn=sn, status=PrintJobStatus.PENDING, priority=priority, created_at=datetime.now(), ) self.session.add(job) self.session.flush() logger.info(f"打印任务入队: job_id={job.id}, cert_id={certificate_id}, sn={sn}") DatabaseManager.commit_with_retry(self.session) return job
[文档] def bulk_enqueue( self, items: Sequence[tuple[int | None, str | None]], priority: int = 0, *, replace_pending: bool = False, ) -> list[PrintJob]: """批量入队(保持传入顺序) Args: items: [(certificate_id, sn), ...] 列表 priority: 统一优先级 replace_pending: 为 True 时,每个 certificate_id 入队前先取消其已有 pending 任务(去重,见 :meth:`enqueue`)。 Returns: list[PrintJob]: 创建的任务列表 """ jobs: list[PrintJob] = [] for cert_id, sn in items: jobs.append( self.enqueue( certificate_id=cert_id, sn=sn, priority=priority, replace_pending=replace_pending, ) ) return jobs
# ============================================================ # 出队 # ============================================================
[文档] def get_next_pending(self) -> PrintJob | None: """取出下一个待打印任务(不出队,仅查询) 出队顺序:priority 降序、created_at 升序。 Returns: PrintJob | None: 下一个 pending 任务;无则返回 None """ return ( self.session.query(PrintJob) .filter(PrintJob.status == PrintJobStatus.PENDING) .order_by(PrintJob.priority.desc(), PrintJob.created_at.asc()) .first() )
[文档] def count_pending(self) -> int: """统计 pending 任务数量 Returns: int: pending 任务数 """ return ( self.session.query(PrintJob).filter(PrintJob.status == PrintJobStatus.PENDING).count() )
# ============================================================ # 状态机流转 # ============================================================
[文档] def get_job(self, job_id: int) -> PrintJob | None: """按 ID 获取任务 Args: job_id: 任务 ID Returns: PrintJob | None: 任务对象或 None """ return self.session.query(PrintJob).filter(PrintJob.id == job_id).first()
def _set_status( self, job_id: int, status: str, *, error: str | None = None, stamp_finished: bool = False, ) -> PrintJob | None: """内部:原子地更新任务状态并落库 Args: job_id: 任务 ID status: 目标状态 error: 错误信息(仅失败/取消时写入) stamp_finished: 是否写入 finished_at 时间戳 Returns: PrintJob | None: 更新后的任务;任务不存在返回 None """ job = self.get_job(job_id) if job is None: logger.warning(f"设置状态失败:任务不存在 job_id={job_id}") return None job.status = status if error is not None: job.error = error if stamp_finished: job.finished_at = datetime.now() DatabaseManager.commit_with_retry(self.session) return job
[文档] def mark_running(self, job_id: int, os_job_id: int | None = None) -> PrintJob | None: """标记任务为运行中(写入 started_at) Args: job_id: 任务 ID os_job_id: 本次打印提交的 OS spooler 作业 ID(偏差 5 作业 ID 捕获链路),非 None 时一并持久化,供 OS 级取消使用。 """ job = self.get_job(job_id) if job is None: return None job.status = PrintJobStatus.RUNNING job.started_at = datetime.now() if os_job_id is not None: job.os_job_id = os_job_id DatabaseManager.commit_with_retry(self.session) return job
[文档] def set_os_job_id(self, job_id: int, os_job_id: int | None) -> PrintJob | None: """补写任务的 OS spooler 作业 ID(打印动作完成后回填) 与 :meth:`mark_running` 解耦:部分场景下 OS 作业 ID 在 ``mark_running`` 之后、打印动作执行时才产生(见 BLUEPRINT §2.3 偏差 5)。 Args: job_id: 任务 ID os_job_id: OS 打印作业 ID;为 None 时忽略(不覆盖已有值)。 Returns: PrintJob | None: 更新后的任务;任务不存在返回 None """ if os_job_id is None: return self.get_job(job_id) job = self.get_job(job_id) if job is None: return None job.os_job_id = os_job_id DatabaseManager.commit_with_retry(self.session) return job
[文档] def mark_paused(self, job_id: int) -> PrintJob | None: """标记任务为已暂停(释放打印机,等待继续)""" return self._set_status(job_id, PrintJobStatus.PAUSED)
[文档] def mark_done(self, job_id: int) -> PrintJob | None: """标记任务为完成""" return self._set_status(job_id, PrintJobStatus.DONE, stamp_finished=True)
[文档] def mark_error(self, job_id: int, error: str) -> PrintJob | None: """标记任务为失败并记录错误""" return self._set_status(job_id, PrintJobStatus.DONE, error=error, stamp_finished=True)
[文档] def mark_cancelled(self, job_id: int, error: str = "用户取消") -> PrintJob | None: """标记任务为已取消""" return self._set_status(job_id, PrintJobStatus.CANCELLED, error=error, stamp_finished=True)
# ============================================================ # 优先级重排 # ============================================================ def _pending_ordered(self) -> list[PrintJob]: """返回当前 pending 任务列表(按出队顺序)""" return ( self.session.query(PrintJob) .filter(PrintJob.status == PrintJobStatus.PENDING) .order_by(PrintJob.priority.desc(), PrintJob.created_at.asc()) .all() ) @staticmethod def _reassign_priorities(ordered: list[PrintJob], session: Session) -> None: """按给定顺序重排 pending 任务优先级(严格递减)""" n = len(ordered) for idx, job in enumerate(ordered): job.priority = (n - idx) * 10 DatabaseManager.commit_with_retry(session)
[文档] def reorder_top(self, job_id: int) -> None: """将指定任务置顶(最高优先级出队)""" ordered = self._pending_ordered() target = next((j for j in ordered if j.id == job_id), None) if target is None: logger.warning(f"置顶失败:任务不在 pending 队列 job_id={job_id}") return ordered.remove(target) self._reassign_priorities([target, *ordered], self.session) logger.info(f"打印任务置顶: job_id={job_id}")
[文档] def reorder_up(self, job_id: int) -> None: """将指定任务上移一位(与上一个 pending 任务交换顺序)""" ordered = self._pending_ordered() ids = [j.id for j in ordered] if job_id not in ids: logger.warning(f"上移失败:任务不在 pending 队列 job_id={job_id}") return idx = ids.index(job_id) if idx == 0: return # 已在最前 ordered[idx], ordered[idx - 1] = ordered[idx - 1], ordered[idx] self._reassign_priorities(ordered, self.session) logger.info(f"打印任务上移: job_id={job_id}")
# ============================================================ # 查询与清理 # ============================================================
[文档] def list_jobs(self, include_terminal: bool = True) -> list[PrintJob]: """列出队列任务 Args: include_terminal: 是否包含已完成/取消的终态任务 Returns: list[PrintJob]: 任务列表(按 id 升序) """ query = self.session.query(PrintJob) if not include_terminal: query = query.filter(PrintJob.status == PrintJobStatus.PENDING) return query.order_by(PrintJob.id.asc()).all()
[文档] def clear_finished(self) -> int: """清理已完成的终态任务(done/cancelled) Returns: int: 删除的任务数量 """ deleted = ( self.session.query(PrintJob) .filter(PrintJob.status.in_(PrintJobStatus.TERMINAL)) .delete() ) DatabaseManager.commit_with_retry(self.session) logger.info(f"清理终态打印任务: {deleted} 条") return deleted
[文档] def recover(self) -> int: """断电/崩溃恢复:将 running/paused 任务重置回 pending 应用启动时调用,避免上次未完成任务卡在 running/paused 状态。 Returns: int: 被恢复(重置为 pending)的任务数量 """ recovered = ( self.session.query(PrintJob) .filter(PrintJob.status.in_([PrintJobStatus.RUNNING, PrintJobStatus.PAUSED])) .update( {PrintJob.status: PrintJobStatus.PENDING}, synchronize_session=False, ) ) DatabaseManager.commit_with_retry(self.session) logger.info(f"打印队列恢复: {recovered} 个任务重置为 pending") return recovered