[C-01][high] 推送与调度循环重构:握手前就推送、同端并发推送、单协程串行循环被慢操作冻结、到点分发每秒最多 100 条 #32

Closed
opened 2026-09-30 13:56:58 +08:00 by nixevol · 2 comments
Owner

编号:C-01 严重级:high 工作线:消息核心(internal/app/message、serve 的 messageLoops) 来源:审查 M-03、P-02、M-11、M-13、M-12、M-15、M-02
依赖:B-02 (#9)(第 5 点删除消息包大帧名额需在其后) 被依赖:C-02 (#33)、C-03 (#34)、D-03 (#29)、L-03 (#22)

结论与统一方案

本 issue 合并审查 M-02、M-03、M-11、M-12、M-13、M-15 与 P-02、P-07(循环部分,原文附在 B-03)。它们都落在 internal/app/message/push.go 与 cmd/nixmsg/serve.go 的 messageLoops,必须由一人按一个设计改完,不能拆给多人。统一设计:

  1. 每连接推送 worker:OnHandshakeComplete 为该连接启动 worker(带容量为 1 的唤醒通道,合并唤醒),OnDisconnect 停止并回收。WakePush 与每秒循环只发信号。未握手、代号不符或已断开的连接不推送、不 claim,修复握手前推送(审查 M-03、P-02)。OnHandshakeComplete 推送前先清掉 pushed_conn=本连接 的旧标记作兜底。
  2. 窗口与顺序:每轮在一个写操作里 claim 本轮全部条目(逐条条件更新,返回成功集合),按 (send_at, seq) 顺序发布;确认超时处理在写操作里再核对 pushed_at。(审查 M-11)
  3. 循环拆分:messageLoops 拆为独立 goroutine,由 WaitGroup 管理(供 L-03 停机等待):
    • 到点分发:按最早 send_at 设定时器,提交定时消息时通知;每轮循环到取空或用完时间预算;有限并发提交写操作让写队列合批;单条失败只记日志并跳过。(审查 M-12)
    • 到期处理:每秒一次。
    • 清理:每小时一次(内容见 C-03)。
    • 指标采样:10–15 秒一次。
    • 每次调用都带 context.WithTimeout。启动恢复只做 SQL 修正,分发交给循环,不再阻塞开始监听。(审查 M-13、P-07)
  4. 上下文:claim 返回 ctx 错误时读一次 pushed_conn 判断是否已写入并清掉;clearPushed 固定用新建的短超时 ctx。(审查 M-15)
  5. 大帧名额:删除 message 包内的 largeSem、largeHeld、acquireLarge、trackLarge、releaseLarge,只依赖 broker 的名额(B-02 合入后做);测试用假下行如需限流,在假实现里模拟。同步更新 DEVIATIONS M2/M3/M4 第 2 条。(审查 M-02)
  6. 与 broker 的契约:PublishDown 返回 ErrBackpressure(B-03)、ErrNotSubscribed、ErrNoConnection(B-06)时一律按发布失败处理(清标记、1 秒后重推),回执与 revoked 路径同样处理。

改动文件

internal/app/message/push.go、app.go、conn.go、session.go、recover.go(RecoverOnStart);cmd/nixmsg/serve.go(messageLoops 改为调用 message 提供的启动与等待接口)。

与其他问题的交互 / 冲突说明

  • 第 5 点依赖 B-02;与 B-03、B-06 以 PublishDown 签名不变为契约,可并行开发。
  • C-02、C-03、D-03、L-03 依赖本条。
  • cmd/nixmsg/uplink.go 的生命周期函数由 B-09 修改;本条只通过 msg.OnHandshakeComplete / msg.OnDisconnect 启停 worker,不改 uplink.go。
  • serve.go 只改 messageLoops。

验收与测试

  • 未就绪时提交消息再唤醒:没有 msg 发布、pushed_conn 为空;握手后才发布。
  • 裸 MQTT 客户端 CONNECT 后延迟 1.5 秒再 SUBSCRIBE 与 hello,仍能在 2 秒内收到排队中的不保留消息(没有被标成 dropped)。
  • 50 个并发唤醒下,任一时刻在途数不超过窗口,同一端到达顺序与 (send_at, seq) 一致。
  • 注入会阻塞的假下行:其他端的推送与定时分发照常。
  • 250 条同时到点的消息 1 秒内分发完;构造一条会失败的消息,不影响其余。
  • 大帧"推送+断线+重推"70 轮后大帧仍能推送。

问题明细(各区审查原文,证据含文件与行号)

以下是本次复审各区审查报告的原文段落。A、M、I、P、S 开头的是原始发现编号(A 管理后台与网页、M 消息核心、I 身份认证群在线、P 传输平台部署、S SDK)。解决方案以本 issue 上方的"结论与统一方案"为准;原文里的方案与之不一致时,按上方执行。

[M-03] 握手前就推送:订阅前推出的帧静默丢失,推送窗口被占 5 分钟,还绕过了 max_receive_bytes

  • 严重级:high
  • 分类:协议一致性 / 逻辑
  • 现象与影响:
    • 连接一建立(CONNACK 之后,SUBSCRIBE 和 hello 之前),OnSessionEstablished 就把它登记进连接表。messageLoops 和各处 WakePush 都会对它调用 PushPending,而 PushPending 不看是否已握手。
    • Clean Start 时 mochi 已经删掉旧订阅,这时向 down 主题发布没有订阅者,Publish 返回 nil,帧静默丢失。投递却已经记成「已推给这个连接」:
      • 占用推送窗口,直到 5 分钟确认超时。
      • 不保留的消息 5 分钟后被判 dropped/not_acked,而接收端其实一直在线。
      • 如果丢的正好是整窗 32 条,这个端 5 分钟内收不到任何消息。
    • 顶号重连更容易命中:旧连接的 OnDisconnect 发现已有新连接时会立即 WakePush,几乎一定早于新连接的 SUBSCRIBE。
    • hello 之前 MaxReceiveBytes=0,超限的大帧会直接发给只声明了 1024 字节的小设备,而不是拒收并回 too_large。
    • 在弱网(2 秒延迟)或服务器重启后大批重连时,几乎每次都会命中。
  • 证据:
func (u *appUplink) OnSessionEstablished(ctx context.Context, conn port.ConnInfo) error {
	u.conns.Set(conn.EndpointID, message.LiveConn{
		ConnID:        conn.ConnID,
		MaxPacketSize: conn.MaxPacketSize,
	})
	return nil
}
  • push.go:88-99:只比对连接代号。
  • push.go:549-560:WakePush 同样不看握手状态。
  • session.go:71-73:旧连接断开时唤醒新连接。
  • mochi server.go:457-476:先继承会话,再发 CONNACK,再调 OnSessionEstablished,之后才开始读包。
  • mochi server.go:1001-1021:没有订阅者时什么都不做,投递失败只记 debug 日志。
  • 文档依据:
    • DEVELOPMENT 6.1「握手完成才算在线,才开始推送」;7.5「每个已握手连接一个推送循环」。
    • PRD F07:接收上限在握手时声明;F08:「对方已完成登录握手:立刻推送」。
  • 为何不是故意设计:7.4 让分发时把「握手中」的连接算作在线,只用于计算 expire_at;推送必须等握手完成。DEVIATIONS 里没有相关说明。
  • 解决方案:
    1. LiveConn 加 Ready bool。OnSessionEstablished 登记为 Ready:false,OnHandshakeComplete 登记为 Ready:true 并带上 MaxReceiveBytes。
    2. PushPending 遇到当前连接不存在、代号不匹配或未就绪时直接返回:不 claim,也不发回执。WakePush 和 messageLoops 同样只处理就绪连接。
    3. dispatchFullTx 的在线判定保持现状(含握手中),符合 7.4。
    4. 可选的纵深防御(连接线):PublishDown 在没有 down 订阅时返回错误,让 message 走清标记重推。
  • 改动文件:internal/app/message/conn.go、push.go;cmd/nixmsg/uplink.go、serve.go;可选 internal/broker/broker.go。
  • 交互/冲突风险:presenceConnTable.IsOnline 用的是同一张表,建议一起改成只认就绪连接,以符合「在线 = 完成握手」(DEVELOPMENT 第 1 节)。
  • 需补测试:
    • 单测:未就绪时提交消息再调 PushPending,没有 msg 发布、pushed_conn 为空;置为就绪后才发布。
    • 集成:裸 MQTT 客户端 CONNECT 后等 1.5 秒再 SUBSCRIBE 和 hello,2 秒内应收到排队中的不保留消息。
    • 集成:hello 声明 1024 字节的连接,对排队中的 2048 字节消息得到 too_large 回执。
  • 置信度:代码阅读确定(已对照 mochi 调用顺序)。

[P-02] 握手完成前就推送:订阅 down 之前推出的消息被 mochi 静默丢弃,却记成"已推送",不保留的消息 5 分钟后按 not_acked 丢弃

  • 严重级:high
  • 分类:逻辑 / 协议一致性 / 数据
  • 现象与影响:
    • appUplink.OnSessionEstablished 在回 CONNACK 之后、客户端发 SUBSCRIBE 之前,就把这个连接登记进了 memConns。
    • 两条推送路径会命中这个窗口:
      • messageLoops 每秒对 conns.Snapshot() 里的每个连接调用 PushPending。
      • message 的 OnDisconnect 在顶号时会立刻 WakePush 新连接。
    • 这时新连接还没订阅 nix/c/{id}/down:
      • server.Publish 找不到订阅者,直接返回 nil(也不触发 OnPublishDropped),PublishDown 于是报成功。
      • 投递已经被条件更新成 pushed_conn=新连接。
      • 握手完成时,message.OnHandshakeComplete 只推 pushed_conn IS NULL 的投递,这批"推进空气"的投递要等 5 分钟确认超时:不保留的按 D19 改成 dropped/not_acked,保留的晚到 5 分钟。
    • 最常见的触发场景是顶号重连:移动端断网后用令牌重连,服务器还没发现旧连接已死。
      • mochi 在给新连接回 CONNACK 之前就关掉旧连接,旧连接的 OnDisconnect 会在新连接的 SUBSCRIBE 到达之前清标记并唤醒新连接,中间至少隔一个往返,窗口内的推送基本必丢。
    • 另外,已订阅但还没发 hello 的窗口里,msg 会先于 hello 的 resp 到达客户端,也违反"握手完成才开始推送"。
  • 证据:
func (u *appUplink) OnSessionEstablished(ctx context.Context, conn port.ConnInfo) error {
	u.conns.Set(conn.EndpointID, message.LiveConn{
		ConnID:        conn.ConnID,
		MaxPacketSize: conn.MaxPacketSize,
	})
	return nil
}
	if live, ok := a.lookupConn(endpointID); ok && live.ConnID != connID {
		a.WakePush(endpointID)
	}

其余位置:

  • serve.go:334-338:对所有连接推送。

  • broker.go:202-244:发布前不检查订阅。

  • mochi server.go:474-476:OnSessionEstablished 在读循环开始之前调用。

  • mochi server.go:985-1022:没有订阅者时什么都不做。

  • message/session.go:12-27:不清本连接已有的推送标记。

  • message/push.go:326-328:不保留的投递确认超时后按 not_acked 丢弃。

  • 文档依据:DEVELOPMENT 6.1「握手完成才算在线,才开始推送」;7.5「每个已握手连接一个推送循环」「握手完成时…然后开始推送」;PRD F10/D4(短暂断线能送到)、D19。

  • 为何不是故意设计:7.4 里"有连接(含握手中)"只是分发时计算 expire_at 用的,不代表可以推送;DEVIATIONS 里没有相关记录。

  • 解决方案:

    1. broker(本区):PublishDown 发布前检查目标连接已订阅 down(复用 hasDownSub(st)),没订阅就返回 ErrNotSubscribed(包一层 ErrNoConnection)。message 包现有的发布失败分支会清标记并 1 秒后重推。
    2. serve(本区):messageLoops 只对 brk.IsHandshook(ep) 为真、且 brk.CurrentConnID(ep) 等于 live.ConnID 的连接调用 PushPending。
    3. message(需 M 线配合):
      • LiveConn 加一个 Handshook 字段,由 appUplink.OnHandshakeComplete 置位。
      • WakePush 和 PushPending 在未握手时直接返回。
      • OnHandshakeComplete 推送前先清掉 pushed_conn=本连接 的标记作兜底,SDK 按 from+id 去重。
  • 改动文件:internal/broker/broker.go、cmd/nixmsg/serve.go、cmd/nixmsg/uplink.go;需要 M 线配合改 internal/app/message/conn.go、push.go、session.go。

  • 与其他模块的交互/冲突风险:握手前的 QoS 0 事件(presence、group_event)会改为发布失败,但它们原本就会被丢,行为不变。Session 发 hello 的 resp 和 fatal 之前已经检查过订阅,不受影响。

  • 需补测试:

    • broker 单测:会话已建立、还没订阅时,PublishDown 返回错误。
    • 集成测:bob 有一条待收的不保留消息;用同编号的新 WS 连接顶号,并把 SUBSCRIBE 推迟 500 毫秒。把 ack_timeout_seconds 调到几秒,确认 hello 之后能收到这条消息,没有被标成 dropped。
  • 置信度:代码阅读确定(时序由 mochi 源码推出),未用测试复现

[M-11] WakePush 每次都新起协程、不合并:同一端并发推送,可能超出推送窗口、打乱顺序

  • 严重级:medium
  • 分类:并发
  • 现象与影响:
    • 同一个编号可能同时有多个 PushPending 在跑:每秒循环、每次唤醒各一个。
    • 各自读出已推数、分别 claim,推送窗口 32 可能被突破。
    • claim 和 publish 分属两步,可能出现后 claim 的先 publish,打破「按发送时刻排序」。
    • 确认超时可能被重复处理,导致重复重推。
    • 高峰时一个群发加上一轮确认,会产生上千个协程和数千次读查询。
    • 每条投递单独一次写提交,推 32 条就要串行 32 次落盘。
  • 证据:push.go:549-560;push.go:110-118(先读已推数再推);push.go:190-208(逐条 claim);push.go:309-338(超时处理没有重查 pushed_at)。
  • 文档依据:DEVELOPMENT 7.5「每个已握手连接一个推送循环」「窗口默认 32」;PRD F08 排序要求。
  • 为何不是故意设计:DEVIATIONS 里没有相关说明。
  • 解决方案:
    1. 每个就绪连接一个推送 worker,用容量为 1 的信号通道合并唤醒;WakePush 只发信号,每秒循环只负责给各 worker 发信号。
    2. 一次写操作 claim 本轮全部条目(逐条条件更新,返回成功的集合),按顺序发布。
    3. 超时处理在写操作里再核对 pushed_at 仍满足超时条件。
  • 改动文件:internal/app/message/push.go、app.go;cmd/nixmsg/serve.go。
  • 交互/冲突风险:是 M-04 在途标记的基础;worker 在握手时启动、断线时回收。
  • 需补测试:50 个并发唤醒下,任一时刻在途数不超过窗口,同一端到达顺序与 (send_at, seq) 一致。
  • 置信度:代码阅读确定(乱序概率低,放进待核实做压测)。

[M-13] messageLoops 单协程串行、没有超时:任何一个慢操作都会冻结全局

  • 严重级:medium
  • 分类:设计 / 并发
  • 现象与影响:
    • 同一个协程里依次执行:到点分发、对所有在线端逐个 PushPending、清理,以及修复 #5 加入的三条 COUNT 统计。
    • 1000 个在线端时,每秒约 5000 次空转查询。
    • 任一端推送变慢(写队列拥塞、M-01 大帧阻塞、一次大量超时处理),都会拖住或冻结全局的分发、清理和确认超时处理。
  • 证据:serve.go:322-348。
  • 文档依据:DEVELOPMENT 7.4 和 7.5 描述的是调度、每连接推送、清理三类独立循环。
  • 为何不是故意设计:DEVIATIONS L-WIRE 第 1 条、M2/M3/M4 第 4 条只讲接线,没有讨论阻塞和串行化。
  • 解决方案:拆成独立协程:分发(定时器)、到期处理(1 秒)、purge(1 小时)、每端推送 worker(M-11)、指标采样(例如 15 秒)。每次调用都带 context.WithTimeout。
  • 改动文件:cmd/nixmsg/serve.go;internal/app/message/app.go。
  • 交互/冲突风险:采样周期变长,指标延迟会略有增加。
  • 需补测试:注入一个会阻塞的假下行,分发和清理仍然每秒执行。
  • 置信度:代码阅读确定。

[M-12] 到点分发每秒最多 100 条、逐条串行提交,积压没有上限;遇到一条失败会卡住整批

  • 严重级:medium
  • 分类:性能 / 与 PRD 不符
  • 现象与影响:
    • 每个 tick 最多取 100 条,每条单独一次写操作、串行等待提交(2 毫秒合批窗口加一次落盘),实际还不到每秒 100 条。
    • 默认延迟(反悔窗口)场景下,按 PRD 每秒 200 条提交,积压每秒增长 100 条以上,「10 秒延迟」会变成几分钟,并越拖越长。
    • 轮询粒度 1 秒,而不是按最早的发送时刻定时唤醒。
    • 任一条分发返回错误,整批立即中止;下一轮这条仍排在最前面,可能永久卡住所有定时消息。
    • 启动恢复在开始监听之前同步分发最多 1000 条;大群定时消息多时,服务要几分钟后才开始接受连接。
  • 证据:serve.go:331;push.go:29-34、55-64;recover.go:37。
  • 文档依据:DEVELOPMENT 7.4「调度循环按最早的 send_at 定时唤醒,提交新消息时也唤醒」;PRD F08「到了发送时刻……立刻推送」;PRD 第 8 节吞吐要求。
  • 为何不是故意设计:DEVIATIONS M2/M3/M4 第 4 条只讲循环由谁启动。
  • 解决方案:
    1. 分发循环独立出来,按最早的 send_at 设定时器,提交定时消息时通知它。
    2. 每轮循环到取空或时间预算用完为止;有限并发(例如 8 路)提交写操作,让写队列合批(排序由推送查询保证,不受影响)。
    3. 单条失败只记日志并跳过,其余继续。
    4. 启动恢复只做 SQL 修正,分发交给循环。
    5. 顺带优化:插入投递用集合语句,接收端配额计数加上限,回执插入不再重复查消息和端。
  • 改动文件:internal/app/message/push.go、dispatch.go、recover.go;cmd/nixmsg/serve.go。
  • 交互/冲突风险:与 M-13 一起重构。
  • 需补测试:250 条同时到点在 1 秒内分发完;构造一条会失败的消息,不影响其余消息;定时误差不超过 100 毫秒。
  • 置信度:代码阅读确定。

[M-15] claim 已提交却没发布时,只能靠 5 分钟确认超时兜底,不保留消息被误判为没确认

  • 严重级:low(M-01 未修时会升到 medium)
  • 分类:并发 / 逻辑
  • 现象与影响:
    • 写队列在调用方 ctx 过期时,即使操作已经提交也会返回 ctx.Err()。PushPending 会当成 claim 失败直接返回,不发布。
    • 发布失败后的 clearPushed 复用已过期的 ctx,写队列会直接拒绝执行。
    • 两种情况都会把「已推送」标记留在一条没发出去的投递上,5 分钟后不保留的消息被判 dropped/not_acked。
  • 证据:store/queue.go:91-104(成功也返回 ctx 错误)、queue.go:71-73(ctx 已过期直接拒绝);push.go:191-205、226-227。
  • 文档依据:DEVELOPMENT 7.5 第 4 步「发布失败……清掉,1 秒后再推」。
  • 为何不是故意设计:DEVIATIONS P2 第 2 条只说明写队列本身的等待语义,调用方没有据此处理。
  • 解决方案:clearPushed 固定用新建的短超时 ctx;claim 返回 ctx 错误时,读一次 pushed_conn,如果已被本次写成当前连接就清掉并安排 1 秒后重推。
  • 改动文件:internal/app/message/push.go。
  • 交互/冲突风险:无。
  • 需补测试:假下行返回 ctx 错误时,pushed_conn 被清空并在 1 秒后重推。
  • 置信度:代码阅读确定。

[M-02] message 包自己那份大帧名额在断线、撤回、作废时不释放,同一投递重推还会重复占用

  • 严重级:high
  • 分类:并发 / 逻辑
  • 现象与影响:
    • 每次推送大帧都会 acquireLarge,再用 trackLarge 把 largeHeld[seq:编号] 置为 true。这是布尔值,不计数。
    • 释放只发生在四种情况:确认结果为 accepted、确认超时、发布失败、OnPublishDropped。
    • 以下路径会永久漏掉名额:
      1. 断线:OnDisconnect 只清推送标记。重连后同一投递重推又占一个名额,确认时只释放一个。每次「大帧在途时断线」漏 1 个。
      2. 撤回:Recall 不释放;之后的 ack 结果是 recalled,而 Ack 只在 accepted 时释放。
      3. 作废:停用、删除、退群、解散由 identity/group 包直接改库,碰不到 message 的内存表。
    • 漏满 64 个后,全进程所有超过 64 KiB 的消息都推不出去,而每秒循环还在反复读正文、序列化、claim 再 clear。弱网下断线频繁,很快就会漏满。
    • 名额要等应用层确认才释放(最长 5 分钟),比文档要求的「PUBACK 即释放」持有得更久,吞吐也被压低。
  • 证据:
		large := len(payload) > largeFrameBytes
		if large {
			if !a.acquireLarge(ctx) {
				_ = a.clearPushed(ctx, it.seq, endpointID, connID, nowMs)
				continue
			}
			a.trackLarge(it.seq, endpointID, true)
		}
  • push.go:508-528:trackLarge 是布尔,releaseLarge 只删一次。
  • ack.go:71-73:只在 accepted 时释放。
  • ack.go:139-162:撤回不释放。
  • session.go:31-75:断线不释放。
  • 文档依据:DEVELOPMENT 7.5,名额在「PUBACK、确认超时或连接断开时释放」。
  • 为何不是故意设计:DEVIATIONS M2/M3/M4 第 2 条只说明「为假下行测试另管一份,确认、超时、清标记时释放」。断线/撤回/作废不释放、重复占用都不在说明范围内。
  • 解决方案:
    • 推荐:删掉 message 内的 largeSem、largeHeld、acquireLarge、trackLarge、releaseLarge,只依赖修好后的 broker 名额(M-01)。测试用的假下行如需限流,在假实现里模拟。同步改 DEVIATIONS M2/M3/M4 第 2 条。
    • 备选(必须保留时):key 改为 (connID, seq, 编号);同一 key 已持有就不再获取;OnDisconnect 释放该连接名下所有 key;Ack 不论结果都释放;撤回和 finishDeliveryTx 对已推送的投递释放;每分钟对账一次,投递不再是「pending 且推给该连接」就释放,以兜住 identity/group 的作废路径。
  • 改动文件:internal/app/message/push.go、ack.go、session.go、app.go。
  • 交互/冲突风险:推荐方案依赖 M-01 先修好。
  • 需补测试:
    • 大帧推送、断线、重推、确认之后,len(a.largeSem)==0。
    • 大帧推送后撤回再 ack,名额归零。
    • 70 轮「推送+断线」后,大帧仍能推送。
  • 置信度:代码阅读确定。

复审基线:main 4059a15(2026-09-30)。编号说明、各工作线的合并顺序、共享文件归属见总览 #7。

**编号**:C-01 **严重级**:high **工作线**:消息核心(internal/app/message、serve 的 messageLoops) **来源**:审查 M-03、P-02、M-11、M-13、M-12、M-15、M-02 **依赖**:B-02 (#9)(第 5 点删除消息包大帧名额需在其后) **被依赖**:C-02 (#33)、C-03 (#34)、D-03 (#29)、L-03 (#22) ### 结论与统一方案 本 issue 合并审查 M-02、M-03、M-11、M-12、M-13、M-15 与 P-02、P-07(循环部分,原文附在 B-03)。它们都落在 `internal/app/message/push.go` 与 `cmd/nixmsg/serve.go` 的 `messageLoops`,必须由一人按一个设计改完,不能拆给多人。统一设计: 1. **每连接推送 worker**:`OnHandshakeComplete` 为该连接启动 worker(带容量为 1 的唤醒通道,合并唤醒),`OnDisconnect` 停止并回收。`WakePush` 与每秒循环只发信号。未握手、代号不符或已断开的连接不推送、不 claim,修复握手前推送(审查 M-03、P-02)。`OnHandshakeComplete` 推送前先清掉 `pushed_conn=本连接` 的旧标记作兜底。 2. **窗口与顺序**:每轮在一个写操作里 claim 本轮全部条目(逐条条件更新,返回成功集合),按 `(send_at, seq)` 顺序发布;确认超时处理在写操作里再核对 `pushed_at`。(审查 M-11) 3. **循环拆分**:`messageLoops` 拆为独立 goroutine,由 WaitGroup 管理(供 L-03 停机等待): - 到点分发:按最早 `send_at` 设定时器,提交定时消息时通知;每轮循环到取空或用完时间预算;有限并发提交写操作让写队列合批;单条失败只记日志并跳过。(审查 M-12) - 到期处理:每秒一次。 - 清理:每小时一次(内容见 C-03)。 - 指标采样:10–15 秒一次。 - 每次调用都带 `context.WithTimeout`。启动恢复只做 SQL 修正,分发交给循环,不再阻塞开始监听。(审查 M-13、P-07) 4. **上下文**:claim 返回 ctx 错误时读一次 `pushed_conn` 判断是否已写入并清掉;`clearPushed` 固定用新建的短超时 ctx。(审查 M-15) 5. **大帧名额**:删除 message 包内的 `largeSem`、`largeHeld`、`acquireLarge`、`trackLarge`、`releaseLarge`,只依赖 broker 的名额(B-02 合入后做);测试用假下行如需限流,在假实现里模拟。同步更新 DEVIATIONS M2/M3/M4 第 2 条。(审查 M-02) 6. **与 broker 的契约**:`PublishDown` 返回 `ErrBackpressure`(B-03)、`ErrNotSubscribed`、`ErrNoConnection`(B-06)时一律按发布失败处理(清标记、1 秒后重推),回执与 revoked 路径同样处理。 ### 改动文件 `internal/app/message/push.go`、`app.go`、`conn.go`、`session.go`、`recover.go`(`RecoverOnStart`);`cmd/nixmsg/serve.go`(`messageLoops` 改为调用 message 提供的启动与等待接口)。 ### 与其他问题的交互 / 冲突说明 - 第 5 点依赖 B-02;与 B-03、B-06 以 `PublishDown` 签名不变为契约,可并行开发。 - C-02、C-03、D-03、L-03 依赖本条。 - `cmd/nixmsg/uplink.go` 的生命周期函数由 B-09 修改;本条只通过 `msg.OnHandshakeComplete` / `msg.OnDisconnect` 启停 worker,不改 uplink.go。 - serve.go 只改 `messageLoops`。 ### 验收与测试 - 未就绪时提交消息再唤醒:没有 msg 发布、`pushed_conn` 为空;握手后才发布。 - 裸 MQTT 客户端 CONNECT 后延迟 1.5 秒再 SUBSCRIBE 与 hello,仍能在 2 秒内收到排队中的不保留消息(没有被标成 dropped)。 - 50 个并发唤醒下,任一时刻在途数不超过窗口,同一端到达顺序与 `(send_at, seq)` 一致。 - 注入会阻塞的假下行:其他端的推送与定时分发照常。 - 250 条同时到点的消息 1 秒内分发完;构造一条会失败的消息,不影响其余。 - 大帧"推送+断线+重推"70 轮后大帧仍能推送。 --- ### 问题明细(各区审查原文,证据含文件与行号) > 以下是本次复审各区审查报告的原文段落。A、M、I、P、S 开头的是原始发现编号(A 管理后台与网页、M 消息核心、I 身份认证群在线、P 传输平台部署、S SDK)。**解决方案以本 issue 上方的"结论与统一方案"为准**;原文里的方案与之不一致时,按上方执行。 #### [M-03] 握手前就推送:订阅前推出的帧静默丢失,推送窗口被占 5 分钟,还绕过了 max_receive_bytes - 严重级:high - 分类:协议一致性 / 逻辑 - 现象与影响: - 连接一建立(CONNACK 之后,SUBSCRIBE 和 hello 之前),`OnSessionEstablished` 就把它登记进连接表。`messageLoops` 和各处 `WakePush` 都会对它调用 `PushPending`,而 `PushPending` 不看是否已握手。 - Clean Start 时 mochi 已经删掉旧订阅,这时向 down 主题发布没有订阅者,`Publish` 返回 nil,帧静默丢失。投递却已经记成「已推给这个连接」: - 占用推送窗口,直到 5 分钟确认超时。 - 不保留的消息 5 分钟后被判 `dropped/not_acked`,而接收端其实一直在线。 - 如果丢的正好是整窗 32 条,这个端 5 分钟内收不到任何消息。 - 顶号重连更容易命中:旧连接的 `OnDisconnect` 发现已有新连接时会立即 `WakePush`,几乎一定早于新连接的 SUBSCRIBE。 - hello 之前 `MaxReceiveBytes=0`,超限的大帧会直接发给只声明了 1024 字节的小设备,而不是拒收并回 `too_large`。 - 在弱网(2 秒延迟)或服务器重启后大批重连时,几乎每次都会命中。 - 证据: ```30:36:e:\code\NixMsg\cmd\nixmsg\uplink.go func (u *appUplink) OnSessionEstablished(ctx context.Context, conn port.ConnInfo) error { u.conns.Set(conn.EndpointID, message.LiveConn{ ConnID: conn.ConnID, MaxPacketSize: conn.MaxPacketSize, }) return nil } ``` - `push.go:88-99`:只比对连接代号。 - `push.go:549-560`:`WakePush` 同样不看握手状态。 - `session.go:71-73`:旧连接断开时唤醒新连接。 - mochi `server.go:457-476`:先继承会话,再发 CONNACK,再调 `OnSessionEstablished`,之后才开始读包。 - mochi `server.go:1001-1021`:没有订阅者时什么都不做,投递失败只记 debug 日志。 - 文档依据: - DEVELOPMENT 6.1「握手完成才算在线,才开始推送」;7.5「每个已握手连接一个推送循环」。 - PRD F07:接收上限在握手时声明;F08:「对方已完成登录握手:立刻推送」。 - 为何不是故意设计:7.4 让分发时把「握手中」的连接算作在线,只用于计算 `expire_at`;推送必须等握手完成。DEVIATIONS 里没有相关说明。 - 解决方案: 1. `LiveConn` 加 `Ready bool`。`OnSessionEstablished` 登记为 `Ready:false`,`OnHandshakeComplete` 登记为 `Ready:true` 并带上 `MaxReceiveBytes`。 2. `PushPending` 遇到当前连接不存在、代号不匹配或未就绪时直接返回:不 claim,也不发回执。`WakePush` 和 `messageLoops` 同样只处理就绪连接。 3. `dispatchFullTx` 的在线判定保持现状(含握手中),符合 7.4。 4. 可选的纵深防御(连接线):`PublishDown` 在没有 down 订阅时返回错误,让 message 走清标记重推。 - 改动文件:`internal/app/message/conn.go`、`push.go`;`cmd/nixmsg/uplink.go`、`serve.go`;可选 `internal/broker/broker.go`。 - 交互/冲突风险:`presenceConnTable.IsOnline` 用的是同一张表,建议一起改成只认就绪连接,以符合「在线 = 完成握手」(DEVELOPMENT 第 1 节)。 - 需补测试: - 单测:未就绪时提交消息再调 `PushPending`,没有 msg 发布、`pushed_conn` 为空;置为就绪后才发布。 - 集成:裸 MQTT 客户端 CONNECT 后等 1.5 秒再 SUBSCRIBE 和 hello,2 秒内应收到排队中的不保留消息。 - 集成:hello 声明 1024 字节的连接,对排队中的 2048 字节消息得到 `too_large` 回执。 - 置信度:代码阅读确定(已对照 mochi 调用顺序)。 #### [P-02] 握手完成前就推送:订阅 down 之前推出的消息被 mochi 静默丢弃,却记成"已推送",不保留的消息 5 分钟后按 not_acked 丢弃 - **严重级**:high - **分类**:逻辑 / 协议一致性 / 数据 - **现象与影响**: - `appUplink.OnSessionEstablished` 在回 CONNACK 之后、客户端发 SUBSCRIBE 之前,就把这个连接登记进了 `memConns`。 - 两条推送路径会命中这个窗口: - `messageLoops` 每秒对 `conns.Snapshot()` 里的每个连接调用 `PushPending`。 - message 的 `OnDisconnect` 在顶号时会立刻 `WakePush` 新连接。 - 这时新连接还没订阅 `nix/c/{id}/down`: - `server.Publish` 找不到订阅者,直接返回 nil(也不触发 OnPublishDropped),`PublishDown` 于是报成功。 - 投递已经被条件更新成 `pushed_conn=新连接`。 - 握手完成时,`message.OnHandshakeComplete` 只推 `pushed_conn IS NULL` 的投递,这批"推进空气"的投递要等 5 分钟确认超时:不保留的按 D19 改成 `dropped/not_acked`,保留的晚到 5 分钟。 - 最常见的触发场景是顶号重连:移动端断网后用令牌重连,服务器还没发现旧连接已死。 - mochi 在给新连接回 CONNACK 之前就关掉旧连接,旧连接的 OnDisconnect 会在新连接的 SUBSCRIBE 到达之前清标记并唤醒新连接,中间至少隔一个往返,窗口内的推送基本必丢。 - 另外,已订阅但还没发 hello 的窗口里,msg 会先于 hello 的 resp 到达客户端,也违反"握手完成才开始推送"。 - **证据**: ```30:36:e:\code\NixMsg\cmd\nixmsg\uplink.go func (u *appUplink) OnSessionEstablished(ctx context.Context, conn port.ConnInfo) error { u.conns.Set(conn.EndpointID, message.LiveConn{ ConnID: conn.ConnID, MaxPacketSize: conn.MaxPacketSize, }) return nil } ``` ```71:73:e:\code\NixMsg\internal\app\message\session.go if live, ok := a.lookupConn(endpointID); ok && live.ConnID != connID { a.WakePush(endpointID) } ``` 其余位置: - `serve.go:334-338`:对所有连接推送。 - `broker.go:202-244`:发布前不检查订阅。 - mochi `server.go:474-476`:OnSessionEstablished 在读循环开始之前调用。 - mochi `server.go:985-1022`:没有订阅者时什么都不做。 - `message/session.go:12-27`:不清本连接已有的推送标记。 - `message/push.go:326-328`:不保留的投递确认超时后按 not_acked 丢弃。 - **文档依据**:DEVELOPMENT 6.1「握手完成才算在线,才开始推送」;7.5「每个已握手连接一个推送循环」「握手完成时…然后开始推送」;PRD F10/D4(短暂断线能送到)、D19。 - **为何不是故意设计**:7.4 里"有连接(含握手中)"只是分发时计算 `expire_at` 用的,不代表可以推送;DEVIATIONS 里没有相关记录。 - **解决方案**: 1. broker(本区):`PublishDown` 发布前检查目标连接已订阅 down(复用 `hasDownSub(st)`),没订阅就返回 `ErrNotSubscribed`(包一层 `ErrNoConnection`)。message 包现有的发布失败分支会清标记并 1 秒后重推。 2. serve(本区):`messageLoops` 只对 `brk.IsHandshook(ep)` 为真、且 `brk.CurrentConnID(ep)` 等于 `live.ConnID` 的连接调用 `PushPending`。 3. message(需 M 线配合): - `LiveConn` 加一个 `Handshook` 字段,由 `appUplink.OnHandshakeComplete` 置位。 - `WakePush` 和 `PushPending` 在未握手时直接返回。 - `OnHandshakeComplete` 推送前先清掉 `pushed_conn=本连接` 的标记作兜底,SDK 按 from+id 去重。 - **改动文件**:`internal/broker/broker.go`、`cmd/nixmsg/serve.go`、`cmd/nixmsg/uplink.go`;需要 M 线配合改 `internal/app/message/conn.go`、`push.go`、`session.go`。 - **与其他模块的交互/冲突风险**:握手前的 QoS 0 事件(presence、group_event)会改为发布失败,但它们原本就会被丢,行为不变。Session 发 hello 的 resp 和 fatal 之前已经检查过订阅,不受影响。 - **需补测试**: - broker 单测:会话已建立、还没订阅时,`PublishDown` 返回错误。 - 集成测:bob 有一条待收的不保留消息;用同编号的新 WS 连接顶号,并把 SUBSCRIBE 推迟 500 毫秒。把 `ack_timeout_seconds` 调到几秒,确认 hello 之后能收到这条消息,没有被标成 dropped。 - **置信度**:代码阅读确定(时序由 mochi 源码推出),未用测试复现 #### [M-11] WakePush 每次都新起协程、不合并:同一端并发推送,可能超出推送窗口、打乱顺序 - 严重级:medium - 分类:并发 - 现象与影响: - 同一个编号可能同时有多个 `PushPending` 在跑:每秒循环、每次唤醒各一个。 - 各自读出已推数、分别 claim,推送窗口 32 可能被突破。 - claim 和 publish 分属两步,可能出现后 claim 的先 publish,打破「按发送时刻排序」。 - 确认超时可能被重复处理,导致重复重推。 - 高峰时一个群发加上一轮确认,会产生上千个协程和数千次读查询。 - 每条投递单独一次写提交,推 32 条就要串行 32 次落盘。 - 证据:`push.go:549-560`;`push.go:110-118`(先读已推数再推);`push.go:190-208`(逐条 claim);`push.go:309-338`(超时处理没有重查 `pushed_at`)。 - 文档依据:DEVELOPMENT 7.5「每个已握手连接一个推送循环」「窗口默认 32」;PRD F08 排序要求。 - 为何不是故意设计:DEVIATIONS 里没有相关说明。 - 解决方案: 1. 每个就绪连接一个推送 worker,用容量为 1 的信号通道合并唤醒;`WakePush` 只发信号,每秒循环只负责给各 worker 发信号。 2. 一次写操作 claim 本轮全部条目(逐条条件更新,返回成功的集合),按顺序发布。 3. 超时处理在写操作里再核对 `pushed_at` 仍满足超时条件。 - 改动文件:`internal/app/message/push.go`、`app.go`;`cmd/nixmsg/serve.go`。 - 交互/冲突风险:是 M-04 在途标记的基础;worker 在握手时启动、断线时回收。 - 需补测试:50 个并发唤醒下,任一时刻在途数不超过窗口,同一端到达顺序与 `(send_at, seq)` 一致。 - 置信度:代码阅读确定(乱序概率低,放进待核实做压测)。 #### [M-13] messageLoops 单协程串行、没有超时:任何一个慢操作都会冻结全局 - 严重级:medium - 分类:设计 / 并发 - 现象与影响: - 同一个协程里依次执行:到点分发、对所有在线端逐个 `PushPending`、清理,以及修复 #5 加入的三条 COUNT 统计。 - 1000 个在线端时,每秒约 5000 次空转查询。 - 任一端推送变慢(写队列拥塞、M-01 大帧阻塞、一次大量超时处理),都会拖住或冻结全局的分发、清理和确认超时处理。 - 证据:`serve.go:322-348`。 - 文档依据:DEVELOPMENT 7.4 和 7.5 描述的是调度、每连接推送、清理三类独立循环。 - 为何不是故意设计:DEVIATIONS L-WIRE 第 1 条、M2/M3/M4 第 4 条只讲接线,没有讨论阻塞和串行化。 - 解决方案:拆成独立协程:分发(定时器)、到期处理(1 秒)、purge(1 小时)、每端推送 worker(M-11)、指标采样(例如 15 秒)。每次调用都带 `context.WithTimeout`。 - 改动文件:`cmd/nixmsg/serve.go`;`internal/app/message/app.go`。 - 交互/冲突风险:采样周期变长,指标延迟会略有增加。 - 需补测试:注入一个会阻塞的假下行,分发和清理仍然每秒执行。 - 置信度:代码阅读确定。 #### [M-12] 到点分发每秒最多 100 条、逐条串行提交,积压没有上限;遇到一条失败会卡住整批 - 严重级:medium - 分类:性能 / 与 PRD 不符 - 现象与影响: - 每个 tick 最多取 100 条,每条单独一次写操作、串行等待提交(2 毫秒合批窗口加一次落盘),实际还不到每秒 100 条。 - 默认延迟(反悔窗口)场景下,按 PRD 每秒 200 条提交,积压每秒增长 100 条以上,「10 秒延迟」会变成几分钟,并越拖越长。 - 轮询粒度 1 秒,而不是按最早的发送时刻定时唤醒。 - 任一条分发返回错误,整批立即中止;下一轮这条仍排在最前面,可能永久卡住所有定时消息。 - 启动恢复在开始监听之前同步分发最多 1000 条;大群定时消息多时,服务要几分钟后才开始接受连接。 - 证据:`serve.go:331`;`push.go:29-34`、`55-64`;`recover.go:37`。 - 文档依据:DEVELOPMENT 7.4「调度循环按最早的 send_at 定时唤醒,提交新消息时也唤醒」;PRD F08「到了发送时刻……立刻推送」;PRD 第 8 节吞吐要求。 - 为何不是故意设计:DEVIATIONS M2/M3/M4 第 4 条只讲循环由谁启动。 - 解决方案: 1. 分发循环独立出来,按最早的 `send_at` 设定时器,提交定时消息时通知它。 2. 每轮循环到取空或时间预算用完为止;有限并发(例如 8 路)提交写操作,让写队列合批(排序由推送查询保证,不受影响)。 3. 单条失败只记日志并跳过,其余继续。 4. 启动恢复只做 SQL 修正,分发交给循环。 5. 顺带优化:插入投递用集合语句,接收端配额计数加上限,回执插入不再重复查消息和端。 - 改动文件:`internal/app/message/push.go`、`dispatch.go`、`recover.go`;`cmd/nixmsg/serve.go`。 - 交互/冲突风险:与 M-13 一起重构。 - 需补测试:250 条同时到点在 1 秒内分发完;构造一条会失败的消息,不影响其余消息;定时误差不超过 100 毫秒。 - 置信度:代码阅读确定。 #### [M-15] claim 已提交却没发布时,只能靠 5 分钟确认超时兜底,不保留消息被误判为没确认 - 严重级:low(M-01 未修时会升到 medium) - 分类:并发 / 逻辑 - 现象与影响: - 写队列在调用方 ctx 过期时,即使操作已经提交也会返回 `ctx.Err()`。`PushPending` 会当成 claim 失败直接返回,不发布。 - 发布失败后的 `clearPushed` 复用已过期的 ctx,写队列会直接拒绝执行。 - 两种情况都会把「已推送」标记留在一条没发出去的投递上,5 分钟后不保留的消息被判 `dropped/not_acked`。 - 证据:`store/queue.go:91-104`(成功也返回 ctx 错误)、`queue.go:71-73`(ctx 已过期直接拒绝);`push.go:191-205`、`226-227`。 - 文档依据:DEVELOPMENT 7.5 第 4 步「发布失败……清掉,1 秒后再推」。 - 为何不是故意设计:DEVIATIONS P2 第 2 条只说明写队列本身的等待语义,调用方没有据此处理。 - 解决方案:`clearPushed` 固定用新建的短超时 ctx;claim 返回 ctx 错误时,读一次 `pushed_conn`,如果已被本次写成当前连接就清掉并安排 1 秒后重推。 - 改动文件:`internal/app/message/push.go`。 - 交互/冲突风险:无。 - 需补测试:假下行返回 ctx 错误时,`pushed_conn` 被清空并在 1 秒后重推。 - 置信度:代码阅读确定。 #### [M-02] message 包自己那份大帧名额在断线、撤回、作废时不释放,同一投递重推还会重复占用 - 严重级:high - 分类:并发 / 逻辑 - 现象与影响: - 每次推送大帧都会 `acquireLarge`,再用 `trackLarge` 把 `largeHeld[seq:编号]` 置为 true。这是布尔值,不计数。 - 释放只发生在四种情况:确认结果为 accepted、确认超时、发布失败、`OnPublishDropped`。 - 以下路径会永久漏掉名额: 1. **断线**:`OnDisconnect` 只清推送标记。重连后同一投递重推又占一个名额,确认时只释放一个。每次「大帧在途时断线」漏 1 个。 2. **撤回**:`Recall` 不释放;之后的 ack 结果是 recalled,而 `Ack` 只在 accepted 时释放。 3. **作废**:停用、删除、退群、解散由 identity/group 包直接改库,碰不到 message 的内存表。 - 漏满 64 个后,全进程所有超过 64 KiB 的消息都推不出去,而每秒循环还在反复读正文、序列化、claim 再 clear。弱网下断线频繁,很快就会漏满。 - 名额要等应用层确认才释放(最长 5 分钟),比文档要求的「PUBACK 即释放」持有得更久,吞吐也被压低。 - 证据: ```210:217:e:\code\NixMsg\internal\app\message\push.go large := len(payload) > largeFrameBytes if large { if !a.acquireLarge(ctx) { _ = a.clearPushed(ctx, it.seq, endpointID, connID, nowMs) continue } a.trackLarge(it.seq, endpointID, true) } ``` - `push.go:508-528`:`trackLarge` 是布尔,`releaseLarge` 只删一次。 - `ack.go:71-73`:只在 accepted 时释放。 - `ack.go:139-162`:撤回不释放。 - `session.go:31-75`:断线不释放。 - 文档依据:DEVELOPMENT 7.5,名额在「PUBACK、确认超时或连接断开时释放」。 - 为何不是故意设计:DEVIATIONS M2/M3/M4 第 2 条只说明「为假下行测试另管一份,确认、超时、清标记时释放」。断线/撤回/作废不释放、重复占用都不在说明范围内。 - 解决方案: - 推荐:删掉 message 内的 `largeSem`、`largeHeld`、`acquireLarge`、`trackLarge`、`releaseLarge`,只依赖修好后的 broker 名额(M-01)。测试用的假下行如需限流,在假实现里模拟。同步改 DEVIATIONS M2/M3/M4 第 2 条。 - 备选(必须保留时):key 改为 `(connID, seq, 编号)`;同一 key 已持有就不再获取;`OnDisconnect` 释放该连接名下所有 key;`Ack` 不论结果都释放;撤回和 `finishDeliveryTx` 对已推送的投递释放;每分钟对账一次,投递不再是「pending 且推给该连接」就释放,以兜住 identity/group 的作废路径。 - 改动文件:`internal/app/message/push.go`、`ack.go`、`session.go`、`app.go`。 - 交互/冲突风险:推荐方案依赖 M-01 先修好。 - 需补测试: - 大帧推送、断线、重推、确认之后,`len(a.largeSem)==0`。 - 大帧推送后撤回再 ack,名额归零。 - 70 轮「推送+断线」后,大帧仍能推送。 - 置信度:代码阅读确定。 --- <sub>复审基线:main `4059a15`(2026-09-30)。编号说明、各工作线的合并顺序、共享文件归属见总览 #7。</sub>
nixevol added the P1-highlane/messagereview-2026-09-30 labels 2026-09-30 13:56:58 +08:00
Author
Owner

C-01 已在 feat/fix-message-c01-c03 提交 eca4c836f1 (#32)

仅向已握手连接推送;StartLoops / WaitLoops 拆开调度循环。停机等待交给 L-03。未合入 main。

C-01 已在 `feat/fix-message-c01-c03` 提交 https://git.asio.asia/nixevol/NixMsg/commit/eca4c836f11eddd72d28ff44a97245ef4cf49f12 (#32) 仅向已握手连接推送;`StartLoops` / `WaitLoops` 拆开调度循环。停机等待交给 L-03。未合入 main。
Author
Owner

已合入 origin/main 0c9b459。落地提交 83521f4 fix: 仅向已握手连接推送并拆分调度循环 (#32)。合入时 8f2ebc7 补唤醒回执推送。

已合入 origin/main `0c9b459`。落地提交 `83521f4` fix: 仅向已握手连接推送并拆分调度循环 (#32)。合入时 `8f2ebc7` 补唤醒回执推送。
Sign in to join this conversation.