第 4 章 · 03 LogEngine 与 Email/Wechat 引擎


文档摘要

第 4 章 · 03 LogEngine 与 Email/Wechat 引擎 本节摘要:本节精读 VeighNa 的三大辅助引擎。LogEngine 极简——订阅 EVENTLOG,把 LogData 用 loguru 的 落盘,通过 总开关控制。EmailEngine 是懒启动典范:构造时不拉线程,首次 才 start;worker 线程每封邮件独立 连接(发完即关),用 Queue 异步投递。WechatEngine(4.4 新增)最复杂——凭据持久化到 wechatsetting.

第 4 章 · 03 LogEngine 与 Email/Wechat 引擎

本节摘要:本节精读 VeighNa 的三大辅助引擎。LogEngine 极简——订阅 EVENT_LOG,把 LogData 用 loguru 的 bind(gateway_name) 落盘,通过 SETTINGS["log.active"] 总开关控制。EmailEngine 是懒启动典范:构造时不拉线程,首次 send_email 才 start;worker 线程每封邮件独立 SMTP_SSL 连接(发完即关),用 Queue 异步投递。WechatEngine(4.4 新增)最复杂——凭据持久化到 wechat_setting.json,启动即拉 worker,核心机制是节流合并(send_interval=60 秒 + pending_msgs 缓冲 + 抽干队列批处理),并设计了 SessionExpired(回填消息+停 worker,等用户重新绑定)与 WeixinError(回填但不停,下次重试)两种异常自愈策略。读完本节,你掌握了"事件驱动 + 异步队列 + 节流重试"的工程范式。

内容来源:原项目源码 vnpy/trader/engine.py(LogEngine L325-357、EmailEngine L590-654、WechatEngine L657-838),精读并套用体系化模板。

学习目标

阅读完本节,你应当能够:

  1. 解释 LogEngine 的 level_map 为什么把整数级别映射成字符串。
  2. 读懂 EmailEngine 的懒启动机制——为什么构造时不 start。
  3. 说清 EmailEngine worker 循环为何"每封邮件独立连接"。
  4. 理解 WechatEngine 的节流合并——send_interval、pending_msgs、抽干队列三件套怎么配合。
  5. 区分 SessionExpired 与 WeixinError 两种异常的不同自愈策略。
  6. 看清 WechatEngine 的凭据持久化(bind/unbind/load_setting/save_setting)闭环。

一、LogEngine:日志落盘的最后一环

回到第 01 节那张图——MainEngine.write_log 只发 EVENT_LOG 事件不写盘,真正写盘的是 LogEngine。engine.py:325-357:

325 class LogEngine(BaseEngine): 330 level_map: dict[int, str] = { 331 DEBUG: "DEBUG", 332 INFO: "INFO", 333 WARNING: "WARNING", 334 ERROR: "ERROR", 335 CRITICAL: "CRITICAL", 336 } 337 338 def __init__(self, main_engine: MainEngine, event_engine: EventEngine) -> None: 340 super().__init__(main_engine, event_engine, "log") 341 342 self.active = SETTINGS["log.active"] 343 344 self.register_log(EVENT_LOG) 345 346 def process_log_event(self, event: Event) -> None: 347 if not self.active: 348 return 349 350 log: LogData = event.data 351 level: str | int = self.level_map.get(log.level, log.level) 352 logger.bind(gateway_name=log.gateway_name).log(level, log.msg) 353 354 def register_log(self, event_type: str) -> None: 356 self.event_engine.register(event_type, self.process_log_event)

1. level_map:整数 → 字符串

level_map 是类属性({DEBUG: "DEBUG", INFO: "INFO", ...})。注意 DEBUG/INFO/WARNING/ERROR/CRITICAL 这些 key 是从 vnpy.trader.logger import 的整数常量(loguru 用整数表示级别),value 是 loguru 接受的字符串名

为什么要这层映射?因为 LogData 的 level 字段可能存整数(策略层用 LogData(msg, level=INFO) 构造),而 loguru 的 .log() 方法既能接受整数也能接受字符串。但字符串更可读——日志文件里看到的是 INFO 而不是 20。所以这里统一转成字符串。

