From 3883d1b4710a59828249bcfb8367ffaf9aa2a567 Mon Sep 17 00:00:00 2001 From: Amirmahdi Nourkazemi Date: Thu, 23 Jul 2026 13:10:19 +0330 Subject: [PATCH] feat: WAL concurrency change is now running --- cmd/loadtest/main.go | 192 +++++++++++++++++++++++++++++ internal/store/concurrency_test.go | 55 +++++++++ internal/store/db.go | 14 ++- 3 files changed, 257 insertions(+), 4 deletions(-) create mode 100644 cmd/loadtest/main.go create mode 100644 internal/store/concurrency_test.go diff --git a/cmd/loadtest/main.go b/cmd/loadtest/main.go new file mode 100644 index 0000000..a2e0bf3 --- /dev/null +++ b/cmd/loadtest/main.go @@ -0,0 +1,192 @@ +// Command loadtest یک استرس‌تستِ واقعیِ WebSocket برای سرورِ حکم است: N بازیکنِ +// مصنوعی را می‌سازد، هرکدام join_queue می‌کند و بازیِ کاملِ حکم را به‌صورتِ خودکار +// (انتخابِ حکم + بازیِ کارتِ قانونی) بازی می‌کند تا بارِ واقعیِ پخشِ وضعیت تولید شود. +// +// توکن‌ها با همان JWT_SECRET سرور ساخته می‌شوند (بدونِ نیاز به OTP/SMS). کاربرها +// باید در DB وجود داشته باشند و برای tierِ انتخابی سکه‌ی کافی داشته باشند. +// +// مثال (روی نمونه‌ی محلی/استیجینگ — نه روی پروداکشنِ واقعی): +// +// go run ./cmd/loadtest -url ws://localhost:8080 -n 200 -tier beginner \ +// -secret "$JWT_SECRET" -startid 1000 -ramp 50ms +// +// خروجی هر ثانیه: اتصال‌های فعال، پیام‌های وضعیت/ثانیه، حرکت‌ها/ثانیه، خطاها. +package main + +import ( + "flag" + "fmt" + "log" + "math/rand" + "net/http" + "strconv" + "strings" + "sync" + "sync/atomic" + "time" + + "github.com/golang-jwt/jwt/v5" + "github.com/gorilla/websocket" +) + +var ( + connected int64 + stateMsgs int64 + moves int64 + errs int64 +) + +type state struct { + Type string `json:"type"` + Phase string `json:"phase"` + YourSeat int `json:"your_seat"` + Hakem int `json:"hakem"` + Turn int `json:"turn"` + TrickDone bool `json:"trick_done"` + YourHand []string `json:"your_hand"` + LeadSuit string `json:"lead_suit"` +} + +var suitWord = map[byte]string{'H': "hearts", 'S': "spades", 'D': "diamonds", 'C': "clubs"} +var wordSuit = map[string]byte{"hearts": 'H', "spades": 'S', "diamonds": 'D', "clubs": 'C'} + +func mintToken(secret string, uid int64, ttl time.Duration) string { + claims := jwt.MapClaims{ + "sub": strconv.FormatInt(uid, 10), + "abl": "play", + "exp": time.Now().Add(ttl).Unix(), + "iat": time.Now().Unix(), + } + t := jwt.NewWithClaims(jwt.SigningMethodHS256, claims) + s, _ := t.SignedString([]byte(secret)) + return s +} + +// pickCard یک کارتِ قانونی برمی‌گرداند (پیرویِ خال در صورتِ داشتن). +func pickCard(hand []string, lead string) string { + if lead != "" { + want := wordSuit[lead] + var same []string + for _, c := range hand { + if c[len(c)-1] == want { + same = append(same, c) + } + } + if len(same) > 0 { + return same[rand.Intn(len(same))] + } + } + return hand[rand.Intn(len(hand))] +} + +func runClient(wsURL, secret, tier string, uid int64, playDelay time.Duration, wg *sync.WaitGroup) { + defer wg.Done() + tok := mintToken(secret, uid, time.Hour) + u := wsURL + "/ws?token=" + tok + c, _, err := websocket.DefaultDialer.Dial(u, http.Header{}) + if err != nil { + atomic.AddInt64(&errs, 1) + return + } + atomic.AddInt64(&connected, 1) + defer func() { atomic.AddInt64(&connected, -1); c.Close() }() + + _ = c.WriteJSON(map[string]any{"type": "join_queue", "tier": tier}) + + for { + var raw map[string]any + if err := c.ReadJSON(&raw); err != nil { + atomic.AddInt64(&errs, 1) + return + } + if raw["type"] != "state" { + continue + } + atomic.AddInt64(&stateMsgs, 1) + // دوباره به‌صورتِ typed decode کن. + var s state + if b, ok := raw["your_hand"].([]any); ok { + for _, x := range b { + s.YourHand = append(s.YourHand, x.(string)) + } + } + s.Phase, _ = raw["phase"].(string) + s.LeadSuit, _ = raw["lead_suit"].(string) + s.YourSeat = intOf(raw["your_seat"]) + s.Hakem = intOf(raw["hakem"]) + s.Turn = intOf(raw["turn"]) + s.TrickDone, _ = raw["trick_done"].(bool) + + if s.Turn != s.YourSeat || s.TrickDone { + continue + } + time.Sleep(playDelay) // شبیه‌سازیِ زمانِ فکرِ انسان + switch s.Phase { + case "choose_trump": + if s.Hakem == s.YourSeat { + _ = c.WriteJSON(map[string]any{"type": "choose_trump", "suit": bestSuit(s.YourHand)}) + atomic.AddInt64(&moves, 1) + } + case "playing": + if len(s.YourHand) > 0 { + _ = c.WriteJSON(map[string]any{"type": "play_card", "card": pickCard(s.YourHand, s.LeadSuit)}) + atomic.AddInt64(&moves, 1) + } + } + } +} + +func bestSuit(hand []string) string { + cnt := map[byte]int{} + for _, c := range hand { + cnt[c[len(c)-1]]++ + } + best, n := byte('S'), -1 + for s, k := range cnt { + if k > n { + best, n = s, k + } + } + return suitWord[best] +} + +func intOf(v any) int { + if f, ok := v.(float64); ok { + return int(f) + } + return 0 +} + +func main() { + url := flag.String("url", "ws://localhost:8080", "ws base url (no /ws)") + n := flag.Int("n", 100, "concurrent players") + tier := flag.String("tier", "beginner", "table tier id") + secret := flag.String("secret", "", "JWT_SECRET of the target server") + startID := flag.Int64("startid", 1000, "first synthetic user id (must exist in DB w/ coins)") + ramp := flag.Duration("ramp", 30*time.Millisecond, "delay between client starts") + play := flag.Duration("play", 300*time.Millisecond, "think time before each move") + flag.Parse() + if *secret == "" { + log.Fatal("-secret (JWT_SECRET) is required") + } + *url = strings.TrimRight(*url, "/") + + var wg sync.WaitGroup + go func() { + for range time.Tick(time.Second) { + fmt.Printf("conns=%d state/s=%d moves/s=%d errs=%d\n", + atomic.LoadInt64(&connected), + atomic.SwapInt64(&stateMsgs, 0), + atomic.SwapInt64(&moves, 0), + atomic.LoadInt64(&errs)) + } + }() + + for i := 0; i < *n; i++ { + wg.Add(1) + go runClient(*url, *secret, *tier, *startID+int64(i), *play, &wg) + time.Sleep(*ramp) + } + fmt.Printf("launched %d clients; Ctrl-C to stop\n", *n) + wg.Wait() +} diff --git a/internal/store/concurrency_test.go b/internal/store/concurrency_test.go new file mode 100644 index 0000000..3fe7363 --- /dev/null +++ b/internal/store/concurrency_test.go @@ -0,0 +1,55 @@ +package store + +import ( + "sync" + "sync/atomic" + "testing" +) + +// TestConcurrentReadWrite تأیید می‌کند که با pool چند-اتصالی + WAL + busy_timeout، +// خواندن و نوشتنِ هم‌زمان بدونِ خطای «database is locked» انجام می‌شود. +func TestConcurrentReadWrite(t *testing.T) { + s, err := Open(t.TempDir() + "/c.db") + if err != nil { + t.Fatal(err) + } + defer s.Close() + if _, err := s.DB.Exec(`CREATE TABLE t(id INTEGER PRIMARY KEY, v INTEGER)`); err != nil { + t.Fatal(err) + } + + const workers = 40 // ۲۰ نویسنده + ۲۰ خواننده هم‌زمان + const iters = 50 + var wg sync.WaitGroup + var writeErr, readErr int64 + for w := 0; w < workers; w++ { + wg.Add(1) + go func(w int) { + defer wg.Done() + for i := 0; i < iters; i++ { + if w%2 == 0 { + if _, err := s.DB.Exec(`INSERT INTO t(v) VALUES(?)`, i); err != nil { + atomic.AddInt64(&writeErr, 1) + t.Logf("write err: %v", err) + } + } else { + var n int + if err := s.DB.QueryRow(`SELECT COUNT(*) FROM t`).Scan(&n); err != nil { + atomic.AddInt64(&readErr, 1) + t.Logf("read err: %v", err) + } + } + } + }(w) + } + wg.Wait() + + if writeErr != 0 || readErr != 0 { + t.Fatalf("concurrency errors: writes=%d reads=%d (expected 0)", writeErr, readErr) + } + var total int + _ = s.DB.QueryRow(`SELECT COUNT(*) FROM t`).Scan(&total) + if want := (workers / 2) * iters; total != want { + t.Fatalf("expected %d rows written, got %d", want, total) + } +} diff --git a/internal/store/db.go b/internal/store/db.go index 8b3171e..b1cd136 100644 --- a/internal/store/db.go +++ b/internal/store/db.go @@ -6,6 +6,7 @@ import ( "fmt" "log/slog" "strings" + "time" _ "modernc.org/sqlite" // درایور SQLite خالص Go (بدون cgo) ) @@ -20,14 +21,19 @@ type Store struct { // Open اتصال SQLite را باز کرده و migrations را اجرا می‌کند. func Open(path string) (*Store, error) { - // pragmaها برای کارایی و یکپارچگی روی سرور کوچک - dsn := fmt.Sprintf("file:%s?_pragma=journal_mode(WAL)&_pragma=busy_timeout(5000)&_pragma=foreign_keys(ON)", path) + // pragmaها برای کارایی و یکپارچگی روی سرور کوچک. busy_timeout بالاتر تا نویسنده‌ی + // دوم به‌جای خطای «database is locked» چند لحظه منتظر بماند. + dsn := fmt.Sprintf("file:%s?_pragma=journal_mode(WAL)&_pragma=busy_timeout(10000)&_pragma=foreign_keys(ON)", path) db, err := sql.Open("sqlite", dsn) if err != nil { return nil, err } - // SQLite تک‌نویسنده است؛ pool کوچک نگه می‌داریم. - db.SetMaxOpenConns(1) + // WAL اجازه می‌دهد چند «خواننده» هم‌زمان باشند و فقط «نویسنده» تک بماند. با + // SetMaxOpenConns(1) این مزیت هدر می‌رفت (همه، حتی خواندن‌ها، در صفِ تک‌اتصالی). + // pool بزرگ‌تر ⇒ خواندن‌های هم‌زمان بدونِ صف؛ نوشتن‌ها با busy_timeout سریالی می‌مانند. + db.SetMaxOpenConns(10) + db.SetMaxIdleConns(10) + db.SetConnMaxIdleTime(5 * time.Minute) if err := db.Ping(); err != nil { return nil, err }