☰
Zeek MQTT 协议分析模块深入解析:mqtt.log 日志生成与 QoS 状态机实现
2026/10/9 4:53:45 网站建设 项目流程
  • 网络安全
  • 网络
  • IDS

【免费下载链接】zeek

Zeek is a powerful network analysis framework that is much different from the typical IDS you may know.

项目地址:https://gitcode.com/gh_mirrors/ze/zeek
点击查看免费下载

导读

本文以 Zeek 仓库中 base/protocols/mqtt/main.zeek 的官方文档(doc/scripts/base/protocols/mqtt/main.zeek.rst)为主体,结合其源码实现、binpac 解析器与 btest 测试用例,完整剖析 Zeek 对 MQTT(v3.1.1)协议的检测能力。读完本文,你将掌握 MQTT 分析模块的日志流设计、ConnectInfo / PublishInfo / SubscribeInfo 三种日志记录的字段语义、QoS 0/1/2 的完成度判定状态机、端口与常量表的可配置方式,以及如何通过策略钩子(PolicyHook)对日志做二次加工。

一、模块概览:功能定位与加载方式

base/protocols/mqtt/main.zeek实现了 Zeek 对MQTT(v3.1.1)协议的基础分析能力,核心产出是三类日志文件。模块归属命名空间MQTT,其自身导入base/protocols/mqtt/consts.zeek(consts.zeek),后者提供了协议常量定义。

整个 MQTT 支持目录的结构如下:

  • scripts/base/protocols/mqtt/load.zeek:加载入口,依次加载consts、main,并通过@load-sigs加载 DPD 签名;
  • scripts/base/protocols/mqtt/main.zeek:日志框架、类型定义、事件处理器;
  • scripts/base/protocols/mqtt/consts.zeek:MQTT 控制报文类型、协议版本、QoS 等级、CONNACK 返回码的常量映射表;
  • scripts/base/protocols/mqtt/dpd.sig:用于动态端口检测(DPD)的签名。

加载方式为:

@load base/protocols/mqtt

__load__.zeek中实际展开为加载consts、main与dpd.sig三部分。DPD 签名非常简洁:

signature dpd_mqtt { payload /^.{4,7}MQ/ enable "mqtt" }

该签名匹配报文前 4~7 字节之后紧跟MQ字样的载荷,用于在没有注册端口的情况下识别 MQTT 流量(mqtt为 mqtt.pac 中定义的 analyzer 名称)。需要说明的是,binpac 解析器(src/analyzer/protocol/mqtt/mqtt.pac)目前标注为支持v3.1.1,不含 v5.0,且流量以 datagram 方式送入MQTT_PDU解析单元。

二、日志设计:三个独立的日志流

模块在zeek_init事件中(优先级 5)注册了三个日志流(main.zeek):

event zeek_init() &priority=5 { Log::create_stream(MQTT::CONNECT_LOG, Log::Stream($columns=ConnectInfo, $ev=log_mqtt, $path="mqtt_connect", $policy=log_policy_connect)); Log::create_stream(MQTT::SUBSCRIBE_LOG, Log::Stream($columns=SubscribeInfo, $path="mqtt_subscribe", $policy=log_policy_subscribe)); Log::create_stream(MQTT::PUBLISH_LOG, Log::Stream($columns=PublishInfo, $path="mqtt_publish", $policy=log_policy_publish)); Analyzer::register_for_ports(Analyzer::ANALYZER_MQTT, ports); }

对应的三个Log::ID枚举在模块顶部通过redef enum Log::ID += { CONNECT_LOG, SUBSCRIBE_LOG, PUBLISH_LOG }扩展(main.zeek)。三者与输出文件、列类型的对应关系如下:

Log::ID日志路径记录类型策略钩子
MQTT::CONNECT_LOGmqtt_connect.logMQTT::ConnectInfoMQTT::log_policy_connect
MQTT::SUBSCRIBE_LOGmqtt_subscribe.logMQTT::SubscribeInfoMQTT::log_policy_subscribe
MQTT::PUBLISH_LOGmqtt_publish.logMQTT::PublishInfoMQTT::log_policy_publish