.get(log.level, log.level) 第二个参数是 log.level 本身——找不到映射时原样返回(容错)。这就是返回类型注解 str | int 的含义:正常情况是 str,异常情况兜底返回原 int。

2. active 开关

self.active = SETTINGS["log.active"]——构造时读一次全局配置。这是一个"假开关":它只控制 process_log_eventif not self.active: return,不取消事件订阅。即便 active=False,事件还是会路由到 process_log_event,只是函数体内立刻返回。这种设计的好处是:运行时改 SETTINGS["log.active"] 后,把 active 也同步改一下,就能动态开关日志,不必反注册事件。

3. loguru 的 bind 用法

logger.bind(gateway_name=log.gateway_name).log(level, log.msg)

logger 是全局 loguru 实例(从 vnpy.trader.logger import)。bind(gateway_name=...) 给这条日志打上来源标签(CTP/IB/MainEngine 等)。loguru 的输出格式串里可以引用 {extra[gateway_name]},从而把来源显示在日志行里。这是 VeighNa 区分日志来源的核心机制——比每个 gateway 自己维护 logger 实例轻量得多。

💡 核心心法:LogEngine 故意把"是否落盘"和"是否产生事件"分离。MainEngine.write_log 永远 put 事件(可能用于 UI 实时显示),而 LogEngine 决定要不要写文件。这样 UI 能看到所有日志,文件却能按需关闭——比如生产环境某些 gateway 噪音太大,关掉文件日志但保留 UI 显示。

二、EmailEngine:懒启动 + 每封独立连接

EmailEngine 是经典的"懒启动 + 异步队列"模式。engine.py:590-654:

590 class EmailEngine(BaseEngine): 594 def __init__(self, main_engine: MainEngine, event_engine: EventEngine) -> None: 597 super().__init__(main_engine, event_engine, "email") 598 599 self.thread: Thread = Thread(target=self.run) 600 self.queue: Queue[EmailMessage] = Queue() 601 self.active: bool = False # 构造时 active=False! 602 603 def send_email(self, subject: str, content: str, receiver: str | None = None) -> None: 606 if not self.active: 607 self.start() # 首次发邮件才 start 608 609 if not receiver: 610 receiver = SETTINGS["email.receiver"] 611 613 msg: EmailMessage = EmailMessage() 614 msg["From"] = SETTINGS["email.sender"] 615 msg["To"] = receiver 616 msg["Subject"] = subject 617 msg.set_content(content) 618 619 self.queue.put(msg)

1. 懒启动:构造时不 start

注意 __init__self.active: bool = False,没有 start。Thread 对象构造了(target=run),但 thread.start() 没调。这与第 01 节的四大内置引擎装配顺序呼应——init_enginesadd_engine(EmailEngine) 时只是构造,worker 线程不跑

为什么懒启动?因为大多数 VeighNa 实例从来不发邮件——很多人跑 CTA 不配置 SMTP。如果构造就拉线程,每个进程都白白占一条线程。懒启动保证"用才起,不用不耗"。

触发点在 send_email(606-607 行):首次调用时检测 if not self.active: self.start(),这时才真正 thread.start() 拉起 worker。后续调用 active 已是 True,直接 put 队列。

2. 队列异步投递

send_email 把构造好的 EmailMessage putself.queue——立刻返回,不等 SMTP 握手。SMTP 握手很慢(几秒),如果同步发,send_notification 会卡住调用线程(往往是策略主线程)。队列解耦让调用方零等待。

3. run worker 循环:每封独立连接

engine.py:621-641:

621 def run(self) -> None: 623 server = SETTINGS["email_server"] 624 port = SETTINGS["email.port"] 625 username = SETTINGS["email.username"] 626 password = SETTINGS["email.password"] 627 628 while self.active: 629 try: 630 msg = self.queue.get(block=True, timeout=1) 631 632 try: 633 with smtplib.SMTP_SSL(server, port) as smtp: 634 smtp.login(username, password) 635 smtp.send_message(msg) 636 smtp.close() 637 except Exception: 638 log_msg = _("邮件发送失败: {}").format(traceback.format_exc()) 639 self.main_engine.write_log(log_msg, "EmailEngine") 640 except Empty: 641 pass

