ForcePilot/backend/package/yuxi/channels/adapters/wechat/polling_lease.py
Kris 05abecc02b feat(wechat): 新增完整微信渠道适配器实现
该提交实现了支持企业微信、微信公众号、个人微信桥接三种模式的完整微信渠道适配器,包含以下核心模块:
1. 基础认证与配置相关:auth_adapter、config_reload、setup_contract等
2. 消息处理与格式转换:format、attachment_adapter、outbound_adapter等
3. 多模式客户端支持:wecom/mp子模块,包含加解密、消息收发能力
4. 辅助能力:限速器、防抖、会话绑定、事件映射、模板渲染等
5. 扩展能力:二维码登录、消息读取、特权用户、心跳监控等

实现了完整的微信生态对接能力,支持消息收发、事件处理、API调用限流、配置热重载等功能。
2026-05-12 00:51:04 +08:00

88 lines
2.7 KiB
Python

from __future__ import annotations
import asyncio
import os
import time
from typing import Any
from yuxi.utils.logging_config import logger
class PollingLease:
def __init__(
self,
channel_id: str = "",
lease_ttl: float = 15.0,
renew_interval: float = 0.5,
):
self._channel_id = channel_id or os.environ.get("CHANNEL_INSTANCE_ID", "wechat-default")
self._lease_ttl = lease_ttl
self._renew_interval = renew_interval
self._acquired_at: float = 0.0
self._last_renew: float = 0.0
self._active = False
self._renew_task: asyncio.Task | None = None
@property
def is_active(self) -> bool:
if not self._active:
return False
if self._lease_ttl > 0 and time.monotonic() - self._last_renew > self._lease_ttl * 1.5:
self._active = False
logger.warning(f"[PollingLease/{self._channel_id}] Lease expired (stale)")
return False
return self._active
@property
def lease_id(self) -> str:
return self._channel_id
async def try_acquire(self, force: bool = False) -> bool:
if self._active and not force:
return True
self._acquired_at = time.monotonic()
self._last_renew = self._acquired_at
self._active = True
logger.info(
f"[PollingLease/{self._channel_id}] Lease acquired "
f"(ttl={self._lease_ttl}s, renew_interval={self._renew_interval}s)"
)
return True
async def start_renew(self) -> None:
if self._renew_task and not self._renew_task.done():
return
self._renew_task = asyncio.create_task(self._renew_loop())
async def stop_renew(self) -> None:
if self._renew_task and not self._renew_task.done():
self._renew_task.cancel()
try:
await self._renew_task
except asyncio.CancelledError:
pass
self._renew_task = None
async def _renew_loop(self) -> None:
while self._active:
try:
await asyncio.sleep(self._renew_interval)
self._last_renew = time.monotonic()
except asyncio.CancelledError:
break
async def release(self) -> None:
await self.stop_renew()
self._active = False
logger.info(f"[PollingLease/{self._channel_id}] Lease released")
def get_snapshot(self) -> dict[str, Any]:
return {
"channel_id": self._channel_id,
"active": self._active,
"acquired_at": self._acquired_at,
"last_renew": self._last_renew,
"lease_ttl": self._lease_ttl,
}