同一连接在三个日志流中各产生一条或多条记录,ts、uid、id三个字段是所有记录共有的基础字段,用于关联同一会话。

三、Redefinable Options:可配置端口

## Well-known ports for MQTT. const ports = { 1883/tcp } &redef;

MQTT::ports类型为set[port],属性为&redef,默认值为{ 1883/tcp }。它在zeek_init中被用于Analyzer::register_for_ports(Analyzer::ANALYZER_MQTT, ports),即告知分析框架:目的端口为 1883/tcp 的流量应交给 MQTT analyzer。若你的 MQTT 服务运行在非标准端口(例如 8883/tcp 的 TLS 端口或内网自定义端口),可以在local.zeek中重新定义:

redef MQTT::ports = { 1883/tcp, 8883/tcp };

注意:由于签名 DPD 机制的存在(见上文dpd.sig),即使不配置端口,Zeek 也可能通过载荷特征识别 MQTT 流量,但显式注册端口是更可靠、开销更低的方式。

四、核心类型:三条日志记录与连接状态

4.1 ConnectInfo:CONNECT/CONNACK 会话信息

MQTT::ConnectInfo记录(main.zeek)保存 MQTT 连接建立阶段的信息,全部字段带&log属性(ts/uid/id之外均为可选):

字段类型语义
tstime事件发生时间戳
uidstring连接唯一 ID
idconn_id连接的四元组(两端地址/端口)
proto_namestring协议名称(如MQTT)
proto_versionstring协议版本(如3.1.1)
client_idstring客户端唯一标识
connect_statusstring服务器对 CONNECT 请求的响应状态(来自 CONNACK 返回码)
will_topicstring"遗愿遗嘱"(Last Will and Testament)消息要发布的主题
will_payloadstring"遗愿遗嘱"消息的载荷

字段填充逻辑在mqtt_connect与mqtt_connack事件中完成(main.zeek):CONNECT 报文解析后填充协议名、版本、client_id、遗嘱主题与载荷;CONNACK 返回码通过return_codes表翻译成可读字符串并写入connect_status,随后Log::write(CONNECT_LOG, info)落盘。

版本与返回码的翻译表定义在 consts.zeek 中:

const versions = { [3] = "3.1", [4] = "3.1.1", [5] = "5.0", } &default = function(n: count): string { return fmt("unknown-version-%d", n); }; const return_codes = { [0] = "Connection Accepted", [1] = "Refused: unacceptable protocol version", [2] = "Refused: identifier rejected", [3] = "Refused: server unavailable", [4] = "Refused: bad user name or password", [5] = "Refused: not authorized", } &default = function(n: count): string { return fmt("unknown-return-code-%d", n); };

两张表都带&default兜底函数,未识别的取值会被格式化为unknown-version-N/unknown-return-code-N,避免脚本报错。

4.2 PublishInfo:PUBLISH 报文详情

MQTT::PublishInfo(main.zeek)记录一次 PUBLISH 消息的完整信息,是三条日志中字段最丰富的一条:

字段类型语义
tstimePUBLISH 消息开始的时间戳
uid/idstring/conn_id连接标识
from_clientbool消息由本连接客户端发布(T),还是服务器推送给客户端(F)
retainbool消息是否要求服务器保留(retained message)
qosstringQoS 等级的文本描述(at most once/at least once/exactly once)
statusstring发布状态,默认"incomplete_qos";QoS 完整交互完成则为"ok"
topicstring发布主题
payloadstring消息载荷(可能按MQTT::max_payload_size截断)
payload_lencount载荷实际长度,用于在payload被截断时还原真实长度
ackbool消息是否被 ACK(默认 F)
recbool服务器是否发送了 QoS2 的 RECEIVED 报文(PUBREC,默认 F)
relbool客户端是否发送了 QoS2 的 RELEASE 报文(PUBREL,默认 F)
compbool服务器是否发送了 QoS2 的 COMPLETE 报文(PUBCOMP,默认 F)
qos_levelcount内部用于比较的数值型 QoS 等级(默认 0)