结构与第 2 章的 _run 主循环完全同构:while self.active + queue.get(timeout=1) + 捕获 Empty。timeout=1 同样是为了优雅退出——close 时设 active=False,worker 最长 1 秒内感知并退出。

每封邮件独立连接:with smtplib.SMTP_SSL(server, port) 每封邮件新建一个 SSL 连接,发完 with 块退出自动关。这种设计看起来浪费(每次都要 TLS 握手),但好处是绝不复用可能已损坏的连接——SMTP 服务器经常踢闲置连接,复用容易踩坑。VeighNa 邮件频率极低(告警级别),性能不是问题,稳定性优先。

异常兜底:except Exception 捕获所有异常(网络断、登录失败、收件人不存在...),把 traceback 通过 main_engine.write_log 写入日志,然后继续循环——一封失败不影响后续邮件。注意这里是"内层 try",不会被外层 except Empty 误捕(因为内层 try 包着 SMTP 调用,不涉及队列)。

4. close 的防御

648 def close(self) -> None: 650 if not self.active: 651 return 652 653 self.active = False 654 self.thread.join()

if not self.active: return——如果 EmailEngine 从未被 send_email 触发过(active 一直是 False),close 直接返回,不会 join 一个从未 start 的线程(否则 RuntimeError)。这是懒启动模式的必要防御。

⚠️ 重要说明:EmailEngine 的凭据(server/port/username/password/sender/receiver)全部从 SETTINGS 全局配置读,而 SETTINGS 从 ~/.vntrader/vt_setting.json 加载。如果用户没配 email 段但调了 send_email,SETTINGS["email.sender"] 会抛 KeyError——所以实践中 send_notification 调用方要确保配置完整,否则会触发 run 循环里的 except Exception,邮件静默失败(只写日志)。

三、WechatEngine:节流合并 + 凭据持久化

WechatEngine(4.4 新增)是三个引擎里最复杂的——它对接微信推送(类似企业微信机器人),要处理凭据绑定、节流、会话过期。engine.py:657-838

1. init:启动即拉 worker(与 EmailEngine 对比)

664 def __init__(self, main_engine: MainEngine, event_engine: EventEngine) -> None: 666 super().__init__(main_engine, event_engine, "wechat") 667 668 self.creds: Credentials | None = None 669 self.user_id: str = "" 670 671 self.send_interval: int = 60 # 节流间隔(秒) 672 self.last_ts: float | None = None # 上次发送时间戳 673 674 self.pending_msgs: list[str] = [] # 待发送消息缓冲 675 676 self.queue: Queue[str] = Queue() 677 self.thread: Thread | None = None 678 self.active: bool = False 679 680 self.load_setting() # 1. 读凭据 682 if self.creds and self.user_id: 683 self.activate() # 2. 凭据齐全才启动

与 EmailEngine 的关键区别:WechatEngine 构造时就会拉 worker(只要凭据齐全)。这是因为微信推送比邮件高频(策略信号、风控告警),且 WechatEngine 的节流机制依赖 worker 持续运转(详第四节)。而邮件是"重大事件才发",懒启动更合适。

activate() 内部调 start(),start 会 new 一个 Thread 并启动(实现见后)。注意 self.thread: Thread | None = None——初始为 None,因为可能不启动(凭据不全)。这与 EmailEngine 的 self.thread: Thread = Thread(target=self.run)(构造时就 new)不同。

2. 凭据持久化闭环

四个方法形成闭环:

方法 作用
load_setting 启动时从 wechat_setting.json 读 bot_id/token/base_url/user_id/send_interval
save_setting 把当前 creds/user_id/send_interval 写回 json
bind(creds, user_id) 扫码绑定:先 deactivate → 存 creds → save_setting → activate
unbind 解绑:deactivate → 清空 creds → 清 pending → save_setting

load_setting(engine.py:685-705)的容错值得学:

687 data = load_json(self.setting_filename) 688 bot_id = data.get("bot_id") or "" 689 token = data.get("token") or "" 690 if bot_id and token: 691 base_url = data.get("base_url") or "https://ilinkai.weixin.qq.com" 692 self.creds = Credentials(bot_id=bot_id, token=token, base_url=base_url) 693 else: 694 self.creds = None 695 self.user_id = str(data.get("user_id") or data.get("chat_id") or "") 697 send_interval = data.get("send_interval") 698 if isinstance(send_interval, int) and send_interval >= 0: 699 self.send_interval = send_interval

