diff --git a/backend/api/conferences.py b/backend/api/conferences.py index 1d75a42..954c378 100644 --- a/backend/api/conferences.py +++ b/backend/api/conferences.py @@ -18,6 +18,8 @@ from schemas.conferences import ( GuestJoinIn, JoinIn, JoinOut, + MuteParticipantIn, + MuteParticipantOut, OccurrenceOut, ResolveOut, ) @@ -34,6 +36,7 @@ from services.conferences import ( InviteeUserNotFoundError, NotConferenceOwnerError, ) +from services.room_control import ParticipantNotInRoomError router = APIRouter(prefix="/api/v1/conferences", tags=["conferences"]) @@ -171,6 +174,37 @@ async def guest_join_conference( ) from exc +@router.post("/{conference_id}/mute-participant", response_model=MuteParticipantOut) +async def mute_participant( + conference_id: uuid.UUID, + data: MuteParticipantIn, + user: Annotated[User, Depends(get_current_user)], + session: Annotated[AsyncSession, Depends(get_session)], +) -> MuteParticipantOut: + """Принудительно выключить микрофон/камеру участника (задача B2) — владелец/администратор. + + Права проверяются ЗАНОВО по владельцу конференции в БД + (`ConferenceService.mute_participant`), а не по метаданным LiveKit-токена + вызывающего — те лишь подсказка для UI и потенциально подделываемы клиентом. + """ + service = ConferenceService(session) + try: + muted = await service.mute_participant( + conference_id, actor=user, target_identity=data.identity, source=data.source + ) + except ConferenceNotFoundError as exc: + raise HTTPException( + status_code=status.HTTP_404_NOT_FOUND, detail="conference_not_found" + ) from exc + except NotConferenceOwnerError as exc: + raise HTTPException(status_code=status.HTTP_403_FORBIDDEN, detail="not_owner") from exc + except ParticipantNotInRoomError as exc: + raise HTTPException( + status_code=status.HTTP_404_NOT_FOUND, detail="participant_not_in_room" + ) from exc + return MuteParticipantOut(muted=muted) + + @router.get("/{conference_id}", response_model=ConferenceOut) async def get_conference( conference_id: uuid.UUID, diff --git a/backend/schemas/conferences.py b/backend/schemas/conferences.py index 8899fa1..5ece9da 100644 --- a/backend/schemas/conferences.py +++ b/backend/schemas/conferences.py @@ -6,6 +6,7 @@ from datetime import UTC, datetime, timedelta from pydantic import BaseModel, EmailStr, Field, field_serializer, field_validator, model_validator from core.plugins.config import SummaryRecipientsMode +from schemas.room_events import ForcedMuteSource from services.recurrence import RecurrenceRule # Допуск в прошлое при плановом создании/правке — небольшой запас на задержку @@ -210,3 +211,21 @@ class GuestJoinIn(BaseModel): display_name: str = Field(min_length=1, max_length=255) email: EmailStr | None = None password: str | None = None + + +class MuteParticipantIn(BaseModel): + """Тело запроса принудительного мьюта участника организатором (задача B2). + + `identity` — тот же формат, что и `Participant.identity` в LiveKit + (`str(user_id)` либо `guest:{id}`); клиент берёт его из `useParticipants()` + LiveKit, не подбирает вручную. + """ + + identity: str = Field(min_length=1) + source: ForcedMuteSource + + +class MuteParticipantOut(BaseModel): + """Ответ на принудительный мьют — `muted=False`, если трек и так не был опубликован.""" + + muted: bool diff --git a/backend/services/conferences.py b/backend/services/conferences.py index d2cfa26..87935ea 100644 --- a/backend/services/conferences.py +++ b/backend/services/conferences.py @@ -33,12 +33,14 @@ from schemas.conferences import ( JoinOut, OccurrenceOut, ) +from services import hand_queue from services.avatars import avatar_url as resolve_avatar_url from services.conference_access import build_join, ensure_joinable from services.conference_ids import generate_number, generate_slug from services.instance_settings import InstanceSettingsService from services.invitations_producer import enqueue_invitations from services.recurrence import RecurrenceRule, expand_occurrences +from services.room_control import MuteSource, mute_participant_track logger = logging.getLogger(__name__) @@ -271,6 +273,36 @@ class ConferenceService: chat_enabled=chat_enabled, ) + async def mute_participant( + self, + conference_id: uuid.UUID, + *, + actor: User, + target_identity: str, + source: MuteSource, + ) -> bool: + """Принудительно замьютить трек участника (задача B2); владелец/администратор. + + Права — ТОЛЬКО отсюда (`_ensure_owner_or_admin` по `conference.owner_id` + в БД), не по метаданным LiveKit-токена вызывающего: те лишь подсказка + для UI (см. `services/conference_access.py::build_join`) и потенциально + подделываемы клиентом. `target_identity` НИКАК не валидируется против + состава участников заранее — если его сейчас нет в комнате LiveKit, + `mute_participant_track` бросит `ParticipantNotInRoomError` (ловит + API-роутер). + """ + conference = await self._get_or_raise(conference_id) + self._ensure_owner_or_admin(conference, actor) + + muted = await mute_participant_track( + conference.slug, identity=target_identity, source=source + ) + if muted: + await hand_queue.publish_forced_mute( + conference.id, identity=target_identity, source=source + ) + return muted + async def update( self, conference_id: uuid.UUID, *, actor: User, data: ConferenceUpdateIn ) -> Conference: diff --git a/backend/services/room_control.py b/backend/services/room_control.py new file mode 100644 index 0000000..22d23cc --- /dev/null +++ b/backend/services/room_control.py @@ -0,0 +1,81 @@ +"""Управление комнатой LiveKit от имени организатора (задача B2): принудительный мьют. + +Тонкая обёртка над `RoomServiceClient` — тот же паттерн, что и +`services/egress.py` (единственная точка мокирования в тестах, свой +`api.LiveKitAPI` на вызов, аутентификация СЕРВЕРНЫМИ `api_key`/`api_secret`, +а не токеном организатора). Именно поэтому организатору не нужен отдельный +LiveKit-грант в собственном access-токене под это действие — мьютит backend +от своего имени, клиент лишь инициирует вызов, а право на это проверяется +по владельцу конференции в БД (`services/conferences.py::mute_participant`), +ДО обращения сюда. +""" + +import logging + +from livekit import api +from livekit.protocol.models import TrackSource + +from core.config import get_settings +from schemas.room_events import ForcedMuteSource + +logger = logging.getLogger(__name__) + +# Переэкспорт под более общим именем — этот модуль не завязан на протокол WS +# (`schemas/room_events.py`), которому концептуально принадлежит `ForcedMuteSource`. +MuteSource = ForcedMuteSource + +_TRACK_SOURCE_BY_NAME: dict[MuteSource, int] = { + "microphone": TrackSource.MICROPHONE, + "camera": TrackSource.CAMERA, +} + + +class ParticipantNotInRoomError(Exception): + """Участника с таким identity сейчас нет в комнате LiveKit (уже вышел/не заходил).""" + + +async def mute_participant_track(room_name: str, *, identity: str, source: MuteSource) -> bool: + """Принудительно замьютить опубликованный трек участника; `True` — трек реально замьючен. + + Если трек данного `source` сейчас не опубликован — не ошибка, а no-op: + искомое состояние («трек не идёт») уже достигнуто. Обычный случай с + 0.0.15 — участники заходят с выключенными микрофоном/камерой (задача + A1), трек попросту не существует, пока человек не включит его сам; + мьютить в этот момент нечего, и это НЕ повод отвечать клиенту ошибкой. + """ + settings = get_settings() + lkapi = api.LiveKitAPI( + settings.livekit_url, + api_key=settings.livekit_api_key, + api_secret=settings.livekit_api_secret, + ) + try: + try: + participant = await lkapi.room.get_participant( + api.RoomParticipantIdentity(room=room_name, identity=identity) + ) + except api.TwirpError as exc: + if exc.status == 404: + raise ParticipantNotInRoomError from exc + raise + + target_source = _TRACK_SOURCE_BY_NAME[source] + track = next((t for t in participant.tracks if t.source == target_source), None) + if track is None or track.muted: + return False + + await lkapi.room.mute_published_track( + api.MuteRoomTrackRequest( + room=room_name, identity=identity, track_sid=track.sid, muted=True + ) + ) + logger.info( + "room_control: принудительный мьют — комната=%s identity=%s source=%s трек=%s", + room_name, + identity, + source, + track.sid, + ) + return True + finally: + await lkapi.aclose() diff --git a/backend/tests/test_mute_participant_api.py b/backend/tests/test_mute_participant_api.py new file mode 100644 index 0000000..4427b3f --- /dev/null +++ b/backend/tests/test_mute_participant_api.py @@ -0,0 +1,245 @@ +"""Тесты `POST /api/v1/conferences/{id}/mute-participant` (задача B2). + +`mute_participant_track` (реальный вызов LiveKit `RoomServiceClient`) мокается +на уровне `services.conferences` — тот же паттерн, что и `start_track_egress` +в `tests/test_livekit_webhook.py`: сетевой вызов к LiveKit в тестах не нужен, +важна только бизнес-логика (права, маршрутизация ошибок, broadcast). +""" + +import asyncio +import uuid +from typing import Any +from unittest.mock import AsyncMock + +import httpx +import pytest +from redis.asyncio.client import PubSub +from sqlalchemy.ext.asyncio import AsyncSession + +import services.conferences as conferences_module +from core.redis import redis_client +from core.security import create_access_token, hash_password +from models.conference import Conference +from models.user import User +from services.conference_ids import generate_number, generate_slug +from services.hand_queue import hand_queue_channel +from services.room_control import ParticipantNotInRoomError + +MUTE_URL = "{base}/mute-participant" + + +async def _receive_within(pubsub: PubSub, *, max_wait: float) -> dict[str, Any] | None: + """Дождаться СОДЕРЖАТЕЛЬНОГО сообщения канала в пределах `max_wait` секунд. + + `ignore_subscribe_messages=True` у `get_message` фильтрует служебное + подтверждение подписки, но при этом всё равно может вернуть `None` для + ЭТОГО конкретного вызова (см. `api/chat.py::_pump_pubsub_to_websocket`, + ровно поэтому там `while True: ... if raw is None: continue`) — здесь тот + же цикл, но с общим дедлайном вместо бесконечного ожидания. + """ + deadline = asyncio.get_event_loop().time() + max_wait + while True: + remaining = deadline - asyncio.get_event_loop().time() + if remaining <= 0: + return None + raw: dict[str, Any] | None = await pubsub.get_message( + ignore_subscribe_messages=True, timeout=remaining + ) + if raw is not None: + return raw + + +async def _make_user(session: AsyncSession, *, role: str = "user") -> User: + user = User( + email=f"{uuid.uuid4()}@example.com", + name_user="Mute Tester", + password_hash=hash_password("password123"), + email_verified=True, + role=role, + ) + session.add(user) + await session.flush() + return user + + +async def _make_conference( + session: AsyncSession, *, owner_id: uuid.UUID | None = None +) -> Conference: + conference = Conference( + number=generate_number(), + slug=generate_slug(), + title="Mute Test", + status="active", + owner_id=owner_id, + ) + session.add(conference) + await session.flush() + return conference + + +def _auth_headers(user: User) -> dict[str, str]: + return {"Authorization": f"Bearer {create_access_token(user.id, user.role)}"} + + +def _url(conference_id: uuid.UUID) -> str: + return MUTE_URL.format(base=f"/api/v1/conferences/{conference_id}") + + +async def test_owner_can_mute_participant_and_broadcast_is_published( + client: httpx.AsyncClient, db_session: AsyncSession, monkeypatch: pytest.MonkeyPatch +) -> None: + owner = await _make_user(db_session) + conference = await _make_conference(db_session, owner_id=owner.id) + await db_session.commit() + + mock_mute = AsyncMock(return_value=True) + monkeypatch.setattr(conferences_module, "mute_participant_track", mock_mute) + + pubsub = redis_client.pubsub() + channel = hand_queue_channel(conference.id) + await pubsub.subscribe(channel) + try: + response = await client.post( + _url(conference.id), + json={"identity": "some-identity", "source": "microphone"}, + headers=_auth_headers(owner), + ) + assert response.status_code == 200, response.text + assert response.json() == {"muted": True} + mock_mute.assert_awaited_once_with( + conference.slug, identity="some-identity", source="microphone" + ) + + raw = await _receive_within(pubsub, max_wait=2) + assert raw is not None + assert raw["data"] == ( + '{"type":"forced_mute","identity":"some-identity","source":"microphone"}' + ) + finally: + await pubsub.unsubscribe(channel) + await pubsub.aclose() # type: ignore[no-untyped-call] + + +async def test_mute_already_off_does_not_broadcast( + client: httpx.AsyncClient, db_session: AsyncSession, monkeypatch: pytest.MonkeyPatch +) -> None: + """Трек не был опубликован (камера/мьют и так выключены, задача A1) — не ошибка. + + Ответ `muted=false` (искомое состояние уже достигнуто), без broadcast'а. + """ + owner = await _make_user(db_session) + conference = await _make_conference(db_session, owner_id=owner.id) + await db_session.commit() + + monkeypatch.setattr( + conferences_module, "mute_participant_track", AsyncMock(return_value=False) + ) + + pubsub = redis_client.pubsub() + channel = hand_queue_channel(conference.id) + await pubsub.subscribe(channel) + try: + response = await client.post( + _url(conference.id), + json={"identity": "some-identity", "source": "camera"}, + headers=_auth_headers(owner), + ) + assert response.status_code == 200, response.text + assert response.json() == {"muted": False} + + raw = await _receive_within(pubsub, max_wait=0.5) + assert raw is None + finally: + await pubsub.unsubscribe(channel) + await pubsub.aclose() # type: ignore[no-untyped-call] + + +async def test_admin_can_mute_participant_of_someone_elses_conference( + client: httpx.AsyncClient, db_session: AsyncSession, monkeypatch: pytest.MonkeyPatch +) -> None: + owner = await _make_user(db_session) + admin = await _make_user(db_session, role="admin") + conference = await _make_conference(db_session, owner_id=owner.id) + await db_session.commit() + + monkeypatch.setattr(conferences_module, "mute_participant_track", AsyncMock(return_value=True)) + + response = await client.post( + _url(conference.id), + json={"identity": "some-identity", "source": "microphone"}, + headers=_auth_headers(admin), + ) + assert response.status_code == 200, response.text + + +async def test_regular_participant_cannot_mute_someone_else( + client: httpx.AsyncClient, db_session: AsyncSession, monkeypatch: pytest.MonkeyPatch +) -> None: + owner = await _make_user(db_session) + other = await _make_user(db_session) + conference = await _make_conference(db_session, owner_id=owner.id) + await db_session.commit() + + mock_mute = AsyncMock() + monkeypatch.setattr(conferences_module, "mute_participant_track", mock_mute) + + response = await client.post( + _url(conference.id), + json={"identity": str(owner.id), "source": "microphone"}, + headers=_auth_headers(other), + ) + assert response.status_code == 403 + assert response.json()["detail"] == "not_owner" + mock_mute.assert_not_awaited() + + +async def test_mute_conference_not_found( + client: httpx.AsyncClient, db_session: AsyncSession +) -> None: + user = await _make_user(db_session) + await db_session.commit() + + response = await client.post( + _url(uuid.uuid4()), + json={"identity": "some-identity", "source": "microphone"}, + headers=_auth_headers(user), + ) + assert response.status_code == 404 + assert response.json()["detail"] == "conference_not_found" + + +async def test_mute_participant_not_in_room_returns_404( + client: httpx.AsyncClient, db_session: AsyncSession, monkeypatch: pytest.MonkeyPatch +) -> None: + owner = await _make_user(db_session) + conference = await _make_conference(db_session, owner_id=owner.id) + await db_session.commit() + + monkeypatch.setattr( + conferences_module, + "mute_participant_track", + AsyncMock(side_effect=ParticipantNotInRoomError()), + ) + + response = await client.post( + _url(conference.id), + json={"identity": "ghost", "source": "microphone"}, + headers=_auth_headers(owner), + ) + assert response.status_code == 404 + assert response.json()["detail"] == "participant_not_in_room" + + +async def test_mute_rejects_invalid_source( + client: httpx.AsyncClient, db_session: AsyncSession +) -> None: + owner = await _make_user(db_session) + conference = await _make_conference(db_session, owner_id=owner.id) + await db_session.commit() + + response = await client.post( + _url(conference.id), + json={"identity": "some-identity", "source": "screen_share"}, + headers=_auth_headers(owner), + ) + assert response.status_code == 422 diff --git a/frontend/src/api/conferences.ts b/frontend/src/api/conferences.ts index fe39a6a..de7d0d0 100644 --- a/frontend/src/api/conferences.ts +++ b/frontend/src/api/conferences.ts @@ -145,6 +145,14 @@ export interface ConferenceGuestJoinPayload { password?: string } +/** Источник трека, который организатор может принудительно выключить (задача B2). */ +export type MuteSource = 'microphone' | 'camera' + +/** Ответ на принудительный мьют — `false`, если трек и так не был опубликован (нечего было мьютить). */ +export interface MuteParticipantResult { + muted: boolean +} + /** Тело частичного обновления конференции — те же поля, что и при создании, все опциональны. */ export type ConferenceUpdatePayload = Partial @@ -202,6 +210,23 @@ export async function guestJoinConference( }) } +/** + * Принудительно выключить микрофон/камеру участника (задача B2) — только + * владелец конференции/администратор, иначе 403 (`not_owner`). 404 + * (`participant_not_in_room`) — участника с таким `identity` сейчас нет в + * комнате LiveKit. + */ +export async function muteParticipant( + conferenceId: string, + identity: string, + source: MuteSource, +): Promise { + return apiRequest(`/conferences/${conferenceId}/mute-participant`, { + method: 'POST', + body: { identity, source }, + }) +} + /** Список «моих» конференций — закреплённые (повторяющиеся) и предстоящие разовые владельца. */ export async function getMyConferences(): Promise { return apiRequest('/conferences/my') diff --git a/frontend/src/components/room/ForcedMuteWatcher.tsx b/frontend/src/components/room/ForcedMuteWatcher.tsx new file mode 100644 index 0000000..dee7923 --- /dev/null +++ b/frontend/src/components/room/ForcedMuteWatcher.tsx @@ -0,0 +1,40 @@ +import { useEffect } from 'react' +import { useLocalParticipant } from '@livekit/components-react' +import { useToast } from '@/components/ui/ToastProvider' +import type { ForcedMuteEvent } from '@/hooks/useChat' + +/** + * Уведомляет ЛОКАЛЬНОГО участника тостом, когда организатор принудительно + * выключил его микрофон/камеру (задача B2). Рендерится безусловно внутри + * `` — `useLocalParticipant` недоступен снаружи (`RoomPage` + * сам вне контекста LiveKit, см. докстринг `useIsOrganizer`). + * + * Само выключение трека организатор делает СЕРВЕРНЫМ вызовом LiveKit API + * (`services/room_control.py`) — тулбарные кнопки (useTrackToggle) сами + * отразят новое состояние по родному событию LiveKit `TrackMuted`, этот + * компонент только поясняет ПОЧЕМУ: без тоста человек не отличил бы + * действие организатора от случайного глюка. Участник может включить себя + * обратно сразу тем же тулбаром — сервер это не блокирует (см. докстринг B2 + * в CHANGELOG/коммите). + */ +export function ForcedMuteWatcher({ event }: { event: ForcedMuteEvent | null }) { + const { localParticipant } = useLocalParticipant() + const toast = useToast() + + useEffect(() => { + if (!event || event.identity !== localParticipant.identity) return + toast.show( + event.source === 'microphone' + ? 'Организатор выключил ваш микрофон' + : 'Организатор выключил вашу камеру', + 'info', + ) + // `event` (включая `nonce`) — единственная зависимость, которая должна + // повторно показывать тост; `localParticipant`/`toast` стабильны в + // рамках подключения и намеренно не входят в список, чтобы их + // пересоздание (если когда-нибудь случится) не дублировало уведомление. + // eslint-disable-next-line react-hooks/exhaustive-deps + }, [event]) + + return null +}