160 lines
6.7 KiB
Python
160 lines
6.7 KiB
Python
"""Общее для сокетов занятия: вход, доставка очереди хаба и выполнение команды.
|
||
|
||
Лестница входа одна для всех каналов, чтобы коды закрытия и тексты отказов
|
||
не расходились: 1008 — чужой origin, 1012 — занятие ушло другому узлу ещё
|
||
до accept, `ErrorEvent` и close — нет прав, нет занятия, не назначено.
|
||
|
||
Команда выполняется одной операцией хаба: сверка `(node_id, epoch)` остаётся
|
||
в её commit перед каждой мутацией, здесь только место вызова.
|
||
"""
|
||
|
||
import asyncio
|
||
import logging
|
||
from collections.abc import Awaitable, Callable, Collection
|
||
from uuid import UUID
|
||
|
||
from fastapi import WebSocket
|
||
from pydantic import BaseModel
|
||
|
||
from app.api.auth import Principal, principal_of, websocket_origin_allowed
|
||
from app.domain.events import CommandAck, ErrorEvent, ErrorKind
|
||
from app.domain.roles import Role
|
||
from app.session.access import can_access
|
||
from app.session.hub import LEASE_FENCED_MESSAGE, hub
|
||
from app.session.state import MAX_STATION_COMMANDS, SessionState
|
||
|
||
log = logging.getLogger(__name__)
|
||
|
||
NOT_STARTED = "Занятие не запущено"
|
||
FORBIDDEN = "Недостаточно прав для этого экрана"
|
||
NOT_ASSIGNED = "Занятие не назначено этому обучающемуся"
|
||
|
||
#: Сколько ждать, пока `pump` доставит событие fencing сам. Очередь могла
|
||
#: переполниться и выпасть из подписки — тогда сокет закрывается здесь.
|
||
_FENCE_DELIVERY_SECONDS = 2.0
|
||
|
||
|
||
def _fenced_event() -> ErrorEvent:
|
||
return ErrorEvent(code=ErrorKind.INTERNAL, message=LEASE_FENCED_MESSAGE)
|
||
|
||
|
||
def _is_fenced_event(item) -> bool:
|
||
return (isinstance(item, ErrorEvent) and item.code is ErrorKind.INTERNAL
|
||
and item.message == LEASE_FENCED_MESSAGE)
|
||
|
||
|
||
async def _refuse(ws: WebSocket, code: ErrorKind, message: str) -> None:
|
||
"""Отказ до входа в цикл: сокет закрывается с объяснением, а не молча."""
|
||
await ws.send_text(ErrorEvent(code=code, message=message).model_dump_json())
|
||
await ws.close()
|
||
|
||
|
||
async def session_socket(
|
||
ws: WebSocket,
|
||
session_id: UUID,
|
||
roles: Collection[Role],
|
||
*,
|
||
require_state: bool = True,
|
||
not_found: str = NOT_STARTED,
|
||
) -> tuple[Principal, SessionState | None] | None:
|
||
"""Принципал и живое занятие — или `None`, если сокет уже закрыт.
|
||
|
||
`require_state=False` — для пульта: занятие создаёт его же `scenario.start`,
|
||
поэтому вход занятие не смотрит, а доступ проверяется на каждой команде.
|
||
"""
|
||
if not websocket_origin_allowed(ws):
|
||
await ws.close(code=1008)
|
||
return None
|
||
if hub.is_lease_fenced(session_id):
|
||
await ws.close(code=1012)
|
||
return None
|
||
await ws.accept()
|
||
|
||
who = principal_of(ws)
|
||
if who is None or who.role not in roles:
|
||
await _refuse(ws, ErrorKind.FORBIDDEN, FORBIDDEN)
|
||
return None
|
||
if not require_state:
|
||
return who, None
|
||
state = hub.get(session_id)
|
||
if state is None:
|
||
await _refuse(ws, ErrorKind.SESSION_NOT_FOUND, not_found)
|
||
return None
|
||
if not can_access(who, state):
|
||
# Чужому преподавателю занятие «не существует»: отказ не выдаёт владельца.
|
||
if who.role is Role.TRAINEE:
|
||
await _refuse(ws, ErrorKind.FORBIDDEN, NOT_ASSIGNED)
|
||
else:
|
||
await _refuse(ws, ErrorKind.SESSION_NOT_FOUND, not_found)
|
||
return None
|
||
return who, state
|
||
|
||
|
||
async def pump(ws: WebSocket, queue: asyncio.Queue) -> None:
|
||
"""Очередь хаба в сокет; событие fencing закрывает его с 1012."""
|
||
while True:
|
||
item = await queue.get()
|
||
# Бинарь — звук звонящего, без обёртки JSON (docs/arch/CONTRACT.md).
|
||
if isinstance(item, bytes):
|
||
await ws.send_bytes(item)
|
||
continue
|
||
await ws.send_text(item.model_dump_json())
|
||
if _is_fenced_event(item):
|
||
await ws.close(code=1012)
|
||
return
|
||
|
||
|
||
async def close_fenced(ws: WebSocket, sender: asyncio.Task | None = None) -> None:
|
||
"""Структурное событие fencing и 1012 вместо безымянного 1006.
|
||
|
||
С очередью событие уже разослал `hub.fence` — его доставляет `pump`.
|
||
Без очереди (пульт) сообщение уходит прямо в сокет.
|
||
"""
|
||
if sender is not None:
|
||
done, _ = await asyncio.wait({sender}, timeout=_FENCE_DELIVERY_SECONDS)
|
||
if done:
|
||
return
|
||
sender.cancel()
|
||
await ws.send_text(_fenced_event().model_dump_json())
|
||
await ws.close(code=1012)
|
||
|
||
|
||
async def run_command(
|
||
ws: WebSocket,
|
||
session_id: UUID,
|
||
handler: Callable[[], Awaitable[None]],
|
||
*,
|
||
command_id: UUID | None = None,
|
||
ack: Callable[[BaseModel], None] | None = None,
|
||
sender: asyncio.Task | None = None,
|
||
) -> bool:
|
||
"""Одна команда — одна операция хаба. `False` — сокет закрыт fencing.
|
||
|
||
С `command_id` повтор уже зафиксированной команды только подтверждается:
|
||
id входит в тот же снимок, что и переход, а `CommandAck` копится в
|
||
операции и уходит после её commit.
|
||
"""
|
||
state = hub.get(session_id)
|
||
if command_id is not None and state is not None:
|
||
if str(command_id) in state.processed_station_commands:
|
||
if ack is not None:
|
||
ack(CommandAck(command_id=command_id))
|
||
return True
|
||
try:
|
||
async with hub.operation(session_id):
|
||
await handler()
|
||
if command_id is not None and state is not None:
|
||
state.processed_station_commands.append(str(command_id))
|
||
del state.processed_station_commands[:-MAX_STATION_COMMANDS]
|
||
if ack is not None:
|
||
ack(CommandAck(command_id=command_id))
|
||
except Exception:
|
||
if not hub.is_lease_fenced(session_id):
|
||
raise
|
||
log.info("закрытие WebSocket после fencing занятия %s", session_id)
|
||
await close_fenced(ws, sender)
|
||
return False
|
||
if hub.is_lease_fenced(session_id):
|
||
await close_fenced(ws, sender)
|
||
return False
|
||
return True
|