131 lines
4.0 KiB
Go
131 lines
4.0 KiB
Go
package ws
|
|
|
|
import (
|
|
"encoding/json"
|
|
"log/slog"
|
|
"sync"
|
|
"time"
|
|
|
|
"github.com/gorilla/websocket"
|
|
)
|
|
|
|
const (
|
|
writeWait = 10 * time.Second
|
|
pongWait = 60 * time.Second
|
|
pingPeriod = (pongWait * 9) / 10
|
|
maxMessageSize = 4096
|
|
// بافرِ بزرگتر تا در سرورِ تکهستهای، هجومِ کوتاهِ وضعیتها (پخشِ کارت + چند حرکتِ
|
|
// بات) روی کلاینتِ کند باعثِ دورریختنِ پیام نشود.
|
|
sendBuffer = 256
|
|
)
|
|
|
|
// Client یک اتصال WebSocket یک بازیکن.
|
|
type Client struct {
|
|
hub *Hub
|
|
conn *websocket.Conn
|
|
send chan []byte
|
|
UserID int64
|
|
Name string
|
|
|
|
// done هنگام قطع اتصال بسته میشود تا writePump خارج شده و
|
|
// trySend دیگر تلاش به ارسال نکند (send هیچگاه close نمیشود تا panic رخ ندهد).
|
|
done chan struct{}
|
|
closeOnce sync.Once
|
|
|
|
// توسط goroutine هاب ست میشوند (تکنویسنده) و فقط توسط آن خوانده میشوند.
|
|
room *Room
|
|
seat int
|
|
tier string // نوع میزی که در صفش است (برای تسویه)
|
|
tableCode string // کدِ میز خصوصیِ در انتظار (پیش از شروع بازی)
|
|
}
|
|
|
|
// close اتصال را یکبار بهصورت امن میبندد.
|
|
func (c *Client) close() {
|
|
c.closeOnce.Do(func() {
|
|
close(c.done)
|
|
c.conn.Close()
|
|
})
|
|
}
|
|
|
|
// readPump پیامهای ورودی را خوانده و به هاب میفرستد.
|
|
func (c *Client) readPump() {
|
|
defer func() {
|
|
c.hub.unregister <- c
|
|
c.close()
|
|
}()
|
|
c.conn.SetReadLimit(maxMessageSize)
|
|
_ = c.conn.SetReadDeadline(time.Now().Add(pongWait))
|
|
c.conn.SetPongHandler(func(string) error {
|
|
return c.conn.SetReadDeadline(time.Now().Add(pongWait))
|
|
})
|
|
|
|
for {
|
|
_, raw, err := c.conn.ReadMessage()
|
|
if err != nil {
|
|
return
|
|
}
|
|
var msg inboundMsg
|
|
if err := json.Unmarshal(raw, &msg); err != nil {
|
|
c.trySend(mustJSON(errorMsg{Type: "error", Message: "invalid message"}))
|
|
continue
|
|
}
|
|
c.hub.inbound <- inbound{client: c, msg: msg}
|
|
}
|
|
}
|
|
|
|
// writePump پیامهای خروجی و ping را به سوکت مینویسد.
|
|
func (c *Client) writePump() {
|
|
ticker := time.NewTicker(pingPeriod)
|
|
defer func() {
|
|
ticker.Stop()
|
|
c.close()
|
|
}()
|
|
for {
|
|
select {
|
|
case <-c.done:
|
|
return
|
|
case msg := <-c.send:
|
|
_ = c.conn.SetWriteDeadline(time.Now().Add(writeWait))
|
|
if err := c.conn.WriteMessage(websocket.TextMessage, msg); err != nil {
|
|
return
|
|
}
|
|
case <-ticker.C:
|
|
_ = c.conn.SetWriteDeadline(time.Now().Add(writeWait))
|
|
if err := c.conn.WriteMessage(websocket.PingMessage, nil); err != nil {
|
|
return
|
|
}
|
|
}
|
|
}
|
|
}
|
|
|
|
// trySend ارسالِ غیرمسدودکننده. اگر بافر پر باشد، بهجای دورریختنِ پیامِ جدید،
|
|
// «قدیمیترین» پیام دور ریخته میشود تا کلاینت روی جدیدترین وضعیت همگرا بماند
|
|
// (هر وضعیتِ کامل، وضعیتِ قبلی را باطل میکند). قبلاً جدیدترین دور ریخته میشد و
|
|
// کلاینت روی حالتِ کهنه گیر میکرد: کارتهای دست دیده نمیشد ولی باتها بازی میکردند.
|
|
// send هیچگاه close نمیشود؛ پس از done صرفاً پیام دور ریخته میشود (بدون panic).
|
|
func (c *Client) trySend(b []byte) {
|
|
for {
|
|
select {
|
|
case <-c.done:
|
|
return
|
|
case c.send <- b:
|
|
return
|
|
default:
|
|
// بافر پر است: یک پیامِ قدیمی را خالی کن و دوباره تلاش کن.
|
|
select {
|
|
case <-c.send:
|
|
slog.Warn("client send buffer full, dropping oldest", "user", c.UserID)
|
|
case <-c.done:
|
|
return
|
|
default:
|
|
// بافر بین دو تلاش خالی شد (تولیدکنندهی دیگر آن را خالی کرد)؛ ادامه بده.
|
|
}
|
|
}
|
|
}
|
|
}
|
|
|
|
func mustJSON(v any) []byte {
|
|
b, _ := json.Marshal(v)
|
|
return b
|
|
}
|