package nixmsg_test import ( "context" "errors" "strings" "sync/atomic" "testing" "time" nixmsg "git.asio.asia/nixevol/NixMsg/sdk/go" ) const ( itestPass = "password12" itestCode = "sdk-go-code1" ) func TestChecklistAgainstRealServer(t *testing.T) { srv := startITestServer(t) admin := newAdminHTTP(t, srv.AdminHTTPBase, srv.AdminPassword) admin.putRegistration(t, true, itestCode) ctx := context.Background() ws := srv.MQTTWS t.Run("01_handshake_limits", func(t *testing.T) { mustRegister(t, ctx, ws, "go01a", "A") c := mustConnect(t, ctx, ws, "go01a", itestPass) defer c.Close() lim := c.Limits() if lim.MaxBodyBytes <= 0 || lim.MaxFrameBytes <= 0 || lim.ServerTimeMs <= 0 { t.Fatalf("handshake limits incomplete: %+v", lim) } if lim.MaxBodyBytes != 262144 { t.Fatalf("max_body_bytes=%d want 262144", lim.MaxBodyBytes) } var tok string c2 := nixmsg.New() c2.OnSession(func(token string) { tok = token }) if err := c2.Connect(ctx, ws, "go01a", nixmsg.Credential{Password: itestPass}, nixmsg.Options{ConnectTimeout: 20 * time.Second}); err != nil { t.Fatal(err) } defer c2.Close() if tok == "" || !strings.HasPrefix(tok, "nst_") { t.Fatalf("session token=%q", tok) } }) t.Run("02_dm_callback_once", func(t *testing.T) { mustRegister(t, ctx, ws, "go02a", "A") mustRegister(t, ctx, ws, "go02b", "B") a := mustConnect(t, ctx, ws, "go02a", itestPass) defer a.Close() b := mustConnect(t, ctx, ws, "go02b", itestPass) defer b.Close() var n atomic.Int32 got := make(chan nixmsg.Message, 4) b.OnMessage(func(msg nixmsg.Message) error { n.Add(1) got <- msg return nil }) res, err := a.Send(ctx, nixmsg.Target{Kind: "endpoint", ID: "go02b"}, nixmsg.Body{Enc: "utf8", Data: "hi-once"}, nixmsg.SendOptions{Delay: delay0()}) if err != nil { t.Fatal(err) } msg := waitMsg(t, got, 8*time.Second) if msg.ID != res.ID || msg.From != "go02a" || msg.Body.Data != "hi-once" { t.Fatalf("msg=%+v res=%+v", msg, res) } time.Sleep(500 * time.Millisecond) if n.Load() != 1 { t.Fatalf("callbacks=%d want 1", n.Load()) } }) t.Run("03_send_while_disconnected", func(t *testing.T) { mustRegister(t, ctx, ws, "go03a", "A") mustRegister(t, ctx, ws, "go03b", "B") a := mustConnect(t, ctx, ws, "go03a", itestPass) defer a.Close() b := mustConnect(t, ctx, ws, "go03b", itestPass) defer b.Close() got := make(chan nixmsg.Message, 4) var n atomic.Int32 b.OnMessage(func(msg nixmsg.Message) error { n.Add(1) got <- msg return nil }) states := make(chan nixmsg.ConnectionState, 16) a.OnConnection(func(ev nixmsg.ConnectionEvent) { states <- ev.State }) admin.kick(t, "go03a") saw := waitConnState(t, states, 15*time.Second, nixmsg.StateReconnecting, nixmsg.StateOnline) resCh := make(chan nixmsg.SendResult, 1) errCh := make(chan error, 1) go func() { res, err := a.Send(ctx, nixmsg.Target{Kind: "endpoint", ID: "go03b"}, nixmsg.Body{Enc: "utf8", Data: "queued"}, nixmsg.SendOptions{Delay: delay0(), Keep: true, ID: "go03-msg-1"}) if err != nil { errCh <- err return } resCh <- res }() if saw != nixmsg.StateOnline { _ = waitConnState(t, states, 30*time.Second, nixmsg.StateOnline) } select { case err := <-errCh: t.Fatal(err) case res := <-resCh: if res.ID != "go03-msg-1" { t.Fatalf("id=%s", res.ID) } case <-time.After(30 * time.Second): t.Fatal("send timeout") } msg := waitMsg(t, got, 15*time.Second) if msg.ID != "go03-msg-1" { t.Fatalf("msg=%+v", msg) } time.Sleep(800 * time.Millisecond) if n.Load() != 1 { t.Fatalf("callbacks=%d want 1", n.Load()) } }) t.Run("04_same_id_retry_and_dedup", func(t *testing.T) { mustRegister(t, ctx, ws, "go04a", "A") mustRegister(t, ctx, ws, "go04b", "B") a := mustConnect(t, ctx, ws, "go04a", itestPass) defer a.Close() b := mustConnect(t, ctx, ws, "go04b", itestPass) defer b.Close() // 同号:断线期间入队发送,重连后以同一消息号送达 var n atomic.Int32 got := make(chan nixmsg.Message, 4) b.OnMessage(func(msg nixmsg.Message) error { n.Add(1) got <- msg return nil }) states := make(chan nixmsg.ConnectionState, 16) a.OnConnection(func(ev nixmsg.ConnectionEvent) { states <- ev.State }) admin.kick(t, "go04a") saw := waitConnState(t, states, 15*time.Second, nixmsg.StateReconnecting, nixmsg.StateOnline) go func() { _, _ = a.Send(ctx, nixmsg.Target{Kind: "endpoint", ID: "go04b"}, nixmsg.Body{Enc: "utf8", Data: "same-id"}, nixmsg.SendOptions{Delay: delay0(), Keep: true, ID: "go04-fixed-id"}) }() if saw != nixmsg.StateOnline { _ = waitConnState(t, states, 30*time.Second, nixmsg.StateOnline) } msg := waitMsg(t, got, 15*time.Second) if msg.ID != "go04-fixed-id" { t.Fatalf("want go04-fixed-id got %s", msg.ID) } time.Sleep(500 * time.Millisecond) if n.Load() != 1 { t.Fatalf("same-id callbacks=%d", n.Load()) } // 模拟未确认:手动确认模式收一次、踢线重连后服务器重推,SDK 不重复回调;再手动 ack b2 := nixmsg.New() var n2 atomic.Int32 got2 := make(chan nixmsg.Message, 4) b2.OnMessage(func(msg nixmsg.Message) error { n2.Add(1) got2 <- msg return nil }) if err := b2.Connect(ctx, ws, "go04b", nixmsg.Credential{Password: itestPass}, nixmsg.Options{ConnectTimeout: 20 * time.Second, ManualAck: true}); err != nil { t.Fatal(err) } defer b2.Close() _, err := a.Send(ctx, nixmsg.Target{Kind: "endpoint", ID: "go04b"}, nixmsg.Body{Enc: "utf8", Data: "repush"}, nixmsg.SendOptions{Delay: delay0(), Keep: true, TTL: int64Ptr(3600), ID: "go04-repush"}) if err != nil { t.Fatal(err) } m1 := waitMsg(t, got2, 10*time.Second) if m1.ID != "go04-repush" { t.Fatalf("m1=%+v", m1) } st2 := make(chan nixmsg.ConnectionState, 16) b2.OnConnection(func(ev nixmsg.ConnectionEvent) { st2 <- ev.State }) admin.kick(t, "go04b") _ = waitConnState(t, st2, 30*time.Second, nixmsg.StateOnline) time.Sleep(2 * time.Second) if n2.Load() != 1 { t.Fatalf("repush callbacks=%d want 1 (no duplicate)", n2.Load()) } if err := b2.Ack(m1); err != nil { t.Fatal(err) } }) t.Run("05_recall_within_delay", func(t *testing.T) { mustRegister(t, ctx, ws, "go05a", "A") mustRegister(t, ctx, ws, "go05b", "B") a := mustConnect(t, ctx, ws, "go05a", itestPass) defer a.Close() b := mustConnect(t, ctx, ws, "go05b", itestPass) defer b.Close() var msgN, revN atomic.Int32 b.OnMessage(func(msg nixmsg.Message) error { msgN.Add(1) return nil }) b.OnRevoked(func(e nixmsg.RevokedEvent) { revN.Add(1) }) d := 10 * time.Second res, err := a.Send(ctx, nixmsg.Target{Kind: "endpoint", ID: "go05b"}, nixmsg.Body{Enc: "utf8", Data: "will-recall"}, nixmsg.SendOptions{Delay: &d, ID: "go05-rec"}) if err != nil { t.Fatal(err) } if res.State != "scheduled" && res.State != "" { // state may be scheduled } if _, err := a.Recall(ctx, "go05-rec"); err != nil { t.Fatal(err) } time.Sleep(1500 * time.Millisecond) if msgN.Load() != 0 { t.Fatalf("bob got msg callbacks=%d", msgN.Load()) } if revN.Load() != 0 { t.Fatalf("bob got revoked=%d (undelivered should be silent)", revN.Load()) } }) t.Run("06_scheduled_about_2s", func(t *testing.T) { mustRegister(t, ctx, ws, "go06a", "A") mustRegister(t, ctx, ws, "go06b", "B") a := mustConnect(t, ctx, ws, "go06a", itestPass) defer a.Close() b := mustConnect(t, ctx, ws, "go06b", itestPass) defer b.Close() got := make(chan time.Time, 1) b.OnMessage(func(msg nixmsg.Message) error { got <- time.Now() return nil }) d := 2 * time.Second start := time.Now() if _, err := a.Send(ctx, nixmsg.Target{Kind: "endpoint", ID: "go06b"}, nixmsg.Body{Enc: "utf8", Data: "later"}, nixmsg.SendOptions{Delay: &d}); err != nil { t.Fatal(err) } select { case at := <-got: elapsed := at.Sub(start) if elapsed < 1500*time.Millisecond || elapsed > 6*time.Second { t.Fatalf("elapsed=%v want ~2s", elapsed) } case <-time.After(10 * time.Second): t.Fatal("timeout waiting scheduled msg") } }) t.Run("07_offline_keep", func(t *testing.T) { mustRegister(t, ctx, ws, "go07a", "A") mustRegister(t, ctx, ws, "go07b", "B") mustRegister(t, ctx, ws, "go07c", "C") a := mustConnect(t, ctx, ws, "go07a", itestPass) defer a.Close() // 晚 1 秒上线能收到 ttl := int64(60) if _, err := a.Send(ctx, nixmsg.Target{Kind: "endpoint", ID: "go07b"}, nixmsg.Body{Enc: "utf8", Data: "keep-ok"}, nixmsg.SendOptions{Delay: delay0(), Keep: true, TTL: &ttl, ID: "go07-keep-ok"}); err != nil { t.Fatal(err) } time.Sleep(1 * time.Second) got := make(chan nixmsg.Message, 2) b := nixmsg.New() b.OnMessage(func(msg nixmsg.Message) error { got <- msg return nil }) if err := b.Connect(ctx, ws, "go07b", nixmsg.Credential{Password: itestPass}, nixmsg.Options{ConnectTimeout: 20 * time.Second}); err != nil { t.Fatal(err) } defer b.Close() msg := waitMsg(t, got, 10*time.Second) if msg.ID != "go07-keep-ok" { t.Fatalf("keep-ok msg=%+v", msg) } // 保留 1 秒,3 秒后上线收不到;发送方收到过期回执 receiptCh := make(chan nixmsg.Receipt, 4) a.OnReceipt(func(r nixmsg.Receipt) { receiptCh <- r }) ttl1 := int64(1) if _, err := a.Send(ctx, nixmsg.Target{Kind: "endpoint", ID: "go07c"}, nixmsg.Body{Enc: "utf8", Data: "expire"}, nixmsg.SendOptions{Delay: delay0(), Keep: true, TTL: &ttl1, ID: "go07-exp", Receipt: boolPtr(true)}); err != nil { t.Fatal(err) } time.Sleep(3 * time.Second) c := mustConnect(t, ctx, ws, "go07c", itestPass) defer c.Close() var cN atomic.Int32 c.OnMessage(func(msg nixmsg.Message) error { cN.Add(1) return nil }) time.Sleep(2 * time.Second) if cN.Load() != 0 { t.Fatalf("expired offline should not deliver, got %d", cN.Load()) } deadline := time.Now().Add(15 * time.Second) for time.Now().Before(deadline) { select { case r := <-receiptCh: if r.ID == "go07-exp" && (r.State == "expired" || r.Reason == "expired" || strings.Contains(r.State, "expir") || strings.Contains(r.Reason, "expir")) { return } case <-time.After(200 * time.Millisecond): } } // 回执可能稍慢;再扫一下 select { case r := <-receiptCh: if r.ID != "go07-exp" { t.Fatalf("unexpected receipt %+v", r) } default: t.Fatal("未收到过期回执") } }) t.Run("08_group_no_echo_to_sender", func(t *testing.T) { mustRegister(t, ctx, ws, "go08a", "A") mustRegister(t, ctx, ws, "go08b", "B") mustRegister(t, ctx, ws, "go08c", "C") a := mustConnect(t, ctx, ws, "go08a", itestPass) defer a.Close() b := mustConnect(t, ctx, ws, "go08b", itestPass) defer b.Close() c := mustConnect(t, ctx, ws, "go08c", itestPass) defer c.Close() if _, err := a.CreateGroup(ctx, "g_go08", "G8", []nixmsg.GroupMemberIn{{ID: "go08b"}, {ID: "go08c"}}); err != nil { t.Fatal(err) } time.Sleep(300 * time.Millisecond) bGot := make(chan nixmsg.Message, 2) cGot := make(chan nixmsg.Message, 2) var aN atomic.Int32 a.OnMessage(func(msg nixmsg.Message) error { aN.Add(1) return nil }) b.OnMessage(func(msg nixmsg.Message) error { bGot <- msg return nil }) c.OnMessage(func(msg nixmsg.Message) error { cGot <- msg return nil }) if _, err := a.Send(ctx, nixmsg.Target{Kind: "group", ID: "g_go08"}, nixmsg.Body{Enc: "utf8", Data: "ghi"}, nixmsg.SendOptions{Delay: delay0(), ID: "go08-g1"}); err != nil { t.Fatal(err) } mb := waitMsg(t, bGot, 10*time.Second) mc := waitMsg(t, cGot, 10*time.Second) if mb.Body.Data != "ghi" || mc.Body.Data != "ghi" { t.Fatalf("b=%+v c=%+v", mb, mc) } time.Sleep(500 * time.Millisecond) if aN.Load() != 0 { t.Fatalf("sender got %d group msgs", aN.Load()) } }) t.Run("09_talk_password", func(t *testing.T) { mustRegister(t, ctx, ws, "go09a", "A") mustRegister(t, ctx, ws, "go09b", "B") a := mustConnect(t, ctx, ws, "go09a", itestPass) defer a.Close() b := mustConnect(t, ctx, ws, "go09b", itestPass) defer b.Close() if err := b.SetTalkPassword(ctx, "talk9"); err != nil { t.Fatal(err) } _, err := a.Send(ctx, nixmsg.Target{Kind: "endpoint", ID: "go09b"}, nixmsg.Body{Enc: "utf8", Data: "x"}, nixmsg.SendOptions{Delay: delay0()}) if err == nil { t.Fatal("want talk_password_required") } var ae *nixmsg.APIError if !errors.As(err, &ae) || ae.Code != "talk_password_required" { t.Fatalf("err=%v", err) } if err := a.Unlock(ctx, "go09b", "talk9"); err != nil { t.Fatal(err) } if _, err := a.Send(ctx, nixmsg.Target{Kind: "endpoint", ID: "go09b"}, nixmsg.Body{Enc: "utf8", Data: "ok"}, nixmsg.SendOptions{Delay: delay0()}); err != nil { t.Fatal(err) } if err := b.SetTalkPassword(ctx, "talk9b"); err != nil { t.Fatal(err) } _, err = a.Send(ctx, nixmsg.Target{Kind: "endpoint", ID: "go09b"}, nixmsg.Body{Enc: "utf8", Data: "stale"}, nixmsg.SendOptions{Delay: delay0()}) if err == nil { t.Fatal("want auth invalid after password change") } if !errors.As(err, &ae) || (ae.Code != "talk_password_required" && ae.Code != "talk_password_invalid") { t.Fatalf("after change err=%v", err) } // 对方先发则可回复:b 给 a 发(a 无对话密码)后,a 可回 b(即使用旧授权已失效,回复授权由 b→a 的发送产生) got := make(chan nixmsg.Message, 2) a.OnMessage(func(msg nixmsg.Message) error { got <- msg return nil }) if _, err := b.Send(ctx, nixmsg.Target{Kind: "endpoint", ID: "go09a"}, nixmsg.Body{Enc: "utf8", Data: "first"}, nixmsg.SendOptions{Delay: delay0()}); err != nil { t.Fatal(err) } _ = waitMsg(t, got, 8*time.Second) if _, err := a.Send(ctx, nixmsg.Target{Kind: "endpoint", ID: "go09b"}, nixmsg.Body{Enc: "utf8", Data: "reply"}, nixmsg.SendOptions{Delay: delay0()}); err != nil { t.Fatalf("reply after peer first send: %v", err) } }) t.Run("10_second_login_kicks_first", func(t *testing.T) { mustRegister(t, ctx, ws, "go10a", "A") c1 := mustConnect(t, ctx, ws, "go10a", itestPass) defer c1.Close() var kicked atomic.Bool c1.OnConnection(func(ev nixmsg.ConnectionEvent) { if ev.State == nixmsg.StateKicked { kicked.Store(true) } }) c2 := mustConnect(t, ctx, ws, "go10a", itestPass) defer c2.Close() deadline := time.Now().Add(15 * time.Second) for time.Now().Before(deadline) { if kicked.Load() { break } time.Sleep(50 * time.Millisecond) } if !kicked.Load() { t.Fatal("first client not kicked") } time.Sleep(2 * time.Second) // 被顶号后不应恢复 online if c1.Limits().ServerTimeMs != 0 { // Limits 仍可能保留旧值;用连接状态回调或再次 Send 探测 } _, err := c1.Send(ctx, nixmsg.Target{Kind: "endpoint", ID: "go10a"}, nixmsg.Body{Enc: "utf8", Data: "x"}, nixmsg.SendOptions{Delay: delay0()}) if err == nil { t.Fatal("kicked client should not send successfully after stop-reconnect") } }) t.Run("11_body_too_large_local", func(t *testing.T) { mustRegister(t, ctx, ws, "go11a", "A") mustRegister(t, ctx, ws, "go11b", "B") a := mustConnect(t, ctx, ws, "go11a", itestPass) defer a.Close() big := strings.Repeat("x", 262144+1) _, err := a.Send(ctx, nixmsg.Target{Kind: "endpoint", ID: "go11b"}, nixmsg.Body{Enc: "utf8", Data: big}, nixmsg.SendOptions{Delay: delay0()}) if err == nil { t.Fatal("want body_too_large") } var ae *nixmsg.APIError if !errors.As(err, &ae) || ae.Code != nixmsg.CodeBodyTooLarge { t.Fatalf("err=%v", err) } }) t.Run("12_registration_switch_and_code", func(t *testing.T) { admin.putRegistration(t, false, itestCode) _, err := nixmsg.Register(ctx, ws, itestCode, nixmsg.RegisterOptions{ID: "go12x", LoginPassword: itestPass, Name: "X"}) if err == nil { t.Fatal("want registration closed") } admin.putRegistration(t, true, itestCode) _, err = nixmsg.Register(ctx, ws, "wrong-code-xx", nixmsg.RegisterOptions{ID: "go12y", LoginPassword: itestPass, Name: "Y"}) if err == nil { t.Fatal("want bad code") } codeNew := "sdk-go-code2" admin.putRegistration(t, true, itestCode) res, err := nixmsg.Register(ctx, ws, itestCode, nixmsg.RegisterOptions{ID: "go12ok", LoginPassword: itestPass, Name: "OK"}) if err != nil { t.Fatal(err) } if res.ID != "go12ok" { t.Fatalf("id=%s", res.ID) } c := mustConnect(t, ctx, ws, "go12ok", itestPass) c.Close() admin.putRegistration(t, true, codeNew) _, err = nixmsg.Register(ctx, ws, itestCode, nixmsg.RegisterOptions{ID: "go12old", LoginPassword: itestPass, Name: "Old"}) if err == nil { t.Fatal("old code should fail") } // 已注册端仍可用旧密码登录 c2 := mustConnect(t, ctx, ws, "go12ok", itestPass) c2.Close() admin.putRegistration(t, true, itestCode) // 恢复供后续子测试 }) t.Run("13_change_login_password", func(t *testing.T) { mustRegister(t, ctx, ws, "go13a", "A") c := mustConnect(t, ctx, ws, "go13a", itestPass) if err := c.ChangeLoginPassword(ctx, itestPass, "password99"); err != nil { t.Fatal(err) } _ = c.Close() cNew := mustConnect(t, ctx, ws, "go13a", "password99") cNew.Close() cBad := nixmsg.New() var authFail atomic.Bool cBad.OnConnection(func(ev nixmsg.ConnectionEvent) { if ev.State == nixmsg.StateAuthFailed { authFail.Store(true) } }) err := cBad.Connect(ctx, ws, "go13a", nixmsg.Credential{Password: itestPass}, nixmsg.Options{ConnectTimeout: 10 * time.Second}) if err == nil { _ = cBad.Close() t.Fatal("old password should fail") } deadline := time.Now().Add(5 * time.Second) for time.Now().Before(deadline) && !authFail.Load() { time.Sleep(50 * time.Millisecond) } time.Sleep(1500 * time.Millisecond) // 不应恢复为 online if !authFail.Load() { t.Log("auth_failed event not observed; connect error present which is enough") } _ = cBad.Close() }) t.Run("15_session_token", func(t *testing.T) { mustRegister(t, ctx, ws, "go15a", "A") var tok1 string c1 := nixmsg.New() c1.OnSession(func(token string) { tok1 = token }) if err := c1.Connect(ctx, ws, "go15a", nixmsg.Credential{Password: itestPass}, nixmsg.Options{ConnectTimeout: 20 * time.Second}); err != nil { t.Fatal(err) } if tok1 == "" { t.Fatal("no token") } _ = c1.Close() cTok := nixmsg.New() if err := cTok.Connect(ctx, ws, "go15a", nixmsg.Credential{SessionToken: tok1}, nixmsg.Options{ConnectTimeout: 20 * time.Second}); err != nil { t.Fatalf("token reconnect: %v", err) } _ = cTok.Close() // 另一处密码登录使旧令牌失效 cPass := mustConnect(t, ctx, ws, "go15a", itestPass) defer cPass.Close() cOld := nixmsg.New() var inv atomic.Bool cOld.OnConnection(func(ev nixmsg.ConnectionEvent) { if ev.State == nixmsg.StateAuthFailed && (ev.Reason == string(nixmsg.AuthSessionInvalid) || strings.Contains(ev.Reason, "session")) { inv.Store(true) } if ev.State == nixmsg.StateAuthFailed { inv.Store(true) } }) err := cOld.Connect(ctx, ws, "go15a", nixmsg.Credential{SessionToken: tok1}, nixmsg.Options{ConnectTimeout: 10 * time.Second}) if err == nil { _ = cOld.Close() t.Fatal("old token should fail after other password login") } _ = cOld.Close() // logout 后令牌失效 var tok2 string c3 := nixmsg.New() c3.OnSession(func(token string) { tok2 = token }) if err := c3.Connect(ctx, ws, "go15a", nixmsg.Credential{Password: itestPass}, nixmsg.Options{ConnectTimeout: 20 * time.Second}); err != nil { t.Fatal(err) } if err := c3.Logout(ctx); err != nil { t.Fatal(err) } _ = c3.Close() c4 := nixmsg.New() err = c4.Connect(ctx, ws, "go15a", nixmsg.Credential{SessionToken: tok2}, nixmsg.Options{ConnectTimeout: 10 * time.Second}) if err == nil { _ = c4.Close() t.Fatal("token after logout should fail") } _ = c4.Close() }) } func mustRegister(t *testing.T, ctx context.Context, ws, id, name string) { t.Helper() _, err := nixmsg.Register(ctx, ws, itestCode, nixmsg.RegisterOptions{ID: id, LoginPassword: itestPass, Name: name}) if err != nil { // 可能已存在(重跑);尝试直接登录 t.Logf("register %s: %v (continue if already exists)", id, err) } } func mustConnect(t *testing.T, ctx context.Context, ws, id, password string) *nixmsg.Client { t.Helper() c := nixmsg.New() cred := nixmsg.Credential{Password: password} if err := c.Connect(ctx, ws, id, cred, nixmsg.Options{ConnectTimeout: 20 * time.Second}); err != nil { t.Fatalf("connect %s: %v", id, err) } return c } func waitMsg(t *testing.T, ch <-chan nixmsg.Message, d time.Duration) nixmsg.Message { t.Helper() select { case m := <-ch: return m case <-time.After(d): t.Fatal("timeout waiting message") return nixmsg.Message{} } } func waitConnState(t *testing.T, states <-chan nixmsg.ConnectionState, d time.Duration, want ...nixmsg.ConnectionState) nixmsg.ConnectionState { t.Helper() wants := map[nixmsg.ConnectionState]bool{} for _, w := range want { wants[w] = true } deadline := time.Now().Add(d) for time.Now().Before(deadline) { select { case st := <-states: if wants[st] { return st } case <-time.After(50 * time.Millisecond): } } t.Fatalf("timeout waiting state among %v", want) return "" } func int64Ptr(v int64) *int64 { return &v } func boolPtr(v bool) *bool { return &v }