JhHardwareWRS_BackPoint/app/services/qrcode_batch_tasks.py
2026-07-10 10:12:03 +08:00

343 lines
12 KiB
Python

import asyncio
import hashlib
import logging
import re
import uuid
import zipfile
from datetime import datetime
from pathlib import Path
from urllib.parse import quote, unquote, urlparse
from sqlalchemy import select
from app.config import settings
from app.database import SessionLocal
from app.models import DeviceQrBatchTask, DeviceQRCode, Product
from app.schemas import DeviceQrBatchTaskOut
from app.services.cleaning import is_cleaning_product
from app.services.display_names import point_mold_display_name
from app.services.misc_work import is_misc_product
from app.services.wechat import create_device_qrcode
from app.timezone import now
QRCODE_TASK_STATUS_PENDING = "pending"
QRCODE_TASK_STATUS_RUNNING = "running"
QRCODE_TASK_STATUS_PAUSED = "paused"
QRCODE_TASK_STATUS_COMPLETED = "completed"
QRCODE_TASK_STATUS_FAILED = "failed"
QRCODE_BATCH_TASK_INTERVAL_SECONDS = 3
QRCODE_TASK_STATUS_NAMES = {
QRCODE_TASK_STATUS_PENDING: "等待生成",
QRCODE_TASK_STATUS_RUNNING: "生成中",
QRCODE_TASK_STATUS_PAUSED: "已停止",
QRCODE_TASK_STATUS_COMPLETED: "已完成",
QRCODE_TASK_STATUS_FAILED: "生成失败",
}
logger = logging.getLogger(__name__)
def _safe_file_stem(value: str) -> str:
normalized = re.sub(r'[\\/:*?"<>|\s]+', "_", value).strip("._")
return normalized[:120] or "模具二维码"
def _unique_zip_name(used_names: set[str], stem: str, suffix: str) -> str:
suffix = suffix if suffix else ".jpg"
name = f"{stem}{suffix}"
index = 2
while name in used_names:
name = f"{stem}_{index}{suffix}"
index += 1
used_names.add(name)
return name
def _qrcode_path_from_url(qr_url: str | None) -> Path | None:
if not qr_url:
return None
filename = unquote(Path(urlparse(qr_url).path).name)
if not filename:
return None
path = settings.upload_path / "qrcodes" / filename
return path if path.is_file() else None
def _scene_token(point_name: str, mold_name: str, process_name: str) -> str:
mold_key = f"{point_name}\0{mold_name}\0{process_name}"
return hashlib.sha1(mold_key.encode("utf-8")).hexdigest()[:16]
def _product_for_mold(db, point_name: str, mold_name: str, process_name: str) -> Product | None:
return db.scalar(
select(Product)
.where(
Product.attendance_point_name == point_name,
Product.product_name == mold_name,
Product.process_name == process_name,
Product.device_no == "",
)
.order_by(Product.project_no.asc())
)
def _public_base_url() -> str:
return settings.public_base_url.rstrip("/") or "http://127.0.0.1:8000"
def qrcode_batch_zip_path(task: DeviceQrBatchTask) -> Path:
return settings.upload_path / "qrcode_zips" / task.file_name
def qrcode_batch_zip_exists(task: DeviceQrBatchTask) -> bool:
return bool(task.file_name and qrcode_batch_zip_path(task).is_file())
def _task_download_url(task: DeviceQrBatchTask, public_base_url: str) -> str:
return (
f"{public_base_url.rstrip('/')}/api/devices/molds/qrcode/batch/tasks/"
f"{task.id}/download"
)
def _task_file_name(items: list[dict]) -> str:
fingerprint = hashlib.sha1(
"|".join(
f"{item.get('attendance_point_name', '')}\0{item.get('product_name', '')}\0{item.get('process_name', '')}"
for item in items
).encode("utf-8")
).hexdigest()[:10]
unique = uuid.uuid4().hex[:8]
return f"mold-qrcodes-{datetime.now().strftime('%Y%m%d%H%M%S')}-{fingerprint}-{unique}.zip"
def _items_signature(items: list[dict]) -> tuple[tuple[str, str, str], ...]:
return tuple(
sorted(
(
str(item.get("attendance_point_name") or "").strip(),
str(item.get("product_name") or "").strip(),
str(item.get("process_name") or "").strip(),
)
for item in items
)
)
def create_qrcode_batch_task(db, *, created_by: str, items: list[dict]) -> DeviceQrBatchTask:
signature = _items_signature(items)
active_tasks = db.scalars(
select(DeviceQrBatchTask)
.where(
DeviceQrBatchTask.created_by == created_by,
DeviceQrBatchTask.status.in_([
QRCODE_TASK_STATUS_PENDING,
QRCODE_TASK_STATUS_RUNNING,
QRCODE_TASK_STATUS_PAUSED,
]),
)
.order_by(DeviceQrBatchTask.created_at.desc(), DeviceQrBatchTask.id.desc())
).all()
for active_task in active_tasks:
if _items_signature(list(active_task.items_json or [])) == signature:
return active_task
task = DeviceQrBatchTask(
file_name=_task_file_name(items),
status=QRCODE_TASK_STATUS_PENDING,
item_count=len(items),
completed_count=0,
failed_count=0,
items_json=items,
created_by=created_by,
)
db.add(task)
db.commit()
db.refresh(task)
return task
def qrcode_batch_task_out(task: DeviceQrBatchTask, public_base_url: str | None = None) -> DeviceQrBatchTaskOut:
item_count = max(0, int(task.item_count or 0))
completed_count = max(0, int(task.completed_count or 0))
if task.status == QRCODE_TASK_STATUS_COMPLETED:
progress_percent = 100.0
elif item_count <= 0:
progress_percent = 0.0
else:
progress_percent = round(min(100, completed_count / item_count * 100), 1)
zip_url = None
if task.status == QRCODE_TASK_STATUS_COMPLETED and qrcode_batch_zip_exists(task):
zip_url = _task_download_url(task, public_base_url) if public_base_url else task.zip_url
return DeviceQrBatchTaskOut(
id=task.id,
file_name=task.file_name,
status=task.status,
status_name=QRCODE_TASK_STATUS_NAMES.get(task.status, task.status),
item_count=item_count,
completed_count=completed_count,
failed_count=max(0, int(task.failed_count or 0)),
progress_percent=progress_percent,
zip_url=zip_url,
error_message=task.error_message,
started_at=task.started_at,
finished_at=task.finished_at,
created_at=task.created_at,
updated_at=task.updated_at,
)
def recover_running_qrcode_batch_tasks() -> None:
with SessionLocal() as db:
tasks = db.scalars(
select(DeviceQrBatchTask).where(DeviceQrBatchTask.status == QRCODE_TASK_STATUS_RUNNING)
).all()
for task in tasks:
task.status = QRCODE_TASK_STATUS_PENDING
task.failed_count = 0
task.error_message = "服务重启后重新排队"
task.finished_at = None
if tasks:
db.commit()
def _claim_next_qrcode_batch_task() -> int | None:
with SessionLocal() as db:
query = (
select(DeviceQrBatchTask)
.where(DeviceQrBatchTask.status == QRCODE_TASK_STATUS_PENDING)
.order_by(DeviceQrBatchTask.created_at.asc(), DeviceQrBatchTask.id.asc())
.limit(1)
)
if db.get_bind().dialect.name in {"mysql", "postgresql"}:
query = query.with_for_update(skip_locked=True)
task = db.scalar(query)
if task is None:
return None
task.status = QRCODE_TASK_STATUS_RUNNING
if task.started_at is None:
task.started_at = now()
task.finished_at = None
task.failed_count = 0
task.error_message = None
db.commit()
return int(task.id)
async def _generate_task_qrcode(db, task: DeviceQrBatchTask, item: dict) -> tuple[str, Path]:
point_name = str(item.get("attendance_point_name") or "").strip()
mold_name = str(item.get("product_name") or "").strip()
process_name = str(item.get("process_name") or "").strip()
product = _product_for_mold(db, point_name, mold_name, process_name)
if product is None:
raise RuntimeError(f"产品清单中已不存在:{point_name} / {mold_name} / {process_name}")
display_name = point_mold_display_name(point_name, mold_name, process_name, product.stamping_method)
scene = f"mold={_scene_token(point_name, mold_name, process_name)}"
qr_url = await create_device_qrcode(
display_name,
page="pages/clock/clock",
scene=scene,
public_base_url=_public_base_url(),
label=f"模具 {display_name}",
)
record = db.get(
DeviceQRCode,
{"attendance_point_name": point_name, "device_no": mold_name, "process_name": process_name},
)
if record is None:
record = DeviceQRCode(
attendance_point_name=point_name,
device_no=mold_name,
process_name=process_name,
created_by=task.created_by,
)
db.add(record)
record.qr_scene = scene
record.qr_url = qr_url
db.commit()
qr_path = _qrcode_path_from_url(qr_url)
if qr_path is None:
raise RuntimeError(f"{display_name} 二维码图片未生成,无法打包")
return display_name, qr_path
async def process_qrcode_batch_task(task_id: int) -> None:
zip_dir = settings.upload_path / "qrcode_zips"
zip_dir.mkdir(parents=True, exist_ok=True)
used_names: set[str] = set()
with SessionLocal() as db:
task = db.get(DeviceQrBatchTask, task_id)
if task is None:
return
items = list(task.items_json or [])
zip_path = zip_dir / task.file_name
completed_count = max(0, min(int(task.completed_count or 0), len(items)))
if completed_count and not zip_path.exists():
completed_count = 0
task.completed_count = 0
db.commit()
zip_mode = "a" if completed_count > 0 else "w"
try:
with zipfile.ZipFile(zip_path, zip_mode, zipfile.ZIP_DEFLATED) as archive:
for index, item in enumerate(items, start=1):
if index <= completed_count:
continue
db.refresh(task)
if task.status == QRCODE_TASK_STATUS_PAUSED:
task.finished_at = None
db.commit()
return
if task.status != QRCODE_TASK_STATUS_RUNNING:
return
display_name, qr_path = await _generate_task_qrcode(db, task, item)
stem = _safe_file_stem(f"{index:03d}_{display_name}")
archive.write(qr_path, _unique_zip_name(set(archive.namelist()), stem, qr_path.suffix))
task.completed_count = index
task.failed_count = 0
task.updated_at = now()
db.commit()
task.status = QRCODE_TASK_STATUS_COMPLETED
task.completed_count = len(items)
task.failed_count = 0
task.zip_url = f"{_public_base_url()}/uploads/qrcode_zips/{quote(task.file_name, safe='')}"
task.error_message = None
task.finished_at = now()
db.commit()
except Exception as exc:
db.rollback()
task = db.get(DeviceQrBatchTask, task_id)
if task is None:
zip_path.unlink(missing_ok=True)
else:
completed_count = int(task.completed_count or 0)
task.status = QRCODE_TASK_STATUS_FAILED
task.failed_count = max(1, int(task.item_count or 0) - completed_count)
task.error_message = str(exc)[:1000]
task.finished_at = now()
db.commit()
raise
async def qrcode_batch_task_loop(interval_seconds: int = QRCODE_BATCH_TASK_INTERVAL_SECONDS) -> None:
await asyncio.to_thread(recover_running_qrcode_batch_tasks)
while True:
try:
task_id = await asyncio.to_thread(_claim_next_qrcode_batch_task)
if task_id is None:
await asyncio.sleep(interval_seconds)
continue
await process_qrcode_batch_task(task_id)
except asyncio.CancelledError:
raise
except Exception:
logger.exception("qrcode batch task failed")
await asyncio.sleep(interval_seconds)