其中ack、rec、rel、comp、qos_level五个字段不带&log属性,属于内部状态跟踪字段,不会输出到日志文件——它们专门服务于下面要讲的 QoS 完成度判定。

QoS 文本翻译表(consts.zeek):

const qos_levels = { [0] = "at most once", [1] = "at least once", [2] = "exactly once", } &default = function(n: count): string { return fmt("unknown-qos-level-%d", n); };

4.3 SubscribeInfo:订阅/取消订阅

MQTT::SubscribeInfo(main.zeek)同时服务于 SUBSCRIBE 与 UNSUBSCRIBE 两类操作,通过action字段区分:

字段类型语义
ts/uid/idtime/string/conn_id基础连接字段
actionMQTT::SubUnsub是订阅(MQTT::SUBSCRIBE)还是取消订阅(MQTT::UNSUBSCRIBE)
topicsstring_vec被订阅的主题(或主题通配符模式)列表
qos_levelsindex_vec各主题请求的 QoS 等级列表
granted_qos_levelcount服务器最终授予的 QoS 等级
ackbool服务器是否 ACK 了该请求(默认 F)

MQTT::SubUnsub是一个可重定义的枚举(&redef):

type MQTT::SubUnsub: enum { MQTT::SUBSCRIBE, MQTT::UNSUBSCRIBE, } &redef;

4.4 State:连接级 pub/sub 状态跟踪

MQTT::State(main.zeek)是挂载在connection记录上的内部数据结构,用于跟踪单个连接的发布/订阅消息状态:

type State: record { publish: table[count] of PublishInfo &optional &write_expire=5secs &expire_func=publish_expire; subscribe: table[count] of SubscribeInfo &optional &write_expire=5secs &expire_func=subscribe_expire; };

两个表均以msg_id(报文标识符)为主键:

  • publish:尚未完成记录/落盘的已发布消息;
  • subscribe:尚未被 ACK 或尚未记录落盘的订阅/取消订阅消息。

两张表都带有&write_expire=5secs与&expire_func属性。若某个msg_id在 5 秒内没有写入更新,将触发过期函数——这构成了 QoS 交互超时的兜底机制。

同时,模块通过redef record connection +=扩展了内置connection记录(main.zeek):

redef record connection += { mqtt: ConnectInfo &optional; mqtt_state: State &optional; };

set_session()函数负责按需初始化这两个字段(main.zeek):首次调用时以network_time()填充ts、c$uid、c$id,并创建空的publish/subscribe表。

五、事件与策略钩子

5.1 MQTT::log_mqtt

global MQTT::log_mqtt: event(rec: ConnectInfo);

该事件在 CONNECT 日志记录被送往日志框架时触发,仅针对ConnectInfo(即mqtt_connect.log对应的流)。通过在脚本中@if/event MQTT::log_mqtt挂钩,可以在记录落盘前访问并修改它。

5.2 三个 PolicyHook

钩子类型对应日志流
MQTT::log_policy_connectLog::PolicyHookmqtt_connect
MQTT::log_policy_publishLog::PolicyHookmqtt_publish
MQTT::log_policy_subscribeLog::PolicyHookmqtt_subscribe

Log::PolicyHook是 Zeek 日志框架的标准过滤机制,典型用法是在local.zeek或策略脚本中:

redef Log::PolicyHook = MQTT::log_policy_publish; # 不可直接赋值,见下 # 正确方式:扩展钩子处理函数 hook MQTT::log_policy_publish(rec: MQTT::PublishInfo, id: string) { # 例如过滤掉内部测试主题 if ( /^_test\// in rec$topic ) break; }

通过break语句可以阻止记录写入,通过修改rec字段可以改写输出内容。这是 Zeek 日志框架中所有Log::PolicyHook的统一语义。

