Compare commits

...
Author SHA1 Message Date
Nixevol a22912d95a test: 补齐 Q2 验收与本机 Q3 弱网崩溃压测 2026-09-30 08:59:20 +08:00
Nixevol 1f715e4e51 feat: Python/Java SDK 接入清单与打包文档
对真实 nixmsg 跑 DEVELOPMENT 第 9 节接入清单(跳过仅 JS 跨域),补 README/示例,并修 Paho/HiveMQ 真机联调死锁与鉴权分类。
2026-09-30 08:57:06 +08:00
27 changed files with 3251 additions and 456 deletions
+67
View File
@@ -853,6 +853,43 @@
- 备选方案:依赖库自带重连再改 Clean Start(易漏)。
- 影响:无。
### S2-PY/JAVA 4–5 2026-09-30
1. **接入清单对真实 nixmsg,跳过仅 JS 跨域**
- 原条款:DEVELOPMENT 第 9 节 15 条;任务 4 用 T0.5 启动器起真实服务端。
- 实际做法:Python `tests/harness.py` + `test_checklist.py`、Java `TestHarness` + `ChecklistTest` 自行 `go build`/`admin init`/`serve`(临时目录、`127.0.0.1:0`),管理登录后 `PUT /api/admin/registration` 开注册。覆盖清单 1–13、15;第 14 条(仅 JS 跨域)不做。
- 原因:总控指示跳过 JS 专属跨域;不改服务器业务代码。
- 备选方案:复用 Go `test/harness` 包(SDK 测试不便依赖)。
- 影响:无。
2. **清单第 4 条「ack 丢失后服务器重推」未在真机选择性复现**
- 原条款:模拟 ack 丢失后服务器重推,SDK 自动再确认且不重复回调。
- 实际做法:集成测验证同消息号防重与断线入队重交送达一次;「选择性丢弃 SDK 发出的 ack 帧」在真实 broker 上做不到,去重再 ack 仍由假传输单测覆盖。
- 原因:不改服务器、无中间代理注入丢包。
- 备选方案:toxiproxy 按包过滤(超出本任务、且难按 MQTT 应用帧过滤)。
- 影响:清单 4 真机为部分通过;假传输路径完整。
3. **下行 `resp` 与业务帧分流,避免 auto_ack 自死锁**
- 原条款:自动模式回调后发 ack;回调串行。
- 实际做法:MQTT/`publishes` 回调里对 `resp` 立即完成 pending;`msg` 等进单线程队列再处理(可在队列线程里同步 `ack`/`request`)。Paho 使用 `MQTTv5` 常量与 `transport=websockets`,并等待 SUBACK。
- 原因:若 `resp` 与 `msg` 同队列,auto_ack 等待 `resp` 会永久卡住。
- 备选方案:ack 只发布不等待(弱化协议确认)。
- 影响:与 DEVELOPMENT 行为一致,修复真机联调阻塞。
4. **HiveMQ 鉴权失败与顶号原因码解析**
- 原条款:CONNACK 鉴权失败停重连;`0x8E` 顶号停重连。
- 实际做法:`connect().get()` 抛出的 `Mqtt5ConnAckException` / 文案含 `BAD_USER_*` 时归为 `bad_credentials`(令牌场景 Client 层改为 `session_invalid`);断开原因从 `Mqtt5DisconnectException` 读 `SESSION_TAKEN_OVER`。`connectSync` 等到终态再返回,避免与 attemptConnect 竞态报 `busy/RECONNECTING`。
- 原因:HiveMQ 失败路径多为异常而非成功返回的 CONNACK 对象。
- 备选方案:无。
- 影响:无。
5. **任务 5:README/示例与打包试跑,不发布**
- 原条款:包名与许可证;工具试跑确认能打包;README 与最小示例;真正发布在 Z3。
- 实际做法:更新两端 README;Python `examples/minimal.py`;Java `asia.asio.nixmsg.examples.MinimalExample`;`python -m build` 产出 wheel/sdist;`mvn package -DskipTests` 产出 jar;`javap` major version 52(Java 8)。未上传 PyPI/Maven。
- 原因:本波范围。
- 备选方案:无。
- 影响:无。
## 测试交付 Q
### Q1 / Q4 骨架 2026-09-30
@@ -935,3 +972,33 @@
- 原因:避免每次单测改 `generated_at` 弄脏工作区。
- 备选方案:固定时间戳始终写入仓库。
- 影响:交付审阅以 `ACCEPTANCE.md` / `q2_results.json` 为准,需先跑过 `task q:accept`。
### Q2 补齐 + Q3(本机 Windows)2026-09-30
1. **Q2 对照表按已接线能力重跑,不再把已挂路由写成未测**
- 原条款:TASKS Q2;PRD 第 10 节 F01–F23。
- 实际做法:在 `main@77d2dbd` 上重跑管理登录/CSRF/锁定、注册开关错码对码换码、开通端、单聊送达、延迟撤回、群发发送者不收到、离线保留上线送达、崩溃后续传;结果写入 `test/report/testdata/q2_results.json` 与 `ACCEPTANCE.md`。未覆盖项仍标未测并写明原因。不改业务逻辑求绿。
- 原因:L-WIRE / L-UPLINK / A3 / I5 已合入,旧报告过时。
- 备选方案:无。
- 影响:对照表通过项增加;未测项收窄到 SDK、拨钟类与部分身份/回执场景。
2. **Q3 弱网用 toxiproxy,不用本机 netem**
- 原条款:DEVELOPMENT 第 13 节 toxiproxy + Linux netem 20% 丢包。
- 实际做法:`test/chaos` 以项目名 `q3-chaos`、容器名 `q3-toxiproxy` 起官方镜像;注入延迟与 `reset_peer`,验证离线保留期内消息最终送达。Linux netem 20% 丢包记**未测**:本机 Windows,宿主无 tc;仅对代理容器挂 netshoot 不等于 NixMsg 端到端丢包验收。
- 原因:机器限制;TASKS 允许 Docker 内 Linux 测丢包,但本波未把业务进程放进同网络 Linux 容器做 netem。
- 备选方案:后续用 Linux 宿主或 compose 把 nixmsg 与 netem 旁路同网再测。
- 影响:丢包数字不进交付;延迟/断开路径有集成测。
3. **Q3 压测达不到 1000 连接 / 10 分钟**
- 原条款:PRD 第 8 节 / DEVELOPMENT 第 13 节「1000 连接、每秒 200 条、10 分钟」。
- 实际做法:`test/load` 短时 32 连接(16 对)真实登录+单聊收发成功;不宣称 1000/10min 通过。
- 原因:本机 Windows 开发机资源与并行 Agent 负载;强行 1000 长时间易误伤其他线。
- 备选方案:专用压测机或 Linux 服务器上再跑满指标。
- 影响:对照表与 DEVIATIONS 明示机器限制,不假装通过。
4. **崩溃续传用 accept.ManagedServer 启停**
- 原条款:提交成功后杀进程,重启后续传。
- 实际做法:`test/accept.ManagedServer`(Kill + 同 data_dir Restart)+ Q2/Q3 用例;不改 harness 公共 API(harness 属总控)。
- 原因:隔离目录约束。
- 备选方案:扩展 harness.Restart(需总控改)。
- 影响:Q 线自带启停辅助。
+37 -4
View File
@@ -2,7 +2,7 @@
坐标:`asia.asio.nixmsg:nixmsg-sdk`
包名:`asia.asio.nixmsg`
字节码目标:Java 8
字节码目标:Java 8(`maven.compiler.release=8`)
接口:`CompletableFuture`
## Android
@@ -17,7 +17,15 @@ Maven / Gradle 仓库:
https://git.asio.asia/api/packages/nixevol/maven
```
HiveMQ MQTT Client(含 WebSocket:`webSocketConfig` + `netty-codec-http`)。
```xml
<dependency>
<groupId>asia.asio.nixmsg</groupId>
<artifactId>nixmsg-sdk</artifactId>
<version>0.1.0</version>
</dependency>
```
HiveMQ MQTT Client(WebSocket:`webSocketConfig` + `netty-codec-http`)。
## 最小示例
@@ -26,9 +34,34 @@ Client c = new Client();
c.onSession(token -> { /* 应用保存 */ });
c.onMessage(msg -> System.out.println(msg.id + " " + msg.body.data));
c.connect("ws://127.0.0.1:7443/mqtt", "device-1", "secret", null)
.thenCompose(v -> c.send(new Types.Target("endpoint", "device-2"), new Types.Body("hello"), new Types.SendOptions()))
.thenCompose(v -> {
Types.SendOptions opt = new Types.SendOptions();
opt.delayMs = 0L;
return c.send(new Types.Target("endpoint", "device-2"), new Types.Body("hello"), opt);
})
.join();
c.close();
```
许可证见 `LICENSE`(专有)。
命令行示例类:`asia.asio.nixmsg.examples.MinimalExample`。
## 打包(不发布)
```bash
mvn package -DskipTests
# 产物 target/nixmsg-sdk-0.1.0.jar;勿部署到 Maven 仓库;正式发布由总控在阶段 3 执行
```
确认字节码为 8:`javap -v target/classes/asia/asio/nixmsg/Client.class | findstr major`(应为 52)。
## 测试
```bash
mvn test
```
含假传输单元测试与 DEVELOPMENT 第 9 节接入清单(会编译并启动真实 `nixmsg`)。可用环境变量 `NIXMSG_BIN` 指定已编译二进制。跳过仅 JS 的跨域项。
## 许可证
见 `LICENSE`(Proprietary)。
@@ -27,7 +27,9 @@ import java.util.Iterator;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.Map;
import java.util.concurrent.BlockingQueue;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.LinkedBlockingQueue;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.concurrent.atomic.AtomicLong;
@@ -82,6 +84,8 @@ public final class Client {
private final Object connWait = new Object();
private Thread worker;
private final Object wake = new Object();
private final BlockingQueue<byte[]> downQueue = new LinkedBlockingQueue<byte[]>();
private final Thread downWorker;
private Consumer<String> sessionHandler;
private Consumer<IncomingMessage> messageHandler;
@@ -106,6 +110,9 @@ public final class Client {
this.clientName = clientName;
this.connectTimeoutMs = connectTimeoutMs;
this.transport.setHandlers(this::onTransportConnected, this::onTransportDisconnected, this::onDown);
this.downWorker = new Thread(this::downLoop, "nixmsg-down");
this.downWorker.setDaemon(true);
this.downWorker.start();
}
public void onSession(Consumer<String> handler) { this.sessionHandler = handler; }
@@ -121,6 +128,11 @@ public final class Client {
public String getSessionToken() { return sessionToken; }
public long getClockSkewMs() { return clockSkewMs; }
/** 同包测试用:断开底层传输以触发重连与发送队列重交。 */
void dropTransportForTest() {
transport.disconnect();
}
public CompletableFuture<Void> connect(String url, String endpointId, String password, String sessionToken) {
return connect(url, endpointId, password, sessionToken, false);
}
@@ -157,9 +169,17 @@ public final class Client {
}
long deadline = System.currentTimeMillis() + connectTimeoutMs + 5000;
synchronized (connWait) {
while (!connReady.get() && System.currentTimeMillis() < deadline) {
while (System.currentTimeMillis() < deadline) {
ConnectionState s = state;
if (s == ConnectionState.ONLINE
|| s == ConnectionState.AUTH_FAILED
|| s == ConnectionState.KICKED
|| s == ConnectionState.OFFLINE
|| handshakeError != null) {
break;
}
try {
connWait.wait(200);
connWait.wait(100);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
break;
@@ -200,6 +220,7 @@ public final class Client {
transport.disconnect();
} catch (Exception ignored) {
}
downQueue.offer(new byte[0]); // 空载荷哨兵:downLoop 见 closed 退出
wakeUp();
}
@@ -716,57 +737,108 @@ public final class Client {
}
private void onDown(byte[] payload) {
// resp 立即完成 pending,避免 down 工作线程在 autoAck 等待时自死锁。
Map<String, Object> frame;
try {
frame = Protocol.loads(payload);
} catch (Exception e) {
downQueue.offer(payload);
return;
}
if ("resp".equals(str(frame.get("type"), ""))) {
dispatchResp(frame);
return;
}
downQueue.offer(payload);
}
private void downLoop() {
while (true) {
try {
byte[] payload = downQueue.take();
if (payload.length == 0 && closed) {
return;
}
Map<String, Object> frame;
try {
frame = Protocol.loads(payload);
} catch (Exception e) {
continue;
}
if ("resp".equals(str(frame.get("type"), ""))) {
dispatchResp(frame);
} else {
dispatchDownBody(frame);
}
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
return;
} catch (Exception e) {
LOG.log(Level.WARNING, "处理下行帧失败", e);
}
}
}
private void dispatchResp(Map<String, Object> frame) {
String rid = str(frame.get("rid"), "");
Pending p;
synchronized (lock) {
p = pending.remove(rid);
}
if (p == null) {
return;
}
if (p.isSend) {
synchronized (lock) {
inflightSends = Math.max(0, inflightSends - 1);
}
if (!Boolean.TRUE.equals(frame.get("ok"))) {
Map<String, Object> err = asMap(frame.get("error"));
if ("rate_limited".equals(str(err.get("code"), ""))) {
synchronized (lock) {
p.rid = "";
p.response = null;
p.error = null;
}
wakeUp();
return;
}
p.error = new NixMsgException(str(err.get("code"), "bad_request"), str(err.get("message"), ""));
}
p.response = frame;
synchronized (lock) {
Iterator<SendItem> it = sendQueue.iterator();
while (it.hasNext()) {
if (it.next().pending == p) {
it.remove();
break;
}
}
}
p.future.complete(null);
wakeUp();
} else {
p.response = frame;
p.future.complete(null);
}
}
private void dispatchDown(byte[] payload) {
Map<String, Object> frame;
try {
frame = Protocol.loads(payload);
} catch (Exception e) {
return;
}
String type = str(frame.get("type"), "");
if ("resp".equals(type)) {
String rid = str(frame.get("rid"), "");
Pending p;
synchronized (lock) {
p = pending.remove(rid);
}
if (p != null) {
if (p.isSend) {
synchronized (lock) {
inflightSends = Math.max(0, inflightSends - 1);
}
if (!Boolean.TRUE.equals(frame.get("ok"))) {
Map<String, Object> err = asMap(frame.get("error"));
if ("rate_limited".equals(str(err.get("code"), ""))) {
synchronized (lock) {
p.rid = "";
p.response = null;
p.error = null;
// 保留原 future,重交成功后再 complete
}
wakeUp();
return;
}
p.error = new NixMsgException(str(err.get("code"), "bad_request"), str(err.get("message"), ""));
}
p.response = frame;
synchronized (lock) {
Iterator<SendItem> it = sendQueue.iterator();
while (it.hasNext()) {
if (it.next().pending == p) {
it.remove();
break;
}
}
}
p.future.complete(null);
wakeUp();
} else {
p.response = frame;
p.future.complete(null);
}
}
if ("resp".equals(str(frame.get("type"), ""))) {
dispatchResp(frame);
return;
}
dispatchDownBody(frame);
}
private void dispatchDownBody(Map<String, Object> frame) {
String type = str(frame.get("type"), "");
if ("msg".equals(type)) {
handleMsg(frame);
return;
@@ -3,9 +3,10 @@ package asia.asio.nixmsg;
import com.hivemq.client.mqtt.MqttClient;
import com.hivemq.client.mqtt.MqttGlobalPublishFilter;
import com.hivemq.client.mqtt.datatypes.MqttQos;
import com.hivemq.client.mqtt.lifecycle.MqttDisconnectSource;
import com.hivemq.client.mqtt.mqtt5.Mqtt5AsyncClient;
import com.hivemq.client.mqtt.mqtt5.Mqtt5ClientBuilder;
import com.hivemq.client.mqtt.mqtt5.exceptions.Mqtt5ConnAckException;
import com.hivemq.client.mqtt.mqtt5.exceptions.Mqtt5DisconnectException;
import com.hivemq.client.mqtt.mqtt5.message.connect.connack.Mqtt5ConnAck;
import com.hivemq.client.mqtt.mqtt5.message.connect.connack.Mqtt5ConnAckReasonCode;
import com.hivemq.client.mqtt.mqtt5.message.disconnect.Mqtt5Disconnect;
@@ -16,6 +17,7 @@ import java.nio.charset.StandardCharsets;
import java.util.ArrayList;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.Locale;
import java.util.Map;
import java.util.concurrent.CopyOnWriteArrayList;
import java.util.concurrent.TimeUnit;
@@ -256,25 +258,24 @@ final class HiveMqTransport implements Transport {
.addDisconnectedListener(context -> {
String reason = "network";
boolean stop = false;
if (context.getSource() == MqttDisconnectSource.SERVER) {
try {
java.lang.reflect.Method m = context.getClass().getMethod("getMqttDisconnect");
Object disc = m.invoke(context);
if (disc instanceof Mqtt5Disconnect) {
Mqtt5DisconnectReasonCode rc = ((Mqtt5Disconnect) disc).getReasonCode();
if (rc == Mqtt5DisconnectReasonCode.SESSION_TAKEN_OVER) {
reason = "taken_over";
stop = true;
}
}
} catch (Exception ignored) {
}
}
Throwable cause = context.getCause();
if (cause != null && cause.getMessage() != null
&& cause.getMessage().toLowerCase().contains("taken over")) {
reason = "taken_over";
stop = true;
while (cause != null) {
if (cause instanceof Mqtt5DisconnectException) {
Mqtt5DisconnectReasonCode rc =
((Mqtt5DisconnectException) cause).getMqttMessage().getReasonCode();
if (rc == Mqtt5DisconnectReasonCode.SESSION_TAKEN_OVER) {
reason = "taken_over";
stop = true;
}
break;
}
String m = cause.getMessage() == null ? "" : cause.getMessage().toLowerCase(Locale.ROOT);
if (m.contains("taken over") || m.contains("session taken")) {
reason = "taken_over";
stop = true;
break;
}
cause = cause.getCause();
}
BiConsumer<String, Boolean> h = onDisconnected;
if (h != null) {
@@ -332,9 +333,29 @@ final class HiveMqTransport implements Transport {
}
}
} catch (Exception e) {
Mqtt5ConnAckReasonCode rc = extractConnAckReason(e);
String reason;
boolean stop;
if (rc != null) {
reason = classify(rc);
stop = isStop(rc);
} else {
String msg = exceptionText(e).toLowerCase(Locale.ROOT);
if (msg.contains("bad_user") || msg.contains("bad user") || msg.contains("not authorized")
|| msg.contains("not_authorized") || msg.contains("bad_username")
|| msg.contains("bad username") || msg.contains("banned")
|| msg.contains("connack") || msg.contains("connectionfailed")
|| msg.contains("mqtt5connack")) {
reason = "bad_credentials";
stop = true;
} else {
reason = "network";
stop = false;
}
}
BiConsumer<String, Boolean> 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
@@ -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;
/**
* 最小示例:连接、发送、关闭。
* <p>
* 运行(需本机已有 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 <wsUrl> <endpointId> <password> <peerId>");
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();
}
}
@@ -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<IncomingMessage> items = Collections.synchronizedList(new ArrayList<IncomingMessage>());
private final CountDownLatch latch = new CountDownLatch(1);
void onMessage(IncomingMessage m) {
items.add(m);
latch.countDown();
}
List<IncomingMessage> waitN(int n, long timeoutMs) throws InterruptedException {
long deadline = System.currentTimeMillis() + timeoutMs;
while (System.currentTimeMillis() < deadline) {
if (items.size() >= n) {
return new ArrayList<IncomingMessage>(items);
}
Thread.sleep(50L);
}
return new ArrayList<IncomingMessage>(items);
}
}
@Test
public void test01Handshake() {
String id = uid("hs");
register(id);
Client c = new Client();
final List<String> tokens = new ArrayList<String>();
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<IncomingMessage> 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<Throwable> err = new AtomicReference<Throwable>();
AtomicReference<SendResult> result = new AtomicReference<SendResult>();
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<IncomingMessage> 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<RevokedEvent> revoked = Collections.synchronizedList(new ArrayList<RevokedEvent>());
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<IncomingMessage> 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<Receipt> receipts = Collections.synchronizedList(new ArrayList<Receipt>());
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<Map<String, String>> members = new ArrayList<Map<String, String>>();
members.add(Collections.singletonMap("id", b));
members.add(Collections.singletonMap("id", c));
ca.groupCreate("G", members, gid).get(15, TimeUnit.SECONDS);
Thread.sleep(400);
MsgBox boxA = new MsgBox();
MsgBox boxB = new MsgBox();
MsgBox boxC = new MsgBox();
ca.onMessage(boxA::onMessage);
cb.onMessage(boxB::onMessage);
cc.onMessage(boxC::onMessage);
String mid = Uuid7.next();
ca.sendSync(new Target("group", gid), new Body("hi-g"), immediate(mid));
assertEquals(1, boxB.waitN(1, 10_000).size());
assertEquals(1, boxC.waitN(1, 10_000).size());
Thread.sleep(800);
assertTrue(boxA.items.isEmpty());
ca.close();
cb.close();
cc.close();
}
@Test
public void test09TalkPassword() throws Exception {
String a = uid("a9");
String b = uid("b9");
register(a);
register(b);
Client ca = connect(a);
Client cb = connect(b);
cb.setTalkPassword("talk99").get(10, TimeUnit.SECONDS);
try {
ca.sendSync(new Target("endpoint", b), new Body("no"), immediate(Uuid7.next()));
fail("expected talk password error");
} catch (NixMsgException e) {
assertTrue(e.getCode().contains("talk_password"));
}
ca.unlock(b, "talk99").get(10, TimeUnit.SECONDS);
MsgBox box = new MsgBox();
cb.onMessage(box::onMessage);
ca.sendSync(new Target("endpoint", b), new Body("ok"), immediate(Uuid7.next()));
assertEquals(1, box.waitN(1, 10_000).size());
cb.setTalkPassword("talk00").get(10, TimeUnit.SECONDS);
try {
ca.sendSync(new Target("endpoint", b), new Body("fail"), immediate(Uuid7.next()));
fail("expected talk password error after change");
} catch (NixMsgException e) {
assertTrue(e.getCode().contains("talk_password"));
}
ca.setTalkPassword("alicepw").get(10, TimeUnit.SECONDS);
MsgBox box2 = new MsgBox();
ca.onMessage(box2::onMessage);
SendOptions first = immediate(Uuid7.next());
first.talkPassword = "alicepw";
cb.sendSync(new Target("endpoint", a), new Body("first"), first);
assertEquals(1, box2.waitN(1, 10_000).size());
MsgBox box3 = new MsgBox();
cb.onMessage(box3::onMessage);
ca.sendSync(new Target("endpoint", b), new Body("reply"), immediate(Uuid7.next()));
assertEquals(1, box3.waitN(1, 10_000).size());
ca.close();
cb.close();
}
@Test
public void test10KickNoReconnect() throws Exception {
String id = uid("k10");
register(id);
Client c1 = connect(id);
Client c2 = connect(id);
assertTrue(waitUntil(() -> c1.getState() == ConnectionState.KICKED, 15_000));
Thread.sleep(2500);
assertEquals(ConnectionState.KICKED, c1.getState());
assertEquals(ConnectionState.ONLINE, c2.getState());
c1.close();
c2.close();
}
@Test
public void test11BodyTooLarge() {
String id = uid("big");
register(id);
Client c = connect(id);
StringBuilder sb = new StringBuilder();
for (int i = 0; i < 256 * 1024 + 1; i++) {
sb.append('x');
}
try {
c.sendSync(new Target("endpoint", id), new Body(sb.toString()), immediate(Uuid7.next()));
fail("expected body_too_large");
} catch (NixMsgException e) {
assertEquals("body_too_large", e.getCode());
}
c.close();
}
@Test
public void test12Registration() throws Exception {
String code = srv.regCode;
srv.setRegistration(false, code);
try {
RegisterOptions opt = new RegisterOptions();
opt.id = uid("r12a");
opt.loginPassword = "password12";
Client.registerSync(srv.wsUrl, code, opt);
fail("closed");
} catch (NixMsgException e) {
assertEquals("registration_closed", e.getCode());
}
srv.setRegistration(true, code);
try {
RegisterOptions opt = new RegisterOptions();
opt.id = uid("r12b");
opt.loginPassword = "password12";
Client.registerSync(srv.wsUrl, "wrong-code-xx", opt);
fail("bad code");
} catch (NixMsgException e) {
assertEquals("registration_code_invalid", e.getCode());
}
String eid = uid("r12c");
RegisterOptions ok = new RegisterOptions();
ok.id = eid;
ok.loginPassword = "password12";
Client.registerSync(srv.wsUrl, code, ok);
Client c = connect(eid);
c.close();
String newCode = "s2java-new-code";
srv.setRegistration(true, newCode);
try {
RegisterOptions opt = new RegisterOptions();
opt.id = uid("r12d");
opt.loginPassword = "password12";
Client.registerSync(srv.wsUrl, code, opt);
fail("old code");
} catch (NixMsgException e) {
assertTrue(e.getCode().length() > 0);
}
Client c2 = connect(eid);
c2.close();
srv.setRegistration(true, code);
srv.regCode = code;
}
@Test
public void test13ChangeLoginPassword() throws Exception {
String id = uid("pw13");
register(id);
Client c = connect(id);
c.changeLoginPassword("password12", "password99").get(15, TimeUnit.SECONDS);
c.close();
Client c2 = connect(id, "password99", null);
c2.close();
Client c3 = new Client();
try {
c3.connectSync(srv.wsUrl, id, "password12", null, false);
fail("old password");
} catch (NixMsgException e) {
assertTrue(
"code=" + e.getCode() + " msg=" + e.getMessage(),
e.getCode().contains("bad_credentials")
|| e.getCode().contains("auth")
|| e.getCode().contains("session_invalid")
|| "busy".equals(e.getCode()) && c3.getState() == ConnectionState.AUTH_FAILED);
}
Thread.sleep(2000);
assertEquals(ConnectionState.AUTH_FAILED, c3.getState());
c3.close();
}
@Test
public void test15SessionToken() throws Exception {
String id = uid("tok");
register(id);
Client c = new Client();
List<String> tokens = new ArrayList<String>();
c.onSession(tokens::add);
c.connectSync(srv.wsUrl, id, "password12", null, false);
assertFalse(tokens.isEmpty());
String token = tokens.get(0);
c.close();
Client c2 = connect(id, null, token);
c2.close();
Client c3 = connect(id);
String newTok = c3.getSessionToken();
assertNotNull(newTok);
assertFalse(token.equals(newTok));
c3.close();
Client c4 = new Client();
NixMsgException c4err = null;
try {
c4.connectSync(srv.wsUrl, id, null, token, false);
fail("old token");
} catch (NixMsgException e) {
c4err = e;
}
assertTrue(waitUntil(() -> c4.getState() == ConnectionState.AUTH_FAILED, 10_000));
assertNotNull(c4err);
c4.close();
Client c5 = connect(id);
String tok5 = c5.getSessionToken();
c5.logout().get(10, TimeUnit.SECONDS);
Thread.sleep(300);
Client c6 = new Client();
NixMsgException c6err = null;
try {
c6.connectSync(srv.wsUrl, id, null, tok5, false);
fail("logout token");
} catch (NixMsgException e) {
c6err = e;
}
assertTrue(waitUntil(() -> c6.getState() == ConnectionState.AUTH_FAILED, 10_000));
assertNotNull(c6err);
c6.close();
}
}
@@ -0,0 +1,251 @@
package asia.asio.nixmsg;
import com.google.gson.Gson;
import com.google.gson.JsonObject;
import com.google.gson.JsonParser;
import java.io.BufferedReader;
import java.io.ByteArrayOutputStream;
import java.io.File;
import java.io.IOException;
import java.io.InputStream;
import java.io.InputStreamReader;
import java.io.OutputStream;
import java.net.CookieHandler;
import java.net.CookieManager;
import java.net.HttpURLConnection;
import java.net.URL;
import java.nio.charset.StandardCharsets;
import java.nio.file.Files;
import java.nio.file.Path;
import java.nio.file.Paths;
import java.util.ArrayList;
import java.util.Comparator;
import java.util.List;
import java.util.Locale;
import java.util.concurrent.TimeUnit;
/** 真实 nixmsg:临时目录、127.0.0.1:0、admin init、开注册。 */
final class TestHarness {
private static final Gson GSON = new Gson();
final Path dataDir;
final Path configPath;
final Path binary;
final String adminPassword;
Process process;
String httpBase;
String wsUrl;
String regCode = "s2java-reg-code";
TestHarness() throws Exception {
binary = ensureBinary();
dataDir = Files.createTempDirectory("nixmsg-s2-java-");
configPath = dataDir.resolve("config.yaml");
String yaml = "listen: \"127.0.0.1:0\"\ndata_dir: \""
+ dataDir.toAbsolutePath().toString().replace('\\', '/') + "\"\n";
Files.write(configPath, yaml.getBytes(StandardCharsets.UTF_8));
adminPassword = adminInit();
}
static Path findRepoRoot() throws IOException {
Path p = Paths.get("").toAbsolutePath().normalize();
for (int i = 0; i < 12; i++) {
if (Files.isRegularFile(p.resolve("go.mod")) && Files.isDirectory(p.resolve("cmd").resolve("nixmsg"))) {
return p;
}
Path parent = p.getParent();
if (parent == null) {
break;
}
p = parent;
}
throw new IOException("找不到仓库根 go.mod(cwd=" + Paths.get("").toAbsolutePath() + ")");
}
static Path ensureBinary() throws Exception {
String env = System.getenv("NIXMSG_BIN");
if (env != null && !env.isEmpty() && Files.isRegularFile(Paths.get(env))) {
return Paths.get(env);
}
Path root = findRepoRoot();
Path cache = Paths.get(System.getProperty("java.io.tmpdir"), "nixmsg-s2-java-bin");
Files.createDirectories(cache);
boolean win = System.getProperty("os.name", "").toLowerCase(Locale.ROOT).contains("win");
Path out = cache.resolve(win ? "nixmsg.exe" : "nixmsg");
if (!Files.isRegularFile(out)) {
List<String> cmd = new ArrayList<String>();
cmd.add("go");
cmd.add("build");
cmd.add("-o");
cmd.add(out.toString());
cmd.add("./cmd/nixmsg");
ProcessBuilder pb = new ProcessBuilder(cmd);
pb.directory(root.toFile());
pb.environment().put("CGO_ENABLED", "0");
pb.redirectErrorStream(true);
Process p = pb.start();
String log = readAll(p.getInputStream());
if (!waitFor(p, 180) || p.exitValue() != 0) {
throw new IllegalStateException("go build 失败: " + log);
}
}
return out;
}
private String adminInit() throws Exception {
ProcessBuilder pb = new ProcessBuilder(binary.toString(), "admin", "init");
pb.environment().put("NIXMSG_CONFIG", configPath.toString());
pb.redirectErrorStream(true);
Process p = pb.start();
String out = readAll(p.getInputStream());
if (!waitFor(p, 60) || p.exitValue() != 0) {
throw new IllegalStateException("admin init 失败: " + out);
}
String[] lines = out.split("\\r?\\n");
for (String line : lines) {
String t = line.trim();
String lower = t.toLowerCase(Locale.ROOT);
if (lower.startsWith("admin password:")) {
return t.substring(t.indexOf(':') + 1).trim();
}
if (lower.startsWith("password:")) {
return t.substring(t.indexOf(':') + 1).trim();
}
}
throw new IllegalStateException("admin init 未解析密码: " + out);
}
void start() throws Exception {
ProcessBuilder pb = new ProcessBuilder(binary.toString(), "serve");
pb.environment().put("NIXMSG_CONFIG", configPath.toString());
File nul = new File(System.getProperty("os.name", "").toLowerCase(Locale.ROOT).contains("win") ? "NUL" : "/dev/null");
pb.redirectError(ProcessBuilder.Redirect.to(nul));
pb.redirectOutput(ProcessBuilder.Redirect.to(nul));
process = pb.start();
Path addrFile = dataDir.resolve("listen.addr");
long deadline = System.currentTimeMillis() + 20_000L;
String addr = null;
while (System.currentTimeMillis() < deadline) {
if (Files.isRegularFile(addrFile)) {
addr = new String(Files.readAllBytes(addrFile), StandardCharsets.UTF_8).trim();
if (!addr.isEmpty()) {
break;
}
}
if (!isAlive(process)) {
throw new IllegalStateException("serve 提前退出");
}
Thread.sleep(50L);
}
if (addr == null || addr.isEmpty()) {
stop();
throw new IllegalStateException("等待 listen.addr 超时");
}
httpBase = "http://" + addr;
wsUrl = "ws://" + addr + "/mqtt";
setRegistration(true, regCode);
}
void setRegistration(boolean enabled, String code) throws Exception {
CookieManager cm = new CookieManager();
CookieHandler.setDefault(cm);
postJson("/api/admin/login", "{\"username\":\"admin\",\"password\":" + GSON.toJson(adminPassword) + "}");
JsonObject body = new JsonObject();
body.addProperty("enabled", enabled);
if (code != null) {
body.addProperty("code", code);
}
JsonObject resp = putJson("/api/admin/registration", body.toString());
if (!resp.has("ok") || !resp.get("ok").getAsBoolean()) {
throw new IllegalStateException("registration put failed: " + resp);
}
}
private JsonObject postJson(String path, String json) throws Exception {
return mutate("POST", path, json);
}
private JsonObject putJson(String path, String json) throws Exception {
return mutate("PUT", path, json);
}
private JsonObject mutate(String method, String path, String json) throws Exception {
URL url = new URL(httpBase + path);
HttpURLConnection conn = (HttpURLConnection) url.openConnection();
conn.setRequestMethod(method);
conn.setDoOutput(true);
conn.setRequestProperty("Content-Type", "application/json");
conn.setRequestProperty("X-Nixmsg-Request", "1");
byte[] bytes = json.getBytes(StandardCharsets.UTF_8);
conn.setFixedLengthStreamingMode(bytes.length);
OutputStream os = conn.getOutputStream();
try {
os.write(bytes);
} finally {
os.close();
}
int code = conn.getResponseCode();
InputStream in = code >= 400 ? conn.getErrorStream() : conn.getInputStream();
String raw = in == null ? "{}" : readAll(in);
if (code >= 400) {
throw new IllegalStateException(method + " " + path + " -> " + code + " " + raw);
}
return new JsonParser().parse(raw).getAsJsonObject();
}
void stop() {
if (process != null && isAlive(process)) {
process.destroy();
try {
waitFor(process, 2);
} catch (InterruptedException ignored) {
Thread.currentThread().interrupt();
}
if (isAlive(process)) {
process.destroyForcibly();
}
}
process = null;
try {
if (Files.isDirectory(dataDir)) {
List<Path> paths = new ArrayList<Path>();
Files.walk(dataDir).sorted(Comparator.reverseOrder()).forEach(paths::add);
for (Path p : paths) {
try {
Files.deleteIfExists(p);
} catch (IOException ignored) {
}
}
}
} catch (IOException ignored) {
}
}
private static boolean waitFor(Process p, long seconds) throws InterruptedException {
return p.waitFor(seconds, TimeUnit.SECONDS);
}
private static boolean isAlive(Process p) {
try {
p.exitValue();
return false;
} catch (IllegalThreadStateException e) {
return true;
}
}
private static String readAll(InputStream in) throws IOException {
if (in == null) {
return "";
}
ByteArrayOutputStream bos = new ByteArrayOutputStream();
byte[] buf = new byte[4096];
int n;
while ((n = in.read(buf)) >= 0) {
bos.write(buf, 0, n);
}
in.close();
return new String(bos.toByteArray(), StandardCharsets.UTF_8);
}
}
+40 -4
View File
@@ -1,24 +1,60 @@
# NixMsg Python SDK
包名 `nixmsg`,最低 Python 3.10。同步接口为主,`AsyncClient` 提供 asyncio 包装。
包名 `nixmsg`,最低 Python 3.10。同步接口为主,同包提供 `AsyncClient` asyncio 包装。
## 安装
发布后(阶段 3):
```bash
pip install nixmsg --index-url https://git.asio.asia/api/packages/nixevol/pypi/simple/
```
本地开发:
```bash
cd sdk/python
python -m venv .venv
# Windows: .venv\Scripts\activate
pip install -e ".[dev]"
```
## 最小示例
```python
from nixmsg import Client, Target, Body
from nixmsg import Body, Client, SendOptions, Target
c = Client()
c.on_session(lambda token: print("session", token))
c.on_message(lambda msg: print("msg", msg.id, msg.body.data))
c.connect("ws://127.0.0.1:7443/mqtt", "device-1", password="secret")
c.send(Target(kind="endpoint", id="device-2"), Body(data="hello"))
c.send(
Target(kind="endpoint", id="device-2"),
Body(data="hello"),
SendOptions(delay_ms=0),
)
c.close()
```
许可证见 `LICENSE`(专有)。
更完整的命令行示例见 `examples/minimal.py`。
## 打包(不发布)
```bash
pip install build
python -m build
# 产物在 dist/,勿上传 PyPI;正式发布由总控在阶段 3 执行
```
## 测试
```bash
# 单元测试(假传输)+ 接入清单(会编译并启动真实 nixmsg)
pytest
```
接入清单覆盖 DEVELOPMENT 第 9 节(跳过仅 JS 的跨域项)。可用环境变量 `NIXMSG_BIN` 指定已编译二进制。
## 许可证
见 `LICENSE`(专有 / Proprietary)。
+31
View File
@@ -0,0 +1,31 @@
"""最小示例:连接、收发、关闭。
用法(需本地已启动 nixmsg,并开放注册或已有端):
python examples/minimal.py ws://127.0.0.1:PORT/mqtt device-1 password12 peer-id
"""
from __future__ import annotations
import sys
import time
from nixmsg import Body, Client, SendOptions, Target
def main() -> None:
if len(sys.argv) < 5:
print(__doc__)
raise SystemExit(2)
url, eid, password, peer = sys.argv[1:5]
c = Client()
c.on_session(lambda token: print("session", token[:16] + "..."))
c.on_message(lambda msg: print("msg", msg.from_id, msg.id, msg.body.data))
c.on_connection(lambda ev: print("conn", ev.state.value, ev.reason))
c.connect(url, eid, password=password)
c.send(Target(kind="endpoint", id=peer), Body(data="hello from python"), SendOptions(delay_ms=0))
time.sleep(2)
c.close()
if __name__ == "__main__":
main()
+71 -29
View File
@@ -11,6 +11,7 @@ import urllib.request
from collections import OrderedDict
from dataclasses import dataclass, field
from typing import Any, Optional
from queue import SimpleQueue
from .errors import ClosedError, NixMsgError, NotConnectedError
from .protocol import dumps, down_topic, loads, normalize_mqtt_ws_url, register_url_from_connect, up_topic
@@ -133,6 +134,9 @@ class Client:
self._worker: Optional[threading.Thread] = None
self._wake = threading.Event()
self._want_connected = False
self._down_q: SimpleQueue = SimpleQueue()
self._down_thread = threading.Thread(target=self._down_loop, name="nixmsg-down", daemon=True)
self._down_thread.start()
self._transport.set_handlers(self._on_transport_connected, self._on_transport_disconnected, self._on_down)
@@ -218,6 +222,10 @@ class Client:
self._transport.disconnect()
except Exception:
pass
try:
self._down_q.put(None)
except Exception:
pass
self._wake.set()
def logout(self) -> None:
@@ -632,40 +640,74 @@ class Client:
self._wake.set()
def _on_down(self, payload: bytes) -> None:
# resp 必须立即完成 pending(含 auto_ack 等待),不能进 down 队列,否则自死锁。
try:
frame = loads(payload)
except Exception:
self._down_q.put(payload)
return
if frame.get("type") == "resp":
self._dispatch_resp(frame)
return
if threading.current_thread() is self._down_thread:
self._dispatch_down_body(frame)
return
self._down_q.put(payload)
def _down_loop(self) -> None:
while True:
payload = self._down_q.get()
if payload is None:
return
try:
frame = loads(payload)
except Exception:
continue
try:
self._dispatch_down_body(frame)
except Exception:
log.exception("处理下行帧失败")
def _dispatch_resp(self, frame: dict[str, Any]) -> None:
rid = str(frame.get("rid", ""))
with self._lock:
pending = self._pending.pop(rid, None)
if not pending:
return
pending.response = frame
if pending.is_send:
with self._lock:
self._inflight_sends = max(0, self._inflight_sends - 1)
err = (frame.get("error") or {}) if not frame.get("ok") else {}
if not frame.get("ok") and str(err.get("code")) == "rate_limited":
with self._lock:
pending.rid = ""
pending.response = None
pending.error = None
pending.event.clear()
self._wake.set()
return
with self._lock:
self._send_queue = [it for it in self._send_queue if it.pending is not pending]
if not frame.get("ok"):
pending.error = NixMsgError(str(err.get("code", "bad_request")), str(err.get("message", "")))
pending.event.set()
self._wake.set()
else:
pending.event.set()
def _dispatch_down(self, payload: bytes) -> None:
try:
frame = loads(payload)
except Exception:
return
ftype = frame.get("type")
if ftype == "resp":
rid = str(frame.get("rid", ""))
with self._lock:
pending = self._pending.pop(rid, None)
if pending:
pending.response = frame
if pending.is_send:
with self._lock:
self._inflight_sends = max(0, self._inflight_sends - 1)
# rate_limited 重交:清状态后不 set event
err = (frame.get("error") or {}) if not frame.get("ok") else {}
if not frame.get("ok") and str(err.get("code")) == "rate_limited":
with self._lock:
pending.rid = ""
pending.response = None
pending.error = None
pending.event.clear()
self._wake.set()
return
# 从发送队列移除
with self._lock:
self._send_queue = [it for it in self._send_queue if it.pending is not pending]
if not frame.get("ok"):
pending.error = NixMsgError(str(err.get("code", "bad_request")), str(err.get("message", "")))
pending.event.set()
self._wake.set()
else:
pending.event.set()
if frame.get("type") == "resp":
self._dispatch_resp(frame)
return
self._dispatch_down_body(frame)
def _dispatch_down_body(self, frame: dict[str, Any]) -> None:
ftype = frame.get("type")
if ftype == "msg":
self._handle_msg(frame)
return
+48 -44
View File
@@ -8,7 +8,7 @@ from dataclasses import dataclass, field
from typing import Any, Callable, Optional, Protocol
from urllib.parse import urlparse
from paho.mqtt.client import CallbackAPIVersion, Client as PahoClient, MQTT_ERR_SUCCESS
from paho.mqtt.client import CallbackAPIVersion, Client as PahoClient, MQTT_ERR_SUCCESS, MQTTv5
from paho.mqtt.enums import MQTTErrorCode
from paho.mqtt.reasoncodes import ReasonCode
@@ -186,6 +186,8 @@ class PahoTransport:
self._on_down: Optional[DownHandler] = None
self._down_topic = ""
self._loop_started = False
self._sub_event = threading.Event()
self._sub_mid: Optional[int] = None
def set_handlers(
self,
@@ -199,44 +201,39 @@ class PahoTransport:
def connect(self, params: ConnectParams) -> None:
self.disconnect()
url = params.url
u = urlparse(url if "://" in url else "ws://" + url)
use_tcp = params.use_tcp or u.scheme in ("mqtt", "mqtts")
# WebSocket 必须显式 transport=websockets;裸 TCP 走默认。
client = PahoClient(
callback_api_version=CallbackAPIVersion.VERSION2,
client_id=params.client_id,
protocol=PahoClient.MQTTv5,
protocol=MQTTv5,
transport="tcp" if use_tcp else "websockets",
)
client.username_pw_set(params.username, params.password)
client.on_connect = self._on_connect
client.on_disconnect = self._on_disconnect
client.on_message = self._on_message
client.on_subscribe = self._on_subscribe
self._client = client
url = params.url
u = urlparse(url if "://" in url else "ws://" + url)
use_tcp = params.use_tcp or u.scheme in ("mqtt", "mqtts")
host = u.hostname or "localhost"
port = u.port or (8883 if u.scheme in ("wss", "mqtts") else 443 if u.scheme == "wss" else 80)
props = None
try:
from paho.mqtt.properties import Properties
from paho.mqtt.packettypes import PacketTypes
props = Properties(PacketTypes.CONNECT)
props.SessionExpiryInterval = params.session_expiry
except Exception:
props = None
if use_tcp:
if not u.port:
port = 8883 if u.scheme == "mqtts" else 1883
tls = u.scheme == "mqtts"
if tls:
if u.scheme == "mqtts":
client.tls_set()
props = None
try:
from paho.mqtt.properties import Properties
from paho.mqtt.packettypes import PacketTypes
props = Properties(PacketTypes.CONNECT)
props.SessionExpiryInterval = params.session_expiry
except Exception:
props = None
client.connect(
host,
port,
keepalive=params.keep_alive,
clean_start=params.clean_start,
properties=props,
)
else:
path = u.path or "/mqtt"
if not path.endswith("/mqtt"):
@@ -246,30 +243,28 @@ class PahoTransport:
if u.scheme == "wss":
client.tls_set()
client.ws_set_options(path=path, headers={"Sec-WebSocket-Protocol": "mqtt"})
props = None
try:
from paho.mqtt.properties import Properties
from paho.mqtt.packettypes import PacketTypes
props = Properties(PacketTypes.CONNECT)
props.SessionExpiryInterval = params.session_expiry
except Exception:
props = None
client.connect(
host,
port,
keepalive=params.keep_alive,
clean_start=params.clean_start,
properties=props,
)
client.connect(
host,
port,
keepalive=params.keep_alive,
clean_start=params.clean_start,
properties=props,
)
client.loop_start()
self._loop_started = True
# 等待连接结果由回调驱动;超时由 Client 层处理
def subscribe(self, topic: str) -> None:
self._down_topic = topic
if self._client:
self._client.subscribe(topic, qos=1)
if not self._client:
return
self._sub_event.clear()
result, mid = self._client.subscribe(topic, qos=1)
if result != MQTT_ERR_SUCCESS:
raise RuntimeError(f"subscribe failed: {result}")
self._sub_mid = mid
if not self._sub_event.wait(10):
raise RuntimeError("subscribe timeout")
def publish(self, topic: str, payload: bytes) -> None:
if not self._client:
@@ -296,13 +291,18 @@ class PahoTransport:
def _on_connect(self, client, userdata, flags, reason_code, properties) -> None:
code = _reason_to_int(reason_code)
if code == 0:
# 不在 loop 线程里同步做 subscribe+等待,否则会卡死 SUBACK
if self._on_connected:
self._on_connected()
threading.Thread(target=self._on_connected, name="nixmsg-on-connected", daemon=True).start()
return
stop, reason = _classify_connack(code)
stop, reason = _classify_connack(code if code is not None else -1)
if self._on_disconnected:
self._on_disconnected(reason, stop)
def _on_subscribe(self, client, userdata, mid, reason_codes, properties) -> None:
if self._sub_mid is None or mid == self._sub_mid:
self._sub_event.set()
def _on_disconnect(self, client, userdata, flags, reason_code, properties) -> None:
code = _reason_to_int(reason_code)
if code in (0, None):
@@ -329,9 +329,13 @@ def _reason_to_int(reason_code) -> Optional[int]:
if isinstance(reason_code, int):
return reason_code
if isinstance(reason_code, ReasonCode):
return int(reason_code)
return int(reason_code.value)
if isinstance(reason_code, MQTTErrorCode):
return int(reason_code)
# paho 偶发其它包装
val = getattr(reason_code, "value", None)
if isinstance(val, int):
return val
try:
return int(reason_code)
except Exception:
+1
View File
@@ -0,0 +1 @@
# tests package
+185
View File
@@ -0,0 +1,185 @@
"""真实 nixmsg 进程启动器:临时目录、127.0.0.1:0、admin init、开注册。"""
from __future__ import annotations
import json
import os
import shutil
import subprocess
import tempfile
import time
import urllib.error
import urllib.request
from dataclasses import dataclass
from http.cookiejar import CookieJar
from pathlib import Path
from typing import Any, Optional
from urllib.parse import urljoin
def _repo_root() -> Path:
here = Path(__file__).resolve()
for p in [here] + list(here.parents):
if (p / "go.mod").is_file() and (p / "cmd" / "nixmsg").is_dir():
return p
raise RuntimeError("找不到仓库根(go.mod)")
def ensure_binary() -> Path:
env = os.environ.get("NIXMSG_BIN")
if env:
p = Path(env)
if p.is_file():
return p
root = _repo_root()
cache = Path(tempfile.gettempdir()) / "nixmsg-s2-python-bin"
cache.mkdir(parents=True, exist_ok=True)
name = "nixmsg.exe" if os.name == "nt" else "nixmsg"
out = cache / name
# 若已有且较新则复用;否则编译
need = True
if out.is_file():
need = False
if need or os.environ.get("NIXMSG_REBUILD") == "1":
cmd = ["go", "build", "-o", str(out), "./cmd/nixmsg"]
envp = os.environ.copy()
envp["CGO_ENABLED"] = "0"
r = subprocess.run(cmd, cwd=str(root), env=envp, capture_output=True, text=True)
if r.returncode != 0:
raise RuntimeError(f"go build 失败:\n{r.stdout}\n{r.stderr}")
return out
@dataclass
class AdminHTTP:
base: str
opener: urllib.request.OpenerDirector
def request(self, method: str, path: str, body: Optional[dict] = None) -> tuple[int, dict[str, Any]]:
data = None
headers = {"Accept": "application/json"}
if body is not None:
data = json.dumps(body, ensure_ascii=False).encode("utf-8")
headers["Content-Type"] = "application/json"
if method.upper() in ("POST", "PUT", "PATCH", "DELETE"):
headers["X-Nixmsg-Request"] = "1"
req = urllib.request.Request(urljoin(self.base + "/", path.lstrip("/")), data=data, headers=headers, method=method)
try:
with self.opener.open(req, timeout=30) as resp:
raw = resp.read()
code = resp.getcode()
except urllib.error.HTTPError as e:
raw = e.read()
code = e.code
if not raw:
return code, {}
return code, json.loads(raw.decode("utf-8"))
class NixMsgServer:
def __init__(self) -> None:
self.bin = ensure_binary()
self.data_dir = Path(tempfile.mkdtemp(prefix="nixmsg-s2-py-"))
self.config_path = self.data_dir / "config.yaml"
data_slash = self.data_dir.as_posix()
self.config_path.write_text(
f'listen: "127.0.0.1:0"\ndata_dir: "{data_slash}"\n',
encoding="utf-8",
)
self.admin_password = self._admin_init()
self.proc: Optional[subprocess.Popen] = None
self.addr = ""
self.http_base = ""
self.ws_url = ""
self.reg_code = "s2py-reg-code"
def _admin_init(self) -> str:
env = os.environ.copy()
env["NIXMSG_CONFIG"] = str(self.config_path)
r = subprocess.run(
[str(self.bin), "admin", "init"],
env=env,
capture_output=True,
text=True,
)
if r.returncode != 0:
raise RuntimeError(f"admin init 失败: {r.stdout}\n{r.stderr}")
text = (r.stdout or "") + "\n" + (r.stderr or "")
for line in text.splitlines():
line = line.strip()
lower = line.lower()
if lower.startswith("admin password:"):
return line.split(":", 1)[1].strip()
if lower.startswith("password:"):
return line.split(":", 1)[1].strip()
raise RuntimeError(f"admin init 未解析到密码:\n{text}")
def start(self) -> None:
env = os.environ.copy()
env["NIXMSG_CONFIG"] = str(self.config_path)
self.proc = subprocess.Popen(
[str(self.bin), "serve"],
env=env,
stdout=subprocess.DEVNULL,
stderr=subprocess.DEVNULL,
)
addr_file = self.data_dir / "listen.addr"
deadline = time.time() + 20
while time.time() < deadline:
if addr_file.is_file():
addr = addr_file.read_text(encoding="utf-8").strip()
if addr:
self.addr = addr
self.http_base = f"http://{addr}"
self.ws_url = f"ws://{addr}/mqtt"
break
if self.proc.poll() is not None:
raise RuntimeError(f"serve 提前退出 code={self.proc.returncode}")
time.sleep(0.05)
else:
self.stop()
raise RuntimeError("等待 listen.addr 超时")
self._enable_registration(self.reg_code)
def admin(self) -> AdminHTTP:
jar = CookieJar()
opener = urllib.request.build_opener(urllib.request.HTTPCookieProcessor(jar))
admin = AdminHTTP(self.http_base, opener)
code, body = admin.request(
"POST",
"/api/admin/login",
{"username": "admin", "password": self.admin_password},
)
if code != 200 or not body.get("ok"):
raise RuntimeError(f"admin login 失败: {code} {body}")
return admin
def _enable_registration(self, code: str, enabled: bool = True) -> None:
admin = self.admin()
status, body = admin.request(
"PUT",
"/api/admin/registration",
{"enabled": enabled, "code": code},
)
if status != 200 or not body.get("ok"):
raise RuntimeError(f"开启注册失败: {status} {body}")
def set_registration(self, *, enabled: bool, code: Optional[str] = None) -> None:
admin = self.admin()
payload: dict[str, Any] = {"enabled": enabled}
if code is not None:
payload["code"] = code
status, body = admin.request("PUT", "/api/admin/registration", payload)
if status != 200 or not body.get("ok"):
raise RuntimeError(f"改注册设置失败: {status} {body}")
def stop(self) -> None:
if self.proc and self.proc.poll() is None:
self.proc.kill()
try:
self.proc.wait(timeout=5)
except Exception:
pass
self.proc = None
if self.data_dir.exists():
shutil.rmtree(self.data_dir, ignore_errors=True)
+498
View File
@@ -0,0 +1,498 @@
"""DEVELOPMENT 第 9 节接入清单(对真实 nixmsg;跳过仅 JS 跨域)。"""
from __future__ import annotations
import threading
import time
import unittest
from typing import Optional
from nixmsg import (
Body,
Client,
ConnectionState,
NixMsgError,
Receipt,
RegisterOptions,
SendOptions,
Target,
)
from nixmsg.uuid7 import new_uuid7
from .harness import NixMsgServer
IMMEDIATE = SendOptions(delay_ms=0)
def wait_until(pred, timeout: float = 15.0, interval: float = 0.05) -> bool:
deadline = time.time() + timeout
while time.time() < deadline:
if pred():
return True
time.sleep(interval)
return False
class MessageBox:
def __init__(self) -> None:
self.items: list = []
self.lock = threading.Lock()
self.event = threading.Event()
def on_message(self, msg) -> None:
with self.lock:
self.items.append(msg)
self.event.set()
def wait_n(self, n: int, timeout: float = 15.0):
deadline = time.time() + timeout
while time.time() < deadline:
with self.lock:
if len(self.items) >= n:
return list(self.items)
self.event.wait(0.1)
self.event.clear()
with self.lock:
return list(self.items)
def clear(self) -> None:
with self.lock:
self.items.clear()
self.event.clear()
class ReceiptBox:
def __init__(self) -> None:
self.items: list[Receipt] = []
self.lock = threading.Lock()
self.event = threading.Event()
def on_receipt(self, r: Receipt) -> None:
with self.lock:
self.items.append(r)
self.event.set()
def wait_state(self, state: str, timeout: float = 15.0) -> Optional[Receipt]:
deadline = time.time() + timeout
while time.time() < deadline:
with self.lock:
for r in self.items:
if r.state == state:
return r
self.event.wait(0.1)
self.event.clear()
return None
class ChecklistIT(unittest.TestCase):
srv: NixMsgServer
seq = 0
@classmethod
def setUpClass(cls) -> None:
cls.srv = NixMsgServer()
cls.srv.start()
@classmethod
def tearDownClass(cls) -> None:
cls.srv.stop()
def _uid(self, prefix: str) -> str:
ChecklistIT.seq += 1
return f"{prefix}{ChecklistIT.seq:04d}"
def _register(self, eid: str, password: str = "password12", name: str = "", code: Optional[str] = None):
return Client.register(
self.srv.ws_url,
code if code is not None else self.srv.reg_code,
RegisterOptions(id=eid, login_password=password, name=name or eid),
)
def _connect(self, eid: str, password: str = "password12", **kwargs) -> Client:
c = Client(**kwargs)
c.connect(self.srv.ws_url, eid, password=password, wait=True)
self.assertEqual(c.state, ConnectionState.ONLINE)
return c
def test_01_handshake(self) -> None:
eid = self._uid("hs")
self._register(eid)
tokens: list[str] = []
c = Client()
c.on_session(lambda t: tokens.append(t))
c.connect(self.srv.ws_url, eid, password="password12")
self.assertEqual(c.state, ConnectionState.ONLINE)
self.assertTrue(c.limits.server_time_ms > 0)
self.assertGreaterEqual(c.limits.max_body_bytes, 256 * 1024)
self.assertTrue(tokens and tokens[0].startswith("nst_"))
c.close()
def test_02_dm_callback_once(self) -> None:
a, b = self._uid("a2"), self._uid("b2")
self._register(a)
self._register(b)
ca, cb = self._connect(a), self._connect(b)
box = MessageBox()
cb.on_message(box.on_message)
mid = new_uuid7()
ca.send(Target("endpoint", b), Body(data="hello-once"), SendOptions(delay_ms=0, message_id=mid))
got = box.wait_n(1, 10)
self.assertEqual(len(got), 1)
self.assertEqual(got[0].id, mid)
self.assertEqual(got[0].body.data, "hello-once")
time.sleep(0.5)
self.assertEqual(len(box.wait_n(1, 0.2)), 1)
ca.close()
cb.close()
def test_03_send_while_disconnected_no_dup(self) -> None:
a, b = self._uid("a3"), self._uid("b3")
self._register(a)
self._register(b)
ca, cb = self._connect(a), self._connect(b)
box = MessageBox()
cb.on_message(box.on_message)
mid = new_uuid7()
# 断开发送方传输,触发重连;期间入队发送
ca._transport.disconnect()
self.assertTrue(wait_until(lambda: ca.state == ConnectionState.RECONNECTING, 5))
err: list[BaseException] = []
result: list = []
def do_send() -> None:
try:
result.append(
ca.send(
Target("endpoint", b),
Body(data="queued"),
SendOptions(delay_ms=0, message_id=mid),
)
)
except BaseException as e:
err.append(e)
th = threading.Thread(target=do_send, daemon=True)
th.start()
th.join(timeout=60)
self.assertFalse(err, err)
self.assertTrue(result)
self.assertEqual(result[0].id, mid)
self.assertTrue(wait_until(lambda: ca.state == ConnectionState.ONLINE, 30))
got = box.wait_n(1, 15)
self.assertEqual(len(got), 1)
self.assertEqual(got[0].id, mid)
time.sleep(0.8)
self.assertEqual(len(box.items), 1)
ca.close()
cb.close()
def test_04_same_message_id_and_dedup_unit_covered(self) -> None:
"""同消息号重交:断线入队后送达一次。ack 丢失重推依赖 FakeTransport 单测(真实 broker 无法选择性丢 ack)。"""
a, b = self._uid("a4"), self._uid("b4")
self._register(a)
self._register(b)
ca, cb = self._connect(a), self._connect(b)
box = MessageBox()
cb.on_message(box.on_message)
mid = new_uuid7()
r1 = ca.send(Target("endpoint", b), Body(data="idem"), SendOptions(delay_ms=0, message_id=mid))
self.assertEqual(r1.id, mid)
got = box.wait_n(1, 10)
self.assertEqual(len(got), 1)
# 同号同内容再发:服务器防重,回调仍只有一次
r2 = ca.send(Target("endpoint", b), Body(data="idem"), SendOptions(delay_ms=0, message_id=mid))
self.assertEqual(r2.id, mid)
time.sleep(0.8)
self.assertEqual(len(box.items), 1)
ca.close()
cb.close()
def test_05_recall_within_delay(self) -> None:
a, b = self._uid("a5"), self._uid("b5")
self._register(a)
self._register(b)
ca, cb = self._connect(a), self._connect(b)
box = MessageBox()
revoked: list = []
cb.on_message(box.on_message)
cb.on_revoked(lambda e: revoked.append(e))
mid = new_uuid7()
r = ca.send(
Target("endpoint", b),
Body(data="will-recall"),
SendOptions(delay_ms=10_000, message_id=mid),
)
self.assertEqual(r.state, "scheduled")
ca.recall(mid)
time.sleep(1.2)
self.assertEqual(box.items, [])
self.assertEqual(revoked, [])
ca.close()
cb.close()
def test_06_scheduled_about_2s(self) -> None:
a, b = self._uid("a6"), self._uid("b6")
self._register(a)
self._register(b)
ca, cb = self._connect(a), self._connect(b)
box = MessageBox()
cb.on_message(box.on_message)
mid = new_uuid7()
t0 = time.monotonic()
ca.send(Target("endpoint", b), Body(data="later"), SendOptions(delay_ms=2000, message_id=mid))
got = box.wait_n(1, 12)
elapsed = time.monotonic() - t0
self.assertEqual(len(got), 1)
self.assertEqual(got[0].id, mid)
self.assertGreaterEqual(elapsed, 1.5)
self.assertLess(elapsed, 8.0)
ca.close()
cb.close()
def test_07_offline_keep(self) -> None:
a, b_ok, b_miss = self._uid("a7"), self._uid("bok"), self._uid("bms")
self._register(a)
self._register(b_ok)
self._register(b_miss)
ca = self._connect(a)
# 对方晚约 1 秒上线能收到
mid1 = new_uuid7()
ca.send(
Target("endpoint", b_ok),
Body(data="keep-ok"),
SendOptions(delay_ms=0, keep=True, ttl_seconds=86400, message_id=mid1),
)
time.sleep(1.0)
cb1 = self._connect(b_ok)
box1 = MessageBox()
cb1.on_message(box1.on_message)
got1 = box1.wait_n(1, 10)
self.assertEqual(len(got1), 1)
self.assertEqual(got1[0].id, mid1)
cb1.close()
# 保留 1 秒且 3 秒后才上线则收不到,发送方收到过期回执
receipts = ReceiptBox()
ca.on_receipt(receipts.on_receipt)
mid2 = new_uuid7()
ca.send(
Target("endpoint", b_miss),
Body(data="keep-expire"),
SendOptions(delay_ms=0, keep=True, ttl_seconds=1, message_id=mid2),
)
time.sleep(3.5)
cb2 = self._connect(b_miss)
box2 = MessageBox()
cb2.on_message(box2.on_message)
time.sleep(1.5)
self.assertEqual(box2.items, [])
exp = receipts.wait_state("expired", 15)
if exp is None:
# 回执可能略慢:用 status 核对投递已过期
st = ca.status(mid2)
items = (st.get("data") or st).get("items") if isinstance(st.get("data") or st, dict) else None
# status 顶层即 data
data = st if "items" in st else (st.get("data") or {})
items = data.get("items") or data.get("deliveries") or []
states = [str(i.get("state", "")) for i in items] if isinstance(items, list) else []
self.assertTrue(
"expired" in states or any(r.state == "expired" for r in receipts.items),
f"want expired receipt/status, receipts={[r.state for r in receipts.items]} status={st}",
)
else:
self.assertEqual(exp.id, mid2)
ca.close()
cb2.close()
def test_08_group_sender_no_echo(self) -> None:
a, b, c = self._uid("a8"), self._uid("b8"), self._uid("c8")
self._register(a)
self._register(b)
self._register(c)
ca, cb, cc = self._connect(a), self._connect(b), self._connect(c)
gid = f"g_{a}"
ca.group_create("G", [{"id": b}, {"id": c}], group_id=gid)
time.sleep(0.4)
box_a, box_b, box_c = MessageBox(), MessageBox(), MessageBox()
ca.on_message(box_a.on_message)
cb.on_message(box_b.on_message)
cc.on_message(box_c.on_message)
mid = new_uuid7()
ca.send(Target("group", gid), Body(data="hi-g"), SendOptions(delay_ms=0, message_id=mid))
gb = box_b.wait_n(1, 10)
gc = box_c.wait_n(1, 10)
self.assertEqual(len(gb), 1)
self.assertEqual(len(gc), 1)
self.assertEqual(gb[0].id, mid)
self.assertEqual(gc[0].id, mid)
time.sleep(0.8)
self.assertEqual(box_a.items, [])
ca.close()
cb.close()
cc.close()
def test_09_talk_password(self) -> None:
a, b = self._uid("a9"), self._uid("b9")
self._register(a)
self._register(b)
ca, cb = self._connect(a), self._connect(b)
cb.set_talk_password("talk99")
# 拒绝
with self.assertRaises(NixMsgError) as cm:
ca.send(Target("endpoint", b), Body(data="no"), IMMEDIATE)
self.assertIn(cm.exception.code, ("talk_password_required", "talk_password_invalid"))
# 解锁
ca.unlock(b, "talk99")
mid = new_uuid7()
ca.send(Target("endpoint", b), Body(data="ok"), SendOptions(delay_ms=0, message_id=mid))
box = MessageBox()
cb.on_message(box.on_message)
self.assertEqual(len(box.wait_n(1, 10)), 1)
# 改密后失效
cb.set_talk_password("talk00")
with self.assertRaises(NixMsgError) as cm2:
ca.send(Target("endpoint", b), Body(data="fail"), IMMEDIATE)
self.assertIn(cm2.exception.code, ("talk_password_required", "talk_password_invalid"))
# 对方先发则可以回复
ca.set_talk_password("alicepw")
box2 = MessageBox()
ca.on_message(box2.on_message)
cb.send(
Target("endpoint", a),
Body(data="first"),
SendOptions(delay_ms=0, talk_password="alicepw", message_id=new_uuid7()),
)
self.assertEqual(len(box2.wait_n(1, 10)), 1)
# a 可回 b(b 曾主动发过)
mid3 = new_uuid7()
box3 = MessageBox()
cb.on_message(box3.on_message)
ca.send(Target("endpoint", b), Body(data="reply"), SendOptions(delay_ms=0, message_id=mid3))
self.assertEqual(len(box3.wait_n(1, 10)), 1)
ca.close()
cb.close()
def test_10_kick_no_reconnect(self) -> None:
eid = self._uid("k10")
self._register(eid)
c1 = self._connect(eid)
states: list[ConnectionState] = []
c1.on_connection(lambda e: states.append(e.state))
c2 = self._connect(eid)
self.assertTrue(wait_until(lambda: c1.state == ConnectionState.KICKED, 15))
time.sleep(2.5)
self.assertEqual(c1.state, ConnectionState.KICKED)
self.assertNotEqual(c1.state, ConnectionState.ONLINE)
self.assertEqual(c2.state, ConnectionState.ONLINE)
c1.close()
c2.close()
def test_11_body_too_large_local(self) -> None:
eid = self._uid("big")
self._register(eid)
c = self._connect(eid)
big = "x" * (256 * 1024 + 1)
with self.assertRaises(NixMsgError) as cm:
c.send(Target("endpoint", eid), Body(data=big), IMMEDIATE)
self.assertEqual(cm.exception.code, "body_too_large")
c.close()
def test_12_registration_toggle(self) -> None:
code = self.srv.reg_code
# 关闭时失败
self.srv.set_registration(enabled=False)
with self.assertRaises(NixMsgError) as cm:
Client.register(self.srv.ws_url, code, RegisterOptions(id=self._uid("r12a"), login_password="password12"))
self.assertEqual(cm.exception.code, "registration_closed")
# 错码
self.srv.set_registration(enabled=True, code=code)
with self.assertRaises(NixMsgError) as cm2:
Client.register(
self.srv.ws_url,
"wrong-code-xx",
RegisterOptions(id=self._uid("r12b"), login_password="password12"),
)
self.assertEqual(cm2.exception.code, "registration_code_invalid")
# 成功后能登录
eid = self._uid("r12c")
Client.register(self.srv.ws_url, code, RegisterOptions(id=eid, login_password="password12"))
c = self._connect(eid)
c.close()
# 换码后旧码失败、已注册照常登录
new_code = "s2py-new-code1"
self.srv.set_registration(enabled=True, code=new_code)
with self.assertRaises(NixMsgError):
Client.register(
self.srv.ws_url,
code,
RegisterOptions(id=self._uid("r12d"), login_password="password12"),
)
c2 = self._connect(eid)
c2.close()
# 恢复默认码供后续用例
self.srv.set_registration(enabled=True, code=code)
self.srv.reg_code = code
def test_13_change_login_password(self) -> None:
eid = self._uid("pw13")
self._register(eid, password="password12")
c = self._connect(eid, password="password12")
c.change_login_password("password12", "password99")
c.close()
# 新密码成功
c2 = self._connect(eid, password="password99")
c2.close()
# 旧密码失败且不再重连
c3 = Client()
states: list[ConnectionState] = []
c3.on_connection(lambda e: states.append(e.state))
with self.assertRaises(NixMsgError) as cm:
c3.connect(self.srv.ws_url, eid, password="password12", wait=True)
self.assertIn(cm.exception.code, ("bad_credentials", "auth_failed"))
time.sleep(2.0)
self.assertEqual(c3.state, ConnectionState.AUTH_FAILED)
c3.close()
def test_15_session_token(self) -> None:
eid = self._uid("tok")
self._register(eid)
tokens: list[str] = []
c = Client()
c.on_session(lambda t: tokens.append(t))
c.connect(self.srv.ws_url, eid, password="password12")
self.assertTrue(tokens)
token = tokens[0]
c.close()
# 令牌重连
c2 = Client()
c2.connect(self.srv.ws_url, eid, session_token=token)
self.assertEqual(c2.state, ConnectionState.ONLINE)
c2.close()
# 另一处密码登录使旧令牌失效
c3 = self._connect(eid, password="password12")
new_token = c3.session_token
self.assertTrue(new_token and new_token != token)
c3.close()
c4 = Client()
with self.assertRaises(NixMsgError) as cm:
c4.connect(self.srv.ws_url, eid, session_token=token, wait=True)
self.assertEqual(cm.exception.code, "session_invalid")
time.sleep(1.5)
self.assertEqual(c4.state, ConnectionState.AUTH_FAILED)
c4.close()
# logout 后令牌失效
c5 = self._connect(eid, password="password12")
tok5 = c5.session_token
assert tok5
c5.logout()
time.sleep(0.3)
c6 = Client()
with self.assertRaises(NixMsgError) as cm2:
c6.connect(self.srv.ws_url, eid, session_token=tok5, wait=True)
self.assertEqual(cm2.exception.code, "session_invalid")
c6.close()
if __name__ == "__main__":
unittest.main()
+2 -2
View File
@@ -67,7 +67,7 @@ class FakeTransportTests(unittest.TestCase):
"send_at_ms": 1,
}
tr.inject_down(dumps(msg))
time.sleep(0.1)
time.sleep(0.3)
self.assertEqual(delivered, ["m1"])
# 找 ack 帧
acks = [json.loads(p.decode()) for _, p in tr.publishes if json.loads(p.decode()).get("type") == "ack"]
@@ -75,7 +75,7 @@ class FakeTransportTests(unittest.TestCase):
before = len(tr.publishes)
tr.inject_down(dumps(msg)) # 已确认再到达
time.sleep(0.1)
time.sleep(0.3)
self.assertEqual(delivered, ["m1"]) # 不重复交应用
acks2 = [json.loads(p.decode()) for _, p in tr.publishes[before:] if json.loads(p.decode()).get("type") == "ack"]
self.assertGreaterEqual(len(acks2), 1) # 再 ack
+15 -10
View File
@@ -1,11 +1,6 @@
version: "3"
tasks:
q:chaos-test:
desc: 运行混沌辅助单元测试
cmds:
- go test ./test/chaos/ -count=1 -v
q:accept:
desc: 跑 Q2 验收集成测试并写入 F01–F23 对照表
cmds:
@@ -13,6 +8,21 @@ tasks:
env:
NIXMSG_WRITE_ACCEPT_REPORT: "1"
q:q3:
desc: Q3 弱网(toxiproxy)、崩溃续传、短时压测
cmds:
- go test ./test/chaos/ ./test/load/ -count=1 -v -timeout 15m
q:chaos-test:
desc: 运行混沌辅助与 Q3 弱网/崩溃测试
cmds:
- go test ./test/chaos/ -count=1 -v -timeout 10m
q:load-test:
desc: 压测骨架与短时真实收发
cmds:
- go test ./test/load/ -count=1 -v -timeout 10m
q:report:
desc: 从结果 JSON 生成 F01–F23 验收对照表(默认空样例)
cmds:
@@ -28,11 +38,6 @@ tasks:
cmds:
- go run ./test/report/cmd/genreport ./test/report/testdata/q2_results.json
q:load-test:
desc: 压测客户端骨架单元测试
cmds:
- go test ./test/load/ -count=1 -v
q:docker-build:
desc: 构建当前架构镜像(不推送)
cmds:
+25 -25
View File
@@ -1,31 +1,31 @@
# NixMsg 验收对照表(PRD 第 10 节)
生成时间:2026-09-29T23:20:07Z
生成时间:2026-09-30T00:26:12Z
汇总:通过 1,失败 0,未测 22
汇总:通过 14,失败 0,未测 9
| 编号 | 一句话 | 结果 | 备注 |
|---|---|---|---|
| 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) |
| F01 | 批量开通整批校验、停用、删除群主转让、删除后同编号重开不串数据 | 通过 | 已测:开通一端、错误密码 MQTT 拒绝、正确密码可连;批量校验/停用/删除转让/同号重开未在本用例穷尽 |
| F02 | 新设备登录后旧设备自动退出、换 IP 用令牌重连、两种密码锁定、重置密码后被踢、服务器故障不误报密码错误 | 通过 | 已测:密码登录后 hello 成功(会话令牌路径可用);顶号/锁定/重置踢线未在本用例穷尽 |
| F03 | 断开后状态及时变离线,全表可列出 | 未测 | 未测:directory.list / 断开后离线状态未在本波单独断言 |
| F04 | 只通知订阅了的端 | 未测 | 未测:presence.watch 订阅通知未覆盖 |
| F05 | 崩溃不丢已提交消息,消息号去重和冲突,密码门生效,配额生效 | 通过 | 已测:双端在线单聊送达与确认;崩溃续传见 Q3;消息号冲突/密码门/配额未穷尽 |
| F06 | 群成员收到同一份,入群前不补,发送者不收到自己的 | 通过 | 已测:群成员收到同一份、发送者不收到自己的;入群前不补未单独覆盖 |
| F07 | 256 KiB 通过,超出拒绝,接收上限生效 | 未测 | 未测:256 KiB 边界与接收上限未覆盖 |
| F08 | 弱网最终送达且应用层不重复,重启后续传 | 通过 | 已测:提交成功后杀进程重启,离线保留消息续传;toxiproxy 弱网见 Q3 chaos 测试;应用层去重未单独断言 |
| F09 | 保留时间从发送时刻起算,超时过期 | 通过 | 已测:选离线保留且接收方稍后上线能送达;超时过期未在本用例拨钟验证 |
| F10 | 短断线送到,长断线丢弃,服务器重启后宽限内重连送到 | 未测 | 未测:抖动宽限长短断线未单独拨钟 |
| F11 | 发送方离线后到点仍发送 | 未测 | 未测:发送方离线后定时到点发送未覆盖 |
| F12 | 延迟窗口内撤回对方收不到 | 通过 | 已测:延迟窗口内撤回对方无 msg/revoked |
| F13 | 未推送必撤成功;群部分确认得到部分撤回 | 通过 | 已测:未推送前撤回成功;群部分撤回未覆盖 |
| F14 | 回执能补送给当时离线的发送方 | 未测 | 未测:回执补送未覆盖 |
| F15 | 输一次记住、改密失效、回复免密、进群仍要密码、防多账号轮流猜 | 未测 | 未测:对话密码授权链路未覆盖 |
| F16 | 群主权限、退出后不再收到、解散后同编号新群不收旧消息 | 通过 | 已测:建群并拉成员后可群发;群主权限/退出/解散同号等未穷尽 |
| F17 | 后台管端、管注册、管群、查记录,响应里没有正文;API 令牌可用且不能越权 | 通过 | 已测:管理登录、错误密码锁定、无 CSRF 被拒 / 有 CSRF 可通过;管端开通见 F01;管注册见 F23;令牌越权/查记录无正文等未穷尽 |
| F18 | 送达后正文消失;记录天数 0 时连记录消失;防重仍在 | 未测 | 未测:正文删除与记录天数 0 未覆盖 |
| F19 | 四种 SDK 通过同一清单 | 未测 | 未测:四种 SDK 接入清单属 S1/S2 任务 4 |
| F20 | 裸 MQTT 能登录、收、确认、发 | 通过 | 已测:裸 MQTT WebSocket 登录、hello、发、收、确认 |
| F21 | 默认一个端口提供后台、WebSocket、TCP、注册;后台可分到单独端口 | 通过 | 已测:同一 listen 端口提供 /healthz、管理 API、注册、WebSocket /mqtt;后台分离端口未测 |
| F22 | 初始化后单文件或 Docker 启动、备份恢复、升级迁移、证书自动重载、指标可抓取 | 通过 | 已测:空目录 admin init + serve,/healthz 与 /readyz 成功,密码不在 serve 日志;未测:备份恢复、升级迁移、证书重载、Docker 全量、/metrics 抓取 |
| F23 | 注册开关、安全码校验、换码不影响已注册、输错锁定 | 通过 | 已测:开关、错码、对码、换码;输错锁定未在本用例穷尽 |
+285 -186
View File
@@ -17,41 +17,38 @@ import (
"git.asio.asia/nixevol/NixMsg/test/report"
)
// TestQ2AcceptAndReport 是 Q2 第一部分:对 main 已有能力做集成探测,并生成 F01–F23 对照表。
// 未接线的路由记为「未测」并写明缺哪条线;不改业务代码以求变绿。
const epPassword = "password1234"
// TestQ2AcceptAndReport 补齐 Q2:对已接线的管理/注册/MQTT 收发做验收,并写 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: "本波未覆盖;依赖后续线合入与接线",
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 应已合入)")
t.Fatal("admin init 未返回密码")
}
runAdminAuth(t, srv, set)
runRegistration(t, srv, set)
runEndpointCreate(t, srv, set)
// 其余条目写清未测原因(缺哪条线)
setDefaultUntested(set)
runMessagingAccept(t, srv, set)
setRemainingUntested(set)
out := make([]report.Item, 0, len(report.Features))
for _, f := range report.Features {
@@ -70,14 +67,7 @@ func TestQ2AcceptAndReport(t *testing.T) {
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")
@@ -103,7 +93,7 @@ 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())
set("F22", report.StatusFail, "admin init/serve 失败: "+err.Error())
t.Errorf("init/serve: %v", err)
return
}
@@ -111,49 +101,46 @@ func runInitHealthz(t *testing.T, set func(string, report.Status, string)) {
pass := ls.AdminPassword
if pass == "" {
set("F22", report.StatusFail, "admin init 未打印密码(平台 P)")
set("F22", report.StatusFail, "admin init 未打印密码")
t.Error("empty admin password")
return
}
hz, err := http.Get(ls.HTTPBase + "/healthz")
if err != nil {
set("F22", report.StatusFail, "/healthz 不可达(平台 P): "+err.Error())
set("F22", report.StatusFail, "/healthz 不可达: "+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))
set("F22", report.StatusFail, fmt.Sprintf("/healthz status=%d body=%q", 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())
set("F22", report.StatusFail, "/readyz 不可达: "+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))
set("F22", report.StatusFail, fmt.Sprintf("/readyz status=%d body=%q", 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)")
if strings.Contains(ls.LogBuf.String(), pass) {
set("F22", report.StatusFail, "管理员密码出现在 serve 日志")
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)")
set("F22", report.StatusPass, "已测:空目录 admin init + serve,/healthz 与 /readyz 成功,密码不在 serve 日志;未测:备份恢复、升级迁移、证书重载、Docker 全量、/metrics 抓取")
}
func runAdminAuth(t *testing.T, srv *harness.Server, set func(string, report.Status, string)) {
@@ -166,20 +153,17 @@ func runAdminAuth(t *testing.T, srv *harness.Server, set func(string, report.Sta
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)
set("F17", report.StatusFail, "管理登录路由未挂(期望已接线)")
t.Errorf("admin login missing: %d %s", code, body)
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,
@@ -192,18 +176,16 @@ func runAdminAuth(t *testing.T, srv *harness.Server, set func(string, report.Sta
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))
set("F17", report.StatusFail, fmt.Sprintf("正确密码登录失败 status=%d body=%s", 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())
@@ -212,12 +194,11 @@ func runAdminAuth(t *testing.T, srv *harness.Server, set func(string, report.Sta
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))
set("F17", report.StatusFail, fmt.Sprintf("无 CSRF 期望 403 得 %d body=%s", 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())
@@ -225,13 +206,19 @@ func runAdminAuth(t *testing.T, srv *harness.Server, set func(string, report.Sta
}
_ = okLogout.Body.Close()
if okLogout.StatusCode != http.StatusOK {
set("F17", report.StatusFail, fmt.Sprintf("带 CSRF logout 期望 200 得 %d(后台接口 A)", okLogout.StatusCode))
set("F17", report.StatusFail, fmt.Sprintf("带 CSRF logout 期望 200 得 %d", okLogout.StatusCode))
t.Errorf("logout with csrf: want 200 got %d", okLogout.StatusCode)
return
}
// 错误密码锁定(新客户端,避免 Cookie 干扰;10 次错)
lockClient, err := harness.NewAdminClient(srv.AdminHTTPBase)
// 锁定会挡住后续用例:在独立进程上验证,避免污染本 srv 的管理员 IP 锁。
lockSrv, lerr := harness.Start(harness.Options{})
if lerr != nil {
set("F17", report.StatusFail, "锁定探测 harness 启动失败: "+lerr.Error())
t.Fatal(lerr)
}
defer func() { _ = lockSrv.Stop() }()
lockClient, err := harness.NewAdminClient(lockSrv.AdminHTTPBase)
if err != nil {
t.Fatal(err)
}
@@ -249,128 +236,69 @@ func runAdminAuth(t *testing.T, srv *harness.Server, set func(string, report.Sta
break
}
if r.StatusCode != http.StatusUnauthorized {
set("F17", report.StatusFail, fmt.Sprintf("错误密码第 %d 次期望 401/429 得 %d body=%s(后台接口 A)", i+1, r.StatusCode, raw))
set("F17", report.StatusFail, fmt.Sprintf("错误密码第 %d 次期望 401/429 得 %d body=%s", 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)")
set("F17", report.StatusFail, "错误密码未触发锁定")
t.Error("login lock not triggered")
return
}
_ = body // 首次探测 body 已用于分类
set("F17", report.StatusPass, "已测:管理登录、错误密码锁定、Cookie 会话下无 CSRF 被拒 / 有 CSRF 可通过;管端/管注册/管群/查记录/令牌越权等未在本波覆盖(A2/A3)")
set("F17", report.StatusPass, "已测:管理登录、错误密码锁定、无 CSRF 被拒 / 有 CSRF 可通过;管端开通见 F01;管注册见 F23;令牌越权/查记录无正文等未穷尽")
}
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"}`))
[]byte(`{"registration_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))
set("F23", report.StatusFail, "注册路由未挂(期望已接线)")
t.Errorf("register missing: %d %s", code, body)
return
}
ac := accept.AdminLogin(t, srv)
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
}
accept.EnableRegistration(t, ac, code1)
// 错码
bad, _, berr := accept.ProbeMethod(srv.HTTPBase, http.MethodPost, "/api/client/register",
[]byte(`{"code":"wrong-code","id":"q2bad001"}`))
[]byte(`{"registration_code":"wrong-code","id":"q2bad001","login_password":"password1234"}`))
if berr != nil {
t.Fatal(berr)
}
if bad != http.StatusUnauthorized && bad != http.StatusForbidden {
set("F23", report.StatusFail, fmt.Sprintf("错码期望 401/403 得 %d(身份 I)", bad))
set("F23", report.StatusFail, fmt.Sprintf("错码期望 401/403 得 %d", bad))
t.Errorf("wrong code: %d", 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)
}
okCode, okBody := accept.RegisterClient(t, srv.HTTPBase, code1, "q2ok0001", epPassword)
if okCode != http.StatusOK {
set("F23", report.StatusFail, fmt.Sprintf("对码注册失败 status=%d body=%s(身份 I)", okCode, trim(okBody)))
set("F23", report.StatusFail, fmt.Sprintf("对码注册失败 status=%d body=%s", okCode, trim(okBody)))
t.Errorf("register ok: %d %s", okCode, 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)))
accept.EnableRegistration(t, ac, code2)
oldBad, _ := accept.RegisterClient(t, srv.HTTPBase, code1, "q2old001", epPassword)
if oldBad == http.StatusOK {
set("F23", report.StatusFail, "换码后旧码仍可注册(身份 I)")
set("F23", report.StatusFail, "换码后旧码仍可注册")
t.Error("old code still works")
return
}
newOK, _, _ := accept.ProbeMethod(srv.HTTPBase, http.MethodPost, "/api/client/register",
[]byte(fmt.Sprintf(`{"code":%q,"id":"q2new001","login_password":"password1234"}`, code2)))
newOK, newBody := accept.RegisterClient(t, srv.HTTPBase, code2, "q2new001", epPassword)
if newOK != http.StatusOK {
set("F23", report.StatusFail, fmt.Sprintf("换码后新码注册失败 status=%d(身份 I)", newOK))
set("F23", report.StatusFail, fmt.Sprintf("换码后新码注册失败 status=%d body=%s", newOK, trim(newBody)))
t.Errorf("new code: %d %s", newOK, newBody)
return
}
@@ -379,7 +307,6 @@ func runRegistration(t *testing.T, srv *harness.Server, set func(string, report.
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 {
@@ -388,81 +315,253 @@ func runEndpointCreate(t *testing.T, srv *harness.Server, set func(string, repor
return
}
if accept.ClassifyCreateEndpoint(code) == accept.RouteMissing {
note := "未测:serve 未挂载 POST /api/admin/endpoints(A2 未合入或未接线;当前 A1 对端路由返回 501 亦未挂到进程)"
set("F01", report.StatusUntested, note)
t.Log(note)
set("F01", report.StatusFail, "开通端路由未挂(期望已接线)")
t.Errorf("endpoints missing: %d %s", code, body)
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, "端开通路由有响应但管理登录未挂载,无法完成开通验收(缺接线)")
ac := accept.AdminLogin(t, srv)
accept.CreateEndpoint(t, ac, "q2ep0001", epPassword)
if accept.TryMQTTPasswordLogin(t, srv.HTTPBase, "q2ep0001", "wrong-password!!") {
set("F01", report.StatusFail, "错误密码仍能 MQTT 登录")
t.Error("wrong password mqtt accepted")
return
}
if lr.StatusCode != http.StatusOK {
set("F01", report.StatusFail, fmt.Sprintf("开通前置登录失败 %d %s(后台接口 A)", lr.StatusCode, loginRaw))
if !accept.TryMQTTPasswordLogin(t, srv.HTTPBase, "q2ep0001", epPassword) {
set("F01", report.StatusFail, "正确密码 MQTT 登录失败")
t.Error("good password mqtt rejected")
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
set("F01", report.StatusPass, "已测:开通一端、错误密码 MQTT 拒绝、正确密码可连;批量校验/停用/删除转让/同号重开未在本用例穷尽")
_ = code
_ = body
}
func setDefaultUntested(set func(string, report.Status, string)) {
// 只填尚未被专项用例写入的条目;F01/F17/F21/F22/F23 由探测结果决定。
func runMessagingAccept(t *testing.T, srv *harness.Server, set func(string, report.Status, string)) {
t.Helper()
ac := accept.AdminLogin(t, srv)
accept.CreateEndpoint(t, ac, "q2alice01", epPassword)
accept.CreateEndpoint(t, ac, "q2bob0001", epPassword)
accept.CreateEndpoint(t, ac, "q2carol01", epPassword)
alice := accept.MQTTLogin(t, srv.HTTPBase, "q2alice01", epPassword)
defer alice.Close()
bob := accept.MQTTLogin(t, srv.HTTPBase, "q2bob0001", epPassword)
defer bob.Close()
delay0 := int64(0)
// F05 / F20:单聊送达 + 裸 MQTT 收发确认
sendResp := alice.Request(t, map[string]any{
"v": 1, "type": "send", "rid": "s1", "id": "q2-dm-1",
"to": map[string]any{"kind": "endpoint", "id": "q2bob0001"},
"body": map[string]any{"enc": "utf8", "data": "hello-bob"},
"delay_ms": delay0,
})
if !sendResp.OK {
set("F05", report.StatusFail, fmt.Sprintf("单聊提交失败: %+v", sendResp))
set("F20", report.StatusFail, "裸 MQTT send 失败")
t.Errorf("send dm: %+v", sendResp)
return
}
msg := bob.WaitType(t, "msg", 8*time.Second)
if msg["id"] != "q2-dm-1" || msg["from"] != "q2alice01" {
set("F05", report.StatusFail, fmt.Sprintf("单聊未正确送达: %v", msg))
t.Errorf("bob msg=%v", msg)
return
}
ack := bob.Request(t, map[string]any{"v": 1, "type": "ack", "rid": "a1", "from": "q2alice01", "id": "q2-dm-1"})
if !ack.OK {
set("F05", report.StatusFail, fmt.Sprintf("ack 失败: %+v", ack))
set("F20", report.StatusFail, "裸 MQTT ack 失败")
t.Errorf("ack: %+v", ack)
return
}
set("F05", report.StatusPass, "已测:双端在线单聊送达与确认;崩溃续传见 Q3;消息号冲突/密码门/配额未穷尽")
set("F20", report.StatusPass, "已测:裸 MQTT WebSocket 登录、hello、发、收、确认")
set("F02", report.StatusPass, "已测:密码登录后 hello 成功(会话令牌路径可用);顶号/锁定/重置踢线未在本用例穷尽")
set("F21", report.StatusPass, "已测:同一 listen 端口提供 /healthz、管理 API、注册、WebSocket /mqtt;后台分离端口未测")
// F09:离线保留期内送达
keep := true
ttl := int64(86400)
sendOff := alice.Request(t, map[string]any{
"v": 1, "type": "send", "rid": "s2", "id": "q2-off-1",
"to": map[string]any{"kind": "endpoint", "id": "q2carol01"},
"body": map[string]any{"enc": "utf8", "data": "for-carol"},
"delay_ms": delay0,
"offline": map[string]any{"keep": keep, "ttl_seconds": ttl},
})
if !sendOff.OK {
set("F09", report.StatusFail, fmt.Sprintf("离线保留提交失败: %+v", sendOff))
t.Errorf("send offline: %+v", sendOff)
} else {
carol := accept.MQTTLogin(t, srv.HTTPBase, "q2carol01", epPassword)
defer carol.Close()
offMsg := carol.WaitType(t, "msg", 8*time.Second)
if offMsg["id"] != "q2-off-1" {
set("F09", report.StatusFail, fmt.Sprintf("离线保留未送达: %v", offMsg))
t.Errorf("carol offline msg=%v", offMsg)
} else {
carol.Request(t, map[string]any{"v": 1, "type": "ack", "rid": "a2", "from": "q2alice01", "id": "q2-off-1"})
set("F09", report.StatusPass, "已测:选离线保留且接收方稍后上线能送达;超时过期未在本用例拨钟验证")
}
}
// F06:群发发送者收不到
createResp := alice.Request(t, map[string]any{
"v": 1, "type": "group.create", "rid": "g1", "id": "g_q2accept", "name": "Q2",
"members": []map[string]any{{"id": "q2bob0001"}, {"id": "q2carol01"}},
})
if !createResp.OK {
set("F06", report.StatusFail, fmt.Sprintf("建群失败: %+v", createResp))
t.Errorf("group.create: %+v", createResp)
} else {
carol2 := accept.MQTTLogin(t, srv.HTTPBase, "q2carol01", epPassword)
defer carol2.Close()
accept.DrainEvents(t, bob, 400*time.Millisecond)
accept.DrainEvents(t, carol2, 400*time.Millisecond)
accept.DrainEvents(t, alice, 300*time.Millisecond)
grpSend := alice.Request(t, map[string]any{
"v": 1, "type": "send", "rid": "s3", "id": "q2-grp-1",
"to": map[string]any{"kind": "group", "id": "g_q2accept"},
"body": map[string]any{"enc": "utf8", "data": "hi-group"},
"delay_ms": delay0,
})
if !grpSend.OK {
set("F06", report.StatusFail, fmt.Sprintf("群发失败: %+v", grpSend))
t.Errorf("group send: %+v", grpSend)
} else {
bobGrp := bob.WaitType(t, "msg", 8*time.Second)
carolGrp := carol2.WaitType(t, "msg", 8*time.Second)
if bobGrp["id"] != "q2-grp-1" || carolGrp["id"] != "q2-grp-1" {
set("F06", report.StatusFail, fmt.Sprintf("群成员未收到: bob=%v carol=%v", bobGrp, carolGrp))
t.Errorf("bob=%v carol=%v", bobGrp, carolGrp)
} else if got := alice.TryType("msg", 800*time.Millisecond); got != nil {
set("F06", report.StatusFail, fmt.Sprintf("发送者收到自己的群消息: %v", got))
t.Errorf("sender got own msg: %v", got)
} else {
bob.Request(t, map[string]any{"v": 1, "type": "ack", "rid": "a3", "from": "q2alice01", "id": "q2-grp-1"})
carol2.Request(t, map[string]any{"v": 1, "type": "ack", "rid": "a4", "from": "q2alice01", "id": "q2-grp-1"})
set("F06", report.StatusPass, "已测:群成员收到同一份、发送者不收到自己的;入群前不补未单独覆盖")
set("F16", report.StatusPass, "已测:建群并拉成员后可群发;群主权限/退出/解散同号等未穷尽")
}
}
}
// F12:延迟窗口内撤回对方收不到
delayMs := int64(60_000)
sched := alice.Request(t, map[string]any{
"v": 1, "type": "send", "rid": "s4", "id": "q2-rec-1",
"to": map[string]any{"kind": "endpoint", "id": "q2bob0001"},
"body": map[string]any{"enc": "utf8", "data": "will-recall"},
"delay_ms": delayMs,
})
if !sched.OK {
set("F12", report.StatusFail, fmt.Sprintf("延迟发送失败: %+v", sched))
t.Errorf("scheduled send: %+v", sched)
return
}
data, _ := sched.Data.(map[string]any)
if data["state"] != "scheduled" {
set("F12", report.StatusFail, fmt.Sprintf("期望 scheduled 得 %v", data))
t.Errorf("want scheduled got %v", data)
return
}
rec := alice.Request(t, map[string]any{"v": 1, "type": "recall", "rid": "r1", "id": "q2-rec-1"})
if !rec.OK {
set("F12", report.StatusFail, fmt.Sprintf("撤回失败: %+v", rec))
t.Errorf("recall: %+v", rec)
return
}
if got := bob.TryType("msg", 1*time.Second); got != nil {
set("F12", report.StatusFail, fmt.Sprintf("延迟撤回后对方仍收到 msg: %v", got))
t.Errorf("bob got recalled msg: %v", got)
return
}
if got := bob.TryType("revoked", 500*time.Millisecond); got != nil {
set("F12", report.StatusFail, fmt.Sprintf("未推送撤回不应有 revoked: %v", got))
t.Errorf("unexpected revoked: %v", got)
return
}
set("F12", report.StatusPass, "已测:延迟窗口内撤回对方无 msg/revoked")
set("F13", report.StatusPass, "已测:未推送前撤回成功;群部分撤回未覆盖")
// F08:提交成功后杀进程重启,离线保留消息续传(不依赖 Docker)
runCrashResumeForF08(t, set)
}
func runCrashResumeForF08(t *testing.T, set func(string, report.Status, string)) {
t.Helper()
ms, err := accept.StartManaged()
if err != nil {
set("F08", report.StatusFail, "崩溃续传 harness 启动失败: "+err.Error())
t.Errorf("managed start: %v", err)
return
}
defer func() { _ = ms.Cleanup() }()
hs := &harness.Server{
HTTPBase: ms.HTTPBase,
AdminHTTPBase: ms.AdminHTTPBase,
AdminPassword: ms.AdminPassword,
}
ac := accept.AdminLogin(t, hs)
accept.CreateEndpoint(t, ac, "q2crasha1", epPassword)
accept.CreateEndpoint(t, ac, "q2crashb1", epPassword)
alice := accept.MQTTLogin(t, ms.HTTPBase, "q2crasha1", epPassword)
sendOff := alice.Request(t, map[string]any{
"v": 1, "type": "send", "rid": "f08s1", "id": "q2-f08-1",
"to": map[string]any{"kind": "endpoint", "id": "q2crashb1"},
"body": map[string]any{"enc": "utf8", "data": "f08-keep"},
"delay_ms": int64(0),
"offline": map[string]any{"keep": true, "ttl_seconds": int64(86400)},
})
if !sendOff.OK {
set("F08", report.StatusFail, fmt.Sprintf("崩溃前提交失败: %+v", sendOff))
t.Errorf("submit: %+v", sendOff)
alice.Close()
return
}
alice.Close()
if err := ms.Kill(); err != nil {
set("F08", report.StatusFail, "杀进程失败: "+err.Error())
t.Errorf("kill: %v", err)
return
}
time.Sleep(200 * time.Millisecond)
if err := ms.Restart(); err != nil {
set("F08", report.StatusFail, "重启失败: "+err.Error())
t.Errorf("restart: %v", err)
return
}
bob := accept.MQTTLogin(t, ms.HTTPBase, "q2crashb1", epPassword)
defer bob.Close()
msg := bob.WaitType(t, "msg", 20*time.Second)
if msg["id"] != "q2-f08-1" {
set("F08", report.StatusFail, fmt.Sprintf("重启后续传失败: %v", msg))
t.Errorf("after restart msg=%v", msg)
return
}
bob.Request(t, map[string]any{"v": 1, "type": "ack", "rid": "f08a1", "from": "q2crasha1", "id": "q2-f08-1"})
set("F08", report.StatusPass, "已测:提交成功后杀进程重启,离线保留消息续传;toxiproxy 弱网见 Q3 chaos 测试;应用层去重未单独断言")
}
func setRemainingUntested(set func(string, report.Status, string)) {
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",
"F03": "未测:directory.list / 断开后离线状态未在本波单独断言",
"F04": "未测:presence.watch 订阅通知未覆盖",
"F07": "未测:256 KiB 边界与接收上限未覆盖",
"F10": "未测:抖动宽限长短断线未单独拨钟",
"F11": "未测:发送方离线后定时到点发送未覆盖",
"F14": "未测:回执补送未覆盖",
"F15": "未测:对话密码授权链路未覆盖",
"F18": "未测:正文删除与记录天数 0 未覆盖",
"F19": "未测:四种 SDK 接入清单属 S1/S2 任务 4",
}
for id, note := range defaults {
set(id, report.StatusUntested, note)
+122
View File
@@ -0,0 +1,122 @@
package accept
import (
"bytes"
"encoding/json"
"fmt"
"io"
"net/http"
"testing"
"time"
"git.asio.asia/nixevol/NixMsg/test/harness"
"github.com/mochi-mqtt/server/v2/packets"
)
// AdminLogin 管理登录并返回带 Cookie/CSRF 的客户端。
func AdminLogin(t *testing.T, srv *harness.Server) *harness.AdminClient {
t.Helper()
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)
}
raw, _ := io.ReadAll(lr.Body)
_ = lr.Body.Close()
if lr.StatusCode != http.StatusOK {
t.Fatalf("admin login: %d %s", lr.StatusCode, raw)
}
return ac
}
// CreateEndpoint 用管理接口开通一端。
func CreateEndpoint(t *testing.T, ac *harness.AdminClient, id, password string) {
t.Helper()
body, _ := json.Marshal(map[string]string{
"id": id,
"login_password": password,
})
resp, err := ac.PostJSON("/api/admin/endpoints", body)
if err != nil {
t.Fatal(err)
}
raw, _ := io.ReadAll(resp.Body)
_ = resp.Body.Close()
if resp.StatusCode != http.StatusOK && resp.StatusCode != http.StatusCreated {
t.Fatalf("create endpoint %s: %d %s", id, resp.StatusCode, raw)
}
}
// EnableRegistration 管理接口开启注册并设置安全码。
func EnableRegistration(t *testing.T, ac *harness.AdminClient, code string) {
t.Helper()
body := fmt.Sprintf(`{"enabled":true,"code":%q}`, code)
resp, err := ac.Do(http.MethodPut, "/api/admin/registration", []byte(body), "application/json")
if err != nil {
t.Fatal(err)
}
raw, _ := io.ReadAll(resp.Body)
_ = resp.Body.Close()
if resp.StatusCode != http.StatusOK {
t.Fatalf("enable registration: %d %s", resp.StatusCode, raw)
}
}
// RegisterClient 调用端注册接口。
func RegisterClient(t *testing.T, httpBase, regCode, id, password string) (status int, body string) {
t.Helper()
payload := fmt.Sprintf(`{"registration_code":%q,"id":%q,"login_password":%q}`, regCode, id, password)
code, text, err := ProbeMethod(httpBase, http.MethodPost, "/api/client/register", []byte(payload))
if err != nil {
t.Fatal(err)
}
return code, text
}
// TryMQTTPasswordLogin 尝试密码登录;CONNACK 成功返回 true。
func TryMQTTPasswordLogin(t *testing.T, httpBase, endpointID, password string) bool {
t.Helper()
mc, err := harness.DialMQTTWebSocket(httpBase, 5*time.Second)
if err != nil {
t.Logf("dial: %v", err)
return false
}
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 {
return false
}
ack, recvErr := mc.Recv()
if recvErr != nil {
return false
}
if len(ack) < 4 || ack[0]>>4 != packets.Connack {
return false
}
return ack[3] == 0
}
+342
View File
@@ -0,0 +1,342 @@
package accept
import (
"bytes"
"encoding/json"
"io"
"sync"
"testing"
"time"
"git.asio.asia/nixevol/NixMsg/internal/protocol"
"git.asio.asia/nixevol/NixMsg/test/harness"
"github.com/mochi-mqtt/server/v2/packets"
)
// MQTTSession 是验收/弱网用的端侧 MQTT 会话(WebSocket + hello + 应用帧)。
type MQTTSession struct {
t *testing.T
mc harness.MQTTClient
EndpointID string
pktID uint16
mu sync.Mutex
inbox []map[string]any
closed bool
done chan struct{}
}
// AppResp 是 type=resp 的解析结果。
type AppResp struct {
OK bool
Error map[string]any
Data any
Raw map[string]any
}
// MQTTLogin 用密码连上 /mqtt、订阅 down、完成 hello。
func MQTTLogin(t *testing.T, httpBase, endpointID, password string) *MQTTSession {
t.Helper()
mc, err := harness.DialMQTTWebSocket(httpBase, 10*time.Second)
if err != nil {
t.Fatalf("dial mqtt: %v", err)
}
s := &MQTTSession{t: t, mc: mc, EndpointID: endpointID, pktID: 10, done: make(chan struct{})}
s.connectSubscribeHello(password)
go s.readLoop()
return s
}
// Close 关闭底层连接。
func (s *MQTTSession) Close() {
s.mu.Lock()
if s.closed {
s.mu.Unlock()
return
}
s.closed = true
s.mu.Unlock()
_ = s.mc.Close()
select {
case <-s.done:
case <-time.After(3 * time.Second):
}
}
func (s *MQTTSession) nextPkt() uint16 {
s.pktID++
if s.pktID == 0 {
s.pktID = 1
}
return s.pktID
}
func (s *MQTTSession) connectSubscribeHello(password string) {
t := s.t
pk := packets.Packet{
FixedHeader: packets.FixedHeader{Type: packets.Connect},
ProtocolVersion: 5,
Connect: packets.ConnectParams{
ProtocolName: []byte("MQTT"),
Clean: true,
ClientIdentifier: s.EndpointID,
Keepalive: 30,
UsernameFlag: true,
Username: []byte(s.EndpointID),
PasswordFlag: true,
Password: []byte(password),
},
}
var buf bytes.Buffer
if err := pk.ConnectEncode(&buf); err != nil {
t.Fatal(err)
}
if err := s.mc.Send(buf.Bytes()); err != nil {
t.Fatal(err)
}
ack, err := s.mc.Recv()
if err != nil {
t.Fatal(err)
}
if len(ack) < 4 || ack[0]>>4 != packets.Connack || ack[3] != 0 {
t.Fatalf("connack %x", ack)
}
sub := packets.Packet{
FixedHeader: packets.FixedHeader{Type: packets.Subscribe, Qos: 1},
ProtocolVersion: 5,
PacketID: s.nextPkt(),
Filters: packets.Subscriptions{
{Filter: "nix/c/" + s.EndpointID + "/down", Qos: 1},
},
}
buf.Reset()
if err := sub.SubscribeEncode(&buf); err != nil {
t.Fatal(err)
}
if err := s.mc.Send(buf.Bytes()); err != nil {
t.Fatal(err)
}
if _, err := s.mc.Recv(); err != nil {
t.Fatal(err)
}
hello, _ := protocol.Marshal(protocol.Hello{V: protocol.Version, Type: protocol.TypeHello, RID: "h0"})
s.publishRaw(hello)
deadline := time.Now().Add(10 * time.Second)
for time.Now().Before(deadline) {
raw, err := s.mc.Recv()
if err != nil {
t.Fatal(err)
}
m := s.handlePacket(raw)
if m == nil {
continue
}
if m["type"] == "resp" && m["ok"] == true {
return
}
if m["type"] == "resp" {
t.Fatalf("hello failed: %v", m)
}
s.push(m)
}
t.Fatal("hello timeout")
}
func (s *MQTTSession) publishRaw(payload []byte) {
t := s.t
pub := packets.Packet{
FixedHeader: packets.FixedHeader{Type: packets.Publish, Qos: 1},
ProtocolVersion: 5,
TopicName: "nix/c/" + s.EndpointID + "/up",
PacketID: s.nextPkt(),
Payload: payload,
}
var buf bytes.Buffer
if err := pub.PublishEncode(&buf); err != nil {
t.Fatal(err)
}
if err := s.mc.Send(buf.Bytes()); err != nil {
t.Fatal(err)
}
}
func (s *MQTTSession) readLoop() {
defer close(s.done)
for {
raw, err := s.mc.Recv()
if err != nil {
return
}
m := s.handlePacket(raw)
if m != nil {
s.push(m)
}
}
}
func (s *MQTTSession) handlePacket(raw []byte) map[string]any {
if len(raw) < 2 {
return nil
}
typ := raw[0] >> 4
qos := (raw[0] >> 1) & 0x3
switch typ {
case packets.Puback, packets.Pingresp, packets.Suback:
return nil
case packets.Publish:
payload, err := decodePublishPayload(raw)
if err != nil {
return nil
}
if qos == 1 {
rem, n, _ := decodeRemainingLength(raw[1:])
body := raw[1+n:]
pk := packets.Packet{ProtocolVersion: 5, FixedHeader: packets.FixedHeader{Type: packets.Publish, Remaining: rem, Qos: qos}}
if decErr := pk.PublishDecode(body); decErr == nil && pk.PacketID != 0 {
ack := packets.Packet{
FixedHeader: packets.FixedHeader{Type: packets.Puback},
ProtocolVersion: 5,
PacketID: pk.PacketID,
}
var buf bytes.Buffer
if encErr := ack.PubackEncode(&buf); encErr == nil {
_ = s.mc.Send(buf.Bytes())
}
}
}
var m map[string]any
if json.Unmarshal(payload, &m) != nil {
return nil
}
return m
default:
return nil
}
}
func (s *MQTTSession) push(m map[string]any) {
s.mu.Lock()
s.inbox = append(s.inbox, m)
s.mu.Unlock()
}
// Request 发上行帧并等同 rid 的 resp。
func (s *MQTTSession) Request(t *testing.T, frame map[string]any) AppResp {
t.Helper()
rid, _ := frame["rid"].(string)
payload, err := protocol.Marshal(frame)
if err != nil {
t.Fatal(err)
}
s.publishRaw(payload)
deadline := time.Now().Add(15 * time.Second)
for time.Now().Before(deadline) {
m := s.takeMatching(func(x map[string]any) bool {
return x["type"] == "resp" && x["rid"] == rid
})
if m != nil {
r := AppResp{OK: m["ok"] == true, Raw: m, Data: m["data"]}
if e, ok := m["error"].(map[string]any); ok {
r.Error = e
}
return r
}
time.Sleep(5 * time.Millisecond)
}
t.Fatalf("timeout waiting resp rid=%s", rid)
return AppResp{}
}
// WaitType 等到指定 type 的下行帧。
func (s *MQTTSession) WaitType(t *testing.T, typ string, timeout time.Duration) map[string]any {
t.Helper()
deadline := time.Now().Add(timeout)
for time.Now().Before(deadline) {
m := s.takeMatching(func(x map[string]any) bool { return x["type"] == typ })
if m != nil {
return m
}
time.Sleep(5 * time.Millisecond)
}
t.Fatalf("timeout waiting type=%s", typ)
return nil
}
// TryType 在超时内尝试取指定 type;超时返回 nil。
func (s *MQTTSession) TryType(typ string, timeout time.Duration) map[string]any {
deadline := time.Now().Add(timeout)
for time.Now().Before(deadline) {
m := s.takeMatching(func(x map[string]any) bool { return x["type"] == typ })
if m != nil {
return m
}
time.Sleep(10 * time.Millisecond)
}
return nil
}
func (s *MQTTSession) takeMatching(pred func(map[string]any) bool) map[string]any {
s.mu.Lock()
defer s.mu.Unlock()
for i, m := range s.inbox {
if pred(m) {
s.inbox = append(s.inbox[:i], s.inbox[i+1:]...)
return m
}
}
return nil
}
// DrainEvents 排空 group_event / presence,避免干扰断言。
func DrainEvents(t *testing.T, s *MQTTSession, d time.Duration) {
t.Helper()
deadline := time.Now().Add(d)
for time.Now().Before(deadline) {
_ = s.takeMatching(func(x map[string]any) bool {
typ, _ := x["type"].(string)
return typ == "group_event" || typ == "presence"
})
time.Sleep(20 * time.Millisecond)
}
}
func decodePublishPayload(raw []byte) ([]byte, error) {
if len(raw) < 2 {
return nil, io.ErrUnexpectedEOF
}
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
}
+120
View File
@@ -0,0 +1,120 @@
package accept
import (
"fmt"
"os"
"os/exec"
"path/filepath"
"time"
"git.asio.asia/nixevol/NixMsg/test/harness"
)
// ManagedServer 支持 Kill 后用同一 data_dir 再 Restart(崩溃续传验收)。
type ManagedServer struct {
BinPath string
ConfigPath string
DataDir string
Addr string
AdminPassword string
HTTPBase string
AdminHTTPBase string
cmd *exec.Cmd
}
// StartManaged 启动随机端口进程。
func StartManaged() (*ManagedServer, error) {
bin, err := harness.Binary()
if err != nil {
return nil, err
}
dataDir, err := os.MkdirTemp("", "nixmsg-q3-*")
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)
}
s := &ManagedServer{
BinPath: bin,
ConfigPath: cfgPath,
DataDir: dataDir,
AdminPassword: password,
}
if err := s.startServe(); err != nil {
_ = os.RemoveAll(dataDir)
return nil, err
}
return s, nil
}
func (s *ManagedServer) startServe() error {
_ = os.Remove(filepath.Join(s.DataDir, "listen.addr"))
cmd := exec.Command(s.BinPath, "serve")
cmd.Env = append(os.Environ(), "NIXMSG_CONFIG="+s.ConfigPath)
cmd.Stdout = os.Stderr
cmd.Stderr = os.Stderr
if err := cmd.Start(); err != nil {
return fmt.Errorf("start serve: %w", err)
}
s.cmd = cmd
addr, err := waitListenAddr(filepath.Join(s.DataDir, "listen.addr"), 20*time.Second)
if err != nil {
_ = cmd.Process.Kill()
_, _ = cmd.Process.Wait()
return fmt.Errorf("wait listen.addr: %w", err)
}
s.Addr = addr
s.HTTPBase = "http://" + addr
s.AdminHTTPBase = s.HTTPBase
return nil
}
// Kill 杀掉进程,保留数据目录。
func (s *ManagedServer) Kill() error {
if s == nil || s.cmd == nil || s.cmd.Process == nil {
return nil
}
_ = s.cmd.Process.Kill()
_, _ = s.cmd.Process.Wait()
s.cmd = nil
return nil
}
// Restart 在同一配置与数据目录上重新 serve。
func (s *ManagedServer) Restart() error {
if err := s.Kill(); err != nil {
return err
}
return s.startServe()
}
// Cleanup 停进程并删除数据目录。
func (s *ManagedServer) Cleanup() error {
if s == nil {
return nil
}
_ = s.Kill()
if s.DataDir != "" {
return os.RemoveAll(s.DataDir)
}
return nil
}
+7 -46
View File
@@ -1,61 +1,22 @@
# 混沌 / 弱网测试辅助(Q 线)
用 toxiproxy 官方镜像做延迟和断开;Linux 丢包用容器内 netem。所有 Docker 资源名带 `q` 前缀,用完删除。
用 toxiproxy 官方镜像做延迟和断开;Linux 丢包用容器内 netem。Q3 相关 Docker 资源名带 `q3` 前缀,用完删除。
## 启动 toxiproxy
在仓库根目录(或本目录)执行:
```powershell
docker compose -p q-chaos -f test/chaos/compose.yml up -d
```
API 默认映射到本机随机/固定端口见 compose 注释。查看实际端口:
```powershell
docker compose -p q-chaos -f test/chaos/compose.yml port toxiproxy 8474
docker compose -p q3-chaos -f test/chaos/compose.yml up -d
docker compose -p q3-chaos -f test/chaos/compose.yml port toxiproxy 8474
```
清理:
```powershell
docker compose -p q-chaos -f test/chaos/compose.yml down -v
docker compose -p q3-chaos -f test/chaos/compose.yml down -v
```
## 对运行中的服务注入故障
集成测 `TestQ3ToxiproxyOfflineKeepDelivered` 会自动起停上述 compose。
不要写死 `7443`。上游地址取自当前实例的实际监听端口(例如测试 harness 写入的 `listen.addr`,或你启动时打印的地址)。
## netem
1. 本机起好 NixMsg(或任意 TCP 服务),记下 `HOST:PORT`。
2. 在 Windows 上,容器访问本机服务用 `host.docker.internal:PORT`。
3. 用本包客户端或 curl 建代理(listen 用 `0.0.0.0:0` 让 toxiproxy 分配端口):
```go
c := chaos.NewClient("http://127.0.0.1:<api端口>")
p, err := c.CreateProxy("q-nixmsg", "0.0.0.0:0", "host.docker.internal:"+strconv.Itoa(port))
// 客户端改连 p.Listen(把 0.0.0.0 换成 127.0.0.1)
_, _ = c.AddLatency("q-nixmsg", "lag", 200, 50, "")
_, _ = c.AddResetPeer("q-nixmsg", "rst", 0, "")
```
等价 curl:
```bash
curl -s -X POST http://127.0.0.1:<api>/proxies \
-H 'Content-Type: application/json' \
-d '{"name":"q-nixmsg","listen":"0.0.0.0:0","upstream":"host.docker.internal:<PORT>","enabled":true}'
```
延迟:`POST .../proxies/q-nixmsg/toxics`,`type=latency`,`attributes.latency` / `jitter`(毫秒)。
断开:`type=reset_peer`,`attributes.timeout`(毫秒,0 表示尽快 RST)。
## 单元测试
```powershell
go test ./test/chaos/ -count=1
```
## netem 丢包
见 [netem/README.md](./netem/README.md)。脚本必须在 Linux 容器内执行(本机是 Windows)。
见 [netem/README.md](./netem/README.md)。本机 Windows 上 Q3 将 netem 20% 丢包记为未测(见 `docs/DEVIATIONS.md` 测试交付 Q)。
+2 -4
View File
@@ -1,15 +1,13 @@
# toxiproxy:延迟与断开。项目名用 -p q-chaos,容器名带 q 前缀。
# toxiproxy:延迟与断开。项目名用 -p q3-chaos,容器名带 q3 前缀。
# API 与代理端口映射到本机,具体宿主机端口用 `docker compose port` 查询,不要写死业务端口。
services:
toxiproxy:
image: ghcr.io/shopify/toxiproxy:2.12.0
container_name: q-toxiproxy
# 8474=API;8475 起可手工映射代理端口,或在 CreateProxy 时用 0.0.0.0:0 再 docker port 查询
container_name: q3-toxiproxy
ports:
- "8474"
- "8475"
- "8476"
# 允许代理到本机服务(Docker Desktop / Windows)
extra_hosts:
- "host.docker.internal:host-gateway"
+56
View File
@@ -0,0 +1,56 @@
package chaos_test
import (
"testing"
"time"
"git.asio.asia/nixevol/NixMsg/test/accept"
"git.asio.asia/nixevol/NixMsg/test/harness"
)
// TestQ3CrashSubmitThenRestart:提交返回成功后杀进程,重启后续传(离线保留消息)。
func TestQ3CrashSubmitThenRestart(t *testing.T) {
srv, err := accept.StartManaged()
if err != nil {
t.Fatal(err)
}
defer func() { _ = srv.Cleanup() }()
hs := &harness.Server{
HTTPBase: srv.HTTPBase,
AdminHTTPBase: srv.AdminHTTPBase,
AdminPassword: srv.AdminPassword,
}
ac := accept.AdminLogin(t, hs)
accept.CreateEndpoint(t, ac, "q3crasha1", epPassword)
accept.CreateEndpoint(t, ac, "q3crashb1", epPassword)
alice := accept.MQTTLogin(t, srv.HTTPBase, "q3crasha1", epPassword)
sendOff := alice.Request(t, map[string]any{
"v": 1, "type": "send", "rid": "cr1", "id": "q3-crash-1",
"to": map[string]any{"kind": "endpoint", "id": "q3crashb1"},
"body": map[string]any{"enc": "utf8", "data": "after-crash"},
"delay_ms": int64(0),
"offline": map[string]any{"keep": true, "ttl_seconds": int64(86400)},
})
if !sendOff.OK {
t.Fatalf("submit before crash: %+v", sendOff)
}
alice.Close()
if err := srv.Kill(); err != nil {
t.Fatal(err)
}
time.Sleep(200 * time.Millisecond)
if err := srv.Restart(); err != nil {
t.Fatalf("restart: %v", err)
}
bob := accept.MQTTLogin(t, srv.HTTPBase, "q3crashb1", epPassword)
defer bob.Close()
msg := bob.WaitType(t, "msg", 20*time.Second)
if msg["id"] != "q3-crash-1" {
t.Fatalf("want q3-crash-1 after restart, got %v", msg)
}
bob.Request(t, map[string]any{"v": 1, "type": "ack", "rid": "cra1", "from": "q3crasha1", "id": "q3-crash-1"})
}
+180
View File
@@ -0,0 +1,180 @@
package chaos_test
import (
"net"
"net/http"
"os"
"os/exec"
"path/filepath"
"strings"
"testing"
"time"
"git.asio.asia/nixevol/NixMsg/test/accept"
"git.asio.asia/nixevol/NixMsg/test/chaos"
"git.asio.asia/nixevol/NixMsg/test/harness"
)
const (
q3ComposeProject = "q3-chaos"
q3ProxyName = "q3-nixmsg"
epPassword = "password1234"
)
// TestQ3ToxiproxyOfflineKeepDelivered:延迟 + 断开后,离线保留期内消息最终送达。
func TestQ3ToxiproxyOfflineKeepDelivered(t *testing.T) {
if _, err := exec.LookPath("docker"); err != nil {
t.Skip("docker 不可用,跳过 toxiproxy 弱网")
}
composeFile := filepath.Join(moduleRoot(t), "test", "chaos", "compose.yml")
down := func() {
cmd := exec.Command("docker", "compose", "-p", q3ComposeProject, "-f", composeFile, "down", "-v", "--remove-orphans")
_ = cmd.Run()
}
down()
t.Cleanup(down)
up := exec.Command("docker", "compose", "-p", q3ComposeProject, "-f", composeFile, "up", "-d")
if out, err := up.CombinedOutput(); err != nil {
t.Fatalf("compose up: %v\n%s", err, out)
}
apiHostPort := waitComposePort(t, composeFile, "8474", 90*time.Second)
listenHostPort := waitComposePort(t, composeFile, "8475", 90*time.Second)
apiURL := "http://127.0.0.1:" + apiHostPort
waitHTTP(t, apiURL+"/version", 90*time.Second)
srv, err := harness.Start(harness.Options{})
if err != nil {
t.Fatal(err)
}
defer func() { _ = srv.Stop() }()
_, nixPort, err := net.SplitHostPort(srv.Addr)
if err != nil {
t.Fatal(err)
}
c := chaos.NewClient(apiURL)
_ = c.DeleteProxy(q3ProxyName)
if _, err := c.CreateProxy(q3ProxyName, "0.0.0.0:8475", "host.docker.internal:"+nixPort); err != nil {
t.Fatalf("CreateProxy: %v", err)
}
t.Cleanup(func() { _ = c.DeleteProxy(q3ProxyName) })
proxyBase := "http://127.0.0.1:" + listenHostPort
if _, err := c.AddLatency(q3ProxyName, "q3-lag", 150, 50, ""); err != nil {
t.Fatalf("AddLatency: %v", err)
}
ac := accept.AdminLogin(t, srv)
accept.CreateEndpoint(t, ac, "q3alice01", epPassword)
accept.CreateEndpoint(t, ac, "q3bob0001", epPassword)
alice := accept.MQTTLogin(t, proxyBase, "q3alice01", epPassword)
defer alice.Close()
sendOff := alice.Request(t, map[string]any{
"v": 1, "type": "send", "rid": "q3s1", "id": "q3-off-keep",
"to": map[string]any{"kind": "endpoint", "id": "q3bob0001"},
"body": map[string]any{"enc": "utf8", "data": "keep-me"},
"delay_ms": int64(0),
"offline": map[string]any{"keep": true, "ttl_seconds": int64(86400)},
})
if !sendOff.OK {
t.Fatalf("send offline keep: %+v", sendOff)
}
if _, err := c.AddResetPeer(q3ProxyName, "q3-rst", 0, ""); err != nil {
t.Fatalf("AddResetPeer: %v", err)
}
time.Sleep(400 * time.Millisecond)
_ = c.RemoveToxic(q3ProxyName, "q3-rst")
_ = c.RemoveToxic(q3ProxyName, "q3-lag")
bob := accept.MQTTLogin(t, proxyBase, "q3bob0001", epPassword)
defer bob.Close()
msg := bob.WaitType(t, "msg", 20*time.Second)
if msg["id"] != "q3-off-keep" {
t.Fatalf("want q3-off-keep got %v", msg)
}
bob.Request(t, map[string]any{"v": 1, "type": "ack", "rid": "q3a1", "from": "q3alice01", "id": "q3-off-keep"})
}
// TestQ3NetemUntested 记录本机 Windows 上 netem 20% 丢包未测原因。
func TestQ3NetemUntested(t *testing.T) {
t.Log("未测:Linux netem 20% 丢包。本机为 Windows;无法在宿主直接用 tc/netem。" +
"若只对 toxiproxy 容器网卡挂 netshoot,不等于 NixMsg 业务端到端丢包验收。" +
"本波用 toxiproxy 延迟/断开覆盖弱网;netem 需 Linux 宿主或把服务放进同网络 Linux 容器后再测。")
}
func moduleRoot(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 waitComposePort(t *testing.T, composeFile, containerPort string, timeout time.Duration) string {
t.Helper()
deadline := time.Now().Add(timeout)
var last string
for time.Now().Before(deadline) {
cmd := exec.Command("docker", "compose", "-p", q3ComposeProject, "-f", composeFile, "port", "toxiproxy", containerPort)
out, err := cmd.CombinedOutput()
last = strings.TrimSpace(string(out))
if err == nil && last != "" {
_, port, perr := net.SplitHostPort(normalizeComposeAddr(last))
if perr == nil && port != "" {
return port
}
}
time.Sleep(400 * time.Millisecond)
}
t.Fatalf("compose port %s timeout, last=%q", containerPort, last)
return ""
}
func normalizeComposeAddr(s string) string {
s = strings.TrimSpace(s)
if strings.HasPrefix(s, "[::]") {
return "127.0.0.1" + s[len("[::]"):]
}
if strings.HasPrefix(s, "0.0.0.0:") {
return "127.0.0.1" + s[len("0.0.0.0"):]
}
return s
}
func waitHTTP(t *testing.T, url string, timeout time.Duration) {
t.Helper()
client := &http.Client{Timeout: 2 * time.Second}
deadline := time.Now().Add(timeout)
var last error
for time.Now().Before(deadline) {
resp, err := client.Get(url)
if err == nil {
_ = resp.Body.Close()
if resp.StatusCode < 500 {
return
}
last = err
} else {
last = err
}
time.Sleep(400 * time.Millisecond)
}
t.Fatalf("wait HTTP %s: %v", url, last)
}
+78
View File
@@ -0,0 +1,78 @@
package load_test
import (
"fmt"
"testing"
"time"
"git.asio.asia/nixevol/NixMsg/test/accept"
"git.asio.asia/nixevol/NixMsg/test/harness"
)
const (
epPassword = "password1234"
// 本机短时压测目标:几十连接收发。不做 1000 / 10 分钟(见 DEVIATIONS)。
q3LivePairs = 16 // 32 个端、16 对收发
)
// TestQ3LiveShortBurst 短时几十连接登录并完成单聊收发。
func TestQ3LiveShortBurst(t *testing.T) {
srv, err := harness.Start(harness.Options{})
if err != nil {
t.Fatal(err)
}
defer func() { _ = srv.Stop() }()
ac := accept.AdminLogin(t, srv)
type pair struct {
a, b string
}
pairs := make([]pair, 0, q3LivePairs)
for i := 0; i < q3LivePairs; i++ {
a := fmt.Sprintf("q3la%04d", i)
b := fmt.Sprintf("q3lb%04d", i)
accept.CreateEndpoint(t, ac, a, epPassword)
accept.CreateEndpoint(t, ac, b, epPassword)
pairs = append(pairs, pair{a: a, b: b})
}
sessions := make([]*accept.MQTTSession, 0, q3LivePairs*2)
defer func() {
for _, s := range sessions {
s.Close()
}
}()
for _, p := range pairs {
sa := accept.MQTTLogin(t, srv.HTTPBase, p.a, epPassword)
sb := accept.MQTTLogin(t, srv.HTTPBase, p.b, epPassword)
sessions = append(sessions, sa, sb)
}
for i, p := range pairs {
sa := sessions[i*2]
sb := sessions[i*2+1]
msgID := fmt.Sprintf("q3-load-%d", i)
resp := sa.Request(t, map[string]any{
"v": 1, "type": "send", "rid": fmt.Sprintf("ls%d", i), "id": msgID,
"to": map[string]any{"kind": "endpoint", "id": p.b},
"body": map[string]any{"enc": "utf8", "data": "burst"},
"delay_ms": int64(0),
})
if !resp.OK {
t.Fatalf("pair %d send: %+v", i, resp)
}
msg := sb.WaitType(t, "msg", 15*time.Second)
if msg["id"] != msgID {
t.Fatalf("pair %d want %s got %v", i, msgID, msg)
}
ack := sb.Request(t, map[string]any{
"v": 1, "type": "ack", "rid": fmt.Sprintf("la%d", i),
"from": p.a, "id": msgID,
})
if !ack.OK {
t.Fatalf("pair %d ack: %+v", i, ack)
}
}
t.Logf("短时压测通过:%d 连接、%d 对单聊收发成功(未做 1000 连接 / 10 分钟)", q3LivePairs*2, q3LivePairs)
}
+37 -37
View File
@@ -1,120 +1,120 @@
{
"generated_at": "2026-09-29T23:20:07Z",
"generated_at": "2026-09-30T00:26:12Z",
"items": [
{
"id": "F01",
"status": "untested",
"note": "未测:serve 未挂载 POST /api/admin/endpoints(A2 未合入或未接线;当前 A1 对端路由返回 501 亦未挂到进程)"
"status": "pass",
"note": "已测:开通一端、错误密码 MQTT 拒绝、正确密码可连;批量校验/停用/删除转让/同号重开未在本用例穷尽"
},
{
"id": "F02",
"status": "untested",
"note": "未测:端登录/会话令牌属连接 N3,main 上 serve 未挂 broker"
"status": "pass",
"note": "已测:密码登录后 hello 成功(会话令牌路径可用);顶号/锁定/重置踢线未在本用例穷尽"
},
{
"id": "F03",
"status": "untested",
"note": "未测:在线状态属身份 I3 + 连接 N3,未接线"
"note": "未测:directory.list / 断开后离线状态未在本波单独断言"
},
{
"id": "F04",
"status": "untested",
"note": "未测:presence.watch 属身份 I3,未接线"
"note": "未测:presence.watch 订阅通知未覆盖"
},
{
"id": "F05",
"status": "untested",
"note": "未测:消息提交属消息 M1,未挂入 broker 上行"
"status": "pass",
"note": "已测:双端在线单聊送达与确认;崩溃续传见 Q3;消息号冲突/密码门/配额未穷尽"
},
{
"id": "F06",
"status": "untested",
"note": "未测:群消息属消息 M + 身份 I4,未接线"
"status": "pass",
"note": "已测:群成员收到同一份、发送者不收到自己的;入群前不补未单独覆盖"
},
{
"id": "F07",
"status": "untested",
"note": "未测:大小限制属消息/连接,未接线"
"note": "未测:256 KiB 边界与接收上限未覆盖"
},
{
"id": "F08",
"status": "untested",
"note": "未测:投递确认属消息 M2,未接线"
"status": "pass",
"note": "已测:提交成功后杀进程重启,离线保留消息续传;toxiproxy 弱网见 Q3 chaos 测试;应用层去重未单独断言"
},
{
"id": "F09",
"status": "untested",
"note": "未测:保留期属消息 M2,未接线"
"status": "pass",
"note": "已测:选离线保留且接收方稍后上线能送达;超时过期未在本用例拨钟验证"
},
{
"id": "F10",
"status": "untested",
"note": "未测:断线策略属消息 M2,未接线"
"note": "未测:抖动宽限长短断线未单独拨钟"
},
{
"id": "F11",
"status": "untested",
"note": "未测:定时发送属消息 M2/M4,未接线"
"note": "未测:发送方离线后定时到点发送未覆盖"
},
{
"id": "F12",
"status": "untested",
"note": "未测:延迟撤回属消息 M3,未接线"
"status": "pass",
"note": "已测:延迟窗口内撤回对方无 msg/revoked"
},
{
"id": "F13",
"status": "untested",
"note": "未测:撤回判定属消息 M3,未接线"
"status": "pass",
"note": "已测:未推送前撤回成功;群部分撤回未覆盖"
},
{
"id": "F14",
"status": "untested",
"note": "未测:回执属消息 M3,未接线"
"note": "未测:回执补送未覆盖"
},
{
"id": "F15",
"status": "untested",
"note": "未测:对话密码属身份 I2,未接线"
"note": "未测:对话密码授权链路未覆盖"
},
{
"id": "F16",
"status": "untested",
"note": "未测:群权限属身份 I4,未接线"
"status": "pass",
"note": "已测:建群并拉成员后可群发;群主权限/退出/解散同号等未穷尽"
},
{
"id": "F17",
"status": "untested",
"note": "未测:serve 未挂载 /api/admin/login(A1 Handler 已实现,缺总控/接线挂到 cmd/nixmsg;DEVIATIONS 后台接口 A §1)"
"status": "pass",
"note": "已测:管理登录、错误密码锁定、无 CSRF 被拒 / 有 CSRF 可通过;管端开通见 F01;管注册见 F23;令牌越权/查记录无正文等未穷尽"
},
{
"id": "F18",
"status": "untested",
"note": "未测:正文清理属消息 M3,未接线"
"note": "未测:正文删除与记录天数 0 未覆盖"
},
{
"id": "F19",
"status": "untested",
"note": "未测:SDK 接入清单属 S1/S2 任务 4,依赖真实服务接线"
"note": "未测:四种 SDK 接入清单属 S1/S2 任务 4"
},
{
"id": "F20",
"status": "untested",
"note": "未测:裸 MQTT 属连接 N,serve 未挂 broker"
"status": "pass",
"note": "已测:裸 MQTT WebSocket 登录、hello、发、收、确认"
},
{
"id": "F21",
"status": "untested",
"note": "仅验证 listen 上 /healthz+/readyz;后台 API、/mqtt、裸 TCP、注册未挂入 serve(缺总控接线 + 连接 N)"
"status": "pass",
"note": "已测:同一 listen 端口提供 /healthz、管理 API、注册、WebSocket /mqtt;后台分离端口未测"
},
{
"id": "F22",
"status": "pass",
"note": "已测:空目录 admin init + serve,/healthz 与 /readyz 成功,密码不在 serve 日志;未测:备份恢复、升级迁移、证书重载、Docker、/metrics"
"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)"
"status": "pass",
"note": "已测:开关、错码、对码、换码;输错锁定未在本用例穷尽"
}
]
}