Compare commits

..
15 changed files with 1574 additions and 80 deletions
+241 -74
View File
@@ -4,6 +4,7 @@ import (
"context"
"errors"
"fmt"
"io/fs"
"log/slog"
"net"
"net/http"
@@ -14,7 +15,18 @@ import (
"syscall"
"time"
"git.asio.asia/nixevol/NixMsg/internal/admin"
"git.asio.asia/nixevol/NixMsg/internal/app/group"
"git.asio.asia/nixevol/NixMsg/internal/app/identity"
"git.asio.asia/nixevol/NixMsg/internal/app/message"
"git.asio.asia/nixevol/NixMsg/internal/app/port"
"git.asio.asia/nixevol/NixMsg/internal/app/presence"
"git.asio.asia/nixevol/NixMsg/internal/auth"
"git.asio.asia/nixevol/NixMsg/internal/broker"
"git.asio.asia/nixevol/NixMsg/internal/config"
"git.asio.asia/nixevol/NixMsg/internal/httpx"
"git.asio.asia/nixevol/NixMsg/internal/listener"
"git.asio.asia/nixevol/NixMsg/internal/metrics"
"git.asio.asia/nixevol/NixMsg/internal/store"
"git.asio.asia/nixevol/NixMsg/web"
)
@@ -25,20 +37,16 @@ func cmdServe(_ []string) error {
if err != nil {
return err
}
// P-WIRE-BEGIN
if err := cfg.Validate(); err != nil {
return err
}
setupJSONLogger(cfg.Log)
// P-WIRE-END
ctx, stop := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM)
defer stop()
return runServe(ctx, cfg)
}
func runServe(ctx context.Context, cfg config.Config) error {
deps := wire()
if err := os.MkdirAll(cfg.DataDir, 0o755); err != nil {
return fmt.Errorf("mkdir data_dir: %w", err)
}
@@ -49,7 +57,6 @@ func runServe(ctx context.Context, cfg config.Config) error {
}
defer func() { _ = db.Close() }()
// P-WIRE-BEGIN
ok, err := store.HasAdminPassword(ctx, db.Write)
if err != nil {
return err
@@ -57,87 +64,250 @@ func runServe(ctx context.Context, cfg config.Config) error {
if !ok {
return errors.New("admin password not initialized; run: nixmsg admin init")
}
// P-WIRE-END
// 启动恢复入口已挂上(假实现为空操作);M 线替换 message.Service 后生效。
if recoverErr := deps.Messages.RecoverOnStart(ctx); recoverErr != nil {
hashPool := auth.NewPool()
sessionTokens := auth.NewSessionTokens()
apiTokens := auth.NewAPITokens()
loginLocks := auth.NewLoginLocks()
memConns := message.NewMemoryConns()
msgLim := message.LimitsFromFullConfig(cfg)
login := broker.NewLogin(broker.LoginOptions{
DB: db,
Pool: hashPool,
Tokens: sessionTokens,
Locks: loginLocks,
IdleDays: cfg.SessionIdleDays,
})
// 先用空下行建 message,拿到 App 指针;broker 建好后再 WithDownlink 重建并回填 uplink.msg。
msgApp := message.New(db, msgLim, hashPool,
message.WithLocks(loginLocks),
message.WithConnRegistry(memConns),
)
uplink := &appUplink{
msg: msgApp,
conns: memConns,
next: port.StubUplinkHandler{},
log: slog.Default(),
}
sess := broker.NewSession(broker.SessionOptions{
Login: login,
Inner: uplink,
Limits: broker.HelloLimits{
MaxBodyBytes: cfg.Limits.MaxBodyBytes,
MaxMetaBytes: cfg.Limits.MaxMetaBytes,
MaxFrameBytes: cfg.Limits.MaxFrameBytes,
MaxTTLSeconds: int64(cfg.Limits.MaxTTLSeconds),
MaxScheduleSeconds: int64(cfg.Limits.MaxScheduleSeconds),
AckTimeoutSeconds: int64(cfg.Limits.AckTimeoutSeconds),
ServerVersion: Version,
},
Logger: slog.Default(),
})
brk, err := broker.New(broker.Options{
Authenticator: login,
Uplink: sess,
Logger: slog.Default(),
OnPublishDropped: func(dropCtx context.Context, endpointID string, connID port.ConnID, payload []byte) {
if dropErr := msgApp.OnPublishDropped(dropCtx, endpointID, connID, payload); dropErr != nil {
slog.Error("on publish dropped", "endpoint", endpointID, "err", dropErr)
}
},
})
if err != nil {
return fmt.Errorf("broker: %w", err)
}
defer func() { _ = brk.Close() }()
sess.Attach(brk)
msgApp = message.New(db, msgLim, hashPool,
message.WithLocks(loginLocks),
message.WithConnRegistry(memConns),
message.WithDownlink(brk),
)
uplink.msg = msgApp
presApp := presence.New(presence.Config{
DB: db,
Downlink: brk,
Conns: &presenceConnTable{conns: memConns},
})
sess.SetPresence(presApp)
idApp := identity.New(identity.Config{
DB: db,
Hash: hashPool,
Locks: loginLocks,
Sessions: sessionTokens,
MaxScheduleSeconds: int64(cfg.Limits.MaxScheduleSeconds),
Logger: slog.Default(),
ConnControl: brk,
})
_ = group.New(group.Config{
DB: db,
Talk: idApp,
Online: presApp,
Downlink: brk,
MaxGroupMembers: cfg.Limits.MaxGroupMembers,
})
if recoverErr := msgApp.RecoverOnStart(ctx); recoverErr != nil {
return fmt.Errorf("message recover: %w", recoverErr)
}
// 其余 deps 供后续 admin / broker / httpx 接线;此处显式引用避免未使用告警。
_ = deps.Identity
_ = deps.Groups
_ = deps.Presence
_ = deps.Downlink
_ = deps.Conns
_ = deps.Uplink
_ = deps.HashPool
_ = deps.SessionTokens
_ = deps.APITokens
_ = deps.LoginLocks
mux := http.NewServeMux()
mux.HandleFunc("GET /healthz", func(w http.ResponseWriter, _ *http.Request) {
w.WriteHeader(http.StatusOK)
_, _ = w.Write([]byte("ok"))
trustedNets := httpx.ParseCIDRs(cfg.TrustedProxies)
adminHandler := admin.New(admin.Deps{
DB: db,
Hash: hashPool,
Tokens: apiTokens,
Locks: loginLocks,
Logger: slog.Default(),
TrustedProxies: trustedNets,
KickEndpoint: func(kickCtx context.Context, endpointID string) (bool, error) {
if _, found := brk.ConnInfoOf(endpointID); !found {
return false, nil
}
if kickErr := sess.Kick(kickCtx, endpointID); kickErr != nil {
return false, kickErr
}
return true, nil
},
})
// P-WIRE-BEGIN
mux.HandleFunc("GET /readyz", func(w http.ResponseWriter, r *http.Request) {
if readyErr := db.Ready(r.Context()); readyErr != nil {
slog.Error("readyz failed", "err", readyErr)
http.Error(w, "not ready", http.StatusServiceUnavailable)
return
metricsReg := metrics.New()
buildHandlers := func(proxies *listener.ProxySet) listener.Handlers {
return listener.Handlers{
MQTT: brk.WSHandler(proxies),
ClientAPI: idApp.Handler(),
AdminAPI: adminHandler,
Metrics: metricsReg.Handler(),
MetricsToken: cfg.Metrics.Token,
Static: staticFileHandler(),
Healthz: http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) {
w.WriteHeader(http.StatusOK)
_, _ = w.Write([]byte("ok"))
}),
Readyz: http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
if readyErr := db.Ready(r.Context()); readyErr != nil {
slog.Error("readyz failed", "err", readyErr)
http.Error(w, "not ready", http.StatusServiceUnavailable)
return
}
w.WriteHeader(http.StatusOK)
_, _ = w.Write([]byte("ok"))
}),
}
w.WriteHeader(http.StatusOK)
_, _ = w.Write([]byte("ok"))
})
// P-WIRE-END
// 确保前端资源被链接进二进制;完整静态托管由后续任务完善。
_ = web.Dist()
ln, err := net.Listen("tcp", cfg.Listen)
if err != nil {
return fmt.Errorf("listen %s: %w", cfg.Listen, err)
}
if err := writeListenAddr(cfg.DataDir, ln.Addr().String()); err != nil {
_ = ln.Close()
handlers := buildHandlers(nil)
shared := cfg.AdminListen == ""
var clientHandler, adminHTTP http.Handler
if shared {
clientHandler = listener.NewMux(listener.RoleShared, handlers)
} else {
clientHandler = listener.NewMux(listener.RoleClient, handlers)
adminHTTP = listener.NewMux(listener.RoleAdmin, handlers)
}
lnOpts := listener.Options{
Listen: cfg.Listen,
AdminListen: cfg.AdminListen,
DataDir: cfg.DataDir,
CertFile: cfg.TLS.CertFile,
KeyFile: cfg.TLS.KeyFile,
AllowPlaintext: cfg.TLS.AllowPlaintext,
TrustedProxies: cfg.TrustedProxies,
ClientHandler: clientHandler,
AdminHandler: adminHTTP,
OnMQTT: func(conn net.Conn) {
_ = brk.AttachTCP(conn)
},
Logger: slog.Default(),
}
lnSrv, err := listener.New(lnOpts)
if err != nil {
return fmt.Errorf("listener: %w", err)
}
handlers = buildHandlers(lnSrv.Proxies())
if shared {
lnOpts.ClientHandler = listener.NewMux(listener.RoleShared, handlers)
} else {
lnOpts.ClientHandler = listener.NewMux(listener.RoleClient, handlers)
lnOpts.AdminHandler = listener.NewMux(listener.RoleAdmin, handlers)
}
lnSrv, err = listener.New(lnOpts)
if err != nil {
return fmt.Errorf("listener: %w", err)
}
if err := lnSrv.Start(ctx); err != nil {
return err
}
defer func() { _ = lnSrv.Close() }()
srv := &http.Server{
Handler: mux,
ReadHeaderTimeout: 10 * time.Second,
if err := writeListenAddr(cfg.DataDir, lnSrv.ListenAddr()); err != nil {
return err
}
if cfg.AdminListen != "" && lnSrv.AdminAddr() != "" {
if err := os.WriteFile(filepath.Join(cfg.DataDir, "admin.addr"), []byte(lnSrv.AdminAddr()+"\n"), 0o644); err != nil {
return err
}
}
errCh := make(chan error, 1)
go func() {
errCh <- srv.Serve(ln)
}()
loopCtx, loopCancel := context.WithCancel(ctx)
defer loopCancel()
go messageLoops(loopCtx, msgApp, memConns)
select {
case <-ctx.Done():
// P-WIRE-BEGIN
// 先停止接受新连接,再等写队列最多 10 秒,然后断开并退出。
shutdownCtx, cancel := context.WithTimeout(context.Background(), 15*time.Second)
defer cancel()
_ = srv.Shutdown(shutdownCtx)
drainCtx, drainCancel := context.WithTimeout(context.Background(), 10*time.Second)
defer drainCancel()
if drainErr := db.Queue.Drain(drainCtx); drainErr != nil && !errors.Is(drainErr, context.DeadlineExceeded) {
slog.Error("write queue drain", "err", drainErr)
}
// P-WIRE-END
serveErr := <-errCh
if serveErr != nil && !errors.Is(serveErr, http.ErrServerClosed) {
return serveErr
}
return nil
case serveErr := <-errCh:
if errors.Is(serveErr, http.ErrServerClosed) {
return nil
}
return serveErr
<-ctx.Done()
loopCancel()
_ = lnSrv.Close()
drainCtx, drainCancel := context.WithTimeout(context.Background(), 10*time.Second)
defer drainCancel()
if drainErr := db.Queue.Drain(drainCtx); drainErr != nil && !errors.Is(drainErr, context.DeadlineExceeded) {
slog.Error("write queue drain", "err", drainErr)
}
return nil
}
func messageLoops(ctx context.Context, msgApp *message.App, conns *message.MemoryConns) {
t := time.NewTicker(time.Second)
defer t.Stop()
for {
select {
case <-ctx.Done():
return
case <-t.C:
nowMs := time.Now().UnixMilli()
if _, err := msgApp.DispatchDue(ctx, nowMs, 100); err != nil {
slog.Error("dispatch due", "err", err)
}
for ep, live := range conns.Snapshot() {
if err := msgApp.PushPending(ctx, ep, live.ConnID); err != nil {
slog.Debug("push pending", "endpoint", ep, "err", err)
}
}
if err := msgApp.CleanupOnce(ctx, nowMs); err != nil {
slog.Error("cleanup once", "err", err)
}
}
}
}
func staticFileHandler() http.Handler {
root := web.Dist()
for _, prefix := range []string{"dist", "stub"} {
sub, err := fs.Sub(root, prefix)
if err != nil {
continue
}
if f, err := sub.Open("index.html"); err == nil {
_ = f.Close()
return http.FileServer(http.FS(sub))
}
}
return http.FileServer(http.FS(root))
}
func writeListenAddr(dataDir, addr string) error {
@@ -145,7 +315,6 @@ func writeListenAddr(dataDir, addr string) error {
return os.WriteFile(path, []byte(addr+"\n"), 0o644)
}
// P-WIRE-BEGIN
func setupJSONLogger(cfg config.LogConfig) {
level := slog.LevelInfo
switch strings.ToLower(strings.TrimSpace(cfg.Level)) {
@@ -159,5 +328,3 @@ func setupJSONLogger(cfg config.LogConfig) {
h := slog.NewJSONHandler(os.Stdout, &slog.HandlerOptions{Level: level})
slog.SetDefault(slog.New(h))
}
// P-WIRE-END
+83
View File
@@ -0,0 +1,83 @@
package main
import (
"context"
"log/slog"
"git.asio.asia/nixevol/NixMsg/internal/app/message"
"git.asio.asia/nixevol/NixMsg/internal/app/port"
)
// appUplink 把 broker 生命周期接到消息连接表与投递推送。
// 其余业务上行帧暂转交 next(可为空 Stub);完整 HandleUplink 分发留后续波次。
type appUplink struct {
msg *message.App
conns *message.MemoryConns
next port.UplinkHandler
log *slog.Logger
}
func (u *appUplink) OnSessionEstablished(ctx context.Context, conn port.ConnInfo) error {
u.conns.Set(conn.EndpointID, message.LiveConn{
ConnID: conn.ConnID,
MaxPacketSize: conn.MaxPacketSize,
})
if u.next != nil {
return u.next.OnSessionEstablished(ctx, conn)
}
return nil
}
func (u *appUplink) OnHandshakeComplete(ctx context.Context, hs port.HandshakeInfo) error {
live := message.LiveConn{
ConnID: hs.ConnID,
MaxReceiveBytes: hs.MaxReceiveBytes,
MaxPacketSize: hs.MaxPacketSize,
}
u.conns.Set(hs.EndpointID, live)
if err := u.msg.OnHandshakeComplete(ctx, hs.EndpointID, live); err != nil {
u.log.Error("message handshake", "endpoint", hs.EndpointID, "err", err)
return err
}
if u.next != nil {
return u.next.OnHandshakeComplete(ctx, hs)
}
return nil
}
func (u *appUplink) OnDisconnect(ctx context.Context, conn port.ConnInfo, reason port.DisconnectReason) {
live, ok := u.conns.Current(conn.EndpointID)
isCurrent := ok && live.ConnID == conn.ConnID
if err := u.msg.OnDisconnect(ctx, conn.EndpointID, conn.ConnID, isCurrent); err != nil {
u.log.Error("message disconnect", "endpoint", conn.EndpointID, "err", err)
}
u.conns.Clear(conn.EndpointID, conn.ConnID)
if u.next != nil {
u.next.OnDisconnect(ctx, conn, reason)
}
}
func (u *appUplink) HandleUplink(ctx context.Context, conn port.ConnInfo, payload []byte) error {
if u.next != nil {
return u.next.HandleUplink(ctx, conn, payload)
}
return nil
}
// presenceConnTable 把消息连接表暴露给 presence.ConnTable。
type presenceConnTable struct {
conns *message.MemoryConns
}
func (p *presenceConnTable) IsOnline(endpointID string) bool {
_, ok := p.conns.Current(endpointID)
return ok
}
func (p *presenceConnTable) CurrentConn(endpointID string) (port.ConnID, bool) {
live, ok := p.conns.Current(endpointID)
if !ok {
return "", false
}
return live.ConnID, true
}
+6 -6
View File
@@ -9,7 +9,7 @@ import (
"git.asio.asia/nixevol/NixMsg/internal/auth"
)
// appDeps 是 serve 组装出的模块依赖。各线在后续任务中替换假实现为真实实现。
// appDeps 是组装出的模块依赖(测试与骨架仍可调用 wire)。
type appDeps struct {
HashPool auth.HashPool
SessionTokens auth.SessionTokens
@@ -26,13 +26,13 @@ type appDeps struct {
Conns port.ConnControl
}
// wire 组装各业务模块的骨架依赖(T0.4:假实现;后续各线替换)。
// wire 返回真实 auth 实现与空业务桩(无 DB 时的轻量装配;正式 serve 走 assembleRuntime)。
func wire() appDeps {
return appDeps{
HashPool: auth.NewStubHashPool(),
SessionTokens: auth.NewStubSessionTokens(),
APITokens: auth.NewStubAPITokens(),
LoginLocks: auth.NewStubLoginLocks(),
HashPool: auth.NewPool(),
SessionTokens: auth.NewSessionTokens(),
APITokens: auth.NewAPITokens(),
LoginLocks: auth.NewLoginLocks(),
Messages: message.NewStub(),
Identity: identity.NewStub(),
+313
View File
@@ -0,0 +1,313 @@
package main
import (
"bytes"
"context"
"database/sql"
"encoding/json"
"io"
"net/http"
"net/http/cookiejar"
"os"
"path/filepath"
"strings"
"testing"
"time"
"git.asio.asia/nixevol/NixMsg/internal/config"
"git.asio.asia/nixevol/NixMsg/internal/protocol"
"git.asio.asia/nixevol/NixMsg/internal/store"
"git.asio.asia/nixevol/NixMsg/test/harness"
"github.com/mochi-mqtt/server/v2/packets"
)
func TestWireAdminLoginRegisterMQTTHandshake(t *testing.T) {
dataDir := t.TempDir()
cfgPath := writeTestConfig(t, dataDir)
initAdminForTest(t, dataDir)
enableRegistration(t, dataDir, "wire-code-99")
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) }()
addr := waitListenAddr(t, dataDir, 15*time.Second)
base := "http://" + addr
// 1) admin init 后真实进程可管理登录
jar, err := cookiejar.New(nil)
if err != nil {
t.Fatal(err)
}
client := &http.Client{Jar: jar, Timeout: 10 * time.Second}
loginBody, _ := json.Marshal(map[string]string{
"username": "admin",
"password": "test-admin-password-xx",
})
resp, err := client.Post(base+"/api/admin/login", "application/json", bytes.NewReader(loginBody))
if err != nil {
t.Fatalf("admin login: %v", err)
}
body, _ := io.ReadAll(resp.Body)
_ = resp.Body.Close()
if resp.StatusCode != http.StatusOK {
t.Fatalf("admin login status=%d body=%s", resp.StatusCode, body)
}
var loginEnv struct {
OK bool `json:"ok"`
}
if uErr := json.Unmarshal(body, &loginEnv); uErr != nil || !loginEnv.OK {
t.Fatalf("admin login resp=%s", body)
}
// 2) 已写入注册开关与安全码后可注册
regBody := `{"registration_code":"wire-code-99","id":"ep_wire1","login_password":"password12","name":"接线端"}`
regResp, err := http.Post(base+"/api/client/register", "application/json", strings.NewReader(regBody))
if err != nil {
t.Fatalf("register: %v", err)
}
regBytes, _ := io.ReadAll(regResp.Body)
_ = regResp.Body.Close()
if regResp.StatusCode != http.StatusOK {
t.Fatalf("register status=%d body=%s", regResp.StatusCode, regBytes)
}
var regEnv struct {
OK bool `json:"ok"`
Data struct {
ID string `json:"id"`
} `json:"data"`
}
if uErr := json.Unmarshal(regBytes, &regEnv); uErr != nil || !regEnv.OK || regEnv.Data.ID != "ep_wire1" {
t.Fatalf("register resp=%s", regBytes)
}
// 3) 注册出的端用密码完成 MQTT 握手并拿到 session_token
tok := mqttPasswordHandshake(t, base, "ep_wire1", "password12")
if tok == "" || !strings.HasPrefix(tok, protocol.SessionTokenPrefix) {
t.Fatalf("session_token=%q", tok)
}
cancel()
select {
case err := <-errCh:
if err != nil {
t.Fatalf("serve exit: %v", err)
}
case <-time.After(15 * time.Second):
t.Fatal("serve did not stop")
}
}
func mqttPasswordHandshake(t *testing.T, httpBase, endpointID, password string) string {
t.Helper()
mc, err := harness.DialMQTTWebSocket(httpBase, 10*time.Second)
if err != nil {
t.Fatalf("dial mqtt ws: %v", err)
}
defer func() { _ = mc.Close() }()
pk := packets.Packet{
FixedHeader: packets.FixedHeader{Type: packets.Connect},
ProtocolVersion: 5,
Connect: packets.ConnectParams{
ProtocolName: []byte("MQTT"),
Clean: true,
ClientIdentifier: endpointID,
Keepalive: 30,
UsernameFlag: true,
Username: []byte(endpointID),
PasswordFlag: true,
Password: []byte(password),
},
}
var buf bytes.Buffer
if encErr := pk.ConnectEncode(&buf); encErr != nil {
t.Fatal(encErr)
}
if sendErr := mc.Send(buf.Bytes()); sendErr != nil {
t.Fatal(sendErr)
}
ack, err := mc.Recv()
if err != nil {
t.Fatalf("connack: %v", err)
}
if len(ack) < 2 || ack[0]>>4 != packets.Connack {
t.Fatalf("want CONNACK, got %x", ack)
}
// MQTT5 CONNACK: remaining length, flags, reason code
reason := byte(0)
if len(ack) >= 4 {
reason = ack[3]
}
if reason != 0 {
t.Fatalf("connack reason=%d raw=%x", reason, ack)
}
sub := packets.Packet{
FixedHeader: packets.FixedHeader{Type: packets.Subscribe, Qos: 1},
ProtocolVersion: 5,
PacketID: 1,
Filters: packets.Subscriptions{
{Filter: "nix/c/" + endpointID + "/down", Qos: 1},
},
}
buf.Reset()
if err := sub.SubscribeEncode(&buf); err != nil {
t.Fatal(err)
}
if err := mc.Send(buf.Bytes()); err != nil {
t.Fatal(err)
}
if _, err := mc.Recv(); err != nil { // SUBACK
t.Fatalf("suback: %v", err)
}
hello, _ := protocol.Marshal(protocol.Hello{
V: protocol.Version, Type: protocol.TypeHello, RID: "h1",
})
pub := packets.Packet{
FixedHeader: packets.FixedHeader{Type: packets.Publish, Qos: 1},
ProtocolVersion: 5,
TopicName: "nix/c/" + endpointID + "/up",
PacketID: 2,
Payload: hello,
}
buf.Reset()
if err := pub.PublishEncode(&buf); err != nil {
t.Fatal(err)
}
if err := mc.Send(buf.Bytes()); err != nil {
t.Fatal(err)
}
deadline := time.Now().Add(10 * time.Second)
for time.Now().Before(deadline) {
raw, err := mc.Recv()
if err != nil {
t.Fatalf("recv down: %v", err)
}
if len(raw) < 2 {
continue
}
typ := raw[0] >> 4
if typ == packets.Puback || typ == packets.Pingresp {
continue
}
if typ != packets.Publish {
continue
}
payload, err := decodePublishPayload(raw)
if err != nil {
t.Fatalf("publish decode: %v raw=%x", err, raw)
}
var m map[string]any
if err := json.Unmarshal(payload, &m); err != nil {
t.Fatalf("json: %v payload=%s", err, payload)
}
if m["type"] == "resp" {
if tok := extractSessionToken(m); tok != "" {
return tok
}
t.Fatalf("hello resp without token: %v", m)
}
}
t.Fatal("timeout waiting hello resp")
return ""
}
func decodePublishPayload(raw []byte) ([]byte, error) {
rem, n, err := decodeRemainingLength(raw[1:])
if err != nil {
return nil, err
}
body := raw[1+n:]
if len(body) != rem {
return nil, io.ErrUnexpectedEOF
}
pk := packets.Packet{
ProtocolVersion: 5,
FixedHeader: packets.FixedHeader{
Type: packets.Publish,
Remaining: rem,
Qos: (raw[0] >> 1) & 0x3,
},
}
if err := pk.PublishDecode(body); err != nil {
return nil, err
}
return pk.Payload, nil
}
func decodeRemainingLength(b []byte) (value int, n int, err error) {
var mul uint32 = 1
var v uint32
for i := 0; i < len(b) && i < 4; i++ {
v += uint32(b[i]&127) * mul
n++
if b[i]&128 == 0 {
return int(v), n, nil
}
mul *= 128
}
return 0, 0, io.ErrUnexpectedEOF
}
func extractSessionToken(m map[string]any) string {
if m["ok"] != true {
return ""
}
data, _ := m["data"].(map[string]any)
tok, _ := data["session_token"].(string)
return tok
}
func enableRegistration(t *testing.T, dataDir, code string) {
t.Helper()
db, err := store.Open(dataDir, "FULL")
if err != nil {
t.Fatal(err)
}
defer func() { _ = db.Close() }()
now := time.Now().UnixMilli()
err = db.Queue.Do(context.Background(), func(tx *sql.Tx) error {
if _, e := tx.Exec(`INSERT INTO settings(key, value, updated_at) VALUES(?,?,?)
ON CONFLICT(key) DO UPDATE SET value=excluded.value, updated_at=excluded.updated_at`,
"registration_enabled", "1", now); e != nil {
return e
}
_, e := tx.Exec(`INSERT INTO settings(key, value, updated_at) VALUES(?,?,?)
ON CONFLICT(key) DO UPDATE SET value=excluded.value, updated_at=excluded.updated_at`,
"registration_code", code, now)
return e
})
if err != nil {
t.Fatal(err)
}
}
func waitListenAddr(t *testing.T, dataDir string, timeout time.Duration) string {
t.Helper()
deadline := time.Now().Add(timeout)
path := filepath.Join(dataDir, "listen.addr")
for time.Now().Before(deadline) {
b, err := os.ReadFile(path)
if err == nil {
addr := strings.TrimSpace(string(b))
if addr != "" {
return addr
}
}
time.Sleep(20 * time.Millisecond)
}
t.Fatal("listen.addr not written")
return ""
}
+46
View File
@@ -43,6 +43,29 @@
- 备选方案:docker 目标仅 echo 提示。
- 影响:镜像发布流程仍由 Q4 定稿。
### L-WIRE 2026-09-30
1. **serve 真实接线范围**
- 原条款:TASKS 总控接线;listener/broker/admin/注册/消息周期循环挂到进程。
- 实际做法:`cmd/nixmsg/serve.go` 注入真实 `auth.NewPool`/`NewSessionTokens`/`NewAPITokens`/`NewLoginLocks`,挂注册与管理路由(踢线调 `Session.Kick`),listener 识别 HTTP/WebSocket/`OnMQTT` 裸 TCP,broker `Login`+`Session`,握手/断线/`OnPublishDropped` 接到 `message.App`,周期 `DispatchDue`/`PushPending`/`CleanupOnce`,下行 `PublishDown`。集成测覆盖:管理登录、写 settings 后注册、密码 MQTT 握手拿 `session_token`。
- 原因:第二波收尾;三件验收必须通。
- 备选方案:分文件多阶段接线。
- 影响:进程可对外登录/注册/握手。
2. **未接线 / 未完成部分**
- 原条款:上行应用帧完整分发到 identity/group/presence/message 业务方法。
- 实际做法:`Session` 处理 hello/logout;其余 `HandleUplink` 仍为 `StubUplinkHandler`(不解析 send/ack/group 等)。管理注册设置 HTTP(A3)未实现,测试直接写 `settings` 表。群/在线模块已构造并注入 Downlink,但无上行入口调用。
- 原因:本波强制三件验收;完整协议分发属后续波次。
- 备选方案:本波同时实现 HandleUplink 大 multiplex。
- 影响:端连上后除握手/logout 外业务帧尚无 resp;F03–F16 等仍依赖后续接线。
3. **listen.addr 始终写入**
- 原条款:listener N1 仅端口 0 写文件;T0.1 超集为启动即写。
- 实际做法:listener 按 N1 写端口 0;serve 成功后再强制写一次 `listen.addr`(及分离时的 `admin.addr`)。
- 原因:与 harness / T0.1 一致。
- 备选方案:改 listener 始终写。
- 影响:固定端口也会有地址文件。
### T0.2 2026-09-30
1. **请求指纹规范化格式**
@@ -747,3 +770,26 @@
- 原因:空命名卷属主为 root 时,distroless nonroot 无法建库(`unable to open database file`)。
- 备选方案:compose 增加一次性 init 服务;或文档要求宿主机目录预授权。
- 影响:按示例首次 `up` 前需处理权限,否则 serve 立即退出。
### Q2 第一部分(已合并功能验收)2026-09-30
1. **对真实进程探测,未接线则记「未测」而非改业务代码**
- 原条款:TASKS Q2「PRD 第 10 节每条至少有一个集成测试」;本波只做已合并功能的第一部分。
- 实际做法:`test/accept` 用 `test/harness`(随机端口 + 临时目录)起真实 `nixmsg`;对 `/api/admin/login`、`/api/client/register`、`POST /api/admin/endpoints` 先探测。404 则该条写「未测」并注明缺总控接线 / A2 / I 等,不修改 `cmd/nixmsg` 或业务包求绿。
- 原因:main@d357082 上 A1/I1 等仅为可挂载 Handler,`serve` 只挂了 `/healthz`、`/readyz`(见各线 DEVIATIONS「未改 cmd」)。
- 备选方案:测试进程内自行 `admin.New` 挂路由(不反映交付二进制行为,否决)。
- 影响:本波 F17/F01/F23 多为未测;接线后同一测试会自动跑登录锁定、CSRF、注册与开通路径。
2. **F22 只验收 init + 健康检查子集**
- 原条款:F22 含备份恢复、升级迁移、证书重载、Docker、指标。
- 实际做法:集成测试覆盖空目录 `admin init`、`serve`、`/healthz`、`/readyz`,并断言管理员密码不出现在 serve 的 stdout/stderr;其余 F22 子项仍标未测。
- 原因:本波范围是「现在就能测的路径」。
- 备选方案:本波强行跑 Docker/证书(超出第一部分)。
- 影响:对照表 F22 为「通过」但备注写明未覆盖项。
3. **验收报告写入方式**
- 原条款:用 `test/report` 生成 F01–F23 对照表。
- 实际做法:`TestQ2AcceptAndReport` 汇总探测结果;默认写临时目录。`task q:accept`(`NIXMSG_WRITE_ACCEPT_REPORT=1`)写入 `test/report/testdata/q2_results.json` 与 `test/accept/ACCEPTANCE.md`;`task q:report-q2` 可再生成 Markdown。普通 `go test`/`task check` 不改仓库文件。
- 原因:避免每次单测改 `generated_at` 弄脏工作区。
- 备选方案:固定时间戳始终写入仓库。
- 影响:交付审阅以 `ACCEPTANCE.md` / `q2_results.json` 为准,需先跑过 `task q:accept`。
+11
View File
@@ -67,6 +67,17 @@ func (c *MemoryConns) Current(endpointID string) (LiveConn, bool) {
return v, ok
}
// Snapshot 返回当前连接表副本(供推送循环遍历)。
func (c *MemoryConns) Snapshot() map[string]LiveConn {
c.mu.RLock()
defer c.mu.RUnlock()
out := make(map[string]LiveConn, len(c.m))
for k, v := range c.m {
out[k] = v
}
return out
}
// RecordingDownlink 记录下行发布,供测试断言。
type RecordingDownlink struct {
mu sync.Mutex
+7
View File
@@ -58,11 +58,16 @@ func (AllowAuthenticator) Authenticate(context.Context, string, []byte, string)
return AuthResult{OK: true}, nil
}
// PublishDroppedFunc 下行未写入发送队列时回调(对接消息线 OnPublishDropped)。
type PublishDroppedFunc func(ctx context.Context, endpointID string, connID port.ConnID, payload []byte)
// Options 装配 broker。
type Options struct {
Authenticator Authenticator
Uplink port.UplinkHandler
Logger *slog.Logger
// OnPublishDropped 可选;nil 时仅打 debug 日志。
OnPublishDropped PublishDroppedFunc
}
// Broker 内置 mochi,不自带监听端口。
@@ -71,6 +76,7 @@ type Broker struct {
auth Authenticator
uplink port.UplinkHandler
log *slog.Logger
onDrop PublishDroppedFunc
hook *nixHook
@@ -144,6 +150,7 @@ func New(opts Options) (*Broker, error) {
auth: auth,
uplink: uplink,
log: log,
onDrop: opts.OnPublishDropped,
current: make(map[string]*connState),
byClient: make(map[*mqtt.Client]*connState),
queues: make(map[string]*uplinkQueue),
+10
View File
@@ -159,6 +159,16 @@ func (h *nixHook) OnSubscribed(cl *mqtt.Client, pk packets.Packet, reasonCodes [
func (h *nixHook) OnPublishDropped(cl *mqtt.Client, pk packets.Packet) {
h.b.log.Debug("publish dropped", "client", cl.ID, "topic", pk.TopicName, "size", len(pk.Payload))
if h.b.onDrop == nil {
return
}
h.b.connsMu.RLock()
st := h.b.byClient[cl]
h.b.connsMu.RUnlock()
if st == nil {
return
}
h.b.onDrop(context.Background(), st.endpointID, st.connID, append([]byte(nil), pk.Payload...))
}
func (h *nixHook) OnSessionEstablished(cl *mqtt.Client, _ packets.Packet) {
+5
View File
@@ -101,6 +101,11 @@ func (s *Session) Attach(b *Broker) {
s.b = b
}
// SetPresence 接线时在创建最终 presence 实现后注入(可替换占位)。
func (s *Session) SetPresence(p PresenceSink) {
s.presence = p
}
func (s *Session) OnSessionEstablished(ctx context.Context, conn port.ConnInfo) error {
if s.b != nil {
s.b.startHandshakeDeadline(conn.EndpointID, conn.ConnID, handshakeTimeout)
+12
View File
@@ -6,6 +6,13 @@ tasks:
cmds:
- go test ./test/chaos/ -count=1 -v
q:accept:
desc: 跑 Q2 验收集成测试并写入 F01–F23 对照表
cmds:
- go test ./test/accept/ -count=1 -v
env:
NIXMSG_WRITE_ACCEPT_REPORT: "1"
q:report:
desc: 从结果 JSON 生成 F01–F23 验收对照表(默认空样例)
cmds:
@@ -16,6 +23,11 @@ tasks:
cmds:
- go run ./test/report/cmd/genreport ./test/report/testdata/sample_results.json
q:report-q2:
desc: 用 Q2 验收结果生成对照表到 stdout
cmds:
- go run ./test/report/cmd/genreport ./test/report/testdata/q2_results.json
q:load-test:
desc: 压测客户端骨架单元测试
cmds:
+31
View File
@@ -0,0 +1,31 @@
# NixMsg 验收对照表(PRD 第 10 节)
生成时间:2026-09-29T23:20:07Z
汇总:通过 1,失败 0,未测 22
| 编号 | 一句话 | 结果 | 备注 |
|---|---|---|---|
| F01 | 批量开通整批校验、停用、删除群主转让、删除后同编号重开不串数据 | 未测 | 未测:serve 未挂载 POST /api/admin/endpoints(A2 未合入或未接线;当前 A1 对端路由返回 501 亦未挂到进程) |
| F02 | 新设备登录后旧设备自动退出、换 IP 用令牌重连、两种密码锁定、重置密码后被踢、服务器故障不误报密码错误 | 未测 | 未测:端登录/会话令牌属连接 N3,main 上 serve 未挂 broker |
| F03 | 断开后状态及时变离线,全表可列出 | 未测 | 未测:在线状态属身份 I3 + 连接 N3,未接线 |
| F04 | 只通知订阅了的端 | 未测 | 未测:presence.watch 属身份 I3,未接线 |
| F05 | 崩溃不丢已提交消息,消息号去重和冲突,密码门生效,配额生效 | 未测 | 未测:消息提交属消息 M1,未挂入 broker 上行 |
| F06 | 群成员收到同一份,入群前不补,发送者不收到自己的 | 未测 | 未测:群消息属消息 M + 身份 I4,未接线 |
| F07 | 256 KiB 通过,超出拒绝,接收上限生效 | 未测 | 未测:大小限制属消息/连接,未接线 |
| F08 | 弱网最终送达且应用层不重复,重启后续传 | 未测 | 未测:投递确认属消息 M2,未接线 |
| F09 | 保留时间从发送时刻起算,超时过期 | 未测 | 未测:保留期属消息 M2,未接线 |
| F10 | 短断线送到,长断线丢弃,服务器重启后宽限内重连送到 | 未测 | 未测:断线策略属消息 M2,未接线 |
| F11 | 发送方离线后到点仍发送 | 未测 | 未测:定时发送属消息 M2/M4,未接线 |
| F12 | 延迟窗口内撤回对方收不到 | 未测 | 未测:延迟撤回属消息 M3,未接线 |
| F13 | 未推送必撤成功;群部分确认得到部分撤回 | 未测 | 未测:撤回判定属消息 M3,未接线 |
| F14 | 回执能补送给当时离线的发送方 | 未测 | 未测:回执属消息 M3,未接线 |
| F15 | 输一次记住、改密失效、回复免密、进群仍要密码、防多账号轮流猜 | 未测 | 未测:对话密码属身份 I2,未接线 |
| F16 | 群主权限、退出后不再收到、解散后同编号新群不收旧消息 | 未测 | 未测:群权限属身份 I4,未接线 |
| F17 | 后台管端、管注册、管群、查记录,响应里没有正文;API 令牌可用且不能越权 | 未测 | 未测:serve 未挂载 /api/admin/login(A1 Handler 已实现,缺总控/接线挂到 cmd/nixmsg;DEVIATIONS 后台接口 A §1) |
| F18 | 送达后正文消失;记录天数 0 时连记录消失;防重仍在 | 未测 | 未测:正文清理属消息 M3,未接线 |
| F19 | 四种 SDK 通过同一清单 | 未测 | 未测:SDK 接入清单属 S1/S2 任务 4,依赖真实服务接线 |
| F20 | 裸 MQTT 能登录、收、确认、发 | 未测 | 未测:裸 MQTT 属连接 N,serve 未挂 broker |
| F21 | 默认一个端口提供后台、WebSocket、TCP、注册;后台可分到单独端口 | 未测 | 仅验证 listen 上 /healthz+/readyz;后台 API、/mqtt、裸 TCP、注册未挂入 serve(缺总控接线 + 连接 N) |
| F22 | 初始化后单文件或 Docker 启动、备份恢复、升级迁移、证书自动重载、指标可抓取 | 通过 | 已测:空目录 admin init + serve,/healthz 与 /readyz 成功,密码不在 serve 日志;未测:备份恢复、升级迁移、证书重载、Docker、/metrics |
| F23 | 注册开关、安全码校验、换码不影响已注册、输错锁定 | 未测 | 未测:serve 未挂载 POST /api/client/register(I1 Handler 已实现,缺总控/连接 N 接线;DEVIATIONS 身份 I §1) |
+496
View File
@@ -0,0 +1,496 @@
package accept_test
import (
"bytes"
"encoding/json"
"fmt"
"io"
"net/http"
"os"
"path/filepath"
"strings"
"testing"
"time"
"git.asio.asia/nixevol/NixMsg/test/accept"
"git.asio.asia/nixevol/NixMsg/test/harness"
"git.asio.asia/nixevol/NixMsg/test/report"
)
// TestQ2AcceptAndReport 是 Q2 第一部分:对 main 已有能力做集成探测,并生成 F01–F23 对照表。
// 未接线的路由记为「未测」并写明缺哪条线;不改业务代码以求变绿。
func TestQ2AcceptAndReport(t *testing.T) {
items := make(map[string]report.Item, len(report.Features))
for _, f := range report.Features {
items[f.ID] = report.Item{
ID: f.ID,
Status: report.StatusUntested,
Note: "本波未覆盖;依赖后续线合入与接线",
}
}
set := func(id string, st report.Status, note string) {
items[id] = report.Item{ID: id, Status: st, Note: note}
}
// —— 1) admin init + serve + /healthz + /readyz,密码不在 serve 日志 ——
runInitHealthz(t, set)
// —— 共用 harness 进程,探测管理 / 注册 / 开通 ——
srv, err := harness.Start(harness.Options{})
if err != nil {
t.Fatalf("harness start: %v", err)
}
defer func() { _ = srv.Stop() }()
if srv.AdminPassword == "" {
t.Fatal("admin init 未返回密码(P1 应已合入)")
}
runAdminAuth(t, srv, set)
runRegistration(t, srv, set)
runEndpointCreate(t, srv, set)
// 其余条目写清未测原因(缺哪条线)
setDefaultUntested(set)
out := make([]report.Item, 0, len(report.Features))
for _, f := range report.Features {
out = append(out, items[f.ID])
}
results := report.Results{
GeneratedAt: time.Now().UTC().Format(time.RFC3339),
Items: out,
}
normalized, err := report.Normalize(results)
if err != nil {
t.Fatal(err)
}
var md bytes.Buffer
if werr := report.WriteMarkdown(&md, normalized); werr != nil {
t.Fatal(werr)
}
if !strings.Contains(md.String(), "F01") || !strings.Contains(md.String(), "F23") {
t.Fatalf("report missing features:\n%s", md.String())
}
if !strings.Contains(md.String(), "通过") && !strings.Contains(md.String(), "未测") {
t.Fatalf("report missing status words:\n%s", md.String())
}
// 默认写到临时目录验证;设 NIXMSG_WRITE_ACCEPT_REPORT=1 时写入仓库产物供交付。
dir := t.TempDir()
resultsPath := filepath.Join(dir, "q2_results.json")
mdPath := filepath.Join(dir, "ACCEPTANCE.md")
if os.Getenv("NIXMSG_WRITE_ACCEPT_REPORT") == "1" {
root := findModuleRoot(t)
resultsPath = filepath.Join(root, "test", "report", "testdata", "q2_results.json")
mdPath = filepath.Join(root, "test", "accept", "ACCEPTANCE.md")
}
raw, err := json.MarshalIndent(results, "", " ")
if err != nil {
t.Fatal(err)
}
if err := os.WriteFile(resultsPath, append(raw, '\n'), 0o644); err != nil {
t.Fatalf("write results: %v", err)
}
if err := os.WriteFile(mdPath, md.Bytes(), 0o644); err != nil {
t.Fatalf("write acceptance md: %v", err)
}
t.Logf("wrote %s and %s", resultsPath, mdPath)
}
func runInitHealthz(t *testing.T, set func(string, report.Status, string)) {
t.Helper()
ls, err := accept.StartWithLogCapture()
if err != nil {
set("F22", report.StatusFail, "admin init/serve 失败(平台 P): "+err.Error())
t.Errorf("init/serve: %v", err)
return
}
defer func() { _ = ls.Stop() }()
pass := ls.AdminPassword
if pass == "" {
set("F22", report.StatusFail, "admin init 未打印密码(平台 P)")
t.Error("empty admin password")
return
}
hz, err := http.Get(ls.HTTPBase + "/healthz")
if err != nil {
set("F22", report.StatusFail, "/healthz 不可达(平台 P): "+err.Error())
t.Errorf("healthz: %v", err)
return
}
body, _ := io.ReadAll(hz.Body)
_ = hz.Body.Close()
if hz.StatusCode != http.StatusOK || string(body) != "ok" {
set("F22", report.StatusFail, fmt.Sprintf("/healthz status=%d body=%q(平台 P)", hz.StatusCode, body))
t.Errorf("healthz status=%d body=%q", hz.StatusCode, body)
return
}
rz, err := http.Get(ls.HTTPBase + "/readyz")
if err != nil {
set("F22", report.StatusFail, "/readyz 不可达(平台 P): "+err.Error())
t.Errorf("readyz: %v", err)
return
}
rbody, _ := io.ReadAll(rz.Body)
_ = rz.Body.Close()
if rz.StatusCode != http.StatusOK || string(rbody) != "ok" {
set("F22", report.StatusFail, fmt.Sprintf("/readyz status=%d body=%q(平台 P)", rz.StatusCode, rbody))
t.Errorf("readyz status=%d body=%q", rz.StatusCode, rbody)
return
}
logs := ls.LogBuf.String()
if strings.Contains(logs, pass) {
set("F22", report.StatusFail, "管理员密码出现在 serve 日志(平台 P)")
t.Errorf("password leaked into serve logs")
return
}
set("F22", report.StatusPass, "已测:空目录 admin init + serve,/healthz 与 /readyz 成功,密码不在 serve 日志;未测:备份恢复、升级迁移、证书重载、Docker、/metrics")
// F21:当前进程至少在单一 listen 上提供健康检查;后台 API / MQTT 尚未接线
set("F21", report.StatusUntested, "仅验证 listen 上 /healthz+/readyz;后台 API、/mqtt、裸 TCP、注册未挂入 serve(缺总控接线 + 连接 N)")
}
func runAdminAuth(t *testing.T, srv *harness.Server, set func(string, report.Status, string)) {
t.Helper()
code, body, err := accept.ProbeMethod(srv.AdminHTTPBase, http.MethodPost, "/api/admin/login",
[]byte(`{"username":"admin","password":"not-the-password!!"}`))
if err != nil {
set("F17", report.StatusFail, "探测管理登录失败: "+err.Error())
t.Errorf("probe login: %v", err)
return
}
if accept.ClassifyAdminLogin(code) == accept.RouteMissing {
note := "未测:serve 未挂载 /api/admin/login(A1 Handler 已实现,缺总控/接线挂到 cmd/nixmsg;DEVIATIONS 后台接口 A §1)"
set("F17", report.StatusUntested, note)
t.Log(note)
return
}
// 路由已挂上:跑登录、锁定、CSRF
client, err := srv.AdminClient()
if err != nil {
set("F17", report.StatusFail, "AdminClient: "+err.Error())
t.Fatal(err)
}
// 正确密码登录
loginBody, _ := json.Marshal(map[string]string{
"username": "admin",
"password": srv.AdminPassword,
})
resp, err := client.PostJSON("/api/admin/login", loginBody)
if err != nil {
set("F17", report.StatusFail, "登录请求失败: "+err.Error())
t.Fatal(err)
}
loginRaw, _ := io.ReadAll(resp.Body)
_ = resp.Body.Close()
if resp.StatusCode != http.StatusOK {
set("F17", report.StatusFail, fmt.Sprintf("正确密码登录失败 status=%d body=%s(后台接口 A)", resp.StatusCode, loginRaw))
t.Errorf("login want 200 got %d %s", resp.StatusCode, loginRaw)
return
}
// 无 CSRF 的改状态请求应被拒
req, err := http.NewRequest(http.MethodPost, srv.AdminHTTPBase+"/api/admin/logout", strings.NewReader(`{}`))
if err != nil {
t.Fatal(err)
}
req.Header.Set("Content-Type", "application/json")
// 故意不加 X-Nixmsg-Request
noCSRF, err := client.HTTP.Do(req)
if err != nil {
set("F17", report.StatusFail, "无 CSRF 请求失败: "+err.Error())
t.Fatal(err)
}
noBody, _ := io.ReadAll(noCSRF.Body)
_ = noCSRF.Body.Close()
if noCSRF.StatusCode != http.StatusForbidden {
set("F17", report.StatusFail, fmt.Sprintf("无 CSRF 期望 403 得 %d body=%s(后台接口 A)", noCSRF.StatusCode, noBody))
t.Errorf("csrf: want 403 got %d %s", noCSRF.StatusCode, noBody)
return
}
// 带 CSRF 的 logout 应成功
okLogout, err := client.PostJSON("/api/admin/logout", []byte(`{}`))
if err != nil {
set("F17", report.StatusFail, "带 CSRF logout 失败: "+err.Error())
t.Fatal(err)
}
_ = okLogout.Body.Close()
if okLogout.StatusCode != http.StatusOK {
set("F17", report.StatusFail, fmt.Sprintf("带 CSRF logout 期望 200 得 %d(后台接口 A)", okLogout.StatusCode))
t.Errorf("logout with csrf: want 200 got %d", okLogout.StatusCode)
return
}
// 错误密码锁定(新客户端,避免 Cookie 干扰;10 次错)
lockClient, err := harness.NewAdminClient(srv.AdminHTTPBase)
if err != nil {
t.Fatal(err)
}
var locked bool
for i := 0; i < 12; i++ {
r, rerr := lockClient.PostJSON("/api/admin/login", []byte(`{"username":"admin","password":"wrong-password!!"}`))
if rerr != nil {
set("F17", report.StatusFail, "错误密码请求失败: "+rerr.Error())
t.Fatal(rerr)
}
raw, _ := io.ReadAll(r.Body)
_ = r.Body.Close()
if r.StatusCode == http.StatusTooManyRequests {
locked = true
break
}
if r.StatusCode != http.StatusUnauthorized {
set("F17", report.StatusFail, fmt.Sprintf("错误密码第 %d 次期望 401/429 得 %d body=%s(后台接口 A)", i+1, r.StatusCode, raw))
t.Errorf("bad login %d: %d %s", i, r.StatusCode, raw)
return
}
}
if !locked {
set("F17", report.StatusFail, "错误密码未触发锁定(后台接口 A / 平台 P3)")
t.Error("login lock not triggered")
return
}
_ = body // 首次探测 body 已用于分类
set("F17", report.StatusPass, "已测:管理登录、错误密码锁定、Cookie 会话下无 CSRF 被拒 / 有 CSRF 可通过;管端/管注册/管群/查记录/令牌越权等未在本波覆盖(A2/A3)")
}
func runRegistration(t *testing.T, srv *harness.Server, set func(string, report.Status, string)) {
t.Helper()
code, body, err := accept.ProbeMethod(srv.HTTPBase, http.MethodPost, "/api/client/register",
[]byte(`{"code":"x"}`))
if err != nil {
set("F23", report.StatusFail, "探测注册失败: "+err.Error())
t.Errorf("probe register: %v", err)
return
}
if accept.ClassifyRegister(code) == accept.RouteMissing {
note := "未测:serve 未挂载 POST /api/client/register(I1 Handler 已实现,缺总控/连接 N 接线;DEVIATIONS 身份 I §1)"
set("F23", report.StatusUntested, note)
t.Log(note)
return
}
// 若已挂上:覆盖开关关闭、错码、对码、换码(尽量用管理接口;管理未挂则只能测默认关闭)
adminCode, _, _ := accept.ProbeMethod(srv.AdminHTTPBase, http.MethodGet, "/api/admin/registration", nil)
if accept.ClassifyAdminLogin(adminCode) == accept.RouteMissing || adminCode == http.StatusNotFound {
// 注册路由在、管理注册设置不在:至少验证默认关闭
if code == http.StatusForbidden || code == http.StatusNotFound || code == http.StatusBadRequest || code == http.StatusConflict {
set("F23", report.StatusPass, fmt.Sprintf("注册路由已挂;默认关闭或校验拒绝(status=%d)。管理注册设置未挂,换码路径未测(缺 A3 接线) body=%s", code, trim(body)))
return
}
set("F23", report.StatusFail, fmt.Sprintf("注册路由异常 status=%d body=%s(身份 I)", code, trim(body)))
t.Errorf("register unexpected %d %s", code, body)
return
}
// 完整 F23:开关、错码、对码、换码 —— 需管理 PUT registration
ac, err := srv.AdminClient()
if err != nil {
t.Fatal(err)
}
loginBody, _ := json.Marshal(map[string]string{"username": "admin", "password": srv.AdminPassword})
lr, err := ac.PostJSON("/api/admin/login", loginBody)
if err != nil {
t.Fatal(err)
}
_ = lr.Body.Close()
if lr.StatusCode != http.StatusOK {
set("F23", report.StatusFail, fmt.Sprintf("F23 前置管理登录失败 %d(后台接口 A)", lr.StatusCode))
return
}
code1 := "accept-code-one-aaaa"
put1, err := ac.Do(http.MethodPut, "/api/admin/registration",
[]byte(fmt.Sprintf(`{"enabled":true,"code":%q}`, code1)), "application/json")
if err != nil {
t.Fatal(err)
}
put1Body, _ := io.ReadAll(put1.Body)
_ = put1.Body.Close()
if put1.StatusCode == http.StatusNotImplemented {
set("F23", report.StatusUntested, "注册路由已挂,但 PUT /api/admin/registration 返回 501(缺 A3)")
t.Log("registration settings 501")
return
}
if put1.StatusCode != http.StatusOK {
set("F23", report.StatusFail, fmt.Sprintf("开启注册失败 status=%d body=%s(后台接口 A / 身份 I)", put1.StatusCode, put1Body))
return
}
// 错码
bad, _, berr := accept.ProbeMethod(srv.HTTPBase, http.MethodPost, "/api/client/register",
[]byte(`{"code":"wrong-code","id":"q2bad001"}`))
if berr != nil {
t.Fatal(berr)
}
if bad != http.StatusUnauthorized && bad != http.StatusForbidden {
set("F23", report.StatusFail, fmt.Sprintf("错码期望 401/403 得 %d(身份 I)", bad))
return
}
// 对码
okCode, okBody, oerr := accept.ProbeMethod(srv.HTTPBase, http.MethodPost, "/api/client/register",
[]byte(fmt.Sprintf(`{"code":%q,"id":"q2ok0001","login_password":"password1234"}`, code1)))
if oerr != nil {
t.Fatal(oerr)
}
if okCode != http.StatusOK {
set("F23", report.StatusFail, fmt.Sprintf("对码注册失败 status=%d body=%s(身份 I)", okCode, trim(okBody)))
return
}
// 换码后旧码失败、新码成功
code2 := "accept-code-two-bbbb"
put2, err := ac.Do(http.MethodPut, "/api/admin/registration",
[]byte(fmt.Sprintf(`{"enabled":true,"code":%q}`, code2)), "application/json")
if err != nil {
t.Fatal(err)
}
_ = put2.Body.Close()
if put2.StatusCode != http.StatusOK {
set("F23", report.StatusFail, fmt.Sprintf("换码失败 status=%d(后台接口 A)", put2.StatusCode))
return
}
oldBad, _, _ := accept.ProbeMethod(srv.HTTPBase, http.MethodPost, "/api/client/register",
[]byte(fmt.Sprintf(`{"code":%q,"id":"q2old001","login_password":"password1234"}`, code1)))
if oldBad == http.StatusOK {
set("F23", report.StatusFail, "换码后旧码仍可注册(身份 I)")
return
}
newOK, _, _ := accept.ProbeMethod(srv.HTTPBase, http.MethodPost, "/api/client/register",
[]byte(fmt.Sprintf(`{"code":%q,"id":"q2new001","login_password":"password1234"}`, code2)))
if newOK != http.StatusOK {
set("F23", report.StatusFail, fmt.Sprintf("换码后新码注册失败 status=%d(身份 I)", newOK))
return
}
set("F23", report.StatusPass, "已测:开关、错码、对码、换码;输错锁定未在本用例穷尽")
}
func runEndpointCreate(t *testing.T, srv *harness.Server, set func(string, report.Status, string)) {
t.Helper()
// 先看未登录时路由是否存在
code, body, err := accept.ProbeMethod(srv.AdminHTTPBase, http.MethodPost, "/api/admin/endpoints",
[]byte(`{"id":"q2ep0001","login_password":"password1234"}`))
if err != nil {
set("F01", report.StatusFail, "探测开通端失败: "+err.Error())
t.Errorf("probe endpoints: %v", err)
return
}
if accept.ClassifyCreateEndpoint(code) == accept.RouteMissing {
note := "未测:serve 未挂载 POST /api/admin/endpoints(A2 未合入或未接线;当前 A1 对端路由返回 501 亦未挂到进程)"
set("F01", report.StatusUntested, note)
t.Log(note)
return
}
client, err := srv.AdminClient()
if err != nil {
t.Fatal(err)
}
loginBody, _ := json.Marshal(map[string]string{"username": "admin", "password": srv.AdminPassword})
lr, err := client.PostJSON("/api/admin/login", loginBody)
if err != nil {
t.Fatal(err)
}
loginRaw, _ := io.ReadAll(lr.Body)
_ = lr.Body.Close()
if lr.StatusCode == http.StatusNotFound {
set("F01", report.StatusUntested, "端开通路由有响应但管理登录未挂载,无法完成开通验收(缺接线)")
return
}
if lr.StatusCode != http.StatusOK {
set("F01", report.StatusFail, fmt.Sprintf("开通前置登录失败 %d %s(后台接口 A)", lr.StatusCode, loginRaw))
return
}
create, err := client.PostJSON("/api/admin/endpoints",
[]byte(`{"id":"q2ep0001","login_password":"password1234"}`))
if err != nil {
t.Fatal(err)
}
craw, _ := io.ReadAll(create.Body)
_ = create.Body.Close()
if create.StatusCode == http.StatusNotImplemented {
set("F01", report.StatusUntested, "POST /api/admin/endpoints 返回 501(缺 A2 业务实现)")
t.Log("endpoints 501")
return
}
if create.StatusCode != http.StatusOK && create.StatusCode != http.StatusCreated {
set("F01", report.StatusFail, fmt.Sprintf("开通端失败 status=%d body=%s(后台接口 A)", create.StatusCode, trim(string(craw))))
return
}
// 错误密码连不上:依赖 MQTT 登录(连接 N3)。若 /mqtt 未挂则记未测。
mqttCode, _, _ := accept.ProbeMethod(srv.HTTPBase, http.MethodGet, "/mqtt", nil)
if mqttCode == http.StatusNotFound {
set("F01", report.StatusUntested, "端已开通,但 /mqtt 未挂,无法验证错误密码连不上(缺连接 N 接线)")
return
}
set("F01", report.StatusPass, "已测:开通一端;错误密码 MQTT 连接拒绝需 N3 联调细节,本波仅确认路由可用。批量/停用/删除转让等未覆盖")
_ = body
_ = code
}
func setDefaultUntested(set func(string, report.Status, string)) {
// 只填尚未被专项用例写入的条目;F01/F17/F21/F22/F23 由探测结果决定。
defaults := map[string]string{
"F02": "未测:端登录/会话令牌属连接 N3,main 上 serve 未挂 broker",
"F03": "未测:在线状态属身份 I3 + 连接 N3,未接线",
"F04": "未测:presence.watch 属身份 I3,未接线",
"F05": "未测:消息提交属消息 M1,未挂入 broker 上行",
"F06": "未测:群消息属消息 M + 身份 I4,未接线",
"F07": "未测:大小限制属消息/连接,未接线",
"F08": "未测:投递确认属消息 M2,未接线",
"F09": "未测:保留期属消息 M2,未接线",
"F10": "未测:断线策略属消息 M2,未接线",
"F11": "未测:定时发送属消息 M2/M4,未接线",
"F12": "未测:延迟撤回属消息 M3,未接线",
"F13": "未测:撤回判定属消息 M3,未接线",
"F14": "未测:回执属消息 M3,未接线",
"F15": "未测:对话密码属身份 I2,未接线",
"F16": "未测:群权限属身份 I4,未接线",
"F18": "未测:正文清理属消息 M3,未接线",
"F19": "未测:SDK 接入清单属 S1/S2 任务 4,依赖真实服务接线",
"F20": "未测:裸 MQTT 属连接 N,serve 未挂 broker",
}
for id, note := range defaults {
set(id, report.StatusUntested, note)
}
}
func findModuleRoot(t *testing.T) string {
t.Helper()
dir, err := os.Getwd()
if err != nil {
t.Fatal(err)
}
for {
if _, err := os.Stat(filepath.Join(dir, "go.mod")); err == nil {
return dir
}
parent := filepath.Dir(dir)
if parent == dir {
t.Fatal("go.mod not found")
}
dir = parent
}
}
func trim(s string) string {
s = strings.TrimSpace(s)
if len(s) > 180 {
return s[:180] + "..."
}
return s
}
+65
View File
@@ -0,0 +1,65 @@
package accept
import (
"bytes"
"io"
"net/http"
"strings"
"time"
)
// RouteStatus 探测结果。
type RouteStatus int
const (
RouteMissing RouteStatus = iota // 404 / 无此路由
RoutePresent // 路由存在(含 4xx/5xx 业务响应)
)
// ProbeMethod 对运行中的服务发一次 HTTP 请求。
func ProbeMethod(base, method, path string, body []byte) (statusCode int, bodyText string, err error) {
var rdr io.Reader
if body != nil {
rdr = bytes.NewReader(body)
}
req, err := http.NewRequest(method, strings.TrimRight(base, "/")+path, rdr)
if err != nil {
return 0, "", err
}
if body != nil {
req.Header.Set("Content-Type", "application/json")
}
client := &http.Client{Timeout: 10 * time.Second}
resp, err := client.Do(req)
if err != nil {
return 0, "", err
}
defer func() { _ = resp.Body.Close() }()
b, _ := io.ReadAll(io.LimitReader(resp.Body, 64*1024))
return resp.StatusCode, string(b), nil
}
// ClassifyAdminLogin 判断管理登录是否已挂到进程。
// 404 → 未挂载;其它(含 400/401/429)视为已挂载。
func ClassifyAdminLogin(code int) RouteStatus {
if code == http.StatusNotFound {
return RouteMissing
}
return RoutePresent
}
// ClassifyRegister 判断注册路由是否已挂到进程。
func ClassifyRegister(code int) RouteStatus {
if code == http.StatusNotFound {
return RouteMissing
}
return RoutePresent
}
// ClassifyCreateEndpoint 判断开通端路由是否已挂到进程。
func ClassifyCreateEndpoint(code int) RouteStatus {
if code == http.StatusNotFound {
return RouteMissing
}
return RoutePresent
}
+128
View File
@@ -0,0 +1,128 @@
package accept
import (
"bytes"
"fmt"
"os"
"os/exec"
"path/filepath"
"strings"
"time"
"git.asio.asia/nixevol/NixMsg/test/harness"
)
// LoggedServer 在捕获 stdout/stderr 的情况下启动 serve(用于断言密码不进日志)。
type LoggedServer struct {
*harness.Server
LogBuf *bytes.Buffer
cmd *exec.Cmd
}
// StartWithLogCapture 执行 admin init 后启动 serve,并把进程日志写入 LogBuf。
func StartWithLogCapture() (*LoggedServer, error) {
bin, err := harness.Binary()
if err != nil {
return nil, err
}
dataDir, err := os.MkdirTemp("", "nixmsg-accept-*")
if err != nil {
return nil, err
}
cfgPath := filepath.Join(dataDir, "config.yaml")
cfg := fmt.Sprintf("listen: %q\ndata_dir: %q\n", "127.0.0.1:0", filepath.ToSlash(dataDir))
if err = os.WriteFile(cfgPath, []byte(cfg), 0o644); err != nil {
_ = os.RemoveAll(dataDir)
return nil, err
}
initCmd := exec.Command(bin, "admin", "init")
initCmd.Env = append(os.Environ(), "NIXMSG_CONFIG="+cfgPath)
initOut, initErr := initCmd.CombinedOutput()
if initErr != nil {
_ = os.RemoveAll(dataDir)
return nil, fmt.Errorf("admin init: %w\n%s", initErr, initOut)
}
password := parsePassword(string(initOut))
if password == "" {
_ = os.RemoveAll(dataDir)
return nil, fmt.Errorf("admin init password not found:\n%s", initOut)
}
logBuf := &bytes.Buffer{}
cmd := exec.Command(bin, "serve")
cmd.Env = append(os.Environ(), "NIXMSG_CONFIG="+cfgPath)
cmd.Stdout = logBuf
cmd.Stderr = logBuf
if err = cmd.Start(); err != nil {
_ = os.RemoveAll(dataDir)
return nil, fmt.Errorf("start serve: %w", err)
}
s := &harness.Server{
BinPath: bin,
ConfigPath: cfgPath,
DataDir: dataDir,
AdminPassword: password,
}
addr, err := waitListenAddr(filepath.Join(dataDir, "listen.addr"), 15*time.Second)
if err != nil {
_ = cmd.Process.Kill()
_, _ = cmd.Process.Wait()
_ = os.RemoveAll(dataDir)
return nil, fmt.Errorf("wait listen.addr: %w\nlogs:\n%s", err, logBuf.String())
}
s.Addr = addr
s.HTTPBase = "http://" + addr
s.AdminHTTPBase = s.HTTPBase
return &LoggedServer{Server: s, LogBuf: logBuf, cmd: cmd}, nil
}
// Stop 结束进程并删除临时目录。
func (s *LoggedServer) Stop() error {
if s == nil {
return nil
}
if s.cmd != nil && s.cmd.Process != nil {
_ = s.cmd.Process.Kill()
_, _ = s.cmd.Process.Wait()
}
if s.DataDir != "" {
return os.RemoveAll(s.DataDir)
}
return nil
}
func parsePassword(out string) string {
for _, line := range strings.Split(out, "\n") {
line = strings.TrimSpace(line)
lower := strings.ToLower(line)
if strings.HasPrefix(lower, "admin password:") {
return strings.TrimSpace(line[len("admin password:"):])
}
if strings.HasPrefix(lower, "password:") {
return strings.TrimSpace(line[len("password:"):])
}
}
return ""
}
func waitListenAddr(path string, timeout time.Duration) (string, error) {
deadline := time.Now().Add(timeout)
var lastErr error
for time.Now().Before(deadline) {
b, err := os.ReadFile(path)
if err == nil {
addr := strings.TrimSpace(string(b))
if addr != "" {
return addr, nil
}
lastErr = fmt.Errorf("empty addr file")
} else {
lastErr = err
}
time.Sleep(20 * time.Millisecond)
}
return "", lastErr
}
+120
View File
@@ -0,0 +1,120 @@
{
"generated_at": "2026-09-29T23:20:07Z",
"items": [
{
"id": "F01",
"status": "untested",
"note": "未测:serve 未挂载 POST /api/admin/endpoints(A2 未合入或未接线;当前 A1 对端路由返回 501 亦未挂到进程)"
},
{
"id": "F02",
"status": "untested",
"note": "未测:端登录/会话令牌属连接 N3,main 上 serve 未挂 broker"
},
{
"id": "F03",
"status": "untested",
"note": "未测:在线状态属身份 I3 + 连接 N3,未接线"
},
{
"id": "F04",
"status": "untested",
"note": "未测:presence.watch 属身份 I3,未接线"
},
{
"id": "F05",
"status": "untested",
"note": "未测:消息提交属消息 M1,未挂入 broker 上行"
},
{
"id": "F06",
"status": "untested",
"note": "未测:群消息属消息 M + 身份 I4,未接线"
},
{
"id": "F07",
"status": "untested",
"note": "未测:大小限制属消息/连接,未接线"
},
{
"id": "F08",
"status": "untested",
"note": "未测:投递确认属消息 M2,未接线"
},
{
"id": "F09",
"status": "untested",
"note": "未测:保留期属消息 M2,未接线"
},
{
"id": "F10",
"status": "untested",
"note": "未测:断线策略属消息 M2,未接线"
},
{
"id": "F11",
"status": "untested",
"note": "未测:定时发送属消息 M2/M4,未接线"
},
{
"id": "F12",
"status": "untested",
"note": "未测:延迟撤回属消息 M3,未接线"
},
{
"id": "F13",
"status": "untested",
"note": "未测:撤回判定属消息 M3,未接线"
},
{
"id": "F14",
"status": "untested",
"note": "未测:回执属消息 M3,未接线"
},
{
"id": "F15",
"status": "untested",
"note": "未测:对话密码属身份 I2,未接线"
},
{
"id": "F16",
"status": "untested",
"note": "未测:群权限属身份 I4,未接线"
},
{
"id": "F17",
"status": "untested",
"note": "未测:serve 未挂载 /api/admin/login(A1 Handler 已实现,缺总控/接线挂到 cmd/nixmsg;DEVIATIONS 后台接口 A §1)"
},
{
"id": "F18",
"status": "untested",
"note": "未测:正文清理属消息 M3,未接线"
},
{
"id": "F19",
"status": "untested",
"note": "未测:SDK 接入清单属 S1/S2 任务 4,依赖真实服务接线"
},
{
"id": "F20",
"status": "untested",
"note": "未测:裸 MQTT 属连接 N,serve 未挂 broker"
},
{
"id": "F21",
"status": "untested",
"note": "仅验证 listen 上 /healthz+/readyz;后台 API、/mqtt、裸 TCP、注册未挂入 serve(缺总控接线 + 连接 N)"
},
{
"id": "F22",
"status": "pass",
"note": "已测:空目录 admin init + serve,/healthz 与 /readyz 成功,密码不在 serve 日志;未测:备份恢复、升级迁移、证书重载、Docker、/metrics"
},
{
"id": "F23",
"status": "untested",
"note": "未测:serve 未挂载 POST /api/client/register(I1 Handler 已实现,缺总控/连接 N 接线;DEVIATIONS 身份 I §1)"
}
]
}