package main import ( "bytes" "context" "encoding/json" "io" "net/http" "strings" "sync" "testing" "time" "git.asio.asia/nixevol/NixMsg/internal/config" "git.asio.asia/nixevol/NixMsg/internal/protocol" "git.asio.asia/nixevol/NixMsg/test/harness" "github.com/mochi-mqtt/server/v2/packets" ) // TestUplinkDMOfflineGroupRecall 用真实进程验证上行分发:单聊、离线保留、群发不含发送者、延迟内撤回。 func TestUplinkDMOfflineGroupRecall(t *testing.T) { dataDir := t.TempDir() cfgPath := writeTestConfig(t, dataDir) initAdminForTest(t, dataDir) enableRegistration(t, dataDir, "uplink-code") cfg, err := config.Load(cfgPath) if err != nil { t.Fatal(err) } if vErr := cfg.Validate(); vErr != nil { t.Fatal(vErr) } ctx, cancel := context.WithCancel(context.Background()) defer cancel() errCh := make(chan error, 1) go func() { errCh <- runServe(ctx, cfg) }() defer func() { cancel() select { case err := <-errCh: if err != nil { t.Errorf("serve exit: %v", err) } case <-time.After(15 * time.Second): t.Error("serve did not stop") } }() addr := waitListenAddr(t, dataDir, 15*time.Second) base := "http://" + addr registerEP(t, base, "alice", "password12", "Alice") registerEP(t, base, "bob", "password12", "Bob") registerEP(t, base, "carol", "password12", "Carol") alice := mqttSessionLogin(t, base, "alice", "password12") defer alice.Close() bob := mqttSessionLogin(t, base, "bob", "password12") defer bob.Close() // 1) 两端在线单聊:bob 收到 msg delay0 := int64(0) sendResp := alice.Request(t, map[string]any{ "v": 1, "type": "send", "rid": "s1", "id": "dm-1", "to": map[string]any{"kind": "endpoint", "id": "bob"}, "body": map[string]any{"enc": "utf8", "data": "hello-bob"}, "delay_ms": delay0, }) if !sendResp.OK { t.Fatalf("send dm: %+v", sendResp) } msg := bob.WaitType(t, "msg", 8*time.Second) if msg["id"] != "dm-1" || msg["from"] != "alice" { t.Fatalf("bob msg=%v", msg) } body, _ := msg["body"].(map[string]any) if body["data"] != "hello-bob" { t.Fatalf("body=%v", body) } // 确认以免窗口占满 bob.Request(t, map[string]any{"v": 1, "type": "ack", "rid": "a1", "from": "alice", "id": "dm-1"}) // 2) 离线保留:carol 离线时 alice 发送,carol 稍后上线能收到 keep := true ttl := int64(86400) sendOff := alice.Request(t, map[string]any{ "v": 1, "type": "send", "rid": "s2", "id": "off-1", "to": map[string]any{"kind": "endpoint", "id": "carol"}, "body": map[string]any{"enc": "utf8", "data": "for-carol"}, "delay_ms": delay0, "offline": map[string]any{"keep": keep, "ttl_seconds": ttl}, }) if !sendOff.OK { t.Fatalf("send offline: %+v", sendOff) } carol := mqttSessionLogin(t, base, "carol", "password12") defer carol.Close() offMsg := carol.WaitType(t, "msg", 8*time.Second) if offMsg["id"] != "off-1" { t.Fatalf("carol offline msg=%v", offMsg) } carol.Request(t, map[string]any{"v": 1, "type": "ack", "rid": "a2", "from": "alice", "id": "off-1"}) // 3) 建群后群发:bob、carol 收到,alice 自己不收到 createResp := alice.Request(t, map[string]any{ "v": 1, "type": "group.create", "rid": "g1", "id": "g_uplink1", "name": "U", "members": []map[string]any{{"id": "bob"}, {"id": "carol"}}, }) if !createResp.OK { t.Fatalf("group.create: %+v", createResp) } // 排空建群带来的 group_event,避免干扰后续断言 drainEvents(t, bob, 500*time.Millisecond) drainEvents(t, carol, 500*time.Millisecond) drainEvents(t, alice, 300*time.Millisecond) grpSend := alice.Request(t, map[string]any{ "v": 1, "type": "send", "rid": "s3", "id": "grp-1", "to": map[string]any{"kind": "group", "id": "g_uplink1"}, "body": map[string]any{"enc": "utf8", "data": "hi-group"}, "delay_ms": delay0, }) if !grpSend.OK { t.Fatalf("group send: %+v", grpSend) } bobGrp := bob.WaitType(t, "msg", 8*time.Second) carolGrp := carol.WaitType(t, "msg", 8*time.Second) if bobGrp["id"] != "grp-1" || carolGrp["id"] != "grp-1" { t.Fatalf("bob=%v carol=%v", bobGrp, carolGrp) } if got := alice.TryType("msg", 800*time.Millisecond); got != nil { t.Fatalf("sender must not receive own group msg: %v", got) } bob.Request(t, map[string]any{"v": 1, "type": "ack", "rid": "a3", "from": "alice", "id": "grp-1"}) carol.Request(t, map[string]any{"v": 1, "type": "ack", "rid": "a4", "from": "alice", "id": "grp-1"}) // 4) 延迟内撤回:对方无 msg / revoked 回调 delayMs := int64(60_000) sched := alice.Request(t, map[string]any{ "v": 1, "type": "send", "rid": "s4", "id": "rec-1", "to": map[string]any{"kind": "endpoint", "id": "bob"}, "body": map[string]any{"enc": "utf8", "data": "will-recall"}, "delay_ms": delayMs, }) if !sched.OK { t.Fatalf("scheduled send: %+v", sched) } data, _ := sched.Data.(map[string]any) if data["state"] != "scheduled" { t.Fatalf("want scheduled got %v", data) } rec := alice.Request(t, map[string]any{"v": 1, "type": "recall", "rid": "r1", "id": "rec-1"}) if !rec.OK { t.Fatalf("recall: %+v", rec) } if got := bob.TryType("msg", 1*time.Second); got != nil { t.Fatalf("bob should not get recalled msg: %v", got) } if got := bob.TryType("revoked", 500*time.Millisecond); got != nil { t.Fatalf("bob should not get revoked for undelivered: %v", got) } } func registerEP(t *testing.T, base, id, password, name string) { t.Helper() body := `{"registration_code":"uplink-code","id":"` + id + `","login_password":"` + password + `","name":"` + name + `"}` resp, err := http.Post(base+"/api/client/register", "application/json", strings.NewReader(body)) if err != nil { t.Fatal(err) } raw, _ := io.ReadAll(resp.Body) _ = resp.Body.Close() if resp.StatusCode != http.StatusOK { t.Fatalf("register %s: %d %s", id, resp.StatusCode, raw) } } type mqttSess struct { t *testing.T mc harness.MQTTClient endpointID string pktID uint16 mu sync.Mutex inbox []map[string]any closed bool done chan struct{} } type appResp struct { OK bool Error map[string]any Data any Raw map[string]any } func mqttSessionLogin(t *testing.T, httpBase, endpointID, password string) *mqttSess { t.Helper() mc, err := harness.DialMQTTWebSocket(httpBase, 10*time.Second) if err != nil { t.Fatalf("dial: %v", err) } s := &mqttSess{t: t, mc: mc, endpointID: endpointID, pktID: 10, done: make(chan struct{})} s.connectSubscribeHello(password) go s.readLoop() return s } func (s *mqttSess) Close() { s.mu.Lock() if s.closed { s.mu.Unlock() return } s.closed = true s.mu.Unlock() _ = s.mc.Close() select { case <-s.done: case <-time.After(3 * time.Second): } } func (s *mqttSess) nextPkt() uint16 { s.pktID++ if s.pktID == 0 { s.pktID = 1 } return s.pktID } func (s *mqttSess) connectSubscribeHello(password string) { t := s.t pk := packets.Packet{ FixedHeader: packets.FixedHeader{Type: packets.Connect}, ProtocolVersion: 5, Connect: packets.ConnectParams{ ProtocolName: []byte("MQTT"), Clean: true, ClientIdentifier: s.endpointID, Keepalive: 30, UsernameFlag: true, Username: []byte(s.endpointID), PasswordFlag: true, Password: []byte(password), }, } var buf bytes.Buffer if err := pk.ConnectEncode(&buf); err != nil { t.Fatal(err) } if err := s.mc.Send(buf.Bytes()); err != nil { t.Fatal(err) } ack, err := s.mc.Recv() if err != nil { t.Fatal(err) } if len(ack) < 4 || ack[0]>>4 != packets.Connack || ack[3] != 0 { t.Fatalf("connack %x", ack) } sub := packets.Packet{ FixedHeader: packets.FixedHeader{Type: packets.Subscribe, Qos: 1}, ProtocolVersion: 5, PacketID: s.nextPkt(), Filters: packets.Subscriptions{ {Filter: "nix/c/" + s.endpointID + "/down", Qos: 1}, }, } buf.Reset() if err := sub.SubscribeEncode(&buf); err != nil { t.Fatal(err) } if err := s.mc.Send(buf.Bytes()); err != nil { t.Fatal(err) } if _, err := s.mc.Recv(); err != nil { t.Fatal(err) } hello, _ := protocol.Marshal(protocol.Hello{V: protocol.Version, Type: protocol.TypeHello, RID: "h0"}) s.publishRaw(hello) deadline := time.Now().Add(10 * time.Second) for time.Now().Before(deadline) { raw, err := s.mc.Recv() if err != nil { t.Fatal(err) } m := s.handlePacket(raw) if m == nil { continue } if m["type"] == "resp" && m["ok"] == true { return } if m["type"] == "resp" { t.Fatalf("hello failed: %v", m) } s.push(m) } t.Fatal("hello timeout") } func (s *mqttSess) publishRaw(payload []byte) { t := s.t pub := packets.Packet{ FixedHeader: packets.FixedHeader{Type: packets.Publish, Qos: 1}, ProtocolVersion: 5, TopicName: "nix/c/" + s.endpointID + "/up", PacketID: s.nextPkt(), Payload: payload, } var buf bytes.Buffer if err := pub.PublishEncode(&buf); err != nil { t.Fatal(err) } if err := s.mc.Send(buf.Bytes()); err != nil { t.Fatal(err) } } func (s *mqttSess) readLoop() { defer close(s.done) for { raw, err := s.mc.Recv() if err != nil { return } m := s.handlePacket(raw) if m != nil { s.push(m) } } } func (s *mqttSess) handlePacket(raw []byte) map[string]any { if len(raw) < 2 { return nil } typ := raw[0] >> 4 qos := (raw[0] >> 1) & 0x3 switch typ { case packets.Puback, packets.Pingresp, packets.Suback: return nil case packets.Publish: payload, err := decodePublishPayload(raw) if err != nil { return nil } if qos == 1 { // 回 PUBACK rem, n, _ := decodeRemainingLength(raw[1:]) body := raw[1+n:] pk := packets.Packet{ProtocolVersion: 5, FixedHeader: packets.FixedHeader{Type: packets.Publish, Remaining: rem, Qos: qos}} if decErr := pk.PublishDecode(body); decErr == nil && pk.PacketID != 0 { ack := packets.Packet{ FixedHeader: packets.FixedHeader{Type: packets.Puback}, ProtocolVersion: 5, PacketID: pk.PacketID, } var buf bytes.Buffer if encErr := ack.PubackEncode(&buf); encErr == nil { _ = s.mc.Send(buf.Bytes()) } } } var m map[string]any if json.Unmarshal(payload, &m) != nil { return nil } return m default: return nil } } func (s *mqttSess) push(m map[string]any) { s.mu.Lock() s.inbox = append(s.inbox, m) s.mu.Unlock() } func (s *mqttSess) Request(t *testing.T, frame map[string]any) appResp { t.Helper() rid, _ := frame["rid"].(string) payload, err := protocol.Marshal(frame) if err != nil { t.Fatal(err) } s.publishRaw(payload) deadline := time.Now().Add(10 * time.Second) for time.Now().Before(deadline) { m := s.takeMatching(func(x map[string]any) bool { return x["type"] == "resp" && x["rid"] == rid }) if m != nil { r := appResp{OK: m["ok"] == true, Raw: m, Data: m["data"]} if e, ok := m["error"].(map[string]any); ok { r.Error = e } return r } time.Sleep(5 * time.Millisecond) } t.Fatalf("timeout waiting resp rid=%s", rid) return appResp{} } func (s *mqttSess) WaitType(t *testing.T, typ string, timeout time.Duration) map[string]any { t.Helper() deadline := time.Now().Add(timeout) for time.Now().Before(deadline) { m := s.takeMatching(func(x map[string]any) bool { return x["type"] == typ }) if m != nil { return m } time.Sleep(5 * time.Millisecond) } t.Fatalf("timeout waiting type=%s", typ) return nil } func (s *mqttSess) TryType(typ string, timeout time.Duration) map[string]any { deadline := time.Now().Add(timeout) for time.Now().Before(deadline) { m := s.takeMatching(func(x map[string]any) bool { return x["type"] == typ }) if m != nil { return m } time.Sleep(10 * time.Millisecond) } return nil } func (s *mqttSess) takeMatching(pred func(map[string]any) bool) map[string]any { s.mu.Lock() defer s.mu.Unlock() for i, m := range s.inbox { if pred(m) { s.inbox = append(s.inbox[:i], s.inbox[i+1:]...) return m } } return nil } func drainEvents(t *testing.T, s *mqttSess, d time.Duration) { t.Helper() deadline := time.Now().Add(d) for time.Now().Before(deadline) { _ = s.takeMatching(func(x map[string]any) bool { typ, _ := x["type"].(string) return typ == "group_event" || typ == "presence" }) time.Sleep(20 * time.Millisecond) } }