From b79bba391149877fc4fef9d374eb4bd0c192b5e3 Mon Sep 17 00:00:00 2001 From: Nixevol Date: Wed, 30 Sep 2026 17:29:09 +0800 Subject: [PATCH] =?UTF-8?q?fix:=20=E8=A1=A5=E4=B8=89=E5=A4=84=E6=8E=A8?= =?UTF-8?q?=E9=80=81=E6=BC=8F=E5=94=A4=E9=86=92=E8=AE=A9=E5=9B=9E=E6=89=A7?= =?UTF-8?q?=E4=B8=8E=E5=90=8E=E7=BB=AD=20pending=20=E7=BB=A7=E7=BB=AD?= =?UTF-8?q?=E6=8E=A8=20(#65)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- docs/DEVIATIONS.md | 10 ++ internal/app/message/ack.go | 5 + internal/app/message/push.go | 8 + internal/app/message/review_r3_65_test.go | 180 ++++++++++++++++++++++ 4 files changed, 203 insertions(+) create mode 100644 internal/app/message/review_r3_65_test.go diff --git a/docs/DEVIATIONS.md b/docs/DEVIATIONS.md index 7deb4c5..cc28178 100644 --- a/docs/DEVIATIONS.md +++ b/docs/DEVIATIONS.md @@ -1718,3 +1718,13 @@ issue #3 未关闭,`feat/fix-3-downlink-deadlock` 未合入 `main`。下面是 - 原因:原先 12 项备注写着未穷尽仍标通过,交付说明写成「通过 23」。 - 备选方案:为每个未测子项补验收用例(本波不做,避免为变绿放松断言)。 - 影响:汇总改为通过 19、部分通过 4(F03/F08/F21/F22)、失败 0。F19 仍引用仓库内 SDK 清单、本波不重跑。 + +### 复审修复 R3-01 + +1. **推送 worker 漏唤醒:定时分发回执、超限后续、撤回腾窗** + - 日期:2026-09-30 + - 原条款:Gitea #65;推送 worker 只在 `WakePush` 时跑一轮 `PushPending`。 + - 实际做法:`dispatchDueBatch` 本条已分发时除 pending 接收方外始终 `WakePush` 发送方;`PushPending` 遇 `too_large` 等跳过且本轮未占满窗口时再 `WakePush` 当前接收方(不在同一次递归扫表);`Recall` 成功改成 `recalled` 的接收方各 `WakePush` 一次,已推送撤回仍走 `flushRevokes`。不改 `PublishDown` 签名。 + - 原因:终态回执、超限后的后续 pending、撤回腾出窗口后都依赖再唤醒,否则空闲在线端收不到。 + - 备选方案:在同一次 `PushPending` 里循环扫完整 pending 表(否决,指令要求合并唤醒下一轮)。 + - 影响:仅 `internal/app/message/`;相关单测见 `review_r3_65_test.go`。 diff --git a/internal/app/message/ack.go b/internal/app/message/ack.go index d6694dd..831e4ca 100644 --- a/internal/app/message/ack.go +++ b/internal/app/message/ack.go @@ -90,6 +90,7 @@ func (a *App) Recall(ctx context.Context, senderID string, req *protocol.Recall) nowMs := a.now().UnixMilli() var data protocol.RecallData var revokes []revokeJob + var wakeReceivers []string err := a.db.Queue.Do(ctx, func(tx *sql.Tx) error { var seq int64 var state string @@ -156,6 +157,7 @@ WHERE seq = ? AND endpoint_id = ? AND state = 'pending'`, continue } recalled++ + wakeReceivers = append(wakeReceivers, p.ep) if p.pushed.Valid && p.pushed.String != "" { revokes = append(revokes, revokeJob{ endpointID: p.ep, @@ -190,6 +192,9 @@ SELECT COUNT(*) FROM deliveries WHERE seq = ? AND state IN ('expired','dropped', a.pendingRevoke = append(a.pendingRevoke, revokes...) a.mu.Unlock() a.flushRevokes(ctx) + for _, ep := range wakeReceivers { + a.WakePush(ep) + } return data, nil } diff --git a/internal/app/message/push.go b/internal/app/message/push.go index b788ac1..89b4dd3 100644 --- a/internal/app/message/push.go +++ b/internal/app/message/push.go @@ -127,6 +127,7 @@ LIMIT ?`, nowMs, limit) } mu.Lock() n++ + wake[d.senderID] = struct{}{} mu.Unlock() rows2, qErr := a.db.Read.QueryContext(ctx, ` SELECT DISTINCT endpoint_id FROM deliveries WHERE seq = ? AND state = 'pending'`, d.seq) @@ -203,8 +204,10 @@ LIMIT ?`, endpointID, room) } var toClaim []pushItem + skipped := false for _, it := range items { if it.body == nil { + skipped = true continue } msg := protocol.Msg{ @@ -227,6 +230,7 @@ LIMIT ?`, endpointID, room) return rejErr } a.WakePush(it.senderID) + skipped = true continue } it.payload = payload @@ -252,6 +256,10 @@ LIMIT ?`, endpointID, room) a.observeDispatchToPush(it.sendAt, nowMs) } } + // 超限等跳过未占满窗口时再唤醒本端,让 worker 下一轮取后续 pending(不在本轮递归扫表)。 + if skipped && len(claimed) < room { + a.WakePush(endpointID) + } return a.pushReceipts(ctx, endpointID, connID, nowMs) } diff --git a/internal/app/message/review_r3_65_test.go b/internal/app/message/review_r3_65_test.go new file mode 100644 index 0000000..b305e4a --- /dev/null +++ b/internal/app/message/review_r3_65_test.go @@ -0,0 +1,180 @@ +package message + +import ( + "context" + "database/sql" + "encoding/json" + "testing" + "time" + + "git.asio.asia/nixevol/NixMsg/internal/app/port" + "git.asio.asia/nixevol/NixMsg/internal/protocol" +) + +// 定时到点且接收方离线成终态时,在线发送方经 WakePush 拿到回执。 +func TestR365DispatchDueWakesSenderForReceipt(t *testing.T) { + t.Parallel() + e := openDeliveryEnv(t, nil) + insertEndpoint(t, e.db, "alice", "", 1, 0) + insertEndpoint(t, e.db, "bob", "", 1, 0) + _ = e.db.Queue.Do(context.Background(), func(tx *sql.Tx) error { + _, err := tx.Exec(`UPDATE endpoints SET offline_since = ? WHERE id=?`, e.nowMs-120_000, "bob") + return err + }) + ctx := context.Background() + aliceConn := LiveConn{ConnID: "c-alice", Ready: true} + e.conns.Set("alice", aliceConn) + if err := e.app.OnHandshakeComplete(ctx, "alice", aliceConn); err != nil { + t.Fatal(err) + } + t.Cleanup(func() { e.app.stopPushWorker("alice", "") }) + + delay := int64(5000) + req := baseSend("due-rcpt", "bob") + req.DelayMs = &delay + if _, err := e.app.Submit(ctx, "alice", port.ConnInfo{}, req); err != nil { + t.Fatal(err) + } + st, _ := e.msgState("alice", "due-rcpt") + if st != StateScheduled { + t.Fatalf("want scheduled got %s", st) + } + if e.down.FilterType(protocol.TypeReceipt) != 0 { + t.Fatal("receipt before due") + } + + e.setNow(e.nowMs + 5000) + if _, err := e.app.DispatchDue(ctx, e.nowMs, 10); err != nil { + t.Fatal(err) + } + deadline := time.Now().Add(2 * time.Second) + for time.Now().Before(deadline) { + if e.down.FilterType(protocol.TypeReceipt) > 0 { + return + } + time.Sleep(20 * time.Millisecond) + } + t.Fatalf("sender got no receipt after scheduled dispatch, receipts=%d", e.down.FilterType(protocol.TypeReceipt)) +} + +// 窗口内首条超限拒收后,同连接后续小消息仍被 worker 推送。 +func TestR365TooLargeThenSmallContinuesPush(t *testing.T) { + t.Parallel() + e := openDeliveryEnv(t, func(l *Limits) { l.DeliveryWindow = 1 }) + insertEndpoint(t, e.db, "alice", "", 1, 0) + insertEndpoint(t, e.db, "bob", "", 1, 0) + ctx := context.Background() + // 整帧上限:小消息能过,大正文整帧超限被拒。 + bobConn := LiveConn{ConnID: "c-bob", Ready: true, MaxReceiveBytes: 200} + e.conns.Set("bob", bobConn) + if err := e.app.OnHandshakeComplete(ctx, "bob", bobConn); err != nil { + t.Fatal(err) + } + t.Cleanup(func() { e.app.stopPushWorker("bob", "") }) + + big := baseSend("big-skip", "bob") + b := make([]byte, 400) + for i := range b { + b[i] = 'A' + } + big.Body.Data = string(b) + if _, err := e.app.Submit(ctx, "alice", port.ConnInfo{}, big); err != nil { + t.Fatal(err) + } + small := baseSend("small-ok", "bob") + small.Body.Data = "hi" + if _, err := e.app.Submit(ctx, "alice", port.ConnInfo{}, small); err != nil { + t.Fatal(err) + } + + deadline := time.Now().Add(2 * time.Second) + for time.Now().Before(deadline) { + foundSmall := false + for _, p := range e.down.Snapshots() { + if payloadType(p.Payload) != protocol.TypeMsg { + continue + } + var m struct { + ID string `json:"id"` + } + _ = json.Unmarshal(p.Payload, &m) + if m.ID == "small-ok" { + foundSmall = true + } + } + if foundSmall { + seqBig := e.seqOf("alice", "big-skip") + st, reason := e.deliveryState(seqBig, "bob") + if st != DeliveryRejected || reason != ReasonTooLarge { + t.Fatalf("big: %s/%s", st, reason) + } + return + } + time.Sleep(20 * time.Millisecond) + } + t.Fatalf("small msg not pushed after too_large; msgs=%d snapshots=%d", + e.down.FilterType(protocol.TypeMsg), len(e.down.Snapshots())) +} + +// 撤回已推在途消息后,同连接更晚的 pending 继续推。 +func TestR365RecallFreesWindowForLaterPending(t *testing.T) { + t.Parallel() + e := openDeliveryEnv(t, func(l *Limits) { l.DeliveryWindow = 1 }) + insertEndpoint(t, e.db, "alice", "", 1, 0) + insertEndpoint(t, e.db, "bob", "", 1, 0) + ctx := context.Background() + bobConn := LiveConn{ConnID: "c-bob", Ready: true} + e.conns.Set("bob", bobConn) + if err := e.app.OnHandshakeComplete(ctx, "bob", bobConn); err != nil { + t.Fatal(err) + } + t.Cleanup(func() { e.app.stopPushWorker("bob", "") }) + + if _, err := e.app.Submit(ctx, "alice", port.ConnInfo{}, baseSend("first-in-flight", "bob")); err != nil { + t.Fatal(err) + } + if _, err := e.app.Submit(ctx, "alice", port.ConnInfo{}, baseSend("second-wait", "bob")); err != nil { + t.Fatal(err) + } + + deadline := time.Now().Add(2 * time.Second) + for time.Now().Before(deadline) { + if e.down.FilterType(protocol.TypeMsg) >= 1 { + break + } + time.Sleep(10 * time.Millisecond) + } + if e.down.FilterType(protocol.TypeMsg) < 1 { + t.Fatal("first msg not pushed") + } + seq1 := e.seqOf("alice", "first-in-flight") + var pushed sql.NullString + _ = e.db.Read.QueryRow(`SELECT pushed_conn FROM deliveries WHERE seq=? AND endpoint_id='bob'`, seq1).Scan(&pushed) + if !pushed.Valid { + t.Fatal("first not in-flight") + } + + if _, err := e.app.Recall(ctx, "alice", &protocol.Recall{ + V: protocol.Version, Type: protocol.TypeRecall, RID: "r1", ID: "first-in-flight", + }); err != nil { + t.Fatal(err) + } + + deadline = time.Now().Add(2 * time.Second) + for time.Now().Before(deadline) { + for _, p := range e.down.Snapshots() { + if payloadType(p.Payload) != protocol.TypeMsg { + continue + } + var m struct { + ID string `json:"id"` + } + _ = json.Unmarshal(p.Payload, &m) + if m.ID == "second-wait" { + return + } + } + time.Sleep(20 * time.Millisecond) + } + t.Fatalf("second pending not pushed after recall; msgs=%d", e.down.FilterType(protocol.TypeMsg)) +}