style: 合入后修 gofmt/govet 与 RegistrationView 类型检查
This commit is contained in:
@@ -178,10 +178,10 @@ func TestAdminSetPasswordClearsSessions(t *testing.T) {
|
|||||||
t.Fatalf("me before set-password status=%d", me.StatusCode)
|
t.Fatalf("me before set-password status=%d", me.StatusCode)
|
||||||
}
|
}
|
||||||
|
|
||||||
if err := cmdAdminSetPassword([]string{"--password", "long-enough-password"}); err != nil {
|
if setErr := cmdAdminSetPassword([]string{"--password", "long-enough-password"}); setErr != nil {
|
||||||
cancel()
|
cancel()
|
||||||
<-errCh
|
<-errCh
|
||||||
t.Fatal(err)
|
t.Fatal(setErr)
|
||||||
}
|
}
|
||||||
|
|
||||||
me2, err := client.Get(base + "/api/admin/me")
|
me2, err := client.Get(base + "/api/admin/me")
|
||||||
|
|||||||
@@ -4,6 +4,9 @@ import (
|
|||||||
"context"
|
"context"
|
||||||
"errors"
|
"errors"
|
||||||
"net/http"
|
"net/http"
|
||||||
|
"net/http/cookiejar"
|
||||||
|
"net/http/httptest"
|
||||||
|
"path/filepath"
|
||||||
"sync"
|
"sync"
|
||||||
"testing"
|
"testing"
|
||||||
|
|
||||||
@@ -11,9 +14,6 @@ import (
|
|||||||
"git.asio.asia/nixevol/NixMsg/internal/app/identity"
|
"git.asio.asia/nixevol/NixMsg/internal/app/identity"
|
||||||
"git.asio.asia/nixevol/NixMsg/internal/auth"
|
"git.asio.asia/nixevol/NixMsg/internal/auth"
|
||||||
"git.asio.asia/nixevol/NixMsg/internal/store"
|
"git.asio.asia/nixevol/NixMsg/internal/store"
|
||||||
"net/http/cookiejar"
|
|
||||||
"net/http/httptest"
|
|
||||||
"path/filepath"
|
|
||||||
)
|
)
|
||||||
|
|
||||||
type failDisableIdentity struct {
|
type failDisableIdentity struct {
|
||||||
|
|||||||
@@ -53,11 +53,11 @@ func TestLoginDeletesExpiredAdminSessions(t *testing.T) {
|
|||||||
})
|
})
|
||||||
srv := httptest.NewServer(h)
|
srv := httptest.NewServer(h)
|
||||||
t.Cleanup(srv.Close)
|
t.Cleanup(srv.Close)
|
||||||
if err := db.Queue.Do(context.Background(), func(tx *sql.Tx) error {
|
if qerr := db.Queue.Do(context.Background(), func(tx *sql.Tx) error {
|
||||||
_, e := tx.Exec(`INSERT INTO admin_sessions(token_hash, created_at, expires_at) VALUES ('expired-hash', 1, 1)`)
|
_, e := tx.Exec(`INSERT INTO admin_sessions(token_hash, created_at, expires_at) VALUES ('expired-hash', 1, 1)`)
|
||||||
return e
|
return e
|
||||||
}); err != nil {
|
}); qerr != nil {
|
||||||
t.Fatal(err)
|
t.Fatal(qerr)
|
||||||
}
|
}
|
||||||
jar, err := cookiejar.New(nil)
|
jar, err := cookiejar.New(nil)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
|
|||||||
@@ -1,7 +1,6 @@
|
|||||||
package admin
|
package admin
|
||||||
|
|
||||||
import (
|
import (
|
||||||
"errors"
|
|
||||||
"log/slog"
|
"log/slog"
|
||||||
"net"
|
"net"
|
||||||
"net/http"
|
"net/http"
|
||||||
@@ -148,9 +147,7 @@ func New(d Deps) *Handler {
|
|||||||
// ServeHTTP 实现 http.Handler。按路由限制请求体大小与读截止时间。
|
// ServeHTTP 实现 http.Handler。按路由限制请求体大小与读截止时间。
|
||||||
func (h *Handler) ServeHTTP(w http.ResponseWriter, r *http.Request) {
|
func (h *Handler) ServeHTTP(w http.ResponseWriter, r *http.Request) {
|
||||||
limit, readFor := requestBodyBudget(r)
|
limit, readFor := requestBodyBudget(r)
|
||||||
if err := http.NewResponseController(w).SetReadDeadline(time.Now().Add(readFor)); err != nil && !errors.Is(err, http.ErrNotSupported) {
|
_ = http.NewResponseController(w).SetReadDeadline(time.Now().Add(readFor))
|
||||||
// 测试用 ResponseRecorder 或不支持截止时间的封装:忽略。
|
|
||||||
}
|
|
||||||
if r.Body != nil {
|
if r.Body != nil {
|
||||||
r.Body = http.MaxBytesReader(w, r.Body, limit)
|
r.Body = http.MaxBytesReader(w, r.Body, limit)
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -302,9 +302,9 @@ func (a *App) ReceiptAck(ctx context.Context, endpointID string, req *protocol.R
|
|||||||
}
|
}
|
||||||
var acked bool
|
var acked bool
|
||||||
err = a.db.Queue.Do(ctx, func(tx *sql.Tx) error {
|
err = a.db.Queue.Do(ctx, func(tx *sql.Tx) error {
|
||||||
res, err := tx.Exec(`UPDATE receipts SET acked = 1 WHERE receipt_id = ? AND sender_id = ? AND acked = 0`, rid, endpointID)
|
res, execErr := tx.Exec(`UPDATE receipts SET acked = 1 WHERE receipt_id = ? AND sender_id = ? AND acked = 0`, rid, endpointID)
|
||||||
if err != nil {
|
if execErr != nil {
|
||||||
return err
|
return execErr
|
||||||
}
|
}
|
||||||
aff, _ := res.RowsAffected()
|
aff, _ := res.RowsAffected()
|
||||||
acked = aff > 0
|
acked = aff > 0
|
||||||
|
|||||||
@@ -191,15 +191,15 @@ LIMIT ?`, endpointID, room)
|
|||||||
var items []pushItem
|
var items []pushItem
|
||||||
for rows.Next() {
|
for rows.Next() {
|
||||||
var it pushItem
|
var it pushItem
|
||||||
if err := rows.Scan(&it.seq, &it.sendAt, &it.keep, &it.expireAt, &it.msgID, &it.senderID, &it.destKind, &it.destID, &it.meta, &it.contentType, &it.bodyEnc, &it.body); err != nil {
|
if scanErr := rows.Scan(&it.seq, &it.sendAt, &it.keep, &it.expireAt, &it.msgID, &it.senderID, &it.destKind, &it.destID, &it.meta, &it.contentType, &it.bodyEnc, &it.body); scanErr != nil {
|
||||||
_ = rows.Close()
|
_ = rows.Close()
|
||||||
return err
|
return scanErr
|
||||||
}
|
}
|
||||||
items = append(items, it)
|
items = append(items, it)
|
||||||
}
|
}
|
||||||
_ = rows.Close()
|
_ = rows.Close()
|
||||||
if err := rows.Err(); err != nil {
|
if rowsErr := rows.Err(); rowsErr != nil {
|
||||||
return err
|
return rowsErr
|
||||||
}
|
}
|
||||||
|
|
||||||
var toClaim []pushItem
|
var toClaim []pushItem
|
||||||
@@ -217,9 +217,9 @@ LIMIT ?`, endpointID, room)
|
|||||||
Meta: decodeMetaJSON(it.meta),
|
Meta: decodeMetaJSON(it.meta),
|
||||||
SendAtMs: it.sendAt,
|
SendAtMs: it.sendAt,
|
||||||
}
|
}
|
||||||
payload, err := protocol.Marshal(msg)
|
payload, marshErr := protocol.Marshal(msg)
|
||||||
if err != nil {
|
if marshErr != nil {
|
||||||
return err
|
return marshErr
|
||||||
}
|
}
|
||||||
limit := effectivePayloadLimit(live.MaxPacketSize, live.MaxReceiveBytes)
|
limit := effectivePayloadLimit(live.MaxPacketSize, live.MaxReceiveBytes)
|
||||||
if limit > 0 && len(payload) > limit {
|
if limit > 0 && len(payload) > limit {
|
||||||
|
|||||||
@@ -198,14 +198,14 @@ func TestC03WALShrinksAfterPurge(t *testing.T) {
|
|||||||
ctx := context.Background()
|
ctx := context.Background()
|
||||||
nowMs := int64(1_700_000_000_000)
|
nowMs := int64(1_700_000_000_000)
|
||||||
err = db.Queue.Do(ctx, func(tx *sql.Tx) error {
|
err = db.Queue.Do(ctx, func(tx *sql.Tx) error {
|
||||||
if _, err := tx.Exec(`INSERT INTO endpoints(id,name,login_hash,enabled,created_at) VALUES('alice','a','x',1,?)`, nowMs); err != nil {
|
if _, insErr := tx.Exec(`INSERT INTO endpoints(id,name,login_hash,enabled,created_at) VALUES('alice','a','x',1,?)`, nowMs); insErr != nil {
|
||||||
return err
|
return insErr
|
||||||
}
|
}
|
||||||
for i := 0; i < 200; i++ {
|
for i := 0; i < 200; i++ {
|
||||||
if _, err := tx.Exec(`INSERT INTO messages(id, sender_id, dest_kind, dest_id, meta, content_type, body_enc, send_at, keep, ttl_seconds, receipt, state, reason, created_at)
|
if _, insErr := tx.Exec(`INSERT INTO messages(id, sender_id, dest_kind, dest_id, meta, content_type, body_enc, send_at, keep, ttl_seconds, receipt, state, reason, created_at)
|
||||||
VALUES(?,?, 'endpoint','bob','{}','text/plain','utf8',?,0,0,0,'completed','',?)`,
|
VALUES(?,?, 'endpoint','bob','{}','text/plain','utf8',?,0,0,0,'completed','',?)`,
|
||||||
"m"+itoa(i), "alice", nowMs-10*24*3600*1000, nowMs-10*24*3600*1000); err != nil {
|
"m"+itoa(i), "alice", nowMs-10*24*3600*1000, nowMs-10*24*3600*1000); insErr != nil {
|
||||||
return err
|
return insErr
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
return nil
|
return nil
|
||||||
@@ -217,8 +217,8 @@ VALUES(?,?, 'endpoint','bob','{}','text/plain','utf8',?,0,0,0,'completed','',?)`
|
|||||||
before, _ := os.Stat(walPath)
|
before, _ := os.Stat(walPath)
|
||||||
app := New(db, Limits{RecordRetentionDays: 1, ReceiptRetentionDays: 1, IdempotencyHours: 1}, nil)
|
app := New(db, Limits{RecordRetentionDays: 1, ReceiptRetentionDays: 1, IdempotencyHours: 1}, nil)
|
||||||
app.nowFn = func() time.Time { return time.UnixMilli(nowMs) }
|
app.nowFn = func() time.Time { return time.UnixMilli(nowMs) }
|
||||||
if err := app.PurgeOnce(ctx, nowMs); err != nil {
|
if purgeErr := app.PurgeOnce(ctx, nowMs); purgeErr != nil {
|
||||||
t.Fatal(err)
|
t.Fatal(purgeErr)
|
||||||
}
|
}
|
||||||
after, err := os.Stat(walPath)
|
after, err := os.Stat(walPath)
|
||||||
if err == nil && before != nil && after.Size() > before.Size() {
|
if err == nil && before != nil && after.Size() > before.Size() {
|
||||||
|
|||||||
@@ -508,7 +508,9 @@ WHERE m.id='late-1' AND d.endpoint_id='bob'`).Scan(&bobN); err != nil {
|
|||||||
app, db := openTestApp(t, lim)
|
app, db := openTestApp(t, lim)
|
||||||
insertEndpoint(t, db, "alice", "", 1, 0)
|
insertEndpoint(t, db, "alice", "", 1, 0)
|
||||||
insertEndpoint(t, db, "bob", "", 1, 0)
|
insertEndpoint(t, db, "bob", "", 1, 0)
|
||||||
if !app.AllowRequest("alice") || !app.AllowRequest("alice") {
|
first := app.AllowRequest("alice")
|
||||||
|
second := app.AllowRequest("alice")
|
||||||
|
if !first || !second {
|
||||||
t.Fatal("burst should allow first two")
|
t.Fatal("burst should allow first two")
|
||||||
}
|
}
|
||||||
if app.AllowRequest("alice") {
|
if app.AllowRequest("alice") {
|
||||||
@@ -544,9 +546,9 @@ func TestU03SubmitTalkLockNoIP(t *testing.T) {
|
|||||||
for i := 0; i < 10; i++ {
|
for i := 0; i < 10; i++ {
|
||||||
bad := baseSend(fmt.Sprintf("w%d", i), "bob")
|
bad := baseSend(fmt.Sprintf("w%d", i), "bob")
|
||||||
bad.TalkPassword = "wrong"
|
bad.TalkPassword = "wrong"
|
||||||
_, err := app.Submit(ctx, "alice", port.ConnInfo{RemoteIP: fmt.Sprintf("10.0.0.%d", i+1)}, bad)
|
_, subErr := app.Submit(ctx, "alice", port.ConnInfo{RemoteIP: fmt.Sprintf("10.0.0.%d", i+1)}, bad)
|
||||||
if protoCode(err) != protocol.CodeTalkPasswordInvalid {
|
if protoCode(subErr) != protocol.CodeTalkPasswordInvalid {
|
||||||
t.Fatalf("i=%d got %v", i, err)
|
t.Fatalf("i=%d got %v", i, subErr)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
empty := baseSend("empty", "bob")
|
empty := baseSend("empty", "bob")
|
||||||
|
|||||||
@@ -142,4 +142,4 @@ VALUES(?,?,?,0,'accepted','',?)`, seq, "bob", nowMs, nowMs)
|
|||||||
if n != 1 {
|
if n != 1 {
|
||||||
t.Fatalf("recently completed message should remain, n=%d", n)
|
t.Fatalf("recently completed message should remain, n=%d", n)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -127,7 +127,6 @@ type connState struct {
|
|||||||
maxPacketSize uint32
|
maxPacketSize uint32
|
||||||
maxRecvBytes int
|
maxRecvBytes int
|
||||||
authOK bool
|
authOK bool
|
||||||
authErr error
|
|
||||||
sessionToken string
|
sessionToken string
|
||||||
handshook bool
|
handshook bool
|
||||||
subscribedDown bool
|
subscribedDown bool
|
||||||
|
|||||||
@@ -132,10 +132,7 @@ func (st *connState) sendOne(b *Broker, item downItem) {
|
|||||||
topic := downTopic(st.endpointID)
|
topic := downTopic(st.endpointID)
|
||||||
before := st.sentPub.Load()
|
before := st.sentPub.Load()
|
||||||
var err error
|
var err error
|
||||||
for {
|
for !b.closed.Load() {
|
||||||
if b.closed.Load() {
|
|
||||||
break
|
|
||||||
}
|
|
||||||
select {
|
select {
|
||||||
case <-st.downStop:
|
case <-st.downStop:
|
||||||
err = ErrNoConnection
|
err = ErrNoConnection
|
||||||
|
|||||||
@@ -33,8 +33,8 @@ func startPlain(t *testing.T, opts Options) *Server {
|
|||||||
}
|
}
|
||||||
ctx, cancel := context.WithCancel(context.Background())
|
ctx, cancel := context.WithCancel(context.Background())
|
||||||
t.Cleanup(cancel)
|
t.Cleanup(cancel)
|
||||||
if err := s.Start(ctx); err != nil {
|
if startErr := s.Start(ctx); startErr != nil {
|
||||||
t.Fatal(err)
|
t.Fatal(startErr)
|
||||||
}
|
}
|
||||||
t.Cleanup(func() { _ = s.Close() })
|
t.Cleanup(func() { _ = s.Close() })
|
||||||
return s
|
return s
|
||||||
@@ -62,8 +62,8 @@ func TestTLSHandshakeTimeout(t *testing.T) {
|
|||||||
}
|
}
|
||||||
ctx, cancel := context.WithCancel(context.Background())
|
ctx, cancel := context.WithCancel(context.Background())
|
||||||
defer cancel()
|
defer cancel()
|
||||||
if err := s.Start(ctx); err != nil {
|
if startErr := s.Start(ctx); startErr != nil {
|
||||||
t.Fatal(err)
|
t.Fatal(startErr)
|
||||||
}
|
}
|
||||||
defer func() { _ = s.Close() }()
|
defer func() { _ = s.Close() }()
|
||||||
|
|
||||||
@@ -72,8 +72,8 @@ func TestTLSHandshakeTimeout(t *testing.T) {
|
|||||||
t.Fatal(err)
|
t.Fatal(err)
|
||||||
}
|
}
|
||||||
defer func() { _ = c.Close() }()
|
defer func() { _ = c.Close() }()
|
||||||
if _, err := c.Write([]byte{0x16}); err != nil {
|
if _, werr := c.Write([]byte{0x16}); werr != nil {
|
||||||
t.Fatal(err)
|
t.Fatal(werr)
|
||||||
}
|
}
|
||||||
start := time.Now()
|
start := time.Now()
|
||||||
_ = c.SetReadDeadline(time.Now().Add(2 * time.Second))
|
_ = c.SetReadDeadline(time.Now().Add(2 * time.Second))
|
||||||
@@ -101,8 +101,8 @@ func TestMQTTConnectHeaderTimeout(t *testing.T) {
|
|||||||
t.Fatal(err)
|
t.Fatal(err)
|
||||||
}
|
}
|
||||||
defer func() { _ = c.Close() }()
|
defer func() { _ = c.Close() }()
|
||||||
if _, err := c.Write([]byte{0x10}); err != nil {
|
if _, werr := c.Write([]byte{0x10}); werr != nil {
|
||||||
t.Fatal(err)
|
t.Fatal(werr)
|
||||||
}
|
}
|
||||||
start := time.Now()
|
start := time.Now()
|
||||||
_ = c.SetReadDeadline(time.Now().Add(2 * time.Second))
|
_ = c.SetReadDeadline(time.Now().Add(2 * time.Second))
|
||||||
@@ -134,8 +134,8 @@ func TestMQTTConnectRemainingTooLarge(t *testing.T) {
|
|||||||
t.Fatal(err)
|
t.Fatal(err)
|
||||||
}
|
}
|
||||||
defer func() { _ = c.Close() }()
|
defer func() { _ = c.Close() }()
|
||||||
if _, err := c.Write([]byte{0x10, 0xFF, 0xFF, 0x2F}); err != nil {
|
if _, werr := c.Write([]byte{0x10, 0xFF, 0xFF, 0x2F}); werr != nil {
|
||||||
t.Fatal(err)
|
t.Fatal(werr)
|
||||||
}
|
}
|
||||||
_ = c.SetReadDeadline(time.Now().Add(2 * time.Second))
|
_ = c.SetReadDeadline(time.Now().Add(2 * time.Second))
|
||||||
buf := make([]byte, 1)
|
buf := make([]byte, 1)
|
||||||
@@ -157,8 +157,8 @@ func TestHTTPIdleTimeoutCloses(t *testing.T) {
|
|||||||
t.Fatal(err)
|
t.Fatal(err)
|
||||||
}
|
}
|
||||||
defer func() { _ = c.Close() }()
|
defer func() { _ = c.Close() }()
|
||||||
if _, err := c.Write([]byte("GET /healthz HTTP/1.1\r\nHost: x\r\n\r\n")); err != nil {
|
if _, werr := c.Write([]byte("GET /healthz HTTP/1.1\r\nHost: x\r\n\r\n")); werr != nil {
|
||||||
t.Fatal(err)
|
t.Fatal(werr)
|
||||||
}
|
}
|
||||||
_ = c.SetReadDeadline(time.Now().Add(2 * time.Second))
|
_ = c.SetReadDeadline(time.Now().Add(2 * time.Second))
|
||||||
buf := make([]byte, 4096)
|
buf := make([]byte, 4096)
|
||||||
@@ -192,8 +192,8 @@ func TestPreHandshakeLimitDropsNewConn(t *testing.T) {
|
|||||||
}
|
}
|
||||||
ctx, cancel := context.WithCancel(context.Background())
|
ctx, cancel := context.WithCancel(context.Background())
|
||||||
defer cancel()
|
defer cancel()
|
||||||
if err := s.Start(ctx); err != nil {
|
if startErr := s.Start(ctx); startErr != nil {
|
||||||
t.Fatal(err)
|
t.Fatal(startErr)
|
||||||
}
|
}
|
||||||
defer func() { _ = s.Close() }()
|
defer func() { _ = s.Close() }()
|
||||||
|
|
||||||
@@ -202,8 +202,8 @@ func TestPreHandshakeLimitDropsNewConn(t *testing.T) {
|
|||||||
t.Fatal(err)
|
t.Fatal(err)
|
||||||
}
|
}
|
||||||
defer func() { _ = c1.Close() }()
|
defer func() { _ = c1.Close() }()
|
||||||
if _, err := c1.Write([]byte{0x16}); err != nil {
|
if _, werr := c1.Write([]byte{0x16}); werr != nil {
|
||||||
t.Fatal(err)
|
t.Fatal(werr)
|
||||||
}
|
}
|
||||||
deadline := time.Now().Add(time.Second)
|
deadline := time.Now().Add(time.Second)
|
||||||
for time.Now().Before(deadline) {
|
for time.Now().Before(deadline) {
|
||||||
@@ -256,8 +256,8 @@ func TestHTTPRequestTLSState(t *testing.T) {
|
|||||||
}
|
}
|
||||||
ctx, cancel := context.WithCancel(context.Background())
|
ctx, cancel := context.WithCancel(context.Background())
|
||||||
defer cancel()
|
defer cancel()
|
||||||
if err := s.Start(ctx); err != nil {
|
if startErr := s.Start(ctx); startErr != nil {
|
||||||
t.Fatal(err)
|
t.Fatal(startErr)
|
||||||
}
|
}
|
||||||
defer func() { _ = s.Close() }()
|
defer func() { _ = s.Close() }()
|
||||||
cl := &http.Client{
|
cl := &http.Client{
|
||||||
@@ -324,8 +324,8 @@ func TestTLSSessionResume(t *testing.T) {
|
|||||||
}
|
}
|
||||||
ctx, cancel := context.WithCancel(context.Background())
|
ctx, cancel := context.WithCancel(context.Background())
|
||||||
defer cancel()
|
defer cancel()
|
||||||
if err := s.Start(ctx); err != nil {
|
if startErr := s.Start(ctx); startErr != nil {
|
||||||
t.Fatal(err)
|
t.Fatal(startErr)
|
||||||
}
|
}
|
||||||
defer func() { _ = s.Close() }()
|
defer func() { _ = s.Close() }()
|
||||||
|
|
||||||
|
|||||||
@@ -28,7 +28,7 @@ func TestFailingMigrationReusesBackup(t *testing.T) {
|
|||||||
name: "9999_fail.sql",
|
name: "9999_fail.sql",
|
||||||
body: "THIS IS NOT VALID SQL",
|
body: "THIS IS NOT VALID SQL",
|
||||||
}}
|
}}
|
||||||
if err := applyPending(w, dir, true, pending); err == nil {
|
if applyErr := applyPending(w, dir, true, pending); applyErr == nil {
|
||||||
t.Fatal("expected first failing migration to error")
|
t.Fatal("expected first failing migration to error")
|
||||||
}
|
}
|
||||||
backupDir := filepath.Join(dir, "backup")
|
backupDir := filepath.Join(dir, "backup")
|
||||||
@@ -40,7 +40,7 @@ func TestFailingMigrationReusesBackup(t *testing.T) {
|
|||||||
t.Fatalf("first backups=%v", dirNames(first))
|
t.Fatalf("first backups=%v", dirNames(first))
|
||||||
}
|
}
|
||||||
|
|
||||||
if err := applyPending(w, dir, true, pending); err == nil {
|
if applyErr := applyPending(w, dir, true, pending); applyErr == nil {
|
||||||
t.Fatal("expected second failing migration to error")
|
t.Fatal("expected second failing migration to error")
|
||||||
}
|
}
|
||||||
second, err := os.ReadDir(backupDir)
|
second, err := os.ReadDir(backupDir)
|
||||||
|
|||||||
@@ -208,8 +208,8 @@ ON CONFLICT(key) DO UPDATE SET value = excluded.value, updated_at = excluded.upd
|
|||||||
}()
|
}()
|
||||||
}
|
}
|
||||||
time.Sleep(2 * time.Millisecond)
|
time.Sleep(2 * time.Millisecond)
|
||||||
if err := db.Queue.Close(); err != nil {
|
if closeErr := db.Queue.Close(); closeErr != nil {
|
||||||
t.Fatal(err)
|
t.Fatal(closeErr)
|
||||||
}
|
}
|
||||||
wg.Wait()
|
wg.Wait()
|
||||||
|
|
||||||
@@ -232,21 +232,21 @@ func TestQueueCheckpointTruncatesWAL(t *testing.T) {
|
|||||||
payload := strings.Repeat("x", 4096)
|
payload := strings.Repeat("x", 4096)
|
||||||
for i := 0; i < 300; i++ {
|
for i := 0; i < 300; i++ {
|
||||||
i := i
|
i := i
|
||||||
if err := db.Queue.Do(ctx, func(tx *sql.Tx) error {
|
if doErr := db.Queue.Do(ctx, func(tx *sql.Tx) error {
|
||||||
_, e := tx.Exec(
|
_, e := tx.Exec(
|
||||||
`INSERT INTO settings(key, value, updated_at) VALUES(?, ?, ?)`,
|
`INSERT INTO settings(key, value, updated_at) VALUES(?, ?, ?)`,
|
||||||
fmt.Sprintf("wal_%d", i), payload, time.Now().UnixMilli(),
|
fmt.Sprintf("wal_%d", i), payload, time.Now().UnixMilli(),
|
||||||
)
|
)
|
||||||
return e
|
return e
|
||||||
}); err != nil {
|
}); doErr != nil {
|
||||||
t.Fatal(err)
|
t.Fatal(doErr)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
walPath := filepath.Join(dir, DBFileName+"-wal")
|
walPath := filepath.Join(dir, DBFileName+"-wal")
|
||||||
before, statErr := os.Stat(walPath)
|
before, statErr := os.Stat(walPath)
|
||||||
if err := db.Checkpoint(ctx); err != nil {
|
if cpErr := db.Checkpoint(ctx); cpErr != nil {
|
||||||
t.Fatal(err)
|
t.Fatal(cpErr)
|
||||||
}
|
}
|
||||||
if statErr != nil {
|
if statErr != nil {
|
||||||
return
|
return
|
||||||
|
|||||||
+2
-2
@@ -91,8 +91,8 @@ func Run(cfg Config) (*Report, error) {
|
|||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, err
|
return nil, err
|
||||||
}
|
}
|
||||||
if err := Provision(cfg, ids); err != nil {
|
if provErr := Provision(cfg, ids); provErr != nil {
|
||||||
return nil, err
|
return nil, provErr
|
||||||
}
|
}
|
||||||
clients, err := ConnectAll(cfg, ids)
|
clients, err := ConnectAll(cfg, ids)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
|
|||||||
+6
-6
@@ -80,8 +80,8 @@ func (c *Client) connectSubscribeHello(password string, timeout time.Duration) e
|
|||||||
if err != nil {
|
if err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
if err := c.mc.Send(pkt); err != nil {
|
if sendErr := c.mc.Send(pkt); sendErr != nil {
|
||||||
return fmt.Errorf("%s CONNECT: %w", c.ID, err)
|
return fmt.Errorf("%s CONNECT: %w", c.ID, sendErr)
|
||||||
}
|
}
|
||||||
ack, err := c.mc.Recv()
|
ack, err := c.mc.Recv()
|
||||||
if err != nil {
|
if err != nil {
|
||||||
@@ -99,11 +99,11 @@ func (c *Client) connectSubscribeHello(password string, timeout time.Duration) e
|
|||||||
if err != nil {
|
if err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
if err := c.mc.Send(sub); err != nil {
|
if sendErr := c.mc.Send(sub); sendErr != nil {
|
||||||
return fmt.Errorf("%s SUBSCRIBE: %w", c.ID, err)
|
return fmt.Errorf("%s SUBSCRIBE: %w", c.ID, sendErr)
|
||||||
}
|
}
|
||||||
if _, err := c.mc.Recv(); err != nil {
|
if _, recvErr := c.mc.Recv(); recvErr != nil {
|
||||||
return fmt.Errorf("%s SUBACK: %w", c.ID, err)
|
return fmt.Errorf("%s SUBACK: %w", c.ID, recvErr)
|
||||||
}
|
}
|
||||||
|
|
||||||
hello := protocol.Hello{V: protocol.Version, Type: protocol.TypeHello, RID: "h0", Client: "mqttbench/0.1"}
|
hello := protocol.Hello{V: protocol.Version, Type: protocol.TypeHello, RID: "h0", Client: "mqttbench/0.1"}
|
||||||
|
|||||||
@@ -108,8 +108,8 @@ func TestMQTT5ConnectFakeBroker(t *testing.T) {
|
|||||||
if err != nil {
|
if err != nil {
|
||||||
t.Fatal(err)
|
t.Fatal(err)
|
||||||
}
|
}
|
||||||
if err := mc.Send(pkt); err != nil {
|
if sendErr := mc.Send(pkt); sendErr != nil {
|
||||||
t.Fatal(err)
|
t.Fatal(sendErr)
|
||||||
}
|
}
|
||||||
ack, err := mc.Recv()
|
ack, err := mc.Recv()
|
||||||
if err != nil {
|
if err != nil {
|
||||||
|
|||||||
@@ -46,7 +46,7 @@ describe("RegistrationView", () => {
|
|||||||
attachTo: document.body,
|
attachTo: document.body,
|
||||||
});
|
});
|
||||||
await flushPromises();
|
await flushPromises();
|
||||||
expect(w.get('[data-testid="reg-enabled"]').exists()).toBe(true);
|
expect(w.find('[data-testid="reg-enabled"]').exists()).toBe(true);
|
||||||
expect(w.find('[data-testid="reg-code-missing"]').exists()).toBe(false);
|
expect(w.find('[data-testid="reg-code-missing"]').exists()).toBe(false);
|
||||||
w.unmount();
|
w.unmount();
|
||||||
});
|
});
|
||||||
|
|||||||
Reference in New Issue
Block a user