第 4 章 · 03 LogEngine 与 Email/Wechat 引擎 本节摘要:本节精读 VeighNa 的三大辅助引擎。LogEngine 极简——订阅 EVENTLOG,把 LogData 用 loguru 的 落盘,通过 总开关控制。EmailEngine 是懒启动典范:构造时不拉线程,首次 才 start;worker 线程每封邮件独立 连接(发完即关),用 Queue 异步投递。WechatEngine(4.4 新增)最复杂——凭据持久化到 wechatsetting.
本节摘要:本节精读 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),精读并套用体系化模板。
阅读完本节,你应当能够:
level_map 为什么把整数级别映射成字符串。回到第 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)
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。
self.active = SETTINGS["log.active"]——构造时读一次全局配置。这是一个"假开关":它只控制 process_log_event 里 if not self.active: return,不取消事件订阅。即便 active=False,事件还是会路由到 process_log_event,只是函数体内立刻返回。这种设计的好处是:运行时改 SETTINGS["log.active"] 后,把 active 也同步改一下,就能动态开关日志,不必反注册事件。
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 是经典的"懒启动 + 异步队列"模式。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)
注意 __init__ 里 self.active: bool = False,没有 start。Thread 对象构造了(target=run),但 thread.start() 没调。这与第 01 节的四大内置引擎装配顺序呼应——init_engines 里 add_engine(EmailEngine) 时只是构造,worker 线程不跑。
为什么懒启动?因为大多数 VeighNa 实例从来不发邮件——很多人跑 CTA 不配置 SMTP。如果构造就拉线程,每个进程都白白占一条线程。懒启动保证"用才起,不用不耗"。
触发点在 send_email(606-607 行):首次调用时检测 if not self.active: self.start(),这时才真正 thread.start() 拉起 worker。后续调用 active 已是 True,直接 put 队列。
send_email 把构造好的 EmailMessage put 到 self.queue——立刻返回,不等 SMTP 握手。SMTP 握手很慢(几秒),如果同步发,send_notification 会卡住调用线程(往往是策略主线程)。队列解耦让调用方零等待。
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 调用,不涉及队列)。
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(4.4 新增)是三个引擎里最复杂的——它对接微信推送(类似企业微信机器人),要处理凭据绑定、节流、会话过期。engine.py:657-838。
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)不同。
四个方法形成闭环:
| 方法 | 作用 |
|---|---|
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_interval 用 isinstance 严格校验类型(防止 json 里被写成字符串)。这种防御式编程让配置文件半残缺也能加载,不会崩。
bind 的顺序很重要:先 deactivate 再赋值再 activate。deactivate 内部 stop 当前 worker 线程(join 等),保证后续 activate 起新线程时不冲突。这就是为什么 self.thread 要支持 None——bind/unbind 会反复置空。
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)。
核心机制(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), ...)
五步:
get(timeout=1) 拿到第一条消息(或超时空过)。timeout=1 同样为优雅退出。while True + get_nowait() 把队列里剩余消息一次性全部取空,都 append 到 pending_msgs。这一步是"合并"的关键——60 秒内积压的多条消息会被攒到一起。continue 跳过后续,回到循环开头。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),用户可配置成实时模式。
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 策略不同:
active=False)。会话过期是"凭据失效",继续重试只会一直失败——必须停 worker 等用户重新扫码绑定(bind 会重新 activate)。日志里写"微信会话已过期,请通过菜单重新扫码绑定",提示用户行动。这种"区分错误严重性"的设计,是异步重试系统的精髓——永久性错误停止重试(避免雪崩),临时性错误继续重试(自愈)。
注意回填用的是 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 行) |
level_map 把整数级别映射成 loguru 字符串名;active 是"假开关"(只控函数返回,不取消订阅);logger.bind(gateway_name=...) 给日志打来源标签。with SMTP_SSL 每封邮件新建连接发完即关,稳定性优先于性能;except Exception 全捕获,失败只写日志不影响后续。current_thread() 自查防死锁。第 4 章结束。你已经看清 VeighNa 的"装配中枢 + 四大引擎"全貌:MainEngine 什么都不做只装不实现,OmsEngine 缓存订单状态,Facade 技巧把子系统接口扁平化,三大辅助引擎示范了"事件驱动 + 异步队列 + 节流重试"的工程范式。下一章(第 5 章)往下走一层——精读 BaseGateway 与 BaseApp 这两个抽象基类,看插件如何被定义、如何被发现、如何与 MainEngine 对接。