data.get(key) or "" 处理 None/空字符串两种"缺失"情况;base_url 缺失给默认值;user_id 兼容旧字段名 chat_id;send_intervalisinstance 严格校验类型(防止 json 里被写成字符串)。这种防御式编程让配置文件半残缺也能加载,不会崩。

bind 的顺序很重要:先 deactivate 再赋值再 activate。deactivate 内部 stop 当前 worker 线程(join 等),保证后续 activate 起新线程时不冲突。这就是为什么 self.thread 要支持 None——bind/unbind 会反复置空。

3. start/stop 的线程安全

753 def start(self) -> None: 755 if self.active: 756 return 757 self.active = True 758 self.thread = Thread(target=self.run) 759 self.thread.start() 760 762 def stop(self) -> None: 764 if not self.active: 765 return 766 self.active = False 767 769 if ( 770 self.thread 771 and self.thread.is_alive() 772 and self.thread is not current_thread() # 关键! 773 ): 774 self.thread.join()

注意 stop 的第三个条件 self.thread is not current_thread()——如果 worker 线程自己调 stop(比如 run 循环里 SessionExpired 时设 active=False),不能 join 自己(否则死锁)。这种"自查 current_thread"的写法是处理"线程自停"的标准技巧,EmailEngine 没有这个需求(它的 run 循环不会自己调 close)。

4. run 工作循环:节流合并五步

核心机制(engine.py:783-834):

783 def run(self) -> None: 785 while self.active: 786 try: 787 msg = self.queue.get(block=True, timeout=1) # 第①步:阻塞读一条 788 self.pending_msgs.append(msg) 789 except Empty: 790 pass 791 792 while True: # 第②步:抽干队列 793 try: 794 self.pending_msgs.append(self.queue.get_nowait()) 795 except Empty: 796 break 797 798 if not self.pending_msgs: # 第③步:空兜底 799 continue 800 801 if not self.creds or not self.user_id: # 第④步:凭据校验 802 continue 803 804 gap = self.send_interval # 第⑤步:节流判断 805 if gap > 0: 806 prior = self.last_ts 807 if prior is not None: 808 wait_secs = prior + gap - time.monotonic() 809 if wait_secs > 0: 810 continue # 没到间隔,跳过本轮 811 812 msgs = list(self.pending_msgs) # 取出快照 813 self.pending_msgs.clear() 814 815 try: 816 self.last_ts = time.monotonic() 817 send_text(self.creds, self.user_id, "\n".join(msgs)) # 合并发送 818 except SessionExpired: 819 self.pending_msgs = msgs + self.pending_msgs # 回填 820 self.active = False # 停 worker 821 self.main_engine.write_log(_("微信会话已过期..."), self.engine_name) 822 except WeixinError as exc: 823 self.pending_msgs = msgs + self.pending_msgs # 回填 824 self.main_engine.write_log(_("微信推送失败:{}").format(exc), ...)

五步:

  1. 阻塞读一条:get(timeout=1) 拿到第一条消息(或超时空过)。timeout=1 同样为优雅退出。
  2. 抽干队列:while True + get_nowait() 把队列里剩余消息一次性全部取空,都 append 到 pending_msgs。这一步是"合并"的关键——60 秒内积压的多条消息会被攒到一起。
  3. 空兜底:如果 pending_msgs 是空(超时没消息),continue 跳过后续,回到循环开头。
  4. 凭据校验:万一 unbind 后 creds 被清空但 worker 还在转(理论上不会,但防御),不发送。
  5. 节流判断:如果距离上次发送不足 send_interval 秒,continue 等下一轮。注意用的是 time.monotonic()(单调时钟,不受系统时间调整影响),不能用 time.time()(NTP 校时可能跳变)。

发送时 "\n".join(msgs)——把多条消息合并成一条发送。这是节流的核心收益:60 秒内的 N 条消息只产生 1 次 API 调用。对微信这种有频控(每分钟限次)的接口至关重要。

💡 核心心法:节流 = 队列 + 缓冲 + 时间窗。VeighNa 的实现很朴素——没有用 token bucket 或 sliding window 算法,就是"上次发送时间 + 间隔"做判断。但这种朴素方案对"消息可合并"的场景最优:不仅限制频率,还减少调用次数。注意 send_interval=0 时禁用节流(805 行 if gap > 0),用户可配置成实时模式。

5. 两种异常的自愈策略

818 except SessionExpired: 819 self.pending_msgs = msgs + self.pending_msgs 820 self.active = False 821 except WeixinError as exc: 822 self.pending_msgs = msgs + self.pending_msgs 823 # 不停 worker

两种异常,回填策略相同(把刚取出的 msgs 拼回 pending_msgs 头部,下次重发),停 worker 策略不同:

  • SessionExpired(会话过期):回填 + 停 worker(active=False)。会话过期是"凭据失效",继续重试只会一直失败——必须停 worker 等用户重新扫码绑定(bind 会重新 activate)。日志里写"微信会话已过期,请通过菜单重新扫码绑定",提示用户行动。
  • WeixinError(一般错误):回填 + 不停 worker。这是临时性错误(网络抖动、API 限流),下一轮(60 秒后)自然会重试。worker 继续跑,等网络恢复就自动发出去了。

这种"区分错误严重性"的设计,是异步重试系统的精髓——永久性错误停止重试(避免雪崩),临时性错误继续重试(自愈)

注意回填用的是 msgs + self.pending_msgs——把这次取出的 msgs 拼到头部(因为 pending_msgs 在 813 行已被 clear,此时为空,等价于赋值,但写成拼接更通用,防止未来代码改动)。这保证了重发顺序与原始一致。

四、三大引擎对比

维度 LogEngine EmailEngine WechatEngine
启动时机 构造即订阅(无 worker) 懒启动(首次 send_email) 凭据齐全即启动
异步机制 同步(事件总线已是异步) Queue + Thread Queue + Thread
节流 send_interval + 合并
凭据来源 SETTINGS SETTINGS wechat_setting.json
异常处理 不调外部 全捕获+继续 区分永久/临时错误
复杂度 最简(30 行) 中等(60 行) 最复杂(180 行)

本节要点回顾

  1. LogEngine:level_map 把整数级别映射成 loguru 字符串名;active 是"假开关"(只控函数返回,不取消订阅);logger.bind(gateway_name=...) 给日志打来源标签。
  2. EmailEngine 懒启动:构造时 active=False 不拉线程;首次 send_email 才 start;Thread 对象一次性不能重启,故懒启动避免无谓线程。
  3. EmailEngine 每封独立连接:with SMTP_SSL 每封邮件新建连接发完即关,稳定性优先于性能;except Exception 全捕获,失败只写日志不影响后续。
  4. WechatEngine 节流合并五步:阻塞读一条 → 抽干队列 → 空兜底 → 凭据校验 → 节流判断(send_interval + monotonic 时钟)→ 合并发送。
  5. 两种异常自愈:SessionExpired 回填+停 worker(永久错误等用户重绑),WeixinError 回填不停(临时错误下轮重试)。
  6. 凭据持久化闭环:load_setting(容错读取)→ bind/deactivate→save→activate → unbind 清空→save;start/stop 用 current_thread() 自查防死锁。
  7. 设计哲学:LogEngine 极简(同步)、EmailEngine 懒启动(低频)、WechatEngine 节流合并(高频限流)——同一套 BaseEngine 框架,根据场景特点做出三种截然不同的实现。

第 4 章结束。你已经看清 VeighNa 的"装配中枢 + 四大引擎"全貌:MainEngine 什么都不做只装不实现,OmsEngine 缓存订单状态,Facade 技巧把子系统接口扁平化,三大辅助引擎示范了"事件驱动 + 异步队列 + 节流重试"的工程范式。下一章(第 5 章)往下走一层——精读 BaseGateway 与 BaseApp 这两个抽象基类,看插件如何被定义、如何被发现、如何与 MainEngine 对接。


发布者: 作者: 灏天文库 转发
评论区 (0)
U