fix(server): name the player on every played word, and bound the process

PlayedWord.player_id was declared and read by the client but never set by
the server, so a room of three or four never showed who played each word.
The chain byline now comes from the room, with a producer-side test.

The process also gains the ceilings it was missing: a cap on live rooms and
on open sockets, a per-connection frame-rate limit so a payload-less frame
is no longer free, and an opt-in trusted-proxy list so the join limiter can
tell players apart behind the documented reverse proxy instead of putting
them in one bucket. The typed word is sanitized before the engine stores it,
since every seat is shown it; a room exiting on its idle clock releases the
sessions still bound to it; the dictionary builder escapes its SQLite path
like the store does and renames over the old database instead of deleting
it first.
This commit is contained in:
tiennm99 committed 2026-09-21 00:38:09 +07:00
1 parent a5052fca3f
commit 442ced32cc
14 files changed
+494 -54

No files matched your search

+5 -2
View File
@@ -189,9 +189,12 @@ Configuration is environment-only; every variable has a working default.
| `NOITU_GRACE` | `30s` | How long a disconnected player's seat is held for a reconnect |
| `NOITU_ALLOWED_ORIGINS` | *(unset)* | Comma-separated origin allowlist. Unset means same-origin only |
| `NOITU_WEB_DIR` | *(unset)* | Built frontend to serve. Unset serves the API alone |
| `NOITU_TRUSTED_PROXIES` | *(unset)* | Comma-separated proxy addresses or CIDRs whose `X-Forwarded-For` is believed. Unset keys limiters on the socket peer |
| `NOITU_MAX_ROOMS` | `1000` | Ceiling on live rooms across the process; a creator past it is told `server_full` |
| `NOITU_MAX_CONNECTIONS` | `2000` | Ceiling on open WebSockets; the next upgrade gets HTTP 503 |
An invalid duration is logged and ignored rather than silently changing the
rules of the game.
An invalid duration or count is logged and ignored rather than silently
changing the rules of the game.
Endpoints: `GET /ws` (Protobuf over binary WebSocket frames), `GET /healthz`,
and — when `NOITU_WEB_DIR` is set — the frontend on everything else, with
+33 -7
View File
@@ -18,9 +18,12 @@ so the image runs with nothing set.
| `NOITU_GRACE` | `30s` | How long a disconnected player's seat is held for a reconnect |
| `NOITU_ALLOWED_ORIGINS` | *(unset)* | Comma-separated origin allowlist. Unset means same-origin only |
| `NOITU_WEB_DIR` | *(unset)* | Built frontend to serve. Unset serves the API alone |
| `NOITU_TRUSTED_PROXIES` | *(unset)* | Comma-separated proxy addresses or CIDRs whose `X-Forwarded-For` is believed. Unset keys limiters on the socket peer |
| `NOITU_MAX_ROOMS` | `1000` | Ceiling on live rooms across the process; a creator past it is told `server_full` |
| `NOITU_MAX_CONNECTIONS` | `2000` | Ceiling on open WebSockets; the next upgrade gets HTTP 503 |
An invalid duration is logged and ignored rather than silently changing the
rules of the game.
An invalid duration or count is logged and ignored rather than silently
changing the rules of the game.
One timing is not configurable: an online room closes after 10 minutes in its
lobby with no game started. It is a fixed constant because nothing about a
@@ -112,11 +115,34 @@ noitu.example {
### The client's own address
Rate limiting counts against `RemoteAddr` and deliberately ignores
`X-Forwarded-For`, because that header is attacker-controlled unless the proxy
overwrites it. Behind a proxy every player therefore shares one bucket. If that
becomes a problem, the fix is to make the proxy the only source of the header
and teach the server to trust it — not to trust it as things stand.
Rate limiting keys on the client's address, and by default that is the
socket's own peer, `RemoteAddr`. `X-Forwarded-For` is attacker-controlled
unless the proxy is known to append to it, so it is ignored until told
otherwise. Behind a proxy every player therefore shares one bucket, which
means one client brute-forcing room codes spends everybody's join budget.
The fix is to name the proxy. Set `NOITU_TRUSTED_PROXIES` to the address, or
CIDR range, the proxy connects from — `127.0.0.1` for the nginx and Caddy
examples above, or the container network's range under Compose — and the
server walks `X-Forwarded-For` from the right, taking the first hop that is
not itself a trusted proxy. Entries a client forged sit to the left of the one
the proxy appended, so they are never reached. A peer that is not on the list
is still keyed on its socket address, header or not.
Both proxies above append the real client to `X-Forwarded-For` by default.
Do not list a range the public can connect from; that is the same as trusting
the header unconditionally.
### Capacity
Two ceilings bound the process as a whole, on top of the per-connection rate
limits: `NOITU_MAX_ROOMS` live rooms and `NOITU_MAX_CONNECTIONS` open
sockets. Past the first, creating a room answers `server_full` and the player
is asked to wait; past the second, the upgrade itself is refused with HTTP 503
so the proxy can count it. Each connection also has a frame-rate ceiling, and a
client past it is disconnected rather than throttled. The defaults are
generous for one binary on a small host; lower them if memory is tight,
because a room is a goroutine and an engine held for up to its idle window.
## Health check
+13 -8
View File
@@ -30,6 +30,7 @@ import (
"time"
"unicode/utf8"
"github.com/tiennm99dev/noitu/server/internal/dictionary"
_ "modernc.org/sqlite"
)
@@ -225,7 +226,7 @@ func parseSenses(cells []string) []sense {
// depends on. The in-memory checks above can only prove what the builder
// intended; this proves what actually landed on disk.
func verify(path string, minWords int, requireCoverage bool) error {
db, err := sql.Open("sqlite", "file:"+path+"?mode=ro")
db, err := sql.Open("sqlite", dictionary.DSN(path, true))
if err != nil {
return err
}
@@ -376,13 +377,17 @@ func write(path string, words map[string]entry, meanings map[string][]sense, ali
return err
}
// os.Rename replaces the destination atomically on POSIX; on Windows it
// fails if the target exists, so clear it first.
if err := os.Remove(path); err != nil && !errors.Is(err, os.ErrNotExist) {
return fmt.Errorf("remove existing output: %w", err)
}
// os.Rename replaces the destination atomically on POSIX, so the good
// database is never missing from disk. Only where that fails — Windows
// refuses to rename over an existing file — is the target cleared first,
// accepting the window there rather than opening it everywhere.
if err := os.Rename(tmp, path); err != nil {
return fmt.Errorf("move temp database into place: %w", err)
if rmErr := os.Remove(path); rmErr != nil && !errors.Is(rmErr, os.ErrNotExist) {
return fmt.Errorf("remove existing output: %w", rmErr)
}
if err := os.Rename(tmp, path); err != nil {
return fmt.Errorf("move temp database into place: %w", err)
}
}
committed = true
@@ -390,7 +395,7 @@ func write(path string, words map[string]entry, meanings map[string][]sense, ali
}
func writeTo(path string, words map[string]entry, meanings map[string][]sense, aliases map[string]string, src sourceSpec) error {
db, err := sql.Open("sqlite", "file:"+path)
db, err := sql.Open("sqlite", dictionary.DSN(path, false))
if err != nil {
return fmt.Errorf("create output: %w", err)
}
+25
View File
@@ -9,6 +9,7 @@ import (
"net/http"
"os"
"os/signal"
"strconv"
"strings"
"syscall"
"time"
@@ -36,6 +37,9 @@ type config struct {
grace time.Duration
allowedOrigins []string
webDir string
trustedProxies []string
maxRooms int
maxConnections int
}
func main() {
@@ -74,6 +78,9 @@ func run() error {
GraceFor: cfg.grace,
AllowedOrigins: cfg.allowedOrigins,
WebDir: cfg.webDir,
TrustedProxies: cfg.trustedProxies,
MaxRooms: cfg.maxRooms,
MaxConnections: cfg.maxConnections,
})
srv := &http.Server{
@@ -114,9 +121,27 @@ func loadConfig() config {
grace: envDuration("NOITU_GRACE", defaultGrace),
allowedOrigins: envList("NOITU_ALLOWED_ORIGINS"),
webDir: env("NOITU_WEB_DIR", ""),
trustedProxies: envList("NOITU_TRUSTED_PROXIES"),
maxRooms: envInt("NOITU_MAX_ROOMS", 0),
maxConnections: envInt("NOITU_MAX_CONNECTIONS", 0),
}
}
// envInt falls back loudly, like envDuration. Zero means "use the built-in
// default", so it is what an unset or invalid value becomes.
func envInt(key string, fallback int) int {
raw := strings.TrimSpace(os.Getenv(key))
if raw == "" {
return fallback
}
n, err := strconv.Atoi(raw)
if err != nil || n < 0 {
slog.Warn("ignoring invalid integer", "key", key, "value", raw, "using", fallback)
return fallback
}
return n
}
func env(key, fallback string) string {
if v := strings.TrimSpace(os.Getenv(key)); v != "" {
return v
+13 -5
View File
@@ -139,11 +139,19 @@ func Open(path string) (*Store, error) {
return s, nil
}
// dsn builds the SQLite URI. The path must be escaped: SQLite reads '#' as a
// URI fragment delimiter, so a bare path containing one silently opens a
// different (usually nonexistent) file and reports a confusing schema error.
func dsn(path string) string {
u := url.URL{Scheme: "file", Opaque: (&url.URL{Path: path}).EscapedPath(), RawQuery: "mode=ro"}
// dsn is the read-only URI the store opens with.
func dsn(path string) string { return DSN(path, true) }
// DSN builds a SQLite URI for path. The path must be escaped: SQLite reads
// '#' as a URI fragment delimiter, so a bare path containing one silently
// opens a different (usually nonexistent) file and reports a confusing schema
// error. Exported so the builder, which writes the file this package reads,
// spells the path the same way.
func DSN(path string, readOnly bool) string {
u := url.URL{Scheme: "file", Opaque: (&url.URL{Path: path}).EscapedPath()}
if readOnly {
u.RawQuery = "mode=ro"
}
return u.String()
}
+1
View File
@@ -109,6 +109,7 @@ func PlayedWord(m game.Move, byMe bool, meanings []dictionary.Sense) *noituv1.Pl
Word: m.Word,
Typed: m.Typed,
ByMe: byMe,
PlayerId: string(m.Player),
Points: uint32(m.Points),
Syllables: uint32(m.Syllables),
Meanings: Senses(meanings),
+6 -1
View File
@@ -170,12 +170,17 @@ func TestDifficultyMapping(t *testing.T) {
// has to be able to show that canonicalization changed the player's text, so
// the raw input must survive onto the wire.
func TestPlayedWordKeepsTypedInput(t *testing.T) {
m := game.Move{Word: "hòa bình", Typed: "hoà bình", Syllables: 2, Points: 2}
m := game.Move{Player: "p2", Word: "hòa bình", Typed: "hoà bình", Syllables: 2, Points: 2}
mine := PlayedWord(m, true, nil)
if mine.GetWord() != "hòa bình" || mine.GetTyped() != "hoà bình" {
t.Errorf("canonical/typed pair lost: word=%q typed=%q", mine.GetWord(), mine.GetTyped())
}
// The byline in a room of three or four is drawn from this field alone;
// by_me only says whether it was the recipient.
if mine.GetPlayerId() != "p2" {
t.Errorf("player_id lost: got %q", mine.GetPlayerId())
}
if !mine.GetByMe() || mine.GetSyllables() != 2 || mine.GetPoints() != 2 {
t.Errorf("unexpected rendering: %+v", mine)
}
+27 -1
View File
@@ -26,8 +26,15 @@ const codeAttempts = 10
var (
errRoomNotFound = errors.New("wsapi: no such room")
errNoRoomCode = errors.New("wsapi: could not allocate a room code")
errServerFull = errors.New("wsapi: room limit reached")
)
// defaultMaxRooms bounds live rooms across the whole process when nothing else
// is configured. Each room is a goroutine, an engine and a registry entry held
// for up to the idle window, so without a ceiling the per-connection limiter
// only sets the rate at which a fleet of connections can fill memory.
const defaultMaxRooms = 1000
// hub owns the registries and nothing else.
//
// It never touches a game: rooms are handed out as pointers whose channels are
@@ -41,6 +48,7 @@ type hub struct {
turnLimit time.Duration
graceFor time.Duration
idleFor time.Duration
maxRooms int
mu sync.Mutex
rooms map[string]*room
@@ -49,13 +57,17 @@ type hub struct {
joinLimiter *keyedLimiter
}
func newHub(ctx context.Context, dict Dictionary, turnLimit, graceFor, idleFor time.Duration) *hub {
func newHub(ctx context.Context, dict Dictionary, turnLimit, graceFor, idleFor time.Duration, maxRooms int) *hub {
if maxRooms <= 0 {
maxRooms = defaultMaxRooms
}
return &hub{
ctx: ctx,
dict: dict,
turnLimit: turnLimit,
graceFor: graceFor,
idleFor: idleFor,
maxRooms: maxRooms,
rooms: map[string]*room{},
sessions: map[string]*session{},
joinLimiter: newKeyedLimiter(joinsPerSecond, joinBurst, limiterIdleFor),
@@ -145,7 +157,14 @@ func (h *hub) newRegisteredRoom() (*room, error) {
r := newRoom(h, code, h.turnLimit, h.graceFor, h.idleFor)
// The ceiling is checked under the same lock that registers the room, so
// two creators racing for the last slot cannot both get it.
h.mu.Lock()
if len(h.rooms) >= h.maxRooms {
h.mu.Unlock()
r.cancel()
return nil, errServerFull
}
h.rooms[code] = r
h.mu.Unlock()
@@ -153,6 +172,13 @@ func (h *hub) newRegisteredRoom() (*room, error) {
return r, nil
}
// roomCount is how many rooms are live right now.
func (h *hub) roomCount() int {
h.mu.Lock()
defer h.mu.Unlock()
return len(h.rooms)
}
// evict removes a finished room. Called by the room goroutine as it exits, so
// a code is reusable the moment its game is done.
func (h *hub) evict(code string) {
+1 -7
View File
@@ -1,9 +1,3 @@
package wsapi
// roomCount reports live rooms, so a test can assert that a room was evicted.
// Nothing in the server depends on it.
func (h *hub) roomCount() int {
h.mu.Lock()
defer h.mu.Unlock()
return len(h.rooms)
}
// roomCount lives in hub.go; tests use it to assert that a room was evicted.
+189
View File
@@ -0,0 +1,189 @@
package wsapi
import (
"net/http"
"net/http/httptest"
"testing"
"time"
"github.com/coder/websocket"
noituv1 "github.com/tiennm99dev/noitu/server/gen/noitu/v1"
"google.golang.org/protobuf/proto"
)
// TestPlayedWordNamesItsPlayerOnTheWire is the producer-side check: the chain
// byline in a room of three or four is drawn from PlayedWord.player_id, and a
// fixture that hand-writes the field proves nothing about the room.
func TestPlayedWordNamesItsPlayerOnTheWire(t *testing.T) {
_, url := newTestServer(t, chainDict(), Config{})
lead, waits, start := pvpGame(t, url)
lead.submit("b c", start.GetTurnSeq())
played := waits.await("turn_update").GetTurnUpdate().GetPlayed()
if played.GetPlayerId() == "" {
t.Fatal("player_id is empty on the wire")
}
if played.GetPlayerId() != start.GetTurnPlayerId() {
t.Errorf("player_id = %q, want the leader %q", played.GetPlayerId(), start.GetTurnPlayerId())
}
}
// TestTypedWordIsSanitizedBeforeItIsEchoed guards the trust boundary the typed
// text crosses: every seat is shown it, so format characters must be gone
// before the engine stores it.
func TestTypedWordIsSanitizedBeforeItIsEchoed(t *testing.T) {
_, url := newTestServer(t, chainDict(), Config{})
lead, waits, start := pvpGame(t, url)
// A zero-width joiner and a bidi override inside an otherwise legal word.
lead.submit("b‍ ‮c", start.GetTurnSeq())
played := waits.await("turn_update").GetTurnUpdate().GetPlayed()
if played.GetWord() != "b c" {
t.Fatalf("word = %q, want the move accepted as %q", played.GetWord(), "b c")
}
if played.GetTyped() != "b c" {
t.Errorf("typed = %q still carries non-printing characters", played.GetTyped())
}
}
// TestRoomCapRefusesTheNextRoom bounds live rooms across the process, not per
// connection: a fleet of connections each under its own limiter must still
// hit a ceiling.
func TestRoomCapRefusesTheNextRoom(t *testing.T) {
_, url := newTestServer(t, chainDict(), Config{MaxRooms: 2})
for range 2 {
c := dial(t, url)
c.hello("Chủ phòng")
c.send(&noituv1.ClientMessage{Payload: &noituv1.ClientMessage_CreateRoom{CreateRoom: &noituv1.CreateRoom{}}})
c.await("room_state")
}
third := dial(t, url)
third.hello("Người thứ ba")
third.send(&noituv1.ClientMessage{Payload: &noituv1.ClientMessage_CreateRoom{CreateRoom: &noituv1.CreateRoom{}}})
if got := third.await("error").GetError().GetCode(); got != "server_full" {
t.Errorf("error code = %q, want server_full", got)
}
}
// TestConnectionCapRefusesBeforeUpgrade: the refusal is an HTTP status a
// client can read, and it costs the server no socket.
func TestConnectionCapRefusesBeforeUpgrade(t *testing.T) {
_, url := newTestServer(t, chainDict(), Config{MaxConnections: 1})
first := dial(t, url)
first.hello("Một")
_, resp, err := websocket.Dial(t.Context(), url+"/ws", nil)
if err == nil {
t.Fatal("second connection was accepted past the cap")
}
if resp == nil || resp.StatusCode != http.StatusServiceUnavailable {
t.Fatalf("want 503 before the upgrade, got %v (err %v)", resp, err)
}
}
// TestFrameFloodClosesTheConnection: a message that matches no dispatch arm
// used to be free at line rate. Now every frame is metered before it is
// decoded, and a flood is closed rather than throttled.
func TestFrameFloodClosesTheConnection(t *testing.T) {
_, url := newTestServer(t, chainDict(), Config{})
c := dial(t, url)
c.hello("Người thử")
empty, _ := proto.Marshal(&noituv1.ClientMessage{})
for range frameBurst + 20 {
if err := c.conn.Write(c.ctx, websocket.MessageBinary, empty); err != nil {
return // closed on us mid-flood, which is the point
}
}
for range 5 {
_, _, err := c.conn.Read(c.ctx)
if err != nil {
return
}
}
t.Error("connection survived a frame flood")
}
// TestClientIPTrustsOnlyConfiguredProxies covers the three shapes the limiter
// key can take: no proxy configured, a trusted proxy carrying a forwarded
// chain, and a forged header from a peer that is not a proxy.
func TestClientIPTrustsOnlyConfiguredProxies(t *testing.T) {
req := func(remote, xff string) *http.Request {
r := httptest.NewRequest(http.MethodGet, "/ws", nil)
r.RemoteAddr = remote
if xff != "" {
r.Header.Set("X-Forwarded-For", xff)
}
return r
}
plain := &Server{}
if got := plain.clientIP(req("203.0.113.9:4000", "10.0.0.1")); got != "203.0.113.9" {
t.Errorf("no proxies configured: got %q, want the peer", got)
}
s := &Server{proxies: parsePrefixes([]string{"127.0.0.1", "10.0.0.0/8", "not-an-address"})}
cases := []struct{ remote, xff, want string }{
// The proxy appended the client; the client forged nothing.
{"127.0.0.1:5000", "198.51.100.7", "198.51.100.7"},
// Two trusted hops behind the peer, then the client.
{"127.0.0.1:5000", "198.51.100.7, 10.1.2.3", "198.51.100.7"},
// The client forged a header; the proxy appended the real address
// after it, and the rightmost untrusted entry wins.
{"10.9.9.9:5000", "1.2.3.4, 198.51.100.7", "198.51.100.7"},
// A peer that is not a proxy is taken at its socket address, header
// or not.
{"203.0.113.9:4000", "198.51.100.7", "203.0.113.9"},
// A garbage hop falls back to the proxy itself rather than keying a
// bucket on whatever was typed.
{"127.0.0.1:5000", "not an ip", "127.0.0.1"},
// An IPv4-mapped peer still matches its IPv4 prefix.
{"[::ffff:127.0.0.1]:5000", "198.51.100.7", "198.51.100.7"},
}
for _, tc := range cases {
if got := s.clientIP(req(tc.remote, tc.xff)); got != tc.want {
t.Errorf("remote=%s xff=%q: got %q, want %q", tc.remote, tc.xff, got, tc.want)
}
}
}
// TestIdleRoomReleasesItsSeats: a room that closes on its idle clock must let
// go of the connections still sitting in it, or they are stuck pointing at a
// goroutine that has exited and can never be seated cleanly again.
func TestIdleRoomReleasesItsSeats(t *testing.T) {
api, url := newTestServer(t, chainDict(), Config{IdleFor: 200 * time.Millisecond})
host := dial(t, url)
host.hello("Chủ phòng")
host.send(&noituv1.ClientMessage{Payload: &noituv1.ClientMessage_CreateRoom{CreateRoom: &noituv1.CreateRoom{}}})
host.await("room_state")
if got := host.await("error").GetError().GetCode(); got != "room_idle_closed" {
t.Fatalf("error code = %q, want room_idle_closed", got)
}
settle()
if n := api.hub.roomCount(); n != 0 {
t.Fatalf("%d rooms still registered after the idle close", n)
}
// The connection is free again: a second room opens and seats it.
host.send(&noituv1.ClientMessage{Payload: &noituv1.ClientMessage_CreateRoom{CreateRoom: &noituv1.CreateRoom{}}})
if host.await("room_state").GetRoomState().GetRoomCode() == "" {
t.Error("could not be seated in a new room after the idle close")
}
}
// pvpGame seats two players, starts the game and sorts them into the one who
// drew the first turn and the one who waits.
func pvpGame(t *testing.T, url string) (lead, waits *testClient, start *noituv1.GameStarted) {
t.Helper()
host, guest, _ := pvpLobby(t, url)
guest.setReady(true)
host.await("room_state")
host.startGame()
return awaitLead(t, host, guest)
}
+28 -3
View File
@@ -46,6 +46,11 @@ const chatHistoryLimit = 20
// maxNicknameRunes is.
const maxChatRunes = 200
// maxWordRunes caps a submitted word before the engine sees it. The longest
// dictionary entries are well under this, so it bounds abuse without ever
// deciding a real move.
const maxWordRunes = 64
// maxChatMarks caps mark stacking in a message, as maxNicknameMarks does for a
// name. A message is ten times longer, so the same stack does ten times more
// damage.
@@ -327,6 +332,11 @@ func (r *room) send(msg any) bool {
func (r *room) run() {
defer r.cancel()
defer r.hub.evict(r.code)
// Whatever ended the room — everybody leaving, the idle window, a server
// shutdown — the connections still seated in it must stop pointing here.
// A session that keeps a dead room would answer every later action with
// "not in a room" and could never be seated anywhere else cleanly.
defer r.detachAll()
var turnTimer, graceTimer, idleTimer *time.Timer
stop := func(t *time.Timer) {
@@ -812,10 +822,16 @@ func (r *room) handleSubmit(m submitInput) {
return
}
// The typed text is echoed back to every seat as PlayedWord.typed, so it
// crosses the same trust boundary a chat line does and gets the same
// filter. The engine's own normalization only lowercases and collapses
// whitespace; it does not drop format characters.
word := sanitizeText(m.word, maxWordRunes, maxNicknameMarks)
before := r.mark()
move, reason := r.engine.Submit(m.player, m.word, time.Now())
move, reason := r.engine.Submit(m.player, word, time.Now())
if reason != game.ReasonNone {
r.sendTo(m.player, moveRejectedMsg(RejectReason(reason), m.word, m.turnSeq))
r.sendTo(m.player, moveRejectedMsg(RejectReason(reason), word, m.turnSeq))
// A rejection for an expired turn also took this player out of the
// game, and everybody has to be told which.
r.applyEliminations(before)
@@ -1338,7 +1354,7 @@ func (r *room) sendChatHistory(s *seat) {
func chatMessageFor(entry chatEntry, id game.PlayerID) *noituv1.ServerMessage {
return &noituv1.ServerMessage{Payload: &noituv1.ServerMessage_ChatMessage{
ChatMessage: &noituv1.ChatMessage{
FromMe: entry.author != "" && entry.author == id,
FromMe: entry.author != "" && entry.author == id,
// Empty together with the name for a vacated seat: a line nobody
// owns must not be coloured as somebody's either.
PlayerId: string(entry.author),
@@ -1482,6 +1498,15 @@ func (r *room) vacate(s *seat) {
}
}
// detachAll releases every connection still bound to this room as it exits.
func (r *room) detachAll() {
for _, s := range r.seats {
if s != nil && s.sess != nil {
s.sess.release(r)
}
}
}
// promote hands the room to whoever is left.
func (r *room) promote() {
for _, s := range r.seats {
+113 -12
View File
@@ -5,9 +5,11 @@ import (
"log/slog"
"net"
"net/http"
"net/netip"
"os"
"path/filepath"
"strings"
"sync/atomic"
"time"
"github.com/coder/websocket"
@@ -30,14 +32,31 @@ type Config struct {
// WebDir is the built frontend. Empty, or missing on disk, serves the API
// alone, which is how the server runs before the frontend has been built.
WebDir string
// TrustedProxies lists the addresses, or CIDR ranges, of reverse proxies
// whose X-Forwarded-For header is believed. Empty means the header is
// ignored and every limiter keys on the socket's own peer address.
TrustedProxies []string
// MaxRooms caps live rooms across the process; zero means a built-in
// default. MaxConnections caps open sockets the same way.
MaxRooms int
MaxConnections int
}
// defaultMaxConnections bounds open WebSockets when nothing else is set. Each
// one is three goroutines and an outbox; the number is generous for one
// binary and small next to what the host can hold.
const defaultMaxConnections = 2000
// Server wires the hub to an HTTP mux.
type Server struct {
hub *hub
mux *http.ServeMux
cancel context.CancelFunc
cfg Config
proxies []netip.Prefix
maxConns int64
conns atomic.Int64
}
// NewServer builds the handler tree.
@@ -45,10 +64,15 @@ func NewServer(ctx context.Context, dict Dictionary, cfg Config) *Server {
ctx, cancel := context.WithCancel(ctx)
s := &Server{
hub: newHub(ctx, dict, cfg.TurnLimit, cfg.GraceFor, cfg.IdleFor),
mux: http.NewServeMux(),
cancel: cancel,
cfg: cfg,
hub: newHub(ctx, dict, cfg.TurnLimit, cfg.GraceFor, cfg.IdleFor, cfg.MaxRooms),
mux: http.NewServeMux(),
cancel: cancel,
cfg: cfg,
proxies: parsePrefixes(cfg.TrustedProxies),
maxConns: int64(cfg.MaxConnections),
}
if s.maxConns <= 0 {
s.maxConns = defaultMaxConnections
}
s.mux.HandleFunc("GET /ws", s.handleWS)
@@ -71,6 +95,15 @@ func (s *Server) Shutdown() {
}
func (s *Server) handleWS(w http.ResponseWriter, r *http.Request) {
// Refused before the upgrade, so a client that is over the line is told
// so in HTTP terms it can read, and never costs a socket.
if s.conns.Add(1) > s.maxConns {
s.conns.Add(-1)
http.Error(w, "server full", http.StatusServiceUnavailable)
return
}
defer s.conns.Add(-1)
conn, err := websocket.Accept(w, r, &websocket.AcceptOptions{
OriginPatterns: s.cfg.AllowedOrigins,
})
@@ -81,7 +114,7 @@ func (s *Server) handleWS(w http.ResponseWriter, r *http.Request) {
return
}
sess := newSession(s.hub.ctx, conn, s.hub, clientIP(r))
sess := newSession(s.hub.ctx, conn, s.hub, s.clientIP(r))
sess.run()
// The token has to outlive the socket by exactly the grace window: that is
@@ -149,18 +182,86 @@ func underRoot(root, path string) bool {
// clientIP is the key the join limiter counts against.
//
// RemoteAddr is deliberately the only source. Behind the reverse proxy this
// deploys under, X-Forwarded-For is attacker-controlled unless the proxy is
// known to overwrite it, and trusting it unconditionally would let one client
// spend everyone else's budget by forging the header.
func clientIP(r *http.Request) string {
host, _, err := net.SplitHostPort(r.RemoteAddr)
// RemoteAddr is the default and the only source when no proxy is trusted:
// X-Forwarded-For is attacker-controlled unless the proxy is known to append
// to it, and trusting it unconditionally would let one client spend everyone
// else's budget by forging the header. When the peer is a configured proxy,
// the header is walked from the right and the first address that is not
// itself a trusted proxy is the client — the entries a client could have
// forged all sit to the left of the one the proxy appended.
func (s *Server) clientIP(r *http.Request) string {
peer := remoteHost(r.RemoteAddr)
if len(s.proxies) == 0 || !s.trusted(peer) {
return peer
}
var hops []string
for _, v := range r.Header.Values("X-Forwarded-For") {
hops = append(hops, strings.Split(v, ",")...)
}
for i := len(hops) - 1; i >= 0; i-- {
hop := strings.TrimSpace(hops[i])
if hop == "" || s.trusted(hop) {
continue
}
if _, err := netip.ParseAddr(hop); err != nil {
// A malformed hop is a header somebody wrote by hand; fall back
// to the proxy's address rather than key a limiter on garbage.
return peer
}
return hop
}
return peer
}
// trusted reports whether host is one of the configured proxies.
func (s *Server) trusted(host string) bool {
addr, err := netip.ParseAddr(host)
if err != nil {
return r.RemoteAddr
return false
}
addr = addr.Unmap()
for _, p := range s.proxies {
if p.Contains(addr) {
return true
}
}
return false
}
// remoteHost strips the port from a RemoteAddr.
func remoteHost(remoteAddr string) string {
host, _, err := net.SplitHostPort(remoteAddr)
if err != nil {
return remoteAddr
}
return host
}
// parsePrefixes reads proxy addresses as CIDR ranges, accepting a bare
// address as a range of one. An entry that parses as neither is logged and
// skipped rather than silently trusting nothing or everything.
func parsePrefixes(raw []string) []netip.Prefix {
var out []netip.Prefix
for _, entry := range raw {
entry = strings.TrimSpace(entry)
if entry == "" {
continue
}
if p, err := netip.ParsePrefix(entry); err == nil {
out = append(out, p.Masked())
continue
}
if a, err := netip.ParseAddr(entry); err == nil {
a = a.Unmap()
out = append(out, netip.PrefixFrom(a, a.BitLen()))
continue
}
slog.Warn("ignoring unparseable trusted proxy", "entry", entry)
}
return out
}
// sweepLimiters keeps the per-key rate limiter from growing without bound.
func (s *Server) sweepLimiters(ctx context.Context) {
ticker := time.NewTicker(limiterIdleFor)
+39 -8
View File
@@ -54,9 +54,20 @@ const (
roomsPerSecond = 0.2
roomBurst = 5
limiterIdleFor = 5 * time.Minute
// framesPerSecond bounds every frame a connection sends, before it is
// routed. The per-action limiters above only meter the actions they know
// about; a Ping, or a ClientMessage with no payload set, matched none of
// them and cost the reader a decode at line rate. A client past this is
// not a player typing, so the connection is closed rather than throttled.
framesPerSecond = 20
frameBurst = 40
)
var errHandshake = errors.New("wsapi: first message must be Hello")
var (
errHandshake = errors.New("wsapi: first message must be Hello")
errFlood = errors.New("wsapi: frame rate exceeded")
)
// session is one WebSocket connection.
//
@@ -108,6 +119,7 @@ type session struct {
submitLimiter *bucket
roomLimiter *bucket
chatLimiter *bucket
frameLimiter *bucket
// greeted marks the handshake done. It is a one-shot transition: a second
// Hello would re-register the session and rewrite its nickname mid-game.
@@ -134,6 +146,7 @@ func newSession(ctx context.Context, conn *websocket.Conn, h *hub, remoteIP stri
submitLimiter: newBucket(submitsPerSecond, submitBurst, time.Now()),
roomLimiter: newBucket(roomsPerSecond, roomBurst, time.Now()),
chatLimiter: newBucket(chatsPerSecond, chatBurst, time.Now()),
frameLimiter: newBucket(framesPerSecond, frameBurst, time.Now()),
}
}
@@ -292,6 +305,10 @@ func (s *session) readLoop() error {
if err != nil {
return err
}
if !s.frameLimiter.allow(time.Now()) {
s.send(errorMsg("too_fast"))
return errFlood
}
msg, err := Decode(typ, raw)
if err != nil {
@@ -407,18 +424,20 @@ func (s *session) dispatch(msg *noituv1.ClientMessage) error {
return s.handleHello(p.Hello)
case *noituv1.ClientMessage_StartBotGame:
// The limiter is charged before the payload is inspected, so a bad
// difficulty costs the same as a good one and cannot be used to probe
// for free.
if !s.roomLimiter.allow(time.Now()) {
s.send(errorMsg("too_many_rooms"))
return nil
}
difficulty, ok := Difficulty(p.StartBotGame.GetDifficulty())
if !ok {
s.send(errorMsg("unknown_difficulty"))
return nil
}
if !s.roomLimiter.allow(time.Now()) {
s.send(errorMsg("too_many_rooms"))
return nil
}
if err := s.hub.startBotRoom(s, difficulty); err != nil {
slog.Error("start bot room", "session", s.id, "err", err)
s.send(errorMsg("room_start_failed"))
s.send(roomCreateError(s.id, err))
}
case *noituv1.ClientMessage_CreateRoom:
@@ -429,7 +448,7 @@ func (s *session) dispatch(msg *noituv1.ClientMessage) error {
return nil
}
if err := s.hub.createRoom(s); err != nil {
s.send(errorMsg("room_start_failed"))
s.send(roomCreateError(s.id, err))
}
case *noituv1.ClientMessage_JoinRoom:
@@ -492,6 +511,18 @@ func (s *session) dispatch(msg *noituv1.ClientMessage) error {
return nil
}
// roomCreateError names the refusal a room could not be opened for. A full
// server is the player's business — they should wait, not retry at once — and
// anything else is the server's, logged here because the client is only told
// that it failed.
func roomCreateError(sessionID string, err error) *noituv1.ServerMessage {
if errors.Is(err, errServerFull) {
return errorMsg("server_full")
}
slog.Error("open room", "session", sessionID, "err", err)
return errorMsg("room_start_failed")
}
// toRoom forwards one lobby action to the room this connection is seated in.
//
// Rate-limited like a submission: every accepted action is broadcast to every
+1
View File
@@ -230,6 +230,7 @@ export const errorMessages = {
room_idle_closed: 'Phòng đã đóng vì không có ván nào được bắt đầu.',
room_not_found: 'Không tìm thấy phòng với mã này.',
room_start_failed: 'Không thể tạo phòng. Hãy thử lại.',
server_full: 'Máy chủ đang quá tải. Hãy thử lại sau ít phút.',
server_restarting: 'Máy chủ đang khởi động lại. Hãy thử lại sau giây lát.',
session_not_resumable: 'Không khôi phục được ván đấu trước.',
too_fast: 'Bạn thao tác quá nhanh. Chậm lại một chút nhé.',