本次提交新增了渠道消息处理的完整核心模块,包含以下核心功能: 1. 新增会话围栏类,实现会话并发控制与过期清理 2. 新增媒体清理器,实现过期媒体文件自动清理 3. 新增熔断器组件,实现服务降级与故障隔离 4. 新增消息处理器,完成渠道消息的完整流转处理 5. 新增限流器组件,实现渠道级和账户级流量控制 6. 新增链路追踪模块,集成Langfuse实现调用链路监控 7. 新增指标统计模块,实现消息处理全链路指标采集 8. 新增统一消息模型,封装全渠道消息格式 9. 新增块回复流水线,实现流式回复的合并与去重 10. 新增本地媒体存储模块,实现媒体文件的本地管理 11. 新增回复分发器,实现回复内容的有序发送与延迟处理 12. 完善__init__.py导出所有核心模块与工具类
53 lines
1.8 KiB
Python
53 lines
1.8 KiB
Python
import asyncio
|
|
import logging
|
|
import time
|
|
|
|
from yuxi.channel.message.models import UnifiedMessage
|
|
|
|
logger = logging.getLogger(__name__)
|
|
|
|
_DEFAULT_FENCE_TTL = 300.0
|
|
|
|
|
|
class ConversationFence:
|
|
def __init__(self, ttl: float = _DEFAULT_FENCE_TTL):
|
|
self._ttl = ttl
|
|
self._versions: dict[str, int] = {}
|
|
self._version_ts: dict[str, float] = {}
|
|
self._locks: dict[str, asyncio.Lock] = {}
|
|
self._abort_events: dict[str, asyncio.Event] = {}
|
|
|
|
@staticmethod
|
|
def key_for(msg: UnifiedMessage) -> str:
|
|
if msg.group and msg.group.id:
|
|
return f"{msg.channel_type}:{msg.account_id}:group:{msg.group.id}"
|
|
return f"{msg.channel_type}:{msg.account_id}:dm:{msg.sender.id}"
|
|
|
|
def enter(self, key: str) -> tuple[int, asyncio.Event]:
|
|
now = time.monotonic()
|
|
self._gc(now)
|
|
self._versions[key] = self._versions.get(key, 0) + 1
|
|
self._version_ts[key] = now
|
|
|
|
old_event = self._abort_events.get(key)
|
|
if old_event is not None and not old_event.is_set():
|
|
old_event.set()
|
|
logger.debug("Foreground fence: aborting previous run for %s", key)
|
|
|
|
new_event = asyncio.Event()
|
|
self._abort_events[key] = new_event
|
|
return self._versions[key], new_event
|
|
|
|
def lock_for(self, key: str) -> asyncio.Lock:
|
|
return self._locks.setdefault(key, asyncio.Lock())
|
|
|
|
def current_version(self, key: str) -> int:
|
|
return self._versions.get(key, 0)
|
|
|
|
def _gc(self, now: float) -> None:
|
|
expired = [k for k, ts in self._version_ts.items() if now - ts > self._ttl]
|
|
for k in expired:
|
|
self._versions.pop(k, None)
|
|
self._version_ts.pop(k, None)
|
|
self._locks.pop(k, None)
|
|
self._abort_events.pop(k, None) |