生命周期管理
ErisPulse 提供统一的钩子/生命周期系统,用于监控系统各组件的运行状态,以及实现审计、统计、自定义逻辑等扩展功能。
系统支持三种触发方式:
await lifecycle.emit("event", data)— 精简版,传递任意数据(to="Owner"时定向投递)lifecycle.emit_sync("event", data)— 同步版(用于非异步上下文)await lifecycle.submit_event("event", ...)— 兼容旧版,自动构建标准事件格式
事件处理机制
注册处理器
from ErisPulse import sdk
# 装饰器模式
@sdk.lifecycle.on("module.load")
async def on_module_load(data):
print(f"模块加载: {data}")
# 编程式注册
sdk.lifecycle.register("module.load", on_module_load, priority=10)
# 取消注册
sdk.lifecycle.unregister("module.load", on_module_load)
# 按所有者批量取消注册(模块/适配器卸载时框架自动调用)
removed = sdk.lifecycle.unregister_by_owner("MyModule")
print(f"清理了 {removed} 个生命周期钩子")
优先级
处理器支持 priority 参数,数值越大越先执行(与模块加载器一致):
@sdk.lifecycle.on("adapter.event.receive", priority=10) # 最先执行
async def first_handler(data):
pass
@sdk.lifecycle.on("adapter.event.receive", priority=0) # 后执行
async def second_handler(data):
pass
点式结构事件
触发具体事件时,也会触发其父级事件:
- 触发
module.load时,也会触发module - 触发
adapter.event.receive时,也会触发adapter.event和adapter
通配符
注册 * 捕获所有事件:
@sdk.lifecycle.on("*")
async def on_anything(data):
print(f"收到事件: {data}")
定向传播(emit to=)
Note
本特性需要 ErisPulse **2.8.0+**。
emit() 指定 to 参数后进入定向传播:事件只分发给以该拥有者(owner)身份注册的
处理器(模块在 on_load 内注册的钩子自动归属本模块),其它模块与通配符 *
处理器不感知。
# 投递方:事件只投给 Chat 模块注册的钩子
await sdk.lifecycle.emit("message_received", {"text": "hi"}, to="Chat")
# 订阅方(Chat 模块内):注册同名钩子,owner 在注册时自动记录
@sdk.lifecycle.on("message_received")
async def on_message_received(data): ...
@sdk.lifecycle.on("message") # 点式父级前缀同样生效(按 owner 过滤)
async def on_any(data): ...
- 目标 owner 无已注册钩子 → 事件不被消费(可用
has_handlers()提前探测) data为 dict 时自动携带_trace_id(不覆盖已有值)emit_sync/submit_event同样支持to=参数- 模块间通信的三层模型(RPC / 定向 / 广播)见 模块间通信
一次性注册(once)
从 2.7.0 起,lifecycle.once() 注册的处理器在触发一次后自动注销,适合"首次就绪"这类一次性钩子:
@sdk.lifecycle.once("core.init.complete")
async def on_first_ready(data):
print("首次就绪,后续不再触发")
- 与
on()同优先级参数语义(priority数值越大越先执行) - 自动注销,无需手动
unregister - 同步/异步处理器均支持
监听者查询(has_handlers)
热路径短路场景可先用 has_handlers() 判断是否有监听者,避免无谓的事件遍历与任务调度:
if sdk.lifecycle.has_handlers("message.sending"):
await sdk.lifecycle.emit("message.sending", send_ctx)
- 覆盖精确事件名、通配符
*、父级事件三种匹配 - 无任何监听者时返回
False,可安全跳过emit
钩子断点一览
一条消息从平台进入框架到处理完成的典型生命周期事件时序:
sequenceDiagram
participant P as 平台
participant A as 适配器
participant F as 框架核心
participant M as 模块处理器
P->>A: 原生事件到达
A->>F: adapter.event.receive(最早期)
F->>F: event.pre_process(处理器执行前)
F->>M: 分发到处理器(命令/消息/通知等)
M->>M: command.matched / command.executed
M->>F: event.reply()
F->>F: message.sending(发送前)
F->>A: SendDSL 发送
A->>P: 发送到平台
A->>F: message.sent(发送完成)
F->>F: adapter.event.dispatched(分发完成)
框架内置了以下钩子断点,用户可以通过 @sdk.lifecycle.on() 监听任意断点实现自定义逻辑。
核心初始化
| 钩子名称 | 触发时机 | 数据 |
|---|---|---|
core.init.start |
SDK 初始化开始 | {} |
core.init.stage |
初始化各阶段开始(后台发射) | {"stage": str},取值 discovery / adapter_register / adapter_start / module_register / module_init / adapter_start_deferred / router_start |
core.init.complete |
SDK 初始化完成 | {"duration": float, "success": bool, "stages": {stage: float}, "adapters": {"enabled": [str], "disabled": [str]}, "modules": {"enabled": [str], "disabled": [str]}, "error": str(仅失败时)} |
core.uninit.complete |
SDK 反初始化完成 | {"duration": float, "success": bool, "adapters_closed": int, "modules_unloaded": int, "module_properties_cleared": int, "module_properties_to_clear": [str], "error": str(仅失败时)} |
示例:启动进度展示
@sdk.lifecycle.on("core.init.stage")
def show_stage(data):
print(f"[启动] 进入阶段: {data['stage']}")
配置变更
| 钩子名称 | 触发时机 | 数据 |
|---|---|---|
config.set |
配置项被修改 | {"key": str, "old_value": Any, "new_value": Any} |
config.updated |
外部编辑 config.toml 后检测到整树变更 | {"old_config": dict, "new_config": dict, "config_file": str} |
示例:配置审计
@sdk.lifecycle.on("config.set")
def audit_config(data):
print(f"[审计] {data['key']}: {data['old_value']} -> {data['new_value']}")
模块生命周期
| 钩子名称 | 触发时机 | 数据 |
|---|---|---|
module.register |
模块类注册到管理器 | {"module_name": str, "success": bool} |
module.load |
模块加载完成(实例化成功) | {"module_name": str, "success": bool} |
module.init |
模块初始化完毕(含懒加载) | {"module_name": str, "success": bool} |
module.unload |
模块卸载 | {"module_name": str, "success": bool} |
module.reload |
模块热重载完成(含级联重载依赖者) | {"module_name": str, "success": bool, "full": bool};全量重载(reload_all)时 module_name 为 "All",payload 额外携带 "results": dict[str, bool] |
适配器生命周期
| 钩子名称 | 触发时机 | 数据 |
|---|---|---|
adapter.load |
适配器注册完成 | {"platform": str, "success": bool} |
adapter.start |
适配器启动 | {"platforms": [str]} |
adapter.status.change |
适配器状态变化 | {"platform": str, "status": str, "retry_count": int, "error": str(仅失败时)};status 完整取值:starting / started / start_failed / stopping / stopped / stop_failed / skipped-dependency / disabled |
adapter.stop |
适配器关闭 | {"platforms": [str]} |
adapter.stopped |
适配器关闭完成 | {"platforms": [str]} |
adapter.bot.online |
Bot 上线 | {"platform": str, "bot_id": str, "info": dict, "status": str} |
adapter.bot.offline |
Bot 下线 | {"platform": str, "bot_id": str, "status": str} |
事件接收与处理
| 钩子名称 | 触发时机 | 数据 |
|---|---|---|
adapter.event.receive |
收到外部平台事件(最早期) | {"platform": str, "event_type": str, "raw_event_type": str} |
adapter.event.blocked |
中间件否决事件(返回 False,事件被丢弃不进入任何处理器) |
{"middleware": str, "platform": str, "event_type": str, "detail_type": str, "event": dict, "_trace_id": str} |
adapter.event.dispatched |
事件分发完成 | {"platform": str, "event_type": str, "raw_event_type": str, "onebot_handlers_count": int} |
event.pre_process |
事件处理器开始执行前 | {"event_type": str, "platform": str, "detail_type": str} |
示例:事件统计
event_counter = {}
@sdk.lifecycle.on("adapter.event.receive")
def count_events(data):
platform = data["platform"]
event_counter[platform] = event_counter.get(platform, 0) + 1
@sdk.lifecycle.on("adapter.event.dispatched")
def log_unhandled(data):
if data["onebot_handlers_count"] == 0:
print(f"[未处理] {data['platform']}/{data['event_type']}")
消息发送
| 钩子名称 | 触发时机 | 数据 |
|---|---|---|
message.sending |
消息即将发送 | {"platform": str, "method": str, "detail_type": str, "target_id": str, "bot_id": str} |
message.sent |
消息发送完成 | {"platform": str, "method": str, "detail_type": str, "target_id": str, "bot_id": str} |
示例:消息发送审计
@sdk.lifecycle.on("message.sending")
def log_sending(data):
print(f"[发送] -> {data['platform']}/{data['detail_type']}/{data['target_id']} via {data['method']}")
命令系统
| 钩子名称 | 触发时机 | 数据 |
|---|---|---|
command.matched |
命令被匹配并即将执行 | {"command": str, "args": list[str], "platform": str, "user_id": str} |
command.executed |
命令执行完成 | {"command": str, "args": list[str], "platform": str, "user_id": str, "success": bool, "error": str(仅失败时)} |
示例:命令统计
@sdk.lifecycle.on("command.matched")
def count_commands(data):
print(f"[命令] /{data['command']} from {data['user_id']}@{data['platform']}")
HTTP 路由
| 钩子名称 | 触发时机 | 数据 |
|---|---|---|
server.request |
HTTP 请求接收 | {"method": str, "path": str, "client_ip": str} |
server.response |
HTTP 响应发送 | {"method": str, "path": str, "status_code": int, "client_ip": str} |
示例:请求日志
@sdk.lifecycle.on("server.response")
def log_http(data):
print(f"[HTTP] {data['method']} {data['path']} -> {data['status_code']}")
WebSocket
| 钩子名称 | 触发时机 | 数据 |
|---|---|---|
server.start |
路由服务器启动 | {"base_url": str, "host": str, "port": int, "success": bool, "error": str(仅失败时)} |
server.stop |
路由服务器停止 | {} |
server.websocket.connect |
WebSocket 连接建立 | {"path": str, "module_name": str, "client_ip": str} |
server.websocket.disconnect |
WebSocket 连接断开 | {"path": str, "module_name": str, "reason": str, "error": str(仅异常时)} |
示例:WebSocket 连接监控
@sdk.lifecycle.on("server.websocket.connect")
def on_ws_connect(data):
print(f"[WS] 连接: {data['path']} from {data['client_ip']}")
@sdk.lifecycle.on("server.websocket.disconnect")
def on_ws_disconnect(data):
print(f"[WS] 断开: {data['path']} ({data['reason']})")
存储连接状态
存储后端连接池的建立、故障与恢复(均后台发射,不阻塞存储操作):
| 钩子名称 | 触发时机 | 数据 |
|---|---|---|
storage.ready |
存储后端连接池就绪(每事件循环首次建池成功) | {"backend": str} |
storage.unreachable |
连接重试耗尽进入冷却期(期间操作快速失败) | {"backend": str, "error": str, "cooldown": float} |
storage.recovered |
冷却结束重连成功,存储恢复可用 | {"backend": str} |
示例:存储故障告警
@sdk.lifecycle.on("storage.unreachable")
def alert_storage_down(data):
print(f"[告警] 存储后端 {data['backend']} 不可达: {data['error']},{data['cooldown']}s 后自动重连")
@sdk.lifecycle.on("storage.recovered")
def notify_storage_back(data):
print(f"[恢复] 存储后端 {data['backend']} 已恢复可用")
HTTP 客户端
sdk.client 的请求与连接事件(均后台发射):
| 钩子名称 | 触发时机 | 数据 |
|---|---|---|
client.request.success |
HTTP 请求成功 | {"method": str, "url": str, "status": int, "elapsed": float} |
client.request.failed |
HTTP 请求重试耗尽最终失败 | {"method": str, "url": str, "error": str, "attempts": int, "elapsed": float} |
client.ws.connect |
WebSocket 连接建立 | {"url": str} |
国际化
| 钩子名称 | 触发时机 | 数据 |
|---|---|---|
i18n.language.changed |
框架语言切换(i18n.set_language) |
{"language": str, "previous": str} |
标准事件定义
STANDARD_EVENTS = {
"core": ["init.start", "init.stage", "init.complete", "uninit.complete"],
"module": ["load", "init", "unload", "register", "reload"],
"adapter": [
"load", "start", "status.change", "stop", "stopped",
"event.receive", "event.dispatched",
"bot.online", "bot.offline",
],
"server": [
"start", "stop",
"request", "response",
"websocket.connect", "websocket.disconnect",
],
"event": ["pre_process"],
"message": ["sending", "sent"],
"command": ["matched", "executed"],
"config": ["set", "updated"],
"storage": ["ready", "unreachable", "recovered"],
"client": ["request.success", "request.failed", "ws.connect"],
"i18n": ["language.changed"],
}
完整 API 参考
注册与取消
| 方法 | 说明 |
|---|---|
@lifecycle.on(event, *, priority=0) |
装饰器注册处理器 |
lifecycle.register(event, handler, *, priority=0) |
编程式注册 |
lifecycle.unregister(event, handler=None) |
取消注册(handler=None 时取消该事件全部处理器) |
触发
| 方法 | 说明 |
|---|---|
await lifecycle.emit(event, data=None, *, to=None) |
异步触发,处理器并行执行(互不阻塞,返回时全部完成),返回非 None 值按优先级顺序回放链式替换 data;to 指定 owner 时定向投递 |
lifecycle.fire(event, data=None, *, to=None) |
后台发射(扔桶即走):处理器在后台任务中并行执行、不等待、无返回值;无监听者时零开销。适用于高频热路径与纯观测事件;关停序列与顺序敏感消费(如 config.set)请用 emit |
lifecycle.emit_sync(event, data=None, *, to=None) |
同步触发,异步处理器以 create_task 调度 |
await lifecycle.submit_event(event_type, *, source, msg, data, to=None, background=False) |
兼容旧版,自动构建标准事件格式;background=True 时走 fire 后台发射 |
工具
| 方法 | 说明 |
|---|---|
lifecycle.start_timer(timer_id) |
开始计时 |
lifecycle.get_duration(timer_id) |
获取已持续时间(秒) |
lifecycle.stop_timer(timer_id) |
停止计时并返回持续时间 |
lifecycle.list_hooks() |
列出所有已注册钩子及处理器数量 |
lifecycle.clear() |
清除所有处理器和计时器 |
模块中使用示例
from ErisPulse.Core.Bases import BaseModule
from ErisPulse import sdk
class Main(BaseModule):
async def on_load(self, event):
# 实现简单的消息统计
self.msg_count = 0
@sdk.lifecycle.on("adapter.event.receive")
async def count(data):
if data["event_type"] == "message":
self.msg_count += 1
# 监控所有命令
@sdk.lifecycle.on("command.matched")
async def log_cmd(data):
sdk.logger.info(f"命令执行: /{data['command']} by {data['user_id']}")
# 配置变更审计
@sdk.lifecycle.on("config.set")
def audit(data):
sdk.logger.info(f"配置变更: {data['key']} = {data['new_value']}")
后台任务归属与自动取消
Note
本特性需要 ErisPulse **2.8.0+**。
模块创建的 asyncio 后台任务若未在 on_unload 中取消,会持有 self 引用导致模块实例无法被回收(热重载后旧实例残留)。框架提供以下兜底机制:
- **
self.spawn(coro)**(模块内推荐):任务自动归属模块名,模块卸载时框架在on_unload之后兜底取消未结束的任务并记录警告 - **
spawn_background(coro)**(ErisPulse.runtime):自动捕获当前owner_scope上下文;cancel_owner_tasks(owner)按归属取消,cancel_all_background_tasks()供sdk.uninit()兜底 - 适配器:关闭时对平台名下的后台任务同样兜底取消
async def on_load(self, event):
# 推荐:后台任务用 self.spawn(),卸载时框架自动兜底取消
self.spawn(self._poll())
async def on_unload(self, event):
# 精细控制的场景仍建议自行取消并等待收尾
if self._poll_task:
self._poll_task.cancel()
await asyncio.gather(self._poll_task, return_exceptions=True)
async def _poll(self):
while True:
await asyncio.sleep(60)
...
Important
框架兜底是强制 cancel(cancel_owner_tasks),它发生在 on_unload 返回之后。因此需要优雅收尾的任务(flush 缓冲、持久化状态、关闭连接)必须在 on_unload 里自行 cancel() + await 完成——别指望兜底能保留收尾逻辑。框架只保证「不残留持有 self 的任务」,不保证「优雅」。需要 await 结果的任务请直接 await,不要丢给后台任务。
注意事项
- 处理器可以是同步或异步:系统自动识别并正确调用
- 数据传递:
emit()模式下,处理器返回非 None 值会修改传递给后续处理器的 data - 事件命名规范:建议使用点式结构命名事件,便于使用父级监听
- 错误隔离:单个处理器异常不会影响其他处理器执行
- 同步触发限制:
emit_sync()中异步处理器以 fire-and-forget 方式调度,返回值无法回传 - 生命周期清理:调用
sdk.uninit()时,所有已注册的处理器和计时器会被清理 - 加载优先性:如需在框架初始化阶段就监听事件,建议设置高优先级并禁用懒加载