feat: 接线上行帧分发到 message/identity/presence/group

This commit is contained in:
Nixevol
2026-09-30 08:07:16 +08:00
parent 0efcfbb8b6
commit b8569b0708
5 changed files with 805 additions and 22 deletions
+6 -2
View File
@@ -88,7 +88,6 @@ func runServe(ctx context.Context, cfg config.Config) error {
uplink := &appUplink{ uplink := &appUplink{
msg: msgApp, msg: msgApp,
conns: memConns, conns: memConns,
next: port.StubUplinkHandler{},
log: slog.Default(), log: slog.Default(),
} }
sess := broker.NewSession(broker.SessionOptions{ sess := broker.NewSession(broker.SessionOptions{
@@ -128,6 +127,7 @@ func runServe(ctx context.Context, cfg config.Config) error {
message.WithDownlink(brk), message.WithDownlink(brk),
) )
uplink.msg = msgApp uplink.msg = msgApp
uplink.down = brk
presApp := presence.New(presence.Config{ presApp := presence.New(presence.Config{
DB: db, DB: db,
@@ -135,6 +135,7 @@ func runServe(ctx context.Context, cfg config.Config) error {
Conns: &presenceConnTable{conns: memConns}, Conns: &presenceConnTable{conns: memConns},
}) })
sess.SetPresence(presApp) sess.SetPresence(presApp)
uplink.presence = presApp
idApp := identity.New(identity.Config{ idApp := identity.New(identity.Config{
DB: db, DB: db,
@@ -145,13 +146,16 @@ func runServe(ctx context.Context, cfg config.Config) error {
Logger: slog.Default(), Logger: slog.Default(),
ConnControl: brk, ConnControl: brk,
}) })
_ = group.New(group.Config{ uplink.identity = idApp
groupApp := group.New(group.Config{
DB: db, DB: db,
Talk: idApp, Talk: idApp,
Online: presApp, Online: presApp,
Downlink: brk, Downlink: brk,
MaxGroupMembers: cfg.Limits.MaxGroupMembers, MaxGroupMembers: cfg.Limits.MaxGroupMembers,
}) })
uplink.groups = groupApp
if recoverErr := msgApp.RecoverOnStart(ctx); recoverErr != nil { if recoverErr := msgApp.RecoverOnStart(ctx); recoverErr != nil {
return fmt.Errorf("message recover: %w", recoverErr) return fmt.Errorf("message recover: %w", recoverErr)
+276 -15
View File
@@ -2,18 +2,26 @@ package main
import ( import (
"context" "context"
"encoding/json"
"errors"
"log/slog" "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/message"
"git.asio.asia/nixevol/NixMsg/internal/app/port" "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 生命周期接到消息连接表与投递推送。 // appUplink 把 broker 生命周期接到消息连接表,并把已握手上行帧分发到各业务服务。
// 其余业务上行帧暂转交 next(可为空 Stub);完整 HandleUplink 分发留后续波次。
type appUplink struct { type appUplink struct {
msg *message.App msg *message.App
identity *identity.App
presence *presence.App
groups *group.App
conns *message.MemoryConns conns *message.MemoryConns
next port.UplinkHandler down port.Downlink
log *slog.Logger log *slog.Logger
} }
@@ -22,9 +30,6 @@ func (u *appUplink) OnSessionEstablished(ctx context.Context, conn port.ConnInfo
ConnID: conn.ConnID, ConnID: conn.ConnID,
MaxPacketSize: conn.MaxPacketSize, MaxPacketSize: conn.MaxPacketSize,
}) })
if u.next != nil {
return u.next.OnSessionEstablished(ctx, conn)
}
return nil 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) u.log.Error("message handshake", "endpoint", hs.EndpointID, "err", err)
return err return err
} }
if u.next != nil {
return u.next.OnHandshakeComplete(ctx, hs)
}
return nil return nil
} }
func (u *appUplink) OnDisconnect(ctx context.Context, conn port.ConnInfo, reason port.DisconnectReason) { 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) live, ok := u.conns.Current(conn.EndpointID)
isCurrent := ok && live.ConnID == conn.ConnID isCurrent := ok && live.ConnID == conn.ConnID
if err := u.msg.OnDisconnect(ctx, conn.EndpointID, conn.ConnID, isCurrent); err != nil { if err := u.msg.OnDisconnect(ctx, conn.EndpointID, conn.ConnID, isCurrent); err != nil {
u.log.Error("message disconnect", "endpoint", conn.EndpointID, "err", err) u.log.Error("message disconnect", "endpoint", conn.EndpointID, "err", err)
} }
u.conns.Clear(conn.EndpointID, conn.ConnID) 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 { func (u *appUplink) HandleUplink(ctx context.Context, conn port.ConnInfo, payload []byte) error {
if u.next != nil { frame, err := protocol.Decode(payload)
return u.next.HandleUplink(ctx, conn, payload) if err != nil {
} u.replyErr(ctx, conn, peekRID(payload), protocol.CodeBadRequest, err.Error())
return nil 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。 // presenceConnTable 把消息连接表暴露给 presence.ConnTable。
type presenceConnTable struct { type presenceConnTable struct {
conns *message.MemoryConns conns *message.MemoryConns
+459
View File
@@ -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)
}
}
+40 -3
View File
@@ -52,12 +52,12 @@
- 备选方案:分文件多阶段接线。 - 备选方案:分文件多阶段接线。
- 影响:进程可对外登录/注册/握手。 - 影响:进程可对外登录/注册/握手。
2. **未接线 / 未完成部分** 2. **未接线 / 未完成部分(已被 L-UPLINK 部分取代)**
- 原条款:上行应用帧完整分发到 identity/group/presence/message 业务方法。 - 原条款:上行应用帧完整分发到 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。 - 备选方案:本波同时实现 HandleUplink 大 multiplex。
- 影响:端连上后除握手/logout 外业务帧尚无 resp;F03–F16 等仍依赖后续接线。 - 影响:见 L-UPLINK。
3. **listen.addr 始终写入** 3. **listen.addr 始终写入**
- 原条款:listener N1 仅端口 0 写文件;T0.1 超集为启动即写。 - 原条款:listener N1 仅端口 0 写文件;T0.1 超集为启动即写。
@@ -66,6 +66,43 @@
- 备选方案:改 listener 始终写。 - 备选方案:改 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 ### T0.2 2026-09-30
1. **请求指纹规范化格式** 1. **请求指纹规范化格式**
+22
View File
@@ -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}
}