From b8569b0708f83d84423e03754194eacbb68fab50 Mon Sep 17 00:00:00 2001 From: Nixevol Date: Wed, 30 Sep 2026 08:07:16 +0800 Subject: [PATCH] =?UTF-8?q?feat:=20=E6=8E=A5=E7=BA=BF=E4=B8=8A=E8=A1=8C?= =?UTF-8?q?=E5=B8=A7=E5=88=86=E5=8F=91=E5=88=B0=20message/identity/presenc?= =?UTF-8?q?e/group?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- cmd/nixmsg/serve.go | 8 +- cmd/nixmsg/uplink.go | 295 ++++++++++++++++- cmd/nixmsg/uplink_integration_test.go | 459 ++++++++++++++++++++++++++ docs/DEVIATIONS.md | 43 ++- internal/app/presence/encode.go | 22 ++ 5 files changed, 805 insertions(+), 22 deletions(-) create mode 100644 cmd/nixmsg/uplink_integration_test.go create mode 100644 internal/app/presence/encode.go diff --git a/cmd/nixmsg/serve.go b/cmd/nixmsg/serve.go index 8ee23fe..192dfe3 100644 --- a/cmd/nixmsg/serve.go +++ b/cmd/nixmsg/serve.go @@ -88,7 +88,6 @@ func runServe(ctx context.Context, cfg config.Config) error { uplink := &appUplink{ msg: msgApp, conns: memConns, - next: port.StubUplinkHandler{}, log: slog.Default(), } sess := broker.NewSession(broker.SessionOptions{ @@ -128,6 +127,7 @@ func runServe(ctx context.Context, cfg config.Config) error { message.WithDownlink(brk), ) uplink.msg = msgApp + uplink.down = brk presApp := presence.New(presence.Config{ DB: db, @@ -135,6 +135,7 @@ func runServe(ctx context.Context, cfg config.Config) error { Conns: &presenceConnTable{conns: memConns}, }) sess.SetPresence(presApp) + uplink.presence = presApp idApp := identity.New(identity.Config{ DB: db, @@ -145,13 +146,16 @@ func runServe(ctx context.Context, cfg config.Config) error { Logger: slog.Default(), ConnControl: brk, }) - _ = group.New(group.Config{ + uplink.identity = idApp + + groupApp := group.New(group.Config{ DB: db, Talk: idApp, Online: presApp, Downlink: brk, MaxGroupMembers: cfg.Limits.MaxGroupMembers, }) + uplink.groups = groupApp if recoverErr := msgApp.RecoverOnStart(ctx); recoverErr != nil { return fmt.Errorf("message recover: %w", recoverErr) diff --git a/cmd/nixmsg/uplink.go b/cmd/nixmsg/uplink.go index e9631a7..9ca0dca 100644 --- a/cmd/nixmsg/uplink.go +++ b/cmd/nixmsg/uplink.go @@ -2,19 +2,27 @@ package main import ( "context" + "encoding/json" + "errors" "log/slog" + "git.asio.asia/nixevol/NixMsg/internal/app/group" + "git.asio.asia/nixevol/NixMsg/internal/app/identity" "git.asio.asia/nixevol/NixMsg/internal/app/message" "git.asio.asia/nixevol/NixMsg/internal/app/port" + "git.asio.asia/nixevol/NixMsg/internal/app/presence" + "git.asio.asia/nixevol/NixMsg/internal/protocol" ) -// appUplink 把 broker 生命周期接到消息连接表与投递推送。 -// 其余业务上行帧暂转交 next(可为空 Stub);完整 HandleUplink 分发留后续波次。 +// appUplink 把 broker 生命周期接到消息连接表,并把已握手上行帧分发到各业务服务。 type appUplink struct { - msg *message.App - conns *message.MemoryConns - next port.UplinkHandler - log *slog.Logger + msg *message.App + identity *identity.App + presence *presence.App + groups *group.App + conns *message.MemoryConns + down port.Downlink + log *slog.Logger } func (u *appUplink) OnSessionEstablished(ctx context.Context, conn port.ConnInfo) error { @@ -22,9 +30,6 @@ func (u *appUplink) OnSessionEstablished(ctx context.Context, conn port.ConnInfo ConnID: conn.ConnID, MaxPacketSize: conn.MaxPacketSize, }) - if u.next != nil { - return u.next.OnSessionEstablished(ctx, conn) - } return nil } @@ -39,31 +44,287 @@ func (u *appUplink) OnHandshakeComplete(ctx context.Context, hs port.HandshakeIn u.log.Error("message handshake", "endpoint", hs.EndpointID, "err", err) return err } - if u.next != nil { - return u.next.OnHandshakeComplete(ctx, hs) - } return nil } func (u *appUplink) OnDisconnect(ctx context.Context, conn port.ConnInfo, reason port.DisconnectReason) { + if u.presence != nil { + u.presence.ClearWatch(conn.ConnID) + } live, ok := u.conns.Current(conn.EndpointID) isCurrent := ok && live.ConnID == conn.ConnID if err := u.msg.OnDisconnect(ctx, conn.EndpointID, conn.ConnID, isCurrent); err != nil { u.log.Error("message disconnect", "endpoint", conn.EndpointID, "err", err) } u.conns.Clear(conn.EndpointID, conn.ConnID) - if u.next != nil { - u.next.OnDisconnect(ctx, conn, reason) - } } func (u *appUplink) HandleUplink(ctx context.Context, conn port.ConnInfo, payload []byte) error { - if u.next != nil { - return u.next.HandleUplink(ctx, conn, payload) + frame, err := protocol.Decode(payload) + if err != nil { + u.replyErr(ctx, conn, peekRID(payload), protocol.CodeBadRequest, err.Error()) + return nil } + + rid, data, callErr := u.dispatch(ctx, conn, frame) + if callErr != nil { + u.replyFromErr(ctx, conn, rid, callErr) + return nil + } + u.replyOK(ctx, conn, rid, data) return nil } +func (u *appUplink) dispatch(ctx context.Context, conn port.ConnInfo, frame any) (rid string, data any, err error) { + switch f := frame.(type) { + case *protocol.Send: + rid = f.RID + res, e := u.msg.Submit(ctx, conn.EndpointID, conn, f) + if e != nil { + return rid, nil, e + } + return rid, protocol.SendData{ID: res.ID, SendAtMs: res.SendAtMs, State: res.State}, nil + + case *protocol.Ack: + rid = f.RID + res, e := u.msg.Ack(ctx, conn.EndpointID, f) + if e != nil { + return rid, nil, e + } + return rid, map[string]any{"result": res.Result}, nil + + case *protocol.ReceiptAck: + rid = f.RID + return rid, nil, u.msg.ReceiptAck(ctx, conn.EndpointID, f) + + case *protocol.Recall: + rid = f.RID + res, e := u.msg.Recall(ctx, conn.EndpointID, f) + if e != nil { + return rid, nil, e + } + return rid, res, nil + + case *protocol.Status: + rid = f.RID + res, e := u.msg.Status(ctx, conn.EndpointID, f) + return rid, res, e + + case *protocol.Unlock: + rid = f.RID + if e := f.Validate(); e != nil { + return rid, nil, e + } + return rid, nil, u.identity.UnlockTalk(ctx, conn.EndpointID, f.EndpointID, f.TalkPassword, conn.RemoteIP) + + case *protocol.SelfGet: + rid = f.RID + info, e := u.identity.SelfGet(ctx, conn.EndpointID) + return rid, info, e + + case *protocol.SelfUpdate: + rid = f.RID + return rid, nil, u.identity.SelfUpdate(ctx, conn.EndpointID, f) + + case *protocol.SelfTalkPassword: + rid = f.RID + if e := f.Validate(); e != nil { + return rid, nil, e + } + return rid, nil, u.identity.SelfSetTalkPassword(ctx, conn.EndpointID, f.TalkPassword) + + case *protocol.SelfLoginPassword: + rid = f.RID + if e := f.Validate(); e != nil { + return rid, nil, e + } + tok, e := u.identity.SelfChangeLoginPassword(ctx, conn.EndpointID, f.OldPassword, f.NewPassword, conn.RemoteIP) + if e != nil { + return rid, nil, e + } + return rid, map[string]any{"session_token": tok}, nil + + case *protocol.PresenceGet: + rid = f.RID + if e := f.Validate(); e != nil { + return rid, nil, e + } + items, e := u.presence.Get(ctx, f.IDs) + if e != nil { + return rid, nil, e + } + return rid, presence.EncodeGetData(items), nil + + case *protocol.DirectoryList: + rid = f.RID + items, next, e := u.presence.Directory(ctx, f) + if e != nil { + return rid, nil, e + } + return rid, map[string]any{"items": items, "next_cursor": next}, nil + + case *protocol.PresenceWatch: + rid = f.RID + return rid, nil, u.presence.Watch(ctx, conn.ConnID, conn.EndpointID, f) + + case *protocol.GroupCreate: + rid = f.RID + res, e := u.groups.Create(ctx, conn.EndpointID, f) + return rid, res, e + + case *protocol.GroupAdd: + rid = f.RID + res, e := u.groups.Add(ctx, conn.EndpointID, f) + return rid, res, e + + case *protocol.GroupRemove: + rid = f.RID + return rid, nil, u.groups.Remove(ctx, conn.EndpointID, f) + + case *protocol.GroupLeave: + rid = f.RID + return rid, nil, u.groups.Leave(ctx, conn.EndpointID, f) + + case *protocol.GroupTransfer: + rid = f.RID + return rid, nil, u.groups.Transfer(ctx, conn.EndpointID, f) + + case *protocol.GroupRename: + rid = f.RID + return rid, nil, u.groups.Rename(ctx, conn.EndpointID, f) + + case *protocol.GroupDissolve: + rid = f.RID + return rid, nil, u.groups.Dissolve(ctx, conn.EndpointID, f) + + case *protocol.GroupList: + rid = f.RID + items, next, e := u.groups.List(ctx, conn.EndpointID, f) + if e != nil { + return rid, nil, e + } + return rid, map[string]any{"items": items, "next_cursor": next}, nil + + case *protocol.GroupGet: + rid = f.RID + res, e := u.groups.Get(ctx, conn.EndpointID, f) + return rid, res, e + + case *protocol.Hello, *protocol.SelfLogout: + // Session 已处理;不应落到此处。 + rid = frameRID(frame) + return rid, nil, &protocol.Error{Code: protocol.CodeBadRequest, Message: "unexpected frame"} + + default: + rid = frameRID(frame) + return rid, nil, &protocol.Error{Code: protocol.CodeBadRequest, Message: "unsupported uplink type"} + } +} + +func (u *appUplink) replyOK(ctx context.Context, conn port.ConnInfo, rid string, data any) { + if rid == "" { + rid = "0" + } + var raw json.RawMessage + if data != nil { + b, err := protocol.Marshal(data) + if err != nil { + u.replyErr(ctx, conn, rid, protocol.CodeBusy, "marshal resp data") + return + } + raw = b + } + resp := protocol.Resp{V: protocol.Version, Type: protocol.TypeResp, RID: rid, OK: true, Data: raw} + u.publishResp(ctx, conn, resp) +} + +func (u *appUplink) replyFromErr(ctx context.Context, conn port.ConnInfo, rid string, err error) { + if rid == "" { + rid = "0" + } + var pe *protocol.Error + if errors.As(err, &pe) && pe != nil { + u.replyErr(ctx, conn, rid, pe.Code, pe.Message) + return + } + u.log.Error("uplink handler", "endpoint", conn.EndpointID, "err", err) + u.replyErr(ctx, conn, rid, protocol.CodeBusy, "internal error") +} + +func (u *appUplink) replyErr(ctx context.Context, conn port.ConnInfo, rid, code, message string) { + if rid == "" { + rid = "0" + } + resp := protocol.Resp{ + V: protocol.Version, + Type: protocol.TypeResp, + RID: rid, + OK: false, + Error: &protocol.ErrorBody{Code: code, Message: message}, + } + u.publishResp(ctx, conn, resp) +} + +func (u *appUplink) publishResp(ctx context.Context, conn port.ConnInfo, resp protocol.Resp) { + if u.down == nil { + return + } + b, err := protocol.Marshal(resp) + if err != nil { + u.log.Error("marshal resp", "err", err) + return + } + if live, ok := u.conns.Current(conn.EndpointID); ok && live.ConnID == conn.ConnID { + limit := respPayloadLimit(live.MaxPacketSize, live.MaxReceiveBytes) + if limit > 0 && len(b) > limit { + tooLarge := protocol.Resp{ + V: protocol.Version, + Type: protocol.TypeResp, + RID: resp.RID, + OK: false, + Error: &protocol.ErrorBody{Code: protocol.CodeResponseTooLarge, Message: "response too large"}, + } + b, err = protocol.Marshal(tooLarge) + if err != nil { + return + } + } + } + if pubErr := u.down.PublishDown(ctx, conn.EndpointID, conn.ConnID, b, port.PublishOpts{QoS: 1}); pubErr != nil { + u.log.Error("publish resp", "endpoint", conn.EndpointID, "err", pubErr) + } +} + +func respPayloadLimit(maxPacketSize uint32, maxRecvBytes int) int { + limit := 0 + if maxRecvBytes > 0 { + limit = maxRecvBytes + } + if maxPacketSize > 0 { + n := int(maxPacketSize) + if limit == 0 || n < limit { + limit = n + } + } + return limit +} + +func peekRID(payload []byte) string { + var peek struct { + RID string `json:"rid"` + } + _ = json.Unmarshal(payload, &peek) + return peek.RID +} + +func frameRID(frame any) string { + b, err := protocol.Marshal(frame) + if err != nil { + return "" + } + return peekRID(b) +} + // presenceConnTable 把消息连接表暴露给 presence.ConnTable。 type presenceConnTable struct { conns *message.MemoryConns diff --git a/cmd/nixmsg/uplink_integration_test.go b/cmd/nixmsg/uplink_integration_test.go new file mode 100644 index 0000000..bed1f90 --- /dev/null +++ b/cmd/nixmsg/uplink_integration_test.go @@ -0,0 +1,459 @@ +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) + } +} diff --git a/docs/DEVIATIONS.md b/docs/DEVIATIONS.md index 5e1163b..00855a6 100644 --- a/docs/DEVIATIONS.md +++ b/docs/DEVIATIONS.md @@ -52,12 +52,12 @@ - 备选方案:分文件多阶段接线。 - 影响:进程可对外登录/注册/握手。 -2. **未接线 / 未完成部分** +2. **未接线 / 未完成部分(已被 L-UPLINK 部分取代)** - 原条款:上行应用帧完整分发到 identity/group/presence/message 业务方法。 - - 实际做法:`Session` 处理 hello/logout;其余 `HandleUplink` 仍为 `StubUplinkHandler`(不解析 send/ack/group 等)。管理注册设置 HTTP(A3)未实现,测试直接写 `settings` 表。群/在线模块已构造并注入 Downlink,但无上行入口调用。 + - 实际做法(L-WIRE 当时):`Session` 处理 hello/logout;其余 `HandleUplink` 仍为 Stub。管理注册设置 HTTP(A3)未实现,测试直接写 `settings` 表。 - 原因:本波强制三件验收;完整协议分发属后续波次。 - 备选方案:本波同时实现 HandleUplink 大 multiplex。 - - 影响:端连上后除握手/logout 外业务帧尚无 resp;F03–F16 等仍依赖后续接线。 + - 影响:见 L-UPLINK。 3. **listen.addr 始终写入** - 原条款:listener N1 仅端口 0 写文件;T0.1 超集为启动即写。 @@ -66,6 +66,43 @@ - 备选方案:改 listener 始终写。 - 影响:固定端口也会有地址文件。 +### L-UPLINK 2026-09-30 + +1. **HandleUplink 分发到已有业务服务** + - 原条款:TASKS 总控接线;DEVELOPMENT 第 6 节上行帧由服务端处理后回 `resp`;群事件/在线通知/消息走下行。 + - 实际做法:`cmd/nixmsg/uplink.go` 的 `appUplink` 在已握手连接上解码并调用 `message`/`identity`/`presence`/`group` 的既有方法;结果编成 `resp` 经 `PublishDown`(QoS 1)回本连接。`Session` 仍独占 `hello`/`self.logout` 与未握手 `not_ready`。`serve` 把四个 App 与 broker Downlink 注入 uplink。不重写业务状态机。 + - 原因:main 上业务实现已齐,缺上行入口。 + - 备选方案:各包自建 MQTT 钩子(与 port 契约不符)。 + - 影响:端侧 send/ack/recall/status/unlock/self.*/presence.*/directory.list/group.* 可走真实进程。 + +2. **presence.get 未知编号编码** + - 原条款:DEVELOPMENT 6.5「未知编号为 not_found」;`StatusItem.NotFound` 为内部标记。 + - 实际做法:`presence.EncodeGetData` 薄封装,resp data 为 `{"items":[...]}`;未知项 `{"id","not_found":true}`,已知项含 `id`/`online`/`since_ms`。 + - 原因:Service 层未规定 JSON 形状,接线需固定可编解码形式。 + - 备选方案:整请求失败 `not_found`;或把 `not_found` 塞进 `online` 旁字符串字段。 + - 影响:SDK 若只认 items 数组两种形状均可(Go SDK 已兼容 wrap/array)。 + +3. **resp 超限改发 response_too_large** + - 原条款:DEVELOPMENT 7.5「resp 超限改发 response_too_large」。 + - 实际做法:发布前按连接表 `MaxReceiveBytes` 与 `MaxPacketSize` 取较小正上限;超限则改发错误 resp(不再发原 data)。未做 MQTT 包头开销扣减(与 message 推送里的 overhead 预算不完全同一函数)。 + - 原因:接线层最小可用检查。 + - 备选方案:复用 `message.effectivePayloadLimit`(未导出)。 + - 影响:接近包上限的大分页可能比推送路径略严或略松。 + +4. **断线清 presence.watch** + - 原条款:presence.watch 断线清空。 + - 实际做法:`OnDisconnect` 额外调 `ClearWatch`;`SetOffline` 路径本身也会清。顶号旧连接未走 `SetOffline(isCurrent)` 时仍能清订阅。 + - 原因:避免旧连接订阅泄漏。 + - 备选方案:仅依赖 Session 对 isCurrent 调 SetOffline。 + - 影响:无。 + +5. **仍未接线 / 本波未覆盖** + - 管理注册设置 HTTP(A3)仍未挂;集成测继续写 `settings` 表开注册。 + - 下行通知帧(`msg`/`receipt`/`revoked`/`presence`/`group_event`/`fatal`)不是上行分发对象,由既有服务经 PublishDown 发出。 + - `self.logout`/`hello` 仍在 Session,不经 appUplink。 + - 身份 I5 停用/删除级联、管理端群等非本任务范围。 + - 集成测覆盖:双端单聊、离线 keep 后上线、群发不含发送者、延迟内撤回对方无回调;未穷尽 presence.watch 通知与全部 group.* 变体。 + ### T0.2 2026-09-30 1. **请求指纹规范化格式** diff --git a/internal/app/presence/encode.go b/internal/app/presence/encode.go new file mode 100644 index 0000000..4df941d --- /dev/null +++ b/internal/app/presence/encode.go @@ -0,0 +1,22 @@ +package presence + +// EncodeGetData 把 Get 结果编成 presence.get 的 resp data(items 数组)。 +// 未知编号项为 {"id":"...","not_found":true},其余含 id/online/since_ms。 +func EncodeGetData(items []StatusItem) map[string]any { + out := make([]map[string]any, 0, len(items)) + for _, it := range items { + if it.NotFound { + out = append(out, map[string]any{ + "id": it.ID, + "not_found": true, + }) + continue + } + out = append(out, map[string]any{ + "id": it.ID, + "online": it.Online, + "since_ms": it.SinceMs, + }) + } + return map[string]any{"items": out} +}