Files
NixMsg/test/chaos/q3_crash_test.go

57 lines
1.6 KiB
Go

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"})
}