h = onDisconnected;
if (h != null) {
- h.accept("network", false);
+ h.accept(reason, stop);
}
}
}
@@ -386,6 +407,30 @@ final class HiveMqTransport implements Transport {
}
}
+ private static Mqtt5ConnAckReasonCode extractConnAckReason(Throwable e) {
+ Throwable t = e;
+ while (t != null) {
+ if (t instanceof Mqtt5ConnAckException) {
+ return ((Mqtt5ConnAckException) t).getMqttMessage().getReasonCode();
+ }
+ t = t.getCause();
+ }
+ return null;
+ }
+
+ private static String exceptionText(Throwable e) {
+ StringBuilder sb = new StringBuilder();
+ Throwable t = e;
+ while (t != null) {
+ sb.append(t.getClass().getName()).append(' ');
+ if (t.getMessage() != null) {
+ sb.append(t.getMessage()).append(' ');
+ }
+ t = t.getCause();
+ }
+ return sb.toString();
+ }
+
private static boolean isStop(Mqtt5ConnAckReasonCode code) {
return code == Mqtt5ConnAckReasonCode.BAD_USER_NAME_OR_PASSWORD
|| code == Mqtt5ConnAckReasonCode.NOT_AUTHORIZED
diff --git a/sdk/java/src/main/java/asia/asio/nixmsg/examples/MinimalExample.java b/sdk/java/src/main/java/asia/asio/nixmsg/examples/MinimalExample.java
new file mode 100644
index 0000000..111604d
--- /dev/null
+++ b/sdk/java/src/main/java/asia/asio/nixmsg/examples/MinimalExample.java
@@ -0,0 +1,38 @@
+package asia.asio.nixmsg.examples;
+
+import asia.asio.nixmsg.Client;
+import asia.asio.nixmsg.Types.Body;
+import asia.asio.nixmsg.Types.SendOptions;
+import asia.asio.nixmsg.Types.Target;
+
+/**
+ * 最小示例:连接、发送、关闭。
+ *
+ * 运行(需本机已有 nixmsg,并准备好端号与密码):
+ * {@code java -cp ... asia.asio.nixmsg.examples.MinimalExample ws://127.0.0.1:PORT/mqtt device-1 password12 peer-id}
+ */
+public final class MinimalExample {
+ private MinimalExample() {}
+
+ public static void main(String[] args) throws Exception {
+ if (args.length < 4) {
+ System.err.println("用法: MinimalExample ");
+ System.exit(2);
+ }
+ String url = args[0];
+ String eid = args[1];
+ String password = args[2];
+ String peer = args[3];
+
+ Client c = new Client();
+ c.onSession(token -> System.out.println("session " + token.substring(0, Math.min(16, token.length())) + "..."));
+ c.onMessage(msg -> System.out.println("msg " + msg.from + " " + msg.id + " " + msg.body.data));
+ c.onConnection(ev -> System.out.println("conn " + ev.state + " " + ev.reason));
+ c.connectSync(url, eid, password, null, false);
+ SendOptions opt = new SendOptions();
+ opt.delayMs = 0L;
+ c.sendSync(new Target("endpoint", peer), new Body("hello from java"), opt);
+ Thread.sleep(2000L);
+ c.close();
+ }
+}
diff --git a/sdk/java/src/test/java/asia/asio/nixmsg/ChecklistTest.java b/sdk/java/src/test/java/asia/asio/nixmsg/ChecklistTest.java
new file mode 100644
index 0000000..a1ee220
--- /dev/null
+++ b/sdk/java/src/test/java/asia/asio/nixmsg/ChecklistTest.java
@@ -0,0 +1,531 @@
+package asia.asio.nixmsg;
+
+import asia.asio.nixmsg.Types.Body;
+import asia.asio.nixmsg.Types.ConnectionEvent;
+import asia.asio.nixmsg.Types.ConnectionState;
+import asia.asio.nixmsg.Types.IncomingMessage;
+import asia.asio.nixmsg.Types.Receipt;
+import asia.asio.nixmsg.Types.RegisterOptions;
+import asia.asio.nixmsg.Types.RevokedEvent;
+import asia.asio.nixmsg.Types.SendOptions;
+import asia.asio.nixmsg.Types.SendResult;
+import asia.asio.nixmsg.Types.Target;
+import org.junit.AfterClass;
+import org.junit.BeforeClass;
+import org.junit.Test;
+
+import java.util.ArrayList;
+import java.util.Collections;
+import java.util.List;
+import java.util.Map;
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicReference;
+
+import static org.junit.Assert.assertEquals;
+import static org.junit.Assert.assertFalse;
+import static org.junit.Assert.assertNotNull;
+import static org.junit.Assert.assertTrue;
+import static org.junit.Assert.fail;
+
+/** DEVELOPMENT 第 9 节接入清单(真实服务器;跳过仅 JS 跨域)。 */
+public class ChecklistTest {
+ private static TestHarness srv;
+ private static int seq;
+
+ @BeforeClass
+ public static void startServer() throws Exception {
+ srv = new TestHarness();
+ srv.start();
+ }
+
+ @AfterClass
+ public static void stopServer() {
+ if (srv != null) {
+ srv.stop();
+ }
+ }
+
+ private static synchronized String uid(String prefix) {
+ seq++;
+ return prefix + String.format("%04d", seq);
+ }
+
+ private static void register(String id) {
+ RegisterOptions opt = new RegisterOptions();
+ opt.id = id;
+ opt.loginPassword = "password12";
+ opt.name = id;
+ Client.registerSync(srv.wsUrl, srv.regCode, opt);
+ }
+
+ private static Client connect(String id) {
+ return connect(id, "password12", null);
+ }
+
+ private static Client connect(String id, String password, String token) {
+ Client c = new Client();
+ c.connectSync(srv.wsUrl, id, password, token, false);
+ assertEquals(ConnectionState.ONLINE, c.getState());
+ return c;
+ }
+
+ private static SendOptions immediate(String messageId) {
+ SendOptions o = new SendOptions();
+ o.delayMs = 0L;
+ o.messageId = messageId;
+ return o;
+ }
+
+ private static boolean waitUntil(Condition cond, long timeoutMs) throws InterruptedException {
+ long deadline = System.currentTimeMillis() + timeoutMs;
+ while (System.currentTimeMillis() < deadline) {
+ if (cond.ok()) {
+ return true;
+ }
+ Thread.sleep(50L);
+ }
+ return cond.ok();
+ }
+
+ private interface Condition {
+ boolean ok();
+ }
+
+ private static final class MsgBox {
+ private final List items = Collections.synchronizedList(new ArrayList());
+ private final CountDownLatch latch = new CountDownLatch(1);
+
+ void onMessage(IncomingMessage m) {
+ items.add(m);
+ latch.countDown();
+ }
+
+ List waitN(int n, long timeoutMs) throws InterruptedException {
+ long deadline = System.currentTimeMillis() + timeoutMs;
+ while (System.currentTimeMillis() < deadline) {
+ if (items.size() >= n) {
+ return new ArrayList(items);
+ }
+ Thread.sleep(50L);
+ }
+ return new ArrayList(items);
+ }
+ }
+
+ @Test
+ public void test01Handshake() {
+ String id = uid("hs");
+ register(id);
+ Client c = new Client();
+ final List tokens = new ArrayList();
+ c.onSession(tokens::add);
+ c.connectSync(srv.wsUrl, id, "password12", null, false);
+ assertEquals(ConnectionState.ONLINE, c.getState());
+ assertTrue(c.getLimits().serverTimeMs > 0);
+ assertTrue(c.getLimits().maxBodyBytes >= 256 * 1024);
+ assertFalse(tokens.isEmpty());
+ assertTrue(tokens.get(0).startsWith("nst_"));
+ c.close();
+ }
+
+ @Test
+ public void test02DmOnce() throws Exception {
+ String a = uid("a2");
+ String b = uid("b2");
+ register(a);
+ register(b);
+ Client ca = connect(a);
+ Client cb = connect(b);
+ MsgBox box = new MsgBox();
+ cb.onMessage(box::onMessage);
+ String mid = Uuid7.next();
+ ca.sendSync(new Target("endpoint", b), new Body("hello-once"), immediate(mid));
+ List got = box.waitN(1, 10_000);
+ assertEquals(1, got.size());
+ assertEquals(mid, got.get(0).id);
+ Thread.sleep(500);
+ assertEquals(1, box.items.size());
+ ca.close();
+ cb.close();
+ }
+
+ @Test
+ public void test03SendWhileDisconnected() throws Exception {
+ String a = uid("a3");
+ String b = uid("b3");
+ register(a);
+ register(b);
+ Client ca = connect(a);
+ Client cb = connect(b);
+ MsgBox box = new MsgBox();
+ cb.onMessage(box::onMessage);
+ String mid = Uuid7.next();
+ ca.dropTransportForTest();
+ assertTrue(waitUntil(() -> ca.getState() == ConnectionState.RECONNECTING, 5_000));
+ AtomicReference err = new AtomicReference();
+ AtomicReference result = new AtomicReference();
+ Thread th = new Thread(() -> {
+ try {
+ result.set(ca.sendSync(new Target("endpoint", b), new Body("queued"), immediate(mid)));
+ } catch (Throwable t) {
+ err.set(t);
+ }
+ });
+ th.start();
+ th.join(60_000L);
+ assertTrue(err.get() == null);
+ assertNotNull(result.get());
+ assertEquals(mid, result.get().id);
+ assertTrue(waitUntil(() -> ca.getState() == ConnectionState.ONLINE, 30_000));
+ List got = box.waitN(1, 15_000);
+ assertEquals(1, got.size());
+ assertEquals(mid, got.get(0).id);
+ Thread.sleep(800);
+ assertEquals(1, box.items.size());
+ ca.close();
+ cb.close();
+ }
+
+ @Test
+ public void test04SameMessageId() throws Exception {
+ String a = uid("a4");
+ String b = uid("b4");
+ register(a);
+ register(b);
+ Client ca = connect(a);
+ Client cb = connect(b);
+ MsgBox box = new MsgBox();
+ cb.onMessage(box::onMessage);
+ String mid = Uuid7.next();
+ ca.sendSync(new Target("endpoint", b), new Body("idem"), immediate(mid));
+ assertEquals(1, box.waitN(1, 10_000).size());
+ ca.sendSync(new Target("endpoint", b), new Body("idem"), immediate(mid));
+ Thread.sleep(800);
+ assertEquals(1, box.items.size());
+ ca.close();
+ cb.close();
+ }
+
+ @Test
+ public void test05RecallWithinDelay() throws Exception {
+ String a = uid("a5");
+ String b = uid("b5");
+ register(a);
+ register(b);
+ Client ca = connect(a);
+ Client cb = connect(b);
+ MsgBox box = new MsgBox();
+ List revoked = Collections.synchronizedList(new ArrayList());
+ cb.onMessage(box::onMessage);
+ cb.onRevoked(revoked::add);
+ String mid = Uuid7.next();
+ SendOptions opt = immediate(mid);
+ opt.delayMs = 10_000L;
+ SendResult r = ca.sendSync(new Target("endpoint", b), new Body("will-recall"), opt);
+ assertEquals("scheduled", r.state);
+ ca.recall(mid).get(10, TimeUnit.SECONDS);
+ Thread.sleep(1200);
+ assertTrue(box.items.isEmpty());
+ assertTrue(revoked.isEmpty());
+ ca.close();
+ cb.close();
+ }
+
+ @Test
+ public void test06Scheduled2s() throws Exception {
+ String a = uid("a6");
+ String b = uid("b6");
+ register(a);
+ register(b);
+ Client ca = connect(a);
+ Client cb = connect(b);
+ MsgBox box = new MsgBox();
+ cb.onMessage(box::onMessage);
+ String mid = Uuid7.next();
+ SendOptions opt = immediate(mid);
+ opt.delayMs = 2000L;
+ long t0 = System.currentTimeMillis();
+ ca.sendSync(new Target("endpoint", b), new Body("later"), opt);
+ List got = box.waitN(1, 12_000);
+ long elapsed = System.currentTimeMillis() - t0;
+ assertEquals(1, got.size());
+ assertTrue(elapsed >= 1500);
+ assertTrue(elapsed < 8000);
+ ca.close();
+ cb.close();
+ }
+
+ @Test
+ public void test07OfflineKeep() throws Exception {
+ String a = uid("a7");
+ String bok = uid("bok");
+ String bms = uid("bms");
+ register(a);
+ register(bok);
+ register(bms);
+ Client ca = connect(a);
+ String mid1 = Uuid7.next();
+ SendOptions keep = immediate(mid1);
+ keep.keep = true;
+ keep.ttlSeconds = 86400L;
+ ca.sendSync(new Target("endpoint", bok), new Body("keep-ok"), keep);
+ Thread.sleep(1000);
+ Client cb1 = connect(bok);
+ MsgBox box1 = new MsgBox();
+ cb1.onMessage(box1::onMessage);
+ assertEquals(1, box1.waitN(1, 10_000).size());
+ cb1.close();
+
+ List receipts = Collections.synchronizedList(new ArrayList());
+ ca.onReceipt(receipts::add);
+ String mid2 = Uuid7.next();
+ SendOptions keepExp = immediate(mid2);
+ keepExp.keep = true;
+ keepExp.ttlSeconds = 1L;
+ ca.sendSync(new Target("endpoint", bms), new Body("keep-expire"), keepExp);
+ Thread.sleep(3200);
+ Client cb2 = connect(bms);
+ MsgBox box2 = new MsgBox();
+ cb2.onMessage(box2::onMessage);
+ Thread.sleep(1500);
+ assertTrue(box2.items.isEmpty());
+ assertTrue(waitUntil(() -> {
+ for (Receipt r : receipts) {
+ if ("expired".equals(r.state) && mid2.equals(r.id)) {
+ return true;
+ }
+ }
+ return false;
+ }, 10_000));
+ ca.close();
+ cb2.close();
+ }
+
+ @Test
+ public void test08GroupNoEcho() throws Exception {
+ String a = uid("a8");
+ String b = uid("b8");
+ String c = uid("c8");
+ register(a);
+ register(b);
+ register(c);
+ Client ca = connect(a);
+ Client cb = connect(b);
+ Client cc = connect(c);
+ String gid = "g_" + a;
+ List