460 lines
12 KiB
Go
460 lines
12 KiB
Go
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)
|
|
}
|
|
}
|