import { describe, expect, it } from "vitest"; import { Client, APIError, buildCleanConnectFlags, register } from "../src/index.js"; import { FakeTransport } from "../src/fake.js"; import { createServer } from "node:http"; async function connectFake(fake: FakeTransport): Promise { const c = new Client(); await c.connect("ws://example.test/mqtt", "ep1", { password: "secret" }, { transport: fake }); return c; } describe("nixmsg sdk", () => { it("Clean Start every connect", async () => { const fake = new FakeTransport(); const c = await connectFake(fake); await fake.simulateReconnect(); await fake.simulateReconnect(); const cs = fake.getConnects(); expect(cs.length).toBeGreaterThanOrEqual(3); for (const x of cs) { expect(x.cleanStart).toBe(true); expect(x.sessionExpiry).toBe(0); } await c.close(); }); it("session token callback", async () => { const fake = new FakeTransport(); fake.helloToken = "nst_abc"; const c = new Client(); let got = ""; c.onSessionHandler((t) => { got = t; }); await c.connect("ws://example.test/mqtt", "ep1", { password: "p" }, { transport: fake }); // 等待串行回调 await new Promise((r) => setTimeout(r, 20)); expect(got).toBe("nst_abc"); await c.close(); }); it("dedup then re-ack", async () => { const fake = new FakeTransport(); const c = await connectFake(fake); let calls = 0; c.onMessageHandler(() => { calls++; }); const replyAcks = setInterval(() => { for (const fr of fake.findUp("ack")) { fake.replyOK(String(fr.rid), { result: "accepted" }); } }, 5); const msg = JSON.stringify({ v: 1, type: "msg", id: "m1", from: "a", to: { kind: "endpoint", id: "ep1" }, body: { enc: "utf8", data: "hi" }, send_at_ms: 1, }); fake.injectDown(msg); fake.injectDown(msg); await new Promise((r) => setTimeout(r, 80)); fake.injectDown(msg); await new Promise((r) => setTimeout(r, 80)); clearInterval(replyAcks); expect(calls).toBe(1); expect(fake.findUp("ack").length).toBeGreaterThanOrEqual(2); await c.close(); }); it("body too large locally", async () => { const fake = new FakeTransport(); fake.maxBodyBytes = 16; const c = await connectFake(fake); await expect( c.send( { kind: "endpoint", id: "b" }, { enc: "utf8", data: "x".repeat(64) }, ), ).rejects.toMatchObject({ code: "body_too_large" }); await c.close(); }); it("resend keeps id and send_at_ms", async () => { const fake = new FakeTransport(); const c = await connectFake(fake); const at = new Date(1_700_000_000_000); let firstId = ""; let firstSendAt: unknown; let replied = false; const timer = setInterval(() => { const sends = fake.findUp("send"); if (!sends.length) return; if (!replied) { replied = true; firstId = String(sends[0].id); firstSendAt = sends[0].send_at_ms; fake.replyErr(String(sends[0].rid), "rate_limited", "slow"); return; } if (sends.length >= 2) { expect(sends[1].id).toBe(firstId); expect(sends[1].send_at_ms).toBe(firstSendAt); fake.replyOK(String(sends[1].rid), { id: firstId, send_at_ms: firstSendAt, state: "scheduled", }); clearInterval(timer); } }, 20); const res = await c.send( { kind: "endpoint", id: "b" }, { enc: "utf8", data: "hi" }, { sendAt: at }, ); expect(res.id).toBe(firstId); await c.close(); }, 10000); it("publishUp transport fail retries with new rid", async () => { const fake = new FakeTransport(); const c = await connectFake(fake); const at = new Date(1_700_000_000_000); const attempts: Array> = []; let failOnce = true; fake.publishUpImpl = async (s) => { const m = JSON.parse(s) as Record; if (m.type !== "send") return; attempts.push(m); if (failOnce) { failOnce = false; throw new Error("transient publish"); } }; const timer = setInterval(() => { const sends = fake.findUp("send"); if (sends.length < 1) return; const last = sends[sends.length - 1]!; fake.replyOK(String(last.rid), { id: last.id, send_at_ms: last.send_at_ms, state: "scheduled", }); }, 5); const res = await c.send( { kind: "endpoint", id: "b" }, { enc: "utf8", data: "hi" }, { sendAt: at }, ); clearInterval(timer); expect(attempts.length).toBeGreaterThanOrEqual(2); expect(String(attempts[1]!.rid)).not.toBe(String(attempts[0]!.rid)); expect(attempts[1]!.id).toBe(attempts[0]!.id); expect(attempts[1]!.send_at_ms).toBe(attempts[0]!.send_at_ms); expect((attempts[1]!.body as { data: string }).data).toBe("hi"); expect(res.id).toBe(String(attempts[0]!.id)); await c.close(); }, 5000); it("register HTTP from ws url", async () => { const srv = createServer((req, res) => { expect(req.url).toBe("/api/client/register"); res.setHeader("content-type", "application/json"); res.end(JSON.stringify({ ok: true, data: { id: "e_1", login_password: "gen" } })); }); await new Promise((r) => srv.listen(0, "127.0.0.1", r)); const addr = srv.address(); if (!addr || typeof addr === "string") throw new Error("addr"); const ws = `ws://127.0.0.1:${addr.port}/mqtt`; const res = await register(ws, "code", { name: "n" }); expect(res.id).toBe("e_1"); expect(res.loginPassword).toBe("gen"); await new Promise((r) => srv.close(() => r())); }); it("buildCleanConnectFlags", () => { expect(buildCleanConnectFlags()).toEqual({ cleanStart: true, sessionExpiry: 0 }); }); it("APIError shape", () => { const e = new APIError("body_too_large", "x"); expect(e.code).toBe("body_too_large"); }); });