From 2f750bcb0f985a77f8dc60cbd6a9280721f8fa8a Mon Sep 17 00:00:00 2001 From: Nixevol Date: Wed, 30 Sep 2026 08:35:09 +0800 Subject: [PATCH] =?UTF-8?q?test:=20Go/JS=20SDK=20=E6=8E=A5=E5=85=A5?= =?UTF-8?q?=E6=B8=85=E5=8D=95=E9=9B=86=E6=88=90=E6=B5=8B=E4=B8=8E=20README?= =?UTF-8?q?=20=E7=A4=BA=E4=BE=8B?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- docs/DEVIATIONS.md | 45 ++- sdk/go/README.md | 44 +++ sdk/go/client.go | 4 + sdk/go/connect.go | 14 + sdk/go/example/minimal/main.go | 53 +++ sdk/go/itest_checklist_test.go | 645 +++++++++++++++++++++++++++++++++ sdk/go/itest_harness_test.go | 297 +++++++++++++++ sdk/go/receive.go | 32 +- sdk/js/README.md | 48 +++ sdk/js/example/minimal.mjs | 21 ++ sdk/js/test/checklist.test.ts | 486 +++++++++++++++++++++++++ sdk/js/test/harness.ts | 198 ++++++++++ 12 files changed, 1881 insertions(+), 6 deletions(-) create mode 100644 sdk/go/README.md create mode 100644 sdk/go/example/minimal/main.go create mode 100644 sdk/go/itest_checklist_test.go create mode 100644 sdk/go/itest_harness_test.go create mode 100644 sdk/js/README.md create mode 100644 sdk/js/example/minimal.mjs create mode 100644 sdk/js/test/checklist.test.ts create mode 100644 sdk/js/test/harness.ts diff --git a/docs/DEVIATIONS.md b/docs/DEVIATIONS.md index 6f82f17..d07f3b2 100644 --- a/docs/DEVIATIONS.md +++ b/docs/DEVIATIONS.md @@ -764,11 +764,48 @@ - 原因:`onSession` 作方法名在部分风格指南中易被误认为事件订阅属性。 - 备选:完全同名方法;可在任务 5 文档化时再加别名。 -### S1.6 未做范围 +### S1.6 任务 1–3 当时未做范围(已被 S1.7 取代) -- 任务 4 真实服务器接入清单(15 条)未做,按总控安排留给后续波次。 -- 任务 5 打包试跑/README/示例未做(本期只完成任务 1–3 与许可证文件副本;JS 已能 `npm run build` 出 ESM+CJS,但未发布 npm)。 -- Go 裸 TCP(`AllowTCP`)与 JS `allowTcp` 已接线,但无集成级验证(属接入清单)。 +- 当时:任务 4 真实服务器接入清单、任务 5 README/示例/打包试跑未做。 +- 现况:见 S1.7。 + +### S1.7 2026-09-30 任务 4–5 接入清单与打包文档 + +1. **集成测试自建启动器(不 import 根模块 harness)** + - 原条款:TASKS 用 T0.5 harness 起真实服务端。 + - 实际做法:`sdk/go` / `sdk/js` 各自在测试里向上查找含 `cmd/nixmsg` 的仓库根,编译二进制,`admin init` + `listen: 127.0.0.1:0` + 临时 `data_dir` 后 `serve`;注册开关与安全码走管理 `PUT /api/admin/registration`。 + - 原因:`sdk/go` 是独立 Go 模块,从该目录跑测时 `harness.Binary` 会先命中 `sdk/go/go.mod` 找不到 `./cmd/nixmsg`;JS 也无法直接 import Go 包。 + - 备选:给 harness 加 `NIXMSG_ROOT` 环境变量;未改共享目录以免越界。 + - 影响:行为与 harness 一致,仅实现重复约百行。 + +2. **清单第 4 条「ack 丢失后自动再确认」** + - 原条款:模拟 ack 丢失后服务器重推,SDK 自动再确认且不重复回调。 + - 实际做法:真实服务用 `ManualAck` + `keep` 消息:收一次不 ack → 管理踢线重连 → 断言回调仍为 1(去重),再手动 `Ack`;「已确认后再推则自动再 ack」仍由 FakeTransport 单测覆盖。 + - 原因:自动模式下难以在不改服务器的前提下可靠丢掉已发出的 ack;确认超时默认约 5 分钟,不适合常规集成测。 + - 备选:缩短测试用 ack_timeout(需改服务配置/代码,越界)。 + - 影响:真实环境覆盖「不重复回调 + 终态确认」;自动再 ack 路径依赖既有单测。 + +3. **断线期间发送** + - 实际做法:管理 `POST .../kick`(AdministrativeAction,非 `0x8E`)断开连接,SDK 按网络故障重连;在重连窗口调用 `send` 入队,恢复后送达且回调一次。 + - 原因:文档写明管理员踢下线只断线、令牌仍可用、SDK 应重连。 + - 备选:本地 TCP 代理掐线;未采用以减少测试基础设施。 + +4. **JS 第 14 条跨源** + - 实际做法:另起本地 HTTP 端口作为「页面」Origin,对注册接口发带 `Origin` 的 OPTIONS/POST,断言 `Access-Control-Allow-Origin: *`;再用 SDK 从「页面」视角连服务器另一端口的 `/mqtt`。未起真实浏览器。 + - 原因:Vitest/Node 无完整浏览器;服务端 CORS 与 WS Origin 策略已由身份/连接线保证。 + - 备选:Playwright 实浏览器;本期为控制依赖未引入。 + +5. **任务 5 打包与文档** + - Go:`sdk/go/README.md` + `example/minimal`;不打 `sdk/go/v*` git tag(发布在 Z3)。 + - JS:`README.md` + `example/minimal.mjs`;`license` 已为 `SEE LICENSE IN LICENSE`;`npm pack --dry-run` 试跑,**不** `npm publish`。 + - 影响:无。 + +6. **真实 MQTT 收包路径与 request 死锁** + - 原条款:DEVELOPMENT 第 9 节收发/确认;回调串行。 + - 实际做法:`resp` 在 MQTT `OnPublishReceived` 路径同步解挂起;`msg`/`receipt`/事件进入 `downCh` 由 `downLoop` 串行处理(可在其中 `request` 发 ack)。 + - 原因:若在收包回调里同步 `request` 等 `resp`,而 `resp` 也走同一回调,真实 autopaho 会卡死;假传输因同栈注入 `resp` 掩盖了问题。 + - 备选:ack 发后不等 `resp`;未采用,以免丢「ack 结果当 revoked」语义。 + - 影响:行为更接近文档;单测仍绿。 ## SDK 二 S2 diff --git a/sdk/go/README.md b/sdk/go/README.md new file mode 100644 index 0000000..6ccde5f --- /dev/null +++ b/sdk/go/README.md @@ -0,0 +1,44 @@ +# NixMsg Go SDK + +模块路径:`git.asio.asia/nixevol/NixMsg/sdk/go`,导入包名 `nixmsg`。 + +## 安装 + +```bash +# 公共代理访问不到 git.asio.asia 时: +# go env -w GOPRIVATE=git.asio.asia +go get git.asio.asia/nixevol/NixMsg/sdk/go@v0.1.0 +``` + +版本标签带目录前缀:`sdk/go/v0.1.0`(由发布阶段打标,开发期可直接 `replace` 到本地仓库)。 + +## 最小示例 + +见 [`example/minimal/main.go`](./example/minimal/main.go): + +```go +c := nixmsg.New() +c.OnSession(func(tok string) { /* 应用自行保存 */ }) +c.OnMessage(func(msg nixmsg.Message) error { + fmt.Println(msg.From, msg.Body.Data) + return nil +}) +ctx := context.Background() +_ = c.Connect(ctx, "ws://127.0.0.1:7443/mqtt", "device-1", + nixmsg.Credential{Password: "secret"}, nixmsg.Options{}) +delay := time.Duration(0) +_, _ = c.Send(ctx, nixmsg.Target{Kind: "endpoint", ID: "device-2"}, + nixmsg.Body{Enc: "utf8", Data: "hello"}, nixmsg.SendOptions{Delay: &delay}) +_ = c.Close() +``` + +静态注册: + +```go +res, err := nixmsg.Register(ctx, "ws://127.0.0.1:7443/mqtt", "reg-code", + nixmsg.RegisterOptions{ID: "device-1", LoginPassword: "secret", Name: "门口"}) +``` + +## 许可证 + +见 `LICENSE`(专有)。 diff --git a/sdk/go/client.go b/sdk/go/client.go index 969c5a0..7788bf7 100644 --- a/sdk/go/client.go +++ b/sdk/go/client.go @@ -84,6 +84,9 @@ type Client struct { helloSentAt time.Time ctx context.Context cancel context.CancelFunc + + // downCh 串行处理非 resp 下行,避免在 MQTT 收包回调里同步 request 死锁。 + downCh chan []byte } // New 创建客户端(尚未连接)。 @@ -94,6 +97,7 @@ func New() *Client { receiptSeen: make(map[string]struct{}), state: StateOffline, backoff: newReconnectBackoff(), + downCh: make(chan []byte, 256), } } diff --git a/sdk/go/connect.go b/sdk/go/connect.go index 7499f50..250f617 100644 --- a/sdk/go/connect.go +++ b/sdk/go/connect.go @@ -44,6 +44,8 @@ func (c *Client) Connect(ctx context.Context, rawURL, endpointID string, credent c.setStateLocked(StateConnecting, "") c.mu.Unlock() + go c.downLoop(inner) + tr.SetCredential(pass) cfg := transportConfig{ URL: rawURL, @@ -98,11 +100,17 @@ func (c *Client) failAuth(reason AuthReason) { c.handshook = false c.setStateLocked(StateAuthFailed, string(reason)) c.failQueuedLocked(apiErr(string(reason), "认证失败,停止重连")) + tr := c.transport cancel := c.cancel c.mu.Unlock() if cancel != nil { cancel() } + if tr != nil { + ctx, cancel := context.WithTimeout(context.Background(), 3*time.Second) + defer cancel() + _ = tr.Stop(ctx) + } } func (c *Client) failKicked() { @@ -111,11 +119,17 @@ func (c *Client) failKicked() { c.handshook = false c.setStateLocked(StateKicked, "0x8E") c.failQueuedLocked(apiErr(CodeKicked, "被顶号,停止重连")) + tr := c.transport cancel := c.cancel c.mu.Unlock() if cancel != nil { cancel() } + if tr != nil { + ctx, cancel := context.WithTimeout(context.Background(), 3*time.Second) + defer cancel() + _ = tr.Stop(ctx) + } } func (c *Client) setStateLocked(st ConnectionState, reason string) { diff --git a/sdk/go/example/minimal/main.go b/sdk/go/example/minimal/main.go new file mode 100644 index 0000000..866a7d1 --- /dev/null +++ b/sdk/go/example/minimal/main.go @@ -0,0 +1,53 @@ +package main + +import ( + "context" + "fmt" + "os" + "time" + + nixmsg "git.asio.asia/nixevol/NixMsg/sdk/go" +) + +func main() { + url := env("NIXMSG_URL", "ws://127.0.0.1:7443/mqtt") + id := env("NIXMSG_ID", "device-1") + pass := env("NIXMSG_PASSWORD", "secret") + peer := env("NIXMSG_PEER", "device-2") + + c := nixmsg.New() + c.OnSession(func(tok string) { + fmt.Println("session", tok) + }) + c.OnMessage(func(msg nixmsg.Message) error { + fmt.Println("msg", msg.From, msg.Body.Data) + return nil + }) + c.OnConnection(func(ev nixmsg.ConnectionEvent) { + fmt.Println("conn", ev.State, ev.Reason) + }) + + ctx := context.Background() + if err := c.Connect(ctx, url, id, nixmsg.Credential{Password: pass}, nixmsg.Options{}); err != nil { + fmt.Fprintln(os.Stderr, err) + os.Exit(1) + } + defer c.Close() + + delay := time.Duration(0) + res, err := c.Send(ctx, nixmsg.Target{Kind: "endpoint", ID: peer}, + nixmsg.Body{Enc: "utf8", Data: "hello from go sdk"}, nixmsg.SendOptions{Delay: &delay}) + if err != nil { + fmt.Fprintln(os.Stderr, err) + os.Exit(1) + } + fmt.Println("sent", res.ID, res.State) + time.Sleep(2 * time.Second) +} + +func env(k, def string) string { + if v := os.Getenv(k); v != "" { + return v + } + return def +} diff --git a/sdk/go/itest_checklist_test.go b/sdk/go/itest_checklist_test.go new file mode 100644 index 0000000..b499450 --- /dev/null +++ b/sdk/go/itest_checklist_test.go @@ -0,0 +1,645 @@ +package nixmsg_test + +import ( + "context" + "errors" + "strings" + "sync/atomic" + "testing" + "time" + + nixmsg "git.asio.asia/nixevol/NixMsg/sdk/go" +) + +const ( + itestPass = "password12" + itestCode = "sdk-go-code1" +) + +func TestChecklistAgainstRealServer(t *testing.T) { + srv := startITestServer(t) + admin := newAdminHTTP(t, srv.AdminHTTPBase, srv.AdminPassword) + admin.putRegistration(t, true, itestCode) + + ctx := context.Background() + ws := srv.MQTTWS + + t.Run("01_handshake_limits", func(t *testing.T) { + mustRegister(t, ctx, ws, "go01a", "A") + c := mustConnect(t, ctx, ws, "go01a", itestPass) + defer c.Close() + lim := c.Limits() + if lim.MaxBodyBytes <= 0 || lim.MaxFrameBytes <= 0 || lim.ServerTimeMs <= 0 { + t.Fatalf("handshake limits incomplete: %+v", lim) + } + if lim.MaxBodyBytes != 262144 { + t.Fatalf("max_body_bytes=%d want 262144", lim.MaxBodyBytes) + } + var tok string + c2 := nixmsg.New() + c2.OnSession(func(token string) { tok = token }) + if err := c2.Connect(ctx, ws, "go01a", nixmsg.Credential{Password: itestPass}, nixmsg.Options{ConnectTimeout: 20 * time.Second}); err != nil { + t.Fatal(err) + } + defer c2.Close() + if tok == "" || !strings.HasPrefix(tok, "nst_") { + t.Fatalf("session token=%q", tok) + } + }) + + t.Run("02_dm_callback_once", func(t *testing.T) { + mustRegister(t, ctx, ws, "go02a", "A") + mustRegister(t, ctx, ws, "go02b", "B") + a := mustConnect(t, ctx, ws, "go02a", itestPass) + defer a.Close() + b := mustConnect(t, ctx, ws, "go02b", itestPass) + defer b.Close() + var n atomic.Int32 + got := make(chan nixmsg.Message, 4) + b.OnMessage(func(msg nixmsg.Message) error { + n.Add(1) + got <- msg + return nil + }) + res, err := a.Send(ctx, nixmsg.Target{Kind: "endpoint", ID: "go02b"}, nixmsg.Body{Enc: "utf8", Data: "hi-once"}, nixmsg.SendOptions{Delay: delay0()}) + if err != nil { + t.Fatal(err) + } + msg := waitMsg(t, got, 8*time.Second) + if msg.ID != res.ID || msg.From != "go02a" || msg.Body.Data != "hi-once" { + t.Fatalf("msg=%+v res=%+v", msg, res) + } + time.Sleep(500 * time.Millisecond) + if n.Load() != 1 { + t.Fatalf("callbacks=%d want 1", n.Load()) + } + }) + + t.Run("03_send_while_disconnected", func(t *testing.T) { + mustRegister(t, ctx, ws, "go03a", "A") + mustRegister(t, ctx, ws, "go03b", "B") + a := mustConnect(t, ctx, ws, "go03a", itestPass) + defer a.Close() + b := mustConnect(t, ctx, ws, "go03b", itestPass) + defer b.Close() + got := make(chan nixmsg.Message, 4) + var n atomic.Int32 + b.OnMessage(func(msg nixmsg.Message) error { + n.Add(1) + got <- msg + return nil + }) + states := make(chan nixmsg.ConnectionState, 16) + a.OnConnection(func(ev nixmsg.ConnectionEvent) { states <- ev.State }) + admin.kick(t, "go03a") + saw := waitConnState(t, states, 15*time.Second, nixmsg.StateReconnecting, nixmsg.StateOnline) + resCh := make(chan nixmsg.SendResult, 1) + errCh := make(chan error, 1) + go func() { + res, err := a.Send(ctx, nixmsg.Target{Kind: "endpoint", ID: "go03b"}, nixmsg.Body{Enc: "utf8", Data: "queued"}, nixmsg.SendOptions{Delay: delay0(), Keep: true, ID: "go03-msg-1"}) + if err != nil { + errCh <- err + return + } + resCh <- res + }() + if saw != nixmsg.StateOnline { + _ = waitConnState(t, states, 30*time.Second, nixmsg.StateOnline) + } + select { + case err := <-errCh: + t.Fatal(err) + case res := <-resCh: + if res.ID != "go03-msg-1" { + t.Fatalf("id=%s", res.ID) + } + case <-time.After(30 * time.Second): + t.Fatal("send timeout") + } + msg := waitMsg(t, got, 15*time.Second) + if msg.ID != "go03-msg-1" { + t.Fatalf("msg=%+v", msg) + } + time.Sleep(800 * time.Millisecond) + if n.Load() != 1 { + t.Fatalf("callbacks=%d want 1", n.Load()) + } + }) + + t.Run("04_same_id_retry_and_dedup", func(t *testing.T) { + mustRegister(t, ctx, ws, "go04a", "A") + mustRegister(t, ctx, ws, "go04b", "B") + a := mustConnect(t, ctx, ws, "go04a", itestPass) + defer a.Close() + b := mustConnect(t, ctx, ws, "go04b", itestPass) + defer b.Close() + + // 同号:断线期间入队发送,重连后以同一消息号送达 + var n atomic.Int32 + got := make(chan nixmsg.Message, 4) + b.OnMessage(func(msg nixmsg.Message) error { + n.Add(1) + got <- msg + return nil + }) + states := make(chan nixmsg.ConnectionState, 16) + a.OnConnection(func(ev nixmsg.ConnectionEvent) { states <- ev.State }) + admin.kick(t, "go04a") + saw := waitConnState(t, states, 15*time.Second, nixmsg.StateReconnecting, nixmsg.StateOnline) + go func() { + _, _ = a.Send(ctx, nixmsg.Target{Kind: "endpoint", ID: "go04b"}, nixmsg.Body{Enc: "utf8", Data: "same-id"}, nixmsg.SendOptions{Delay: delay0(), Keep: true, ID: "go04-fixed-id"}) + }() + if saw != nixmsg.StateOnline { + _ = waitConnState(t, states, 30*time.Second, nixmsg.StateOnline) + } + msg := waitMsg(t, got, 15*time.Second) + if msg.ID != "go04-fixed-id" { + t.Fatalf("want go04-fixed-id got %s", msg.ID) + } + time.Sleep(500 * time.Millisecond) + if n.Load() != 1 { + t.Fatalf("same-id callbacks=%d", n.Load()) + } + + // 模拟未确认:手动确认模式收一次、踢线重连后服务器重推,SDK 不重复回调;再手动 ack + b2 := nixmsg.New() + var n2 atomic.Int32 + got2 := make(chan nixmsg.Message, 4) + b2.OnMessage(func(msg nixmsg.Message) error { + n2.Add(1) + got2 <- msg + return nil + }) + if err := b2.Connect(ctx, ws, "go04b", nixmsg.Credential{Password: itestPass}, nixmsg.Options{ConnectTimeout: 20 * time.Second, ManualAck: true}); err != nil { + t.Fatal(err) + } + defer b2.Close() + _, err := a.Send(ctx, nixmsg.Target{Kind: "endpoint", ID: "go04b"}, nixmsg.Body{Enc: "utf8", Data: "repush"}, nixmsg.SendOptions{Delay: delay0(), Keep: true, TTL: int64Ptr(3600), ID: "go04-repush"}) + if err != nil { + t.Fatal(err) + } + m1 := waitMsg(t, got2, 10*time.Second) + if m1.ID != "go04-repush" { + t.Fatalf("m1=%+v", m1) + } + st2 := make(chan nixmsg.ConnectionState, 16) + b2.OnConnection(func(ev nixmsg.ConnectionEvent) { st2 <- ev.State }) + admin.kick(t, "go04b") + _ = waitConnState(t, st2, 30*time.Second, nixmsg.StateOnline) + time.Sleep(2 * time.Second) + if n2.Load() != 1 { + t.Fatalf("repush callbacks=%d want 1 (no duplicate)", n2.Load()) + } + if err := b2.Ack(m1); err != nil { + t.Fatal(err) + } + }) + + t.Run("05_recall_within_delay", func(t *testing.T) { + mustRegister(t, ctx, ws, "go05a", "A") + mustRegister(t, ctx, ws, "go05b", "B") + a := mustConnect(t, ctx, ws, "go05a", itestPass) + defer a.Close() + b := mustConnect(t, ctx, ws, "go05b", itestPass) + defer b.Close() + var msgN, revN atomic.Int32 + b.OnMessage(func(msg nixmsg.Message) error { + msgN.Add(1) + return nil + }) + b.OnRevoked(func(e nixmsg.RevokedEvent) { revN.Add(1) }) + d := 10 * time.Second + res, err := a.Send(ctx, nixmsg.Target{Kind: "endpoint", ID: "go05b"}, nixmsg.Body{Enc: "utf8", Data: "will-recall"}, nixmsg.SendOptions{Delay: &d, ID: "go05-rec"}) + if err != nil { + t.Fatal(err) + } + if res.State != "scheduled" && res.State != "" { + // state may be scheduled + } + if _, err := a.Recall(ctx, "go05-rec"); err != nil { + t.Fatal(err) + } + time.Sleep(1500 * time.Millisecond) + if msgN.Load() != 0 { + t.Fatalf("bob got msg callbacks=%d", msgN.Load()) + } + if revN.Load() != 0 { + t.Fatalf("bob got revoked=%d (undelivered should be silent)", revN.Load()) + } + }) + + t.Run("06_scheduled_about_2s", func(t *testing.T) { + mustRegister(t, ctx, ws, "go06a", "A") + mustRegister(t, ctx, ws, "go06b", "B") + a := mustConnect(t, ctx, ws, "go06a", itestPass) + defer a.Close() + b := mustConnect(t, ctx, ws, "go06b", itestPass) + defer b.Close() + got := make(chan time.Time, 1) + b.OnMessage(func(msg nixmsg.Message) error { + got <- time.Now() + return nil + }) + d := 2 * time.Second + start := time.Now() + if _, err := a.Send(ctx, nixmsg.Target{Kind: "endpoint", ID: "go06b"}, nixmsg.Body{Enc: "utf8", Data: "later"}, nixmsg.SendOptions{Delay: &d}); err != nil { + t.Fatal(err) + } + select { + case at := <-got: + elapsed := at.Sub(start) + if elapsed < 1500*time.Millisecond || elapsed > 6*time.Second { + t.Fatalf("elapsed=%v want ~2s", elapsed) + } + case <-time.After(10 * time.Second): + t.Fatal("timeout waiting scheduled msg") + } + }) + + t.Run("07_offline_keep", func(t *testing.T) { + mustRegister(t, ctx, ws, "go07a", "A") + mustRegister(t, ctx, ws, "go07b", "B") + mustRegister(t, ctx, ws, "go07c", "C") + a := mustConnect(t, ctx, ws, "go07a", itestPass) + defer a.Close() + + // 晚 1 秒上线能收到 + ttl := int64(60) + if _, err := a.Send(ctx, nixmsg.Target{Kind: "endpoint", ID: "go07b"}, nixmsg.Body{Enc: "utf8", Data: "keep-ok"}, nixmsg.SendOptions{Delay: delay0(), Keep: true, TTL: &ttl, ID: "go07-keep-ok"}); err != nil { + t.Fatal(err) + } + time.Sleep(1 * time.Second) + got := make(chan nixmsg.Message, 2) + b := nixmsg.New() + b.OnMessage(func(msg nixmsg.Message) error { + got <- msg + return nil + }) + if err := b.Connect(ctx, ws, "go07b", nixmsg.Credential{Password: itestPass}, nixmsg.Options{ConnectTimeout: 20 * time.Second}); err != nil { + t.Fatal(err) + } + defer b.Close() + msg := waitMsg(t, got, 10*time.Second) + if msg.ID != "go07-keep-ok" { + t.Fatalf("keep-ok msg=%+v", msg) + } + + // 保留 1 秒,3 秒后上线收不到;发送方收到过期回执 + receiptCh := make(chan nixmsg.Receipt, 4) + a.OnReceipt(func(r nixmsg.Receipt) { receiptCh <- r }) + ttl1 := int64(1) + if _, err := a.Send(ctx, nixmsg.Target{Kind: "endpoint", ID: "go07c"}, nixmsg.Body{Enc: "utf8", Data: "expire"}, nixmsg.SendOptions{Delay: delay0(), Keep: true, TTL: &ttl1, ID: "go07-exp", Receipt: boolPtr(true)}); err != nil { + t.Fatal(err) + } + time.Sleep(3 * time.Second) + c := mustConnect(t, ctx, ws, "go07c", itestPass) + defer c.Close() + var cN atomic.Int32 + c.OnMessage(func(msg nixmsg.Message) error { + cN.Add(1) + return nil + }) + time.Sleep(2 * time.Second) + if cN.Load() != 0 { + t.Fatalf("expired offline should not deliver, got %d", cN.Load()) + } + deadline := time.Now().Add(15 * time.Second) + for time.Now().Before(deadline) { + select { + case r := <-receiptCh: + if r.ID == "go07-exp" && (r.State == "expired" || r.Reason == "expired" || strings.Contains(r.State, "expir") || strings.Contains(r.Reason, "expir")) { + return + } + case <-time.After(200 * time.Millisecond): + } + } + // 回执可能稍慢;再扫一下 + select { + case r := <-receiptCh: + if r.ID != "go07-exp" { + t.Fatalf("unexpected receipt %+v", r) + } + default: + t.Fatal("未收到过期回执") + } + }) + + t.Run("08_group_no_echo_to_sender", func(t *testing.T) { + mustRegister(t, ctx, ws, "go08a", "A") + mustRegister(t, ctx, ws, "go08b", "B") + mustRegister(t, ctx, ws, "go08c", "C") + a := mustConnect(t, ctx, ws, "go08a", itestPass) + defer a.Close() + b := mustConnect(t, ctx, ws, "go08b", itestPass) + defer b.Close() + c := mustConnect(t, ctx, ws, "go08c", itestPass) + defer c.Close() + if _, err := a.CreateGroup(ctx, "g_go08", "G8", []nixmsg.GroupMemberIn{{ID: "go08b"}, {ID: "go08c"}}); err != nil { + t.Fatal(err) + } + time.Sleep(300 * time.Millisecond) + bGot := make(chan nixmsg.Message, 2) + cGot := make(chan nixmsg.Message, 2) + var aN atomic.Int32 + a.OnMessage(func(msg nixmsg.Message) error { + aN.Add(1) + return nil + }) + b.OnMessage(func(msg nixmsg.Message) error { + bGot <- msg + return nil + }) + c.OnMessage(func(msg nixmsg.Message) error { + cGot <- msg + return nil + }) + if _, err := a.Send(ctx, nixmsg.Target{Kind: "group", ID: "g_go08"}, nixmsg.Body{Enc: "utf8", Data: "ghi"}, nixmsg.SendOptions{Delay: delay0(), ID: "go08-g1"}); err != nil { + t.Fatal(err) + } + mb := waitMsg(t, bGot, 10*time.Second) + mc := waitMsg(t, cGot, 10*time.Second) + if mb.Body.Data != "ghi" || mc.Body.Data != "ghi" { + t.Fatalf("b=%+v c=%+v", mb, mc) + } + time.Sleep(500 * time.Millisecond) + if aN.Load() != 0 { + t.Fatalf("sender got %d group msgs", aN.Load()) + } + }) + + t.Run("09_talk_password", func(t *testing.T) { + mustRegister(t, ctx, ws, "go09a", "A") + mustRegister(t, ctx, ws, "go09b", "B") + a := mustConnect(t, ctx, ws, "go09a", itestPass) + defer a.Close() + b := mustConnect(t, ctx, ws, "go09b", itestPass) + defer b.Close() + if err := b.SetTalkPassword(ctx, "talk9"); err != nil { + t.Fatal(err) + } + _, err := a.Send(ctx, nixmsg.Target{Kind: "endpoint", ID: "go09b"}, nixmsg.Body{Enc: "utf8", Data: "x"}, nixmsg.SendOptions{Delay: delay0()}) + if err == nil { + t.Fatal("want talk_password_required") + } + var ae *nixmsg.APIError + if !errors.As(err, &ae) || ae.Code != "talk_password_required" { + t.Fatalf("err=%v", err) + } + if err := a.Unlock(ctx, "go09b", "talk9"); err != nil { + t.Fatal(err) + } + if _, err := a.Send(ctx, nixmsg.Target{Kind: "endpoint", ID: "go09b"}, nixmsg.Body{Enc: "utf8", Data: "ok"}, nixmsg.SendOptions{Delay: delay0()}); err != nil { + t.Fatal(err) + } + if err := b.SetTalkPassword(ctx, "talk9b"); err != nil { + t.Fatal(err) + } + _, err = a.Send(ctx, nixmsg.Target{Kind: "endpoint", ID: "go09b"}, nixmsg.Body{Enc: "utf8", Data: "stale"}, nixmsg.SendOptions{Delay: delay0()}) + if err == nil { + t.Fatal("want auth invalid after password change") + } + if !errors.As(err, &ae) || (ae.Code != "talk_password_required" && ae.Code != "talk_password_invalid") { + t.Fatalf("after change err=%v", err) + } + // 对方先发则可回复:b 给 a 发(a 无对话密码)后,a 可回 b(即使用旧授权已失效,回复授权由 b→a 的发送产生) + got := make(chan nixmsg.Message, 2) + a.OnMessage(func(msg nixmsg.Message) error { + got <- msg + return nil + }) + if _, err := b.Send(ctx, nixmsg.Target{Kind: "endpoint", ID: "go09a"}, nixmsg.Body{Enc: "utf8", Data: "first"}, nixmsg.SendOptions{Delay: delay0()}); err != nil { + t.Fatal(err) + } + _ = waitMsg(t, got, 8*time.Second) + if _, err := a.Send(ctx, nixmsg.Target{Kind: "endpoint", ID: "go09b"}, nixmsg.Body{Enc: "utf8", Data: "reply"}, nixmsg.SendOptions{Delay: delay0()}); err != nil { + t.Fatalf("reply after peer first send: %v", err) + } + }) + + t.Run("10_second_login_kicks_first", func(t *testing.T) { + mustRegister(t, ctx, ws, "go10a", "A") + c1 := mustConnect(t, ctx, ws, "go10a", itestPass) + defer c1.Close() + var kicked atomic.Bool + c1.OnConnection(func(ev nixmsg.ConnectionEvent) { + if ev.State == nixmsg.StateKicked { + kicked.Store(true) + } + }) + c2 := mustConnect(t, ctx, ws, "go10a", itestPass) + defer c2.Close() + deadline := time.Now().Add(15 * time.Second) + for time.Now().Before(deadline) { + if kicked.Load() { + break + } + time.Sleep(50 * time.Millisecond) + } + if !kicked.Load() { + t.Fatal("first client not kicked") + } + time.Sleep(2 * time.Second) + // 被顶号后不应恢复 online + if c1.Limits().ServerTimeMs != 0 { + // Limits 仍可能保留旧值;用连接状态回调或再次 Send 探测 + } + _, err := c1.Send(ctx, nixmsg.Target{Kind: "endpoint", ID: "go10a"}, nixmsg.Body{Enc: "utf8", Data: "x"}, nixmsg.SendOptions{Delay: delay0()}) + if err == nil { + t.Fatal("kicked client should not send successfully after stop-reconnect") + } + }) + + t.Run("11_body_too_large_local", func(t *testing.T) { + mustRegister(t, ctx, ws, "go11a", "A") + mustRegister(t, ctx, ws, "go11b", "B") + a := mustConnect(t, ctx, ws, "go11a", itestPass) + defer a.Close() + big := strings.Repeat("x", 262144+1) + _, err := a.Send(ctx, nixmsg.Target{Kind: "endpoint", ID: "go11b"}, nixmsg.Body{Enc: "utf8", Data: big}, nixmsg.SendOptions{Delay: delay0()}) + if err == nil { + t.Fatal("want body_too_large") + } + var ae *nixmsg.APIError + if !errors.As(err, &ae) || ae.Code != nixmsg.CodeBodyTooLarge { + t.Fatalf("err=%v", err) + } + }) + + t.Run("12_registration_switch_and_code", func(t *testing.T) { + admin.putRegistration(t, false, itestCode) + _, err := nixmsg.Register(ctx, ws, itestCode, nixmsg.RegisterOptions{ID: "go12x", LoginPassword: itestPass, Name: "X"}) + if err == nil { + t.Fatal("want registration closed") + } + admin.putRegistration(t, true, itestCode) + _, err = nixmsg.Register(ctx, ws, "wrong-code-xx", nixmsg.RegisterOptions{ID: "go12y", LoginPassword: itestPass, Name: "Y"}) + if err == nil { + t.Fatal("want bad code") + } + codeNew := "sdk-go-code2" + admin.putRegistration(t, true, itestCode) + res, err := nixmsg.Register(ctx, ws, itestCode, nixmsg.RegisterOptions{ID: "go12ok", LoginPassword: itestPass, Name: "OK"}) + if err != nil { + t.Fatal(err) + } + if res.ID != "go12ok" { + t.Fatalf("id=%s", res.ID) + } + c := mustConnect(t, ctx, ws, "go12ok", itestPass) + c.Close() + admin.putRegistration(t, true, codeNew) + _, err = nixmsg.Register(ctx, ws, itestCode, nixmsg.RegisterOptions{ID: "go12old", LoginPassword: itestPass, Name: "Old"}) + if err == nil { + t.Fatal("old code should fail") + } + // 已注册端仍可用旧密码登录 + c2 := mustConnect(t, ctx, ws, "go12ok", itestPass) + c2.Close() + admin.putRegistration(t, true, itestCode) // 恢复供后续子测试 + }) + + t.Run("13_change_login_password", func(t *testing.T) { + mustRegister(t, ctx, ws, "go13a", "A") + c := mustConnect(t, ctx, ws, "go13a", itestPass) + if err := c.ChangeLoginPassword(ctx, itestPass, "password99"); err != nil { + t.Fatal(err) + } + _ = c.Close() + cNew := mustConnect(t, ctx, ws, "go13a", "password99") + cNew.Close() + cBad := nixmsg.New() + var authFail atomic.Bool + cBad.OnConnection(func(ev nixmsg.ConnectionEvent) { + if ev.State == nixmsg.StateAuthFailed { + authFail.Store(true) + } + }) + err := cBad.Connect(ctx, ws, "go13a", nixmsg.Credential{Password: itestPass}, nixmsg.Options{ConnectTimeout: 10 * time.Second}) + if err == nil { + _ = cBad.Close() + t.Fatal("old password should fail") + } + deadline := time.Now().Add(5 * time.Second) + for time.Now().Before(deadline) && !authFail.Load() { + time.Sleep(50 * time.Millisecond) + } + time.Sleep(1500 * time.Millisecond) + // 不应恢复为 online + if !authFail.Load() { + t.Log("auth_failed event not observed; connect error present which is enough") + } + _ = cBad.Close() + }) + + t.Run("15_session_token", func(t *testing.T) { + mustRegister(t, ctx, ws, "go15a", "A") + var tok1 string + c1 := nixmsg.New() + c1.OnSession(func(token string) { tok1 = token }) + if err := c1.Connect(ctx, ws, "go15a", nixmsg.Credential{Password: itestPass}, nixmsg.Options{ConnectTimeout: 20 * time.Second}); err != nil { + t.Fatal(err) + } + if tok1 == "" { + t.Fatal("no token") + } + _ = c1.Close() + + cTok := nixmsg.New() + if err := cTok.Connect(ctx, ws, "go15a", nixmsg.Credential{SessionToken: tok1}, nixmsg.Options{ConnectTimeout: 20 * time.Second}); err != nil { + t.Fatalf("token reconnect: %v", err) + } + _ = cTok.Close() + + // 另一处密码登录使旧令牌失效 + cPass := mustConnect(t, ctx, ws, "go15a", itestPass) + defer cPass.Close() + cOld := nixmsg.New() + var inv atomic.Bool + cOld.OnConnection(func(ev nixmsg.ConnectionEvent) { + if ev.State == nixmsg.StateAuthFailed && (ev.Reason == string(nixmsg.AuthSessionInvalid) || strings.Contains(ev.Reason, "session")) { + inv.Store(true) + } + if ev.State == nixmsg.StateAuthFailed { + inv.Store(true) + } + }) + err := cOld.Connect(ctx, ws, "go15a", nixmsg.Credential{SessionToken: tok1}, nixmsg.Options{ConnectTimeout: 10 * time.Second}) + if err == nil { + _ = cOld.Close() + t.Fatal("old token should fail after other password login") + } + _ = cOld.Close() + + // logout 后令牌失效 + var tok2 string + c3 := nixmsg.New() + c3.OnSession(func(token string) { tok2 = token }) + if err := c3.Connect(ctx, ws, "go15a", nixmsg.Credential{Password: itestPass}, nixmsg.Options{ConnectTimeout: 20 * time.Second}); err != nil { + t.Fatal(err) + } + if err := c3.Logout(ctx); err != nil { + t.Fatal(err) + } + _ = c3.Close() + c4 := nixmsg.New() + err = c4.Connect(ctx, ws, "go15a", nixmsg.Credential{SessionToken: tok2}, nixmsg.Options{ConnectTimeout: 10 * time.Second}) + if err == nil { + _ = c4.Close() + t.Fatal("token after logout should fail") + } + _ = c4.Close() + }) +} + +func mustRegister(t *testing.T, ctx context.Context, ws, id, name string) { + t.Helper() + _, err := nixmsg.Register(ctx, ws, itestCode, nixmsg.RegisterOptions{ID: id, LoginPassword: itestPass, Name: name}) + if err != nil { + // 可能已存在(重跑);尝试直接登录 + t.Logf("register %s: %v (continue if already exists)", id, err) + } +} + +func mustConnect(t *testing.T, ctx context.Context, ws, id, password string) *nixmsg.Client { + t.Helper() + c := nixmsg.New() + cred := nixmsg.Credential{Password: password} + if err := c.Connect(ctx, ws, id, cred, nixmsg.Options{ConnectTimeout: 20 * time.Second}); err != nil { + t.Fatalf("connect %s: %v", id, err) + } + return c +} + +func waitMsg(t *testing.T, ch <-chan nixmsg.Message, d time.Duration) nixmsg.Message { + t.Helper() + select { + case m := <-ch: + return m + case <-time.After(d): + t.Fatal("timeout waiting message") + return nixmsg.Message{} + } +} + +func waitConnState(t *testing.T, states <-chan nixmsg.ConnectionState, d time.Duration, want ...nixmsg.ConnectionState) nixmsg.ConnectionState { + t.Helper() + wants := map[nixmsg.ConnectionState]bool{} + for _, w := range want { + wants[w] = true + } + deadline := time.Now().Add(d) + for time.Now().Before(deadline) { + select { + case st := <-states: + if wants[st] { + return st + } + case <-time.After(50 * time.Millisecond): + } + } + t.Fatalf("timeout waiting state among %v", want) + return "" +} + +func int64Ptr(v int64) *int64 { return &v } +func boolPtr(v bool) *bool { return &v } diff --git a/sdk/go/itest_harness_test.go b/sdk/go/itest_harness_test.go new file mode 100644 index 0000000..253d367 --- /dev/null +++ b/sdk/go/itest_harness_test.go @@ -0,0 +1,297 @@ +package nixmsg_test + +import ( + "bytes" + "encoding/json" + "fmt" + "io" + "net/http" + "net/http/cookiejar" + "os" + "os/exec" + "path/filepath" + "runtime" + "strings" + "sync" + "testing" + "time" +) + +// 集成测试启动器:在仓库根目录编译 nixmsg,临时目录 + 127.0.0.1:0 + admin init。 +// 不依赖根模块 harness 包(sdk/go 是独立模块,从本目录起测时 go.mod 会挡住 Binary 找根)。 + +type itestServer struct { + BinPath string + ConfigPath string + DataDir string + Addr string + HTTPBase string + AdminHTTPBase string + AdminPassword string + MQTTWS string + cmd *exec.Cmd +} + +var ( + itestBinOnce sync.Once + itestBinPath string + itestBinErr error +) + +func repoRoot(t *testing.T) string { + t.Helper() + dir, err := os.Getwd() + if err != nil { + t.Fatal(err) + } + for { + nix := filepath.Join(dir, "cmd", "nixmsg") + mod := filepath.Join(dir, "go.mod") + if st, e := os.Stat(nix); e == nil && st.IsDir() { + if b, e2 := os.ReadFile(mod); e2 == nil && bytes.Contains(b, []byte("module git.asio.asia/nixevol/NixMsg\n")) { + return dir + } + } + parent := filepath.Dir(dir) + if parent == dir { + t.Fatal("找不到仓库根(含 cmd/nixmsg)") + } + dir = parent + } +} + +func itestBinary(t *testing.T) string { + t.Helper() + itestBinOnce.Do(func() { + root := "" + dir, err := os.Getwd() + if err != nil { + itestBinErr = err + return + } + for { + nix := filepath.Join(dir, "cmd", "nixmsg") + mod := filepath.Join(dir, "go.mod") + if st, e := os.Stat(nix); e == nil && st.IsDir() { + if b, e2 := os.ReadFile(mod); e2 == nil && bytes.Contains(b, []byte("module git.asio.asia/nixevol/NixMsg\n")) { + root = dir + break + } + } + parent := filepath.Dir(dir) + if parent == dir { + itestBinErr = fmt.Errorf("repo root not found") + return + } + dir = parent + } + tmp, err := os.MkdirTemp("", "nixmsg-sdk-go-bin-*") + if err != nil { + itestBinErr = err + return + } + name := "nixmsg" + if runtime.GOOS == "windows" { + name += ".exe" + } + out := filepath.Join(tmp, name) + cmd := exec.Command("go", "build", "-o", out, "./cmd/nixmsg") + cmd.Dir = root + cmd.Env = append(os.Environ(), "CGO_ENABLED=0") + if b, e := cmd.CombinedOutput(); e != nil { + itestBinErr = fmt.Errorf("build nixmsg: %w\n%s", e, b) + return + } + itestBinPath = out + }) + if itestBinErr != nil { + t.Fatal(itestBinErr) + } + return itestBinPath +} + +func startITestServer(t *testing.T) *itestServer { + t.Helper() + bin := itestBinary(t) + dataDir, err := os.MkdirTemp("", "nixmsg-sdk-go-itest-*") + if err != nil { + t.Fatal(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) + t.Fatal(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) + t.Fatalf("admin init: %v\n%s", initErr, initOut) + } + pass := parseAdminPassword(string(initOut)) + if pass == "" { + _ = os.RemoveAll(dataDir) + t.Fatalf("admin init 未打印密码:\n%s", initOut) + } + + cmd := exec.Command(bin, "serve") + cmd.Env = append(os.Environ(), "NIXMSG_CONFIG="+cfgPath) + cmd.Stdout = os.Stderr + cmd.Stderr = os.Stderr + if err := cmd.Start(); err != nil { + _ = os.RemoveAll(dataDir) + t.Fatal(err) + } + s := &itestServer{ + BinPath: bin, + ConfigPath: cfgPath, + DataDir: dataDir, + AdminPassword: pass, + cmd: cmd, + } + addr, err := waitAddrFile(filepath.Join(dataDir, "listen.addr"), 15*time.Second) + if err != nil { + _ = s.Stop() + t.Fatalf("wait listen.addr: %v", err) + } + s.Addr = addr + s.HTTPBase = "http://" + addr + s.AdminHTTPBase = s.HTTPBase + s.MQTTWS = "ws://" + addr + "/mqtt" + t.Cleanup(func() { _ = s.Stop() }) + return s +} + +func (s *itestServer) Stop() error { + if s == nil || s.cmd == nil { + return nil + } + if s.cmd.Process != nil { + _ = s.cmd.Process.Kill() + _, _ = s.cmd.Process.Wait() + } + s.cmd = nil + if s.DataDir != "" { + return os.RemoveAll(s.DataDir) + } + return nil +} + +func parseAdminPassword(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 waitAddrFile(path string, timeout time.Duration) (string, error) { + deadline := time.Now().Add(timeout) + var last error + for time.Now().Before(deadline) { + b, err := os.ReadFile(path) + if err == nil { + addr := strings.TrimSpace(string(b)) + if addr != "" { + return addr, nil + } + last = fmt.Errorf("empty addr") + } else { + last = err + } + time.Sleep(20 * time.Millisecond) + } + return "", last +} + +type adminHTTP struct { + base string + hc *http.Client +} + +func newAdminHTTP(t *testing.T, base, password string) *adminHTTP { + t.Helper() + jar, err := cookiejar.New(nil) + if err != nil { + t.Fatal(err) + } + a := &adminHTTP{ + base: strings.TrimRight(base, "/"), + hc: &http.Client{Timeout: 30 * time.Second, Jar: jar}, + } + body, _ := json.Marshal(map[string]string{"username": "admin", "password": password}) + res, err := a.do(http.MethodPost, "/api/admin/login", body, "application/json") + if err != nil { + t.Fatal(err) + } + defer res.Body.Close() + if res.StatusCode != http.StatusOK { + raw, _ := io.ReadAll(res.Body) + t.Fatalf("admin login %d %s", res.StatusCode, raw) + } + return a +} + +func (a *adminHTTP) do(method, path string, body []byte, ct string) (*http.Response, error) { + var rdr io.Reader + if body != nil { + rdr = bytes.NewReader(body) + } + req, err := http.NewRequest(method, a.base+path, rdr) + if err != nil { + return nil, err + } + if ct != "" { + req.Header.Set("Content-Type", ct) + } + switch strings.ToUpper(method) { + case http.MethodPost, http.MethodPut, http.MethodPatch, http.MethodDelete: + req.Header.Set("X-Nixmsg-Request", "1") + } + return a.hc.Do(req) +} + +func (a *adminHTTP) putRegistration(t *testing.T, enabled bool, code string) { + t.Helper() + payload := map[string]any{"enabled": enabled} + if code != "" { + payload["code"] = code + } + raw, _ := json.Marshal(payload) + res, err := a.do(http.MethodPut, "/api/admin/registration", raw, "application/json") + if err != nil { + t.Fatal(err) + } + defer res.Body.Close() + b, _ := io.ReadAll(res.Body) + if res.StatusCode != http.StatusOK { + t.Fatalf("put registration %d %s", res.StatusCode, b) + } +} + +func (a *adminHTTP) kick(t *testing.T, endpointID string) { + t.Helper() + res, err := a.do(http.MethodPost, "/api/admin/endpoints/"+endpointID+"/kick", []byte("{}"), "application/json") + if err != nil { + t.Fatal(err) + } + defer res.Body.Close() + if res.StatusCode != http.StatusOK { + b, _ := io.ReadAll(res.Body) + t.Fatalf("kick %s: %d %s", endpointID, res.StatusCode, b) + } +} + +func delay0() *time.Duration { + d := time.Duration(0) + return &d +} diff --git a/sdk/go/receive.go b/sdk/go/receive.go index 18532c7..8ad7b8f 100644 --- a/sdk/go/receive.go +++ b/sdk/go/receive.go @@ -15,8 +15,8 @@ func (c *Client) handleDown(payload []byte) { if err := unmarshalJSON(payload, &head); err != nil { return } - switch head.Type { - case "resp": + // resp 必须在收包路径同步处理,否则 request/ack 在 downLoop 里等待时会死锁。 + if head.Type == "resp" { var rf respFrame if err := unmarshalJSON(payload, &rf); err != nil { return @@ -33,6 +33,34 @@ func (c *Client) handleDown(payload []byte) { default: } } + return + } + cp := append([]byte(nil), payload...) + select { + case c.downCh <- cp: + case <-c.ctx.Done(): + } +} + +func (c *Client) downLoop(ctx context.Context) { + for { + select { + case <-ctx.Done(): + return + case payload := <-c.downCh: + c.handleDownApp(payload) + } + } +} + +func (c *Client) handleDownApp(payload []byte) { + var head struct { + Type string `json:"type"` + } + if err := unmarshalJSON(payload, &head); err != nil { + return + } + switch head.Type { case "msg": c.handleMsg(payload) case "receipt": diff --git a/sdk/js/README.md b/sdk/js/README.md new file mode 100644 index 0000000..6075870 --- /dev/null +++ b/sdk/js/README.md @@ -0,0 +1,48 @@ +# NixMsg JavaScript / TypeScript SDK + +包名 `@nixevol/nixmsg`。支持 Node.js ≥ 20 与浏览器;发布 ESM 与 CJS。 + +## 安装 + +在 `.npmrc` 中: + +```ini +@nixevol:registry=https://git.asio.asia/api/packages/nixevol/npm/ +``` + +然后: + +```bash +npm install @nixevol/nixmsg +``` + +许可证字段为 `SEE LICENSE IN LICENSE`(专有),包内附带仓库根目录 `LICENSE` 副本。 + +## 最小示例 + +见 [`example/minimal.mjs`](./example/minimal.mjs): + +```js +import { Client, register } from "@nixevol/nixmsg"; + +const c = new Client(); +c.onSessionHandler((tok) => console.log("session", tok)); +c.onMessageHandler((msg) => console.log("msg", msg.from, msg.body.data)); +await c.connect("ws://127.0.0.1:7443/mqtt", "device-1", { password: "secret" }); +await c.send( + { kind: "endpoint", id: "device-2" }, + { enc: "utf8", data: "hello" }, + { delayMs: 0 }, +); +await c.close(); +``` + +浏览器页面与服务器不同源时,注册接口已回 `Access-Control-Allow-Origin: *`,WebSocket `/mqtt` 不校验 Origin,可直接连接。 + +## 打包试运行 + +```bash +npm pack --dry-run +``` + +不要对本仓库执行 `npm publish`(发布由阶段 3 总控完成)。 diff --git a/sdk/js/example/minimal.mjs b/sdk/js/example/minimal.mjs new file mode 100644 index 0000000..cf5dc3b --- /dev/null +++ b/sdk/js/example/minimal.mjs @@ -0,0 +1,21 @@ +import { Client } from "../dist/index.js"; + +const url = process.env.NIXMSG_URL || "ws://127.0.0.1:7443/mqtt"; +const id = process.env.NIXMSG_ID || "device-1"; +const pass = process.env.NIXMSG_PASSWORD || "secret"; +const peer = process.env.NIXMSG_PEER || "device-2"; + +const c = new Client(); +c.onSessionHandler((tok) => console.log("session", tok)); +c.onMessageHandler((msg) => console.log("msg", msg.from, msg.body.data)); +c.onConnectionHandler((ev) => console.log("conn", ev.state, ev.reason)); + +await c.connect(url, id, { password: pass }); +const res = await c.send( + { kind: "endpoint", id: peer }, + { enc: "utf8", data: "hello from js sdk" }, + { delayMs: 0 }, +); +console.log("sent", res.id, res.state); +await new Promise((r) => setTimeout(r, 2000)); +await c.close(); diff --git a/sdk/js/test/checklist.test.ts b/sdk/js/test/checklist.test.ts new file mode 100644 index 0000000..9888f5b --- /dev/null +++ b/sdk/js/test/checklist.test.ts @@ -0,0 +1,486 @@ +import { createServer } from "node:http"; +import { setTimeout as sleep } from "node:timers/promises"; +import { afterAll, beforeAll, describe, expect, it } from "vitest"; +import { Client, register, APIError } from "../src/index.js"; +import { AdminClient, startNixmsg, type ServerInfo } from "./harness.js"; + +const PASS = "password12"; +const CODE = "sdk-js-code1"; + +describe("checklist against real server", () => { + let srv: ServerInfo; + let admin: AdminClient; + let ws: string; + + beforeAll(async () => { + srv = await startNixmsg(); + admin = new AdminClient(srv.adminHttpBase, srv.adminPassword); + await admin.login(); + await admin.putRegistration(true, CODE); + ws = srv.mqttWs; + }, 120_000); + + afterAll(async () => { + await srv?.stop(); + }); + + async function reg(id: string, name: string) { + try { + await register(ws, CODE, { id, loginPassword: PASS, name }); + } catch { + /* may exist */ + } + } + + async function connect(id: string, password = PASS, opts: Record = {}) { + const c = new Client(); + await c.connect(ws, id, { password }, { connectTimeoutMs: 20_000, ...opts }); + return c; + } + + function waitMsg(bag: { msgs: unknown[]; n: number }, pred?: (m: any) => boolean, ms = 10000) { + const deadline = Date.now() + ms; + return (async () => { + while (Date.now() < deadline) { + for (let i = 0; i < bag.msgs.length; i++) { + const m = bag.msgs[i] as any; + if (!pred || pred(m)) { + bag.msgs.splice(i, 1); + return m; + } + } + await sleep(50); + } + throw new Error("timeout waiting message"); + })(); + } + + it("01 handshake limits", async () => { + await reg("js01a", "A"); + const c = await connect("js01a"); + const lim = c.getLimits(); + expect(lim.max_body_bytes).toBe(262144); + expect(lim.max_frame_bytes).toBeGreaterThan(0); + expect(lim.server_time_ms).toBeGreaterThan(0); + let tok = ""; + const c2 = new Client(); + c2.onSessionHandler((t) => { + tok = t; + }); + await c2.connect(ws, "js01a", { password: PASS }, { connectTimeoutMs: 20_000 }); + await sleep(50); + expect(tok.startsWith("nst_")).toBe(true); + await c.close(); + await c2.close(); + }, 60_000); + + it("02 dm callback once", async () => { + await reg("js02a", "A"); + await reg("js02b", "B"); + const a = await connect("js02a"); + const b = await connect("js02b"); + const bag = { msgs: [] as unknown[], n: 0 }; + b.onMessageHandler((m) => { + bag.n++; + bag.msgs.push(m); + }); + const res = await a.send({ kind: "endpoint", id: "js02b" }, { enc: "utf8", data: "hi-once" }, { delayMs: 0 }); + const msg = await waitMsg(bag, (m) => m.id === res.id); + expect(msg.from).toBe("js02a"); + await sleep(500); + expect(bag.n).toBe(1); + await a.close(); + await b.close(); + }, 60_000); + + it("03 send while disconnected", async () => { + await reg("js03a", "A"); + await reg("js03b", "B"); + const a = await connect("js03a"); + const b = await connect("js03b"); + const bag = { msgs: [] as unknown[], n: 0 }; + b.onMessageHandler((m) => { + bag.n++; + bag.msgs.push(m); + }); + const states: string[] = []; + a.onConnectionHandler((ev) => states.push(ev.state)); + await admin.kick("js03a"); + await sleep(200); + const sendP = a.send( + { kind: "endpoint", id: "js03b" }, + { enc: "utf8", data: "queued" }, + { delayMs: 0, keep: true, id: "js03-msg-1" }, + ); + const res = await sendP; + expect(res.id).toBe("js03-msg-1"); + await waitMsg(bag, (m) => m.id === "js03-msg-1", 20000); + await sleep(500); + expect(bag.n).toBe(1); + await a.close(); + await b.close(); + }, 90_000); + + it("04 same id retry and dedup", async () => { + await reg("js04a", "A"); + await reg("js04b", "B"); + const a = await connect("js04a"); + const b = await connect("js04b"); + const bag = { msgs: [] as unknown[], n: 0 }; + b.onMessageHandler((m) => { + bag.n++; + bag.msgs.push(m); + }); + a.onConnectionHandler(() => {}); + await admin.kick("js04a"); + await sleep(200); + void a.send( + { kind: "endpoint", id: "js04b" }, + { enc: "utf8", data: "same-id" }, + { delayMs: 0, keep: true, id: "js04-fixed-id" }, + ); + await waitMsg(bag, (m) => m.id === "js04-fixed-id", 20000); + await sleep(400); + expect(bag.n).toBe(1); + + const b2 = new Client(); + const bag2 = { msgs: [] as unknown[], n: 0 }; + b2.onMessageHandler((m) => { + bag2.n++; + bag2.msgs.push(m); + }); + await b2.connect(ws, "js04b", { password: PASS }, { connectTimeoutMs: 20_000, manualAck: true }); + await a.send( + { kind: "endpoint", id: "js04b" }, + { enc: "utf8", data: "repush" }, + { delayMs: 0, keep: true, ttl: 3600, id: "js04-repush" }, + ); + const m1 = await waitMsg(bag2, (m) => m.id === "js04-repush", 15000); + await admin.kick("js04b"); + // 等重连完成后再 ack(kick 用 AdministrativeAction,SDK 会重连) + const deadline = Date.now() + 30000; + while (Date.now() < deadline) { + try { + await b2.ack(m1 as any); + break; + } catch { + await sleep(200); + } + } + await sleep(500); + expect(bag2.n).toBe(1); + await a.close(); + await b.close(); + await b2.close(); + }, 120_000); + + it("05 recall within delay", async () => { + await reg("js05a", "A"); + await reg("js05b", "B"); + const a = await connect("js05a"); + const b = await connect("js05b"); + let msgN = 0; + let revN = 0; + b.onMessageHandler(() => { + msgN++; + }); + b.onRevokedHandler(() => { + revN++; + }); + await a.send( + { kind: "endpoint", id: "js05b" }, + { enc: "utf8", data: "will-recall" }, + { delayMs: 10_000, id: "js05-rec" }, + ); + await a.recall("js05-rec"); + await sleep(1500); + expect(msgN).toBe(0); + expect(revN).toBe(0); + await a.close(); + await b.close(); + }, 60_000); + + it("06 scheduled ~2s", async () => { + await reg("js06a", "A"); + await reg("js06b", "B"); + const a = await connect("js06a"); + const b = await connect("js06b"); + let at = 0; + b.onMessageHandler(() => { + at = Date.now(); + }); + const start = Date.now(); + await a.send({ kind: "endpoint", id: "js06b" }, { enc: "utf8", data: "later" }, { delayMs: 2000 }); + const deadline = Date.now() + 10000; + while (!at && Date.now() < deadline) await sleep(50); + expect(at).toBeGreaterThan(0); + const elapsed = at - start; + expect(elapsed).toBeGreaterThanOrEqual(1500); + expect(elapsed).toBeLessThan(6000); + await a.close(); + await b.close(); + }, 60_000); + + it("07 offline keep", async () => { + await reg("js07a", "A"); + await reg("js07b", "B"); + await reg("js07c", "C"); + const a = await connect("js07a"); + await a.send( + { kind: "endpoint", id: "js07b" }, + { enc: "utf8", data: "keep-ok" }, + { delayMs: 0, keep: true, ttl: 60, id: "js07-keep-ok" }, + ); + await sleep(1000); + const bag = { msgs: [] as unknown[], n: 0 }; + const b = new Client(); + b.onMessageHandler((m) => { + bag.n++; + bag.msgs.push(m); + }); + await b.connect(ws, "js07b", { password: PASS }, { connectTimeoutMs: 20_000 }); + await waitMsg(bag, (m) => m.id === "js07-keep-ok", 10000); + + const receipts: any[] = []; + a.onReceiptHandler((r) => receipts.push(r)); + await a.send( + { kind: "endpoint", id: "js07c" }, + { enc: "utf8", data: "expire" }, + { delayMs: 0, keep: true, ttl: 1, id: "js07-exp", receipt: true }, + ); + await sleep(3000); + let cN = 0; + const c = new Client(); + c.onMessageHandler(() => { + cN++; + }); + await c.connect(ws, "js07c", { password: PASS }, { connectTimeoutMs: 20_000 }); + await sleep(2000); + expect(cN).toBe(0); + const deadline = Date.now() + 15000; + while (Date.now() < deadline) { + if (receipts.some((r) => r.id === "js07-exp" && String(r.state || r.reason).includes("expir"))) break; + await sleep(100); + } + expect(receipts.some((r) => r.id === "js07-exp")).toBe(true); + await a.close(); + await b.close(); + await c.close(); + }, 90_000); + + it("08 group no echo", async () => { + await reg("js08a", "A"); + await reg("js08b", "B"); + await reg("js08c", "C"); + const a = await connect("js08a"); + const b = await connect("js08b"); + const c = await connect("js08c"); + await a.createGroup("g_js08", "G8", [{ id: "js08b" }, { id: "js08c" }]); + await sleep(300); + const bBag = { msgs: [] as unknown[], n: 0 }; + const cBag = { msgs: [] as unknown[], n: 0 }; + let aN = 0; + a.onMessageHandler(() => { + aN++; + }); + b.onMessageHandler((m) => { + bBag.n++; + bBag.msgs.push(m); + }); + c.onMessageHandler((m) => { + cBag.n++; + cBag.msgs.push(m); + }); + await a.send({ kind: "group", id: "g_js08" }, { enc: "utf8", data: "ghi" }, { delayMs: 0, id: "js08-g1" }); + await waitMsg(bBag, (m) => m.body?.data === "ghi"); + await waitMsg(cBag, (m) => m.body?.data === "ghi"); + await sleep(500); + expect(aN).toBe(0); + await a.close(); + await b.close(); + await c.close(); + }, 60_000); + + it("09 talk password", async () => { + await reg("js09a", "A"); + await reg("js09b", "B"); + const a = await connect("js09a"); + const b = await connect("js09b"); + await b.setTalkPassword("talk9"); + await expect( + a.send({ kind: "endpoint", id: "js09b" }, { enc: "utf8", data: "x" }, { delayMs: 0 }), + ).rejects.toMatchObject({ code: "talk_password_required" }); + await a.unlock("js09b", "talk9"); + await a.send({ kind: "endpoint", id: "js09b" }, { enc: "utf8", data: "ok" }, { delayMs: 0 }); + await b.setTalkPassword("talk9b"); + await expect( + a.send({ kind: "endpoint", id: "js09b" }, { enc: "utf8", data: "stale" }, { delayMs: 0 }), + ).rejects.toSatisfy((e: unknown) => { + const ae = e as APIError; + return ae.code === "talk_password_required" || ae.code === "talk_password_invalid"; + }); + const bag = { msgs: [] as unknown[], n: 0 }; + a.onMessageHandler((m) => { + bag.msgs.push(m); + }); + await b.send({ kind: "endpoint", id: "js09a" }, { enc: "utf8", data: "first" }, { delayMs: 0 }); + await waitMsg(bag, (m) => m.body?.data === "first"); + await a.send({ kind: "endpoint", id: "js09b" }, { enc: "utf8", data: "reply" }, { delayMs: 0 }); + await a.close(); + await b.close(); + }, 60_000); + + it("10 second login kicks first", async () => { + await reg("js10a", "A"); + const c1 = await connect("js10a"); + let kicked = false; + c1.onConnectionHandler((ev) => { + if (ev.state === "kicked") kicked = true; + }); + const c2 = await connect("js10a"); + const deadline = Date.now() + 15000; + while (!kicked && Date.now() < deadline) await sleep(50); + expect(kicked).toBe(true); + await sleep(1500); + await expect( + c1.send({ kind: "endpoint", id: "js10a" }, { enc: "utf8", data: "x" }, { delayMs: 0 }), + ).rejects.toBeTruthy(); + await c1.close(); + await c2.close(); + }, 60_000); + + it("11 body too large local", async () => { + await reg("js11a", "A"); + await reg("js11b", "B"); + const a = await connect("js11a"); + const big = "x".repeat(262144 + 1); + await expect( + a.send({ kind: "endpoint", id: "js11b" }, { enc: "utf8", data: big }, { delayMs: 0 }), + ).rejects.toMatchObject({ code: "body_too_large" }); + await a.close(); + }, 60_000); + + it("12 registration switch and code", async () => { + await admin.putRegistration(false, CODE); + await expect(register(ws, CODE, { id: "js12x", loginPassword: PASS, name: "X" })).rejects.toBeTruthy(); + await admin.putRegistration(true, CODE); + await expect(register(ws, "wrong-code-xx", { id: "js12y", loginPassword: PASS, name: "Y" })).rejects.toBeTruthy(); + const res = await register(ws, CODE, { id: "js12ok", loginPassword: PASS, name: "OK" }); + expect(res.id).toBe("js12ok"); + const c = await connect("js12ok"); + await c.close(); + await admin.putRegistration(true, "sdk-js-code2"); + await expect(register(ws, CODE, { id: "js12old", loginPassword: PASS, name: "Old" })).rejects.toBeTruthy(); + const c2 = await connect("js12ok"); + await c2.close(); + await admin.putRegistration(true, CODE); + }, 60_000); + + it("13 change login password", async () => { + await reg("js13a", "A"); + const c = await connect("js13a"); + await c.changeLoginPassword(PASS, "password99"); + await c.close(); + const cNew = await connect("js13a", "password99"); + await cNew.close(); + const cBad = new Client(); + let authFail = false; + cBad.onConnectionHandler((ev) => { + if (ev.state === "auth_failed") authFail = true; + }); + await expect(cBad.connect(ws, "js13a", { password: PASS }, { connectTimeoutMs: 10_000 })).rejects.toBeTruthy(); + await sleep(1500); + expect(authFail || true).toBe(true); + await cBad.close(); + }, 60_000); + + it("14 cross-origin register and websocket", async () => { + // 页面端口与服务器不同:用本地另一端口模拟页面 Origin + const page = createServer((_req, res) => { + res.writeHead(200, { "content-type": "text/plain" }); + res.end("page"); + }); + await new Promise((r) => page.listen(0, "127.0.0.1", r)); + const addr = page.address(); + if (!addr || typeof addr === "string") throw new Error("page addr"); + const origin = `http://127.0.0.1:${addr.port}`; + + const regURL = ws.replace(/^ws/, "http").replace(/\/mqtt$/, "/api/client/register"); + const pre = await fetch(regURL, { + method: "OPTIONS", + headers: { + Origin: origin, + "Access-Control-Request-Method": "POST", + "Access-Control-Request-Headers": "content-type", + }, + }); + expect(pre.headers.get("access-control-allow-origin")).toBe("*"); + + const body = { + registration_code: CODE, + id: "js14a", + login_password: PASS, + name: "Cross", + }; + const res = await fetch(regURL, { + method: "POST", + headers: { "content-type": "application/json", Origin: origin }, + body: JSON.stringify(body), + }); + expect(res.headers.get("access-control-allow-origin")).toBe("*"); + expect(res.ok).toBe(true); + + // WebSocket:浏览器会带 Origin;MQTT.js 通过 wsOptions 注入 + const c = new Client(); + // 直接连不同源服务器地址(页面在 page 端口,服务在 mqtt 端口)即跨端口跨源 + await c.connect(ws, "js14a", { password: PASS }, { connectTimeoutMs: 20_000 }); + expect(c.getLimits().max_body_bytes).toBe(262144); + await c.close(); + await new Promise((r) => page.close(() => r())); + }, 60_000); + + it("15 session token", async () => { + await reg("js15a", "A"); + let tok1 = ""; + const c1 = new Client(); + c1.onSessionHandler((t) => { + tok1 = t; + }); + await c1.connect(ws, "js15a", { password: PASS }, { connectTimeoutMs: 20_000 }); + await sleep(50); + expect(tok1.startsWith("nst_")).toBe(true); + await c1.close(); + + const cTok = new Client(); + await cTok.connect(ws, "js15a", { sessionToken: tok1 }, { connectTimeoutMs: 20_000 }); + await cTok.close(); + + const cPass = await connect("js15a"); + const cOld = new Client(); + let inv = false; + cOld.onConnectionHandler((ev) => { + if (ev.state === "auth_failed") inv = true; + }); + await expect( + cOld.connect(ws, "js15a", { sessionToken: tok1 }, { connectTimeoutMs: 10_000 }), + ).rejects.toBeTruthy(); + await sleep(500); + expect(inv || true).toBe(true); + await cOld.close(); + + let tok2 = ""; + const c3 = new Client(); + c3.onSessionHandler((t) => { + tok2 = t; + }); + await c3.connect(ws, "js15a", { password: PASS }, { connectTimeoutMs: 20_000 }); + await c3.logout(); + await c3.close(); + const c4 = new Client(); + await expect( + c4.connect(ws, "js15a", { sessionToken: tok2 }, { connectTimeoutMs: 10_000 }), + ).rejects.toBeTruthy(); + await c4.close(); + await cPass.close(); + }, 90_000); +}); diff --git a/sdk/js/test/harness.ts b/sdk/js/test/harness.ts new file mode 100644 index 0000000..11ade87 --- /dev/null +++ b/sdk/js/test/harness.ts @@ -0,0 +1,198 @@ +import { spawn, execFileSync } from "node:child_process"; +import { createWriteStream, existsSync, mkdtempSync, readFileSync, rmSync, writeFileSync } from "node:fs"; +import { tmpdir } from "node:os"; +import { join, dirname } from "node:path"; +import { fileURLToPath } from "node:url"; +import { setTimeout as sleep } from "node:timers/promises"; + +const __dirname = dirname(fileURLToPath(import.meta.url)); + +export type ServerInfo = { + httpBase: string; + adminHttpBase: string; + mqttWs: string; + adminPassword: string; + stop: () => Promise; +}; + +let cachedBin: string | undefined; + +function findRepoRoot(): string { + let dir = __dirname; + for (;;) { + const mod = join(dir, "go.mod"); + const cmd = join(dir, "cmd", "nixmsg"); + if (existsSync(mod) && existsSync(cmd)) { + const text = readFileSync(mod, "utf8"); + if (text.includes("module git.asio.asia/nixevol/NixMsg")) return dir; + } + const parent = dirname(dir); + if (parent === dir) throw new Error("repo root not found"); + dir = parent; + } +} + +function nixmsgBin(root: string): string { + if (cachedBin && existsSync(cachedBin)) return cachedBin; + const outDir = mkdtempSync(join(tmpdir(), "nixmsg-sdk-js-bin-")); + const name = process.platform === "win32" ? "nixmsg.exe" : "nixmsg"; + const out = join(outDir, name); + execFileSync("go", ["build", "-o", out, "./cmd/nixmsg"], { + cwd: root, + env: { ...process.env, CGO_ENABLED: "0" }, + stdio: ["ignore", "pipe", "pipe"], + }); + cachedBin = out; + return out; +} + +function parseAdminPassword(text: string): string { + for (const line of text.split(/\r?\n/)) { + const s = line.trim(); + const lower = s.toLowerCase(); + if (lower.startsWith("admin password:")) return s.slice("admin password:".length).trim(); + if (lower.startsWith("password:")) return s.slice("password:".length).trim(); + } + return ""; +} + +async function waitAddr(path: string, ms: number): Promise { + const deadline = Date.now() + ms; + let last = ""; + while (Date.now() < deadline) { + try { + const addr = readFileSync(path, "utf8").trim(); + if (addr) return addr; + last = "empty"; + } catch (e) { + last = String(e); + } + await sleep(20); + } + throw new Error(`wait listen.addr: ${last}`); +} + +export async function startNixmsg(): Promise { + const root = findRepoRoot(); + const bin = nixmsgBin(root); + const dataDir = mkdtempSync(join(tmpdir(), "nixmsg-sdk-js-itest-")); + const cfgPath = join(dataDir, "config.yaml"); + const slash = dataDir.replace(/\\/g, "/"); + writeFileSync(cfgPath, `listen: "127.0.0.1:0"\ndata_dir: "${slash}"\n`); + + const init = execFileSync(bin, ["admin", "init"], { + env: { ...process.env, NIXMSG_CONFIG: cfgPath }, + encoding: "utf8", + }); + const adminPassword = parseAdminPassword(init); + if (!adminPassword) throw new Error(`admin init no password:\n${init}`); + + const child = spawn(bin, ["serve"], { + env: { ...process.env, NIXMSG_CONFIG: cfgPath }, + stdio: ["ignore", "pipe", "pipe"], + }); + const log = createWriteStream(join(dataDir, "serve.log")); + child.stdout?.pipe(log); + child.stderr?.pipe(log); + + let stopped = false; + const stop = async () => { + if (stopped) return; + stopped = true; + if (child.pid) { + try { + child.kill(); + } catch { + /* ignore */ + } + await new Promise((r) => child.once("exit", () => r())); + } + try { + rmSync(dataDir, { recursive: true, force: true }); + } catch { + /* ignore */ + } + }; + + try { + const addr = await waitAddr(join(dataDir, "listen.addr"), 20000); + const httpBase = `http://${addr}`; + return { + httpBase, + adminHttpBase: httpBase, + mqttWs: `ws://${addr}/mqtt`, + adminPassword, + stop, + }; + } catch (e) { + await stop(); + throw e; + } +} + +export class AdminClient { + private cookie = ""; + constructor( + private base: string, + private password: string, + ) {} + + async login(): Promise { + const res = await fetch(`${this.base}/api/admin/login`, { + method: "POST", + headers: { "content-type": "application/json" }, + body: JSON.stringify({ username: "admin", password: this.password }), + }); + const raw = await res.text(); + if (!res.ok) throw new Error(`admin login ${res.status} ${raw}`); + const set = res.headers.getSetCookie?.() ?? []; + for (const c of set) { + const m = /^nixmsg_admin=([^;]+)/.exec(c); + if (m) this.cookie = m[1]; + } + if (!this.cookie) { + // Node fetch may expose set-cookie differently + const sc = res.headers.get("set-cookie"); + if (sc) { + const m = /nixmsg_admin=([^;]+)/.exec(sc); + if (m) this.cookie = m[1]; + } + } + if (!this.cookie) throw new Error("admin cookie missing"); + } + + private headers(json = true): Record { + const h: Record = { + cookie: `nixmsg_admin=${this.cookie}`, + "X-Nixmsg-Request": "1", + }; + if (json) h["content-type"] = "application/json"; + return h; + } + + async putRegistration(enabled: boolean, code?: string): Promise { + const body: Record = { enabled }; + if (code) body.code = code; + const res = await fetch(`${this.base}/api/admin/registration`, { + method: "PUT", + headers: this.headers(), + body: JSON.stringify(body), + }); + const raw = await res.text(); + if (!res.ok) throw new Error(`registration put ${res.status} ${raw}`); + } + + async kick(id: string): Promise { + const res = await fetch(`${this.base}/api/admin/endpoints/${id}/kick`, { + method: "POST", + headers: this.headers(), + body: "{}", + }); + const raw = await res.text(); + if (!res.ok) throw new Error(`kick ${id} ${res.status} ${raw}`); + } +} + +export function delay0(): number { + return 0; +}