fix: 回执按连接去重在途且不改可重复约定
This commit is contained in:
+1
-1
@@ -1021,7 +1021,7 @@ close()
|
||||
|
||||
## 10. 无 SDK 的设备
|
||||
|
||||
使用 MQTT 5 客户端(例如 ESP-IDF 的 mqtt,走 WebSocket 或 TCP)。ClientID 和用户名是端编号,密码是登录密码。订阅 `nix/c/{编号}/down`,向 `nix/c/{编号}/up` 发第 6 节的 JSON。收到 `msg` 后处理完发 `ack`。用「from + id」去重,已确认过的再到达时再回一次 `ack`。在 `hello` 里按自己的接收缓冲声明 `max_receive_bytes`(不小于 1024)。
|
||||
使用 MQTT 5 客户端(例如 ESP-IDF 的 mqtt,走 WebSocket 或 TCP)。ClientID 和用户名是端编号,密码是登录密码。订阅 `nix/c/{编号}/down`,向 `nix/c/{编号}/up` 发第 6 节的 JSON。收到 `msg` 后处理完发 `ack`。用「from + id」去重,已确认过的再到达时再回一次 `ack`。在 `hello` 里按自己的接收缓冲声明 `max_receive_bytes`(不小于 1024)。裸设备要么实现 `receipt_ack`,要么发送时带 `receipt:false`,否则未确认的回执会占满推送窗口。
|
||||
|
||||
不要开持久会话,不要订阅通配符,不要发保留消息。能发 HTTP 请求的设备可以按第 6.9 节自助注册,不能的由管理员开通。设备可以每次都用登录密码连接(每次都算新登录、会换令牌),也可以把握手拿到的 `session_token` 存起来,重连时当作密码用。3.1.1 客户端看不到被顶号的原因;每次都用密码登录的设备不要和别的设备共用编号,否则会来回互踢。
|
||||
|
||||
|
||||
@@ -0,0 +1,187 @@
|
||||
package message
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"strconv"
|
||||
"sync"
|
||||
"testing"
|
||||
|
||||
"git.asio.asia/nixevol/NixMsg/internal/app/port"
|
||||
"git.asio.asia/nixevol/NixMsg/internal/protocol"
|
||||
)
|
||||
|
||||
func TestC02ReceiptWindowAndDedupe(t *testing.T) {
|
||||
t.Parallel()
|
||||
e := openDeliveryEnv(t, func(l *Limits) { l.ReceiptWindow = 2 })
|
||||
insertEndpoint(t, e.db, "alice", "", 1, 0)
|
||||
insertEndpoint(t, e.db, "bob", "", 1, 0)
|
||||
e.online("alice", "c-alice")
|
||||
e.online("bob", "c-bob")
|
||||
ctx := context.Background()
|
||||
for i := 0; i < 5; i++ {
|
||||
id := "rct" + strconv.Itoa(i)
|
||||
if _, err := e.app.Submit(ctx, "alice", port.ConnInfo{}, baseSend(id, "bob")); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := e.app.PushPending(ctx, "bob", "c-bob"); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if _, err := e.app.Ack(ctx, "bob", &protocol.Ack{V: 1, Type: protocol.TypeAck, RID: "1", From: "alice", ID: id}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
e.down.mu.Lock()
|
||||
e.down.Published = nil
|
||||
e.down.mu.Unlock()
|
||||
for i := 0; i < 3; i++ {
|
||||
if err := e.app.PushPending(ctx, "alice", "c-alice"); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
if n := e.down.FilterType(protocol.TypeReceipt); n != 2 {
|
||||
t.Fatalf("3 pushes with window 2: got %d receipts want 2", n)
|
||||
}
|
||||
var firstID string
|
||||
for _, p := range e.down.Snapshots() {
|
||||
if payloadType(p.Payload) == protocol.TypeReceipt {
|
||||
var head struct {
|
||||
ReceiptID string `json:"receipt_id"`
|
||||
}
|
||||
_ = json.Unmarshal(p.Payload, &head)
|
||||
firstID = head.ReceiptID
|
||||
break
|
||||
}
|
||||
}
|
||||
if firstID == "" {
|
||||
t.Fatal("no receipt_id")
|
||||
}
|
||||
if err := e.app.ReceiptAck(ctx, "alice", &protocol.ReceiptAck{V: 1, Type: protocol.TypeReceiptAck, ReceiptID: firstID}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
before := e.down.FilterType(protocol.TypeReceipt)
|
||||
if err := e.app.PushPending(ctx, "alice", "c-alice"); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
after := e.down.FilterType(protocol.TypeReceipt)
|
||||
if after-before != 1 {
|
||||
t.Fatalf("after ack new receipts=%d want 1", after-before)
|
||||
}
|
||||
}
|
||||
|
||||
func TestC02ConcurrentPushReceiptsOnce(t *testing.T) {
|
||||
t.Parallel()
|
||||
e := openDeliveryEnv(t, func(l *Limits) { l.ReceiptWindow = 64 })
|
||||
insertEndpoint(t, e.db, "alice", "", 1, 0)
|
||||
insertEndpoint(t, e.db, "bob", "", 1, 0)
|
||||
e.online("alice", "c-alice")
|
||||
e.online("bob", "c-bob")
|
||||
ctx := context.Background()
|
||||
if _, err := e.app.Submit(ctx, "alice", port.ConnInfo{}, baseSend("once1", "bob")); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
_ = e.app.PushPending(ctx, "bob", "c-bob")
|
||||
if _, err := e.app.Ack(ctx, "bob", &protocol.Ack{V: 1, Type: protocol.TypeAck, RID: "1", From: "alice", ID: "once1"}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
e.down.mu.Lock()
|
||||
e.down.Published = nil
|
||||
e.down.mu.Unlock()
|
||||
var wg sync.WaitGroup
|
||||
for i := 0; i < 10; i++ {
|
||||
wg.Add(1)
|
||||
go func() {
|
||||
defer wg.Done()
|
||||
_ = e.app.PushPending(ctx, "alice", "c-alice")
|
||||
}()
|
||||
}
|
||||
wg.Wait()
|
||||
if n := e.down.FilterType(protocol.TypeReceipt); n != 1 {
|
||||
t.Fatalf("concurrent push receipts=%d want 1", n)
|
||||
}
|
||||
}
|
||||
|
||||
func TestC02DisconnectResendsUnacked(t *testing.T) {
|
||||
t.Parallel()
|
||||
e := openDeliveryEnv(t, func(l *Limits) { l.ReceiptWindow = 8 })
|
||||
insertEndpoint(t, e.db, "alice", "", 1, 0)
|
||||
insertEndpoint(t, e.db, "bob", "", 1, 0)
|
||||
e.online("alice", "c-alice")
|
||||
e.online("bob", "c-bob")
|
||||
ctx := context.Background()
|
||||
if _, err := e.app.Submit(ctx, "alice", port.ConnInfo{}, baseSend("rs1", "bob")); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
_ = e.app.PushPending(ctx, "bob", "c-bob")
|
||||
if _, err := e.app.Ack(ctx, "bob", &protocol.Ack{V: 1, Type: protocol.TypeAck, RID: "1", From: "alice", ID: "rs1"}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
_ = e.app.PushPending(ctx, "alice", "c-alice")
|
||||
if n := e.down.FilterType(protocol.TypeReceipt); n != 1 {
|
||||
t.Fatalf("first send %d", n)
|
||||
}
|
||||
if err := e.app.OnDisconnect(ctx, "alice", "c-alice", true); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
e.conns.Clear("alice", "c-alice")
|
||||
live := LiveConn{ConnID: "c-alice2", Ready: true}
|
||||
e.conns.Set("alice", live)
|
||||
e.down.mu.Lock()
|
||||
e.down.Published = nil
|
||||
e.down.mu.Unlock()
|
||||
if err := e.app.PushPending(ctx, "alice", "c-alice2"); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if n := e.down.FilterType(protocol.TypeReceipt); n != 1 {
|
||||
t.Fatalf("reconnect resend %d want 1", n)
|
||||
}
|
||||
}
|
||||
|
||||
func TestC02DuplicateAckDoesNotWakeReceipt(t *testing.T) {
|
||||
t.Parallel()
|
||||
e := openDeliveryEnv(t, nil)
|
||||
insertEndpoint(t, e.db, "alice", "", 1, 0)
|
||||
insertEndpoint(t, e.db, "bob", "", 1, 0)
|
||||
e.online("alice", "c-alice")
|
||||
e.online("bob", "c-bob")
|
||||
ctx := context.Background()
|
||||
if _, err := e.app.Submit(ctx, "alice", port.ConnInfo{}, baseSend("dupack", "bob")); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
_ = e.app.PushPending(ctx, "bob", "c-bob")
|
||||
if _, err := e.app.Ack(ctx, "bob", &protocol.Ack{V: 1, Type: protocol.TypeAck, RID: "1", From: "alice", ID: "dupack"}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
_ = e.app.PushPending(ctx, "alice", "c-alice")
|
||||
before := e.down.FilterType(protocol.TypeReceipt)
|
||||
if _, err := e.app.Ack(ctx, "bob", &protocol.Ack{V: 1, Type: protocol.TypeAck, RID: "2", From: "alice", ID: "dupack"}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
_ = e.app.PushPending(ctx, "alice", "c-alice")
|
||||
after := e.down.FilterType(protocol.TypeReceipt)
|
||||
if after != before {
|
||||
t.Fatalf("dup ack added receipts %d -> %d", before, after)
|
||||
}
|
||||
}
|
||||
|
||||
func TestC02SelfMessageSingleReceipt(t *testing.T) {
|
||||
t.Parallel()
|
||||
e := openDeliveryEnv(t, nil)
|
||||
insertEndpoint(t, e.db, "alice", "", 1, 0)
|
||||
e.online("alice", "c-alice")
|
||||
ctx := context.Background()
|
||||
if _, err := e.app.Submit(ctx, "alice", port.ConnInfo{}, baseSend("self1", "alice")); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
_ = e.app.PushPending(ctx, "alice", "c-alice")
|
||||
if _, err := e.app.Ack(ctx, "alice", &protocol.Ack{V: 1, Type: protocol.TypeAck, RID: "1", From: "alice", ID: "self1"}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
e.down.mu.Lock()
|
||||
e.down.Published = nil
|
||||
e.down.mu.Unlock()
|
||||
_ = e.app.PushPending(ctx, "alice", "c-alice")
|
||||
if n := e.down.FilterType(protocol.TypeReceipt); n != 1 {
|
||||
t.Fatalf("self receipt frames=%d want 1", n)
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user