Compare commits
6
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
57c6ffb032
|
||
|
|
c10d45523a
|
||
|
|
658e3d6ec7
|
||
|
|
5e4e1310b4
|
||
|
|
db7ba30d13
|
||
|
|
f7fa8e78e5
|
@@ -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.
|
||||||
|
|||||||
+18
-3
@@ -60,6 +60,7 @@ 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 |
|
||||||
| `APP_HOST` / `APP_PORT` | `0.0.0.0` / `8080` | Bind address |
|
| `APP_HOST` / `APP_PORT` | `0.0.0.0` / `8080` | Bind address |
|
||||||
|
|
||||||
@@ -111,6 +112,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 +140,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 +188,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 +212,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 +266,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
|
||||||
|
|||||||
@@ -28,6 +28,7 @@ from kaya.session.redis import RedisSessionStore
|
|||||||
from redis.asyncio import Redis
|
from redis.asyncio import Redis
|
||||||
|
|
||||||
from .config import settings
|
from .config import 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
|
||||||
@@ -76,7 +77,8 @@ 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])
|
app = KayaApp(mixins=[session_mixin, oidc_mixin, tortoise_mixin, openapi_mixin,
|
||||||
|
DeadlineSchedulerMixin(game_store)])
|
||||||
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,
|
||||||
|
|||||||
@@ -76,6 +76,11 @@ 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]
|
||||||
@@ -113,6 +118,7 @@ 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"),
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|||||||
@@ -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)
|
||||||
@@ -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,
|
||||||
|
|||||||
@@ -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", [])],
|
||||||
|
|||||||
@@ -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))
|
||||||
|
|||||||
@@ -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
@@ -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)
|
||||||
|
|||||||
@@ -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()
|
||||||
@@ -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)
|
||||||
|
|||||||
@@ -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)
|
||||||
|
|||||||
@@ -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
@@ -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,2 +1,3 @@
|
|||||||
pub mod card;
|
pub mod card;
|
||||||
pub mod summary;
|
pub mod summary;
|
||||||
|
pub mod toast;
|
||||||
|
|||||||
@@ -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;
|
||||||
|
|||||||
@@ -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) } }))
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -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
@@ -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)
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -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! {
|
||||||
|
|||||||
@@ -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
@@ -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
@@ -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 {
|
||||||
|
|||||||
@@ -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 {
|
||||||
|
|||||||
Reference in New Issue
Block a user