六、QoS 状态机:status 字段如何被判定为 ok

PublishInfo$status字段的默认值是"incomplete_qos",只有完整观察到对应 QoS 级别的往返交互才会被置为"ok"。这一逻辑在mqtt_publish、mqtt_puback、mqtt_pubrec、mqtt_pubrel、mqtt_pubcomp等事件的成对(&priority=5/&priority=-5)处理器中实现(main.zeek)。

各 QoS 级别的判定条件:

  • QoS 0(at most once):mqtt_publish优先级 5 的处理器中直接pi$status="ok";优先级 -5 的处理器随即Log::write(PUBLISH_LOG, pi)并删除表中条目。QoS 0 无 ACK 交互,因此立即记录。
  • QoS 1(at least once):需要mqtt_puback确认。优先级 5 的处理器收到 PUBACK 后将pi$ack = T,若qos_level == 1则status = "ok";优先级 -5 的处理器在status == "ok"时落盘并删除条目。
  • QoS 2(exactly once):需要完整的 PUBREC → PUBREL → PUBCOMP 三次握手。mqtt_pubrec置rec = T,mqtt_pubrel置rel = T,mqtt_pubcomp在pi$qos_level == 2 && pi$rec && pi$rel && pi$comp时置status = "ok",随后优先级 -5 的处理器落盘。

若某个 QoS 1/2 消息在 5 秒内未完成交互,&write_expire=5secs会触发过期函数:

function publish_expire(tbl: table[count] of PublishInfo, idx: count): interval { Log::write(PUBLISH_LOG, tbl[idx]); return 0sec; }

publish_expire与subscribe_expire(main.zeek)的逻辑一致:把尚未落盘的记录直接写入对应日志流,返回0sec表示立即删除条目。这样即使客户端断开或丢包,观察到的半程 QoS 交互也会以"incomplete_qos"状态出现在日志中,不会静默丢失。subscribe_expire同理,兜底写入SUBSCRIBE_LOG。

订阅侧的状态转换则相对简单:mqtt_subscribe将请求存入subscribe表,mqtt_suback在收到服务器授予的 QoS 后写入granted_qos_level、置ack = T并立即落盘删除;mqtt_unsubscribe/mqtt_unsuback走同样的流程,仅action为MQTT::UNSUBSCRIBE(main.zeek)。

七、底层解析:binpac 解析器与事件生成

协议解析层位于 src/analyzer/protocol/mqtt/mqtt.pac。它声明了analyzer MQTT withcontext,连接由双向流组成,每个方向以MQTT_PDU(is_orig)作为 datagram 解析单元,并依次%include了connect.pac、connack.pac、publish.pac、puback.pac、pubrec.pac、pubrel.pac、pubcomp.pac、subscribe.pac、suback.pac、unsuback.pac、unsubscribe.pac、disconnect.pac、pingreq.pac、pingresp.pac等命令定义文件(位于src/analyzer/protocol/mqtt/commands/目录)。解析器在识别出各类型控制报文后,向脚本层抛出mqtt_connect、mqtt_connack、mqtt_publish、mqtt_puback、mqtt_pubrec、mqtt_pubrel、mqtt_pubcomp、mqtt_subscribe、mqtt_suback、mqtt_unsubscribe、mqtt_unsuback等事件,脚本层的事件处理器再驱动上文描述的日志逻辑。

mqtt-protocol.pac(src/analyzer/protocol/mqtt/mqtt-protocol.pac)提供MQTT_PDU等核心解析结构定义,负责解析 MQTT 固定报头、剩余长度编码以及各类型报文的可变头部与载荷。

八、常量总览(consts.zeek)

consts.zeek(consts.zeek)共定义四张常量表,覆盖 MQTT 报文类型(1~14 对应 connect 到 disconnect 的全部 14 种控制报文)、协议版本(3→3.1、4→3.1.1、5→5.0)、QoS 等级(0→at most once、1→at least once、2→exactly once)以及 CONNACK 返回码(0~5)。所有表均带&default匿名函数兜底,保证解析到未知数值时仍能安全输出形如unknown-msg-type-N的占位文本。

