Author SHA1 Message Date
woggioni ebeccc937d optimized dockerfile for caching
CI / Build and push docker image (push) Successful in 1m14s
2026-09-18 19:28:33 +08:00
woggioni e0896b4a95 Configure CORS headers from environment variables
CI / Build and push docker image (push) Successful in 3m24s
2026-09-18 19:18:27 +08:00
woggioni e2091ee4df updated Docker builder 2026-09-18 19:18:04 +08:00
woggioni 61f3539c4e Reconnect the game websocket after connectivity loss
The socket had no recovery path: a mid-game drop left the table showing
stale state with no indication, and plays were silently swallowed by the
dead channel. Surface the server close code from ws::connect and add a
reconnect driver in the game page: transient losses retry with
exponential backoff (capped, then a manual Retry), while deliberate
closes (session expired, game gone) stop retrying. Reconnects are free
resyncs because the server pushes a full state snapshot on connect.

Also gate card clicks while disconnected, show a connection banner, and
drop the ticking interval and pending retries on unmount.
2026-09-18 19:16:39 +08:00
woggioni 8fea4fac74 Fix deadline consumer startup under granian RSGI
The mixin received the event loop from Kaya but ignored it, calling
asyncio.get_running_loop() instead. Under granian RSGI __rsgi_init__ runs
before the loop starts, so that raised RuntimeError and killed the
worker. Use the loop passed to setup(), falling back to the running loop
for the lazy calls from sync_deadline().
2026-09-18 19:16:31 +08:00
woggioni c5b6c84408 Show own captured and scopa counts in game view 2026-09-18 19:16:17 +08:00
woggioni 2d1a13f663 Drive timeouts from a shared Redis deadline queue
Turn auto-play and hand-end auto-continue were process-local asyncio
tasks armed only by client connects and state broadcasts: with no
sockets connected the next turn's timer was never armed, a hand-end
timer died with its worker, and neither survived a pod restart.

Deadlines are now driven by the absolute timestamps persisted on the
game state and enqueued in a shared Redis sorted set. Every worker runs
a consumer that fires due entries under the per-game lock after
revalidating them against the live state, so timeouts no longer depend
on any player being connected and survive the death of any worker.
Delivery is at-least-once: entries are removed only after processing,
and revalidation makes duplicate deliveries no-ops.

