169 lines
5.8 KiB
Python
169 lines
5.8 KiB
Python
"""
|
|
报价 PDF 任务壳:入队后由 LibreOffice **单槽**消费。
|
|
|
|
本文件职责:enqueue_pdf_job / process 调用 LibreOfficeConverter。
|
|
禁止:在 http_jobs 槽跑 soffice;禁止回调线程转换。
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
import logging
|
|
from dataclasses import dataclass
|
|
from typing import Any, Optional
|
|
|
|
from agent.jobs.libreoffice import LibreOfficeConverter, LibreOfficeResult
|
|
from agent.redis_coord import get_redis_runtime, try_acquire_idempotency
|
|
from agent.redis_coord.keys import STREAM_PDF
|
|
from agent.redis_coord.runtime import RedisRuntime
|
|
from agent.redis_coord.stream_queue import StreamQueue
|
|
|
|
logger = logging.getLogger(__name__)
|
|
|
|
|
|
@dataclass
|
|
class PdfEnqueueResult:
|
|
accepted: bool
|
|
entry_id: str = ""
|
|
reason: str = ""
|
|
|
|
|
|
def pdf_stream(redis: Optional[RedisRuntime] = None) -> StreamQueue:
|
|
rt = redis or get_redis_runtime()
|
|
return StreamQueue(client=rt.client, stream_suffix=STREAM_PDF)
|
|
|
|
|
|
def enqueue_pdf_job(
|
|
*,
|
|
xlsx_path: str,
|
|
inquiry_no: str = "",
|
|
quote_version: int = 0,
|
|
idempotency_key: str = "",
|
|
redis: Optional[RedisRuntime] = None,
|
|
touser: str = "",
|
|
filename: str = "",
|
|
quote: Optional[dict[str, Any]] = None,
|
|
facts: Optional[dict[str, Any]] = None,
|
|
) -> PdfEnqueueResult:
|
|
"""Excel→PDF 任务入队;真正转换只在 LibreOffice 槽。"""
|
|
rt = redis or get_redis_runtime()
|
|
q = pdf_stream(rt)
|
|
q.ensure_group()
|
|
if idempotency_key:
|
|
ok = try_acquire_idempotency(
|
|
rt.client,
|
|
scope="pdf",
|
|
idempotency_key=idempotency_key,
|
|
ttl_seconds=600,
|
|
)
|
|
if not ok:
|
|
return PdfEnqueueResult(accepted=False, reason="duplicate_idempotency_key")
|
|
payload = {
|
|
"xlsx_path": xlsx_path,
|
|
"inquiry_no": inquiry_no,
|
|
"quote_version": quote_version,
|
|
"touser": touser,
|
|
"filename": filename,
|
|
"quote": dict(quote or {}),
|
|
"facts": dict(facts or {}),
|
|
}
|
|
entry_id = q.enqueue(payload=payload, idempotency_key=idempotency_key)
|
|
logger.info("pdf.enqueue id=%s inquiry=%s", entry_id, inquiry_no or "-")
|
|
return PdfEnqueueResult(accepted=True, entry_id=entry_id)
|
|
|
|
|
|
def process_pdf_job(
|
|
data: dict[str, Any],
|
|
*,
|
|
converter: Optional[LibreOfficeConverter] = None,
|
|
) -> LibreOfficeResult:
|
|
"""
|
|
处理一条 PDF 任务(调用方必须已占用 libreoffice 槽)。
|
|
|
|
若 payload 含 template_xlsx + mapping_json + facts,先 quote_fill 再转 PDF。
|
|
缺路径或本机无 soffice 时返回 ok=False,不抛崩 Worker。
|
|
"""
|
|
conv = converter or LibreOfficeConverter.from_settings()
|
|
template_xlsx = str(data.get("template_xlsx") or "")
|
|
mapping_json = data.get("mapping_json")
|
|
facts = data.get("facts")
|
|
work_dir = str(data.get("work_dir") or "")
|
|
if template_xlsx and mapping_json is not None and isinstance(facts, dict):
|
|
from pathlib import Path
|
|
import tempfile
|
|
|
|
try:
|
|
from agent.jobs.quote_fill import fill_then_convert_pdf
|
|
except ImportError:
|
|
return LibreOfficeResult(ok=False, error="quote_fill_missing")
|
|
|
|
wd = Path(work_dir) if work_dir else Path(tempfile.mkdtemp(prefix="quote_fill_"))
|
|
fill, pdf = fill_then_convert_pdf(
|
|
source_xlsx=template_xlsx,
|
|
work_dir=wd,
|
|
mapping=mapping_json if not isinstance(mapping_json, str) else mapping_json,
|
|
facts=facts,
|
|
converter=conv,
|
|
)
|
|
if not fill.ok:
|
|
return LibreOfficeResult(ok=False, error=f"fill_failed:{fill.error}")
|
|
assert pdf is not None
|
|
_deliver_quote_file(data, pdf)
|
|
return pdf
|
|
|
|
path = str(data.get("xlsx_path") or "")
|
|
if not path:
|
|
return LibreOfficeResult(ok=False, error="missing_xlsx_path")
|
|
result = conv.convert_xlsx_to_pdf(path)
|
|
_deliver_quote_file(data, result)
|
|
return result
|
|
|
|
|
|
def _deliver_quote_file(data: dict[str, Any], result: LibreOfficeResult) -> None:
|
|
"""
|
|
转换成功后把 PDF 写入 outbox,由企微通道发文件。
|
|
|
|
无 touser 则只落盘(旧任务);失败不抛,避免卡 LibreOffice 槽。
|
|
"""
|
|
if not result.ok:
|
|
return
|
|
touser = str(data.get("touser") or "").strip()
|
|
if not touser:
|
|
return
|
|
wo = str(data.get("inquiry_no") or "")
|
|
filename = str(data.get("filename") or "").strip() or f"{wo}_quote.pdf"
|
|
try:
|
|
from agent.channel.queue import get_message_store
|
|
from agent.policy.inquiry_copy import file_ready
|
|
|
|
store = get_message_store()
|
|
store.enqueue_outbound(
|
|
touser=touser,
|
|
content=file_ready(kind="PDF", work_order_no=wo, filename=filename),
|
|
dedupe_key=f"pdf-file:{wo}:{filename}",
|
|
payload={
|
|
"msgtype": "file",
|
|
"filepath": result.pdf_path,
|
|
"filename": filename,
|
|
},
|
|
)
|
|
# 文件先入队,成交卡后入队;出站单线程保证销售先看到文件
|
|
quote = data.get("quote") if isinstance(data.get("quote"), dict) else {}
|
|
facts = data.get("facts") if isinstance(data.get("facts"), dict) else {}
|
|
if wo:
|
|
from agent.policy import inquiry_copy as copy
|
|
|
|
store.enqueue_outbound(
|
|
touser=touser,
|
|
content=copy.deal_card(work_order_no=wo, quote=quote, facts=facts),
|
|
dedupe_key=f"deal-card:{wo}:{filename}",
|
|
payload=copy.deal_wecom_payload(work_order_no=wo, quote=quote, facts=facts),
|
|
)
|
|
try:
|
|
from agent.policy.air_text_flow import get_air_text_flow
|
|
|
|
get_air_text_flow().mark_wait_deal(touser, wo)
|
|
except Exception:
|
|
logger.exception("PDF 出文件后未能切到成交跟进 inquiry=%s", wo)
|
|
except Exception:
|
|
logger.exception("PDF 出站入队失败 inquiry=%s", wo)
|