九、测试验证与实战演示

仓库自带三个 btest 用例(testing/btest/scripts/base/protocols/mqtt/):

  • mqtt.test:核心功能测试。执行zeek -b -r $TRACES/mqtt.pcap %INPUT,并对mqtt_connect.log、mqtt_subscribe.log、mqtt_publish.log三个输出做btest-diff基线比对;
  • mqtt-payload-cap.test:验证MQTT::max_payload_size对载荷截断的行为(对应payload_len字段的设计意图);
  • mqtt-payload-cap-dynamic.test:动态调整载荷上限后的行为验证。

本地复现 MQTT 分析的完整命令:

zeek -r mqtt.pcap base/protocols/mqtt

运行后会在当前目录生成三个文件:mqtt_connect.log、mqtt_subscribe.log、mqtt_publish.log。其中mqtt_publish.log是信息量最大的输出,典型的记录行(TSV 格式字段顺序对应PublishInfo定义)包含ts、uid、id、from_client、retain、qos、status、topic、payload、payload_len等列。若观察到的 QoS 1/2 交互不完整,status列会显示incomplete_qos;若握手完整,则为ok。

十、扩展实践:二次加工日志的策略示例

结合log_mqtt事件与log_policy_*钩子,可以在不修改核心脚本的前提下扩展 MQTT 分析。例如,将client_id提取为独立字段或按主题过滤:

@load base/protocols/mqtt # 1) 在 mqtt_connect 记录落盘前补充自定义字段的示例思路: # ConnectInfo 本身不带 &log 的自定义字段时,建议通过 Log::create_stream 新建流, # 或直接用 log_policy 钩子改写现有字段。 # 2) 丢弃针对内部测试主题的发布记录 hook MQTT::log_policy_publish(rec: MQTT::PublishInfo, id: string) { if ( /^internal\.test\./ in rec$topic ) break; } # 3) 关注高价值主题的 QoS 2 消息 event MQTT::log_mqtt(rec: MQTT::ConnectInfo) { if ( rec?$client_id ) NOTICE([$note=MQTT::... ]); # 示例示意,实际需按 Notice 框架完整构造 }

三个策略钩子的处理函数签名均为 Zeek 标准 PolicyHook 形式:hook(rec: <对应记录类型>, id: string)。所有扩展都应放置在local.zeek或自定义策略脚本中,经@load引入,以保持base/目录原始脚本的纯净。

十一、小结

base/protocols/mqtt/main.zeek是 Zeek 内置 MQTT v3.1.1 分析能力的核心脚本层,与 mqtt.pac 解析器、consts.zeek 常量定义和 dpd.sig 检测签名协同工作。其设计要点可归纳为:

  1. 三条日志流分别覆盖连接建立、发布、订阅三类 MQTT 活动,且共享ts/uid/id连接标识便于关联分析;
  2. QoS 状态机通过State记录中的&write_expire=5secs表与publish_expire/subscribe_expire兜底函数,确保不完整交互也能被记录,并明确标记incomplete_qos;
  3. 策略钩子log_policy_connect/log_policy_publish/log_policy_subscribe与事件log_mqtt为上层策略脚本提供了标准的日志改写与过滤入口;
  4. 可配置性集中在MQTT::ports(&redef端口集合)与consts.zeek的常量表,扩展非标准端口只需一行redef。

以上结论均以当前仓库 scripts/base/protocols/mqtt/main.zeek、consts.zeek 及对应测试用例为直接依据,读者可对照源码逐行验证。

  • 网络安全
  • 网络
  • IDS

【免费下载链接】zeek

Zeek is a powerful network analysis framework that is much different from the typical IDS you may know.

项目地址:https://gitcode.com/gh_mirrors/ze/zeek
点击查看免费下载
上一篇:Apache Paimon核心技术概念解析
下一篇:3种简单方法:免费提升macOS鼠标体验的终极指南

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询