编号: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,必须由一人按一个设计改完,不能拆给多人。统一设计:
internal/app/message/push.go
cmd/nixmsg/serve.go
messageLoops
OnHandshakeComplete
OnDisconnect
WakePush
pushed_conn=本连接
(send_at, seq)
pushed_at
send_at
context.WithTimeout
pushed_conn
clearPushed
largeSem
largeHeld
acquireLarge
trackLarge
releaseLarge
PublishDown
ErrBackpressure
ErrNotSubscribed
ErrNoConnection
internal/app/message/push.go、app.go、conn.go、session.go、recover.go(RecoverOnStart);cmd/nixmsg/serve.go(messageLoops 改为调用 message 提供的启动与等待接口)。
app.go
conn.go
session.go
recover.go
RecoverOnStart
cmd/nixmsg/uplink.go
msg.OnHandshakeComplete
msg.OnDisconnect
以下是本次复审各区审查报告的原文段落。A、M、I、P、S 开头的是原始发现编号(A 管理后台与网页、M 消息核心、I 身份认证群在线、P 传输平台部署、S SDK)。解决方案以本 issue 上方的"结论与统一方案"为准;原文里的方案与之不一致时,按上方执行。
OnSessionEstablished
PushPending
Publish
dropped/not_acked
MaxReceiveBytes=0
too_large
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
session.go:71-73
server.go:457-476
server.go:1001-1021
expire_at
LiveConn
Ready bool
Ready:false
Ready:true
MaxReceiveBytes
dispatchFullTx
internal/app/message/conn.go
push.go
serve.go
internal/broker/broker.go
presenceConnTable.IsOnline
appUplink.OnSessionEstablished
memConns
conns.Snapshot()
nix/c/{id}/down
server.Publish
pushed_conn=新连接
message.OnHandshakeComplete
pushed_conn IS NULL
if live, ok := a.lookupConn(endpointID); ok && live.ConnID != connID { a.WakePush(endpointID) }
其余位置:
serve.go:334-338:对所有连接推送。
serve.go:334-338
broker.go:202-244:发布前不检查订阅。
broker.go:202-244
mochi server.go:474-476:OnSessionEstablished 在读循环开始之前调用。
server.go:474-476
mochi server.go:985-1022:没有订阅者时什么都不做。
server.go:985-1022
message/session.go:12-27:不清本连接已有的推送标记。
message/session.go:12-27
message/push.go:326-328:不保留的投递确认超时后按 not_acked 丢弃。
message/push.go:326-328
文档依据:DEVELOPMENT 6.1「握手完成才算在线,才开始推送」;7.5「每个已握手连接一个推送循环」「握手完成时…然后开始推送」;PRD F10/D4(短暂断线能送到)、D19。
为何不是故意设计:7.4 里"有连接(含握手中)"只是分发时计算 expire_at 用的,不代表可以推送;DEVIATIONS 里没有相关记录。
解决方案:
hasDownSub(st)
brk.IsHandshook(ep)
brk.CurrentConnID(ep)
live.ConnID
Handshook
appUplink.OnHandshakeComplete
改动文件: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 之前已经检查过订阅,不受影响。
需补测试:
ack_timeout_seconds
置信度:代码阅读确定(时序由 mochi 源码推出),未用测试复现
push.go:110-118
push.go:190-208
push.go:309-338
serve.go:322-348
internal/app/message/app.go
serve.go:331
push.go:29-34
55-64
recover.go:37
dispatch.go
ctx.Err()
store/queue.go:91-104
queue.go:71-73
push.go:191-205
226-227
largeHeld[seq:编号]
OnPublishDropped
Recall
Ack
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
ack.go:71-73
ack.go:139-162
session.go:31-75
(connID, seq, 编号)
finishDeliveryTx
ack.go
len(a.largeSem)==0
复审基线:main 4059a15(2026-09-30)。编号说明、各工作线的合并顺序、共享文件归属见总览 #7。
4059a15
C-01 已在 feat/fix-message-c01-c03 提交 eca4c836f1 (#32)
feat/fix-message-c01-c03
eca4c836f1
仅向已握手连接推送;StartLoops / WaitLoops 拆开调度循环。停机等待交给 L-03。未合入 main。
StartLoops
WaitLoops
已合入 origin/main 0c9b459。落地提交 83521f4 fix: 仅向已握手连接推送并拆分调度循环 (#32)。合入时 8f2ebc7 补唤醒回执推送。
0c9b459
83521f4
8f2ebc7
No dependencies set.
The note is not visible to the blocked user.
编号: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,必须由一人按一个设计改完,不能拆给多人。统一设计:OnHandshakeComplete为该连接启动 worker(带容量为 1 的唤醒通道,合并唤醒),OnDisconnect停止并回收。WakePush与每秒循环只发信号。未握手、代号不符或已断开的连接不推送、不 claim,修复握手前推送(审查 M-03、P-02)。OnHandshakeComplete推送前先清掉pushed_conn=本连接的旧标记作兜底。(send_at, seq)顺序发布;确认超时处理在写操作里再核对pushed_at。(审查 M-11)messageLoops拆为独立 goroutine,由 WaitGroup 管理(供 L-03 停机等待):send_at设定时器,提交定时消息时通知;每轮循环到取空或用完时间预算;有限并发提交写操作让写队列合批;单条失败只记日志并跳过。(审查 M-12)context.WithTimeout。启动恢复只做 SQL 修正,分发交给循环,不再阻塞开始监听。(审查 M-13、P-07)pushed_conn判断是否已写入并清掉;clearPushed固定用新建的短超时 ctx。(审查 M-15)largeSem、largeHeld、acquireLarge、trackLarge、releaseLarge,只依赖 broker 的名额(B-02 合入后做);测试用假下行如需限流,在假实现里模拟。同步更新 DEVIATIONS M2/M3/M4 第 2 条。(审查 M-02)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 提供的启动与等待接口)。与其他问题的交互 / 冲突说明
PublishDown签名不变为契约,可并行开发。cmd/nixmsg/uplink.go的生命周期函数由 B-09 修改;本条只通过msg.OnHandshakeComplete/msg.OnDisconnect启停 worker,不改 uplink.go。messageLoops。验收与测试
pushed_conn为空;握手后才发布。(send_at, seq)一致。问题明细(各区审查原文,证据含文件与行号)
[M-03] 握手前就推送:订阅前推出的帧静默丢失,推送窗口被占 5 分钟,还绕过了 max_receive_bytes
OnSessionEstablished就把它登记进连接表。messageLoops和各处WakePush都会对它调用PushPending,而PushPending不看是否已握手。Publish返回 nil,帧静默丢失。投递却已经记成「已推给这个连接」:dropped/not_acked,而接收端其实一直在线。OnDisconnect发现已有新连接时会立即WakePush,几乎一定早于新连接的 SUBSCRIBE。MaxReceiveBytes=0,超限的大帧会直接发给只声明了 1024 字节的小设备,而不是拒收并回too_large。push.go:88-99:只比对连接代号。push.go:549-560:WakePush同样不看握手状态。session.go:71-73:旧连接断开时唤醒新连接。server.go:457-476:先继承会话,再发 CONNACK,再调OnSessionEstablished,之后才开始读包。server.go:1001-1021:没有订阅者时什么都不做,投递失败只记 debug 日志。expire_at;推送必须等握手完成。DEVIATIONS 里没有相关说明。LiveConn加Ready bool。OnSessionEstablished登记为Ready:false,OnHandshakeComplete登记为Ready:true并带上MaxReceiveBytes。PushPending遇到当前连接不存在、代号不匹配或未就绪时直接返回:不 claim,也不发回执。WakePush和messageLoops同样只处理就绪连接。dispatchFullTx的在线判定保持现状(含握手中),符合 7.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为空;置为就绪后才发布。too_large回执。[P-02] 握手完成前就推送:订阅 down 之前推出的消息被 mochi 静默丢弃,却记成"已推送",不保留的消息 5 分钟后按 not_acked 丢弃
appUplink.OnSessionEstablished在回 CONNACK 之后、客户端发 SUBSCRIBE 之前,就把这个连接登记进了memConns。messageLoops每秒对conns.Snapshot()里的每个连接调用PushPending。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 分钟。其余位置:
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 里没有相关记录。解决方案:
PublishDown发布前检查目标连接已订阅 down(复用hasDownSub(st)),没订阅就返回ErrNotSubscribed(包一层ErrNoConnection)。message 包现有的发布失败分支会清标记并 1 秒后重推。messageLoops只对brk.IsHandshook(ep)为真、且brk.CurrentConnID(ep)等于live.ConnID的连接调用PushPending。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 之前已经检查过订阅,不受影响。
需补测试:
PublishDown返回错误。ack_timeout_seconds调到几秒,确认 hello 之后能收到这条消息,没有被标成 dropped。置信度:代码阅读确定(时序由 mochi 源码推出),未用测试复现
[M-11] WakePush 每次都新起协程、不合并:同一端并发推送,可能超出推送窗口、打乱顺序
PushPending在跑:每秒循环、每次唤醒各一个。push.go:549-560;push.go:110-118(先读已推数再推);push.go:190-208(逐条 claim);push.go:309-338(超时处理没有重查pushed_at)。WakePush只发信号,每秒循环只负责给各 worker 发信号。pushed_at仍满足超时条件。internal/app/message/push.go、app.go;cmd/nixmsg/serve.go。(send_at, seq)一致。[M-13] messageLoops 单协程串行、没有超时:任何一个慢操作都会冻结全局
PushPending、清理,以及修复 #5 加入的三条 COUNT 统计。serve.go:322-348。context.WithTimeout。cmd/nixmsg/serve.go;internal/app/message/app.go。[M-12] 到点分发每秒最多 100 条、逐条串行提交,积压没有上限;遇到一条失败会卡住整批
serve.go:331;push.go:29-34、55-64;recover.go:37。send_at设定时器,提交定时消息时通知它。internal/app/message/push.go、dispatch.go、recover.go;cmd/nixmsg/serve.go。[M-15] claim 已提交却没发布时,只能靠 5 分钟确认超时兜底,不保留消息被误判为没确认
ctx.Err()。PushPending会当成 claim 失败直接返回,不发布。clearPushed复用已过期的 ctx,写队列会直接拒绝执行。dropped/not_acked。store/queue.go:91-104(成功也返回 ctx 错误)、queue.go:71-73(ctx 已过期直接拒绝);push.go:191-205、226-227。clearPushed固定用新建的短超时 ctx;claim 返回 ctx 错误时,读一次pushed_conn,如果已被本次写成当前连接就清掉并安排 1 秒后重推。internal/app/message/push.go。pushed_conn被清空并在 1 秒后重推。[M-02] message 包自己那份大帧名额在断线、撤回、作废时不释放,同一投递重推还会重复占用
acquireLarge,再用trackLarge把largeHeld[seq:编号]置为 true。这是布尔值,不计数。OnPublishDropped。OnDisconnect只清推送标记。重连后同一投递重推又占一个名额,确认时只释放一个。每次「大帧在途时断线」漏 1 个。Recall不释放;之后的 ack 结果是 recalled,而Ack只在 accepted 时释放。push.go:508-528:trackLarge是布尔,releaseLarge只删一次。ack.go:71-73:只在 accepted 时释放。ack.go:139-162:撤回不释放。session.go:31-75:断线不释放。largeSem、largeHeld、acquireLarge、trackLarge、releaseLarge,只依赖修好后的 broker 名额(M-01)。测试用的假下行如需限流,在假实现里模拟。同步改 DEVIATIONS M2/M3/M4 第 2 条。(connID, seq, 编号);同一 key 已持有就不再获取;OnDisconnect释放该连接名下所有 key;Ack不论结果都释放;撤回和finishDeliveryTx对已推送的投递释放;每分钟对账一次,投递不再是「pending 且推给该连接」就释放,以兜住 identity/group 的作废路径。internal/app/message/push.go、ack.go、session.go、app.go。len(a.largeSem)==0。复审基线:main
4059a15(2026-09-30)。编号说明、各工作线的合并顺序、共享文件归属见总览 #7。C-01 已在
feat/fix-message-c01-c03提交eca4c836f1(#32)仅向已握手连接推送;
StartLoops/WaitLoops拆开调度循环。停机等待交给 L-03。未合入 main。已合入 origin/main
0c9b459。落地提交83521f4fix: 仅向已握手连接推送并拆分调度循环 (#32)。合入时8f2ebc7补唤醒回执推送。