diff --git a/README.md b/README.md index e487d10..93371d2 100644 --- a/README.md +++ b/README.md @@ -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 diff --git a/docs/deployment.md b/docs/deployment.md index 1edac91..c474edd 100644 --- a/docs/deployment.md +++ b/docs/deployment.md @@ -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 diff --git a/server/cmd/build-dictionary/main.go b/server/cmd/build-dictionary/main.go index 4049f5c..dcd6f77 100644 --- a/server/cmd/build-dictionary/main.go +++ b/server/cmd/build-dictionary/main.go @@ -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) } diff --git a/server/cmd/noitu-server/main.go b/server/cmd/noitu-server/main.go index 5171d91..9b9fac0 100644 --- a/server/cmd/noitu-server/main.go +++ b/server/cmd/noitu-server/main.go @@ -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 diff --git a/server/internal/dictionary/store.go b/server/internal/dictionary/store.go index 91c23b0..8c44b30 100644 --- a/server/internal/dictionary/store.go +++ b/server/internal/dictionary/store.go @@ -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() } diff --git a/server/internal/wsapi/convert.go b/server/internal/wsapi/convert.go index bb46c0f..b2dfa5a 100644 --- a/server/internal/wsapi/convert.go +++ b/server/internal/wsapi/convert.go @@ -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), diff --git a/server/internal/wsapi/convert_test.go b/server/internal/wsapi/convert_test.go index f65d4a9..a3c4e7d 100644 --- a/server/internal/wsapi/convert_test.go +++ b/server/internal/wsapi/convert_test.go @@ -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) } diff --git a/server/internal/wsapi/hub.go b/server/internal/wsapi/hub.go index 75d82ba..b299798 100644 --- a/server/internal/wsapi/hub.go +++ b/server/internal/wsapi/hub.go @@ -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) { diff --git a/server/internal/wsapi/hub_test.go b/server/internal/wsapi/hub_test.go index 1cbaa69..0e9124d 100644 --- a/server/internal/wsapi/hub_test.go +++ b/server/internal/wsapi/hub_test.go @@ -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. diff --git a/server/internal/wsapi/limits_test.go b/server/internal/wsapi/limits_test.go new file mode 100644 index 0000000..f37b2e6 --- /dev/null +++ b/server/internal/wsapi/limits_test.go @@ -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) +} diff --git a/server/internal/wsapi/room.go b/server/internal/wsapi/room.go index dd84408..34eb6fd 100644 --- a/server/internal/wsapi/room.go +++ b/server/internal/wsapi/room.go @@ -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 { diff --git a/server/internal/wsapi/server.go b/server/internal/wsapi/server.go index b94f682..2fc6da8 100644 --- a/server/internal/wsapi/server.go +++ b/server/internal/wsapi/server.go @@ -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) diff --git a/server/internal/wsapi/session.go b/server/internal/wsapi/session.go index 1707eb6..17b194f 100644 --- a/server/internal/wsapi/session.go +++ b/server/internal/wsapi/session.go @@ -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 diff --git a/web/src/lib/i18n/vi.js b/web/src/lib/i18n/vi.js index d42558e..6b43b6b 100644 --- a/web/src/lib/i18n/vi.js +++ b/web/src/lib/i18n/vi.js @@ -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é.',