Queue entries carry the deadline as integer epoch milliseconds, which
also serves as the revalidation token, and the score derives from the
same value.
2026-09-18 19:16:07 +08:00
woggioni c1e7f70a3b Add configurable napola rule with instant win on a full denari sweep 2026-09-18 19:14:54 +08:00
woggioni bf04a9b38d Auto-dismiss error toasts after 10 seconds 2026-09-18 19:14:49 +08:00
34 changed files with 1477 additions and 188 deletions
+1
View File
@@ -32,6 +32,7 @@ jobs:
uses: docker/build-push-action@v6 uses: docker/build-push-action@v6
with: with:
context: . context: .
builder: multiplatform-builder
file: server/Dockerfile file: server/Dockerfile
platforms: linux/amd64 platforms: linux/amd64
push: true push: true
+2 -1
View File
@@ -32,7 +32,8 @@ echo "127.0.0.1 mockoauth" | sudo tee -a /etc/hosts
When a hand ends but the match is not decided, the game pauses on a When a hand ends but the match is not decided, the game pauses on a
**scoring summary screen**: every player sees how each category was won **scoring summary screen**: every player sees how each category was won
(carte, denara, settebello, primiera, scope) with the running totals and (carte, denara, settebello, primiera, scope — plus napola when enabled)
with the running totals and
must click "Understood" before the next hand is dealt. If someone is away must click "Understood" before the next hand is dealt. If someone is away
the next hand is dealt automatically after `HAND_ACK_TIMEOUT_SECONDS` the next hand is dealt automatically after `HAND_ACK_TIMEOUT_SECONDS`
(default 30s). The match-ending hand is explained on the final screen. (default 30s). The match-ending hand is explained on the final screen.
+10
View File
@@ -61,6 +61,16 @@ data:
GAME_TTL_SECONDS: "86400" GAME_TTL_SECONDS: "86400"
HAND_ACK_TIMEOUT_SECONDS: "30" HAND_ACK_TIMEOUT_SECONDS: "30"
TURN_TIMEOUT_SECONDS: "30" TURN_TIMEOUT_SECONDS: "30"
# CORS (kaya-cors' CorsMixin). Disabled unless CORS_ALLOW_ORIGINS or
# CORS_ALLOW_ORIGIN_REGEX is set — unneeded when the SPA and the API are
# served from the same origin. See server/.env.example for details.
# CORS_ALLOW_ORIGINS: "https://example.com,https://app.example.com" # or "*"
# CORS_ALLOW_ORIGIN_REGEX: 'https://tavolo-[a-z0-9-]+\.vercel\.app'
# CORS_ALLOW_METHODS: "GET,POST" # default: GET; "*" = all
# CORS_ALLOW_HEADERS: "Authorization,Content-Type" # "*" mirrors the request
# CORS_ALLOW_CREDENTIALS: "false"
# CORS_EXPOSE_HEADERS: ""
# CORS_MAX_AGE: "600"
# OIDC (provider lives in another namespace). # OIDC (provider lives in another namespace).
OIDC_CLIENT_ID: tavolo OIDC_CLIENT_ID: tavolo
OIDC_POST_LOGIN_REDIRECT: / OIDC_POST_LOGIN_REDIRECT: /
+9
View File
@@ -113,6 +113,15 @@ services:
REDIS_URL: redis://redis:6379/0 REDIS_URL: redis://redis:6379/0
HAND_ACK_TIMEOUT_SECONDS: ${HAND_ACK_TIMEOUT_SECONDS:-30} HAND_ACK_TIMEOUT_SECONDS: ${HAND_ACK_TIMEOUT_SECONDS:-30}
TURN_TIMEOUT_SECONDS: ${TURN_TIMEOUT_SECONDS:-30} TURN_TIMEOUT_SECONDS: ${TURN_TIMEOUT_SECONDS:-30}
# CORS is disabled unless CORS_ALLOW_ORIGINS or CORS_ALLOW_ORIGIN_REGEX
# is set (see server/.env.example for the full list of options).
CORS_ALLOW_ORIGINS: ${CORS_ALLOW_ORIGINS:-}
CORS_ALLOW_ORIGIN_REGEX: ${CORS_ALLOW_ORIGIN_REGEX:-}
CORS_ALLOW_METHODS: ${CORS_ALLOW_METHODS:-}
CORS_ALLOW_HEADERS: ${CORS_ALLOW_HEADERS:-}
CORS_ALLOW_CREDENTIALS: ${CORS_ALLOW_CREDENTIALS:-}
CORS_EXPOSE_HEADERS: ${CORS_EXPOSE_HEADERS:-}
CORS_MAX_AGE: ${CORS_MAX_AGE:-}
ports: ports:
- "127.0.0.1:${APP_PORT:-8080}:8080" - "127.0.0.1:${APP_PORT:-8080}:8080"
+22
View File
@@ -39,6 +39,28 @@ TURN_TIMEOUT_SECONDS=30
# schema). Unset logs DEBUG to the console. # schema). Unset logs DEBUG to the console.
#LOGGING_CONFIG=/path/to/logging.yaml #LOGGING_CONFIG=/path/to/logging.yaml
# CORS (via kaya-cors' CorsMixin; same semantics as Starlette's
# CORSMiddleware). Disabled unless CORS_ALLOW_ORIGINS or
# CORS_ALLOW_ORIGIN_REGEX is set — the app serves the SPA and the API from
# the same origin, so no CORS headers are needed by default.
# Comma-separated list of origins allowed to make cross-origin requests,
# or "*" for any origin:
#CORS_ALLOW_ORIGINS=https://example.com,https://app.example.com
# Optional regex (fullmatch) allowed origins are additionally checked
# against — handy for dynamic preview URLs:
#CORS_ALLOW_ORIGIN_REGEX=https://tavolo-[a-z0-9-]+\.vercel\.app
# Comma-separated allowed methods, or "*" for all (default GET):
#CORS_ALLOW_METHODS=GET,POST
# Comma-separated allowed request headers, or "*" to mirror back whatever
# the browser requests (default: only the CORS-safelisted headers):
#CORS_ALLOW_HEADERS=Authorization,Content-Type
# Allow cookies/credentials on cross-origin requests (1/true/yes/on):
#CORS_ALLOW_CREDENTIALS=false
# Comma-separated response headers exposed to the browser:
#CORS_EXPOSE_HEADERS=
# Seconds browsers may cache the preflight response (default 600):
#CORS_MAX_AGE=600
# App server # App server
APP_HOST=0.0.0.0 APP_HOST=0.0.0.0
APP_PORT=8080 APP_PORT=8080
+10 -6
View File
@@ -48,16 +48,20 @@ RUN --mount=type=cache,target=/var/cache/apk \
WORKDIR /build WORKDIR /build
COPY server/pyproject.toml server/README.md server/requirements.txt ./
COPY server/src/ ./src/
# aerich migration files are a release artifact: the db-migrate compose
# service runs `aerich upgrade` from this image before the app starts.
COPY server/migrations/ ./migrations/
COPY server/requirements.txt ./
RUN --mount=type=cache,target=/root/.cache/pip \ RUN --mount=type=cache,target=/root/.cache/pip \
python3 -m venv /opt/venv \ python3 -m venv /opt/venv \
&& /opt/venv/bin/pip install --upgrade pip \ && /opt/venv/bin/pip install --upgrade pip \
&& /opt/venv/bin/pip install -r requirements.txt . && /opt/venv/bin/pip install -r requirements.txt
COPY server/pyproject.toml server/README.md ./
COPY server/src/ ./src/
RUN --mount=type=cache,target=/root/.cache/pip \
/opt/venv/bin/pip install .
# aerich migration files are a release artifact: the db-migrate compose
# service runs `aerich upgrade` from this image before the app starts.
COPY server/migrations/ ./migrations/
# --- Runtime --------------------------------------------------------------- # --- Runtime ---------------------------------------------------------------
FROM alpine:3.24 FROM alpine:3.24
+25 -3
View File
@@ -60,7 +60,15 @@ All configuration comes from environment variables (see `.env.example`):
| `GAME_TTL_SECONDS` | `86400` | Sliding TTL of a live game in Redis | | `GAME_TTL_SECONDS` | `86400` | Sliding TTL of a live game in Redis |
| `HAND_ACK_TIMEOUT_SECONDS` | `30` | Seconds the between-hands scoring summary waits for acknowledgements | | `HAND_ACK_TIMEOUT_SECONDS` | `30` | Seconds the between-hands scoring summary waits for acknowledgements |
| `TURN_TIMEOUT_SECONDS` | `30` | Seconds a player has to play before the server plays a random legal card for them | | `TURN_TIMEOUT_SECONDS` | `30` | Seconds a player has to play before the server plays a random legal card for them |
| `DEADLINE_HEARTBEAT_MS` | `1000` | Upper bound on the deadline consumer's poll interval (locally enqueued deadlines fire on time regardless) |
| `LOGGING_CONFIG` | unset | Path to a YAML logging configuration file (see below). Unset logs DEBUG to the console | | `LOGGING_CONFIG` | unset | Path to a YAML logging configuration file (see below). Unset logs DEBUG to the console |
| `CORS_ALLOW_ORIGINS` | unset | Comma-separated origins allowed for cross-origin requests, or `*` for any. CORS is disabled unless this or `CORS_ALLOW_ORIGIN_REGEX` is set |
| `CORS_ALLOW_ORIGIN_REGEX` | unset | Regex (fullmatch) additionally matched against request origins, e.g. `https://tavolo-[a-z0-9-]+\.vercel\.app` |
| `CORS_ALLOW_METHODS` | `GET` | Comma-separated methods allowed for cross-origin requests, or `*` for all |
| `CORS_ALLOW_HEADERS` | unset | Comma-separated request headers allowed in cross-origin requests, or `*` to mirror back the requested ones. The CORS-safelisted headers are always allowed |
| `CORS_ALLOW_CREDENTIALS` | `false` | `1`/`true`/`yes`/`on` allow cookies/credentials on cross-origin requests |
| `CORS_EXPOSE_HEADERS` | unset | Comma-separated response headers exposed to the browser |
| `CORS_MAX_AGE` | `600` | Seconds browsers may cache the preflight response |
| `APP_HOST` / `APP_PORT` | `0.0.0.0` / `8080` | Bind address | | `APP_HOST` / `APP_PORT` | `0.0.0.0` / `8080` | Bind address |
## Logging ## Logging
@@ -111,6 +119,12 @@ loggers:
- `tavolo:game:<uuid>:events` — a pub/sub channel carrying "state changed" - `tavolo:game:<uuid>:events` — a pub/sub channel carrying "state changed"
signals; every open WebSocket reloads the state and pushes the signals; every open WebSocket reloads the state and pushes the
personalized view to its player. personalized view to its player.
- `tavolo:deadlines` — a sorted set (score = due timestamp) of pending
timeouts: turn auto-plays and hand-end auto-continues. Every worker runs
a consumer that fires due entries under the per-game lock, so timeouts
do not depend on any player being connected and survive the death of
any worker (delivery is at-least-once; entries are revalidated against
the live state before firing).
### Postgres (statistics, via Tortoise ORM + aerich migrations) ### Postgres (statistics, via Tortoise ORM + aerich migrations)
@@ -133,7 +147,7 @@ All endpoints except `/api/health`, `/api/docs`, `/api/openapi.json`,
| Method | Path | Description | | Method | Path | Description |
|---|---|---| |---|---|---|
| `GET` | `/api/game-types` | The card games the platform can host (for the creation dropdown) | | `GET` | `/api/game-types` | The card games the platform can host (for the creation dropdown) |
| `POST` | `/api/games` | Create a lobby game. Optional body `{"game_type": "scopone_scientifico", "target_score": 11}`. Returns `{id, join_code}` | | `POST` | `/api/games` | Create a lobby game. Optional body `{"game_type": "scopone_scientifico", "target_score": 11, "napola": true}`. Returns `{id, join_code}` |
| `POST` | `/api/games/join` | Join with `{"code": "ABC123"}`. The fourth player triggers the deal | | `POST` | `/api/games/join` | Join with `{"code": "ABC123"}`. The fourth player triggers the deal |
| `GET` | `/api/games/{id}` | Personalized snapshot (only your own hand is visible) | | `GET` | `/api/games/{id}` | Personalized snapshot (only your own hand is visible) |
| `GET` | `/api/me/matches` | Cursor-paginated match history with final scores (`?limit=&cursor=&game_type=`) | | `GET` | `/api/me/matches` | Cursor-paginated match history with final scores (`?limit=&cursor=&game_type=`) |
@@ -181,7 +195,9 @@ player on turn does not move before it, the server plays a random legal
card for them (picking one of the legal captures at random when a capture card for them (picking one of the legal captures at random when a capture
is required), so a disconnected or idle player cannot stall the match. The is required), so a disconnected or idle player cannot stall the match. The
timeout is `TURN_TIMEOUT_SECONDS` (default 30); the auto-played move is timeout is `TURN_TIMEOUT_SECONDS` (default 30); the auto-played move is
broadcast like any other. broadcast like any other. Deadlines fire from the shared `tavolo:deadlines`
queue (see above), not from timers tied to client connections, so the
match keeps progressing even with every player disconnected.
### Hand-end summary ### Hand-end summary
@@ -203,6 +219,11 @@ must dismiss. A play attempted in this phase is rejected with an
- Hand points: `carte` (most captured cards), `denara` (most diamonds), - Hand points: `carte` (most captured cards), `denara` (most diamonds),
`settebello` (7♦), `primiera` (best 7/6/5/4 per suit, all four suits `settebello` (7♦), `primiera` (best 7/6/5/4 per suit, all four suits
required), plus one point per scopa. Ties award nothing. required), plus one point per scopa. Ties award nothing.
- Optional *napola* rule (per-game `napola` flag on `POST /api/games`,
default on): the longest run of consecutive denari starting from the
ace scores one point per card once it reaches three cards (A-2-3 = 3,
A-2-3-4 = 4, …). A team that captures the whole denari suit (ace to
king) wins the match instantly, regardless of the score.
- The match ends when a team reaches the target score (default 11, - The match ends when a team reaches the target score (default 11,
configurable per game) with a clear lead; a tie at or above the target is configurable per game) with a clear lead; a tie at or above the target is
broken by another hand. broken by another hand.
@@ -252,7 +273,8 @@ src/tavolo/
├── aerich_config.py # aerich CLI configuration ├── aerich_config.py # aerich CLI configuration
├── models.py # Match, MatchPlayer (Postgres) ├── models.py # Match, MatchPlayer (Postgres)
├── stats.py # finished match -> Postgres persistence ├── stats.py # finished match -> Postgres persistence
├── store.py # Redis / in-memory live-game store ├── store.py # Redis / in-memory live-game store (+ deadline queue)
├── deadlines.py # connection-independent timeout scheduler
├── ws.py # WebSocket live-play endpoint ├── ws.py # WebSocket live-play endpoint
├── game/ ├── game/
│ ├── state.py # GameState / PlayerState / Card, JSON (de)serialization │ ├── state.py # GameState / PlayerState / Card, JSON (de)serialization
+1
View File
@@ -10,6 +10,7 @@ readme = "README.md"
requires-python = ">=3.10" requires-python = ">=3.10"
dependencies = [ dependencies = [
"kaya-core", "kaya-core",
"kaya-cors",
"kaya-session", "kaya-session",
"kaya-session-redis", "kaya-session-redis",
"kaya-oidc", "kaya-oidc",
+3
View File
@@ -52,11 +52,14 @@ iso8601==2.1.0
# via tortoise-orm # via tortoise-orm
kaya-core==0.0.3 kaya-core==0.0.3
# via # via
# kaya-cors
# kaya-oidc # kaya-oidc
# kaya-openapi # kaya-openapi
# kaya-rsgi # kaya-rsgi
# kaya-session # kaya-session
# tavolo (pyproject.toml) # tavolo (pyproject.toml)
kaya-cors==0.0.3
# via tavolo (pyproject.toml)
kaya-oidc==0.0.3 kaya-oidc==0.0.3
# via tavolo (pyproject.toml) # via tavolo (pyproject.toml)
kaya-openapi==0.0.3 kaya-openapi==0.0.3
+44 -3
View File
@@ -11,6 +11,9 @@ Assembles the :class:`~kaya.core.KayaApp` with four mixins:
- :class:`~kaya.openapi.OpenAPIMixin` (serves the OpenAPI document at - :class:`~kaya.openapi.OpenAPIMixin` (serves the OpenAPI document at
``/api/openapi.json`` and a Swagger UI at ``/api/docs``) ``/api/openapi.json`` and a Swagger UI at ``/api/docs``)
A :class:`~kaya.cors.CorsMixin` is prepended when CORS is configured via the
``CORS_*`` environment variables (see :mod:`tavolo.config`).
Live games are kept in :data:`game_store` (Redis when configured, in-memory Live games are kept in :data:`game_store` (Redis when configured, in-memory
otherwise). Routes and the websocket handlers are registered by importing otherwise). Routes and the websocket handlers are registered by importing
their modules at the bottom; imports must happen after ``app`` is built. their modules at the bottom; imports must happen after ``app`` is built.
@@ -19,15 +22,18 @@ from __future__ import annotations
from importlib.metadata import version as _pkg_version from importlib.metadata import version as _pkg_version
from logging import getLogger from logging import getLogger
from typing import Optional
from kaya.core import KayaApp from kaya.core import KayaApp, KayaMixin
from kaya.cors import CorsMixin
from kaya.oidc import OIDCConfig, OIDCMixin from kaya.oidc import OIDCConfig, OIDCMixin
from kaya.openapi import OpenAPIMixin from kaya.openapi import OpenAPIMixin
from kaya.session import InMemorySessionStore, SessionMixin, SessionStore from kaya.session import InMemorySessionStore, SessionMixin, SessionStore
from kaya.session.redis import RedisSessionStore from kaya.session.redis import RedisSessionStore
from redis.asyncio import Redis from redis.asyncio import Redis
from .config import settings from .config import Settings, settings
from .deadlines import DeadlineSchedulerMixin
from .logging_config import configure_logging from .logging_config import configure_logging
from .store import GameStore, InMemoryGameStore, RedisGameStore from .store import GameStore, InMemoryGameStore, RedisGameStore
from .tortoise_mixin import TortoiseMixin from .tortoise_mixin import TortoiseMixin
@@ -35,6 +41,27 @@ from .tortoise_mixin import TortoiseMixin
configure_logging(settings.logging_config) configure_logging(settings.logging_config)
log = getLogger(__name__) log = getLogger(__name__)
def cors_mixin_from_settings(settings: Settings) -> Optional[CorsMixin]:
"""Build a :class:`~kaya.cors.CorsMixin` from the CORS settings.
Returns ``None`` — CORS disabled — unless at least one of
``CORS_ALLOW_ORIGINS`` / ``CORS_ALLOW_ORIGIN_REGEX`` is configured.
Settings left unset fall back to the mixin's own defaults.
"""
if settings.cors_allow_origins is None and settings.cors_allow_origin_regex is None:
return None
return CorsMixin(
allow_origins=settings.cors_allow_origins or (),
allow_origin_regex=settings.cors_allow_origin_regex,
allow_methods=settings.cors_allow_methods or ("GET",),
allow_headers=settings.cors_allow_headers or (),
allow_credentials=settings.cors_allow_credentials,
expose_headers=settings.cors_expose_headers or (),
max_age=settings.cors_max_age,
)
session_store: SessionStore session_store: SessionStore
if settings.redis_url is not None: if settings.redis_url is not None:
# Lazy client: no connection is opened until a session is actually # Lazy client: no connection is opened until a session is actually
@@ -76,7 +103,21 @@ tortoise_mixin = TortoiseMixin(
skip_paths=frozenset({"/api/health", "/api/docs", "/api/openapi.json"}), skip_paths=frozenset({"/api/health", "/api/docs", "/api/openapi.json"}),
) )
app = KayaApp(mixins=[session_mixin, oidc_mixin, tortoise_mixin, openapi_mixin]) mixins: list[KayaMixin] = [session_mixin, oidc_mixin, tortoise_mixin, openapi_mixin,
DeadlineSchedulerMixin(game_store)]
cors_mixin = cors_mixin_from_settings(settings)
if cors_mixin is not None:
# First in the list: preflight requests are answered before the session
# and OIDC hooks run.
mixins.insert(0, cors_mixin)
log.info(
"CORS enabled: origins=%s origin_regex=%s credentials=%s",
settings.cors_allow_origins,
settings.cors_allow_origin_regex,
settings.cors_allow_credentials,
)
app = KayaApp(mixins=mixins)
log.debug( log.debug(
"timeouts: hand_ack=%ds turn=%ds", "timeouts: hand_ack=%ds turn=%ds",
settings.hand_ack_timeout_seconds, settings.hand_ack_timeout_seconds,
+46 -1
View File
@@ -7,7 +7,7 @@ from __future__ import annotations
import os import os
from dataclasses import dataclass from dataclasses import dataclass
from typing import Optional from typing import Optional, Tuple
from urllib.parse import quote from urllib.parse import quote
@@ -20,6 +20,25 @@ def _env(name: str, default: Optional[str] = None) -> str:
return value return value
def _env_list(name: str) -> Optional[Tuple[str, ...]]:
"""Parse a comma-separated environment variable into a tuple of values.
Items are stripped and empty items dropped. Unset or empty variables
yield ``None``.
"""
value = os.environ.get(name)
if value is None or value.strip() == "":
return None
return tuple(part.strip() for part in value.split(",") if part.strip())
def _env_bool(name: str, default: bool = False) -> bool:
value = os.environ.get(name)
if value is None or value == "":
return default
return value.strip().lower() in ("1", "true", "yes", "on")
def _database_url_from_parts(engine: str, def _database_url_from_parts(engine: str,
user: str, user: str,
password: Optional[str], password: Optional[str],
@@ -76,9 +95,24 @@ class Settings:
# Seconds a player has to play before the server plays a random legal # Seconds a player has to play before the server plays a random legal
# card for them (covering disconnects and idle players). # card for them (covering disconnects and idle players).
turn_timeout_seconds: int turn_timeout_seconds: int
# Upper bound on how long the deadline consumer sleeps between polls.
# Locally enqueued deadlines wake the consumer immediately; the
# heartbeat only bounds the discovery delay for deadlines enqueued by
# other workers.
deadline_heartbeat_ms: int
# Path to a YAML logging configuration file (logging.config.dictConfig # Path to a YAML logging configuration file (logging.config.dictConfig
# schema). Unset uses the built-in default: DEBUG to the console. # schema). Unset uses the built-in default: DEBUG to the console.
logging_config: Optional[str] logging_config: Optional[str]
# CORS (kaya-cors' CorsMixin). Disabled unless CORS_ALLOW_ORIGINS or
# CORS_ALLOW_ORIGIN_REGEX is set; the app serves the SPA and the API
# from the same origin, so no CORS headers are needed by default.
cors_allow_origins: Optional[Tuple[str, ...]]
cors_allow_origin_regex: Optional[str]
cors_allow_methods: Optional[Tuple[str, ...]]
cors_allow_headers: Optional[Tuple[str, ...]]
cors_allow_credentials: bool
cors_expose_headers: Optional[Tuple[str, ...]]
cors_max_age: int
@staticmethod @staticmethod
def from_env() -> "Settings": def from_env() -> "Settings":
@@ -113,7 +147,18 @@ class Settings:
static_dir=_env("STATIC_DIR", "web/dist"), static_dir=_env("STATIC_DIR", "web/dist"),
hand_ack_timeout_seconds=int(_env("HAND_ACK_TIMEOUT_SECONDS", "30")), hand_ack_timeout_seconds=int(_env("HAND_ACK_TIMEOUT_SECONDS", "30")),
turn_timeout_seconds=int(_env("TURN_TIMEOUT_SECONDS", "30")), turn_timeout_seconds=int(_env("TURN_TIMEOUT_SECONDS", "30")),
deadline_heartbeat_ms=int(_env("DEADLINE_HEARTBEAT_MS", "1000")),
logging_config=os.environ.get("LOGGING_CONFIG"), logging_config=os.environ.get("LOGGING_CONFIG"),
# CORS is disabled unless CORS_ALLOW_ORIGINS (a comma-separated
# list of origins, or "*" for any) or CORS_ALLOW_ORIGIN_REGEX
# is set.
cors_allow_origins=_env_list("CORS_ALLOW_ORIGINS"),
cors_allow_origin_regex=os.environ.get("CORS_ALLOW_ORIGIN_REGEX") or None,
cors_allow_methods=_env_list("CORS_ALLOW_METHODS"),
cors_allow_headers=_env_list("CORS_ALLOW_HEADERS"),
cors_allow_credentials=_env_bool("CORS_ALLOW_CREDENTIALS"),
cors_expose_headers=_env_list("CORS_EXPOSE_HEADERS"),
cors_max_age=int(_env("CORS_MAX_AGE", "600")),
) )
+296
View File
@@ -0,0 +1,296 @@
"""Deadline-driven timeouts, independent of player connections.
Both in-match timeouts — the per-turn auto-play (``turn_deadline``) and
the hand-end summary auto-continue (``hand_end_deadline``) — are driven by
the absolute deadlines persisted on the game state, never by which players
(or whether any players) are connected.
Every mutation that sets a deadline enqueues an entry in the store's
shared deadline queue (a Redis sorted set in production, see
:mod:`tavolo.store`), and a background consumer running on **every**
worker polls the queue for due entries. An entry records the phase, hand,
turn and deadline (as integer epoch milliseconds) it was enqueued for;
before acting, the consumer revalidates all of it against the live state
under the per-game lock, so entries that were overtaken by events (a play
landed in time, the hand was acknowledged, the deadline moved) are simply
discarded.
Delivery is at-least-once: an entry is removed from the queue only after
it has been processed. If a worker dies mid-processing, the entry stays in
Redis and another worker's consumer picks it up — the lock plus
revalidation make the duplicate delivery a no-op. Entries whose game has
expired are dropped the first time they fire, so the queue is
self-cleaning.
"""
from __future__ import annotations
import asyncio
import json
import time
from datetime import datetime
from logging import getLogger
from typing import Any, Dict, Optional
from kaya.core import KayaApp, KayaMixin
from .config import settings
from .game import engine
from .game.errors import GameError
from .game.state import PHASE_HAND_END, PHASE_PLAYING, PHASE_FINISHED, GameState
from .stats import save_match_result
from .store import GameStore
log = getLogger(__name__)
# Entry kinds enqueued in the deadline queue.
KIND_TURN = "turn"
KIND_HAND_END = "hand_end"
# One consumer task and its wake-up event per event loop (tests run each
# test on a fresh loop).
_consumers: Dict[asyncio.AbstractEventLoop, asyncio.Task] = {}
_wake_events: Dict[asyncio.AbstractEventLoop, asyncio.Event] = {}
def encode(entry: Dict[str, Any]) -> str:
"""Canonical queue-member encoding for a deadline entry."""
return json.dumps(entry, sort_keys=True)
def _decode(member: Any) -> Optional[Dict[str, Any]]:
if isinstance(member, bytes):
member = member.decode("utf-8")
if not isinstance(member, str):
return None
try:
entry = json.loads(member)
except ValueError:
return None
return entry if isinstance(entry, dict) else None
def _deadline_ms(iso: Optional[str]) -> Optional[int]:
"""Epoch milliseconds for an ISO-8601 deadline, ``None`` when absent
or unparseable. Queue entries carry this integer (never the ISO
string) as their revalidation token."""
if not iso:
return None
try:
return int(datetime.fromisoformat(iso).timestamp() * 1000)
except ValueError:
return None
async def sync_deadline(store: GameStore, state: GameState) -> None:
"""Enqueue the deadline the current state carries, if any.
Called after every mutation that can set a deadline (plays, acks, game
start) and as a backstop when a client connects. Enqueueing is
idempotent: an identical entry is already queued with the same due
time, so re-adding it changes nothing.
"""
entry: Optional[Dict[str, Any]] = None
due_ms: Optional[int] = None
if state.phase == PHASE_PLAYING and state.turn_deadline:
due_ms = _deadline_ms(state.turn_deadline)
entry = {
"game_id": state.id,
"kind": KIND_TURN,
"hand": state.hand_number,
"turn": state.turn,
"deadline": due_ms,
}
elif state.phase == PHASE_HAND_END and state.hand_end_deadline:
due_ms = _deadline_ms(state.hand_end_deadline)
entry = {
"game_id": state.id,
"kind": KIND_HAND_END,
"hand": state.hand_number,
"deadline": due_ms,
}
if entry is None or due_ms is None:
if entry is not None:
log.warning("game %s: unparseable deadline", state.id)
return
ensure_consumer(store)
# The score derives from the same value carried in the member, so the
# two can never disagree.
await store.add_deadline(encode(entry), due_ms / 1000)
wake = _wake_events.get(asyncio.get_running_loop())
if wake is not None:
wake.set()
async def finalize_mutation(store: GameStore, state: GameState) -> None:
"""Persist a successful mutation, notify subscribers and enqueue the
next deadline.
Callers must hold the per-game lock. Handles the terminal transition:
the match result is written to Postgres once (guarded by
``stats_saved``).
"""
if state.phase == PHASE_FINISHED:
await save_match_result(state)
log.info(
"game %s finished: team %s wins %d-%d",
state.id,
"A" if state.winner == 0 else "B",
state.scores[0],
state.scores[1],
)
await store.save(state)
await store.publish(state.id)
await sync_deadline(store, state)
async def process_due(store: GameStore, member: Any) -> None:
"""Fire a single due deadline entry.
Revalidates the entry against the live state under the per-game lock;
stale or foreign entries are discarded without effect. The entry is
removed from the queue once handled (including "nothing to do"); if
handling fails (e.g. the lock cannot be acquired), the entry is left
in the queue so another consumer retries it.
"""
entry = _decode(member)
if entry is None:
log.warning("deadline consumer: dropping malformed entry %r", member)
await store.remove_deadline(member)
return
game_id = entry.get("game_id")
kind = entry.get("kind")
if not isinstance(game_id, str):
await store.remove_deadline(member)
return
async with store.lock(game_id):
state = await store.load(game_id)
if state is not None:
if kind == KIND_TURN:
await _fire_turn(store, state, entry)
elif kind == KIND_HAND_END:
await _fire_hand_end(store, state, entry)
await store.remove_deadline(member)
async def _fire_turn(store: GameStore, state: GameState, entry: Dict[str, Any]) -> None:
if (
state.phase != PHASE_PLAYING
or state.hand_number != entry.get("hand")
or state.turn != entry.get("turn")
or _deadline_ms(state.turn_deadline) != entry.get("deadline")
):
return
seat = state.turn
try:
engine.auto_play(state)
except GameError:
return
log.info(
"game %s: auto-played for %s (turn timeout, hand %d)",
state.id,
state.players[seat].sub if seat < len(state.players) else "?",
entry.get("hand"),
)
await finalize_mutation(store, state)
async def _fire_hand_end(store: GameStore, state: GameState, entry: Dict[str, Any]) -> None:
if (
state.phase != PHASE_HAND_END
or state.hand_number != entry.get("hand")
or _deadline_ms(state.hand_end_deadline) != entry.get("deadline")
):
return
for player in state.players:
engine.acknowledge_hand(state, player.sub)
log.info(
"game %s: hand %d auto-advanced after the acknowledgement timeout",
state.id,
entry.get("hand"),
)
await finalize_mutation(store, state)
# --- consumer lifecycle -------------------------------------------------------
def ensure_consumer(
store: GameStore, loop: Optional[asyncio.AbstractEventLoop] = None
) -> None:
"""Start the deadline consumer on the given (or running) loop if not
yet running.
Called lazily whenever a deadline is enqueued (the ASGI test transport
never fires the lifespan hooks, so the mixin's ``setup`` alone is not
enough) and on application startup. The explicit ``loop`` matters at
startup: under RSGI granian calls ``setup`` before the loop runs, so
``asyncio.get_running_loop()`` would fail there.
"""
if loop is None:
loop = asyncio.get_running_loop()
for old in list(_consumers):
if old.is_closed():
_consumers.pop(old, None)
_wake_events.pop(old, None)
task = _consumers.get(loop)
if task is None or task.done():
_wake_events[loop] = asyncio.Event()
_consumers[loop] = loop.create_task(_run(store, loop))
log.debug("deadline consumer started")
def stop_consumer(loop: asyncio.AbstractEventLoop) -> None:
task = _consumers.pop(loop, None)
_wake_events.pop(loop, None)
if task is not None:
task.cancel()
async def _run(store: GameStore, loop: asyncio.AbstractEventLoop) -> None:
wake = _wake_events[loop]
heartbeat = settings.deadline_heartbeat_ms / 1000
while True:
# Clear before polling so an enqueue racing the poll re-wakes us.
wake.clear()
delay = heartbeat
try:
for member in await store.due_deadlines(time.time()):
try:
await process_due(store, member)
except asyncio.CancelledError:
raise
except Exception:
# Left in the queue; retried on the next pass.
log.exception("deadline consumer: failed to process %r", member)
next_due = await store.next_deadline()
if next_due is not None:
delay = max(0.0, min(heartbeat, next_due - time.time()))
except asyncio.CancelledError:
raise
except Exception:
log.exception("deadline consumer: poll failed; retrying")
try:
await asyncio.wait_for(wake.wait(), timeout=delay)
except asyncio.TimeoutError:
pass
class DeadlineSchedulerMixin(KayaMixin):
"""Run the deadline consumer for the whole app lifetime.
Every worker (and every pod) runs the same consumer; coordination
happens exclusively through the shared deadline queue and the per-game
locks, so any worker may fire any game's deadline.
"""
def __init__(self, store: GameStore) -> None:
self._store = store
def apply(self, app: KayaApp) -> None:
pass
def setup(self, loop: asyncio.AbstractEventLoop) -> None:
ensure_consumer(self._store, loop)
def shutdown(self, loop: asyncio.AbstractEventLoop) -> None:
stop_consumer(loop)
+44
View File
@@ -22,6 +22,10 @@ Rules implemented
cards), ``settebello`` (the 7 of diamonds), ``primiera`` (best cards), ``settebello`` (the 7 of diamonds), ``primiera`` (best
seven/five/four/three card of each suit, all four suits required), plus seven/five/four/three card of each suit, all four suits required), plus
one point per ``scopa``. Ties on carte/denara/primiera award nothing. one point per ``scopa``. Ties on carte/denara/primiera award nothing.
* Optional ``napola`` rule (enabled by default): the longest run of
consecutive denari starting from the ace scores one point per card when
it reaches at least three cards (A-2-3 = 3, A-2-3-4 = 4, ...). A team
capturing the whole denari suit (ace to king) wins the match instantly.
* The match ends when a team reaches the target score with a clear lead; a * The match ends when a team reaches the target score with a clear lead; a
tie at or above the target is broken by playing another hand. tie at or above the target is broken by playing another hand.
""" """
@@ -135,6 +139,7 @@ def create_game(
hand_ack_timeout: int = DEFAULT_HAND_ACK_TIMEOUT_SECONDS, hand_ack_timeout: int = DEFAULT_HAND_ACK_TIMEOUT_SECONDS,
turn_timeout: int = DEFAULT_TURN_TIMEOUT_SECONDS, turn_timeout: int = DEFAULT_TURN_TIMEOUT_SECONDS,
game_type: str = "scopone_scientifico", game_type: str = "scopone_scientifico",
napola: bool = True,
) -> GameState: ) -> GameState:
"""Create a lobby game with the creator seated first.""" """Create a lobby game with the creator seated first."""
if target_score < 1 or target_score > 100: if target_score < 1 or target_score > 100:
@@ -145,6 +150,7 @@ def create_game(
creator_sub=creator_sub, creator_sub=creator_sub,
game_type=game_type, game_type=game_type,
target_score=target_score, target_score=target_score,
napola=napola,
phase=PHASE_LOBBY, phase=PHASE_LOBBY,
players=[PlayerState(sub=creator_sub, name=creator_name, seat=0)], players=[PlayerState(sub=creator_sub, name=creator_name, seat=0)],
hand_ack_timeout=hand_ack_timeout, hand_ack_timeout=hand_ack_timeout,
@@ -350,6 +356,17 @@ def _end_hand(state: GameState) -> None:
a, a,
b, b,
) )
# A full napola (the whole denari suit) wins the match outright,
# regardless of the score.
napola = details.get("napola")
if isinstance(napola, dict):
for team, name in enumerate(TEAM_NAMES):
if napola.get(name) == 10:
state.phase = PHASE_FINISHED
state.winner = team
state.finished_at = datetime.now(timezone.utc).isoformat()
log.info("game %s: team %s swept the denari (napola) and wins", state.id, name)
return
reached = max(a, b) >= state.target_score reached = max(a, b) >= state.target_score
if reached and a != b: if reached and a != b:
state.phase = PHASE_FINISHED state.phase = PHASE_FINISHED
@@ -404,6 +421,22 @@ def primiera_score(captured: Sequence[Card]) -> int:
return sum(best.values()) return sum(best.values())
def napola_score(captured: Sequence[Card]) -> int:
"""Return the napola value of a capture pile.
The longest run of consecutive denari starting from the ace scores one
point per card once it reaches three cards (A-2-3 = 3, A-2-3-4 = 4,
...), so the whole suit (ace to king) is worth 10. Shorter runs score
nothing. Only one team can score a napola: the ace of denari belongs
to exactly one capture pile.
"""
ranks = {card.rank for card in captured if card.suit == "D"}
run = 0
while run + 1 in ranks:
run += 1
return run if run >= 3 else 0
def hand_points(state: GameState) -> Tuple[List[int], Dict[str, object]]: def hand_points(state: GameState) -> Tuple[List[int], Dict[str, object]]:
"""Compute the hand points for both teams (index 0 = team A).""" """Compute the hand points for both teams (index 0 = team A)."""
piles: List[List[Card]] = [[], []] piles: List[List[Card]] = [[], []]
@@ -460,6 +493,16 @@ def hand_points(state: GameState) -> Tuple[List[int], Dict[str, object]]:
"scope": {"A": scope[0], "B": scope[1]}, "scope": {"A": scope[0], "B": scope[1]},
"award": award, "award": award,
} }
# Napola (optional rule): consecutive denari from the ace. A run of 10
# means the team swept the whole suit and wins the match instantly.
if state.napola:
napola = [napola_score(piles[t]) for t in (0, 1)]
for team in (0, 1):
points[team] += napola[team]
award["napola"] = next(
(TEAM_NAMES[t] for t in (0, 1) if napola[t] > 0), None
)
details["napola"] = {"A": napola[0], "B": napola[1]}
return points, details return points, details
@@ -493,6 +536,7 @@ def state_for_player(state: GameState, sub: str) -> Dict[str, object]:
"game_type": state.game_type, "game_type": state.game_type,
"phase": state.phase, "phase": state.phase,
"target_score": state.target_score, "target_score": state.target_score,
"napola": state.napola,
"hand_number": state.hand_number, "hand_number": state.hand_number,
"dealer": state.dealer, "dealer": state.dealer,
"turn": state.turn, "turn": state.turn,
+6
View File
@@ -152,6 +152,10 @@ class GameState:
# Defaults so states serialized before game types existed still load. # Defaults so states serialized before game types existed still load.
game_type: str = "scopone_scientifico" game_type: str = "scopone_scientifico"
target_score: int = DEFAULT_TARGET_SCORE target_score: int = DEFAULT_TARGET_SCORE
# Whether the napola rule is scored (denari run from the ace; a full
# suit wins the match instantly). Default on, also for states
# serialized before the option existed.
napola: bool = True
phase: str = PHASE_LOBBY phase: str = PHASE_LOBBY
players: List[PlayerState] = field(default_factory=list) players: List[PlayerState] = field(default_factory=list)
table: List[Card] = field(default_factory=list) table: List[Card] = field(default_factory=list)
@@ -189,6 +193,7 @@ class GameState:
"creator_sub": self.creator_sub, "creator_sub": self.creator_sub,
"game_type": self.game_type, "game_type": self.game_type,
"target_score": self.target_score, "target_score": self.target_score,
"napola": self.napola,
"phase": self.phase, "phase": self.phase,
"players": [p.to_json() for p in self.players], "players": [p.to_json() for p in self.players],
"table": [c.to_json() for c in self.table], "table": [c.to_json() for c in self.table],
@@ -218,6 +223,7 @@ class GameState:
creator_sub=str(data.get("creator_sub", "")), creator_sub=str(data.get("creator_sub", "")),
game_type=str(data.get("game_type", "scopone_scientifico")), game_type=str(data.get("game_type", "scopone_scientifico")),
target_score=int(data.get("target_score", DEFAULT_TARGET_SCORE)), target_score=int(data.get("target_score", DEFAULT_TARGET_SCORE)),
napola=bool(data.get("napola", True)),
phase=str(data.get("phase", PHASE_LOBBY)), phase=str(data.get("phase", PHASE_LOBBY)),
players=[PlayerState.from_json(p) for p in data.get("players", [])], players=[PlayerState.from_json(p) for p in data.get("players", [])],
table=[Card.from_json(c) for c in data.get("table", [])], table=[Card.from_json(c) for c in data.get("table", [])],
+17 -1
View File
@@ -16,7 +16,7 @@ from typing import Any, Dict, Optional
from kaya.core import HttpContext from kaya.core import HttpContext
from kaya.openapi import operation from kaya.openapi import operation
from .. import auth from .. import auth, deadlines
from ..app import app, game_store, oidc_mixin from ..app import app, game_store, oidc_mixin
from ..auth import require_auth from ..auth import require_auth
from ..config import settings from ..config import settings
@@ -52,6 +52,7 @@ def _lobby_payload(state: GameState) -> Dict[str, Any]:
"join_code": state.join_code, "join_code": state.join_code,
"game_type": state.game_type, "game_type": state.game_type,
"target_score": state.target_score, "target_score": state.target_score,
"napola": state.napola,
"phase": state.phase, "phase": state.phase,
"players": [ "players": [
{"sub": p.sub, "name": p.name, "seat": p.seat, "team": "A" if p.seat % 2 == 0 else "B"} {"sub": p.sub, "name": p.name, "seat": p.seat, "team": "A" if p.seat % 2 == 0 else "B"}
@@ -94,6 +95,12 @@ async def list_game_types(ctx: HttpContext) -> None:
"description": "One of the ids from GET /api/game-types", "description": "One of the ids from GET /api/game-types",
}, },
"target_score": {"type": "integer", "minimum": 1, "maximum": 100}, "target_score": {"type": "integer", "minimum": 1, "maximum": 100},
"napola": {
"type": "boolean",
"default": True,
"description": "Score the napola rule; a full "
"denari sweep wins the match",
},
}, },
} }
} }
@@ -123,6 +130,11 @@ async def create_game(ctx: HttpContext) -> None:
await send_error(ctx, 400, f"unknown game_type: {game_type!r}") await send_error(ctx, 400, f"unknown game_type: {game_type!r}")
return return
napola: Any = body.get("napola", True)
if not isinstance(napola, bool):
await send_error(ctx, 400, "napola must be a boolean")
return
user = oidc_mixin.get_user(ctx) user = oidc_mixin.get_user(ctx)
assert user is not None # enforced by @require_auth assert user is not None # enforced by @require_auth
game_id = str(uuid.uuid4()) game_id = str(uuid.uuid4())
@@ -137,6 +149,7 @@ async def create_game(ctx: HttpContext) -> None:
hand_ack_timeout=settings.hand_ack_timeout_seconds, hand_ack_timeout=settings.hand_ack_timeout_seconds,
turn_timeout=settings.turn_timeout_seconds, turn_timeout=settings.turn_timeout_seconds,
game_type=game_type, game_type=game_type,
napola=napola,
) )
except GameError as exc: except GameError as exc:
await send_error(ctx, 400, str(exc)) await send_error(ctx, 400, str(exc))
@@ -209,6 +222,9 @@ async def join_game(ctx: HttpContext) -> None:
return return
await game_store.save(state) await game_store.save(state)
await game_store.publish(state.id) await game_store.publish(state.id)
# When the fourth join started the match, the first turn deadline
# was armed; queue it so it fires even if nobody ever connects.
await deadlines.sync_deadline(game_store, state)
seat = next(p.seat for p in state.players if p.sub == user.sub) seat = next(p.seat for p in state.players if p.sub == user.sub)
if state.phase == PHASE_LOBBY: if state.phase == PHASE_LOBBY:
log.info("%s joined game %s (seat %d, %d/4 players)", user.sub, state.id, seat, len(state.players)) log.info("%s joined game %s (seat %d, %d/4 players)", user.sub, state.id, seat, len(state.players))
+64 -1
View File
@@ -18,6 +18,14 @@ channel as a simple "something changed" signal; every open websocket
reloads the state and renders the personalized view. Publishing only a reloads the state and renders the personalized view. Publishing only a
signal (never the state) means updated state reaches connections on every signal (never the state) means updated state reaches connections on every
worker without leaking hidden hands into the channel. worker without leaking hidden hands into the channel.
Timeouts (turn auto-play, hand-end auto-continue) are driven by a shared
delayed-deadline queue: producers enqueue an opaque ``member`` string with
a due timestamp, and a consumer on every worker polls for due entries.
Delivery is at-least-once — entries are removed only after they are
processed — so a worker dying mid-processing cannot lose a deadline;
consumers revalidate entries against the live state under the per-game
lock, which makes duplicate deliveries harmless.
""" """
from __future__ import annotations from __future__ import annotations
@@ -26,7 +34,7 @@ import contextlib
import json import json
from abc import ABC, abstractmethod from abc import ABC, abstractmethod
from logging import getLogger from logging import getLogger
from typing import AsyncContextManager, AsyncIterator, Dict, Optional, Set from typing import AsyncContextManager, AsyncIterator, Dict, List, Optional, Set, cast
from redis.asyncio import Redis from redis.asyncio import Redis
@@ -37,6 +45,7 @@ log = getLogger(__name__)
GAME_KEY_PREFIX = "tavolo:game:" GAME_KEY_PREFIX = "tavolo:game:"
CODE_KEY_PREFIX = "tavolo:code:" CODE_KEY_PREFIX = "tavolo:code:"
CHANNEL_PREFIX = "tavolo:game:" CHANNEL_PREFIX = "tavolo:game:"
DEADLINES_KEY = "tavolo:deadlines"
# Sentinel pushed into in-memory subscriber queues to signal a change. # Sentinel pushed into in-memory subscriber queues to signal a change.
_BUMP = b"update" _BUMP = b"update"
@@ -69,6 +78,26 @@ class GameStore(ABC):
async def publish(self, game_id: str) -> None: async def publish(self, game_id: str) -> None:
"""Signal that the state of ``game_id`` changed.""" """Signal that the state of ``game_id`` changed."""
@abstractmethod
async def add_deadline(self, member: str, due_at: float) -> None:
"""Enqueue ``member`` to fire at ``due_at`` (epoch seconds).
Idempotent for identical members: re-adding an existing member only
updates its due time.
"""
@abstractmethod
async def due_deadlines(self, now: float, limit: int = 32) -> List[str]:
"""Return up to ``limit`` enqueued members due at or before ``now``."""
@abstractmethod
async def next_deadline(self) -> Optional[float]:
"""Return the earliest pending due time (epoch seconds), if any."""
@abstractmethod
async def remove_deadline(self, member: str) -> None:
"""Remove ``member`` from the queue; a no-op when absent."""
def _channel(game_id: str) -> str: def _channel(game_id: str) -> str:
return f"{CHANNEL_PREFIX}{game_id}:events" return f"{CHANNEL_PREFIX}{game_id}:events"
@@ -128,6 +157,25 @@ class RedisGameStore(GameStore):
await self._redis.publish(_channel(game_id), "update") await self._redis.publish(_channel(game_id), "update")
log.debug("redis publish %s", game_id) log.debug("redis publish %s", game_id)
async def add_deadline(self, member: str, due_at: float) -> None:
await self._redis.zadd(DEADLINES_KEY, {member: due_at})
async def due_deadlines(self, now: float, limit: int = 32) -> List[str]:
members = cast(
list,
await self._redis.zrangebyscore(
DEADLINES_KEY, "-inf", now, start=0, num=limit
),
)
return [m.decode("utf-8") if isinstance(m, bytes) else m for m in members]
async def next_deadline(self) -> Optional[float]:
earliest = await self._redis.zrange(DEADLINES_KEY, 0, 0, withscores=True)
return float(earliest[0][1]) if earliest else None
async def remove_deadline(self, member: str) -> None:
await self._redis.zrem(DEADLINES_KEY, member)
async def _redis_events(pubsub) -> AsyncIterator[None]: async def _redis_events(pubsub) -> AsyncIterator[None]:
async for message in pubsub.listen(): async for message in pubsub.listen():
@@ -143,6 +191,7 @@ class InMemoryGameStore(GameStore):
self._codes: Dict[str, str] = {} self._codes: Dict[str, str] = {}
self._locks: Dict[str, asyncio.Lock] = {} self._locks: Dict[str, asyncio.Lock] = {}
self._subscribers: Dict[str, Set[asyncio.Queue]] = {} self._subscribers: Dict[str, Set[asyncio.Queue]] = {}
self._deadlines: Dict[str, float] = {}
def _lock_for(self, game_id: str) -> asyncio.Lock: def _lock_for(self, game_id: str) -> asyncio.Lock:
lock = self._locks.get(game_id) lock = self._locks.get(game_id)
@@ -187,6 +236,20 @@ class InMemoryGameStore(GameStore):
for queue in list(self._subscribers.get(game_id, ())): for queue in list(self._subscribers.get(game_id, ())):
queue.put_nowait(_BUMP) queue.put_nowait(_BUMP)
async def add_deadline(self, member: str, due_at: float) -> None:
self._deadlines[member] = due_at
async def due_deadlines(self, now: float, limit: int = 32) -> List[str]:
due = [m for m, due_at in self._deadlines.items() if due_at <= now]
due.sort(key=self._deadlines.__getitem__)
return due[:limit]
async def next_deadline(self) -> Optional[float]:
return min(self._deadlines.values(), default=None)
async def remove_deadline(self, member: str) -> None:
self._deadlines.pop(member, None)
async def _queue_events(queue: asyncio.Queue) -> AsyncIterator[None]: async def _queue_events(queue: asyncio.Queue) -> AsyncIterator[None]:
while True: while True:
+17 -133
View File
@@ -30,29 +30,28 @@ state is saved to Redis and a change signal is published. Every connected
websocket is subscribed to that signal and re-renders the state, so all websocket is subscribed to that signal and re-renders the state, so all
players see the move immediately (and consistently across workers). players see the move immediately (and consistently across workers).
If a player does not move before the per-game ``turn_timeout``, the server Timeouts do not depend on anyone being connected: both the per-turn
plays a random card (with a random legal capture when one is required) for auto-play and the hand-end auto-continue are driven by the absolute
them, so a disconnected or idle player cannot stall the match. The timer is deadlines persisted on the game state, via the shared deadline queue
re-armed by every client connection and state broadcast, and fires drained by a consumer on every worker (see :mod:`tavolo.deadlines`). A
immediately when a reconnect finds the deadline already past. disconnected or idle player therefore cannot stall the match, and a
worker dying cannot either.
""" """
from __future__ import annotations from __future__ import annotations
import asyncio import asyncio
import json import json
from contextlib import suppress from contextlib import suppress
from datetime import datetime, timezone
from logging import getLogger from logging import getLogger
from typing import Any, Awaitable, Callable, Dict, Optional from typing import Any, Awaitable, Callable, Dict
from kaya.core import WebSocket from kaya.core import WebSocket
from . import auth from . import auth, deadlines
from .app import app, game_store from .app import app, game_store
from .game import engine from .game import engine
from .game.errors import GameError from .game.errors import GameError
from .game.state import PHASE_FINISHED, PHASE_HAND_END, PHASE_PLAYING, GameState from .game.state import PHASE_FINISHED, GameState
from .stats import save_match_result
log = getLogger(__name__) log = getLogger(__name__)
@@ -95,7 +94,10 @@ async def game_socket(ws: WebSocket, game_id: str) -> None:
await ws.send_text(json.dumps(payload)) await ws.send_text(json.dumps(payload))
await send(_state_message(state, user.sub)) await send(_state_message(state, user.sub))
schedule_turn_timer(game_id, state) # Backstop: make sure the current phase's deadline is queued even if
# its entry was lost (e.g. the queue was flushed while the game lived
# on thanks to its sliding TTL).
await deadlines.sync_deadline(game_store, state)
async with game_store.subscribe(game_id) as events: async with game_store.subscribe(game_id) as events:
forward = asyncio.create_task( forward = asyncio.create_task(
@@ -126,7 +128,6 @@ async def _forward(
state = await game_store.load(game_id) state = await game_store.load(game_id)
if state is None: if state is None:
return return
schedule_turn_timer(game_id, state)
await send(_state_message(state, sub)) await send(_state_message(state, sub))
if state.phase == PHASE_FINISHED: if state.phase == PHASE_FINISHED:
await send( await send(
@@ -166,10 +167,6 @@ async def _handle_message(send: Send, game_id: str, sub: str, raw: str) -> None:
# --- hand-end acknowledgement ------------------------------------------------ # --- hand-end acknowledgement ------------------------------------------------
# Running auto-continue timers, keyed by (game_id, hand_number), so a hand's
# timeout is scheduled only once even when several clients are connected.
_hand_end_timers: Dict[tuple, asyncio.Task] = {}
async def _handle_ack(send: Send, game_id: str, sub: str) -> None: async def _handle_ack(send: Send, game_id: str, sub: str) -> None:
async with game_store.lock(game_id): async with game_store.lock(game_id):
@@ -185,122 +182,9 @@ async def _handle_ack(send: Send, game_id: str, sub: str) -> None:
log.debug("game %s: %s acknowledged hand %d", game_id, sub, state.hand_number) log.debug("game %s: %s acknowledged hand %d", game_id, sub, state.hand_number)
await game_store.save(state) await game_store.save(state)
await game_store.publish(game_id) await game_store.publish(game_id)
# The fourth ack deals the next hand, which arms a new turn
# deadline; earlier acks change nothing and this is a no-op.
def schedule_hand_end_timer(game_id: str, hand_number: int, timeout: int) -> None: await deadlines.sync_deadline(game_store, state)
"""Deal the next hand after the acknowledgement timeout, even if not
everyone has clicked. Fizzles if the hand already advanced."""
key = (game_id, hand_number)
if key in _hand_end_timers:
return
async def _auto_advance() -> None:
try:
await asyncio.sleep(timeout)
async with game_store.lock(game_id):
state = await game_store.load(game_id)
if (
state is None
or state.phase != engine.PHASE_HAND_END
or state.hand_number != hand_number
):
return
for player in state.players:
engine.acknowledge_hand(state, player.sub)
await game_store.save(state)
await game_store.publish(game_id)
log.info(
"game %s: hand %d auto-advanced after the acknowledgement timeout",
game_id,
hand_number,
)
finally:
_hand_end_timers.pop(key, None)
_hand_end_timers[key] = asyncio.create_task(_auto_advance())
# --- auto-play on turn timeout ------------------------------------------------
# Running turn timers, keyed by (game_id, hand_number, turn, deadline), so a
# turn's timeout is scheduled only once even when several clients are
# connected. Including the deadline means a re-arm after a reconnect cannot
# duplicate a timer for a turn that was already auto-played.
_turn_timers: Dict[tuple, asyncio.Task] = {}
def schedule_turn_timer(game_id: str, state: GameState) -> None:
"""Auto-play a random legal card if the player on turn misses the
deadline. Fizzles if the turn already advanced."""
if state.phase != PHASE_PLAYING or not state.turn_deadline:
return
key = (game_id, state.hand_number, state.turn, state.turn_deadline)
if key in _turn_timers:
return
hand_number = state.hand_number
turn = state.turn
deadline_raw = state.turn_deadline
try:
deadline = datetime.fromisoformat(deadline_raw)
except ValueError:
return
async def _auto_play() -> None:
try:
delay = (deadline - datetime.now(timezone.utc)).total_seconds()
await asyncio.sleep(max(delay, 0))
async with game_store.lock(game_id):
state = await game_store.load(game_id)
if (
state is None
or state.phase != PHASE_PLAYING
or state.hand_number != hand_number
or state.turn != turn
or state.turn_deadline != deadline_raw
):
# The turn moved on (or the game ended) without this
# timer firing: make sure the current turn is armed.
if state is not None:
schedule_turn_timer(game_id, state)
return
try:
engine.auto_play(state)
except GameError:
return
log.info(
"game %s: auto-played for %s (turn timeout, hand %d)",
game_id,
state.players[turn].sub if turn < len(state.players) else "?",
hand_number,
)
await _after_play(state, game_id)
finally:
_turn_timers.pop(key, None)
_turn_timers[key] = asyncio.create_task(_auto_play())
async def _after_play(state: GameState, game_id: str) -> None:
"""Persist a successful move and notify every connected player.
Callers must hold the per-game lock. Handles the two terminal
transitions: the match result is written to Postgres once, and a
hand-end summary schedules the auto-continue timeout.
"""
if state.phase == PHASE_FINISHED:
await save_match_result(state)
log.info(
"game %s finished: team %s wins %d-%d",
game_id,
"A" if state.winner == 0 else "B",
state.scores[0],
state.scores[1],
)
elif state.phase == PHASE_HAND_END:
schedule_hand_end_timer(game_id, state.hand_number, state.hand_ack_timeout)
await game_store.save(state)
await game_store.publish(game_id)
async def _handle_play( async def _handle_play(
@@ -335,4 +219,4 @@ async def _handle_play(
return return
log.debug("game %s: %s played %s (capture: %s)", game_id, sub, card, capture or "-") log.debug("game %s: %s played %s (capture: %s)", game_id, sub, card, capture or "-")
await _after_play(state, game_id) await deadlines.finalize_mutation(game_store, state)
+49
View File
@@ -81,5 +81,54 @@ class DatabaseUrlTests(unittest.TestCase):
) )
class CorsSettingsTests(unittest.TestCase):
def test_cors_disabled_by_default(self):
settings = _settings({})
self.assertIsNone(settings.cors_allow_origins)
self.assertIsNone(settings.cors_allow_origin_regex)
self.assertIsNone(settings.cors_allow_methods)
self.assertIsNone(settings.cors_allow_headers)
self.assertFalse(settings.cors_allow_credentials)
self.assertIsNone(settings.cors_expose_headers)
self.assertEqual(600, settings.cors_max_age)
def test_allow_origins_parses_comma_separated_list(self):
settings = _settings({
"CORS_ALLOW_ORIGINS": "https://a.example, https://b.example ,,https://c.example",
})
self.assertEqual(
("https://a.example", "https://b.example", "https://c.example"),
settings.cors_allow_origins,
)
def test_allow_origins_star_is_passed_through(self):
settings = _settings({"CORS_ALLOW_ORIGINS": "*"})
self.assertEqual(("*",), settings.cors_allow_origins)
def test_allow_origin_regex_is_passed_through(self):
settings = _settings({"CORS_ALLOW_ORIGIN_REGEX": r"https://.*\.example\.com"})
self.assertEqual(r"https://.*\.example\.com", settings.cors_allow_origin_regex)
def test_allow_methods_and_headers_parse_as_lists(self):
settings = _settings({
"CORS_ALLOW_METHODS": "GET,POST",
"CORS_ALLOW_HEADERS": "Authorization, X-Custom-Header",
"CORS_EXPOSE_HEADERS": "X-Total-Count",
})
self.assertEqual(("GET", "POST"), settings.cors_allow_methods)
self.assertEqual(("Authorization", "X-Custom-Header"), settings.cors_allow_headers)
self.assertEqual(("X-Total-Count",), settings.cors_expose_headers)
def test_allow_credentials_parses_boolean(self):
for value in ("1", "true", "TRUE", "yes", "on"):
self.assertTrue(_settings({"CORS_ALLOW_CREDENTIALS": value}).cors_allow_credentials)
for value in ("0", "false", "no", "off", "anything-else"):
self.assertFalse(_settings({"CORS_ALLOW_CREDENTIALS": value}).cors_allow_credentials)
def test_max_age_parses_int(self):
settings = _settings({"CORS_MAX_AGE": "3600"})
self.assertEqual(3600, settings.cors_max_age)
if __name__ == "__main__": if __name__ == "__main__":
unittest.main() unittest.main()
+129
View File
@@ -0,0 +1,129 @@
"""Integration tests for the CORS configuration in :mod:`tavolo.app`.
The mixin under test is kaya-cors' :class:`~kaya.cors.CorsMixin`; these
tests only verify that :func:`tavolo.app.cors_mixin_from_settings` maps the
environment-driven :class:`~tavolo.config.Settings` onto it correctly. A
minimal ``KayaApp`` is used instead of the global ``app`` so the tests do
not depend on the environment the suite was imported with.
"""
from __future__ import annotations
import os
import unittest
from unittest.mock import patch
from httpx import ASGITransport, AsyncClient
from kaya.core import HttpContext, KayaApp
from pwo import async_test
from tavolo.app import cors_mixin_from_settings
from tavolo.config import Settings
ORIGIN = "https://cards.example"
def _settings(env: dict) -> Settings:
with patch.dict(os.environ, env, clear=True):
return Settings.from_env()
def _app(settings: Settings) -> KayaApp:
mixin = cors_mixin_from_settings(settings)
assert mixin is not None
app = KayaApp(mixins=[mixin])
@app.GET("/api/health")
async def health(ctx: HttpContext) -> None:
await ctx.send_str(200, "ok")
return app
class CorsMixinFromSettingsTests(unittest.TestCase):
def test_disabled_when_unconfigured(self):
self.assertIsNone(cors_mixin_from_settings(_settings({})))
def test_enabled_by_allow_origins(self):
self.assertIsNotNone(cors_mixin_from_settings(
_settings({"CORS_ALLOW_ORIGINS": ORIGIN})))
def test_enabled_by_allow_origin_regex_alone(self):
self.assertIsNotNone(cors_mixin_from_settings(
_settings({"CORS_ALLOW_ORIGIN_REGEX": r"https://.*\.example\.com"})))
class CorsBehaviorTests(unittest.TestCase):
@async_test
async def test_request_without_origin_is_untouched(self) -> None:
app = _app(_settings({"CORS_ALLOW_ORIGINS": ORIGIN}))
async with AsyncClient(transport=ASGITransport(app=app),
base_url="http://127.0.0.1") as client:
response = await client.get("/api/health")
self.assertEqual(200, response.status_code)
self.assertNotIn("access-control-allow-origin", response.headers)
@async_test
async def test_simple_request_with_allowed_origin(self) -> None:
app = _app(_settings({"CORS_ALLOW_ORIGINS": ORIGIN}))
async with AsyncClient(transport=ASGITransport(app=app),
base_url="http://127.0.0.1") as client:
response = await client.get("/api/health", headers={"Origin": ORIGIN})
self.assertEqual(200, response.status_code)
self.assertEqual(ORIGIN, response.headers["access-control-allow-origin"])
@async_test
async def test_simple_request_with_disallowed_origin(self) -> None:
app = _app(_settings({"CORS_ALLOW_ORIGINS": ORIGIN}))
async with AsyncClient(transport=ASGITransport(app=app),
base_url="http://127.0.0.1") as client:
response = await client.get(
"/api/health", headers={"Origin": "https://mallory.example"})
self.assertEqual(200, response.status_code)
self.assertNotIn("access-control-allow-origin", response.headers)
@async_test
async def test_preflight_allowed(self) -> None:
app = _app(_settings({
"CORS_ALLOW_ORIGINS": ORIGIN,
"CORS_ALLOW_METHODS": "GET,POST",
"CORS_MAX_AGE": "3600",
}))
async with AsyncClient(transport=ASGITransport(app=app),
base_url="http://127.0.0.1") as client:
response = await client.options("/api/health", headers={
"Origin": ORIGIN,
"Access-Control-Request-Method": "POST",
})
self.assertEqual(200, response.status_code)
self.assertEqual(ORIGIN, response.headers["access-control-allow-origin"])
self.assertEqual("GET, POST", response.headers["access-control-allow-methods"])
self.assertEqual("3600", response.headers["access-control-max-age"])
@async_test
async def test_preflight_disallowed_origin(self) -> None:
app = _app(_settings({"CORS_ALLOW_ORIGINS": ORIGIN}))
async with AsyncClient(transport=ASGITransport(app=app),
base_url="http://127.0.0.1") as client:
response = await client.options("/api/health", headers={
"Origin": "https://mallory.example",
"Access-Control-Request-Method": "GET",
})
self.assertEqual(400, response.status_code)
self.assertIn("Disallowed CORS", response.text)
@async_test
async def test_credentials_echo_origin_and_set_flag(self) -> None:
app = _app(_settings({
"CORS_ALLOW_ORIGINS": "*",
"CORS_ALLOW_CREDENTIALS": "true",
}))
async with AsyncClient(transport=ASGITransport(app=app),
base_url="http://127.0.0.1") as client:
response = await client.get("/api/health", headers={"Origin": ORIGIN})
self.assertEqual(200, response.status_code)
self.assertEqual(ORIGIN, response.headers["access-control-allow-origin"])
self.assertEqual("true", response.headers["access-control-allow-credentials"])
if __name__ == "__main__":
unittest.main()
+191
View File
@@ -0,0 +1,191 @@
"""Deadline-queue timeout tests.
Timeouts must be driven by the persisted deadlines and the shared queue,
not by connected sockets: these tests seed games, queue their deadlines
and let the background consumer fire them without a single websocket.
"""
from __future__ import annotations
import asyncio
import unittest
from datetime import datetime, timedelta, timezone
from typing import Optional
from pwo import async_test
from tavolo import deadlines
from tavolo.app import game_store
from tavolo.game import engine
from tavolo.game.state import GameState, PlayerState
PLAYERS = ("alice", "bob", "carol", "dave")
def _ms(iso: str) -> int:
"""Epoch milliseconds for an ISO-8601 timestamp (the queue-entry form)."""
return int(datetime.fromisoformat(iso).timestamp() * 1000)
def _hand_end_state(game_id: str, deadline: str) -> GameState:
"""A game paused on the hand-end summary, waiting for acks."""
state = GameState(
id=game_id,
join_code="DLhend",
creator_sub="alice",
target_score=11,
phase="hand_end",
# Long turn timeout: the next hand's auto-play must not interfere
# with later tests sharing this store.
turn_timeout=3600,
)
state.players = [
PlayerState(sub=name, name=name.capitalize(), seat=i)
for i, name in enumerate(PLAYERS)
]
state.hand_end_deadline = deadline
return state
async def _wait_for(predicate, timeout: float = 5.0) -> Optional[GameState]:
"""Poll the store until ``predicate`` holds for the loaded state."""
deadline = asyncio.get_running_loop().time() + timeout
while asyncio.get_running_loop().time() < deadline:
state = await predicate()
if state is not None:
return state
await asyncio.sleep(0.05)
return None
class ConnectionIndependenceTest(unittest.TestCase):
@async_test
async def test_turn_timeout_fires_with_no_connections(self) -> None:
state = engine.create_game(
"dl-turn-1", "DLT001", "alice", "Alice",
target_score=11, turn_timeout=1,
)
for name in PLAYERS[1:]:
engine.join_game(state, name, name.capitalize())
assert state.turn_deadline is not None
await game_store.save(state)
await deadlines.sync_deadline(game_store, state)
# Nobody ever connects: the consumer must still auto-play for Bob
# (seat 1, first to act).
result = await _wait_for(
lambda: _turn_is(state.id, 2),
)
self.assertIsNotNone(result, "turn deadline never fired")
assert result is not None
self.assertEqual(1, result.last_move.seat if result.last_move else None)
# Defuse the follow-on turn deadlines so this game cannot keep
# auto-playing while later tests run.
result.turn_timeout = 3600
await game_store.save(result)
@async_test
async def test_hand_end_timeout_fires_with_no_connections(self) -> None:
deadline = (datetime.now(timezone.utc) + timedelta(seconds=1)).isoformat()
state = _hand_end_state("dl-handend-1", deadline)
await game_store.save(state)
await deadlines.sync_deadline(game_store, state)
# Nobody acks (nobody is even connected): the deadline must deal
# the next hand.
result = await _wait_for(
lambda: _phase_is("dl-handend-1", "playing"),
)
self.assertIsNotNone(result, "hand-end deadline never fired")
assert result is not None
self.assertEqual(2, result.hand_number)
self.assertEqual([], result.acked)
async def _turn_is(game_id: str, turn: int) -> Optional[GameState]:
state = await game_store.load(game_id)
return state if state is not None and state.turn == turn else None
async def _phase_is(game_id: str, phase: str) -> Optional[GameState]:
state = await game_store.load(game_id)
return state if state is not None and state.phase == phase else None
class ProcessDueTest(unittest.TestCase):
"""Direct ``process_due`` behaviour: revalidation and idempotency."""
@async_test
async def test_processing_twice_is_a_no_op(self) -> None:
# Simulates a worker dying after firing but before removing the
# entry: another worker re-delivers the same entry.
deadline = (datetime.now(timezone.utc) - timedelta(seconds=1)).isoformat()
state = _hand_end_state("dl-idem-1", deadline)
await game_store.save(state)
member = deadlines.encode({
"game_id": state.id,
"kind": deadlines.KIND_HAND_END,
"hand": state.hand_number,
"deadline": _ms(deadline),
})
await deadlines.process_due(game_store, member)
await deadlines.process_due(game_store, member)
result = await game_store.load(state.id)
assert result is not None
# Advanced exactly once: hand 2, not hand 3.
self.assertEqual("playing", result.phase)
self.assertEqual(2, result.hand_number)
@async_test
async def test_stale_entry_is_discarded(self) -> None:
# A turn entry enqueued before a play landed in time: the state's
# deadline has moved, so the entry must not fire.
state = engine.create_game(
"dl-stale-1", "DLS001", "alice", "Alice",
target_score=11, turn_timeout=3600,
)
for name in PLAYERS[1:]:
engine.join_game(state, name, name.capitalize())
await game_store.save(state)
member = deadlines.encode({
"game_id": state.id,
"kind": deadlines.KIND_TURN,
"hand": state.hand_number,
"turn": state.turn,
# Not the live deadline (epoch milliseconds).
"deadline": 946684800000,
})
await game_store.add_deadline(member, due_at=0.0)
await deadlines.process_due(game_store, member)
result = await game_store.load(state.id)
assert result is not None
self.assertEqual(state.turn, result.turn)
# The entry was removed after processing.
self.assertNotIn(member, await game_store.due_deadlines(float("inf")))
@async_test
async def test_entry_for_expired_game_is_dropped(self) -> None:
member = deadlines.encode({
"game_id": "dl-gone",
"kind": deadlines.KIND_TURN,
"hand": 1,
"turn": 0,
"deadline": 946684800000,
})
await game_store.add_deadline(member, due_at=0.0)
await deadlines.process_due(game_store, member)
self.assertNotIn(member, await game_store.due_deadlines(float("inf")))
@async_test
async def test_malformed_entry_is_dropped(self) -> None:
await game_store.add_deadline("not json", due_at=0.0)
await deadlines.process_due(game_store, "not json")
self.assertNotIn("not json", await game_store.due_deadlines(float("inf")))
if __name__ == "__main__":
unittest.main()
+89
View File
@@ -239,6 +239,95 @@ class ScoringTest(unittest.TestCase):
self.assertEqual([0, 0], points) self.assertEqual([0, 0], points)
class NapolaTest(unittest.TestCase):
def test_napola_score_runs(self) -> None:
self.assertEqual(0, engine.napola_score(
[card(c) for c in ["02D", "03D", "04D"]])) # no ace
self.assertEqual(0, engine.napola_score(
[card(c) for c in ["01D", "02D"]])) # too short
self.assertEqual(3, engine.napola_score(
[card(c) for c in ["03D", "01D", "02D"]])) # order-independent
self.assertEqual(4, engine.napola_score(
[card(c) for c in ["01D", "02D", "03D", "04D", "07C"]]))
self.assertEqual(3, engine.napola_score(
[card(c) for c in ["01D", "02D", "03D", "05D"]])) # broken run
self.assertEqual(10, engine.napola_score(
[card(f"{rank:02d}D") for rank in range(1, 11)]))
def test_hand_points_napola(self) -> None:
state = make_state(
[[], [], [], []],
table=[],
captured=[
["01D", "02D", "03D", "04C"], # seat 0, team A
["05D", "06D", "07D", "08D"], # seat 1, team B
["09D", "10D", "01C", "02C"], # seat 2, team A
["03C", "05C", "06C", "07C"], # seat 3, team B
],
)
points, details = engine.hand_points(state)
# Team A has the ace-led run 01D-03D (3 points); team B's denari
# start at the 5, so no napola. Carte tie (8 each), denara to A
# (5 vs 4), settebello to B, primiere tied at 0 (missing suits).
self.assertEqual({"A": 3, "B": 0}, details["napola"])
self.assertEqual("A", details["award"]["napola"])
self.assertEqual([4, 1], points)
def test_napola_disabled(self) -> None:
state = make_state(
[[], [], [], []],
table=[],
captured=[
["01D", "02D", "03D", "04C"],
["05D", "06D", "07D", "08D"],
["09D", "10D", "01C", "02C"],
["03C", "05C", "06C", "07C"],
],
)
state.napola = False
points, details = engine.hand_points(state)
self.assertNotIn("napola", details)
self.assertEqual([1, 1], points)
def test_full_denari_sweep_wins_match_instantly(self) -> None:
# Team A already captured the whole denari suit; the last play of
# the hand cannot capture. Team B leads 50-0, yet the napola ends
# the match in team A's favour, well below the target of 100.
state = make_state(
[["02C"], [], [], []],
table=[],
target=100,
captured=[
[f"{rank:02d}D" for rank in range(1, 11)],
[],
[],
[],
],
)
state.scores = [0, 50]
engine.play(state, "p0", "02C")
self.assertEqual(PHASE_FINISHED, state.phase)
self.assertEqual(0, state.winner)
self.assertLess(state.scores[0], 100)
self.assertEqual(10, state.hand_scores[-1]["napola"]["A"])
def test_napola_serialization_roundtrip(self) -> None:
state = make_state([["02D"], [], [], []], table=[])
self.assertTrue(state.napola)
state.napola = False
self.assertFalse(GameState.from_json(state.to_json()).napola)
# States serialized before the option existed default to enabled.
data = state.to_json()
del data["napola"]
self.assertTrue(GameState.from_json(data).napola)
def test_create_game_napola_default_and_override(self) -> None:
self.assertTrue(engine.create_game("g", "CODE42", "p0", "p0").napola)
self.assertFalse(
engine.create_game("g", "CODE42", "p0", "p0", napola=False).napola
)
class MatchFlowTest(unittest.TestCase): class MatchFlowTest(unittest.TestCase):
def test_join_starts_when_full(self) -> None: def test_join_starts_when_full(self) -> None:
state = engine.create_game("g", "CODE42", "p0", "p0", target_score=11) state = engine.create_game("g", "CODE42", "p0", "p0", target_score=11)
+18
View File
@@ -71,6 +71,24 @@ class GamesRouteTest(unittest.TestCase):
self.assertNotIn("hand", state["players"][0]) self.assertNotIn("hand", state["players"][0])
self.assertEqual(1, state["turn"]) self.assertEqual(1, state["turn"])
@async_test
async def test_create_napola_option(self) -> None:
transport = ASGITransport(app=app)
async with AsyncClient(transport=transport, base_url="http://127.0.0.1") as client:
with oidc_user("alice"):
default = await client.post("/api/games", json={})
self.assertEqual(201, default.status_code)
self.assertTrue(default.json()["napola"])
with oidc_user("alice"):
disabled = await client.post("/api/games", json={"napola": False})
self.assertEqual(201, disabled.status_code)
self.assertFalse(disabled.json()["napola"])
with oidc_user("alice"):
invalid = await client.post("/api/games", json={"napola": "yes"})
self.assertEqual(400, invalid.status_code)
@async_test @async_test
async def test_join_errors(self) -> None: async def test_join_errors(self) -> None:
transport = ASGITransport(app=app) transport = ASGITransport(app=app)
+23
View File
@@ -108,6 +108,29 @@ class InMemoryGameStoreTest(unittest.TestCase):
["holder-enter", "holder-exit", "contender"], order ["holder-enter", "holder-exit", "contender"], order
) )
@async_test
async def test_deadline_queue(self) -> None:
store = InMemoryGameStore()
self.assertIsNone(await store.next_deadline())
self.assertEqual([], await store.due_deadlines(now=100.0))
await store.add_deadline("b", due_at=50.0)
await store.add_deadline("a", due_at=10.0)
await store.add_deadline("c", due_at=200.0)
# Re-adding an existing member only updates its due time.
await store.add_deadline("b", due_at=60.0)
self.assertEqual(10.0, await store.next_deadline())
self.assertEqual(["a"], await store.due_deadlines(now=10.0))
self.assertEqual(["a", "b"], await store.due_deadlines(now=100.0))
# Due entries come out in due-time order and stay queued until removed.
self.assertEqual(["a", "b"], await store.due_deadlines(now=100.0))
await store.remove_deadline("a")
await store.remove_deadline("a") # removing twice is a no-op
self.assertEqual(60.0, await store.next_deadline())
self.assertEqual(["b"], await store.due_deadlines(now=100.0))
if __name__ == "__main__": if __name__ == "__main__":
unittest.main() unittest.main()
+2 -2
View File
@@ -35,9 +35,9 @@ pub async fn game_types() -> Result<Vec<GameTypeInfo>, String> {
Ok(page.results) Ok(page.results)
} }
pub async fn create_game(game_type: &str, target_score: i32) -> Result<GameView, String> { pub async fn create_game(game_type: &str, target_score: i32, napola: bool) -> Result<GameView, String> {
let resp = Request::post("/api/games") let resp = Request::post("/api/games")
.json(&serde_json::json!({ "game_type": game_type, "target_score": target_score })) .json(&serde_json::json!({ "game_type": game_type, "target_score": target_score, "napola": napola }))
.map_err(|e| e.to_string())? .map_err(|e| e.to_string())?
.send() .send()
.await .await
+1
View File
@@ -1,2 +1,3 @@
pub mod card; pub mod card;
pub mod summary; pub mod summary;
pub mod toast;
+39
View File
@@ -113,6 +113,45 @@ pub fn summary_rows(summary: HandSummary) -> View {
summary.award.primiera.clone(), summary.award.primiera.clone(),
"+1", "+1",
), ),
// Napola (only when the rule is enabled for this match)
match summary.napola {
None => view! {},
Some(napola) => {
let n = match summary.award.napola.as_deref() {
Some("A") => napola.a,
Some("B") => napola.b,
_ => 0,
};
let text = match (&summary.award.napola, n) {
(Some(t), 10) => format!(
"Team {t} swept the whole denari suit — napola! Instant match win"
),
(Some(t), n) => format!(
"Team {t} captured {n} consecutive denari from the ace"
),
(None, _) => "No napola this hand".to_string(),
};
let chip = match &summary.award.napola {
Some(t) => format!("Team {t} +{n}"),
None => "tie".to_string(),
};
let cls = match summary.award.napola.as_deref() {
Some("A") => "score-row team-a",
Some("B") => "score-row team-b",
_ => "score-row tie",
};
view! {
div(class=cls) {
div(class="score-icon") { (card_img("01D".to_string(), "score-mini")) }
div(class="score-body") {
div(class="score-title") { "Napola" }
div(class="score-text") { (text) }
}
div(class="score-points") { (chip) }
}
}
}
},
// Scope // Scope
{ {
let a = summary.scope.a; let a = summary.scope.a;
+20
View File
@@ -0,0 +1,20 @@
//! Auto-dismissing error toast.
use gloo_timers::callback::Timeout;
use sycamore::prelude::*;
const TOAST_MS: u32 = 10_000;
/// Renders the error from `error` as a toast; hides it after `TOAST_MS`.
/// A new error replaces the message and restarts the timer.
pub fn toast(error: Signal<Option<String>>) -> View {
create_effect(move || {
if error.get_clone().is_some() {
// Held until cleanup; dropped (cancelled) when the effect re-runs.
let timeout = Timeout::new(TOAST_MS, move || error.set(None));
on_cleanup(move || drop(timeout));
}
});
view! {
(move || error.get_clone().map(|e| view! { div(class="toast") { (e) } }))
}
}
+13
View File
@@ -2,6 +2,10 @@
use serde::Deserialize; use serde::Deserialize;
use std::collections::HashMap; use std::collections::HashMap;
fn default_true() -> bool {
true
}
#[derive(Debug, Clone, Deserialize)] #[derive(Debug, Clone, Deserialize)]
#[allow(dead_code)] #[allow(dead_code)]
pub struct User { pub struct User {
@@ -75,6 +79,8 @@ pub struct Award {
pub settebello: Option<String>, pub settebello: Option<String>,
#[serde(default)] #[serde(default)]
pub primiera: Option<String>, pub primiera: Option<String>,
#[serde(default)]
pub napola: Option<String>,
} }
/// The scoring breakdown of one completed hand. /// The scoring breakdown of one completed hand.
@@ -86,6 +92,10 @@ pub struct HandSummary {
pub settebello: TeamBools, pub settebello: TeamBools,
pub primiera: TeamCounts, pub primiera: TeamCounts,
pub scope: TeamCounts, pub scope: TeamCounts,
/// Napola run lengths per team; absent when the rule is disabled (or
/// the summary predates the option).
#[serde(default)]
pub napola: Option<TeamCounts>,
pub award: Award, pub award: Award,
#[serde(default)] #[serde(default)]
pub hand: i32, pub hand: i32,
@@ -104,6 +114,9 @@ pub struct GameView {
/// Which card game this match is (id from /api/game-types). /// Which card game this match is (id from /api/game-types).
#[serde(default)] #[serde(default)]
pub game_type: String, pub game_type: String,
/// Whether the napola rule is scored in this match.
#[serde(default = "default_true")]
pub napola: bool,
pub phase: String, pub phase: String,
#[serde(default)] #[serde(default)]
pub target_score: i32, pub target_score: i32,
+224 -30
View File
@@ -1,8 +1,12 @@
//! Live game page: table view over the websocket. //! Live game page: table view over the websocket.
use std::cell::Cell;
use std::rc::Rc;
use sycamore::prelude::*; use sycamore::prelude::*;
use crate::components::card::{card_back, card_img}; use crate::components::card::{card_back, card_img};
use crate::components::summary::{hand_summary_modal, summary_rows}; use crate::components::summary::{hand_summary_modal, summary_rows};
use crate::components::toast::toast;
use crate::model::{card_label, GameView, MoveView, PlayerView, Scores, ServerMessage}; use crate::model::{card_label, GameView, MoveView, PlayerView, Scores, ServerMessage};
use crate::ws::{self, GameSocket}; use crate::ws::{self, GameSocket};
@@ -69,6 +73,143 @@ fn move_banner(mv: MoveView) -> View {
} }
} }
/// Signals shared by the websocket connection and its reconnect attempts.
#[derive(Clone, Copy)]
struct ConnCtx {
socket: Signal<Option<GameSocket>>,
game: Signal<Option<GameView>>,
over: Signal<Option<(Scores, Option<String>)>>,
error: Signal<Option<String>>,
closed: Signal<bool>,
/// Reconnect attempts exhausted; only a manual retry resumes.
gave_up: Signal<bool>,
/// The server closed the connection deliberately (auth or game gone);
/// retrying is pointless.
fatal: Signal<bool>,
attempts: Signal<u32>,
capture_choice: Signal<Option<(String, Vec<Vec<String>>)>>,
selected: Signal<Option<String>>,
}
/// Reconnect attempts: 1s, 2s, 4s, … capped at 30s, at most this many.
const MAX_RECONNECT_ATTEMPTS: u32 = 10;
fn backoff_ms(attempt: u32) -> u32 {
(1000u32 << attempt.min(5)).min(30_000)
}
/// Connect the game websocket, wiring state updates and reconnects.
///
/// The server pushes a full state snapshot on connect, so a reconnect is
/// also a resync: no client-side state merging is needed.
fn start_connect(id: Rc<String>, ctx: ConnCtx, alive: Rc<Cell<bool>>) {
let on_message = {
let alive = alive.clone();
move |msg: ServerMessage| {
if !alive.get() {
// The page is unmounted; its signals are disposed.
return;
}
match msg {
ServerMessage::State { game: g } => {
ctx.capture_choice.set(None);
ctx.selected.set(None);
// A received state proves the (re)connection works.
ctx.attempts.set(0);
ctx.gave_up.set(false);
ctx.closed.set(false);
ctx.game.set(Some(g));
}
ServerMessage::GameOver { scores, winner } => {
ctx.over.set(Some((scores, winner)))
}
ServerMessage::Error { message, .. } => ctx.error.set(Some(message)),
}
}
};
let on_close = {
let id = id.clone();
let alive = alive.clone();
move |code: Option<u16>| {
if !alive.get() {
return;
}
ctx.closed.set(true);
match code {
Some(4401) => {
ctx.fatal.set(true);
ctx.error
.set(Some("Session expired — please log in again.".to_string()));
}
Some(4403) | Some(4404) => {
ctx.fatal.set(true);
ctx.error
.set(Some("This game is no longer available.".to_string()));
}
_ => schedule_retry(id.clone(), ctx, alive.clone()),
}
}
};
match ws::connect(&id, on_message, on_close) {
Some(s) => ctx.socket.set(Some(s)),
// WebSocket::open failed synchronously: treat as a transient loss.
None if alive.get() => {
ctx.closed.set(true);
schedule_retry(id, ctx, alive);
}
None => {}
}
}
/// Retry `start_connect` with exponential backoff, unless we gave up.
fn schedule_retry(id: Rc<String>, ctx: ConnCtx, alive: Rc<Cell<bool>>) {
let attempt = ctx.attempts.get();
if attempt >= MAX_RECONNECT_ATTEMPTS {
ctx.gave_up.set(true);
return;
}
ctx.attempts.set(attempt + 1);
gloo_timers::callback::Timeout::new(backoff_ms(attempt), move || {
if alive.get() {
start_connect(id, ctx, alive);
}
})
.forget();
}
/// Slim banner shown over the table while the socket is down.
fn conn_banner(
closed: bool,
gave_up: bool,
fatal: bool,
has_game: bool,
reconnect: Rc<dyn Fn()>,
) -> View {
if !closed || !has_game {
return view! {};
}
if fatal {
view! {
div(class="conn-banner") {
"Connection closed by the server. "
a(href="/") { "Back to lobby" }
}
}
} else if gave_up {
view! {
div(class="conn-banner") {
"Connection lost."
button(class="button", on:click=move |_| reconnect()) { "Retry now" }
a(href="/") { "Back to lobby" }
}
}
} else {
view! {
div(class="conn-banner") { "Connection lost — reconnecting…" }
}
}
}
#[component(inline_props)] #[component(inline_props)]
pub fn GamePage(id: String) -> View { pub fn GamePage(id: String) -> View {
let game = create_signal(Option::<GameView>::None); let game = create_signal(Option::<GameView>::None);
@@ -77,32 +218,51 @@ pub fn GamePage(id: String) -> View {
let selected = create_signal(Option::<String>::None); let selected = create_signal(Option::<String>::None);
let over = create_signal(Option::<(Scores, Option<String>)>::None); let over = create_signal(Option::<(Scores, Option<String>)>::None);
let closed = create_signal(false); let closed = create_signal(false);
let gave_up = create_signal(false);
let fatal = create_signal(false);
let attempts = create_signal(0u32);
let socket = create_signal(Option::<GameSocket>::None); let socket = create_signal(Option::<GameSocket>::None);
// Ticking clock driving the hand-end countdown display. // Ticking clock driving the hand-end countdown display.
let now = create_signal(js_sys::Date::now()); let now = create_signal(js_sys::Date::now());
gloo_timers::callback::Interval::new(500, move || now.set(js_sys::Date::now())).forget(); let ticker = gloo_timers::callback::Interval::new(500, move || now.set(js_sys::Date::now()));
{ // Stops the ticker and any pending reconnect once the page unmounts.
let on_message = move |msg: ServerMessage| match msg { let alive = Rc::new(Cell::new(true));
ServerMessage::State { game: g } => { on_cleanup({
capture_choice.set(None); let alive = alive.clone();
selected.set(None); move || {
game.set(Some(g)); alive.set(false);
} drop(ticker);
ServerMessage::GameOver { scores, winner } => {
over.set(Some((scores, winner)));
}
ServerMessage::Error { message, .. } => error.set(Some(message)),
};
let on_close = move || closed.set(true);
match ws::connect(&id, on_message, on_close) {
Some(s) => socket.set(Some(s)),
None => error.set(Some("Could not connect to the game".to_string())),
} }
} });
let id = Rc::new(id);
let ctx = ConnCtx {
socket,
game,
over,
error,
closed,
gave_up,
fatal,
attempts,
capture_choice,
selected,
};
start_connect(id.clone(), ctx, alive.clone());
let reconnect: Rc<dyn Fn()> = Rc::new(move || {
ctx.attempts.set(0);
ctx.gave_up.set(false);
ctx.closed.set(false);
start_connect(id.clone(), ctx, alive.clone());
});
// Clicking a card in the player's own hand. // Clicking a card in the player's own hand.
let on_hand_card = move |code: String| { let on_hand_card = move |code: String| {
if closed.get() {
// A dead socket would swallow the play silently.
return;
}
let Some(g) = game.get_clone() else { return }; let Some(g) = game.get_clone() else { return };
if g.your_turn != Some(true) { if g.your_turn != Some(true) {
return; return;
@@ -121,20 +281,50 @@ pub fn GamePage(id: String) -> View {
} }
}; };
let reconnect_banner = reconnect.clone();
view! { view! {
div(class="game-page") { div(class="game-page") {
(move || error.get_clone().map(|e| view! { div(class="toast") { (e) } })) (toast(error))
(move || conn_banner(
closed.get(),
gave_up.get(),
fatal.get(),
game.get_clone().is_some(),
reconnect_banner.clone(),
))
(move || match game.get_clone() { (move || match game.get_clone() {
None => { None => {
let status = if closed.get() { if fatal.get() {
"Connection closed." view! {
div(class="panel status-panel") {
p { "Connection closed." }
p { a(href="/") { "Back to lobby" } }
}
}
} else if gave_up.get() {
let reconnect = reconnect.clone();
view! {
div(class="panel status-panel") {
p { "Connection lost." }
p {
button(class="button primary", on:click=move |_| reconnect()) {
"Retry now"
}
}
p { a(href="/") { "Back to lobby" } }
}
}
} else { } else {
"Connecting to the game…" let status = if closed.get() {
}; "Connection lost — reconnecting…"
view! { } else {
div(class="panel status-panel") { "Connecting to the game…"
p { (status) } };
p { a(href="/") { "Back to lobby" } } view! {
div(class="panel status-panel") {
p { (status) }
p { a(href="/") { "Back to lobby" } }
}
} }
} }
} }
@@ -246,9 +436,10 @@ fn table_view(
let hand_number = game.hand_number; let hand_number = game.hand_number;
let target_score = game.target_score; let target_score = game.target_score;
let hand = player_for_seat(&game, viewer_seat) let viewer = player_for_seat(&game, viewer_seat);
.and_then(|p| p.hand) let my_captured = viewer.as_ref().map(|p| p.captured_count).unwrap_or(0);
.unwrap_or_default(); let my_scope = viewer.as_ref().map(|p| p.scope).unwrap_or(0);
let hand = viewer.and_then(|p| p.hand).unwrap_or_default();
let current_selection = selected.get_clone(); let current_selection = selected.get_clone();
let hand_cards = hand let hand_cards = hand
.into_iter() .into_iter()
@@ -300,6 +491,9 @@ fn table_view(
} }
(right) (right)
div(class="seat-bottom") { div(class="seat-bottom") {
div(class="seat-stats own-stats") {
(my_captured) " captured · " (my_scope) " scope"
}
div(class="hand") { (hand_cards) } div(class="hand") { (hand_cards) }
(hint) (hint)
} }
+2 -1
View File
@@ -3,6 +3,7 @@ use wasm_bindgen_futures::spawn_local;
use sycamore::prelude::*; use sycamore::prelude::*;
use crate::api; use crate::api;
use crate::components::toast::toast;
use crate::model::MatchesPage; use crate::model::MatchesPage;
#[component] #[component]
@@ -45,7 +46,7 @@ pub fn HistoryPage() -> View {
a(href="/leaderboard") { "Leaderboard" } a(href="/leaderboard") { "Leaderboard" }
} }
h1 { "My matches" } h1 { "My matches" }
(move || error.get_clone().map(|e| view! { div(class="toast") { (e) } })) (toast(error))
(move || match page.get_clone() { (move || match page.get_clone() {
None => view! { p(class="status") { "Loading…" } }, None => view! { p(class="status") { "Loading…" } },
Some(_) if rows.get_clone().is_empty() => view! { Some(_) if rows.get_clone().is_empty() => view! {
+2 -1
View File
@@ -3,6 +3,7 @@ use wasm_bindgen_futures::spawn_local;
use sycamore::prelude::*; use sycamore::prelude::*;
use crate::api; use crate::api;
use crate::components::toast::toast;
use crate::model::LeaderboardPage; use crate::model::LeaderboardPage;
#[component] #[component]
@@ -24,7 +25,7 @@ pub fn LeaderboardPage() -> View {
a(href="/history") { "My matches" } a(href="/history") { "My matches" }
} }
h1 { "Leaderboard" } h1 { "Leaderboard" }
(move || error.get_clone().map(|e| view! { div(class="toast") { (e) } })) (toast(error))
(move || match page.get_clone() { (move || match page.get_clone() {
None => view! { p(class="status") { "Loading…" } }, None => view! { p(class="status") { "Loading…" } },
Some(p) => { Some(p) => {
+14 -2
View File
@@ -4,6 +4,7 @@ use sycamore::prelude::*;
use sycamore_router::navigate; use sycamore_router::navigate;
use crate::api; use crate::api;
use crate::components::toast::toast;
use crate::model::{GameTypeInfo, User}; use crate::model::{GameTypeInfo, User};
/// Used when the game-types fetch fails: match creation must still work. /// Used when the game-types fetch fails: match creation must still work.
@@ -23,6 +24,7 @@ pub fn LobbyPage() -> View {
let code = create_signal(String::new()); let code = create_signal(String::new());
let game_types = create_signal(fallback_game_types()); let game_types = create_signal(fallback_game_types());
let selected_game = create_signal("scopone_scientifico".to_string()); let selected_game = create_signal("scopone_scientifico".to_string());
let napola = create_signal(true);
spawn_local(async move { spawn_local(async move {
match api::me().await { match api::me().await {
@@ -46,8 +48,9 @@ pub fn LobbyPage() -> View {
let on_create = move |target: i32| { let on_create = move |target: i32| {
let game_type = selected_game.get_clone(); let game_type = selected_game.get_clone();
let napola = napola.get();
spawn_local(async move { spawn_local(async move {
match api::create_game(&game_type, target).await { match api::create_game(&game_type, target, napola).await {
Ok(game) => navigate(&format!("/game/{}", game.id)), Ok(game) => navigate(&format!("/game/{}", game.id)),
Err(e) => error.set(Some(e)), Err(e) => error.set(Some(e)),
} }
@@ -70,7 +73,7 @@ pub fn LobbyPage() -> View {
view! { view! {
div(class="lobby") { div(class="lobby") {
h1 { "Scopone scientifico" } h1 { "Scopone scientifico" }
(move || error.get_clone().map(|e| view! { div(class="toast") { (e) } })) (toast(error))
(move || match user.get_clone() { (move || match user.get_clone() {
None => view! { p(class="status") { "Loading…" } }, None => view! { p(class="status") { "Loading…" } },
Some(None) => view! { Some(None) => view! {
@@ -100,6 +103,15 @@ pub fn LobbyPage() -> View {
key=|g| g.id.clone(), key=|g| g.id.clone(),
) )
} }
label(class="check", r#for="napola") {
input(id="napola", r#type="checkbox", bind:checked=napola)
" Napola"
}
p(class="hint") {
"A-2-3 of denari scores 3 points, plus 1 per extra "
"consecutive denari card; sweeping the whole suit "
"wins the match instantly."
}
p { "First team to reach the target score wins." } p { "First team to reach the target score wins." }
div(class="target-buttons") { div(class="target-buttons") {
button(class="button", on:click=move |_| on_create(11)) { "Target 11" } button(class="button", on:click=move |_| on_create(11)) { "Target 11" }
+12 -3
View File
@@ -4,7 +4,7 @@ use std::rc::Rc;
use futures::channel::mpsc; use futures::channel::mpsc;
use futures::{SinkExt, StreamExt}; use futures::{SinkExt, StreamExt};
use gloo_net::websocket::{futures::WebSocket, Message}; use gloo_net::websocket::{futures::WebSocket, Message, WebSocketError};
use wasm_bindgen_futures::spawn_local; use wasm_bindgen_futures::spawn_local;
use crate::model::ServerMessage; use crate::model::ServerMessage;
@@ -53,10 +53,14 @@ impl GameSocket {
/// Open the websocket for `game_id` and forward parsed server messages to /// Open the websocket for `game_id` and forward parsed server messages to
/// `on_message`. Returns the socket handle, or `None` if the connection /// `on_message`. Returns the socket handle, or `None` if the connection
/// could not be created. /// could not be created.
///
/// `on_close` fires exactly once when the connection ends; it receives the
/// server close code when one was sent (e.g. 4401 unauthenticated, 4403 not
/// seated, 4404 unknown game) or `None` for an abnormal network loss.
pub fn connect( pub fn connect(
game_id: &str, game_id: &str,
on_message: impl Fn(ServerMessage) + 'static, on_message: impl Fn(ServerMessage) + 'static,
on_close: impl Fn() + 'static, on_close: impl Fn(Option<u16>) + 'static,
) -> Option<GameSocket> { ) -> Option<GameSocket> {
let ws = WebSocket::open(&ws_url(game_id)).ok()?; let ws = WebSocket::open(&ws_url(game_id)).ok()?;
let (mut write, mut read) = ws.split(); let (mut write, mut read) = ws.split();
@@ -72,6 +76,7 @@ pub fn connect(
}); });
spawn_local(async move { spawn_local(async move {
let mut close_code = None;
while let Some(msg) = read.next().await { while let Some(msg) = read.next().await {
match msg { match msg {
Ok(Message::Text(text)) => { Ok(Message::Text(text)) => {
@@ -80,10 +85,14 @@ pub fn connect(
} }
} }
Ok(Message::Bytes(_)) => {} Ok(Message::Bytes(_)) => {}
Err(WebSocketError::ConnectionClose(e)) => {
close_code = Some(e.code);
break;
}
Err(_) => break, Err(_) => break,
} }
} }
on_close(); on_close(close_code);
}); });
Some(GameSocket { Some(GameSocket {
+32
View File
@@ -123,6 +123,12 @@ body {
gap: 0.5rem; gap: 0.5rem;
} }
.check {
display: flex;
align-items: center;
gap: 0.4rem;
}
.join-form { .join-form {
display: flex; display: flex;
gap: 0.5rem; gap: 0.5rem;
@@ -304,6 +310,11 @@ table.matches td.lost {
margin-top: auto; margin-top: auto;
} }
/* The viewer's own stats sit above the hand; no flex context here. */
.seat-bottom .seat-stats {
margin-bottom: 0.35rem;
}
.center { .center {
grid-area: center; grid-area: center;
display: flex; display: flex;
@@ -407,6 +418,27 @@ table.matches td.lost {
margin-left: 0.25rem; margin-left: 0.25rem;
} }
/* ---------- connection banner ---------- */
.conn-banner {
display: flex;
align-items: center;
justify-content: center;
gap: 0.75rem;
background: rgba(232, 197, 71, 0.15);
border: 1px solid var(--accent);
border-radius: 8px;
color: var(--accent);
padding: 0.4rem 1rem;
margin: 0.5rem auto 0;
width: fit-content;
}
.conn-banner a {
color: var(--accent);
text-decoration: underline;
}
/* ---------- overlays ---------- */ /* ---------- overlays ---------- */
.overlay { .overlay {