fix: 回执按连接去重在途且不改可重复约定

This commit is contained in:
Nixevol
2026-09-30 16:22:43 +08:00
parent 83521f4b87
commit bd550f313f
2 changed files with 188 additions and 1 deletions
+187
View File
@@ -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)
}
}