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) } }