成交走询价消息服务;陆运/海运私聊下一句能记下原因;协商中回执写明顺延三个工作日和下次 9:00。 Co-authored-by: Cursor <cursoragent@cursor.com>
323 lines
12 KiB
Python
323 lines
12 KiB
Python
"""
|
||
任务槽与认领循环:识别、LibreOffice、Archive、wake、LLM Stream。
|
||
|
||
本包约束(硬):LibreOffice 仅 1 槽;HTTP 类 2~3 槽;inbox/outbox 不得与 soffice 串同一阻塞线。
|
||
禁止:一把大锁串行所有工单;禁止 wecom:worker:singleton。
|
||
"""
|
||
|
||
from __future__ import annotations
|
||
|
||
import logging
|
||
import threading
|
||
import time
|
||
from dataclasses import dataclass, field
|
||
from typing import Callable, Optional
|
||
|
||
from agent.config import Settings, get_settings
|
||
from agent.graph.runtime import GraphRuntime
|
||
from agent.jobs.archive import archive_stream, process_archive_stub
|
||
from agent.jobs.libreoffice import LibreOfficeConverter
|
||
from agent.jobs.pdf_job import pdf_stream, process_pdf_job
|
||
from agent.jobs.quote_fill import fill_quote_workbook
|
||
from agent.jobs.recognition import process_recognition
|
||
from agent.jobs.wake_consumer import claim_and_apply_wake, wake_stream
|
||
from agent.llm import LlmGateway, LlmMode
|
||
from agent.redis_coord.runtime import RedisRuntime
|
||
from agent.routing import DeepSeekARouter
|
||
|
||
logger = logging.getLogger(__name__)
|
||
|
||
|
||
@dataclass
|
||
class SlotPool:
|
||
"""
|
||
进程内有限槽线程池骨架。
|
||
|
||
用途:限制同类慢活并发(尤其 LibreOffice=1)。
|
||
"""
|
||
|
||
name: str
|
||
size: int
|
||
_sem: threading.BoundedSemaphore = field(init=False, repr=False)
|
||
|
||
def __post_init__(self) -> None:
|
||
if self.size < 1:
|
||
raise ValueError(f"槽位数必须 >=1: {self.name}")
|
||
self._sem = threading.BoundedSemaphore(self.size)
|
||
|
||
def run(self, fn: Callable[[], None], *, label: str = "") -> None:
|
||
"""占用一个槽执行 fn;勿在 HTTP 回调线程调用。"""
|
||
acquired = self._sem.acquire(timeout=30)
|
||
if not acquired:
|
||
logger.warning("槽超时未获取 name=%s label=%s", self.name, label)
|
||
return
|
||
try:
|
||
fn()
|
||
finally:
|
||
self._sem.release()
|
||
|
||
|
||
@dataclass
|
||
class WorkerRuntime:
|
||
"""
|
||
Worker 运行时:分槽 + wake / LLM / 识别 / 存档 / LibreOffice 循环。
|
||
|
||
启动时机:worker_app main。
|
||
"""
|
||
|
||
libreoffice: SlotPool
|
||
http_jobs: SlotPool
|
||
settings: Settings
|
||
graph: Optional[GraphRuntime] = None
|
||
redis: Optional[RedisRuntime] = None
|
||
llm: Optional[LlmGateway] = None
|
||
_stop: threading.Event = field(default_factory=threading.Event)
|
||
_threads: list[threading.Thread] = field(default_factory=list)
|
||
|
||
@classmethod
|
||
def from_settings(
|
||
cls,
|
||
libreoffice_slots: int,
|
||
http_job_slots: int,
|
||
*,
|
||
graph: Optional[GraphRuntime] = None,
|
||
redis: Optional[RedisRuntime] = None,
|
||
llm: Optional[LlmGateway] = None,
|
||
settings: Optional[Settings] = None,
|
||
) -> "WorkerRuntime":
|
||
"""按配置创建槽;强制 LibreOffice 槽为 1。"""
|
||
cfg = settings or get_settings()
|
||
lo = 1 if libreoffice_slots != 1 else libreoffice_slots
|
||
if libreoffice_slots != 1:
|
||
logger.warning(
|
||
"LibreOffice 槽必须为 1(收到 %s),已强制为 1,避免本机双 soffice 互抢",
|
||
libreoffice_slots,
|
||
)
|
||
return cls(
|
||
libreoffice=SlotPool("libreoffice", lo),
|
||
http_jobs=SlotPool("http_jobs", max(2, http_job_slots)),
|
||
settings=cfg,
|
||
graph=graph,
|
||
redis=redis,
|
||
llm=llm,
|
||
)
|
||
|
||
def start(self) -> None:
|
||
"""启动分槽消费循环。"""
|
||
self._stop.clear()
|
||
self._threads = [
|
||
threading.Thread(target=self._wake_loop, name="wake-loop", daemon=True),
|
||
threading.Thread(target=self._llm_loop, name="llm-loop", daemon=True),
|
||
threading.Thread(target=self._recognition_loop, name="recog-loop", daemon=True),
|
||
threading.Thread(target=self._archive_loop, name="archive-loop", daemon=True),
|
||
threading.Thread(target=self._wshoto_poll_loop, name="wshoto-poll", daemon=True),
|
||
threading.Thread(target=self._deal_followup_loop, name="deal-followup", daemon=True),
|
||
threading.Thread(target=self._pdf_loop, name="pdf-loop", daemon=True),
|
||
]
|
||
for t in self._threads:
|
||
t.start()
|
||
lo = LibreOfficeConverter.from_settings(self.settings)
|
||
logger.info(
|
||
"WorkerRuntime 已启动 libreoffice_slots=%s http_job_slots=%s soffice=%s",
|
||
self.libreoffice.size,
|
||
self.http_jobs.size,
|
||
"ok" if lo.available() else "missing",
|
||
)
|
||
|
||
def stop(self) -> None:
|
||
self._stop.set()
|
||
|
||
def _wake_loop(self) -> None:
|
||
"""HTTP 类槽消费结构化 wake → Graph。"""
|
||
if not self.redis or not self.graph:
|
||
while not self._stop.is_set():
|
||
time.sleep(2.0)
|
||
return
|
||
q = wake_stream(self.redis.client)
|
||
consumer = f"wake-{threading.get_ident()}"
|
||
while not self._stop.is_set():
|
||
|
||
def _job() -> None:
|
||
claim_and_apply_wake(queue=q, graph=self.graph, consumer=consumer)
|
||
|
||
self.http_jobs.run(_job, label="wake")
|
||
time.sleep(0.2)
|
||
|
||
def _llm_loop(self) -> None:
|
||
"""
|
||
HTTP 类槽消费 LLM Stream。
|
||
|
||
ordered_actions:走 DeepSeekARouter(默认 stub;LLM_ALLOW_NETWORK 才可能真网)。
|
||
其它 mode:骨架 stub。
|
||
"""
|
||
if not self.redis or not self.llm:
|
||
while not self._stop.is_set():
|
||
time.sleep(2.0)
|
||
return
|
||
stream = self.redis.llm_stream
|
||
stream.ensure_group()
|
||
consumer = f"llm-{threading.get_ident()}"
|
||
router = DeepSeekARouter.from_settings(
|
||
self.settings,
|
||
allow_network=bool(self.settings.llm_allow_network),
|
||
)
|
||
while not self._stop.is_set():
|
||
|
||
def _job() -> None:
|
||
claimed = stream.claim(consumer=consumer, count=1, block_ms=200)
|
||
if not claimed:
|
||
return
|
||
entry_id, data = claimed[0]
|
||
mode = str(data.get("mode") or "ordered_actions")
|
||
payload = dict(data.get("payload") or {})
|
||
try:
|
||
if mode == LlmMode.ORDERED_ACTIONS.value:
|
||
decision = router.route(
|
||
text=str(payload.get("text") or ""),
|
||
sender_id=str(payload.get("sender_id") or ""),
|
||
)
|
||
logger.info(
|
||
"llm 槽 A 路由 id=%s intent=%s",
|
||
entry_id,
|
||
decision.intent,
|
||
)
|
||
else:
|
||
result = self.llm.invoke_stub(mode=mode, payload=payload)
|
||
logger.info(
|
||
"llm 槽处理 id=%s mode=%s stub=%s",
|
||
entry_id,
|
||
mode,
|
||
result.get("stub"),
|
||
)
|
||
finally:
|
||
stream.ack(entry_id)
|
||
|
||
self.http_jobs.run(_job, label="llm")
|
||
time.sleep(0.2)
|
||
|
||
def _recognition_loop(self) -> None:
|
||
"""HTTP 类槽消费识别 Stream。"""
|
||
if not self.redis:
|
||
while not self._stop.is_set():
|
||
time.sleep(2.0)
|
||
return
|
||
stream = self.redis.recognition_stream
|
||
stream.ensure_group()
|
||
consumer = f"recog-{threading.get_ident()}"
|
||
while not self._stop.is_set():
|
||
|
||
def _job() -> None:
|
||
claimed = stream.claim(consumer=consumer, count=1, block_ms=200)
|
||
if not claimed:
|
||
return
|
||
entry_id, data = claimed[0]
|
||
try:
|
||
process_recognition(
|
||
data,
|
||
allow_network=bool(self.settings.llm_allow_network),
|
||
)
|
||
finally:
|
||
stream.ack(entry_id)
|
||
|
||
self.http_jobs.run(_job, label="recognition")
|
||
time.sleep(0.2)
|
||
|
||
def _archive_loop(self) -> None:
|
||
"""HTTP 类槽消费存档 Stream(与 soffice 隔离)。"""
|
||
if not self.redis:
|
||
while not self._stop.is_set():
|
||
time.sleep(2.0)
|
||
return
|
||
stream = archive_stream(self.redis)
|
||
stream.ensure_group()
|
||
consumer = f"arch-{threading.get_ident()}"
|
||
while not self._stop.is_set():
|
||
|
||
def _job() -> None:
|
||
claimed = stream.claim(consumer=consumer, count=1, block_ms=200)
|
||
if not claimed:
|
||
return
|
||
entry_id, data = claimed[0]
|
||
try:
|
||
process_archive_stub(data)
|
||
finally:
|
||
stream.ack(entry_id)
|
||
|
||
self.http_jobs.run(_job, label="archive")
|
||
time.sleep(0.2)
|
||
|
||
def _wshoto_poll_loop(self) -> None:
|
||
"""
|
||
独立线程轮询微盛群消息,写入 inbox。
|
||
|
||
不占回调、不占 HTTP 槽、不与 soffice 串队。建群或刚有人说话会立刻再扫。
|
||
未配钥匙时只空等,避免打无密钥请求。
|
||
"""
|
||
from agent.channel.wshoto.client import wshoto_ready
|
||
from agent.channel.wshoto.poll import poll_once
|
||
from agent.channel.wshoto.rooms import has_hot_room, wait_wshoto_poll
|
||
|
||
interval = float(getattr(self.settings, "wshoto_poll_interval_sec", 2.0) or 2.0)
|
||
interval = min(30.0, max(1.0, interval))
|
||
logged_skip = False
|
||
while not self._stop.is_set():
|
||
if not wshoto_ready(self.settings):
|
||
if not logged_skip:
|
||
logger.info("微盛轮询未开启:缺少 WSHOTO_APP_ID/SECRET 或开关已关")
|
||
logged_skip = True
|
||
self._stop.wait(15.0)
|
||
continue
|
||
try:
|
||
poll_once(settings=self.settings, redis=self.redis.client if self.redis else None)
|
||
except Exception:
|
||
logger.exception("微盛轮询失败")
|
||
wait = 1.0 if has_hot_room() else interval
|
||
wait_wshoto_poll(wait)
|
||
|
||
def _deal_followup_loop(self) -> None:
|
||
"""
|
||
HTTP 槽轮询协商再问。独立线程,不占回调、不与 soffice 串队。
|
||
"""
|
||
from agent.jobs.deal_followup_poll import poll_due_negotiate
|
||
|
||
while not self._stop.is_set():
|
||
try:
|
||
self.http_jobs.run(
|
||
lambda: poll_due_negotiate(settings=self.settings),
|
||
label="deal-followup",
|
||
)
|
||
except Exception:
|
||
logger.exception("协商再问轮询失败")
|
||
self._stop.wait(60.0)
|
||
|
||
def _pdf_loop(self) -> None:
|
||
"""仅 LibreOffice 单槽消费 PDF Stream。"""
|
||
if not self.redis:
|
||
while not self._stop.is_set():
|
||
time.sleep(2.0)
|
||
return
|
||
stream = pdf_stream(self.redis)
|
||
stream.ensure_group()
|
||
converter = LibreOfficeConverter.from_settings(self.settings)
|
||
consumer = f"pdf-{threading.get_ident()}"
|
||
while not self._stop.is_set():
|
||
|
||
def _job() -> None:
|
||
claimed = stream.claim(consumer=consumer, count=1, block_ms=200)
|
||
if not claimed:
|
||
return
|
||
entry_id, data = claimed[0]
|
||
try:
|
||
result = process_pdf_job(data, converter=converter)
|
||
logger.info(
|
||
"pdf 槽 id=%s ok=%s err=%s",
|
||
entry_id,
|
||
result.ok,
|
||
result.error or "-",
|
||
)
|
||
finally:
|
||
stream.ack(entry_id)
|
||
|
||
self.libreoffice.run(_job, label="pdf")
|
||
time.sleep(0.3)
|