mirror of
https://github.com/zhenxun-org/zhenxun_bot.git
synced 2026-09-28 16:20:56 +08:00
* ✨ feat!(llm): 重构并升级大语言模型服务为全新 AI 智能体框架 - 【重构】将原 services/llm 重构并迁移至全新的 services/ai 架构,提供向下兼容垫片 - 【新增】引入 Agent、Team、Workflow 三大智能体与工作流编排范式 - 【新增】引入基于 RAG 的长期向量记忆与中期槽位记忆系统 - 【新增】引入基于 Docker 的安全代码执行沙箱环境 - 【新增】支持 MCP 协议,允许动态管理和调用 MCP 服务 - 【新增】引入输入输出安全合规护栏与自愈反思机制 - 【优化】重构并优化多厂商 API 适配器 (Gemini, OpenAI, DeepSeek, GLM 等) - 【优化】优化日志脱敏与 Token 预估机制 - 【移除】移除旧版 llm default 和 llm reset-key 命令,新增 llm mcp 管理命令 * 🔧 chore(deps): 更新项目依赖与配置 - 添加 mcp、jieba 和 aiodocker 依赖到配置文件及 requirements.txt - 在 pyright 配置中设置 reportMissingImports 为 none - 调整 .gitignore 中 resources 目录的忽略规则 * ♻️ refactor(tools): 重构工具终止机制并清理知识库日志输出 - 统一使用 `context.state["__end_run__"]` 替代 `EndRunResult` 控制任务结束 - 移除文件系统和向量知识库检索工具中 `ToolResult` 的 `.with_log` 调用 - 调整指令处理器(Directive)的返回值为 `tool_res.output` - 修复部分类型检查警告并优化联合类型判断语法 * ♻️ refactor(tools): 重构工具副作用指令与控制流熔断机制 - 引入 `DirectivePayload` 及 `ToolResult` 的子类以结构化表达工具副作用 - 移除通过 `context.state` 传递魔术变量的隐式控制流设计 - 重构 `DirectiveManager` 处理器接口,直接在处理器中修改 `AgentState` 并构建 `AgentRunResult` - 在 `StandardAgentExecutor` 中统一通过 `directive_manager` 调度工具返回的副作用指令 - 补全 `MessageBuilder` 中部分核心方法的文档注释 * 🐛 fix(sandbox): 修复 Docker 沙箱容器状态检测与会话清理逻辑 -【修复】修正 `is_alive` 中直接读取私有属性的问题,改用 `show()` 返回值 -【修复】解决 `execute_code` 中缓存的执行器与当前会话不一致的问题 -【优化】在清理工作区前增加容器存活检测,避免向已死容器发送请求 -【优化】创建容器时增加运行状态校验,若已停止则自动从缓存中移除并重建 -【优化】优化容器销毁和清理逻辑,静默处理容器不存在 (404) 的异常 * 📝 docs(core): 补充核心模块初始化方法的文档注释 * 🚨 auto fix by pre-commit hooks --------- Co-authored-by: webjoin111 <455457521@qq.com> Co-authored-by: pre-commit-ci[bot] <66853113+pre-commit-ci[bot]@users.noreply.github.com>
182 lines
6.6 KiB
Python
182 lines
6.6 KiB
Python
import asyncio
|
|
from collections.abc import Awaitable, Callable
|
|
import heapq
|
|
import time
|
|
from typing import Any
|
|
|
|
from nonebot.utils import is_coroutine_callable
|
|
|
|
from zhenxun.services.log import logger
|
|
|
|
|
|
class LifespanManager:
|
|
"""
|
|
高精度资源生命周期调度器
|
|
"""
|
|
|
|
def __init__(self):
|
|
self._resources: dict[
|
|
Any,
|
|
tuple[
|
|
float,
|
|
float,
|
|
Callable[[Any], Awaitable[Any] | Any],
|
|
Callable[[], Awaitable[bool] | bool] | None,
|
|
],
|
|
] = {}
|
|
self._heap: list[tuple[float, Any]] = []
|
|
self._lock = asyncio.Lock()
|
|
self._wakeup_event = asyncio.Event()
|
|
self._watchdog_task: asyncio.Task | None = None
|
|
|
|
def _ensure_watchdog(self):
|
|
"""确保后台看门狗任务正在运行"""
|
|
if self._watchdog_task is None or self._watchdog_task.done():
|
|
self._watchdog_task = asyncio.create_task(self._watchdog_loop())
|
|
|
|
async def register(
|
|
self,
|
|
resource_id: Any,
|
|
ttl: float,
|
|
cleanup_callback: Callable,
|
|
is_busy_callback: Callable[[], Awaitable[bool] | bool] | None = None,
|
|
):
|
|
"""
|
|
将资源注册到生命周期管理器中。
|
|
|
|
参数:
|
|
resource_id: 资源的唯一标识符 (可以是字符串、数字或其他 Hashable 对象)。
|
|
ttl: 资源的存活时间 (秒)。
|
|
cleanup_callback: 资源过期时触发的回调函数,接收 resource_id 作为唯一参数。
|
|
is_busy_callback: (可选) 延迟存活探针。在触发清理前调用,若返回 True 则放弃清理并自动续期。
|
|
""" # noqa: E501
|
|
if ttl <= 0:
|
|
await self.unregister(resource_id)
|
|
return
|
|
async with self._lock:
|
|
expire_time = time.time() + ttl
|
|
self._resources[resource_id] = (
|
|
expire_time,
|
|
ttl,
|
|
cleanup_callback,
|
|
is_busy_callback,
|
|
)
|
|
heapq.heappush(self._heap, (expire_time, resource_id))
|
|
self._wakeup_event.set()
|
|
|
|
self._ensure_watchdog()
|
|
|
|
async def touch(self, resource_id: Any, ttl: float):
|
|
"""刷新资源的存活时间,为其续命"""
|
|
if ttl <= 0:
|
|
await self.unregister(resource_id)
|
|
return
|
|
async with self._lock:
|
|
if resource_id in self._resources:
|
|
_, original_ttl, cb, is_busy = self._resources[resource_id]
|
|
expire_time = time.time() + ttl
|
|
self._resources[resource_id] = (expire_time, original_ttl, cb, is_busy)
|
|
heapq.heappush(self._heap, (expire_time, resource_id))
|
|
self._wakeup_event.set()
|
|
|
|
async def unregister(self, resource_id: Any):
|
|
"""主动从管理器中注销资源 (不再触发超时回收)"""
|
|
async with self._lock:
|
|
self._resources.pop(resource_id, None)
|
|
|
|
async def _watchdog_loop(self):
|
|
"""核心看门狗循环"""
|
|
try:
|
|
while True:
|
|
await self._wakeup_event.wait()
|
|
self._wakeup_event.clear()
|
|
|
|
while True:
|
|
async with self._lock:
|
|
if not self._heap:
|
|
break
|
|
expire_time, res_id = self._heap[0]
|
|
|
|
if (
|
|
res_id not in self._resources
|
|
or self._resources[res_id][0] != expire_time
|
|
):
|
|
heapq.heappop(self._heap)
|
|
continue
|
|
|
|
now = time.time()
|
|
sleep_time = expire_time - now
|
|
|
|
if sleep_time > 0:
|
|
try:
|
|
await asyncio.wait_for(
|
|
self._wakeup_event.wait(), timeout=sleep_time
|
|
)
|
|
self._wakeup_event.clear()
|
|
continue
|
|
except asyncio.TimeoutError:
|
|
pass
|
|
|
|
async with self._lock:
|
|
if (
|
|
res_id not in self._resources
|
|
or self._resources[res_id][0] != expire_time
|
|
):
|
|
continue
|
|
_, original_ttl, cb, is_busy = self._resources[res_id]
|
|
|
|
is_active = False
|
|
if is_busy is not None:
|
|
try:
|
|
res = is_busy()
|
|
if isinstance(res, Awaitable):
|
|
is_active = await res
|
|
else:
|
|
is_active = res
|
|
except Exception as e:
|
|
logger.error(
|
|
f"执行资源存活探针失败: {e}", command="LifespanManager"
|
|
)
|
|
|
|
if is_active:
|
|
logger.debug(
|
|
f"探针检测到资源 '{res_id}' 仍在忙碌,"
|
|
f"已自动续期 ({original_ttl}s)。",
|
|
command="LifespanManager",
|
|
)
|
|
await self.touch(res_id, original_ttl)
|
|
continue
|
|
|
|
async with self._lock:
|
|
if (
|
|
res_id in self._resources
|
|
and self._resources[res_id][0] == expire_time
|
|
):
|
|
self._resources.pop(res_id)
|
|
heapq.heappop(self._heap)
|
|
else:
|
|
continue
|
|
|
|
logger.info(
|
|
f"♻️ 资源 '{res_id}' 闲置超时,触发自动回收。",
|
|
command="LifespanManager",
|
|
)
|
|
try:
|
|
if is_coroutine_callable(cb):
|
|
await cb(res_id)
|
|
else:
|
|
cb(res_id)
|
|
except Exception as e:
|
|
logger.error(
|
|
f"回收资源 '{res_id}' 时发生业务异常: {e}",
|
|
command="LifespanManager",
|
|
)
|
|
|
|
except asyncio.CancelledError:
|
|
pass
|
|
|
|
async def stop(self):
|
|
"""停止生命周期管理器"""
|
|
if self._watchdog_task:
|
|
self._watchdog_task.cancel()
|