feat(room): принудительный мьют участника организатором
Some checks failed
CI / backend (push) Has been cancelled
CI / frontend (push) Has been cancelled

Новый эндпоинт POST /conferences/{id}/mute-participant: права проверяются
ЗАНОВО по владельцу конференции в БД (ConferenceService.mute_participant),
не по метаданным LiveKit-токена вызывающего — те лишь подсказка для UI и
потенциально подделываемы клиентом. Обычный участник получает 403, чужая/
несуществующая конференция — 404, участник не в комнате LiveKit — отдельный
404 (participant_not_in_room).

Само выключение — серверный вызов api.LiveKitAPI (services/room_control.py,
тот же паттерн, что services/egress.py): backend аутентифицируется
СОБСТВЕННЫМИ api_key/api_secret, а не токеном организатора, поэтому
дополнительный LiveKit-грант в токене организатора не нужен — мьютит сервер
от своего имени. Если трек данного source не опубликован (с 0.0.15 участники
заходят с выключенными микрофоном/камерой) — не ошибка, а no-op: искомое
состояние уже достигнуто, ответ muted:false.

Уведомление участника — тот же общий канал комнаты, что и очередь рук
(hand_queue_channel): рассылается всем, получатель сам сверяет identity
(ForcedMuteWatcher, рендерится внутри LiveKitRoom). Само выключение трека
участник видит сразу через штатный useTrackToggle (LiveKit сам присылает
TrackMuted), тост только поясняет причину — иначе не отличить от глюка.
Включить себя обратно можно сразу тем же тулбаром, сервер это не блокирует.

Кнопки — на чужой плитке камеры, видны только организатору по наведению
(на тач-устройствах — всегда, как и булавка закрепления).

Тесты: владелец мьютит успешно и публикует broadcast, уже-выключенный трек —
muted:false без broadcast, администратор мьютит чужую конференцию, обычный
участник получает 403 без обращения к LiveKit, конференция не найдена и
участник не в комнате — соответствующие 404.
This commit is contained in:
2026-08-01 22:07:06 +03:00
parent 8e5eda88a2
commit 4c60e092e5
7 changed files with 476 additions and 0 deletions

View File

@@ -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,

View File

@@ -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

View File

@@ -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:

View File

@@ -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()

View File

@@ -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