Compare commits
1
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
68e790be60 |
+1
-6
@@ -88,6 +88,7 @@ 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{
|
||||||
@@ -127,7 +128,6 @@ 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,7 +135,6 @@ 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,
|
||||||
@@ -146,8 +145,6 @@ func runServe(ctx context.Context, cfg config.Config) error {
|
|||||||
Logger: slog.Default(),
|
Logger: slog.Default(),
|
||||||
ConnControl: brk,
|
ConnControl: brk,
|
||||||
})
|
})
|
||||||
uplink.identity = idApp
|
|
||||||
|
|
||||||
groupApp := group.New(group.Config{
|
groupApp := group.New(group.Config{
|
||||||
DB: db,
|
DB: db,
|
||||||
Talk: idApp,
|
Talk: idApp,
|
||||||
@@ -155,7 +152,6 @@ func runServe(ctx context.Context, cfg config.Config) error {
|
|||||||
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)
|
||||||
@@ -169,7 +165,6 @@ func runServe(ctx context.Context, cfg config.Config) error {
|
|||||||
Locks: loginLocks,
|
Locks: loginLocks,
|
||||||
Logger: slog.Default(),
|
Logger: slog.Default(),
|
||||||
TrustedProxies: trustedNets,
|
TrustedProxies: trustedNets,
|
||||||
Identity: idApp,
|
|
||||||
Groups: groupApp,
|
Groups: groupApp,
|
||||||
Config: cfg,
|
Config: cfg,
|
||||||
Version: Version,
|
Version: Version,
|
||||||
|
|||||||
+15
-276
@@ -2,26 +2,18 @@ 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
|
||||||
down port.Downlink
|
next port.UplinkHandler
|
||||||
log *slog.Logger
|
log *slog.Logger
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -30,6 +22,9 @@ 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
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -44,287 +39,31 @@ 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 {
|
||||||
frame, err := protocol.Decode(payload)
|
if u.next != nil {
|
||||||
if err != nil {
|
return u.next.HandleUplink(ctx, conn, payload)
|
||||||
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
|
||||||
|
|||||||
@@ -1,459 +0,0 @@
|
|||||||
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)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
+3
-63
@@ -52,12 +52,12 @@
|
|||||||
- 备选方案:分文件多阶段接线。
|
- 备选方案:分文件多阶段接线。
|
||||||
- 影响:进程可对外登录/注册/握手。
|
- 影响:进程可对外登录/注册/握手。
|
||||||
|
|
||||||
2. **未接线 / 未完成部分(已被 L-UPLINK 部分取代)**
|
2. **未接线 / 未完成部分**
|
||||||
- 原条款:上行应用帧完整分发到 identity/group/presence/message 业务方法。
|
- 原条款:上行应用帧完整分发到 identity/group/presence/message 业务方法。
|
||||||
- 实际做法(L-WIRE 当时):`Session` 处理 hello/logout;其余 `HandleUplink` 仍为 Stub。管理注册设置 HTTP(A3)未实现,测试直接写 `settings` 表。
|
- 实际做法:`Session` 处理 hello/logout;其余 `HandleUplink` 仍为 `StubUplinkHandler`(不解析 send/ack/group 等)。管理注册设置 HTTP(A3)未实现,测试直接写 `settings` 表。群/在线模块已构造并注入 Downlink,但无上行入口调用。
|
||||||
- 原因:本波强制三件验收;完整协议分发属后续波次。
|
- 原因:本波强制三件验收;完整协议分发属后续波次。
|
||||||
- 备选方案:本波同时实现 HandleUplink 大 multiplex。
|
- 备选方案:本波同时实现 HandleUplink 大 multiplex。
|
||||||
- 影响:见 L-UPLINK。
|
- 影响:端连上后除握手/logout 外业务帧尚无 resp;F03–F16 等仍依赖后续接线。
|
||||||
|
|
||||||
3. **listen.addr 始终写入**
|
3. **listen.addr 始终写入**
|
||||||
- 原条款:listener N1 仅端口 0 写文件;T0.1 超集为启动即写。
|
- 原条款:listener N1 仅端口 0 写文件;T0.1 超集为启动即写。
|
||||||
@@ -66,43 +66,6 @@
|
|||||||
- 备选方案:改 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. **请求指纹规范化格式**
|
||||||
@@ -559,29 +522,6 @@
|
|||||||
- 备选方案:测试里注入在线连接表使默认不保留也走 pending(与「离线不保留」场景重复覆盖)。
|
- 备选方案:测试里注入在线连接表使默认不保留也走 pending(与「离线不保留」场景重复覆盖)。
|
||||||
- 影响:仅测试期望;产品行为不变。
|
- 影响:仅测试期望;产品行为不变。
|
||||||
|
|
||||||
### I5 2026-09-30
|
|
||||||
|
|
||||||
1. **停用/删除级联在 identity 包内完成,admin 可选注入**
|
|
||||||
- 原条款:DEVELOPMENT 7.6 / PRD F01;TASKS I5「提供给 A 线调用」。
|
|
||||||
- 实际做法:`identity.App.Disable`/`Enable`/`Delete` 在一个写操作里完成启停、清令牌、作废消息/投递、发 `revoked`(可选 Downlink)、退群/转让群主/解散、清授权/回执/防重/发出记录。`admin.Deps.Identity` 非空时,`setEndpointEnabled`/`deleteEndpointBasic`(及 PATCH enabled)委托上述方法;为空时保留 A2 仅改库行为。未改 `cmd/nixmsg`、未改消息上行分发。
|
|
||||||
- 原因:隔离要求不改接线;A 线已有路由,只接级联。
|
|
||||||
- 备选方案:在 admin 内复制级联 SQL;或强制 Identity 必填。
|
|
||||||
- 影响:生产须在挂载 admin 时注入 `identity.App`(及 KickEndpoint/ConnControl/Downlink),否则停用/删除仍无完整作废。
|
|
||||||
|
|
||||||
2. **消息级作废回执 state 用 `rejected`**
|
|
||||||
- 原条款:DEVELOPMENT 6.4 消息级作废写 `endpoint_id` 空、`state=rejected`;I4 解散群对 scheduled 曾写 `completed`。
|
|
||||||
- 实际做法:I5 对「发给停用/删除端的 scheduled 单聊」及删除时解散群的 scheduled,回执 `state=rejected`,原因分别为 `endpoint_*` / `group_dissolved`。
|
|
||||||
- 原因:与 6.4 字面一致。
|
|
||||||
- 备选方案:与 I4 一样写 `completed`。
|
|
||||||
- 影响:后台/SDK 若按 state 过滤回执需同时认 rejected。
|
|
||||||
|
|
||||||
3. **删除时群主转让按 `joined_at` 最早,并列按编号**
|
|
||||||
- 原条款:转给最早加入的其他成员。
|
|
||||||
- 实际做法:`ORDER BY joined_at ASC, endpoint_id ASC LIMIT 1`。
|
|
||||||
- 原因:同时加入时需稳定次序。
|
|
||||||
- 备选方案:仅按 joined_at。
|
|
||||||
- 影响:同毫秒加入时编号小者优先。
|
|
||||||
|
|
||||||
## 后台接口 A
|
## 后台接口 A
|
||||||
|
|
||||||
### A1 2026-09-30
|
### A1 2026-09-30
|
||||||
|
|||||||
@@ -320,11 +320,7 @@ func (h *Handler) handleEndpointPatch(w http.ResponseWriter, r *http.Request) {
|
|||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
hasMeta := req.Name != nil || req.Remark != nil || req.DefaultDelaySeconds != nil
|
wasEnabled, err := h.patchEndpoint(r.Context(), id, req.Name, req.Remark, req.DefaultDelaySeconds, req.Enabled)
|
||||||
var wasEnabled bool
|
|
||||||
var err error
|
|
||||||
if hasMeta || (req.Enabled != nil && h.identity == nil) {
|
|
||||||
wasEnabled, err = h.patchEndpoint(r.Context(), id, req.Name, req.Remark, req.DefaultDelaySeconds, req.Enabled)
|
|
||||||
if err != nil {
|
if err != nil {
|
||||||
if errors.Is(err, sql.ErrNoRows) {
|
if errors.Is(err, sql.ErrNoRows) {
|
||||||
h.audit(actorString(p), "endpoint_patch", id, "not_found", ip)
|
h.audit(actorString(p), "endpoint_patch", id, "not_found", ip)
|
||||||
@@ -335,24 +331,7 @@ func (h *Handler) handleEndpointPatch(w http.ResponseWriter, r *http.Request) {
|
|||||||
httpx.WriteError(w, http.StatusInternalServerError, "internal", "内部错误")
|
httpx.WriteError(w, http.StatusInternalServerError, "internal", "内部错误")
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
}
|
if req.Enabled != nil && !*req.Enabled && wasEnabled {
|
||||||
if req.Enabled != nil && h.identity != nil {
|
|
||||||
var found bool
|
|
||||||
found, err = h.setEndpointEnabled(r.Context(), id, *req.Enabled)
|
|
||||||
if err != nil {
|
|
||||||
h.audit(actorString(p), "endpoint_patch", id, "error", ip)
|
|
||||||
httpx.WriteError(w, http.StatusInternalServerError, "internal", "内部错误")
|
|
||||||
return
|
|
||||||
}
|
|
||||||
if !found {
|
|
||||||
h.audit(actorString(p), "endpoint_patch", id, "not_found", ip)
|
|
||||||
httpx.WriteError(w, http.StatusNotFound, "not_found", "端不存在")
|
|
||||||
return
|
|
||||||
}
|
|
||||||
if !*req.Enabled {
|
|
||||||
_, _ = h.kickEndpoint(r.Context(), id)
|
|
||||||
}
|
|
||||||
} else if req.Enabled != nil && !*req.Enabled && wasEnabled {
|
|
||||||
_, _ = h.kickEndpoint(r.Context(), id)
|
_, _ = h.kickEndpoint(r.Context(), id)
|
||||||
}
|
}
|
||||||
row, err := h.getEndpoint(r.Context(), id)
|
row, err := h.getEndpoint(r.Context(), id)
|
||||||
|
|||||||
@@ -308,21 +308,6 @@ func (h *Handler) patchEndpoint(ctx context.Context, id string, name, remark *st
|
|||||||
}
|
}
|
||||||
|
|
||||||
func (h *Handler) setEndpointEnabled(ctx context.Context, id string, enabled bool) (found bool, err error) {
|
func (h *Handler) setEndpointEnabled(ctx context.Context, id string, enabled bool) (found bool, err error) {
|
||||||
if h.identity != nil {
|
|
||||||
var opErr error
|
|
||||||
if enabled {
|
|
||||||
opErr = h.identity.Enable(ctx, id)
|
|
||||||
} else {
|
|
||||||
opErr = h.identity.Disable(ctx, id)
|
|
||||||
}
|
|
||||||
if opErr != nil {
|
|
||||||
if isEndpointNotFound(opErr) {
|
|
||||||
return false, nil
|
|
||||||
}
|
|
||||||
return false, opErr
|
|
||||||
}
|
|
||||||
return true, nil
|
|
||||||
}
|
|
||||||
err = h.db.Queue.Do(ctx, func(tx *sql.Tx) error {
|
err = h.db.Queue.Do(ctx, func(tx *sql.Tx) error {
|
||||||
v := 0
|
v := 0
|
||||||
if enabled {
|
if enabled {
|
||||||
@@ -346,18 +331,9 @@ func (h *Handler) setEndpointEnabled(ctx context.Context, id string, enabled boo
|
|||||||
return found, err
|
return found, err
|
||||||
}
|
}
|
||||||
|
|
||||||
// deleteEndpointBasic 删除端;注入 Identity 时走 I5 完整级联。
|
// deleteEndpointBasic 删除端行并清掉与之相关的 talk_grants。
|
||||||
|
// 作废消息、群主转让等完整级联留给 I5。
|
||||||
func (h *Handler) deleteEndpointBasic(ctx context.Context, id string) (found bool, err error) {
|
func (h *Handler) deleteEndpointBasic(ctx context.Context, id string) (found bool, err error) {
|
||||||
if h.identity != nil {
|
|
||||||
opErr := h.identity.Delete(ctx, id)
|
|
||||||
if opErr != nil {
|
|
||||||
if isEndpointNotFound(opErr) {
|
|
||||||
return false, nil
|
|
||||||
}
|
|
||||||
return false, opErr
|
|
||||||
}
|
|
||||||
return true, nil
|
|
||||||
}
|
|
||||||
err = h.db.Queue.Do(ctx, func(tx *sql.Tx) error {
|
err = h.db.Queue.Do(ctx, func(tx *sql.Tx) error {
|
||||||
if _, e := tx.ExecContext(ctx, `DELETE FROM talk_grants WHERE sender_id = ? OR target_id = ?`, id, id); e != nil {
|
if _, e := tx.ExecContext(ctx, `DELETE FROM talk_grants WHERE sender_id = ? OR target_id = ?`, id, id); e != nil {
|
||||||
return e
|
return e
|
||||||
@@ -373,11 +349,6 @@ func (h *Handler) deleteEndpointBasic(ctx context.Context, id string) (found boo
|
|||||||
return found, err
|
return found, err
|
||||||
}
|
}
|
||||||
|
|
||||||
func isEndpointNotFound(err error) bool {
|
|
||||||
var pe *protocol.Error
|
|
||||||
return errors.As(err, &pe) && pe.Code == protocol.CodeNotFound
|
|
||||||
}
|
|
||||||
|
|
||||||
func (h *Handler) resetLoginPassword(ctx context.Context, id, loginHash string) (bool, error) {
|
func (h *Handler) resetLoginPassword(ctx context.Context, id, loginHash string) (bool, error) {
|
||||||
var found bool
|
var found bool
|
||||||
err := h.db.Queue.Do(ctx, func(tx *sql.Tx) error {
|
err := h.db.Queue.Do(ctx, func(tx *sql.Tx) error {
|
||||||
|
|||||||
@@ -8,7 +8,6 @@ import (
|
|||||||
"time"
|
"time"
|
||||||
|
|
||||||
"git.asio.asia/nixevol/NixMsg/internal/app/group"
|
"git.asio.asia/nixevol/NixMsg/internal/app/group"
|
||||||
"git.asio.asia/nixevol/NixMsg/internal/app/identity"
|
|
||||||
"git.asio.asia/nixevol/NixMsg/internal/auth"
|
"git.asio.asia/nixevol/NixMsg/internal/auth"
|
||||||
"git.asio.asia/nixevol/NixMsg/internal/config"
|
"git.asio.asia/nixevol/NixMsg/internal/config"
|
||||||
"git.asio.asia/nixevol/NixMsg/internal/store"
|
"git.asio.asia/nixevol/NixMsg/internal/store"
|
||||||
@@ -41,8 +40,6 @@ type Deps struct {
|
|||||||
SecureCookies bool
|
SecureCookies bool
|
||||||
// KickEndpoint 踢下线钩子(只断开连接);nil 时踢线为 no-op。
|
// KickEndpoint 踢下线钩子(只断开连接);nil 时踢线为 no-op。
|
||||||
KickEndpoint EndpointKickFunc
|
KickEndpoint EndpointKickFunc
|
||||||
// Identity 端停用/启用/删除级联(I5);nil 时回退为仅改 enabled/删行。
|
|
||||||
Identity identity.Service
|
|
||||||
|
|
||||||
// Groups 群服务;nil 时群写操作返回 busy。列表/详情可只读库。
|
// Groups 群服务;nil 时群写操作返回 busy。列表/详情可只读库。
|
||||||
Groups group.Service
|
Groups group.Service
|
||||||
@@ -63,7 +60,6 @@ type Handler struct {
|
|||||||
ttl time.Duration
|
ttl time.Duration
|
||||||
forceSec bool
|
forceSec bool
|
||||||
kick EndpointKickFunc
|
kick EndpointKickFunc
|
||||||
identity identity.Service
|
|
||||||
groups group.Service
|
groups group.Service
|
||||||
cfg config.Config
|
cfg config.Config
|
||||||
version string
|
version string
|
||||||
@@ -105,7 +101,6 @@ func New(d Deps) *Handler {
|
|||||||
ttl: ttl,
|
ttl: ttl,
|
||||||
forceSec: d.SecureCookies,
|
forceSec: d.SecureCookies,
|
||||||
kick: d.KickEndpoint,
|
kick: d.KickEndpoint,
|
||||||
identity: d.Identity,
|
|
||||||
groups: d.Groups,
|
groups: d.Groups,
|
||||||
cfg: cfg,
|
cfg: cfg,
|
||||||
version: ver,
|
version: ver,
|
||||||
|
|||||||
@@ -32,10 +32,8 @@ type Config struct {
|
|||||||
Sessions auth.SessionTokens
|
Sessions auth.SessionTokens
|
||||||
// MaxScheduleSeconds 限制 self.update 的 default_delay_ms。
|
// MaxScheduleSeconds 限制 self.update 的 default_delay_ms。
|
||||||
MaxScheduleSeconds int64
|
MaxScheduleSeconds int64
|
||||||
// ConnControl 可选:logout / 停用 / 删除后踢线;未接线时为 nil。
|
// ConnControl 可选:logout 后踢线;未接线时为 nil。
|
||||||
ConnControl port.ConnControl
|
ConnControl port.ConnControl
|
||||||
// Downlink 可选:停用/删除时发 revoked 与群事件;未接线时为 nil。
|
|
||||||
Downlink port.Downlink
|
|
||||||
}
|
}
|
||||||
|
|
||||||
// App 实现 identity.Service(含 I1 注册与 I2 self/对话密码)。
|
// App 实现 identity.Service(含 I1 注册与 I2 self/对话密码)。
|
||||||
@@ -47,7 +45,6 @@ type App struct {
|
|||||||
sessions auth.SessionTokens
|
sessions auth.SessionTokens
|
||||||
maxScheduleSeconds int64
|
maxScheduleSeconds int64
|
||||||
connCtrl port.ConnControl
|
connCtrl port.ConnControl
|
||||||
down port.Downlink
|
|
||||||
nowFn func() time.Time
|
nowFn func() time.Time
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -84,7 +81,6 @@ func New(cfg Config) *App {
|
|||||||
sessions: cfg.Sessions,
|
sessions: cfg.Sessions,
|
||||||
maxScheduleSeconds: cfg.MaxScheduleSeconds,
|
maxScheduleSeconds: cfg.MaxScheduleSeconds,
|
||||||
connCtrl: cfg.ConnControl,
|
connCtrl: cfg.ConnControl,
|
||||||
down: cfg.Downlink,
|
|
||||||
nowFn: cfg.Now,
|
nowFn: cfg.Now,
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -1,664 +0,0 @@
|
|||||||
package identity
|
|
||||||
|
|
||||||
import (
|
|
||||||
"bytes"
|
|
||||||
"context"
|
|
||||||
"database/sql"
|
|
||||||
"errors"
|
|
||||||
|
|
||||||
"git.asio.asia/nixevol/NixMsg/internal/app/port"
|
|
||||||
"git.asio.asia/nixevol/NixMsg/internal/protocol"
|
|
||||||
)
|
|
||||||
|
|
||||||
const (
|
|
||||||
reasonEndpointDisabled = "endpoint_disabled"
|
|
||||||
reasonEndpointDeleted = "endpoint_deleted"
|
|
||||||
reasonSenderDisabled = "sender_disabled"
|
|
||||||
reasonSenderDeleted = "sender_deleted"
|
|
||||||
reasonLeftGroup = "left_group"
|
|
||||||
reasonGroupDissolved = "group_dissolved"
|
|
||||||
|
|
||||||
eventOwnerChanged = "owner_changed"
|
|
||||||
eventLeft = "left"
|
|
||||||
eventMemberRemoved = "member_removed"
|
|
||||||
eventDissolved = "dissolved"
|
|
||||||
)
|
|
||||||
|
|
||||||
type revokeItem struct {
|
|
||||||
endpointID string
|
|
||||||
msgID string
|
|
||||||
fromID string
|
|
||||||
reason string
|
|
||||||
}
|
|
||||||
|
|
||||||
type groupNotify struct {
|
|
||||||
recipients []string
|
|
||||||
groupID string
|
|
||||||
event string
|
|
||||||
endpointID string
|
|
||||||
atMs int64
|
|
||||||
}
|
|
||||||
|
|
||||||
// Disable 停用端:清令牌、作废相关消息/投递,并踢连接(DEVELOPMENT 7.6)。
|
|
||||||
func (a *App) Disable(ctx context.Context, endpointID string) error {
|
|
||||||
return a.disableOrDelete(ctx, endpointID, false)
|
|
||||||
}
|
|
||||||
|
|
||||||
// Enable 仅恢复 enabled=1;已作废消息不恢复。
|
|
||||||
func (a *App) Enable(ctx context.Context, endpointID string) error {
|
|
||||||
err := a.db.Queue.Do(ctx, func(tx *sql.Tx) error {
|
|
||||||
res, e := tx.ExecContext(ctx, `UPDATE endpoints SET enabled = 1 WHERE id = ?`, endpointID)
|
|
||||||
if e != nil {
|
|
||||||
return e
|
|
||||||
}
|
|
||||||
n, _ := res.RowsAffected()
|
|
||||||
if n == 0 {
|
|
||||||
return errCode(protocol.CodeNotFound, "endpoint not found")
|
|
||||||
}
|
|
||||||
return nil
|
|
||||||
})
|
|
||||||
return err
|
|
||||||
}
|
|
||||||
|
|
||||||
// Delete 删除端:停用效果(原因改为 deleted)+ 退群/转让群主 + 清授权与发出记录。
|
|
||||||
func (a *App) Delete(ctx context.Context, endpointID string) error {
|
|
||||||
return a.disableOrDelete(ctx, endpointID, true)
|
|
||||||
}
|
|
||||||
|
|
||||||
func (a *App) disableOrDelete(ctx context.Context, endpointID string, hardDelete bool) error {
|
|
||||||
nowMs := a.now().UnixMilli()
|
|
||||||
recvReason := reasonEndpointDisabled
|
|
||||||
sendReason := reasonSenderDisabled
|
|
||||||
if hardDelete {
|
|
||||||
recvReason = reasonEndpointDeleted
|
|
||||||
sendReason = reasonSenderDeleted
|
|
||||||
}
|
|
||||||
|
|
||||||
var revokes []revokeItem
|
|
||||||
var notifies []groupNotify
|
|
||||||
err := a.db.Queue.Do(ctx, func(tx *sql.Tx) error {
|
|
||||||
var one int
|
|
||||||
if e := tx.QueryRowContext(ctx, `SELECT 1 FROM endpoints WHERE id = ?`, endpointID).Scan(&one); e != nil {
|
|
||||||
if errors.Is(e, sql.ErrNoRows) {
|
|
||||||
return errCode(protocol.CodeNotFound, "endpoint not found")
|
|
||||||
}
|
|
||||||
return e
|
|
||||||
}
|
|
||||||
|
|
||||||
if _, e := tx.ExecContext(ctx, `
|
|
||||||
UPDATE endpoints SET enabled = 0,
|
|
||||||
session_hash = NULL, session_issued_at = NULL, session_used_at = NULL
|
|
||||||
WHERE id = ?`, endpointID); e != nil {
|
|
||||||
return e
|
|
||||||
}
|
|
||||||
|
|
||||||
if e := voidEndpointMessagesTx(tx, endpointID, recvReason, sendReason, nowMs, &revokes); e != nil {
|
|
||||||
return e
|
|
||||||
}
|
|
||||||
|
|
||||||
if !hardDelete {
|
|
||||||
return nil
|
|
||||||
}
|
|
||||||
|
|
||||||
if e := leaveAllGroupsTx(tx, endpointID, nowMs, &revokes, ¬ifies); e != nil {
|
|
||||||
return e
|
|
||||||
}
|
|
||||||
if _, e := tx.ExecContext(ctx, `DELETE FROM talk_grants WHERE sender_id = ? OR target_id = ?`, endpointID, endpointID); e != nil {
|
|
||||||
return e
|
|
||||||
}
|
|
||||||
if _, e := tx.ExecContext(ctx, `DELETE FROM receipts WHERE sender_id = ?`, endpointID); e != nil {
|
|
||||||
return e
|
|
||||||
}
|
|
||||||
if _, e := tx.ExecContext(ctx, `DELETE FROM send_keys WHERE sender_id = ?`, endpointID); e != nil {
|
|
||||||
return e
|
|
||||||
}
|
|
||||||
if _, e := tx.ExecContext(ctx, `DELETE FROM messages WHERE sender_id = ?`, endpointID); e != nil {
|
|
||||||
return e
|
|
||||||
}
|
|
||||||
_, e := tx.ExecContext(ctx, `DELETE FROM endpoints WHERE id = ?`, endpointID)
|
|
||||||
return e
|
|
||||||
})
|
|
||||||
if err != nil {
|
|
||||||
return err
|
|
||||||
}
|
|
||||||
|
|
||||||
a.publishRevokes(ctx, revokes)
|
|
||||||
a.publishGroupEvents(ctx, notifies)
|
|
||||||
if a.connCtrl != nil {
|
|
||||||
_ = a.connCtrl.Disconnect(ctx, endpointID, "", port.DisconnectFatal)
|
|
||||||
}
|
|
||||||
return nil
|
|
||||||
}
|
|
||||||
|
|
||||||
func voidEndpointMessagesTx(tx *sql.Tx, endpointID, recvReason, sendReason string, nowMs int64, revokes *[]revokeItem) error {
|
|
||||||
// 发给 X 的 pending → rejected
|
|
||||||
rows, err := tx.Query(`
|
|
||||||
SELECT d.seq, d.pushed_at, m.id, m.sender_id, m.receipt
|
|
||||||
FROM deliveries d
|
|
||||||
JOIN messages m ON m.seq = d.seq
|
|
||||||
WHERE d.endpoint_id = ? AND d.state = 'pending'`, endpointID)
|
|
||||||
if err != nil {
|
|
||||||
return err
|
|
||||||
}
|
|
||||||
type pendRow struct {
|
|
||||||
seq int64
|
|
||||||
pushed sql.NullInt64
|
|
||||||
msgID string
|
|
||||||
senderID string
|
|
||||||
receipt int
|
|
||||||
}
|
|
||||||
var pending []pendRow
|
|
||||||
for rows.Next() {
|
|
||||||
var r pendRow
|
|
||||||
if scanErr := rows.Scan(&r.seq, &r.pushed, &r.msgID, &r.senderID, &r.receipt); scanErr != nil {
|
|
||||||
_ = rows.Close()
|
|
||||||
return scanErr
|
|
||||||
}
|
|
||||||
pending = append(pending, r)
|
|
||||||
}
|
|
||||||
_ = rows.Close()
|
|
||||||
if err = rows.Err(); err != nil {
|
|
||||||
return err
|
|
||||||
}
|
|
||||||
finalSeqs := map[int64]struct{}{}
|
|
||||||
for _, r := range pending {
|
|
||||||
if _, execErr := tx.Exec(`
|
|
||||||
UPDATE deliveries SET state = 'rejected', reason = ?, updated_at = ?
|
|
||||||
WHERE seq = ? AND endpoint_id = ? AND state = 'pending'`,
|
|
||||||
recvReason, nowMs, r.seq, endpointID); execErr != nil {
|
|
||||||
return execErr
|
|
||||||
}
|
|
||||||
if r.receipt != 0 {
|
|
||||||
if e := insertReceiptIfWantedTx(tx, r.senderID, r.msgID, endpointID, "rejected", recvReason, nowMs, true); e != nil {
|
|
||||||
return e
|
|
||||||
}
|
|
||||||
}
|
|
||||||
if r.pushed.Valid && revokes != nil {
|
|
||||||
*revokes = append(*revokes, revokeItem{
|
|
||||||
endpointID: endpointID, msgID: r.msgID, fromID: r.senderID, reason: recvReason,
|
|
||||||
})
|
|
||||||
}
|
|
||||||
finalSeqs[r.seq] = struct{}{}
|
|
||||||
}
|
|
||||||
|
|
||||||
// 发给 X 的 scheduled 单聊 → completed,要回执则写
|
|
||||||
srows, err := tx.Query(`
|
|
||||||
SELECT seq, id, sender_id, receipt FROM messages
|
|
||||||
WHERE dest_kind = 'endpoint' AND dest_id = ? AND state = 'scheduled'`, endpointID)
|
|
||||||
if err != nil {
|
|
||||||
return err
|
|
||||||
}
|
|
||||||
type schedRow struct {
|
|
||||||
seq int64
|
|
||||||
msgID string
|
|
||||||
senderID string
|
|
||||||
receipt int
|
|
||||||
}
|
|
||||||
var scheduledTo []schedRow
|
|
||||||
for srows.Next() {
|
|
||||||
var r schedRow
|
|
||||||
if scanErr := srows.Scan(&r.seq, &r.msgID, &r.senderID, &r.receipt); scanErr != nil {
|
|
||||||
_ = srows.Close()
|
|
||||||
return scanErr
|
|
||||||
}
|
|
||||||
scheduledTo = append(scheduledTo, r)
|
|
||||||
}
|
|
||||||
_ = srows.Close()
|
|
||||||
if err = srows.Err(); err != nil {
|
|
||||||
return err
|
|
||||||
}
|
|
||||||
for _, r := range scheduledTo {
|
|
||||||
if _, execErr := tx.Exec(`
|
|
||||||
UPDATE messages SET state = 'completed', reason = ? WHERE seq = ? AND state = 'scheduled'`,
|
|
||||||
recvReason, r.seq); execErr != nil {
|
|
||||||
return execErr
|
|
||||||
}
|
|
||||||
if _, execErr := tx.Exec(`DELETE FROM message_bodies WHERE seq = ?`, r.seq); execErr != nil {
|
|
||||||
return execErr
|
|
||||||
}
|
|
||||||
if r.receipt != 0 {
|
|
||||||
// 消息级作废:endpoint_id 空,state=rejected(DEVELOPMENT 6.4)
|
|
||||||
if e := insertReceiptIfWantedTx(tx, r.senderID, r.msgID, "", "rejected", recvReason, nowMs, true); e != nil {
|
|
||||||
return e
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
// X 发出的 scheduled → completed(sender_*),不写回执
|
|
||||||
outSched, err := tx.Query(`SELECT seq FROM messages WHERE sender_id = ? AND state = 'scheduled'`, endpointID)
|
|
||||||
if err != nil {
|
|
||||||
return err
|
|
||||||
}
|
|
||||||
var outSeqs []int64
|
|
||||||
for outSched.Next() {
|
|
||||||
var seq int64
|
|
||||||
if scanErr := outSched.Scan(&seq); scanErr != nil {
|
|
||||||
_ = outSched.Close()
|
|
||||||
return scanErr
|
|
||||||
}
|
|
||||||
outSeqs = append(outSeqs, seq)
|
|
||||||
}
|
|
||||||
_ = outSched.Close()
|
|
||||||
if err = outSched.Err(); err != nil {
|
|
||||||
return err
|
|
||||||
}
|
|
||||||
for _, seq := range outSeqs {
|
|
||||||
if _, execErr := tx.Exec(`
|
|
||||||
UPDATE messages SET state = 'completed', reason = ? WHERE seq = ? AND state = 'scheduled'`,
|
|
||||||
sendReason, seq); execErr != nil {
|
|
||||||
return execErr
|
|
||||||
}
|
|
||||||
if _, execErr := tx.Exec(`DELETE FROM message_bodies WHERE seq = ?`, seq); execErr != nil {
|
|
||||||
return execErr
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
// X 发出的消息的 pending 投递 → rejected(sender_*),不写回执
|
|
||||||
drows, err := tx.Query(`
|
|
||||||
SELECT d.seq, d.endpoint_id, d.pushed_at, m.id, m.sender_id
|
|
||||||
FROM deliveries d
|
|
||||||
JOIN messages m ON m.seq = d.seq
|
|
||||||
WHERE m.sender_id = ? AND d.state = 'pending'`, endpointID)
|
|
||||||
if err != nil {
|
|
||||||
return err
|
|
||||||
}
|
|
||||||
type outPend struct {
|
|
||||||
seq int64
|
|
||||||
endpointID string
|
|
||||||
pushed sql.NullInt64
|
|
||||||
msgID string
|
|
||||||
senderID string
|
|
||||||
}
|
|
||||||
var outPending []outPend
|
|
||||||
for drows.Next() {
|
|
||||||
var r outPend
|
|
||||||
if scanErr := drows.Scan(&r.seq, &r.endpointID, &r.pushed, &r.msgID, &r.senderID); scanErr != nil {
|
|
||||||
_ = drows.Close()
|
|
||||||
return scanErr
|
|
||||||
}
|
|
||||||
outPending = append(outPending, r)
|
|
||||||
}
|
|
||||||
_ = drows.Close()
|
|
||||||
if err = drows.Err(); err != nil {
|
|
||||||
return err
|
|
||||||
}
|
|
||||||
for _, r := range outPending {
|
|
||||||
if _, execErr := tx.Exec(`
|
|
||||||
UPDATE deliveries SET state = 'rejected', reason = ?, updated_at = ?
|
|
||||||
WHERE seq = ? AND endpoint_id = ? AND state = 'pending'`,
|
|
||||||
sendReason, nowMs, r.seq, r.endpointID); execErr != nil {
|
|
||||||
return execErr
|
|
||||||
}
|
|
||||||
if r.pushed.Valid && revokes != nil {
|
|
||||||
*revokes = append(*revokes, revokeItem{
|
|
||||||
endpointID: r.endpointID, msgID: r.msgID, fromID: r.senderID, reason: sendReason,
|
|
||||||
})
|
|
||||||
}
|
|
||||||
finalSeqs[r.seq] = struct{}{}
|
|
||||||
}
|
|
||||||
|
|
||||||
for seq := range finalSeqs {
|
|
||||||
if e := tryFinalizeTx(tx, seq); e != nil {
|
|
||||||
return e
|
|
||||||
}
|
|
||||||
}
|
|
||||||
return nil
|
|
||||||
}
|
|
||||||
|
|
||||||
func leaveAllGroupsTx(tx *sql.Tx, endpointID string, nowMs int64, revokes *[]revokeItem, notifies *[]groupNotify) error {
|
|
||||||
grows, err := tx.Query(`
|
|
||||||
SELECT g.id, g.owner_id
|
|
||||||
FROM groups g
|
|
||||||
JOIN group_members gm ON gm.group_id = g.id
|
|
||||||
WHERE gm.endpoint_id = ?`, endpointID)
|
|
||||||
if err != nil {
|
|
||||||
return err
|
|
||||||
}
|
|
||||||
type grow struct {
|
|
||||||
id string
|
|
||||||
owner string
|
|
||||||
}
|
|
||||||
var groups []grow
|
|
||||||
for grows.Next() {
|
|
||||||
var g grow
|
|
||||||
if scanErr := grows.Scan(&g.id, &g.owner); scanErr != nil {
|
|
||||||
_ = grows.Close()
|
|
||||||
return scanErr
|
|
||||||
}
|
|
||||||
groups = append(groups, g)
|
|
||||||
}
|
|
||||||
_ = grows.Close()
|
|
||||||
if err = grows.Err(); err != nil {
|
|
||||||
return err
|
|
||||||
}
|
|
||||||
|
|
||||||
for _, g := range groups {
|
|
||||||
members, memErr := loadMembersTx(tx, g.id)
|
|
||||||
if memErr != nil {
|
|
||||||
return memErr
|
|
||||||
}
|
|
||||||
if g.owner == endpointID {
|
|
||||||
others := withoutMember(members, endpointID)
|
|
||||||
if len(others) == 0 {
|
|
||||||
if e := voidGroupAllTx(tx, g.id, nowMs, revokes); e != nil {
|
|
||||||
return e
|
|
||||||
}
|
|
||||||
if _, e := tx.Exec(`DELETE FROM group_members WHERE group_id = ?`, g.id); e != nil {
|
|
||||||
return e
|
|
||||||
}
|
|
||||||
if _, e := tx.Exec(`DELETE FROM groups WHERE id = ?`, g.id); e != nil {
|
|
||||||
return e
|
|
||||||
}
|
|
||||||
if notifies != nil {
|
|
||||||
*notifies = append(*notifies, groupNotify{
|
|
||||||
recipients: members, groupID: g.id, event: eventDissolved, atMs: nowMs,
|
|
||||||
})
|
|
||||||
}
|
|
||||||
continue
|
|
||||||
}
|
|
||||||
newOwner, ownErr := earliestOtherMemberTx(tx, g.id, endpointID)
|
|
||||||
if ownErr != nil {
|
|
||||||
return ownErr
|
|
||||||
}
|
|
||||||
if _, e := tx.Exec(`UPDATE groups SET owner_id = ? WHERE id = ?`, newOwner, g.id); e != nil {
|
|
||||||
return e
|
|
||||||
}
|
|
||||||
if _, e := tx.Exec(`DELETE FROM group_members WHERE group_id = ? AND endpoint_id = ?`, g.id, endpointID); e != nil {
|
|
||||||
return e
|
|
||||||
}
|
|
||||||
if e := voidMemberDeliveriesTx(tx, g.id, endpointID, reasonLeftGroup, nowMs, revokes); e != nil {
|
|
||||||
return e
|
|
||||||
}
|
|
||||||
left := withoutMember(members, endpointID)
|
|
||||||
if notifies != nil {
|
|
||||||
*notifies = append(*notifies,
|
|
||||||
groupNotify{recipients: left, groupID: g.id, event: eventOwnerChanged, endpointID: newOwner, atMs: nowMs},
|
|
||||||
groupNotify{recipients: append(append([]string{}, left...), endpointID), groupID: g.id, event: eventMemberRemoved, endpointID: endpointID, atMs: nowMs},
|
|
||||||
)
|
|
||||||
}
|
|
||||||
continue
|
|
||||||
}
|
|
||||||
|
|
||||||
if _, e := tx.Exec(`DELETE FROM group_members WHERE group_id = ? AND endpoint_id = ?`, g.id, endpointID); e != nil {
|
|
||||||
return e
|
|
||||||
}
|
|
||||||
if e := voidMemberDeliveriesTx(tx, g.id, endpointID, reasonLeftGroup, nowMs, revokes); e != nil {
|
|
||||||
return e
|
|
||||||
}
|
|
||||||
left := withoutMember(members, endpointID)
|
|
||||||
if notifies != nil {
|
|
||||||
*notifies = append(*notifies, groupNotify{
|
|
||||||
recipients: append(append([]string{}, left...), endpointID),
|
|
||||||
groupID: g.id, event: eventLeft, endpointID: endpointID, atMs: nowMs,
|
|
||||||
})
|
|
||||||
}
|
|
||||||
}
|
|
||||||
return nil
|
|
||||||
}
|
|
||||||
|
|
||||||
func loadMembersTx(tx *sql.Tx, groupID string) ([]string, error) {
|
|
||||||
rows, err := tx.Query(`SELECT endpoint_id FROM group_members WHERE group_id = ? ORDER BY joined_at ASC, endpoint_id ASC`, groupID)
|
|
||||||
if err != nil {
|
|
||||||
return nil, err
|
|
||||||
}
|
|
||||||
defer func() { _ = rows.Close() }()
|
|
||||||
var out []string
|
|
||||||
for rows.Next() {
|
|
||||||
var id string
|
|
||||||
if e := rows.Scan(&id); e != nil {
|
|
||||||
return nil, e
|
|
||||||
}
|
|
||||||
out = append(out, id)
|
|
||||||
}
|
|
||||||
return out, rows.Err()
|
|
||||||
}
|
|
||||||
|
|
||||||
func earliestOtherMemberTx(tx *sql.Tx, groupID, exceptID string) (string, error) {
|
|
||||||
var id string
|
|
||||||
err := tx.QueryRow(`
|
|
||||||
SELECT endpoint_id FROM group_members
|
|
||||||
WHERE group_id = ? AND endpoint_id != ?
|
|
||||||
ORDER BY joined_at ASC, endpoint_id ASC
|
|
||||||
LIMIT 1`, groupID, exceptID).Scan(&id)
|
|
||||||
return id, err
|
|
||||||
}
|
|
||||||
|
|
||||||
func withoutMember(ids []string, drop string) []string {
|
|
||||||
out := make([]string, 0, len(ids))
|
|
||||||
for _, id := range ids {
|
|
||||||
if id != drop {
|
|
||||||
out = append(out, id)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
return out
|
|
||||||
}
|
|
||||||
|
|
||||||
// voidMemberDeliveriesTx 与 group 包同语义:退群成员的 pending 群投递改 rejected。
|
|
||||||
func voidMemberDeliveriesTx(tx *sql.Tx, groupID, endpointID, reason string, nowMs int64, revokes *[]revokeItem) error {
|
|
||||||
rows, err := tx.Query(`
|
|
||||||
SELECT d.seq, d.pushed_at, m.id, m.sender_id
|
|
||||||
FROM deliveries d
|
|
||||||
JOIN messages m ON m.seq = d.seq
|
|
||||||
WHERE d.endpoint_id = ? AND d.state = 'pending'
|
|
||||||
AND m.dest_kind = 'group' AND m.dest_id = ?`, endpointID, groupID)
|
|
||||||
if err != nil {
|
|
||||||
return err
|
|
||||||
}
|
|
||||||
type row struct {
|
|
||||||
seq int64
|
|
||||||
pushed sql.NullInt64
|
|
||||||
msgID string
|
|
||||||
senderID string
|
|
||||||
}
|
|
||||||
var list []row
|
|
||||||
for rows.Next() {
|
|
||||||
var r row
|
|
||||||
if scanErr := rows.Scan(&r.seq, &r.pushed, &r.msgID, &r.senderID); scanErr != nil {
|
|
||||||
_ = rows.Close()
|
|
||||||
return scanErr
|
|
||||||
}
|
|
||||||
list = append(list, r)
|
|
||||||
}
|
|
||||||
_ = rows.Close()
|
|
||||||
if err = rows.Err(); err != nil {
|
|
||||||
return err
|
|
||||||
}
|
|
||||||
for _, r := range list {
|
|
||||||
if _, execErr := tx.Exec(`
|
|
||||||
UPDATE deliveries SET state = 'rejected', reason = ?, updated_at = ?
|
|
||||||
WHERE seq = ? AND endpoint_id = ? AND state = 'pending'`,
|
|
||||||
reason, nowMs, r.seq, endpointID); execErr != nil {
|
|
||||||
return execErr
|
|
||||||
}
|
|
||||||
if r.pushed.Valid && revokes != nil {
|
|
||||||
*revokes = append(*revokes, revokeItem{
|
|
||||||
endpointID: endpointID, msgID: r.msgID, fromID: r.senderID, reason: reason,
|
|
||||||
})
|
|
||||||
}
|
|
||||||
if e := tryFinalizeTx(tx, r.seq); e != nil {
|
|
||||||
return e
|
|
||||||
}
|
|
||||||
}
|
|
||||||
return nil
|
|
||||||
}
|
|
||||||
|
|
||||||
func voidGroupAllTx(tx *sql.Tx, groupID string, nowMs int64, revokes *[]revokeItem) error {
|
|
||||||
rows, err := tx.Query(`
|
|
||||||
SELECT d.seq, d.endpoint_id, d.pushed_at, m.id, m.sender_id
|
|
||||||
FROM deliveries d
|
|
||||||
JOIN messages m ON m.seq = d.seq
|
|
||||||
WHERE d.state = 'pending' AND m.dest_kind = 'group' AND m.dest_id = ?`, groupID)
|
|
||||||
if err != nil {
|
|
||||||
return err
|
|
||||||
}
|
|
||||||
type drow struct {
|
|
||||||
seq int64
|
|
||||||
endpointID string
|
|
||||||
pushed sql.NullInt64
|
|
||||||
msgID string
|
|
||||||
senderID string
|
|
||||||
}
|
|
||||||
var dlist []drow
|
|
||||||
for rows.Next() {
|
|
||||||
var r drow
|
|
||||||
if scanErr := rows.Scan(&r.seq, &r.endpointID, &r.pushed, &r.msgID, &r.senderID); scanErr != nil {
|
|
||||||
_ = rows.Close()
|
|
||||||
return scanErr
|
|
||||||
}
|
|
||||||
dlist = append(dlist, r)
|
|
||||||
}
|
|
||||||
_ = rows.Close()
|
|
||||||
if err = rows.Err(); err != nil {
|
|
||||||
return err
|
|
||||||
}
|
|
||||||
for _, r := range dlist {
|
|
||||||
if _, execErr := tx.Exec(`
|
|
||||||
UPDATE deliveries SET state = 'rejected', reason = ?, updated_at = ?
|
|
||||||
WHERE seq = ? AND endpoint_id = ? AND state = 'pending'`,
|
|
||||||
reasonGroupDissolved, nowMs, r.seq, r.endpointID); execErr != nil {
|
|
||||||
return execErr
|
|
||||||
}
|
|
||||||
if r.pushed.Valid && revokes != nil {
|
|
||||||
*revokes = append(*revokes, revokeItem{
|
|
||||||
endpointID: r.endpointID, msgID: r.msgID, fromID: r.senderID, reason: reasonGroupDissolved,
|
|
||||||
})
|
|
||||||
}
|
|
||||||
if e := tryFinalizeTx(tx, r.seq); e != nil {
|
|
||||||
return e
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
srows, err := tx.Query(`
|
|
||||||
SELECT seq, id, sender_id, receipt FROM messages
|
|
||||||
WHERE dest_kind = 'group' AND dest_id = ? AND state = 'scheduled'`, groupID)
|
|
||||||
if err != nil {
|
|
||||||
return err
|
|
||||||
}
|
|
||||||
type srow struct {
|
|
||||||
seq int64
|
|
||||||
msgID string
|
|
||||||
senderID string
|
|
||||||
receipt int
|
|
||||||
}
|
|
||||||
var slist []srow
|
|
||||||
for srows.Next() {
|
|
||||||
var r srow
|
|
||||||
if scanErr := srows.Scan(&r.seq, &r.msgID, &r.senderID, &r.receipt); scanErr != nil {
|
|
||||||
_ = srows.Close()
|
|
||||||
return scanErr
|
|
||||||
}
|
|
||||||
slist = append(slist, r)
|
|
||||||
}
|
|
||||||
_ = srows.Close()
|
|
||||||
if err = srows.Err(); err != nil {
|
|
||||||
return err
|
|
||||||
}
|
|
||||||
for _, r := range slist {
|
|
||||||
if _, execErr := tx.Exec(`
|
|
||||||
UPDATE messages SET state = 'completed', reason = ? WHERE seq = ? AND state = 'scheduled'`,
|
|
||||||
reasonGroupDissolved, r.seq); execErr != nil {
|
|
||||||
return execErr
|
|
||||||
}
|
|
||||||
if _, execErr := tx.Exec(`DELETE FROM message_bodies WHERE seq = ?`, r.seq); execErr != nil {
|
|
||||||
return execErr
|
|
||||||
}
|
|
||||||
if r.receipt != 0 {
|
|
||||||
if e := insertReceiptIfWantedTx(tx, r.senderID, r.msgID, "", "rejected", reasonGroupDissolved, nowMs, true); e != nil {
|
|
||||||
return e
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
return nil
|
|
||||||
}
|
|
||||||
|
|
||||||
func insertReceiptIfWantedTx(tx *sql.Tx, senderID, msgID, endpointID, state, reason string, nowMs int64, alreadyWanted bool) error {
|
|
||||||
if !alreadyWanted {
|
|
||||||
return nil
|
|
||||||
}
|
|
||||||
var one int
|
|
||||||
err := tx.QueryRow(`SELECT 1 FROM endpoints WHERE id = ?`, senderID).Scan(&one)
|
|
||||||
if errors.Is(err, sql.ErrNoRows) {
|
|
||||||
return nil
|
|
||||||
}
|
|
||||||
if err != nil {
|
|
||||||
return err
|
|
||||||
}
|
|
||||||
_, err = tx.Exec(`
|
|
||||||
INSERT INTO receipts(sender_id, msg_id, endpoint_id, state, reason, created_at, acked)
|
|
||||||
VALUES(?,?,?,?,?,?,0)`, senderID, msgID, endpointID, state, reason, nowMs)
|
|
||||||
return err
|
|
||||||
}
|
|
||||||
|
|
||||||
func tryFinalizeTx(tx *sql.Tx, seq int64) error {
|
|
||||||
var n int
|
|
||||||
if err := tx.QueryRow(`SELECT COUNT(*) FROM deliveries WHERE seq = ? AND state = 'pending'`, seq).Scan(&n); err != nil {
|
|
||||||
return err
|
|
||||||
}
|
|
||||||
if n > 0 {
|
|
||||||
return nil
|
|
||||||
}
|
|
||||||
var state string
|
|
||||||
if err := tx.QueryRow(`SELECT state FROM messages WHERE seq = ?`, seq).Scan(&state); err != nil {
|
|
||||||
if errors.Is(err, sql.ErrNoRows) {
|
|
||||||
return nil
|
|
||||||
}
|
|
||||||
return err
|
|
||||||
}
|
|
||||||
if state == "completed" {
|
|
||||||
return nil
|
|
||||||
}
|
|
||||||
if _, err := tx.Exec(`UPDATE messages SET state = 'completed' WHERE seq = ?`, seq); err != nil {
|
|
||||||
return err
|
|
||||||
}
|
|
||||||
_, err := tx.Exec(`DELETE FROM message_bodies WHERE seq = ?`, seq)
|
|
||||||
return err
|
|
||||||
}
|
|
||||||
|
|
||||||
func (a *App) publishRevokes(ctx context.Context, items []revokeItem) {
|
|
||||||
if a.down == nil || len(items) == 0 {
|
|
||||||
return
|
|
||||||
}
|
|
||||||
for _, it := range items {
|
|
||||||
frame := protocol.Revoked{
|
|
||||||
V: protocol.Version, Type: protocol.TypeRevoked,
|
|
||||||
ID: it.msgID, From: it.fromID, Reason: it.reason,
|
|
||||||
}
|
|
||||||
payload, encErr := encodeFrame(frame)
|
|
||||||
if encErr != nil {
|
|
||||||
continue
|
|
||||||
}
|
|
||||||
_ = a.down.PublishDown(ctx, it.endpointID, "", payload, port.PublishOpts{QoS: 1})
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
func (a *App) publishGroupEvents(ctx context.Context, items []groupNotify) {
|
|
||||||
if a.down == nil || len(items) == 0 {
|
|
||||||
return
|
|
||||||
}
|
|
||||||
for _, it := range items {
|
|
||||||
frame := protocol.GroupEvent{
|
|
||||||
V: protocol.Version, Type: protocol.TypeGroupEvent,
|
|
||||||
GroupID: it.groupID, Event: it.event, EndpointID: it.endpointID, AtMs: it.atMs,
|
|
||||||
}
|
|
||||||
payload, encErr := encodeFrame(frame)
|
|
||||||
if encErr != nil {
|
|
||||||
continue
|
|
||||||
}
|
|
||||||
seen := map[string]struct{}{}
|
|
||||||
for _, id := range it.recipients {
|
|
||||||
if _, ok := seen[id]; ok {
|
|
||||||
continue
|
|
||||||
}
|
|
||||||
seen[id] = struct{}{}
|
|
||||||
_ = a.down.PublishDown(ctx, id, "", payload, port.PublishOpts{QoS: 0})
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
func encodeFrame(v any) ([]byte, error) {
|
|
||||||
var buf bytes.Buffer
|
|
||||||
if err := protocol.Encode(&buf, v); err != nil {
|
|
||||||
return nil, err
|
|
||||||
}
|
|
||||||
return buf.Bytes(), nil
|
|
||||||
}
|
|
||||||
@@ -1,386 +0,0 @@
|
|||||||
package identity_test
|
|
||||||
|
|
||||||
import (
|
|
||||||
"context"
|
|
||||||
"database/sql"
|
|
||||||
"io"
|
|
||||||
"net/http"
|
|
||||||
"net/http/cookiejar"
|
|
||||||
"net/http/httptest"
|
|
||||||
"path/filepath"
|
|
||||||
"strings"
|
|
||||||
"testing"
|
|
||||||
"time"
|
|
||||||
|
|
||||||
"git.asio.asia/nixevol/NixMsg/internal/admin"
|
|
||||||
"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/auth"
|
|
||||||
"git.asio.asia/nixevol/NixMsg/internal/config"
|
|
||||||
"git.asio.asia/nixevol/NixMsg/internal/protocol"
|
|
||||||
"git.asio.asia/nixevol/NixMsg/internal/store"
|
|
||||||
)
|
|
||||||
|
|
||||||
func openLifecycle(t *testing.T) (*identity.App, *message.App, *store.DB) {
|
|
||||||
t.Helper()
|
|
||||||
dir := t.TempDir()
|
|
||||||
db, err := store.Open(filepath.Join(dir, "data"), "FULL")
|
|
||||||
if err != nil {
|
|
||||||
t.Fatal(err)
|
|
||||||
}
|
|
||||||
t.Cleanup(func() { _ = db.Close() })
|
|
||||||
fixed := time.UnixMilli(1_700_000_000_000)
|
|
||||||
idApp := identity.New(identity.Config{
|
|
||||||
DB: db,
|
|
||||||
Hash: auth.NewStubHashPool(),
|
|
||||||
Locks: auth.NewStubLoginLocks(),
|
|
||||||
Sessions: auth.NewSessionTokens(),
|
|
||||||
MaxScheduleSeconds: int64(config.Default().Limits.MaxScheduleSeconds),
|
|
||||||
Now: func() time.Time { return fixed },
|
|
||||||
ConnControl: &port.StubConnControl{},
|
|
||||||
Downlink: &port.StubDownlink{},
|
|
||||||
})
|
|
||||||
lim := message.LimitsFromConfig(config.Default().Limits)
|
|
||||||
lim.RequestsPerSecond = 0
|
|
||||||
lim.RecordRetentionDays = 7
|
|
||||||
msgApp := message.New(db, lim, auth.NewStubHashPool(),
|
|
||||||
message.WithNow(func() time.Time { return fixed }),
|
|
||||||
message.WithLocks(auth.NewStubLoginLocks()),
|
|
||||||
)
|
|
||||||
return idApp, msgApp, db
|
|
||||||
}
|
|
||||||
|
|
||||||
func insertEPFull(t *testing.T, db *store.DB, id string) {
|
|
||||||
t.Helper()
|
|
||||||
err := db.Queue.Do(context.Background(), func(tx *sql.Tx) error {
|
|
||||||
_, e := tx.Exec(`
|
|
||||||
INSERT INTO endpoints(id, name, login_hash, talk_hash, talk_version, default_delay_ms, enabled, created_at, offline_since)
|
|
||||||
VALUES(?,?,?,?,0,0,1,?,?)`, id, id, "stub$login", nil, 1_700_000_000_000, 1_700_000_000_000)
|
|
||||||
return e
|
|
||||||
})
|
|
||||||
if err != nil {
|
|
||||||
t.Fatal(err)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
func TestF01DisableVoidsScheduledAndRejectsNew(t *testing.T) {
|
|
||||||
t.Parallel()
|
|
||||||
idApp, msgApp, db := openLifecycle(t)
|
|
||||||
ctx := context.Background()
|
|
||||||
insertEPFull(t, db, "alice")
|
|
||||||
insertEPFull(t, db, "bob")
|
|
||||||
|
|
||||||
future := int64(1_700_000_000_000 + 3600_000)
|
|
||||||
receipt := true
|
|
||||||
toBob := &protocol.Send{
|
|
||||||
V: protocol.Version, Type: protocol.TypeSend, RID: "1", ID: "m-to-bob",
|
|
||||||
To: protocol.Target{Kind: protocol.TargetEndpoint, ID: "bob"},
|
|
||||||
Body: protocol.Body{Enc: protocol.EncUTF8, Data: "hi"},
|
|
||||||
SendAtMs: &future,
|
|
||||||
Receipt: &receipt,
|
|
||||||
}
|
|
||||||
if _, err := msgApp.Submit(ctx, "alice", port.ConnInfo{EndpointID: "alice"}, toBob); err != nil {
|
|
||||||
t.Fatal(err)
|
|
||||||
}
|
|
||||||
fromBob := &protocol.Send{
|
|
||||||
V: protocol.Version, Type: protocol.TypeSend, RID: "2", ID: "m-from-bob",
|
|
||||||
To: protocol.Target{Kind: protocol.TargetEndpoint, ID: "alice"},
|
|
||||||
Body: protocol.Body{Enc: protocol.EncUTF8, Data: "bye"},
|
|
||||||
SendAtMs: &future,
|
|
||||||
}
|
|
||||||
if _, err := msgApp.Submit(ctx, "bob", port.ConnInfo{EndpointID: "bob"}, fromBob); err != nil {
|
|
||||||
t.Fatal(err)
|
|
||||||
}
|
|
||||||
|
|
||||||
if err := idApp.Disable(ctx, "bob"); err != nil {
|
|
||||||
t.Fatal(err)
|
|
||||||
}
|
|
||||||
|
|
||||||
var enabled int
|
|
||||||
if err := db.Read.QueryRow(`SELECT enabled FROM endpoints WHERE id='bob'`).Scan(&enabled); err != nil {
|
|
||||||
t.Fatal(err)
|
|
||||||
}
|
|
||||||
if enabled != 0 {
|
|
||||||
t.Fatalf("enabled=%d", enabled)
|
|
||||||
}
|
|
||||||
|
|
||||||
var toState, toReason string
|
|
||||||
if err := db.Read.QueryRow(`SELECT state, reason FROM messages WHERE sender_id='alice' AND id='m-to-bob'`).
|
|
||||||
Scan(&toState, &toReason); err != nil {
|
|
||||||
t.Fatal(err)
|
|
||||||
}
|
|
||||||
if toState != "completed" || toReason != "endpoint_disabled" {
|
|
||||||
t.Fatalf("to bob: state=%s reason=%s", toState, toReason)
|
|
||||||
}
|
|
||||||
var fromState, fromReason string
|
|
||||||
if err := db.Read.QueryRow(`SELECT state, reason FROM messages WHERE sender_id='bob' AND id='m-from-bob'`).
|
|
||||||
Scan(&fromState, &fromReason); err != nil {
|
|
||||||
t.Fatal(err)
|
|
||||||
}
|
|
||||||
if fromState != "completed" || fromReason != "sender_disabled" {
|
|
||||||
t.Fatalf("from bob: state=%s reason=%s", fromState, fromReason)
|
|
||||||
}
|
|
||||||
|
|
||||||
newSend := &protocol.Send{
|
|
||||||
V: protocol.Version, Type: protocol.TypeSend, RID: "3", ID: "m-new",
|
|
||||||
To: protocol.Target{Kind: protocol.TargetEndpoint, ID: "bob"},
|
|
||||||
Body: protocol.Body{Enc: protocol.EncUTF8, Data: "x"},
|
|
||||||
}
|
|
||||||
_, err := msgApp.Submit(ctx, "alice", port.ConnInfo{EndpointID: "alice"}, newSend)
|
|
||||||
if protoCode(err) != protocol.CodeEndpointDisabled {
|
|
||||||
t.Fatalf("want endpoint_disabled got %v", err)
|
|
||||||
}
|
|
||||||
|
|
||||||
if err := idApp.Enable(ctx, "bob"); err != nil {
|
|
||||||
t.Fatal(err)
|
|
||||||
}
|
|
||||||
if err := db.Read.QueryRow(`SELECT state FROM messages WHERE id='m-to-bob'`).Scan(&toState); err != nil {
|
|
||||||
t.Fatal(err)
|
|
||||||
}
|
|
||||||
if toState != "completed" {
|
|
||||||
t.Fatalf("voided message restored? %s", toState)
|
|
||||||
}
|
|
||||||
okSend := &protocol.Send{
|
|
||||||
V: protocol.Version, Type: protocol.TypeSend, RID: "4", ID: "m-after",
|
|
||||||
To: protocol.Target{Kind: protocol.TargetEndpoint, ID: "bob"},
|
|
||||||
Body: protocol.Body{Enc: protocol.EncUTF8, Data: "ok"},
|
|
||||||
}
|
|
||||||
if _, err := msgApp.Submit(ctx, "alice", port.ConnInfo{EndpointID: "alice"}, okSend); err != nil {
|
|
||||||
t.Fatal(err)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
func TestF01DeleteOwnerTransfersEarliest(t *testing.T) {
|
|
||||||
t.Parallel()
|
|
||||||
idApp, _, db := openLifecycle(t)
|
|
||||||
ctx := context.Background()
|
|
||||||
insertEPFull(t, db, "owner")
|
|
||||||
insertEPFull(t, db, "early")
|
|
||||||
insertEPFull(t, db, "late")
|
|
||||||
|
|
||||||
err := db.Queue.Do(ctx, func(tx *sql.Tx) error {
|
|
||||||
if _, e := tx.Exec(`INSERT INTO groups(id, name, owner_id, created_at) VALUES('g1','群','owner',?)`, 1_700_000_000_000); e != nil {
|
|
||||||
return e
|
|
||||||
}
|
|
||||||
_, e := tx.Exec(`INSERT INTO group_members(group_id, endpoint_id, joined_at) VALUES
|
|
||||||
('g1','owner',100),('g1','early',200),('g1','late',300)`)
|
|
||||||
return e
|
|
||||||
})
|
|
||||||
if err != nil {
|
|
||||||
t.Fatal(err)
|
|
||||||
}
|
|
||||||
|
|
||||||
if err := idApp.Delete(ctx, "owner"); err != nil {
|
|
||||||
t.Fatal(err)
|
|
||||||
}
|
|
||||||
|
|
||||||
var owner string
|
|
||||||
if err := db.Read.QueryRow(`SELECT owner_id FROM groups WHERE id='g1'`).Scan(&owner); err != nil {
|
|
||||||
t.Fatal(err)
|
|
||||||
}
|
|
||||||
if owner != "early" {
|
|
||||||
t.Fatalf("want earliest other member early, got %s", owner)
|
|
||||||
}
|
|
||||||
var n int
|
|
||||||
if err := db.Read.QueryRow(`SELECT COUNT(*) FROM group_members WHERE group_id='g1' AND endpoint_id='owner'`).Scan(&n); err != nil {
|
|
||||||
t.Fatal(err)
|
|
||||||
}
|
|
||||||
if n != 0 {
|
|
||||||
t.Fatal("owner should have left the group")
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
func TestF01DeleteReopenNoOldReceipts(t *testing.T) {
|
|
||||||
t.Parallel()
|
|
||||||
idApp, msgApp, db := openLifecycle(t)
|
|
||||||
ctx := context.Background()
|
|
||||||
insertEPFull(t, db, "alice")
|
|
||||||
insertEPFull(t, db, "bob")
|
|
||||||
|
|
||||||
receipt := true
|
|
||||||
send := &protocol.Send{
|
|
||||||
V: protocol.Version, Type: protocol.TypeSend, RID: "1", ID: "m1",
|
|
||||||
To: protocol.Target{Kind: protocol.TargetEndpoint, ID: "alice"},
|
|
||||||
Body: protocol.Body{Enc: protocol.EncUTF8, Data: "hi"},
|
|
||||||
Receipt: &receipt,
|
|
||||||
}
|
|
||||||
if _, err := msgApp.Submit(ctx, "bob", port.ConnInfo{EndpointID: "bob"}, send); err != nil {
|
|
||||||
t.Fatal(err)
|
|
||||||
}
|
|
||||||
err := db.Queue.Do(ctx, func(tx *sql.Tx) error {
|
|
||||||
_, e := tx.Exec(`
|
|
||||||
INSERT INTO receipts(sender_id, msg_id, endpoint_id, state, reason, created_at, acked)
|
|
||||||
VALUES('bob','m1','alice','accepted','',?,0)`, 1_700_000_000_000)
|
|
||||||
return e
|
|
||||||
})
|
|
||||||
if err != nil {
|
|
||||||
t.Fatal(err)
|
|
||||||
}
|
|
||||||
|
|
||||||
if err := idApp.Delete(ctx, "bob"); err != nil {
|
|
||||||
t.Fatal(err)
|
|
||||||
}
|
|
||||||
|
|
||||||
var n int
|
|
||||||
if err := db.Read.QueryRow(`SELECT COUNT(*) FROM receipts WHERE sender_id='bob'`).Scan(&n); err != nil {
|
|
||||||
t.Fatal(err)
|
|
||||||
}
|
|
||||||
if n != 0 {
|
|
||||||
t.Fatalf("old receipts should be gone, got %d", n)
|
|
||||||
}
|
|
||||||
|
|
||||||
insertEPFull(t, db, "bob")
|
|
||||||
if err := db.Read.QueryRow(`SELECT COUNT(*) FROM receipts WHERE sender_id='bob'`).Scan(&n); err != nil {
|
|
||||||
t.Fatal(err)
|
|
||||||
}
|
|
||||||
if n != 0 {
|
|
||||||
t.Fatalf("reopened endpoint must not inherit receipts, got %d", n)
|
|
||||||
}
|
|
||||||
if err := db.Read.QueryRow(`SELECT COUNT(*) FROM messages WHERE sender_id='bob'`).Scan(&n); err != nil {
|
|
||||||
t.Fatal(err)
|
|
||||||
}
|
|
||||||
if n != 0 {
|
|
||||||
t.Fatalf("reopened endpoint must not inherit messages, got %d", n)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
func TestAdminDisableDeleteHTTP(t *testing.T) {
|
|
||||||
t.Parallel()
|
|
||||||
dir := t.TempDir()
|
|
||||||
db, err := store.Open(filepath.Join(dir, "data"), "FULL")
|
|
||||||
if err != nil {
|
|
||||||
t.Fatal(err)
|
|
||||||
}
|
|
||||||
t.Cleanup(func() { _ = db.Close() })
|
|
||||||
|
|
||||||
fixed := time.UnixMilli(1_700_000_000_000)
|
|
||||||
hash := auth.NewStubHashPool()
|
|
||||||
if seedErr := admin.SeedAdminPassword(context.Background(), db, hash, "adminpassword1"); seedErr != nil {
|
|
||||||
t.Fatal(seedErr)
|
|
||||||
}
|
|
||||||
idApp := identity.New(identity.Config{
|
|
||||||
DB: db, Hash: hash, Locks: auth.NewStubLoginLocks(), Sessions: auth.NewSessionTokens(),
|
|
||||||
MaxScheduleSeconds: 86400, Now: func() time.Time { return fixed },
|
|
||||||
})
|
|
||||||
lim := message.LimitsFromConfig(config.Default().Limits)
|
|
||||||
lim.RequestsPerSecond = 0
|
|
||||||
msgApp := message.New(db, lim, hash,
|
|
||||||
message.WithNow(func() time.Time { return fixed }),
|
|
||||||
message.WithLocks(auth.NewStubLoginLocks()),
|
|
||||||
)
|
|
||||||
|
|
||||||
kick := &lifecycleKick{}
|
|
||||||
h := admin.New(admin.Deps{
|
|
||||||
DB: db, Hash: hash, Tokens: admin.NewRandomAPITokens(),
|
|
||||||
Locks: admin.NewMemoryLoginLocks(), KickEndpoint: kick.Kick, Identity: idApp,
|
|
||||||
})
|
|
||||||
srv := httptest.NewServer(h)
|
|
||||||
t.Cleanup(srv.Close)
|
|
||||||
|
|
||||||
jar, _ := cookiejar.New(nil)
|
|
||||||
client := &http.Client{Jar: jar}
|
|
||||||
loginRes, err := client.Post(srv.URL+"/api/admin/login", "application/json",
|
|
||||||
strings.NewReader(`{"username":"admin","password":"adminpassword1"}`))
|
|
||||||
if err != nil {
|
|
||||||
t.Fatal(err)
|
|
||||||
}
|
|
||||||
_ = loginRes.Body.Close()
|
|
||||||
if loginRes.StatusCode != 200 {
|
|
||||||
t.Fatalf("login %d", loginRes.StatusCode)
|
|
||||||
}
|
|
||||||
|
|
||||||
createEP := func(id string) {
|
|
||||||
t.Helper()
|
|
||||||
req, _ := http.NewRequest(http.MethodPost, srv.URL+"/api/admin/endpoints",
|
|
||||||
strings.NewReader(`{"id":"`+id+`","name":"`+id+`","login_password":"password12"}`))
|
|
||||||
req.Header.Set("Content-Type", "application/json")
|
|
||||||
req.Header.Set("X-Nixmsg-Request", "1")
|
|
||||||
res, e := client.Do(req)
|
|
||||||
if e != nil {
|
|
||||||
t.Fatal(e)
|
|
||||||
}
|
|
||||||
raw, _ := io.ReadAll(res.Body)
|
|
||||||
_ = res.Body.Close()
|
|
||||||
if res.StatusCode != 200 {
|
|
||||||
t.Fatalf("create %s: %d %s", id, res.StatusCode, raw)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
createEP("alice")
|
|
||||||
createEP("bob")
|
|
||||||
createEP("carol")
|
|
||||||
|
|
||||||
ctx := context.Background()
|
|
||||||
future := fixed.UnixMilli() + 3600_000
|
|
||||||
toBob := &protocol.Send{
|
|
||||||
V: protocol.Version, Type: protocol.TypeSend, RID: "1", ID: "sched-bob",
|
|
||||||
To: protocol.Target{Kind: protocol.TargetEndpoint, ID: "bob"},
|
|
||||||
Body: protocol.Body{Enc: protocol.EncUTF8, Data: "x"},
|
|
||||||
SendAtMs: &future,
|
|
||||||
}
|
|
||||||
if _, subErr := msgApp.Submit(ctx, "alice", port.ConnInfo{EndpointID: "alice"}, toBob); subErr != nil {
|
|
||||||
t.Fatal(subErr)
|
|
||||||
}
|
|
||||||
|
|
||||||
req, _ := http.NewRequest(http.MethodPost, srv.URL+"/api/admin/endpoints/batch",
|
|
||||||
strings.NewReader(`{"ids":["bob"],"action":"disable"}`))
|
|
||||||
req.Header.Set("Content-Type", "application/json")
|
|
||||||
req.Header.Set("X-Nixmsg-Request", "1")
|
|
||||||
res, err := client.Do(req)
|
|
||||||
if err != nil {
|
|
||||||
t.Fatal(err)
|
|
||||||
}
|
|
||||||
body, _ := io.ReadAll(res.Body)
|
|
||||||
_ = res.Body.Close()
|
|
||||||
if res.StatusCode != 200 {
|
|
||||||
t.Fatalf("disable: %d %s", res.StatusCode, body)
|
|
||||||
}
|
|
||||||
var st string
|
|
||||||
if scanErr := db.Read.QueryRow(`SELECT state FROM messages WHERE id='sched-bob'`).Scan(&st); scanErr != nil {
|
|
||||||
t.Fatal(scanErr)
|
|
||||||
}
|
|
||||||
if st != "completed" {
|
|
||||||
t.Fatalf("scheduled should be voided, got %s", st)
|
|
||||||
}
|
|
||||||
|
|
||||||
err = db.Queue.Do(ctx, func(tx *sql.Tx) error {
|
|
||||||
if _, e := tx.Exec(`INSERT INTO groups(id, name, owner_id, created_at) VALUES('ghttp','G','bob',?)`, fixed.UnixMilli()); e != nil {
|
|
||||||
return e
|
|
||||||
}
|
|
||||||
_, e := tx.Exec(`INSERT INTO group_members(group_id, endpoint_id, joined_at) VALUES
|
|
||||||
('ghttp','bob',1),('ghttp','alice',2),('ghttp','carol',3)`)
|
|
||||||
return e
|
|
||||||
})
|
|
||||||
if err != nil {
|
|
||||||
t.Fatal(err)
|
|
||||||
}
|
|
||||||
|
|
||||||
// bob 已停用,需先启用才能作为「仍存在的群主」再删除?删除不要求 enabled。
|
|
||||||
req, _ = http.NewRequest(http.MethodDelete, srv.URL+"/api/admin/endpoints/bob", nil)
|
|
||||||
req.Header.Set("X-Nixmsg-Request", "1")
|
|
||||||
res, err = client.Do(req)
|
|
||||||
if err != nil {
|
|
||||||
t.Fatal(err)
|
|
||||||
}
|
|
||||||
raw, _ := io.ReadAll(res.Body)
|
|
||||||
_ = res.Body.Close()
|
|
||||||
if res.StatusCode != 200 {
|
|
||||||
t.Fatalf("delete: %d %s", res.StatusCode, raw)
|
|
||||||
}
|
|
||||||
var owner string
|
|
||||||
if err := db.Read.QueryRow(`SELECT owner_id FROM groups WHERE id='ghttp'`).Scan(&owner); err != nil {
|
|
||||||
t.Fatal(err)
|
|
||||||
}
|
|
||||||
if owner != "alice" {
|
|
||||||
t.Fatalf("want alice as new owner, got %s", owner)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
type lifecycleKick struct {
|
|
||||||
Calls []string
|
|
||||||
}
|
|
||||||
|
|
||||||
func (k *lifecycleKick) Kick(_ context.Context, endpointID string) (bool, error) {
|
|
||||||
k.Calls = append(k.Calls, endpointID)
|
|
||||||
return true, nil
|
|
||||||
}
|
|
||||||
@@ -215,3 +215,8 @@ WHERE id = ?`, endpointID)
|
|||||||
}
|
}
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// Disable / Enable / Delete 属 I5,此处保留未实现。
|
||||||
|
func (a *App) Disable(context.Context, string) error { return ErrNotImplemented }
|
||||||
|
func (a *App) Enable(context.Context, string) error { return ErrNotImplemented }
|
||||||
|
func (a *App) Delete(context.Context, string) error { return ErrNotImplemented }
|
||||||
|
|||||||
@@ -1,22 +0,0 @@
|
|||||||
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}
|
|
||||||
}
|
|
||||||
Reference in New Issue
Block a user