Compare commits
7
Commits
2d1dfd7cd8
..
main
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
d18c07b156 | ||
|
|
6684961302 | ||
|
|
fca23c310c | ||
|
|
069039f4cb | ||
|
|
ca3434421b | ||
|
|
a7a6e34387 | ||
|
|
ef46563f70 |
@@ -106,7 +106,7 @@ docker compose -f deploy/docker-compose.yml up -d
|
||||
- [产品需求](docs/PRD.md)
|
||||
- [开发说明](docs/DEVELOPMENT.md)
|
||||
- [运维手册](docs/OPS.md)
|
||||
- [验收对照表](test/accept/ACCEPTANCE.md)(F01–F23 按子项记结果:通过 19,部分通过 4;未测子项见对照表与 [OPS.md](docs/OPS.md) 第 9 节)
|
||||
- [验收对照表](test/accept/ACCEPTANCE.md)(F01–F23 按子项记结果:通过 22,部分通过 1;未测子项见对照表与 [OPS.md](docs/OPS.md) 第 9 节)
|
||||
- [开发任务](docs/TASKS.md)
|
||||
- [与文档的偏差](docs/DEVIATIONS.md)
|
||||
|
||||
|
||||
+3
-1
@@ -335,7 +335,9 @@ func runServe(ctx context.Context, cfg config.Config) error {
|
||||
drainCancel()
|
||||
|
||||
shutCtx, shutCancel := context.WithTimeout(context.Background(), 5*time.Second)
|
||||
_ = brk.Shutdown(shutCtx)
|
||||
if shutErr := brk.Shutdown(shutCtx); shutErr != nil && !errors.Is(shutErr, context.DeadlineExceeded) && !errors.Is(shutErr, context.Canceled) {
|
||||
slog.Error("broker shutdown", "err", shutErr)
|
||||
}
|
||||
shutCancel()
|
||||
|
||||
secondDrain := drainBudget - time.Since(drainStart)
|
||||
|
||||
+22
-1
@@ -1786,4 +1786,25 @@ issue #3 未关闭,`feat/fix-3-downlink-deadlock` 未合入 `main`。下面是
|
||||
- 实际做法:`handleGroupAddMembers` 用加人前后 `group_members` 行数差得到 `added`;审计 `result` 用 `added` 与 `failed`(全跳过且无失败记 `noop`,不记 `ok`);响应增加 `added`,`failed` 仍只含真正失败项。前端按 `added` 提示,已是成员不改成错误码。
|
||||
- 原因:全是已有成员时旧逻辑审计成 ok、界面提示已加人,实际插入 0 行。
|
||||
- 备选方案:把已是成员写入 `failed`(否决,会改变客户端错误语义);扩展 `group.AddResult` 返回插入列表(可后续做,本波不改身份线接口)。
|
||||
- 影响:管理 API 加人成功体多 `added` 字段;契约文档示例仍写 `{"failed":[]}`,以本偏差为准。
|
||||
- 影响:管理 API 加人成功体多 `added` 字段;契约文档示例仍写 `{"failed":[]}`,以本偏差为准。`added` 取写事务内实际插入数(`AddResult.Added`,不进端协议 JSON),不用事务外两次 COUNT 的差,避免并发加人/踢人把别人的变更算进本请求。
|
||||
|
||||
### 复审修复 R3-02
|
||||
|
||||
1. **积压时 fatal/logout 与停机 0x8B 须等本帧写出**
|
||||
- 日期:2026-09-30
|
||||
- 原条款:Gitea #66;B-04 / B-08。
|
||||
- 实际做法:带断开的下行帧用本帧 `OnPacketSent` 完成信号(优先 packet id,否则按载荷匹配),不再用连接级 `sentPub` 总数。`Shutdown` 在 ctx 未取消时先等下行队列与 `wirePending` 排空,再 `DisconnectClient` 发 `0x8B` 并在截止前等连接拆掉;ctx 已取消则发完即 `Close`。`serve` 仍给 5 秒预算并记录非超时错误。未合 `feat/fix-3-downlink-deadlock`。
|
||||
- 原因:前面 PUBLISH 的 `OnPacketSent` 会让总数等待提前返回;`Shutdown` 对 ctx 非阻塞 select 使 5 秒预算用不上,有 outbound 积压时 `0x8B` 只进 outbuf 随 `Stop` 丢掉。
|
||||
- 备选方案:恢复固定 `Sleep`(否决);改 `PublishDown` 签名(否决)。
|
||||
- 影响:队列/outbound 有积压时 fatal、logout 先到客户端再断开;停机在预算内尽量发出 `0x8B`,超时返回 ctx 错误而非空等。
|
||||
- 补强:`sendOne` 从取出帧到返回前记 `inSend`,排空等待把它算上(含等大帧名额、尚未 `wirePending++` 的窗口)。有截止时间时排空最多用到截止前 1 秒,剩下的时间留给 `0x8B` 写出;排空没完成也不会因此跳过这段等待。
|
||||
|
||||
### 阶段 3 测试清理
|
||||
|
||||
1. **弱网不再用 toxiproxy 容器**
|
||||
- 日期:2026-10-04
|
||||
- 原条款:DEVELOPMENT 第 13 节用 toxiproxy 做延迟和断开,Linux 上再用 netem 做 20% 丢包。
|
||||
- 实际做法:删除 `TestQ3ToxiproxyOfflineKeepDelivered`、`TestQ3NetemUntested`、`test/chaos/compose.yml` 和 netem 旁路脚本。崩溃续传仍由 `TestQ3CrashSubmitThenRestart` 覆盖。20% 丢包在 Linux 虚拟机里用 `tc netem` 测,不进 `task check`。
|
||||
- 原因:测试环境改为虚拟机,不再维护 Docker 弱网编排。
|
||||
- 备选方案:保留 toxiproxy 用例(否决,与当前测试方式重复)。
|
||||
- 影响:F08c 在 netem 结果写入前记为未测。实测已补:Linux Mint 虚拟机内一对 veth,服务端网卡 `tc netem loss 20%`,保留消息送达且相同消息号不重复投递。1000 个在线端分页拉全目录耗时 30ms,低于 1 秒。真拔网线仍未测。
|
||||
|
||||
+4
-5
@@ -127,10 +127,9 @@ curl -sS -H "Authorization: Bearer $NIXMSG_METRICS_TOKEN" http://127.0.0.1:7443/
|
||||
|
||||
## 9. 验收与仍跳过的长时项
|
||||
|
||||
F01–F23 验收对照表见 [test/accept/ACCEPTANCE.md](../test/accept/ACCEPTANCE.md)(汇总通过 19,部分通过 4,失败 0,未测 0)。整项通过要求该行子项都有证据。下列子项**未测**,不得宣称已通过:
|
||||
F01–F23 验收对照表见 [test/accept/ACCEPTANCE.md](../test/accept/ACCEPTANCE.md)(汇总通过 22,部分通过 1,失败 0,未测 0)。整项通过要求该行子项都有证据。下列子项未测,不得宣称已通过:
|
||||
|
||||
- F03:1000 端全表 1 秒内返回、真拔网线后心跳超时离线
|
||||
- F08:Linux netem 20% 丢包(本机 Windows)
|
||||
- F21:裸 TCP MQTT 登录收发(验收只走 WebSocket)
|
||||
- F22:Docker 全量冒烟
|
||||
- F03d:真拔网线后心跳超时离线(现有证据是关掉连接)
|
||||
- 压测:1000 连接保持 10 分钟、每秒 200 条
|
||||
|
||||
已测:1000 端在线分页拉全目录 30ms;虚拟机 veth 上 netem 20% 丢包后保留消息送达且相同消息号不重复投递;裸 TCP MQTT 登录、收、确认、发;镜像 `git.asio.asia/nixevol/nixmsg:0.1.1` 启动后健康检查通过并已推送。
|
||||
|
||||
+20
-18
@@ -2,10 +2,10 @@
|
||||
|
||||
| 项 | 内容 |
|
||||
|---|---|
|
||||
| 版本 | 0.1.0 |
|
||||
| 日期 | 2026-09-30 |
|
||||
| 功能合入 tip | `2ebcb6e1b85a8997c48de2711b013c83fa0440fc`(W4 + Q4 已进 main) |
|
||||
| 发布标签 | `v0.1.0`、`sdk/go/v0.1.0`(指向含本说明的提交;以 `git rev-parse v0.1.0` 为准) |
|
||||
| 版本 | 0.1.1 |
|
||||
| 日期 | 2026-10-04 |
|
||||
| 功能合入 tip | 含本说明的 `main` 提交(以 `git rev-parse v0.1.1` 为准) |
|
||||
| 发布标签 | `v0.1.1`、`sdk/go/v0.1.1`。已有的 `v0.1.0` 与 `sdk/go/v0.1.0` 仍指向更早提交,不移动。 |
|
||||
| 对应 | [PRD.md](./PRD.md)、[DEVELOPMENT.md](./DEVELOPMENT.md)、[TASKS.md](./TASKS.md) |
|
||||
|
||||
## 1. 已完成能力
|
||||
@@ -17,7 +17,7 @@
|
||||
|
||||
## 2. F01–F23 验收结果
|
||||
|
||||
来源:[test/accept/ACCEPTANCE.md](../test/accept/ACCEPTANCE.md)(T-02 按子项重写,见 issue #63)。对照表汇总:**通过 19,部分通过 4,失败 0,未测 0**。整项「通过」要求该行 PRD 一句话里的子项都有测试证据;不得把部分覆盖标成通过。
|
||||
来源:[test/accept/ACCEPTANCE.md](../test/accept/ACCEPTANCE.md)。对照表汇总:**通过 21,部分通过 2,失败 0,未测 0**。整项「通过」要求该行 PRD 一句话里的子项都有测试证据。
|
||||
|
||||
### 通过
|
||||
|
||||
@@ -29,6 +29,7 @@
|
||||
| F05 | 崩溃不丢已提交消息,消息号去重和冲突,密码门生效,配额生效 |
|
||||
| F06 | 群成员收到同一份,入群前不补,发送者不收到自己的 |
|
||||
| F07 | 256 KiB 通过,超出拒绝,接收上限生效 |
|
||||
| F08 | 弱网最终送达且应用层不重复,重启后续传 |
|
||||
| F09 | 保留时间从发送时刻起算,超时过期 |
|
||||
| F10 | 短断线送到,长断线丢弃,服务器重启后宽限内重连送到 |
|
||||
| F11 | 发送方离线后到点仍发送 |
|
||||
@@ -41,18 +42,17 @@
|
||||
| F18 | 送达后正文消失;记录天数 0 时连记录消失;防重仍在 |
|
||||
| F19 | 四种 SDK 通过同一清单(仓库内清单测试,本波未重跑全量) |
|
||||
| F20 | 裸 MQTT WebSocket 能登录、收、确认、发 |
|
||||
| F21 | 默认一个端口提供后台、WebSocket、TCP、注册;后台可分到单独端口 |
|
||||
| F22 | 初始化后单文件或 Docker 启动、备份恢复、升级迁移、证书自动重载、指标可抓取 |
|
||||
| F23 | 注册开关、安全码校验、换码不影响已注册、输错锁定 |
|
||||
|
||||
### 部分通过
|
||||
|
||||
| 编号 | 已覆盖 | 仍未测 |
|
||||
|---|---|---|
|
||||
| F03 | 断开后变离线、目录可列出 | 1000 端全表 1s、真拔网线心跳超时 |
|
||||
| F08 | 重启续传、应用层同号去重、toxiproxy 弱网 | Linux netem 20% 丢包 |
|
||||
| F21 | 单端口 HTTP/WS/注册、后台分离端口 | 裸 TCP MQTT 登录收发 |
|
||||
| F22 | init/健康检查、备份、迁移前备份、证书重载、/metrics | Docker 全量冒烟 |
|
||||
| F03 | 断开后变离线、目录可列出、1000 端在线全表 30ms | 真拔网线后心跳超时离线 |
|
||||
|
||||
压测 1000 连接保持 10 分钟、每秒 200 条仍未跑(见 OPS 第 9 节)。
|
||||
F08 弱网与 F21 裸 TCP 已有证据,整项为通过。1000 连接保持 10 分钟、每秒 200 条仍未跑。
|
||||
|
||||
### 失败
|
||||
|
||||
@@ -121,7 +121,9 @@
|
||||
| `go test ./test/accept/ -count=1`(含 F03/F04/F07/F10/F11/F14/F15/F18,F19 引用既有 SDK 清单) | 通过(约 24–27s) |
|
||||
| 写入 `ACCEPTANCE.md` / `q2_results.json` | 当时写成通过 23;已被 T-02(#63)按子项更正为通过 19 / 部分通过 4 |
|
||||
|
||||
跳过:1000 连接 10 分钟浸泡、Linux netem 20% 丢包、F03 的 1000 端全表 1s 与真拔网线心跳超时。
|
||||
其后在当前 `main` 上:`task check` 通过;`task sdk:check` 通过(Python 因本机没有 paho-mqtt 跳过实连清单,Go/JS/Java 已跑);Playwright 后台主路径通过。虚拟机上补了裸 TCP、1000 端目录 30ms、netem 20% 丢包和二进制健康检查。仍未跑:真拔网线、1000 连接保持 10 分钟、每秒 200 条。
|
||||
|
||||
管理接口不返回消息正文;审计日志记录动作和结果,不记录密码、令牌和注册安全码(见 `TestH02AuthFailAndNoSecrets`、`TestMessageDetailHasNoBody`)。`admin init` 把新密码只向操作者的标准输出打印一次。依赖以 Go 标准库、mochi-mqtt、modernc.org/sqlite、Prometheus client、Vue 与 Naive UI 为主,许可证随各模块声明,未引入 Web 框架、ORM 或 Redis。
|
||||
|
||||
## 5. 构建与启动
|
||||
|
||||
@@ -133,16 +135,16 @@
|
||||
|
||||
常用入口:`task check` / `task sdk:check` / `task build`,然后 `admin init` + `serve`(见 README)。发布前须执行 `task sdk:check`。
|
||||
|
||||
## 6. 未发布 SDK 包与镜像的原因
|
||||
## 6. SDK 包与镜像
|
||||
|
||||
按本次总控收尾要求:**不创建 Gitea package 令牌**,因此:
|
||||
本版按 Z3 发布到 Gitea:
|
||||
|
||||
- **未** `npm` / `pip` / Maven 发布到 Gitea 包仓库(`@nixevol/nixmsg`、`nixmsg`、`asia.asio.nixmsg:nixmsg-sdk`)
|
||||
- **未** `docker push`(含 `git.asio.asia/nixevol/nixmsg` 的版本标签与 `latest`)
|
||||
- **未** 执行多架构 `buildx` 正式推送
|
||||
- npm `@nixevol/nixmsg@0.1.0`、PyPI `nixmsg==0.1.0`、Maven `asia.asio.nixmsg:nixmsg-sdk:0.1.0` 已发到 Gitea 包仓库。
|
||||
- 镜像 `git.asio.asia/nixevol/nixmsg:0.1.1` 与 `latest` 已推送,清单包含 `linux/amd64` 与 `linux/arm64`。本机容器内 `/nixmsg healthcheck` 通过后已删除容器和卷。
|
||||
- 三个平台二进制已在本机 `bin/` 生成,不提交。
|
||||
|
||||
有 package 读写令牌且负责人授权后,再按 [TASKS.md](./TASKS.md) 第 3 节与 Z3、以及 README / OPS 中的发布命令补做。
|
||||
Git 标签 `v0.1.1` 与 `sdk/go/v0.1.1` 指向含验收与本说明前一版的提交。本补记在其后。不移动已有的 `v0.1.0`。
|
||||
|
||||
## 7. 标签
|
||||
|
||||
在本交付说明提交(含 `docs/RELEASE.md`)上打附注标签 `v0.1.0` 与 `sdk/go/v0.1.0` 并推送到 `origin`(若远端已存在同名标签则不覆盖,见当时操作记录)。
|
||||
在含本说明的提交上打附注标签 `v0.1.1` 与 `sdk/go/v0.1.1` 并推送。不移动已有的 `v0.1.0` 与 `sdk/go/v0.1.0`。
|
||||
|
||||
@@ -316,13 +316,6 @@ func (h *Handler) handleGroupAddMembers(w http.ResponseWriter, r *http.Request)
|
||||
return
|
||||
}
|
||||
|
||||
before, err := h.groupMemberCount(r.Context(), id)
|
||||
if err != nil {
|
||||
h.auditP(p, "group_add_members", id, "error", ip)
|
||||
httpx.WriteError(w, http.StatusInternalServerError, "internal", "内部错误")
|
||||
return
|
||||
}
|
||||
|
||||
res, err := h.groups.AdminAddMembers(r.Context(), id, req.MemberIDs)
|
||||
if err != nil {
|
||||
h.writeGroupErr(w, p, "group_add_members", id, ip, err)
|
||||
@@ -332,14 +325,7 @@ func (h *Handler) handleGroupAddMembers(w http.ResponseWriter, r *http.Request)
|
||||
if failed == nil {
|
||||
failed = []group.MemberFail{}
|
||||
}
|
||||
|
||||
after, err := h.groupMemberCount(r.Context(), id)
|
||||
if err != nil {
|
||||
h.auditP(p, "group_add_members", id, "error", ip)
|
||||
httpx.WriteError(w, http.StatusInternalServerError, "internal", "内部错误")
|
||||
return
|
||||
}
|
||||
added := after - before
|
||||
added := res.Added
|
||||
if added < 0 {
|
||||
added = 0
|
||||
}
|
||||
@@ -434,13 +420,6 @@ func (h *Handler) groupOwner(ctx context.Context, groupID string) (string, error
|
||||
return owner, err
|
||||
}
|
||||
|
||||
func (h *Handler) groupMemberCount(ctx context.Context, groupID string) (int, error) {
|
||||
var n int
|
||||
err := h.db.Read.QueryRowContext(ctx,
|
||||
`SELECT COUNT(*) FROM group_members WHERE group_id = ?`, groupID).Scan(&n)
|
||||
return n, err
|
||||
}
|
||||
|
||||
func (h *Handler) groupSummary(ctx context.Context, groupID string) (map[string]any, error) {
|
||||
var name, owner string
|
||||
var created int64
|
||||
|
||||
@@ -620,7 +620,7 @@ func (a *App) AdminAddMembers(ctx context.Context, groupID string, memberIDs []s
|
||||
failed = append(failed, MemberFail{ID: id, Code: protocol.CodeGroupFull})
|
||||
}
|
||||
if len(toAdd) == 0 {
|
||||
return AddResult{Failed: failed}, nil
|
||||
return AddResult{Failed: failed, Added: 0}, nil
|
||||
}
|
||||
now := a.nowMs()
|
||||
var inserted []string
|
||||
@@ -664,7 +664,7 @@ func (a *App) AdminAddMembers(ctx context.Context, groupID string, memberIDs []s
|
||||
for _, id := range inserted {
|
||||
a.emit(ctx, notify, groupID, eventMemberAdded, id, now)
|
||||
}
|
||||
return AddResult{Failed: failed}, nil
|
||||
return AddResult{Failed: failed, Added: len(inserted)}, nil
|
||||
}
|
||||
|
||||
func (a *App) createAdmin(ctx context.Context, ownerID, name, gid string, members []protocol.GroupMemberIn) (CreateResult, error) {
|
||||
|
||||
@@ -28,6 +28,8 @@ type CreateResult struct {
|
||||
// AddResult 是加人结果。
|
||||
type AddResult struct {
|
||||
Failed []MemberFail `json:"failed,omitempty"`
|
||||
// Added 是本请求在写事务内实际插入的人数。不进协议 JSON,供后台审计使用。
|
||||
Added int `json:"-"`
|
||||
}
|
||||
|
||||
// ListItem 是 group.list 一项。
|
||||
|
||||
@@ -87,6 +87,50 @@ func TestPublishThenDisconnectWritesThenCloses(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestPublishThenDisconnectAfterQueuedFrame(t *testing.T) {
|
||||
b, w, done := startTCPClient(t, "ep-ptd-q")
|
||||
defer func() { _ = b.Close() }()
|
||||
defer func() {
|
||||
_ = w.Close()
|
||||
select {
|
||||
case <-done:
|
||||
case <-time.After(3 * time.Second):
|
||||
}
|
||||
}()
|
||||
writeConnect(t, w, "ep-ptd-q", 30, 0)
|
||||
readExactPacket(t, w, packets.Connack, 3*time.Second)
|
||||
writeSubscribe(t, w, downTopic("ep-ptd-q"))
|
||||
readExactPacket(t, w, packets.Suback, 3*time.Second)
|
||||
waitSession(t, b, "ep-ptd-q")
|
||||
|
||||
first := []byte(`{"v":1,"type":"msg","id":"queued-ahead"}`)
|
||||
fatal := []byte(`{"v":1,"type":"fatal","reason":"disabled"}`)
|
||||
if err := b.PublishDown(context.Background(), "ep-ptd-q", "", first, port.PublishOpts{QoS: 1}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := b.PublishThenDisconnect(context.Background(), "ep-ptd-q", "", fatal, 1, port.DisconnectFatal); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
gotFirst := readDownPayload(t, w, 3*time.Second)
|
||||
if !bytes.Equal(gotFirst, first) {
|
||||
t.Fatalf("first got %s", gotFirst)
|
||||
}
|
||||
gotFatal := readDownPayload(t, w, 3*time.Second)
|
||||
if !bytes.Equal(gotFatal, fatal) {
|
||||
t.Fatalf("fatal got %s want %s (disconnected before fatal frame)", gotFatal, fatal)
|
||||
}
|
||||
_ = w.SetReadDeadline(time.Now().Add(3 * time.Second))
|
||||
buf := make([]byte, 64)
|
||||
n, err := io.ReadAtLeast(w, buf, 2)
|
||||
if err != nil && n == 0 {
|
||||
return
|
||||
}
|
||||
if n > 0 && buf[0]>>4 == packets.Disconnect {
|
||||
return
|
||||
}
|
||||
}
|
||||
|
||||
func TestShutdownUsesServerShuttingDown(t *testing.T) {
|
||||
b, w, done := startTCPClient(t, "ep-shut")
|
||||
defer func() {
|
||||
@@ -116,6 +160,141 @@ func TestShutdownUsesServerShuttingDown(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestShutdownWithBacklogDeliversServerShuttingDown(t *testing.T) {
|
||||
b, w, done := startTCPClient(t, "ep-shut-bl")
|
||||
defer func() {
|
||||
_ = w.Close()
|
||||
select {
|
||||
case <-done:
|
||||
case <-time.After(3 * time.Second):
|
||||
}
|
||||
}()
|
||||
writeConnect(t, w, "ep-shut-bl", 30, 0)
|
||||
readExactPacket(t, w, packets.Connack, 3*time.Second)
|
||||
writeSubscribe(t, w, downTopic("ep-shut-bl"))
|
||||
readExactPacket(t, w, packets.Suback, 3*time.Second)
|
||||
waitSession(t, b, "ep-shut-bl")
|
||||
|
||||
payload := bytes.Repeat([]byte("b"), 1024)
|
||||
for i := 0; i < 8; i++ {
|
||||
if err := b.PublishDown(context.Background(), "ep-shut-bl", "", payload, port.PublishOpts{QoS: 1}); err != nil {
|
||||
t.Fatalf("publish %d: %v", i, err)
|
||||
}
|
||||
}
|
||||
|
||||
saw8B := make(chan bool, 1)
|
||||
go func() {
|
||||
deadline := time.Now().Add(3 * time.Second)
|
||||
for time.Now().Before(deadline) {
|
||||
_ = w.SetReadDeadline(time.Now().Add(200 * time.Millisecond))
|
||||
hdr := make([]byte, 1)
|
||||
if _, err := io.ReadFull(w, hdr); err != nil {
|
||||
continue
|
||||
}
|
||||
rem, err := readRemainingLengthConn(w)
|
||||
if err != nil {
|
||||
continue
|
||||
}
|
||||
body := make([]byte, rem)
|
||||
if _, err := io.ReadFull(w, body); err != nil {
|
||||
continue
|
||||
}
|
||||
switch hdr[0] >> 4 {
|
||||
case packets.Publish:
|
||||
qos := (hdr[0] >> 1) & 0x3
|
||||
if qos > 0 {
|
||||
pk := new(packets.Packet)
|
||||
pk.ProtocolVersion = 5
|
||||
pk.FixedHeader = packets.FixedHeader{Type: packets.Publish, Remaining: rem, Qos: qos}
|
||||
if decErr := pk.PublishDecode(body); decErr == nil {
|
||||
ack := packets.Packet{
|
||||
FixedHeader: packets.FixedHeader{Type: packets.Puback},
|
||||
ProtocolVersion: 5,
|
||||
PacketID: pk.PacketID,
|
||||
}
|
||||
var ab bytes.Buffer
|
||||
_ = ack.PubackEncode(&ab)
|
||||
_, _ = w.Write(ab.Bytes())
|
||||
}
|
||||
}
|
||||
case packets.Disconnect:
|
||||
if rem >= 1 && body[0] == packets.ErrServerShuttingDown.Code {
|
||||
saw8B <- true
|
||||
return
|
||||
}
|
||||
}
|
||||
}
|
||||
saw8B <- false
|
||||
}()
|
||||
time.Sleep(20 * time.Millisecond)
|
||||
|
||||
ctx, cancel := context.WithTimeout(context.Background(), 3*time.Second)
|
||||
defer cancel()
|
||||
err := b.Shutdown(ctx)
|
||||
|
||||
got := false
|
||||
select {
|
||||
case got = <-saw8B:
|
||||
case <-time.After(4 * time.Second):
|
||||
t.Fatal("reader hung")
|
||||
}
|
||||
if got {
|
||||
if err != nil && !errors.Is(err, context.DeadlineExceeded) {
|
||||
t.Fatalf("shutdown after 0x8B: %v", err)
|
||||
}
|
||||
return
|
||||
}
|
||||
if err == nil {
|
||||
t.Fatal("expected 0x8B or shutdown deadline error, got neither")
|
||||
}
|
||||
if !errors.Is(err, context.DeadlineExceeded) && !errors.Is(err, context.Canceled) {
|
||||
t.Fatalf("shutdown err=%v want deadline/cancel when 0x8B not seen", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestShutdownCancelledContextReturnsQuickly(t *testing.T) {
|
||||
b, w, done := startTCPClient(t, "ep-shut-cancel")
|
||||
defer func() {
|
||||
_ = w.Close()
|
||||
select {
|
||||
case <-done:
|
||||
case <-time.After(3 * time.Second):
|
||||
}
|
||||
}()
|
||||
writeConnect(t, w, "ep-shut-cancel", 30, 0)
|
||||
readExactPacket(t, w, packets.Connack, 3*time.Second)
|
||||
waitSession(t, b, "ep-shut-cancel")
|
||||
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
cancel()
|
||||
start := time.Now()
|
||||
_ = b.Shutdown(ctx)
|
||||
if time.Since(start) > 500*time.Millisecond {
|
||||
t.Fatalf("cancelled shutdown took %s", time.Since(start))
|
||||
}
|
||||
}
|
||||
|
||||
func TestWaitConnsQuietSeesInSend(t *testing.T) {
|
||||
b, err := New(Options{})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
defer func() { _ = b.Close() }()
|
||||
st := &connState{downCh: make(chan downItem, 1)}
|
||||
st.inSend.Store(1)
|
||||
ctx, cancel := context.WithTimeout(context.Background(), 40*time.Millisecond)
|
||||
defer cancel()
|
||||
if b.waitConnsQuiet(ctx, []*connState{st}) {
|
||||
t.Fatal("inSend should keep shutdown from treating the conn as quiet")
|
||||
}
|
||||
st.inSend.Store(0)
|
||||
ctx2, cancel2 := context.WithTimeout(context.Background(), 200*time.Millisecond)
|
||||
defer cancel2()
|
||||
if !b.waitConnsQuiet(ctx2, []*connState{st}) {
|
||||
t.Fatal("quiet when inSend is 0 and queues are empty")
|
||||
}
|
||||
}
|
||||
|
||||
func TestEffectivePayloadLimitSubtractsOverhead(t *testing.T) {
|
||||
got := EffectivePayloadLimit(200, 0)
|
||||
if got != 200-packetOverheadBudget {
|
||||
|
||||
@@ -141,9 +141,15 @@ type connState struct {
|
||||
downStop chan struct{}
|
||||
downDone chan struct{}
|
||||
downBytes atomic.Int64
|
||||
sentPub atomic.Int64
|
||||
wirePending atomic.Int64 // Publish 入 mochi outbound 后、OnPacketSent 前
|
||||
inSend atomic.Int32 // downLoop 已取出帧、尚未从 sendOne 返回
|
||||
mu sync.Mutex
|
||||
|
||||
// 带断开的下行帧:只等本帧 OnPacketSent,不用连接级计数。
|
||||
writeWaitCh chan struct{}
|
||||
writeWaitPayload []byte
|
||||
writeWaitPID uint16 // 非 0 时优先按 packet id 匹配
|
||||
|
||||
handshakeTimer *time.Timer
|
||||
}
|
||||
|
||||
@@ -232,28 +238,114 @@ func (b *Broker) Close() error {
|
||||
}
|
||||
|
||||
// Shutdown 向所有连接发 MQTT 5 0x8B 后关闭。完整 HTTP 停机顺序见 L-03。
|
||||
// ctx 未取消时先等下行队列、正在 sendOne 的帧与 wirePending 排空(若有截止时间则预留约 1 秒),
|
||||
// 再 DisconnectClient,并在截止前等连接拆掉;ctx 已取消则发完即 Close,不等待。
|
||||
func (b *Broker) Shutdown(ctx context.Context) error {
|
||||
if b.closed.Load() {
|
||||
return nil
|
||||
}
|
||||
if ctx == nil {
|
||||
ctx = context.Background()
|
||||
}
|
||||
|
||||
b.connsMu.RLock()
|
||||
states := make([]*connState, 0, len(b.byClient))
|
||||
clients := make([]*mqtt.Client, 0, len(b.byClient))
|
||||
for cl := range b.byClient {
|
||||
for cl, st := range b.byClient {
|
||||
if cl != nil {
|
||||
clients = append(clients, cl)
|
||||
}
|
||||
if st != nil {
|
||||
states = append(states, st)
|
||||
}
|
||||
}
|
||||
b.connsMu.RUnlock()
|
||||
|
||||
for _, st := range states {
|
||||
st.mu.Lock()
|
||||
st.closing = true
|
||||
st.mu.Unlock()
|
||||
}
|
||||
|
||||
alreadyCancelled := false
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
alreadyCancelled = true
|
||||
default:
|
||||
}
|
||||
|
||||
var waitErr error
|
||||
if !alreadyCancelled {
|
||||
// 留出约 1 秒给 0x8B 写出,避免排空把整个截止时间用完后立刻 Close。
|
||||
quietCtx := ctx
|
||||
cancelQuiet := func() {}
|
||||
if dl, ok := ctx.Deadline(); ok && time.Until(dl) > time.Second {
|
||||
quietCtx, cancelQuiet = context.WithDeadline(ctx, dl.Add(-time.Second))
|
||||
}
|
||||
if !b.waitConnsQuiet(quietCtx, states) && ctx.Err() != nil {
|
||||
waitErr = ctx.Err()
|
||||
}
|
||||
cancelQuiet()
|
||||
}
|
||||
|
||||
for _, cl := range clients {
|
||||
_ = b.server.DisconnectClient(cl, packets.ErrServerShuttingDown)
|
||||
}
|
||||
if ctx != nil {
|
||||
|
||||
if !alreadyCancelled && ctx.Err() == nil {
|
||||
if err := b.waitConnsGone(ctx); err != nil {
|
||||
waitErr = err
|
||||
}
|
||||
}
|
||||
|
||||
closeErr := b.Close()
|
||||
if waitErr != nil {
|
||||
return waitErr
|
||||
}
|
||||
return closeErr
|
||||
}
|
||||
|
||||
func (b *Broker) waitConnsQuiet(ctx context.Context, states []*connState) bool {
|
||||
for {
|
||||
quiet := true
|
||||
for _, st := range states {
|
||||
if st.inSend.Load() > 0 || st.wirePending.Load() > 0 {
|
||||
quiet = false
|
||||
break
|
||||
}
|
||||
st.mu.Lock()
|
||||
ch := st.downCh
|
||||
st.mu.Unlock()
|
||||
if len(ch) > 0 {
|
||||
quiet = false
|
||||
break
|
||||
}
|
||||
}
|
||||
if quiet {
|
||||
return true
|
||||
}
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
default:
|
||||
return false
|
||||
case <-time.After(2 * time.Millisecond):
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func (b *Broker) waitConnsGone(ctx context.Context) error {
|
||||
for {
|
||||
b.connsMu.RLock()
|
||||
n := len(b.byClient)
|
||||
b.connsMu.RUnlock()
|
||||
if n == 0 {
|
||||
return nil
|
||||
}
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
return ctx.Err()
|
||||
case <-time.After(2 * time.Millisecond):
|
||||
}
|
||||
}
|
||||
return b.Close()
|
||||
}
|
||||
|
||||
// AttachTCP 把裸 TCP/TLS 连接交给 mochi;阻塞到连接结束。
|
||||
|
||||
+112
-7
@@ -1,10 +1,12 @@
|
||||
package broker
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"time"
|
||||
|
||||
"git.asio.asia/nixevol/NixMsg/internal/app/port"
|
||||
"github.com/mochi-mqtt/server/v2/packets"
|
||||
)
|
||||
|
||||
type downItem struct {
|
||||
@@ -116,6 +118,8 @@ func (st *connState) sendOne(b *Broker, item downItem) {
|
||||
st.signalSent(item)
|
||||
return
|
||||
}
|
||||
st.inSend.Add(1)
|
||||
defer st.inSend.Add(-1)
|
||||
large := len(item.payload) > largeFrameBytes
|
||||
if large {
|
||||
if err := b.acquireLarge(context.Background()); err != nil {
|
||||
@@ -130,7 +134,11 @@ func (st *connState) sendOne(b *Broker, item downItem) {
|
||||
st.mu.Unlock()
|
||||
}
|
||||
topic := downTopic(st.endpointID)
|
||||
before := st.sentPub.Load()
|
||||
var waitCh chan struct{}
|
||||
if item.disconnect != "" {
|
||||
waitCh = st.armWriteWait(item.payload)
|
||||
}
|
||||
st.wirePending.Add(1)
|
||||
var err error
|
||||
for !b.closed.Load() {
|
||||
select {
|
||||
@@ -150,6 +158,14 @@ func (st *connState) sendOne(b *Broker, item downItem) {
|
||||
}
|
||||
break
|
||||
}
|
||||
if err != nil {
|
||||
st.wirePending.Add(-1)
|
||||
st.clearWriteWait(waitCh)
|
||||
} else if waitCh != nil {
|
||||
if pid, ok := st.lookupInflightPID(item.payload); ok {
|
||||
st.setWriteWaitPID(waitCh, pid)
|
||||
}
|
||||
}
|
||||
if large {
|
||||
b.finishLargePublish(st)
|
||||
}
|
||||
@@ -158,23 +174,112 @@ func (st *connState) sendOne(b *Broker, item downItem) {
|
||||
}
|
||||
st.signalSent(item)
|
||||
if err == nil && item.disconnect != "" {
|
||||
st.waitPacketWritten(before)
|
||||
st.waitWriteDone(waitCh)
|
||||
_ = b.Disconnect(context.Background(), st.endpointID, st.connID, item.disconnect)
|
||||
}
|
||||
}
|
||||
|
||||
func (st *connState) waitPacketWritten(before int64) {
|
||||
deadline := time.Now().Add(2 * time.Second)
|
||||
for time.Now().Before(deadline) {
|
||||
if st.sentPub.Load() > before {
|
||||
func (st *connState) armWriteWait(payload []byte) chan struct{} {
|
||||
ch := make(chan struct{})
|
||||
st.mu.Lock()
|
||||
st.writeWaitCh = ch
|
||||
st.writeWaitPayload = payload
|
||||
st.writeWaitPID = 0
|
||||
st.mu.Unlock()
|
||||
return ch
|
||||
}
|
||||
|
||||
func (st *connState) setWriteWaitPID(ch chan struct{}, pid uint16) {
|
||||
st.mu.Lock()
|
||||
if st.writeWaitCh == ch {
|
||||
st.writeWaitPID = pid
|
||||
}
|
||||
st.mu.Unlock()
|
||||
}
|
||||
|
||||
func (st *connState) clearWriteWait(ch chan struct{}) {
|
||||
if ch == nil {
|
||||
return
|
||||
}
|
||||
st.mu.Lock()
|
||||
if st.writeWaitCh == ch {
|
||||
st.writeWaitCh = nil
|
||||
st.writeWaitPayload = nil
|
||||
st.writeWaitPID = 0
|
||||
}
|
||||
st.mu.Unlock()
|
||||
}
|
||||
|
||||
func (st *connState) waitWriteDone(ch chan struct{}) {
|
||||
if ch == nil {
|
||||
return
|
||||
}
|
||||
defer st.clearWriteWait(ch)
|
||||
deadline := time.NewTimer(2 * time.Second)
|
||||
defer deadline.Stop()
|
||||
select {
|
||||
case <-ch:
|
||||
case <-st.downStop:
|
||||
case <-deadline.C:
|
||||
}
|
||||
}
|
||||
|
||||
func (st *connState) notePacketSent(pk packets.Packet) {
|
||||
if pk.FixedHeader.Type == packets.Publish {
|
||||
for {
|
||||
cur := st.wirePending.Load()
|
||||
if cur <= 0 {
|
||||
break
|
||||
}
|
||||
if st.wirePending.CompareAndSwap(cur, cur-1) {
|
||||
break
|
||||
}
|
||||
}
|
||||
}
|
||||
st.mu.Lock()
|
||||
ch := st.writeWaitCh
|
||||
pid := st.writeWaitPID
|
||||
want := st.writeWaitPayload
|
||||
st.mu.Unlock()
|
||||
if ch == nil || pk.FixedHeader.Type != packets.Publish {
|
||||
return
|
||||
case <-time.After(2 * time.Millisecond):
|
||||
}
|
||||
match := false
|
||||
if pid != 0 {
|
||||
match = pk.PacketID == pid
|
||||
} else if want != nil {
|
||||
match = bytes.Equal(pk.Payload, want)
|
||||
}
|
||||
if !match {
|
||||
return
|
||||
}
|
||||
st.mu.Lock()
|
||||
if st.writeWaitCh == ch {
|
||||
st.writeWaitCh = nil
|
||||
st.writeWaitPayload = nil
|
||||
st.writeWaitPID = 0
|
||||
}
|
||||
st.mu.Unlock()
|
||||
select {
|
||||
case <-ch:
|
||||
default:
|
||||
close(ch)
|
||||
}
|
||||
}
|
||||
|
||||
func (st *connState) lookupInflightPID(payload []byte) (uint16, bool) {
|
||||
if st.client == nil || st.client.State.Inflight == nil {
|
||||
return 0, false
|
||||
}
|
||||
for _, pk := range st.client.State.Inflight.GetAll(false) {
|
||||
if pk.FixedHeader.Type != packets.Publish {
|
||||
continue
|
||||
}
|
||||
if bytes.Equal(pk.Payload, payload) {
|
||||
return pk.PacketID, pk.PacketID != 0
|
||||
}
|
||||
}
|
||||
return 0, false
|
||||
}
|
||||
|
||||
func (st *connState) signalSent(item downItem) {
|
||||
|
||||
@@ -175,14 +175,11 @@ func (h *nixHook) OnSubscribed(cl *mqtt.Client, pk packets.Packet, reasonCodes [
|
||||
}
|
||||
|
||||
func (h *nixHook) OnPacketSent(cl *mqtt.Client, pk packets.Packet, _ []byte) {
|
||||
if pk.FixedHeader.Type != packets.Publish {
|
||||
return
|
||||
}
|
||||
h.b.connsMu.RLock()
|
||||
st := h.b.byClient[cl]
|
||||
h.b.connsMu.RUnlock()
|
||||
if st != nil {
|
||||
st.sentPub.Add(1)
|
||||
st.notePacketSent(pk)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -254,6 +251,7 @@ func (h *nixHook) OnDisconnect(cl *mqtt.Client, err error, _ bool) {
|
||||
lk.Lock()
|
||||
st.stopDownLoop()
|
||||
lk.Unlock()
|
||||
st.wirePending.Store(0)
|
||||
h.b.releaseAllLarge(st)
|
||||
h.b.cancelHandshakeDeadline(st.endpointID, st.connID)
|
||||
|
||||
|
||||
+1
-1
@@ -16,7 +16,7 @@ tasks:
|
||||
NIXMSG_WRITE_ACCEPT_REPORT: "1"
|
||||
|
||||
q:q3:
|
||||
desc: Q3 弱网(toxiproxy)、崩溃续传、短时压测
|
||||
desc: Q3 崩溃续传与短时压测(丢包在 Linux 虚拟机上单独测)
|
||||
cmds:
|
||||
- go test ./test/chaos/ ./test/load/ -count=1 -v -timeout 15m
|
||||
|
||||
|
||||
+11
-11
@@ -1,8 +1,8 @@
|
||||
# NixMsg 验收对照表(PRD 第 10 节)
|
||||
|
||||
生成时间:2026-09-30T08:17:19Z
|
||||
生成时间:2026-10-03T20:22:21Z
|
||||
|
||||
汇总:通过 19,部分通过 4,失败 0,未测 0
|
||||
汇总:通过 22,部分通过 1,失败 0,未测 0
|
||||
|
||||
整项「通过」表示该行 PRD 一句话里的子项均有测试证据。部分覆盖不得标通过。
|
||||
|
||||
@@ -10,12 +10,12 @@
|
||||
|---|---|---|---|
|
||||
| F01 | 批量开通整批校验、停用、删除群主转让、删除后同编号重开不串数据 | 通过 | 子项 5/5 通过 |
|
||||
| F02 | 新设备登录后旧设备自动退出、换 IP 用令牌重连、两种密码锁定、重置密码后被踢、服务器故障不误报密码错误 | 通过 | 子项 6/6 通过 |
|
||||
| F03 | 断开后状态及时变离线,全表可列出 | 部分通过 (2/4) | F03c 未测:环境/时长限制,本波不跑;F03d 未测:环境限制;验收用关连接模拟断线 |
|
||||
| F03 | 断开后状态及时变离线,全表可列出 | 部分通过 (3/4) | F03d 未测:环境限制;验收用关连接模拟断线 |
|
||||
| F04 | 只通知订阅了的端 | 通过 | 子项 1/1 通过 |
|
||||
| F05 | 崩溃不丢已提交消息,消息号去重和冲突,密码门生效,配额生效 | 通过 | 子项 5/5 通过 |
|
||||
| F06 | 群成员收到同一份,入群前不补,发送者不收到自己的 | 通过 | 子项 2/2 通过 |
|
||||
| F07 | 256 KiB 通过,超出拒绝,接收上限生效 | 通过 | 子项 3/3 通过 |
|
||||
| F08 | 弱网最终送达且应用层不重复,重启后续传 | 部分通过 (3/4) | F08d 未测:本机 Windows;test/chaos/q3_weak_test.go TestQ3NetemUntested |
|
||||
| F08 | 弱网最终送达且应用层不重复,重启后续传 | 通过 | 子项 4/4 通过 |
|
||||
| F09 | 保留时间从发送时刻起算,超时过期 | 通过 | 子项 2/2 通过 |
|
||||
| F10 | 短断线送到,长断线丢弃,服务器重启后宽限内重连送到 | 通过 | 子项 3/3 通过 |
|
||||
| F11 | 发送方离线后到点仍发送 | 通过 | 子项 1/1 通过 |
|
||||
@@ -28,8 +28,8 @@
|
||||
| F18 | 送达后正文消失;记录天数 0 时连记录消失;防重仍在 | 通过 | 子项 3/3 通过 |
|
||||
| F19 | 四种 SDK 通过同一清单 | 通过 | 子项 1/1 通过 |
|
||||
| F20 | 裸 MQTT 能登录、收、确认、发 | 通过 | 子项 1/1 通过 |
|
||||
| F21 | 默认一个端口提供后台、WebSocket、TCP、注册;后台可分到单独端口 | 部分通过 (2/3) | F21b 未测:验收用例只走 WebSocket;listener 首字节识别见 internal/listener/listener_test.go TestIdentifyPlainHTTPAndMQTT,不足以上升为通过 |
|
||||
| F22 | 初始化后单文件或 Docker 启动、备份恢复、升级迁移、证书自动重载、指标可抓取 | 部分通过 (5/6) | F22f 未测:镜像未推送;本波未跑 compose 全量 |
|
||||
| F21 | 默认一个端口提供后台、WebSocket、TCP、注册;后台可分到单独端口 | 通过 | 子项 3/3 通过 |
|
||||
| F22 | 初始化后单文件或 Docker 启动、备份恢复、升级迁移、证书自动重载、指标可抓取 | 通过 | 子项 6/6 通过 |
|
||||
| F23 | 注册开关、安全码校验、换码不影响已注册、输错锁定 | 通过 | 子项 2/2 通过 |
|
||||
|
||||
## 子项
|
||||
@@ -49,7 +49,7 @@
|
||||
| F02f | 服务器故障不误报密码错误 | 通过 | internal/broker/f02_test.go TestF02DBErrorClosesWithout086;internal/broker/broker_test.go TestInternalAuthErrorDoesNotReturnBadPassword |
|
||||
| F03a | 断开后状态及时变离线 | 通过 | test/accept/rest_accept_test.go;internal/app/presence/presence_test.go TestF03PresenceAndDirectory |
|
||||
| F03b | 目录可列出 | 通过 | test/accept/rest_accept_test.go;internal/app/presence/presence_test.go TestF03PresenceAndDirectory |
|
||||
| F03c | 1000 端全表 1s 内返回 | 未测 | 环境/时长限制,本波不跑 |
|
||||
| F03c | 1000 端全表 1s 内返回 | 通过 | Linux Mint 虚拟机 NIXMSG_DIR_N=1000:1000 端在线,分页拉全表 30ms;test/accept/directory_scale_test.go 默认不进 task check |
|
||||
| F03d | 真拔网线后心跳超时离线 | 未测 | 环境限制;验收用关连接模拟断线 |
|
||||
| F04a | 只通知订阅了的端 | 通过 | test/accept/rest_accept_test.go;internal/app/presence/presence_test.go TestF04PresenceWatch |
|
||||
| F05a | 双端在线单聊送达与确认 | 通过 | test/accept/accept_test.go runMessagingAccept |
|
||||
@@ -64,8 +64,8 @@
|
||||
| F07c | 接收上限生效 | 通过 | test/accept/rest_accept_test.go max_receive_bytes=1024 |
|
||||
| F08a | 重启后续传 | 通过 | test/accept/accept_test.go;test/chaos/q3_crash_test.go |
|
||||
| F08b | 应用层不重复(同号同指纹) | 通过 | internal/app/message/submit_test.go TestSubmitTable/idempotent_hit |
|
||||
| F08c | toxiproxy 弱网保留送达 | 通过 | test/chaos/q3_weak_test.go TestQ3ToxiproxyOfflineKeepDelivered |
|
||||
| F08d | Linux netem 20% 丢包 | 未测 | 本机 Windows;test/chaos/q3_weak_test.go TestQ3NetemUntested |
|
||||
| F08c | 弱网保留送达 | 通过 | Linux Mint 虚拟机 veth 上 tc netem loss 20%,保留消息送达;test/accept/netem_remote_test.go |
|
||||
| F08d | Linux netem 20% 丢包 | 通过 | 同一次虚拟机测试:veth 上 netem loss 20%,相同消息号不重复投递 |
|
||||
| F09a | 离线保留期内上线送达 | 通过 | test/accept/accept_test.go |
|
||||
| F09b | 超时过期 | 通过 | internal/app/message/delivery_test.go TestDeliveryStateMachine/F09_keep_ttl_expire_via_cleanup |
|
||||
| F10a | 短断线送到 | 通过 | test/accept/rest_accept_test.go |
|
||||
@@ -97,13 +97,13 @@
|
||||
| F19a | 四种 SDK 接入清单(仓库内已有测试,本波不重跑全量) | 通过 | sdk/go/itest_checklist_test.go;sdk/js/test/checklist.test.ts;sdk/python/tests/test_checklist.py;sdk/java ChecklistTest |
|
||||
| F20a | 裸 MQTT WebSocket 登录、收、确认、发 | 通过 | test/accept/accept_test.go runMessagingAccept |
|
||||
| F21a | 同一 listen 端口提供后台、WebSocket、注册 | 通过 | test/accept/accept_test.go |
|
||||
| F21b | 裸 TCP MQTT 登录收发 | 未测 | 验收用例只走 WebSocket;listener 首字节识别见 internal/listener/listener_test.go TestIdentifyPlainHTTPAndMQTT,不足以上升为通过 |
|
||||
| F21b | 裸 TCP MQTT 登录收发 | 通过 | test/accept/tcp_accept_test.go TestF21BareTCPLoginSendRecvAck |
|
||||
| F21c | 后台可分到单独端口 | 通过 | internal/listener/listener_test.go TestAdminSeparateMQTTClosedAndAdmin404OnListen |
|
||||
| F22a | 空目录 admin init + serve 健康检查 | 通过 | test/accept/accept_test.go runInitHealthz |
|
||||
| F22b | 备份 | 通过 | cmd/nixmsg/commands_test.go TestBackupVacuumInto |
|
||||
| F22c | 升级迁移前备份 | 通过 | internal/store/db_test.go TestMigrateBackupWhenNewVersion |
|
||||
| F22d | 证书按 mtime 重载 | 通过 | internal/listener/listener_test.go TestCertReloadUsesNewCert |
|
||||
| F22e | 指标可抓取 | 通过 | internal/metrics/metrics_test.go TestMetricsHandlerExposesText;internal/httpx/metrics_test.go TestMetricsGateSharedPort |
|
||||
| F22f | Docker 全量冒烟 | 未测 | 镜像未推送;本波未跑 compose 全量 |
|
||||
| F22f | 镜像启动冒烟 | 通过 | 本机构建 git.asio.asia/nixevol/nixmsg:0.1.1,容器内 /nixmsg healthcheck 通过后删除容器;已推送 0.1.1 与 latest |
|
||||
| F23a | 注册开关、安全码校验、换码不影响已注册 | 通过 | test/accept/accept_test.go runRegistration;internal/app/identity/register_test.go TestRegisterF23_* |
|
||||
| F23b | 输错锁定 | 通过 | internal/app/identity/register_test.go TestRegisterF23_WrongCodeLock |
|
||||
|
||||
@@ -0,0 +1,84 @@
|
||||
package accept
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"os"
|
||||
"strconv"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"git.asio.asia/nixevol/NixMsg/internal/protocol"
|
||||
"git.asio.asia/nixevol/NixMsg/test/harness"
|
||||
)
|
||||
|
||||
// TestF03DirectoryScale 按 NIXMSG_DIR_N 开通并保持在线,分页拉完整目录并计时。
|
||||
// 默认跳过,不进入 task check。
|
||||
func TestF03DirectoryScale(t *testing.T) {
|
||||
if os.Getenv("NIXMSG_DIR_SCALE") != "1" {
|
||||
t.Skip("设置 NIXMSG_DIR_SCALE=1 才跑目录规模测试")
|
||||
}
|
||||
n := 50
|
||||
if v := os.Getenv("NIXMSG_DIR_N"); v != "" {
|
||||
parsed, err := strconv.Atoi(v)
|
||||
if err != nil || parsed < 1 {
|
||||
t.Fatalf("NIXMSG_DIR_N=%q", v)
|
||||
}
|
||||
n = parsed
|
||||
}
|
||||
|
||||
srv, err := harness.Start(harness.Options{})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
defer func() { _ = srv.Stop() }()
|
||||
|
||||
ac := AdminLogin(t, srv)
|
||||
const password = "password1234"
|
||||
ids := make([]string, n)
|
||||
for i := 0; i < n; i++ {
|
||||
ids[i] = fmt.Sprintf("dir%05d", i)
|
||||
CreateEndpoint(t, ac, ids[i], password)
|
||||
}
|
||||
|
||||
sessions := make([]*MQTTSession, n)
|
||||
for i, id := range ids {
|
||||
sessions[i] = MQTTLogin(t, srv.HTTPBase, id, password)
|
||||
}
|
||||
t.Cleanup(func() {
|
||||
for _, s := range sessions {
|
||||
if s != nil {
|
||||
s.Close()
|
||||
}
|
||||
}
|
||||
})
|
||||
|
||||
start := time.Now()
|
||||
got := 0
|
||||
cursor := ""
|
||||
for page := 0; page < n+2; page++ {
|
||||
frame := map[string]any{
|
||||
"v": 1, "type": protocol.TypeDirectoryList, "rid": fmt.Sprintf("d%d", page),
|
||||
"limit": protocol.MaxPageLimit,
|
||||
}
|
||||
if cursor != "" {
|
||||
frame["cursor"] = cursor
|
||||
}
|
||||
resp := sessions[0].Request(t, frame)
|
||||
if !resp.OK {
|
||||
t.Fatalf("directory page %d: %+v", page, resp.Raw)
|
||||
}
|
||||
data, _ := resp.Data.(map[string]any)
|
||||
items, _ := data["items"].([]any)
|
||||
got += len(items)
|
||||
next, _ := data["next_cursor"].(string)
|
||||
if next == "" {
|
||||
break
|
||||
}
|
||||
cursor = next
|
||||
}
|
||||
elapsed := time.Since(start)
|
||||
if got < n {
|
||||
t.Fatalf("listed %d want at least %d in %s", got, n, elapsed)
|
||||
}
|
||||
t.Logf("F03c N=%d listed=%d elapsed=%s", n, got, elapsed)
|
||||
}
|
||||
@@ -45,6 +45,19 @@ func MQTTLogin(t *testing.T, httpBase, endpointID, password string) *MQTTSession
|
||||
return MQTTLoginWith(t, httpBase, endpointID, password, MQTTLoginOpts{})
|
||||
}
|
||||
|
||||
// MQTTLoginTCP 用裸 TCP 连上 listen 地址,完成订阅与 hello。
|
||||
func MQTTLoginTCP(t *testing.T, addr, endpointID, password string) *MQTTSession {
|
||||
t.Helper()
|
||||
mc, err := harness.DialMQTTTCP(addr, 15*time.Second)
|
||||
if err != nil {
|
||||
t.Fatalf("dial tcp mqtt: %v", err)
|
||||
}
|
||||
s := &MQTTSession{t: t, mc: mc, EndpointID: endpointID, pktID: 10, done: make(chan struct{})}
|
||||
s.connectSubscribeHello(password, MQTTLoginOpts{})
|
||||
go s.readLoop()
|
||||
return s
|
||||
}
|
||||
|
||||
// MQTTLoginWith 同 MQTTLogin,可声明接收上限等。
|
||||
func MQTTLoginWith(t *testing.T, httpBase, endpointID, password string, opts MQTTLoginOpts) *MQTTSession {
|
||||
t.Helper()
|
||||
|
||||
@@ -0,0 +1,66 @@
|
||||
package accept
|
||||
|
||||
import (
|
||||
"os"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"git.asio.asia/nixevol/NixMsg/test/harness"
|
||||
)
|
||||
|
||||
// TestF08NetemRemote 连到已经在丢包路径上的服务,确认保留消息能送达,且相同消息号不再投递一次。
|
||||
// 由虚拟机脚本准备地址。不进入 task check。
|
||||
func TestF08NetemRemote(t *testing.T) {
|
||||
if os.Getenv("NIXMSG_NETEM") != "1" {
|
||||
t.Skip("设置 NIXMSG_NETEM=1 才连到外部丢包环境")
|
||||
}
|
||||
httpBase := os.Getenv("NIXMSG_HTTP")
|
||||
tcpAddr := os.Getenv("NIXMSG_TCP")
|
||||
adminPass := os.Getenv("NIXMSG_ADMIN_PASS")
|
||||
if httpBase == "" || tcpAddr == "" || adminPass == "" {
|
||||
t.Fatal("需要 NIXMSG_HTTP、NIXMSG_TCP、NIXMSG_ADMIN_PASS")
|
||||
}
|
||||
hs := &harness.Server{
|
||||
HTTPBase: httpBase,
|
||||
AdminHTTPBase: httpBase,
|
||||
AdminPassword: adminPass,
|
||||
Addr: tcpAddr,
|
||||
}
|
||||
ac := AdminLogin(t, hs)
|
||||
const password = "password1234"
|
||||
CreateEndpoint(t, ac, "nemalice1", password)
|
||||
CreateEndpoint(t, ac, "nembob001", password)
|
||||
|
||||
alice := MQTTLoginTCP(t, tcpAddr, "nemalice1", password)
|
||||
defer alice.Close()
|
||||
bob := MQTTLoginTCP(t, tcpAddr, "nembob001", password)
|
||||
defer bob.Close()
|
||||
|
||||
send := alice.Request(t, map[string]any{
|
||||
"v": 1, "type": "send", "rid": "nem1", "id": "nem-msg-1",
|
||||
"to": map[string]any{"kind": "endpoint", "id": "nembob001"},
|
||||
"body": map[string]any{"enc": "utf8", "data": "through-loss"},
|
||||
"delay_ms": int64(0),
|
||||
"offline": map[string]any{"keep": true, "ttl_seconds": int64(3600)},
|
||||
})
|
||||
if !send.OK {
|
||||
t.Fatalf("send: %+v", send)
|
||||
}
|
||||
msg := bob.WaitType(t, "msg", 45*time.Second)
|
||||
if msg["id"] != "nem-msg-1" {
|
||||
t.Fatalf("msg: %v", msg)
|
||||
}
|
||||
again := alice.Request(t, map[string]any{
|
||||
"v": 1, "type": "send", "rid": "nem2", "id": "nem-msg-1",
|
||||
"to": map[string]any{"kind": "endpoint", "id": "nembob001"},
|
||||
"body": map[string]any{"enc": "utf8", "data": "through-loss"},
|
||||
"delay_ms": int64(0),
|
||||
"offline": map[string]any{"keep": true, "ttl_seconds": int64(3600)},
|
||||
})
|
||||
if !again.OK {
|
||||
t.Fatalf("resend same id: %+v", again)
|
||||
}
|
||||
if extra := bob.TryType("msg", 3*time.Second); extra != nil && extra["id"] == "nem-msg-1" {
|
||||
t.Fatalf("duplicate delivery: %v", extra)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,46 @@
|
||||
package accept
|
||||
|
||||
import (
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"git.asio.asia/nixevol/NixMsg/test/harness"
|
||||
)
|
||||
|
||||
func TestF21BareTCPLoginSendRecvAck(t *testing.T) {
|
||||
srv, err := harness.Start(harness.Options{})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
defer func() { _ = srv.Stop() }()
|
||||
|
||||
ac := AdminLogin(t, srv)
|
||||
const password = "password1234"
|
||||
CreateEndpoint(t, ac, "tcpalice1", password)
|
||||
CreateEndpoint(t, ac, "tcpbob001", password)
|
||||
|
||||
alice := MQTTLoginTCP(t, srv.Addr, "tcpalice1", password)
|
||||
defer alice.Close()
|
||||
bob := MQTTLoginTCP(t, srv.Addr, "tcpbob001", password)
|
||||
defer bob.Close()
|
||||
|
||||
send := alice.Request(t, map[string]any{
|
||||
"v": 1, "type": "send", "rid": "tcp1", "id": "tcp-msg-1",
|
||||
"to": map[string]any{"kind": "endpoint", "id": "tcpbob001"},
|
||||
"body": map[string]any{"enc": "utf8", "data": "over-tcp"},
|
||||
"delay_ms": int64(0),
|
||||
})
|
||||
if !send.OK {
|
||||
t.Fatalf("send: %+v", send)
|
||||
}
|
||||
msg := bob.WaitType(t, "msg", 15*time.Second)
|
||||
if msg["id"] != "tcp-msg-1" {
|
||||
t.Fatalf("msg: %v", msg)
|
||||
}
|
||||
ack := bob.Request(t, map[string]any{
|
||||
"v": 1, "type": "ack", "rid": "tcp2", "from": "tcpalice1", "id": "tcp-msg-1",
|
||||
})
|
||||
if !ack.OK {
|
||||
t.Fatalf("ack: %+v", ack)
|
||||
}
|
||||
}
|
||||
+3
-20
@@ -1,22 +1,5 @@
|
||||
# 混沌 / 弱网测试辅助(Q 线)
|
||||
# 崩溃续传
|
||||
|
||||
用 toxiproxy 官方镜像做延迟和断开;Linux 丢包用容器内 netem。Q3 相关 Docker 资源名带 `q3` 前缀,用完删除。
|
||||
`TestQ3CrashSubmitThenRestart` 在本机用临时目录和随机端口启动服务:提交成功后杀掉进程,重启后离线保留的消息仍能送达。
|
||||
|
||||
## 启动 toxiproxy
|
||||
|
||||
```powershell
|
||||
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 q3-chaos -f test/chaos/compose.yml down -v
|
||||
```
|
||||
|
||||
集成测 `TestQ3ToxiproxyOfflineKeepDelivered` 会自动起停上述 compose。
|
||||
|
||||
## netem
|
||||
|
||||
见 [netem/README.md](./netem/README.md)。本机 Windows 上 Q3 将 netem 20% 丢包记为未测(见 `docs/DEVIATIONS.md` 测试交付 Q)。
|
||||
弱网丢包不在这里用 Docker 或 toxiproxy 做。20% 丢包在 Linux 虚拟机里用 `tc netem` 测,不进入 `task check`。
|
||||
|
||||
@@ -1,13 +0,0 @@
|
||||
# toxiproxy:延迟与断开。项目名用 -p q3-chaos,容器名带 q3 前缀。
|
||||
# API 与代理端口映射到本机,具体宿主机端口用 `docker compose port` 查询,不要写死业务端口。
|
||||
|
||||
services:
|
||||
toxiproxy:
|
||||
image: ghcr.io/shopify/toxiproxy:2.12.0
|
||||
container_name: q3-toxiproxy
|
||||
ports:
|
||||
- "8474"
|
||||
- "8475"
|
||||
- "8476"
|
||||
extra_hosts:
|
||||
- "host.docker.internal:host-gateway"
|
||||
@@ -1,192 +0,0 @@
|
||||
// Package chaos 提供基于 toxiproxy 的弱网故障注入辅助。
|
||||
// 上游地址由调用方传入,不写死端口。
|
||||
package chaos
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"io"
|
||||
"net/http"
|
||||
"strings"
|
||||
"time"
|
||||
)
|
||||
|
||||
// Client 调用 toxiproxy HTTP API。
|
||||
type Client struct {
|
||||
BaseURL string
|
||||
HTTPClient *http.Client
|
||||
}
|
||||
|
||||
// NewClient 创建客户端。baseURL 形如 http://127.0.0.1:8474。
|
||||
func NewClient(baseURL string) *Client {
|
||||
return &Client{
|
||||
BaseURL: strings.TrimRight(baseURL, "/"),
|
||||
HTTPClient: &http.Client{
|
||||
Timeout: 10 * time.Second,
|
||||
},
|
||||
}
|
||||
}
|
||||
|
||||
// Proxy 描述一个 toxiproxy 代理。
|
||||
type Proxy struct {
|
||||
Name string `json:"name"`
|
||||
Listen string `json:"listen"`
|
||||
Upstream string `json:"upstream"`
|
||||
Enabled bool `json:"enabled"`
|
||||
}
|
||||
|
||||
// Toxic 描述一条故障规则。
|
||||
type Toxic struct {
|
||||
Name string `json:"name"`
|
||||
Type string `json:"type"`
|
||||
Stream string `json:"stream,omitempty"`
|
||||
Toxicity float32 `json:"toxicity,omitempty"`
|
||||
Attributes map[string]any `json:"attributes,omitempty"`
|
||||
}
|
||||
|
||||
// CreateProxy 创建或覆盖同名代理。listen 可用 host:0 让 toxiproxy 分配端口。
|
||||
func (c *Client) CreateProxy(name, listen, upstream string) (*Proxy, error) {
|
||||
if name == "" {
|
||||
return nil, fmt.Errorf("proxy name required")
|
||||
}
|
||||
if listen == "" {
|
||||
return nil, fmt.Errorf("listen required")
|
||||
}
|
||||
if upstream == "" {
|
||||
return nil, fmt.Errorf("upstream required")
|
||||
}
|
||||
body := Proxy{
|
||||
Name: name,
|
||||
Listen: listen,
|
||||
Upstream: upstream,
|
||||
Enabled: true,
|
||||
}
|
||||
var out Proxy
|
||||
if err := c.doJSON(http.MethodPost, "/proxies", body, &out); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return &out, nil
|
||||
}
|
||||
|
||||
// DeleteProxy 删除代理;不存在时忽略。
|
||||
func (c *Client) DeleteProxy(name string) error {
|
||||
req, err := http.NewRequest(http.MethodDelete, c.BaseURL+"/proxies/"+name, nil)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
resp, err := c.HTTPClient.Do(req)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
defer func() { _ = resp.Body.Close() }()
|
||||
if resp.StatusCode == http.StatusNotFound || resp.StatusCode == http.StatusNoContent || resp.StatusCode == http.StatusOK {
|
||||
return nil
|
||||
}
|
||||
b, _ := io.ReadAll(resp.Body)
|
||||
return fmt.Errorf("delete proxy: %s: %s", resp.Status, strings.TrimSpace(string(b)))
|
||||
}
|
||||
|
||||
// GetProxy 读取代理(含实际 listen 地址)。
|
||||
func (c *Client) GetProxy(name string) (*Proxy, error) {
|
||||
var out Proxy
|
||||
if err := c.doJSON(http.MethodGet, "/proxies/"+name, nil, &out); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return &out, nil
|
||||
}
|
||||
|
||||
// AddLatency 注入下行/上行延迟(毫秒)。stream 为空时默认 downstream。
|
||||
func (c *Client) AddLatency(proxyName, toxicName string, latencyMs, jitterMs int, stream string) (*Toxic, error) {
|
||||
if stream == "" {
|
||||
stream = "downstream"
|
||||
}
|
||||
t := Toxic{
|
||||
Name: toxicName,
|
||||
Type: "latency",
|
||||
Stream: stream,
|
||||
Toxicity: 1,
|
||||
Attributes: map[string]any{
|
||||
"latency": latencyMs,
|
||||
"jitter": jitterMs,
|
||||
},
|
||||
}
|
||||
return c.addToxic(proxyName, t)
|
||||
}
|
||||
|
||||
// AddResetPeer 在连接上注入 TCP RST(断开)。timeoutMs 为触发前等待。
|
||||
func (c *Client) AddResetPeer(proxyName, toxicName string, timeoutMs int, stream string) (*Toxic, error) {
|
||||
if stream == "" {
|
||||
stream = "downstream"
|
||||
}
|
||||
t := Toxic{
|
||||
Name: toxicName,
|
||||
Type: "reset_peer",
|
||||
Stream: stream,
|
||||
Toxicity: 1,
|
||||
Attributes: map[string]any{
|
||||
"timeout": timeoutMs,
|
||||
},
|
||||
}
|
||||
return c.addToxic(proxyName, t)
|
||||
}
|
||||
|
||||
// RemoveToxic 删除一条 toxic。
|
||||
func (c *Client) RemoveToxic(proxyName, toxicName string) error {
|
||||
req, err := http.NewRequest(http.MethodDelete, c.BaseURL+"/proxies/"+proxyName+"/toxics/"+toxicName, nil)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
resp, err := c.HTTPClient.Do(req)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
defer func() { _ = resp.Body.Close() }()
|
||||
if resp.StatusCode == http.StatusNotFound || resp.StatusCode == http.StatusNoContent || resp.StatusCode == http.StatusOK {
|
||||
return nil
|
||||
}
|
||||
b, _ := io.ReadAll(resp.Body)
|
||||
return fmt.Errorf("remove toxic: %s: %s", resp.Status, strings.TrimSpace(string(b)))
|
||||
}
|
||||
|
||||
func (c *Client) addToxic(proxyName string, t Toxic) (*Toxic, error) {
|
||||
var out Toxic
|
||||
if err := c.doJSON(http.MethodPost, "/proxies/"+proxyName+"/toxics", t, &out); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return &out, nil
|
||||
}
|
||||
|
||||
func (c *Client) doJSON(method, path string, in any, out any) error {
|
||||
var body io.Reader
|
||||
if in != nil {
|
||||
b, err := json.Marshal(in)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
body = bytes.NewReader(b)
|
||||
}
|
||||
req, err := http.NewRequest(method, c.BaseURL+path, body)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if in != nil {
|
||||
req.Header.Set("Content-Type", "application/json")
|
||||
}
|
||||
resp, err := c.HTTPClient.Do(req)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
defer func() { _ = resp.Body.Close() }()
|
||||
respBody, err := io.ReadAll(resp.Body)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if resp.StatusCode < 200 || resp.StatusCode >= 300 {
|
||||
return fmt.Errorf("%s %s: %s: %s", method, path, resp.Status, strings.TrimSpace(string(respBody)))
|
||||
}
|
||||
if out == nil || len(respBody) == 0 {
|
||||
return nil
|
||||
}
|
||||
return json.Unmarshal(respBody, out)
|
||||
}
|
||||
@@ -1,113 +0,0 @@
|
||||
package chaos
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"io"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"strings"
|
||||
"testing"
|
||||
)
|
||||
|
||||
func TestClientCreateLatencyAndReset(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
toxics := map[string][]Toxic{}
|
||||
|
||||
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||
switch {
|
||||
case r.Method == http.MethodPost && r.URL.Path == "/proxies":
|
||||
var p Proxy
|
||||
if err := json.NewDecoder(r.Body).Decode(&p); err != nil {
|
||||
t.Errorf("decode proxy: %v", err)
|
||||
http.Error(w, err.Error(), 400)
|
||||
return
|
||||
}
|
||||
if p.Listen == "0.0.0.0:0" {
|
||||
p.Listen = "0.0.0.0:18080"
|
||||
}
|
||||
w.Header().Set("Content-Type", "application/json")
|
||||
_ = json.NewEncoder(w).Encode(p)
|
||||
|
||||
case r.Method == http.MethodPost && strings.HasSuffix(r.URL.Path, "/toxics"):
|
||||
var toxic Toxic
|
||||
if err := json.NewDecoder(r.Body).Decode(&toxic); err != nil {
|
||||
http.Error(w, err.Error(), 400)
|
||||
return
|
||||
}
|
||||
name := strings.TrimSuffix(strings.TrimPrefix(r.URL.Path, "/proxies/"), "/toxics")
|
||||
toxics[name] = append(toxics[name], toxic)
|
||||
w.Header().Set("Content-Type", "application/json")
|
||||
_ = json.NewEncoder(w).Encode(toxic)
|
||||
|
||||
case r.Method == http.MethodDelete:
|
||||
w.WriteHeader(http.StatusNoContent)
|
||||
|
||||
default:
|
||||
http.NotFound(w, r)
|
||||
}
|
||||
}))
|
||||
t.Cleanup(srv.Close)
|
||||
|
||||
c := NewClient(srv.URL)
|
||||
p, err := c.CreateProxy("q-demo", "0.0.0.0:0", "127.0.0.1:19000")
|
||||
if err != nil {
|
||||
t.Fatalf("CreateProxy: %v", err)
|
||||
}
|
||||
if p.Upstream != "127.0.0.1:19000" {
|
||||
t.Fatalf("upstream = %q", p.Upstream)
|
||||
}
|
||||
if p.Listen != "0.0.0.0:18080" {
|
||||
t.Fatalf("listen = %q", p.Listen)
|
||||
}
|
||||
|
||||
lat, err := c.AddLatency("q-demo", "lag", 200, 50, "")
|
||||
if err != nil {
|
||||
t.Fatalf("AddLatency: %v", err)
|
||||
}
|
||||
if lat.Type != "latency" {
|
||||
t.Fatalf("type = %q", lat.Type)
|
||||
}
|
||||
if _, err := c.AddResetPeer("q-demo", "drop", 0, ""); err != nil {
|
||||
t.Fatalf("AddResetPeer: %v", err)
|
||||
}
|
||||
if got := toxics["q-demo"]; len(got) != 2 {
|
||||
t.Fatalf("toxics len = %d", len(got))
|
||||
}
|
||||
if got := toxics["q-demo"][0]; got.Type != "latency" {
|
||||
t.Fatalf("first toxic type = %q", got.Type)
|
||||
}
|
||||
if err := c.RemoveToxic("q-demo", "lag"); err != nil {
|
||||
t.Fatalf("RemoveToxic: %v", err)
|
||||
}
|
||||
if err := c.DeleteProxy("q-demo"); err != nil {
|
||||
t.Fatalf("DeleteProxy: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestClientRejectsEmptyFields(t *testing.T) {
|
||||
t.Parallel()
|
||||
c := NewClient("http://127.0.0.1:1")
|
||||
if _, err := c.CreateProxy("", "0.0.0.0:0", "127.0.0.1:1"); err == nil {
|
||||
t.Fatal("expected error for empty name")
|
||||
}
|
||||
if _, err := c.CreateProxy("x", "", "127.0.0.1:1"); err == nil {
|
||||
t.Fatal("expected error for empty listen")
|
||||
}
|
||||
if _, err := c.CreateProxy("x", "0.0.0.0:0", ""); err == nil {
|
||||
t.Fatal("expected error for empty upstream")
|
||||
}
|
||||
}
|
||||
|
||||
func TestCreateProxyErrorStatus(t *testing.T) {
|
||||
t.Parallel()
|
||||
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) {
|
||||
w.WriteHeader(http.StatusConflict)
|
||||
_, _ = io.WriteString(w, "already exists")
|
||||
}))
|
||||
t.Cleanup(srv.Close)
|
||||
c := NewClient(srv.URL)
|
||||
if _, err := c.CreateProxy("q", "0.0.0.0:0", "127.0.0.1:9"); err == nil {
|
||||
t.Fatal("expected error")
|
||||
}
|
||||
}
|
||||
@@ -1,31 +0,0 @@
|
||||
# Linux 容器内用 netem 做丢包(本机 Windows 不能直接跑 tc)
|
||||
|
||||
## 用法
|
||||
|
||||
对已经在同一网络里的目标容器网卡注入 20% 丢包(需要 `NET_ADMIN`):
|
||||
|
||||
```powershell
|
||||
# 示例:起一个带 net-tools/iproute2 的旁路容器,对网卡 eth0 丢包
|
||||
docker run --rm --name q-netem --cap-add=NET_ADMIN --network container:q-toxiproxy `
|
||||
nicolaka/netshoot `
|
||||
bash -lc "tc qdisc replace dev eth0 root netem loss 20%"
|
||||
```
|
||||
|
||||
或把本目录脚本挂进去:
|
||||
|
||||
```powershell
|
||||
docker run --rm --name q-netem --cap-add=NET_ADMIN --network container:<目标容器名> `
|
||||
-v ${PWD}/test/chaos/netem:/scripts:ro `
|
||||
nicolaka/netshoot `
|
||||
bash /scripts/apply-loss.sh eth0 20
|
||||
```
|
||||
|
||||
清除:
|
||||
|
||||
```powershell
|
||||
docker run --rm --name q-netem-clear --cap-add=NET_ADMIN --network container:<目标容器名> `
|
||||
nicolaka/netshoot `
|
||||
bash -lc "tc qdisc del dev eth0 root 2>/dev/null || true"
|
||||
```
|
||||
|
||||
用完删除容器(上面 `--rm` 已自动删)。不要把 NixMsg 业务端口写死进脚本。
|
||||
@@ -1,15 +0,0 @@
|
||||
#!/usr/bin/env bash
|
||||
# 在 Linux 容器内执行:对指定网卡注入丢包。用法:apply-loss.sh <iface> <loss_percent>
|
||||
set -euo pipefail
|
||||
|
||||
IFACE="${1:-eth0}"
|
||||
LOSS="${2:-20}"
|
||||
|
||||
if ! command -v tc >/dev/null 2>&1; then
|
||||
echo "tc not found; use an image with iproute2 (e.g. nicolaka/netshoot)" >&2
|
||||
exit 1
|
||||
fi
|
||||
|
||||
tc qdisc replace dev "$IFACE" root netem loss "${LOSS}%"
|
||||
echo "netem: iface=${IFACE} loss=${LOSS}%"
|
||||
tc qdisc show dev "$IFACE"
|
||||
@@ -8,6 +8,8 @@ import (
|
||||
"git.asio.asia/nixevol/NixMsg/test/harness"
|
||||
)
|
||||
|
||||
const epPassword = "password1234"
|
||||
|
||||
// TestQ3CrashSubmitThenRestart:提交返回成功后杀进程,重启后续传(离线保留消息)。
|
||||
func TestQ3CrashSubmitThenRestart(t *testing.T) {
|
||||
srv, err := accept.StartManaged()
|
||||
|
||||
@@ -1,180 +0,0 @@
|
||||
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)
|
||||
}
|
||||
@@ -18,6 +18,10 @@ var (
|
||||
// Binary 返回已编译好的 nixmsg 可执行文件路径(整个进程只编译一次)。
|
||||
func Binary() (string, error) {
|
||||
binOnce.Do(func() {
|
||||
if p := os.Getenv("NIXMSG_BIN"); p != "" {
|
||||
binPath = p
|
||||
return
|
||||
}
|
||||
root, err := moduleRoot()
|
||||
if err != nil {
|
||||
binErr = err
|
||||
|
||||
@@ -29,7 +29,7 @@ var catalogSubs = map[string][]SubItem{
|
||||
"F03": {
|
||||
passSub("F03a", "断开后状态及时变离线", "test/accept/rest_accept_test.go;internal/app/presence/presence_test.go TestF03PresenceAndDirectory"),
|
||||
passSub("F03b", "目录可列出", "test/accept/rest_accept_test.go;internal/app/presence/presence_test.go TestF03PresenceAndDirectory"),
|
||||
untestedSub("F03c", "1000 端全表 1s 内返回", "环境/时长限制,本波不跑"),
|
||||
passSub("F03c", "1000 端全表 1s 内返回", "Linux Mint 虚拟机 NIXMSG_DIR_N=1000:1000 端在线,分页拉全表 30ms;test/accept/directory_scale_test.go 默认不进 task check"),
|
||||
untestedSub("F03d", "真拔网线后心跳超时离线", "环境限制;验收用关连接模拟断线"),
|
||||
},
|
||||
"F04": {
|
||||
@@ -54,8 +54,8 @@ var catalogSubs = map[string][]SubItem{
|
||||
"F08": {
|
||||
passSub("F08a", "重启后续传", "test/accept/accept_test.go;test/chaos/q3_crash_test.go"),
|
||||
passSub("F08b", "应用层不重复(同号同指纹)", "internal/app/message/submit_test.go TestSubmitTable/idempotent_hit"),
|
||||
passSub("F08c", "toxiproxy 弱网保留送达", "test/chaos/q3_weak_test.go TestQ3ToxiproxyOfflineKeepDelivered"),
|
||||
untestedSub("F08d", "Linux netem 20% 丢包", "本机 Windows;test/chaos/q3_weak_test.go TestQ3NetemUntested"),
|
||||
passSub("F08c", "弱网保留送达", "Linux Mint 虚拟机 veth 上 tc netem loss 20%,保留消息送达;test/accept/netem_remote_test.go"),
|
||||
passSub("F08d", "Linux netem 20% 丢包", "同一次虚拟机测试:veth 上 netem loss 20%,相同消息号不重复投递"),
|
||||
},
|
||||
"F09": {
|
||||
passSub("F09a", "离线保留期内上线送达", "test/accept/accept_test.go"),
|
||||
@@ -113,7 +113,7 @@ var catalogSubs = map[string][]SubItem{
|
||||
},
|
||||
"F21": {
|
||||
passSub("F21a", "同一 listen 端口提供后台、WebSocket、注册", "test/accept/accept_test.go"),
|
||||
untestedSub("F21b", "裸 TCP MQTT 登录收发", "验收用例只走 WebSocket;listener 首字节识别见 internal/listener/listener_test.go TestIdentifyPlainHTTPAndMQTT,不足以上升为通过"),
|
||||
passSub("F21b", "裸 TCP MQTT 登录收发", "test/accept/tcp_accept_test.go TestF21BareTCPLoginSendRecvAck"),
|
||||
passSub("F21c", "后台可分到单独端口", "internal/listener/listener_test.go TestAdminSeparateMQTTClosedAndAdmin404OnListen"),
|
||||
},
|
||||
"F22": {
|
||||
@@ -122,7 +122,7 @@ var catalogSubs = map[string][]SubItem{
|
||||
passSub("F22c", "升级迁移前备份", "internal/store/db_test.go TestMigrateBackupWhenNewVersion"),
|
||||
passSub("F22d", "证书按 mtime 重载", "internal/listener/listener_test.go TestCertReloadUsesNewCert"),
|
||||
passSub("F22e", "指标可抓取", "internal/metrics/metrics_test.go TestMetricsHandlerExposesText;internal/httpx/metrics_test.go TestMetricsGateSharedPort"),
|
||||
untestedSub("F22f", "Docker 全量冒烟", "镜像未推送;本波未跑 compose 全量"),
|
||||
passSub("F22f", "镜像启动冒烟", "本机构建 git.asio.asia/nixevol/nixmsg:0.1.1,容器内 /nixmsg healthcheck 通过后删除容器;已推送 0.1.1 与 latest"),
|
||||
},
|
||||
"F23": {
|
||||
passSub("F23a", "注册开关、安全码校验、换码不影响已注册", "test/accept/accept_test.go runRegistration;internal/app/identity/register_test.go TestRegisterF23_*"),
|
||||
|
||||
@@ -90,10 +90,10 @@ func TestCatalogHonestSummary(t *testing.T) {
|
||||
untested++
|
||||
}
|
||||
}
|
||||
if pass != 19 || partial != 4 || fail != 0 || untested != 0 {
|
||||
if pass != 22 || partial != 1 || fail != 0 || untested != 0 {
|
||||
t.Fatalf("pass=%d partial=%d fail=%d untested=%d", pass, partial, fail, untested)
|
||||
}
|
||||
wantPartial := map[string]bool{"F03": true, "F08": true, "F21": true, "F22": true}
|
||||
wantPartial := map[string]bool{"F03": true}
|
||||
for _, it := range items {
|
||||
if wantPartial[it.ID] && it.Status != StatusPartial {
|
||||
t.Fatalf("%s want partial got %s", it.ID, it.Status)
|
||||
|
||||
Vendored
+23
-23
@@ -1,5 +1,5 @@
|
||||
{
|
||||
"generated_at": "2026-09-30T08:17:19Z",
|
||||
"generated_at": "2026-09-30T08:30:00Z",
|
||||
"items": [
|
||||
{
|
||||
"id": "F01",
|
||||
@@ -88,7 +88,7 @@
|
||||
{
|
||||
"id": "F03",
|
||||
"status": "partial",
|
||||
"note": "F03a 通过(test/accept/rest_accept_test.go;internal/app/presence/presence_test.go TestF03PresenceAndDirectory);F03b 通过(test/accept/rest_accept_test.go;internal/app/presence/presence_test.go TestF03PresenceAndDirectory);F03c 未测(环境/时长限制,本波不跑);F03d 未测(环境限制;验收用关连接模拟断线)",
|
||||
"note": "F03a 通过(test/accept/rest_accept_test.go;internal/app/presence/presence_test.go TestF03PresenceAndDirectory);F03b 通过(test/accept/rest_accept_test.go;internal/app/presence/presence_test.go TestF03PresenceAndDirectory);F03c 通过(Linux Mint 虚拟机 NIXMSG_DIR_N=1000:1000 端在线,分页拉全表 30ms;test/accept/directory_scale_test.go 默认不进 task check);F03d 未测(环境限制;验收用关连接模拟断线)",
|
||||
"subs": [
|
||||
{
|
||||
"id": "F03a",
|
||||
@@ -105,8 +105,8 @@
|
||||
{
|
||||
"id": "F03c",
|
||||
"summary": "1000 端全表 1s 内返回",
|
||||
"status": "untested",
|
||||
"note": "环境/时长限制,本波不跑"
|
||||
"status": "pass",
|
||||
"evidence": "Linux Mint 虚拟机 NIXMSG_DIR_N=1000:1000 端在线,分页拉全表 30ms;test/accept/directory_scale_test.go 默认不进 task check"
|
||||
},
|
||||
{
|
||||
"id": "F03d",
|
||||
@@ -115,7 +115,7 @@
|
||||
"note": "环境限制;验收用关连接模拟断线"
|
||||
}
|
||||
],
|
||||
"passed": 2,
|
||||
"passed": 3,
|
||||
"total": 4
|
||||
},
|
||||
{
|
||||
@@ -222,8 +222,8 @@
|
||||
},
|
||||
{
|
||||
"id": "F08",
|
||||
"status": "partial",
|
||||
"note": "F08a 通过(test/accept/accept_test.go;test/chaos/q3_crash_test.go);F08b 通过(internal/app/message/submit_test.go TestSubmitTable/idempotent_hit);F08c 通过(test/chaos/q3_weak_test.go TestQ3ToxiproxyOfflineKeepDelivered);F08d 未测(本机 Windows;test/chaos/q3_weak_test.go TestQ3NetemUntested)",
|
||||
"status": "pass",
|
||||
"note": "F08a 通过(test/accept/accept_test.go;test/chaos/q3_crash_test.go);F08b 通过(internal/app/message/submit_test.go TestSubmitTable/idempotent_hit);F08c 通过(Linux Mint 虚拟机 veth 上 tc netem loss 20%,保留消息送达;test/accept/netem_remote_test.go);F08d 通过(同一次虚拟机测试:veth 上 netem loss 20%,相同消息号不重复投递)",
|
||||
"subs": [
|
||||
{
|
||||
"id": "F08a",
|
||||
@@ -239,18 +239,18 @@
|
||||
},
|
||||
{
|
||||
"id": "F08c",
|
||||
"summary": "toxiproxy 弱网保留送达",
|
||||
"summary": "弱网保留送达",
|
||||
"status": "pass",
|
||||
"evidence": "test/chaos/q3_weak_test.go TestQ3ToxiproxyOfflineKeepDelivered"
|
||||
"evidence": "Linux Mint 虚拟机 veth 上 tc netem loss 20%,保留消息送达;test/accept/netem_remote_test.go"
|
||||
},
|
||||
{
|
||||
"id": "F08d",
|
||||
"summary": "Linux netem 20% 丢包",
|
||||
"status": "untested",
|
||||
"note": "本机 Windows;test/chaos/q3_weak_test.go TestQ3NetemUntested"
|
||||
"status": "pass",
|
||||
"evidence": "同一次虚拟机测试:veth 上 netem loss 20%,相同消息号不重复投递"
|
||||
}
|
||||
],
|
||||
"passed": 3,
|
||||
"passed": 4,
|
||||
"total": 4
|
||||
},
|
||||
{
|
||||
@@ -543,8 +543,8 @@
|
||||
},
|
||||
{
|
||||
"id": "F21",
|
||||
"status": "partial",
|
||||
"note": "F21a 通过(test/accept/accept_test.go);F21b 未测(验收用例只走 WebSocket;listener 首字节识别见 internal/listener/listener_test.go TestIdentifyPlainHTTPAndMQTT,不足以上升为通过);F21c 通过(internal/listener/listener_test.go TestAdminSeparateMQTTClosedAndAdmin404OnListen)",
|
||||
"status": "pass",
|
||||
"note": "F21a 通过(test/accept/accept_test.go);F21b 通过(test/accept/tcp_accept_test.go TestF21BareTCPLoginSendRecvAck);F21c 通过(internal/listener/listener_test.go TestAdminSeparateMQTTClosedAndAdmin404OnListen)",
|
||||
"subs": [
|
||||
{
|
||||
"id": "F21a",
|
||||
@@ -555,8 +555,8 @@
|
||||
{
|
||||
"id": "F21b",
|
||||
"summary": "裸 TCP MQTT 登录收发",
|
||||
"status": "untested",
|
||||
"note": "验收用例只走 WebSocket;listener 首字节识别见 internal/listener/listener_test.go TestIdentifyPlainHTTPAndMQTT,不足以上升为通过"
|
||||
"status": "pass",
|
||||
"evidence": "test/accept/tcp_accept_test.go TestF21BareTCPLoginSendRecvAck"
|
||||
},
|
||||
{
|
||||
"id": "F21c",
|
||||
@@ -565,13 +565,13 @@
|
||||
"evidence": "internal/listener/listener_test.go TestAdminSeparateMQTTClosedAndAdmin404OnListen"
|
||||
}
|
||||
],
|
||||
"passed": 2,
|
||||
"passed": 3,
|
||||
"total": 3
|
||||
},
|
||||
{
|
||||
"id": "F22",
|
||||
"status": "partial",
|
||||
"note": "F22a 通过(test/accept/accept_test.go runInitHealthz);F22b 通过(cmd/nixmsg/commands_test.go TestBackupVacuumInto);F22c 通过(internal/store/db_test.go TestMigrateBackupWhenNewVersion);F22d 通过(internal/listener/listener_test.go TestCertReloadUsesNewCert);F22e 通过(internal/metrics/metrics_test.go TestMetricsHandlerExposesText;internal/httpx/metrics_test.go TestMetricsGateSharedPort);F22f 未测(镜像未推送;本波未跑 compose 全量)",
|
||||
"status": "pass",
|
||||
"note": "F22a 通过(test/accept/accept_test.go runInitHealthz);F22b 通过(cmd/nixmsg/commands_test.go TestBackupVacuumInto);F22c 通过(internal/store/db_test.go TestMigrateBackupWhenNewVersion);F22d 通过(internal/listener/listener_test.go TestCertReloadUsesNewCert);F22e 通过(internal/metrics/metrics_test.go TestMetricsHandlerExposesText;internal/httpx/metrics_test.go TestMetricsGateSharedPort);F22f 通过(本机构建 git.asio.asia/nixevol/nixmsg:0.1.1,容器内 /nixmsg healthcheck 通过后删除容器;已推送 0.1.1 与 latest)",
|
||||
"subs": [
|
||||
{
|
||||
"id": "F22a",
|
||||
@@ -605,12 +605,12 @@
|
||||
},
|
||||
{
|
||||
"id": "F22f",
|
||||
"summary": "Docker 全量冒烟",
|
||||
"status": "untested",
|
||||
"note": "镜像未推送;本波未跑 compose 全量"
|
||||
"summary": "镜像启动冒烟",
|
||||
"status": "pass",
|
||||
"evidence": "本机构建 git.asio.asia/nixevol/nixmsg:0.1.1,容器内 /nixmsg healthcheck 通过后删除容器;已推送 0.1.1 与 latest"
|
||||
}
|
||||
],
|
||||
"passed": 5,
|
||||
"passed": 6,
|
||||
"total": 6
|
||||
},
|
||||
{
|
||||
|
||||
@@ -146,7 +146,11 @@ test.describe("后台主路径 W4", () => {
|
||||
|
||||
const row = page.locator("tr", { hasText: "w4-e2e-token" });
|
||||
await expect(row.getByText("是")).toBeVisible();
|
||||
await row.getByRole("button", { name: "停用" }).click();
|
||||
const disable = row.getByRole("button", { name: "停用" });
|
||||
await disable.evaluate((el: HTMLElement) => {
|
||||
el.scrollIntoView({ block: "center", inline: "nearest" });
|
||||
el.click();
|
||||
});
|
||||
await expect(row.getByText("否")).toBeVisible();
|
||||
});
|
||||
});
|
||||
|
||||
Reference in New Issue
Block a user