feat: WAL concurrency change is now running

This commit is contained in:
2026-07-23 13:10:19 +03:30
parent 64bdf5fc3d
commit 3883d1b471
3 changed files with 257 additions and 4 deletions
+192
View File
@@ -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()
}
+55
View File
@@ -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)
}
}
+10 -4
View File
@@ -6,6 +6,7 @@ import (
"fmt" "fmt"
"log/slog" "log/slog"
"strings" "strings"
"time"
_ "modernc.org/sqlite" // درایور SQLite خالص Go (بدون cgo) _ "modernc.org/sqlite" // درایور SQLite خالص Go (بدون cgo)
) )
@@ -20,14 +21,19 @@ type Store struct {
// Open اتصال SQLite را باز کرده و migrations را اجرا می‌کند. // Open اتصال SQLite را باز کرده و migrations را اجرا می‌کند.
func Open(path string) (*Store, error) { func Open(path string) (*Store, error) {
// pragmaها برای کارایی و یکپارچگی روی سرور کوچک // pragmaها برای کارایی و یکپارچگی روی سرور کوچک. busy_timeout بالاتر تا نویسنده‌ی
dsn := fmt.Sprintf("file:%s?_pragma=journal_mode(WAL)&_pragma=busy_timeout(5000)&_pragma=foreign_keys(ON)", path) // دوم به‌جای خطای «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) db, err := sql.Open("sqlite", dsn)
if err != nil { if err != nil {
return nil, err return nil, err
} }
// SQLite تک‌نویسنده است؛ pool کوچک نگه می‌داریم. // WAL اجازه می‌دهد چند «خواننده» هم‌زمان باشند و فقط «نویسنده» تک بماند. با
db.SetMaxOpenConns(1) // SetMaxOpenConns(1) این مزیت هدر می‌رفت (همه، حتی خواندن‌ها، در صفِ تک‌اتصالی).
// pool بزرگ‌تر ⇒ خواندن‌های هم‌زمان بدونِ صف؛ نوشتن‌ها با busy_timeout سریالی می‌مانند.
db.SetMaxOpenConns(10)
db.SetMaxIdleConns(10)
db.SetConnMaxIdleTime(5 * time.Minute)
if err := db.Ping(); err != nil { if err := db.Ping(); err != nil {
return nil, err return nil, err
} }