ForcePilot/backend/package/yuxi/channels/adapters/arq_queue_adapter.py
Kris b88c0ae29e feat(channels): 批量新增多渠道网关限界上下文基础代码与契约
新增完整的 channels 限界上下文模块,包含契约层、领域核心层、应用服务、管道编排、插件体系、基础设施组合根等全层级代码,新增飞书与微信 iLink 渠道插件基础结构,补充各类 DTO、端口协议与领域服务实现。
2026-07-02 03:22:12 +08:00

209 lines
7.8 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

"""ARQQueueAdapter实现 QueuePort,复用现有 arq.enqueue_job。
- enqueue 必须复用 arq.enqueue_job
- 不得另起重试 Worker
- ARQ 不可用时 Outbox 必须暂停重试并触发优雅降级
依赖边界:只依赖 yuxi.channels.contract端口 + DTO + 错误)。
ARQ 连接池与 ``LoggerPort`` 通过构造函数注入,由外部 factory 管理生命周期,
不依赖 yuxi.servicesADP-005
"""
from __future__ import annotations
import asyncio
from typing import TYPE_CHECKING, Any
from yuxi.channels.contract.dtos.queue import EnqueueCmd, JobId, JobStatus
from yuxi.channels.contract.errors import DependencyError, NotFoundError, ValidationError
from yuxi.channels.contract.errors.base import Error
from yuxi.channels.contract.ports.driven.logger_port import LoggerPort
from yuxi.channels.contract.ports.driven.queue_port import QueuePort
if TYPE_CHECKING:
from arq import ArqRedis
class ARQQueueAdapter(QueuePort):
"""实现 QueuePort,复用现有 arq.enqueue_job。
通过构造函数注入的 ``ArqRedis`` 连接池复用 ``arq.enqueue_job`` 入队,
通过 ``arq.jobs.Job`` 查询状态与取消任务。ARQ 不可用时抛出
``DependencyError``,由调用方暂停 Outbox 重试并触发优雅降级。
ARQ 连接池与 ``LoggerPort`` 由外部 factory 创建并注入,适配器仅持有
引用,不负责其生命周期管理(参见 ``close``)。
"""
def __init__(self, arq_pool: ArqRedis, logger: LoggerPort) -> None:
"""初始化适配器,注入 ARQ 连接池与日志端口。
Args:
arq_pool: ARQ 连接池实例(``arq.ArqRedis``,由外部 factory
创建并管理生命周期。适配器仅持有引用用于入队、状态查询
与取消任务。
logger: 日志被驱动端口,用于 ``ping`` 故障时记录告警
(硬约束:所有适配器必须使用注入的 ``LoggerPort`` 而非
全局 logger 工具)。
"""
self._arq_pool = arq_pool
self._logger = logger
async def enqueue(self, cmd: EnqueueCmd) -> JobId:
"""入队异步任务,复用 arq.enqueue_job。
校验 task_name 与 payload 非空后,通过 ``arq.enqueue_job`` 入队。
幂等键重复入队时返回首次入队的 JobIdARQ 不可用时抛出
``DependencyError``。
Args:
cmd: 入队命令,携带任务名称、负载与幂等键。
Returns:
任务 ID。
Raises:
ValidationError: task_name 或 payload 为空。
DependencyError: ARQ 不可用或 enqueue_job 返回 None 且无幂等键。
"""
try:
if not cmd.task_name:
raise ValidationError("task_name", "must not be empty")
if not cmd.payload:
raise ValidationError("payload", "must not be empty")
queue = self._arq_pool
# 消费 scheduled_at 字段,非空时通过 _defer_until 交给 ARQ 延迟投递
enqueue_kwargs: dict[str, Any] = {"_job_id": cmd.idempotency_key}
if cmd.scheduled_at is not None:
enqueue_kwargs["_defer_until"] = cmd.scheduled_at
job = await queue.enqueue_job(cmd.task_name, **cmd.payload, **enqueue_kwargs)
if job is None:
if cmd.idempotency_key:
return JobId(cmd.idempotency_key)
raise DependencyError("arq", Error("enqueue_job returned None"))
return JobId(job.job_id)
except (ValidationError, DependencyError):
raise
except Exception as exc:
raise DependencyError("arq", exc) from exc
async def getJobStatus(self, job_id: JobId) -> JobStatus:
"""查询任务状态。
通过 ``arq.jobs.Job.status()`` 获取 ARQ 状态并映射为契约层
``JobStatus`` 枚举。任务不存在时抛出 ``NotFoundError``ARQ 不可用
时抛出 ``DependencyError``。
状态映射:
- arq ``deferred`` / ``queued`` → ``JobStatus.QUEUED``
- arq ``in_progress`` → ``JobStatus.RUNNING``
- arq ``complete`` → 根据 ``result_info.success`` 区分
``COMPLETED`` / ``FAILED`` / ``CANCELLED``
Args:
job_id: 任务 ID。
Returns:
任务状态枚举值。
Raises:
NotFoundError: 任务不存在。
DependencyError: ARQ 不可用。
"""
try:
from arq.jobs import Job as ArqJob
from arq.jobs import JobStatus as ArqJobStatus
queue = self._arq_pool
job = ArqJob(
job_id.value,
redis=queue,
_queue_name=queue.default_queue_name,
_deserializer=queue.job_deserializer,
)
status = await job.status()
if status == ArqJobStatus.not_found:
raise NotFoundError("job", job_id.value)
if status == ArqJobStatus.complete:
result_info = await job.result_info()
if result_info is not None and not result_info.success:
if isinstance(result_info.result, asyncio.CancelledError):
return JobStatus.CANCELLED
return JobStatus.FAILED
return JobStatus.COMPLETED
if status == ArqJobStatus.in_progress:
return JobStatus.RUNNING
return JobStatus.QUEUED
except NotFoundError:
raise
except Exception as exc:
raise DependencyError("arq", Error(str(exc))) from exc
async def cancelJob(self, job_id: JobId) -> bool:
"""取消任务,用于 Outbox 重试取消场景。
通过 ``arq.jobs.Job.abort()`` 取消任务。任务已取消返回 ``True``
任务已完成或不存在返回 ``False``ARQ 不可用时抛出
``DependencyError``。
Args:
job_id: 任务 ID。
Returns:
任务已取消返回 ``True``,已完成或不存在返回 ``False``。
Raises:
DependencyError: ARQ 不可用。
"""
try:
from arq.jobs import Job as ArqJob
queue = self._arq_pool
job = ArqJob(
job_id.value,
redis=queue,
_queue_name=queue.default_queue_name,
_deserializer=queue.job_deserializer,
)
return await job.abort()
except Exception as exc:
raise DependencyError("arq", Error(str(exc))) from exc
async def getWorkerStatus(self) -> dict[str, Any]:
"""获取 ARQ Worker 状态。
返回 ARQ 队列的关键指标,供 ``DiagnosticsExporter`` 填充诊断包。
Raises:
DependencyError: ARQ 不可用。
"""
try:
queue = self._arq_pool
return {
"queue_name": queue.default_queue_name,
"available": True,
}
except Exception as exc:
raise DependencyError("arq", Error(str(exc))) from exc
async def ping(self) -> bool:
"""主动探测 ARQ 连接可用性,故障时返回 False降级不阻断
执行 ARQ 连接池 ``ping`` 验证连接可用性。故障时返回 False 并通过
注入的 ``LoggerPort`` 记录 warning 日志,不抛异常,供 ``HostBootstrap``
启动期连通性检查使用。
Returns:
True 表示连接可用False 表示故障。
"""
try:
await self._arq_pool.ping()
return True
except Exception as exc:
await self._logger.warn(
"arq queue ping failed",
error_type=type(exc).__name__,
error=str(exc),
)
return False