diff --git a/.gitignore b/.gitignore index e333f2320b4..fdda7a37623 100644 --- a/.gitignore +++ b/.gitignore @@ -79,3 +79,9 @@ book/ # Don't include users' poetry configs /poetry.toml +admin_token.txt +vod_service_token.txt +google_api_key.txt +admin_token.txt +vod_service_token.txt +google_api_key.txt diff --git a/homeserver.yaml.example b/homeserver.yaml.example index 992a7d1dc46..6119dd0ad28 100644 --- a/homeserver.yaml.example +++ b/homeserver.yaml.example @@ -48,10 +48,24 @@ admins: - "@admin:localhost" modules: + - module: modules.vod_service.module.VodServiceModule + config: + object_storage_base_url: "https://objectstorage..oraclecloud.com" + namespace: "" + bucket: "vod-teste" - module: modules.room_service.module.RoomServiceModule config: - admin_user_id: "@admin:localhost" - admin_token: "syt_" - homeserver: "http://localhost:XXXX" + admin_user_id: "@admin:localhost" # IMPORTANTE!!! PRECISA ADAPTAR CASO A CASO + admin_token_file: "/data/admin_token.txt" + homeserver: "http://localhost:3000" + - module: modules.room_features_service.module.RoomFeaturesServiceModule + config: + admin_user_id: "@admin:localhost" # IMPORTANTE!!! PRECISA ADAPTAR CASO A CASO + admin_token_file: "/data/admin_token.txt" + homeserver: "http://localhost:3000" + - module: modules.schedule_service.module.ScheduleServiceModule + config: + google_api_key_file: "/data/google_api_key.txt" + timezone: "America/Sao_Paulo" - module: modules.bundle_service.module.BundleServiceModule # vim:ft=yaml diff --git a/infra/synapse-docker/data/homeserver.yaml b/infra/synapse-docker/data/homeserver.yaml index 7216261af71..4fa31171fcc 100644 --- a/infra/synapse-docker/data/homeserver.yaml +++ b/infra/synapse-docker/data/homeserver.yaml @@ -72,10 +72,24 @@ admins: - "@admin:localhost" # IMPORTANTE!!! PRECISA ADAPTAR CASO A CASO modules: + - module: modules.vod_service.module.VodServiceModule + config: + object_storage_base_url: "https://objectstorage..oraclecloud.com" + namespace: "" + bucket: "vod-teste" - module: modules.room_service.module.RoomServiceModule config: admin_user_id: "@admin:localhost" # IMPORTANTE!!! PRECISA ADAPTAR CASO A CASO admin_token_file: "/data/admin_token.txt" homeserver: "http://synapse:3000" + - module: modules.room_features_service.module.RoomFeaturesServiceModule + config: + admin_user_id: "@admin:localhost" # IMPORTANTE!!! PRECISA ADAPTAR CASO A CASO + admin_token_file: "/data/admin_token.txt" + homeserver: "http://synapse:3000" + - module: modules.schedule_service.module.ScheduleServiceModule + config: + google_api_key_file: "/home/livia/BuzzLabs/synapse/google_api_key.txt" + timezone: "America/Sao_Paulo" - module: modules.bundle_service.module.BundleServiceModule # vim:ft=yaml diff --git a/infra/synapse-docker/migrations/001_init.sql b/infra/synapse-docker/migrations/001_init.sql index f1ad13b846e..3c4ec8e75b5 100644 --- a/infra/synapse-docker/migrations/001_init.sql +++ b/infra/synapse-docker/migrations/001_init.sql @@ -28,3 +28,49 @@ CREATE TABLE public.room_business ( keyword text UNIQUE, created_at int8 NOT NULL ); + +CREATE TABLE public.room_features ( + room_id text NOT NULL, + feature text NOT NULL, + enabled boolean NOT NULL DEFAULT false, + created_at int8 NOT NULL, + updated_at int8 NOT NULL, + PRIMARY KEY (room_id, feature), + FOREIGN KEY (room_id) + REFERENCES public.room_business(room_id) + ON DELETE CASCADE +); + +CREATE INDEX room_features_room_idx ON public.room_features (room_id); + +CREATE TABLE public.streams ( + id BIGSERIAL PRIMARY KEY, + stream_id TEXT, + room_id TEXT NOT NULL, + title TEXT, + category_id BIGINT, + recording_path TEXT, + recording_duration_ms BIGINT, + recording_started_at BIGINT, + recording_ended_at BIGINT, + started_at BIGINT, + ended_at BIGINT, + created_at BIGINT, + updated_at BIGINT, + FOREIGN KEY (room_id) + REFERENCES public.room_business(room_id) + ON DELETE CASCADE +); + +CREATE INDEX streams_room_started_idx ON public.streams (room_id, started_at DESC); +CREATE INDEX streams_stream_id_idx ON public.streams (stream_id); + +CREATE TABLE public.room_calendars ( + room_id text NOT NULL, + calendar_id text NOT NULL, + created_at int8 NOT NULL, + updated_at int8 NOT NULL, + PRIMARY KEY (room_id), + FOREIGN KEY (room_id) REFERENCES public.room_business(room_id) ON DELETE CASCADE +); +ALTER TABLE room_calendars OWNER TO ; \ No newline at end of file diff --git a/infra/k8s/admin_token.txt b/modules/bundle_service/resources/__init__.py similarity index 100% rename from infra/k8s/admin_token.txt rename to modules/bundle_service/resources/__init__.py diff --git a/infra/synapse-docker/data/admin_token.txt b/modules/bundle_service/tests/__init__.py similarity index 100% rename from infra/synapse-docker/data/admin_token.txt rename to modules/bundle_service/tests/__init__.py diff --git a/modules/room_features_service/README b/modules/room_features_service/README new file mode 100644 index 00000000000..934139cfd81 --- /dev/null +++ b/modules/room_features_service/README @@ -0,0 +1,184 @@ +# Room Features Service Module for Matrix Synapse + +This module exposes **HTTP endpoints** to enable/disable optional per-room +features (feature flags), stored in the `room_features` table. + +Each feature toggles an optional tab/section in the room UI (e.g. the VODs +tab, the events tab). A feature that is not set is treated as **disabled**. + +This is a **UI/capability toggle**, not access control: it decides whether a +room *shows* a feature, not who may read the room's content. Access to the +room itself is handled elsewhere (membership, paywall, etc.). + +--- + +## Features + +- store per-room feature flags in a dedicated table (`room_features`) +- read a single flag, or all flags for a room +- toggle a flag (global Synapse admins only) +- missing flag = disabled (reads never 404 on an unset feature) +- generic by design: any feature key works, so new features (events, cuts, + polls, ...) don't require code changes to this module + +--- + +## How it works + +1. Synapse loads the module on startup +2. The module initializes a `RoomFeaturesService` +3. Custom HTTP endpoints are registered under: + /\_synapse/room_features/\* +4. The frontend calls these endpoints: + - on room load, to decide which tabs to show + - from room settings, when an admin toggles a feature +5. The module reads/writes the `room_features` table and returns JSON + +--- + +## Registered endpoints + +All endpoints are **POST** with a JSON body (same convention as +`room_service`). + +- /\_synapse/room_features/get +- /\_synapse/room_features/list +- /\_synapse/room_features/set + +--- + +### Get one feature + +POST /get + +Returns whether a single feature is enabled for a room. An unset feature +returns `enabled: false` (never 404). + +**Body:** + +- `room_id` (required) +- `feature` (required) – e.g. `"vods"`, `"events"` + +**Response:** + +- `roomId` +- `feature` +- `enabled` + +**Errors:** + +- `400` – room_id or feature missing + +--- + +### List all features + +POST /list + +Returns every feature flag set for a room, as a map. Absent features are +simply not present in the map (treated as disabled by the caller). + +**Body:** + +- `room_id` (required) + +**Response:** + +- `roomId` +- `features` – object like `{ "vods": true, "events": false }` + +**Errors:** + +- `400` – room_id missing + +--- + +### Set a feature (admin only) + +POST /set + +Enables or disables a feature for a room. Requires a **global Synapse admin** +(the requester is identified from the access token; non-admins get `403`). + +Upsert: creates the row if it doesn't exist, updates it otherwise. + +**Body:** + +- `room_id` (required) +- `feature` (required) +- `enabled` (required, boolean) + +**Response:** + +- `roomId` +- `feature` +- `enabled` + +**Errors:** + +- `400` – room_id/feature/enabled missing or `enabled` not a boolean +- `403` – requester is not a global admin + +--- + +## Database + +Table `room_features`: + +```sql +CREATE TABLE public.room_features ( + room_id text NOT NULL, + feature text NOT NULL, + enabled boolean NOT NULL DEFAULT false, + created_at int8 NOT NULL, + updated_at int8 NOT NULL, + PRIMARY KEY (room_id, feature), + FOREIGN KEY (room_id) + REFERENCES public.room_business(room_id) + ON DELETE CASCADE +); +CREATE INDEX room_features_room_idx ON public.room_features (room_id); +``` + +Notes: + +- `(room_id, feature)` is the primary key: one row per feature per room. +- The FK to `room_business` means a room must exist there before it can have + features. Timestamps are epoch milliseconds (`int8`), matching the other + modules' tables. +- The Synapse database user needs privileges on this table: + +```sql +ALTER TABLE room_features OWNER TO ; +``` + +--- + +## Configuration + +Example `homeserver.yaml` configuration: + +```yaml +modules: + - module: modules.room_features_service.module.RoomFeaturesServiceModule + config: + admin_user_id: "@admin:localhost" + admin_token_file: "/path/to/admin_token.txt" + homeserver: "http://localhost:3000" +``` + +`admin_token_file` is read at startup (waiting briefly for the file to +appear, like `room_service`). `admin_token` (inline) is also accepted as a +fallback. + +--- + +## Testing + +``` +PYTHONPATH=. pytest -vv modules/room_features_service/tests/ +``` + +Covers: missing-flag-returns-false, enabled/disabled reads, the full list +map, admin enforcement on writes (403 for non-admins), the upsert, and +input validation. \ No newline at end of file diff --git a/modules/room_features_service/__init__.py b/modules/room_features_service/__init__.py new file mode 100644 index 00000000000..e69de29bb2d diff --git a/modules/room_features_service/db.py b/modules/room_features_service/db.py new file mode 100644 index 00000000000..75bb13d7a82 --- /dev/null +++ b/modules/room_features_service/db.py @@ -0,0 +1,37 @@ +def get_feature(txn, room_id: str, feature: str): + txn.execute( + """ + SELECT enabled + FROM room_features + WHERE room_id = ? AND feature = ? + """, + (room_id, feature), + ) + return txn.fetchone() + + +def list_features(txn, room_id: str): + txn.execute( + """ + SELECT feature, enabled + FROM room_features + WHERE room_id = ? + """, + (room_id,), + ) + return txn.fetchall() + + +def set_feature(txn, room_id: str, feature: str, enabled: bool, now_ms: int): + """ + Upsert: liga/desliga a feature da sala. Cria a linha se nao existir. + """ + txn.execute( + """ + INSERT INTO room_features (room_id, feature, enabled, created_at, updated_at) + VALUES (?, ?, ?, ?, ?) + ON CONFLICT (room_id, feature) + DO UPDATE SET enabled = EXCLUDED.enabled, updated_at = EXCLUDED.updated_at + """, + (room_id, feature, enabled, now_ms, now_ms), + ) \ No newline at end of file diff --git a/modules/room_features_service/module.py b/modules/room_features_service/module.py new file mode 100644 index 00000000000..e8ddb557966 --- /dev/null +++ b/modules/room_features_service/module.py @@ -0,0 +1,63 @@ +# module.py +import logging +import time +import os + +from synapse.module_api import ModuleApi + +from .service import RoomFeaturesService +from .resources.get_feature import GetFeatureResource +from .resources.list_features import ListFeaturesResource +from .resources.set_feature import SetFeatureResource + +logger = logging.getLogger(__name__) + + +def _read_admin_token(config: dict): + """ + Le o admin token, esperando o arquivo aparecer (~30s), igual ao room_service. + Aceita admin_token_file (preferido) ou admin_token direto. + """ + admin_token_file = config.get("admin_token_file") + if admin_token_file: + for _ in range(30): # tenta por ~30s + if os.path.exists(admin_token_file): + with open(admin_token_file, "r") as f: + return f.read().strip() + time.sleep(1) + return None + + return config.get("admin_token") + + +class RoomFeaturesServiceModule: + def __init__(self, config, api: ModuleApi): + self.hs = api._hs + + admin_user_id = config["admin_user_id"] + homeserver = config["homeserver"] + admin_token = _read_admin_token(config) + + service = RoomFeaturesService( + api=api, + admin_user_id=admin_user_id, + admin_token=admin_token, + homeserver=homeserver, + ) + + self.hs.room_features_service = service + + api.register_web_resource( + "/_synapse/room_features/get", + GetFeatureResource(api, service), + ) + api.register_web_resource( + "/_synapse/room_features/list", + ListFeaturesResource(api, service), + ) + api.register_web_resource( + "/_synapse/room_features/set", + SetFeatureResource(api, service), + ) + + logger.info("RoomFeaturesServiceModule carregado") \ No newline at end of file diff --git a/modules/room_features_service/resources/__init__.py b/modules/room_features_service/resources/__init__.py new file mode 100644 index 00000000000..e69de29bb2d diff --git a/modules/room_features_service/resources/get_feature.py b/modules/room_features_service/resources/get_feature.py new file mode 100644 index 00000000000..49d289e3e1e --- /dev/null +++ b/modules/room_features_service/resources/get_feature.py @@ -0,0 +1,24 @@ +from synapse.api.errors import SynapseError +from synapse.http.server import DirectServeJsonResource +from synapse.http.servlet import parse_json_object_from_request + + +class GetFeatureResource(DirectServeJsonResource): + isLeaf = True + + def __init__(self, api, service): + super().__init__() + self.api = api + self.service = service + + async def _async_render_POST(self, request): + content = parse_json_object_from_request(request) + room_id = content.get("room_id") + feature = content.get("feature") + if not room_id: + raise SynapseError(400, "room_id is required") + if not feature: + raise SynapseError(400, "feature is required") + + result = await self.service.get_feature(room_id=room_id, feature=feature) + return 200, result \ No newline at end of file diff --git a/modules/room_features_service/resources/list_features.py b/modules/room_features_service/resources/list_features.py new file mode 100644 index 00000000000..9448a36d35c --- /dev/null +++ b/modules/room_features_service/resources/list_features.py @@ -0,0 +1,21 @@ +from synapse.api.errors import SynapseError +from synapse.http.server import DirectServeJsonResource +from synapse.http.servlet import parse_json_object_from_request + + +class ListFeaturesResource(DirectServeJsonResource): + isLeaf = True + + def __init__(self, api, service): + super().__init__() + self.api = api + self.service = service + + async def _async_render_POST(self, request): + content = parse_json_object_from_request(request) + room_id = content.get("room_id") + if not room_id: + raise SynapseError(400, "room_id is required") + + result = await self.service.list_features(room_id=room_id) + return 200, result \ No newline at end of file diff --git a/modules/room_features_service/resources/set_feature.py b/modules/room_features_service/resources/set_feature.py new file mode 100644 index 00000000000..952f1a2cf34 --- /dev/null +++ b/modules/room_features_service/resources/set_feature.py @@ -0,0 +1,33 @@ +# set_feature.py +from synapse.api.errors import SynapseError +from synapse.http.server import DirectServeJsonResource +from synapse.http.servlet import parse_json_object_from_request + + +class SetFeatureResource(DirectServeJsonResource): + isLeaf = True + + def __init__(self, api, service): + super().__init__() + self.api = api + self.service = service + + async def _async_render_POST(self, request): + requester = await self.api.get_user_by_req(request) + content = parse_json_object_from_request(request) + + room_id = content.get("room_id") + feature = content.get("feature") + enabled = content.get("enabled") + + if not room_id: + raise SynapseError(400, "room_id is required") + if not feature: + raise SynapseError(400, "feature is required") + if enabled is None: + raise SynapseError(400, "enabled is required") + + result = await self.service.set_feature( + requester=requester, room_id=room_id, feature=feature, enabled=enabled, + ) + return 200, result \ No newline at end of file diff --git a/modules/room_features_service/service.py b/modules/room_features_service/service.py new file mode 100644 index 00000000000..e025a838877 --- /dev/null +++ b/modules/room_features_service/service.py @@ -0,0 +1,105 @@ +# service.py +import logging +import time + +from synapse.module_api import ModuleApi +from synapse.api.errors import SynapseError + +from . import db + +logger = logging.getLogger(__name__) + + +class RoomFeaturesService: + """ + Liga/desliga features por sala (tabela room_features). + + Leitura: qualquer um pode consultar se uma feature esta ligada. + Escrita: apenas admin global (reusa is_admin, igual ao room_service). + + Feature ausente = desligada (get retorna enabled=False, nunca 404). + """ + + def __init__( + self, + api: ModuleApi, + admin_user_id: str, + admin_token: str, + homeserver: str, + ): + self.api = api + self.hs = api._hs + self.admin_user_id = admin_user_id + self.admin_token = admin_token + self.homeserver = homeserver + + self.store = self.hs.get_datastores().main + + # ---------------- ADMIN CHECK ---------------- + async def assert_is_admin(self, user_id: str): + is_admin = await self.api.is_user_admin(user_id) + if not is_admin: + raise SynapseError(403, "Only Synapse admins can perform this action") + + # ---------------- READ ---------------- + async def get_feature(self, room_id: str, feature: str) -> dict: + logger.info("get_feature: room_id=%s feature=%s", room_id, feature) + + if not room_id: + raise SynapseError(400, "room_id is required") + if not feature: + raise SynapseError(400, "feature is required") + + row = await self.store.db_pool.runInteraction( + "get_room_feature", + db.get_feature, + room_id, + feature, + ) + + # ausente = desligada, nunca 404 + enabled = bool(row[0]) if row else False + + return { + "roomId": room_id, + "feature": feature, + "enabled": enabled, + } + + async def list_features(self, room_id: str) -> dict: + logger.info("list_features: room_id=%s", room_id) + + if not room_id: + raise SynapseError(400, "room_id is required") + + rows = await self.store.db_pool.runInteraction( + "list_room_features", + db.list_features, + room_id, + ) + + return { + "roomId": room_id, + "features": {feature: bool(enabled) for (feature, enabled) in rows}, + } + + # ---------------- WRITE ---------------- + async def set_feature(self, *, requester, room_id: str, feature: str, enabled: bool) -> dict: + user_id = requester.user.to_string() + logger.info("set_feature: requester=%s room_id=%s feature=%s enabled=%s", + user_id, room_id, feature, enabled) + + if not room_id: + raise SynapseError(400, "room_id is required") + if not feature: + raise SynapseError(400, "feature is required") + if not isinstance(enabled, bool): + raise SynapseError(400, "enabled must be a boolean") + + await self.assert_is_admin(user_id) + + now_ms = int(time.time() * 1000) + await self.store.db_pool.runInteraction( + "set_room_feature", db.set_feature, room_id, feature, enabled, now_ms, + ) + return {"roomId": room_id, "feature": feature, "enabled": enabled} \ No newline at end of file diff --git a/modules/room_features_service/tests/__init__.py b/modules/room_features_service/tests/__init__.py new file mode 100644 index 00000000000..e69de29bb2d diff --git a/modules/room_features_service/tests/test_room_features_service.py b/modules/room_features_service/tests/test_room_features_service.py new file mode 100644 index 00000000000..8f77610635f --- /dev/null +++ b/modules/room_features_service/tests/test_room_features_service.py @@ -0,0 +1,266 @@ +import pytest +from unittest.mock import MagicMock, AsyncMock + +from synapse.api.errors import SynapseError +from modules.room_features_service.service import RoomFeaturesService +from modules.room_features_service import db + +# testa +# 1. get retorna enabled=False quando a feature nao existe (nunca 404) +# 2. get retorna o valor gravado quando existe +# 3. set exige admin global +# 4. set faz upsert com os parametros certos +# 5. validacao de room_id / feature / enabled + +# para testar: PYTHONPATH=. pytest -vv modules/room_features_service/tests/ + +ROOM_ID = "!room:localhost" +ADMIN = "@admin:localhost" +USER = "@user:localhost" + + +def _make_service(row=None, rows=None, is_admin=True): + api = MagicMock() + hs = MagicMock() + api._hs = hs + + api.is_user_admin = AsyncMock(return_value=is_admin) + + store = MagicMock() + + async def fake_run_interaction(desc, func, *args): + if desc == "get_room_feature": + return row + if desc == "list_room_features": + return rows if rows is not None else [] + if desc == "set_room_feature": + return None + raise AssertionError(f"runInteraction inesperado: {desc}") + + store.db_pool.runInteraction = AsyncMock(side_effect=fake_run_interaction) + hs.get_datastores.return_value.main = store + + service = RoomFeaturesService( + api=api, + admin_user_id=ADMIN, + admin_token="token", + homeserver="http://localhost:8008", + ) + service._store_mock = store + return service + +def _make_requester(user_id: str): + requester = MagicMock() + requester.user.to_string.return_value = user_id + return requester + +def _make_txn(fetchone=None, fetchall=None): + txn = MagicMock() + txn.fetchone.return_value = fetchone + txn.fetchall.return_value = fetchall if fetchall is not None else [] + return txn + + +# ---------------- GET ---------------- + +@pytest.mark.asyncio +async def test_get_feature_ausente_retorna_false(): + """ + Feature sem linha na tabela deve retornar enabled=False, nao 404. + """ + service = _make_service(row=None) + + result = await service.get_feature(room_id=ROOM_ID, feature="vods") + + assert result["enabled"] is False + assert result["roomId"] == ROOM_ID + assert result["feature"] == "vods" + + +@pytest.mark.asyncio +async def test_get_feature_ligada(): + """ + Feature ligada deve retornar enabled=True. + """ + service = _make_service(row=(True,)) + + result = await service.get_feature(room_id=ROOM_ID, feature="vods") + + assert result["enabled"] is True + + +@pytest.mark.asyncio +async def test_get_feature_desligada(): + """ + Feature com linha mas enabled=False deve retornar False. + """ + service = _make_service(row=(False,)) + + result = await service.get_feature(room_id=ROOM_ID, feature="vods") + + assert result["enabled"] is False + + +@pytest.mark.asyncio +async def test_get_feature_sem_room_id(): + """ + room_id vazio deve retornar 400. + """ + service = _make_service() + + with pytest.raises(SynapseError) as err: + await service.get_feature(room_id="", feature="vods") + + assert err.value.code == 400 + + +@pytest.mark.asyncio +async def test_get_feature_sem_feature(): + """ + feature vazia deve retornar 400. + """ + service = _make_service() + + with pytest.raises(SynapseError) as err: + await service.get_feature(room_id=ROOM_ID, feature="") + + assert err.value.code == 400 + + +# ---------------- LIST ---------------- + +@pytest.mark.asyncio +async def test_list_features_mapeia_todas(): + """ + list deve devolver um dict feature -> bool. + """ + service = _make_service(rows=[("vods", True), ("events", False)]) + + result = await service.list_features(room_id=ROOM_ID) + + assert result["features"] == {"vods": True, "events": False} + + +@pytest.mark.asyncio +async def test_list_features_vazio(): + """ + Sala sem features deve devolver dict vazio. + """ + service = _make_service(rows=[]) + + result = await service.list_features(room_id=ROOM_ID) + + assert result["features"] == {} + + +# ---------------- SET (admin) ---------------- + +@pytest.mark.asyncio +async def test_set_feature_exige_admin(): + """ + Nao-admin deve receber 403 e nada deve ser gravado. + """ + service = _make_service(is_admin=False) + + with pytest.raises(SynapseError) as err: + await service.set_feature( + requester=_make_requester(USER), + room_id=ROOM_ID, + feature="vods", + enabled=True, + ) + + assert err.value.code == 403 + # nao deve ter chamado o set no banco + calls = [c.args[0] for c in service._store_mock.db_pool.runInteraction.await_args_list] + assert "set_room_feature" not in calls + + +@pytest.mark.asyncio +async def test_set_feature_admin_grava(): + """ + Admin deve conseguir ligar a feature. + """ + service = _make_service(is_admin=True) + + result = await service.set_feature( + requester=_make_requester(ADMIN), + room_id=ROOM_ID, + feature="vods", + enabled=True, + ) + + assert result["enabled"] is True + + call = service._store_mock.db_pool.runInteraction.await_args + # (desc, func, room_id, feature, enabled, now_ms) + assert call.args[0] == "set_room_feature" + assert call.args[2] == ROOM_ID + assert call.args[3] == "vods" + assert call.args[4] is True + + +@pytest.mark.asyncio +async def test_set_feature_enabled_nao_booleano(): + """ + enabled que nao e bool deve retornar 400 (antes de tocar no banco). + """ + service = _make_service(is_admin=True) + + with pytest.raises(SynapseError) as err: + await service.set_feature( + requester=_make_requester(ADMIN), + room_id=ROOM_ID, + feature="vods", + enabled="sim", + ) + + assert err.value.code == 400 + + +@pytest.mark.asyncio +async def test_set_feature_sem_room_id(): + """ + room_id vazio deve retornar 400. + """ + service = _make_service(is_admin=True) + + with pytest.raises(SynapseError) as err: + await service.set_feature( + requester=_make_requester(ADMIN), + room_id="", + feature="vods", + enabled=True, + ) + + assert err.value.code == 400 + + +# ---------------- DB ---------------- + +def test_db_get_feature_query(): + """ + get_feature deve filtrar por room_id e feature. + """ + txn = _make_txn(fetchone=(True,)) + + db.get_feature(txn, room_id=ROOM_ID, feature="vods") + + sql = txn.execute.call_args.args[0] + params = txn.execute.call_args.args[1] + assert "WHERE room_id = ? AND feature = ?" in sql + assert params == (ROOM_ID, "vods") + + +def test_db_set_feature_upsert(): + """ + set_feature deve usar upsert (ON CONFLICT). + """ + txn = _make_txn() + + db.set_feature(txn, room_id=ROOM_ID, feature="vods", enabled=True, now_ms=123) + + sql = txn.execute.call_args.args[0] + params = txn.execute.call_args.args[1] + assert "ON CONFLICT" in sql + assert params == (ROOM_ID, "vods", True, 123, 123) diff --git a/modules/room_service/resources/__init__.py b/modules/room_service/resources/__init__.py new file mode 100644 index 00000000000..e69de29bb2d diff --git a/modules/room_service/service.py b/modules/room_service/service.py index 8816f4db864..705e548d16f 100644 --- a/modules/room_service/service.py +++ b/modules/room_service/service.py @@ -767,6 +767,9 @@ async def _admin_join(self, room_id: str, user_id: str): payload = json.dumps({"user_id": user_id}).encode() + logger.info(f"ADMIN TOKEN RAW: {self.admin_token}") + logger.info(f"AUTH HEADER: Bearer {self.admin_token}") + response = await self.agent.request( b"POST", url.encode(), diff --git a/modules/schedule_service/README b/modules/schedule_service/README new file mode 100644 index 00000000000..c6ced14223d --- /dev/null +++ b/modules/schedule_service/README @@ -0,0 +1,204 @@ +# Schedule Service Module for Matrix Synapse + +This module exposes **HTTP endpoints** to show upcoming events for a room, +backed by **Google Calendar**. Each room points to its own calendar, and the +module fetches that calendar's events server-side. + +The Google API key stays on the server (in `homeserver.yaml`) and is never +sent to the frontend: the frontend asks Synapse for a room's events, and +Synapse calls Google. + +--- + +## Features + +- per-room events: each room maps to a Google Calendar (`room_calendars` table) +- fetch upcoming events for a room from Google Calendar +- admins configure which calendar a room uses +- the Google API key never leaves the server + +--- + +## How it works + +1. Synapse loads the module on startup +2. The module initializes a `ScheduleService` +3. Custom HTTP endpoints are registered under: + /\_synapse/schedule_service/\* +4. Callers use these endpoints: + - an admin sets which calendar a room uses (set_calendar) + - the frontend reads the room's calendar id (get_calendar) and its upcoming + events (list_events) +5. For events, the module calls the Google Calendar API server-side using the + configured API key, filters the result, and returns JSON + +The module stores only the room → calendar mapping. The events themselves live +in Google Calendar, not in the database. + +--- + +## Registered endpoints + +All endpoints are **POST** with a JSON body. + +- /\_synapse/schedule_service/list_events +- /\_synapse/schedule_service/get_calendar +- /\_synapse/schedule_service/set_calendar + +--- + +### List events + +POST /list_events + +Returns the upcoming events for a room's calendar (from today onward, ordered +by start time). A room with no calendar configured returns an empty list — not +an error. + +**Body:** + +- `room_id` (required) +- `limit` (optional, default 6, max 50) +- `time_min` (optional ISO timestamp; defaults to the start of today, UTC) + +**Response:** + +- `roomId` +- `items` — array of events, each with `id`, `status`, `summary`, + `description`, `location`, `start`, `end`, `htmlLink` + +Cancelled events and events without a start are filtered out. + +**Errors:** + +- `400` – room_id missing or invalid limit +- `502` – the Google Calendar request failed + +--- + +### Get calendar + +POST /get_calendar + +Returns which calendar a room is configured to use. Cheap (no Google call) — +meant to pre-fill the admin form. + +**Body:** + +- `room_id` (required) + +**Response:** + +- `roomId` +- `calendarId` – the configured calendar id, or `null` if none + +**Errors:** + +- `400` – room_id missing + +--- + +### Set calendar (admin only) + +POST /set_calendar + +Sets which Google Calendar a room uses. Requires a **global Synapse admin** +(the requester is identified from the access token; non-admins get `403`). + +Upsert: creates the mapping if it doesn't exist, updates it otherwise. + +**Body:** + +- `room_id` (required) +- `calendar_id` (required) – e.g. `abc123@group.calendar.google.com` + +**Response:** + +- `roomId` +- `calendarId` + +**Errors:** + +- `400` – room_id or calendar_id missing +- `403` – requester is not a global admin + +--- + +## Google setup + +The module needs a Google API key with the **Google Calendar API enabled** on +the project: + +1. Google Cloud Console → APIs & Services → Library → enable "Google Calendar API" +2. Credentials → create an API key +3. Restrict the key: Application restrictions → IP addresses (the server's IP) + in production; API restrictions → restrict to Google Calendar API only + +The calendar a room reads must be accessible to that key — a public calendar, +or one shared appropriately. For a quick test, the public Brazilian holidays +calendar works: `pt.brazilian#holiday@group.v.calendar.google.com`. + +--- + +## Database + +Table `room_calendars`: + +```sql +CREATE TABLE public.room_calendars ( + room_id text NOT NULL, + calendar_id text NOT NULL, + created_at int8 NOT NULL, + updated_at int8 NOT NULL, + PRIMARY KEY (room_id), + FOREIGN KEY (room_id) + REFERENCES public.room_business(room_id) + ON DELETE CASCADE +); +``` + +The Synapse database user needs privileges on this table: + +```sql +ALTER TABLE room_calendars OWNER TO ; +``` + +Notes: + +- one row per room (`room_id` is the primary key). +- the FK to `room_business` means a room must exist there before it can get a + calendar. Timestamps are epoch milliseconds (`int8`), matching the other + modules' tables. + +--- + +## Configuration + +Example `homeserver.yaml` configuration: + +```yaml +modules: + - module: modules.schedule_service.module.ScheduleServiceModule + config: + google_api_key_file: "/path/to/google_api_key.txt" + timezone: "America/Sao_Paulo" +``` + +- `google_api_key_file` (or `google_api_key` inline) is **required** — the + module fails to load without it. +- `timezone` is optional (defaults to `America/Sao_Paulo`). + +Keep the API key file out of version control (add it to `.gitignore`). + +--- + +## Testing + +``` +PYTHONPATH=. pytest -vv modules/schedule_service/tests/ +``` + +Covers: room without a configured calendar (returns empty, no Google call), +the admin check on writes (403 for non-admins), the event filtering +(cancelled/no-start dropped), limit validation, and a Google failure surfacing +as 502. \ No newline at end of file diff --git a/modules/schedule_service/__init__.py b/modules/schedule_service/__init__.py new file mode 100644 index 00000000000..e69de29bb2d diff --git a/modules/schedule_service/db.py b/modules/schedule_service/db.py new file mode 100644 index 00000000000..d02cb425920 --- /dev/null +++ b/modules/schedule_service/db.py @@ -0,0 +1,36 @@ +def get_calendar(txn, room_id: str): + txn.execute( + """ + SELECT calendar_id + FROM room_calendars + WHERE room_id = ? + """, + (room_id,), + ) + return txn.fetchone() + + +def set_calendar(txn, room_id: str, calendar_id: str, now_ms: int): + """ + Upsert: define qual Google Calendar a sala usa. + """ + txn.execute( + """ + INSERT INTO room_calendars (room_id, calendar_id, created_at, updated_at) + VALUES (?, ?, ?, ?) + ON CONFLICT (room_id) + DO UPDATE SET calendar_id = EXCLUDED.calendar_id, + updated_at = EXCLUDED.updated_at + """, + (room_id, calendar_id, now_ms, now_ms), + ) + + +def delete_calendar(txn, room_id: str): + txn.execute( + """ + DELETE FROM room_calendars + WHERE room_id = ? + """, + (room_id,), + ) \ No newline at end of file diff --git a/modules/schedule_service/module.py b/modules/schedule_service/module.py new file mode 100644 index 00000000000..e667ca3fb36 --- /dev/null +++ b/modules/schedule_service/module.py @@ -0,0 +1,59 @@ +# module.py +import logging +import os + +from synapse.module_api import ModuleApi + +from .service import ScheduleService +from .resources.list_events import ListEventsResource +from .resources.get_calendar import GetCalendarResource +from .resources.set_calendar import SetCalendarResource + +logger = logging.getLogger(__name__) + +def _read_google_api_key(config: dict): + """ + Le a Google API key de arquivo (preferido) ou direto do config. + A key fica no servidor e nunca vai para o front. + """ + key_file = config.get("google_api_key_file") + if key_file and os.path.exists(key_file): + with open(key_file, "r") as f: + return f.read().strip() + return config.get("google_api_key") + + +class ScheduleServiceModule: + def __init__(self, config, api: ModuleApi): + self.hs = api._hs + + google_api_key = _read_google_api_key(config) + if not google_api_key: + raise Exception( + "schedule_service: 'google_api_key_file' or 'google_api_key' is required" + ) + + timezone_name = config.get("timezone", "America/Sao_Paulo") + + service = ScheduleService( + api=api, + google_api_key=google_api_key, + timezone_name=timezone_name, + ) + + self.hs.schedule_service = service + + api.register_web_resource( + "/_synapse/schedule_service/list_events", + ListEventsResource(api, service), + ) + api.register_web_resource( + "/_synapse/schedule_service/get_calendar", + GetCalendarResource(api, service), + ) + api.register_web_resource( + "/_synapse/schedule_service/set_calendar", + SetCalendarResource(api, service), + ) + + logger.info("ScheduleServiceModule carregado") \ No newline at end of file diff --git a/modules/schedule_service/resources/__init__.py b/modules/schedule_service/resources/__init__.py new file mode 100644 index 00000000000..e69de29bb2d diff --git a/modules/schedule_service/resources/get_calendar.py b/modules/schedule_service/resources/get_calendar.py new file mode 100644 index 00000000000..b4e3e3dfe2f --- /dev/null +++ b/modules/schedule_service/resources/get_calendar.py @@ -0,0 +1,23 @@ +# get_calendar.py +from synapse.api.errors import SynapseError +from synapse.http.server import DirectServeJsonResource +from synapse.http.servlet import parse_json_object_from_request + + +class GetCalendarResource(DirectServeJsonResource): + isLeaf = True + + def __init__(self, api, service): + super().__init__() + self.api = api + self.service = service + + async def _async_render_POST(self, request): + content = parse_json_object_from_request(request) + + room_id = content.get("room_id") + if not room_id: + raise SynapseError(400, "room_id is required") + + result = await self.service.get_calendar(room_id=room_id) + return 200, result \ No newline at end of file diff --git a/modules/schedule_service/resources/list_events.py b/modules/schedule_service/resources/list_events.py new file mode 100644 index 00000000000..76a33297421 --- /dev/null +++ b/modules/schedule_service/resources/list_events.py @@ -0,0 +1,30 @@ +# list_events.py +from synapse.api.errors import SynapseError +from synapse.http.server import DirectServeJsonResource +from synapse.http.servlet import parse_json_object_from_request + + +class ListEventsResource(DirectServeJsonResource): + isLeaf = True + + def __init__(self, api, service): + super().__init__() + self.api = api + self.service = service + + async def _async_render_POST(self, request): + content = parse_json_object_from_request(request) + + room_id = content.get("room_id") + if not room_id: + raise SynapseError(400, "room_id is required") + + limit = content.get("limit", 6) + time_min = content.get("time_min") + + result = await self.service.list_events( + room_id=room_id, + limit=limit, + time_min=time_min, + ) + return 200, result \ No newline at end of file diff --git a/modules/schedule_service/resources/set_calendar.py b/modules/schedule_service/resources/set_calendar.py new file mode 100644 index 00000000000..06cefbb8351 --- /dev/null +++ b/modules/schedule_service/resources/set_calendar.py @@ -0,0 +1,34 @@ +# set_calendar.py +from synapse.api.errors import SynapseError +from synapse.http.server import DirectServeJsonResource +from synapse.http.servlet import parse_json_object_from_request + + +class SetCalendarResource(DirectServeJsonResource): + isLeaf = True + + def __init__(self, api, service): + super().__init__() + self.api = api + self.service = service + + async def _async_render_POST(self, request): + # quem chama (para o admin check) - requester inteiro, igual room_features + requester = await self.api.get_user_by_req(request) + + content = parse_json_object_from_request(request) + + room_id = content.get("room_id") + calendar_id = content.get("calendar_id") + + if not room_id: + raise SynapseError(400, "room_id is required") + if not calendar_id: + raise SynapseError(400, "calendar_id is required") + + result = await self.service.set_calendar( + requester=requester, + room_id=room_id, + calendar_id=calendar_id, + ) + return 200, result \ No newline at end of file diff --git a/modules/schedule_service/service.py b/modules/schedule_service/service.py new file mode 100644 index 00000000000..062a8219794 --- /dev/null +++ b/modules/schedule_service/service.py @@ -0,0 +1,157 @@ +# service.py +import logging +import time +from urllib.parse import quote +from datetime import datetime, timezone + +from synapse.module_api import ModuleApi +from synapse.api.errors import SynapseError + +from . import db + +logger = logging.getLogger(__name__) + +GOOGLE_CALENDAR_API = "https://www.googleapis.com/calendar/v3/calendars" + + +class ScheduleService: + """ + Eventos por sala via Google Calendar. + + Cada sala aponta para um calendar_id (tabela room_calendars). O front pede + os eventos de uma sala; o modulo busca no Google server-side (a API key + fica no homeserver.yaml e nunca vai para o front) e devolve o resultado. + + Leitura de eventos: qualquer um pode consultar. + Configurar o calendar da sala: apenas admin global. + """ + + def __init__( + self, + api: ModuleApi, + google_api_key: str, + timezone_name: str = "America/Sao_Paulo", + ): + self.api = api + self.hs = api._hs + self.google_api_key = google_api_key + self.timezone_name = timezone_name + + self.store = self.hs.get_datastores().main + # cliente HTTP do proprio Synapse (nao usamos requests/fetch) + self.http_client = self.hs.get_proxied_http_client() + + # ---------------- ADMIN ---------------- + async def assert_is_admin(self, user_id: str) -> None: + is_admin = await self.api.is_user_admin(user_id) + if not is_admin: + raise SynapseError(403, "Only Synapse admins can perform this action") + + # ---------------- CALENDAR CONFIG ---------------- + async def get_calendar(self, room_id: str) -> dict: + if not room_id: + raise SynapseError(400, "room_id is required") + + row = await self.store.db_pool.runInteraction( + "get_room_calendar", + db.get_calendar, + room_id, + ) + + return { + "roomId": room_id, + "calendarId": row[0] if row else None, + } + + async def set_calendar(self, *, requester, room_id: str, calendar_id: str) -> dict: + user_id = requester.user.to_string() + logger.info( + "set_calendar: requester=%s room_id=%s calendar_id=%s", + user_id, + room_id, + calendar_id, + ) + + if not room_id: + raise SynapseError(400, "room_id is required") + if not calendar_id: + raise SynapseError(400, "calendar_id is required") + + await self.assert_is_admin(user_id) + + now_ms = int(time.time() * 1000) + await self.store.db_pool.runInteraction( + "set_room_calendar", + db.set_calendar, + room_id, + calendar_id, + now_ms, + ) + + return {"roomId": room_id, "calendarId": calendar_id} + + # ---------------- EVENTS (Google) ---------------- + def _today_start_utc_iso(self) -> str: + """ + Inicio do dia de hoje em UTC ISO. Porte do luxon startOf('day').toUTC() + da POC, simplificado para meia-noite UTC. + """ + now = datetime.now(timezone.utc) + start = now.replace(hour=0, minute=0, second=0, microsecond=0) + return start.isoformat().replace("+00:00", "Z") + + async def list_events(self, room_id: str, limit: int = 6, time_min: str = None) -> dict: + logger.info("list_events: room_id=%s limit=%s", room_id, limit) + + if not room_id: + raise SynapseError(400, "room_id is required") + if limit < 1 or limit > 50: + raise SynapseError(400, "limit must be between 1 and 50") + + # descobre o calendar da sala + row = await self.store.db_pool.runInteraction( + "get_room_calendar", + db.get_calendar, + room_id, + ) + if not row or not row[0]: + # sala sem calendar configurado: lista vazia, nao e erro + return {"roomId": room_id, "items": []} + + calendar_id = row[0] + effective_time_min = time_min or self._today_start_utc_iso() + + url = ( + f"{GOOGLE_CALENDAR_API}/{quote(calendar_id, safe='')}/events" + f"?singleEvents=true&orderBy=startTime" + f"&timeMin={quote(effective_time_min, safe='')}" + f"&maxResults={limit}" + f"&key={self.google_api_key}" + ) + + try: + data = await self.http_client.get_json(url) + except Exception as e: + logger.exception("list_events: falha ao chamar Google Calendar") + raise SynapseError(502, f"Google Calendar request failed: {e}") + + # filtra cancelados e sem start, igual a POC + items = [ + self._serialize_event(e) + for e in (data.get("items") or []) + if e.get("status") != "cancelled" and e.get("start") + ] + + return {"roomId": room_id, "items": items} + + def _serialize_event(self, e: dict) -> dict: + return { + "id": e.get("id"), + "status": e.get("status"), + "summary": e.get("summary"), + "description": e.get("description"), + "location": e.get("location"), + "start": e.get("start"), + "end": e.get("end"), + "htmlLink": e.get("htmlLink"), + } \ No newline at end of file diff --git a/modules/schedule_service/tests/__init__.py b/modules/schedule_service/tests/__init__.py new file mode 100644 index 00000000000..e69de29bb2d diff --git a/modules/schedule_service/tests/test_schedule_service.py b/modules/schedule_service/tests/test_schedule_service.py new file mode 100644 index 00000000000..eca96bcb262 --- /dev/null +++ b/modules/schedule_service/tests/test_schedule_service.py @@ -0,0 +1,176 @@ +import pytest +from unittest.mock import MagicMock, AsyncMock + +from synapse.api.errors import SynapseError +from modules.schedule_service.service import ScheduleService +from modules.schedule_service import db + +# testa +# 1. get_calendar de sala sem config retorna None +# 2. set_calendar exige admin +# 3. list_events de sala sem calendar retorna vazio (nao erro) +# 4. list_events chama o Google e filtra cancelados/sem start +# 5. validacao de room_id / calendar_id / limit + +# para testar: PYTHONPATH=. pytest -vv modules/schedule_service/tests/ + +ROOM_ID = "!room:localhost" +ADMIN = "@admin:localhost" +USER = "@user:localhost" + + +def _make_service(calendar_row=None, google_response=None, is_admin=True): + api = MagicMock() + hs = MagicMock() + api._hs = hs + api.is_user_admin = AsyncMock(return_value=is_admin) + + store = MagicMock() + + async def fake_run(desc, func, *args): + if desc == "get_room_calendar": + return calendar_row + if desc == "set_room_calendar": + return None + raise AssertionError(f"runInteraction inesperado: {desc}") + + store.db_pool.runInteraction = AsyncMock(side_effect=fake_run) + hs.get_datastores.return_value.main = store + + http = MagicMock() + http.get_json = AsyncMock(return_value=google_response or {"items": []}) + hs.get_proxied_http_client.return_value = http + + service = ScheduleService(api=api, google_api_key="KEY") + service._store_mock = store + service._http_mock = http + return service + + +def _make_txn(fetchone=None): + txn = MagicMock() + txn.fetchone.return_value = fetchone + return txn + + +def _requester(user_id): + r = MagicMock() + r.user.to_string.return_value = user_id + return r + + +# ---------------- GET CALENDAR ---------------- + +@pytest.mark.asyncio +async def test_get_calendar_sem_config(): + service = _make_service(calendar_row=None) + result = await service.get_calendar(room_id=ROOM_ID) + assert result["calendarId"] is None + + +@pytest.mark.asyncio +async def test_get_calendar_configurado(): + service = _make_service(calendar_row=("cal123@group.calendar.google.com",)) + result = await service.get_calendar(room_id=ROOM_ID) + assert result["calendarId"] == "cal123@group.calendar.google.com" + + +# ---------------- SET CALENDAR ---------------- + +@pytest.mark.asyncio +async def test_set_calendar_exige_admin(): + service = _make_service(is_admin=False) + with pytest.raises(SynapseError) as err: + await service.set_calendar( + requester=_requester(USER), room_id=ROOM_ID, calendar_id="cal" + ) + assert err.value.code == 403 + + +@pytest.mark.asyncio +async def test_set_calendar_admin_grava(): + service = _make_service(is_admin=True) + result = await service.set_calendar( + requester=_requester(ADMIN), room_id=ROOM_ID, calendar_id="cal123" + ) + assert result["calendarId"] == "cal123" + call = service._store_mock.db_pool.runInteraction.await_args + assert call.args[0] == "set_room_calendar" + assert call.args[2] == ROOM_ID + assert call.args[3] == "cal123" + + +@pytest.mark.asyncio +async def test_set_calendar_sem_id(): + service = _make_service(is_admin=True) + with pytest.raises(SynapseError) as err: + await service.set_calendar( + requester=_requester(ADMIN), room_id=ROOM_ID, calendar_id="" + ) + assert err.value.code == 400 + + +# ---------------- LIST EVENTS ---------------- + +@pytest.mark.asyncio +async def test_list_events_sem_calendar(): + """Sala sem calendar configurado retorna lista vazia, nao erro.""" + service = _make_service(calendar_row=None) + result = await service.list_events(room_id=ROOM_ID) + assert result["items"] == [] + # nao deve ter chamado o Google + service._http_mock.get_json.assert_not_awaited() + + +@pytest.mark.asyncio +async def test_list_events_chama_google_e_filtra(): + google = { + "items": [ + {"id": "1", "status": "confirmed", "summary": "Live", "start": {"dateTime": "2026-08-20T20:00:00Z"}}, + {"id": "2", "status": "cancelled", "summary": "Cancelado", "start": {"dateTime": "2026-08-21T20:00:00Z"}}, + {"id": "3", "status": "confirmed", "summary": "Sem start"}, # sem start -> filtrado + ] + } + service = _make_service(calendar_row=("cal123",), google_response=google) + + result = await service.list_events(room_id=ROOM_ID) + + # so o evento 1 sobra (2 cancelado, 3 sem start) + assert len(result["items"]) == 1 + assert result["items"][0]["id"] == "1" + service._http_mock.get_json.assert_awaited_once() + + +@pytest.mark.asyncio +async def test_list_events_limit_invalido(): + service = _make_service(calendar_row=("cal123",)) + with pytest.raises(SynapseError) as err: + await service.list_events(room_id=ROOM_ID, limit=999) + assert err.value.code == 400 + + +@pytest.mark.asyncio +async def test_list_events_google_falha(): + service = _make_service(calendar_row=("cal123",)) + service._http_mock.get_json = AsyncMock(side_effect=Exception("timeout")) + with pytest.raises(SynapseError) as err: + await service.list_events(room_id=ROOM_ID) + assert err.value.code == 502 + + +# ---------------- DB ---------------- + +def test_db_set_calendar_upsert(): + txn = _make_txn() + db.set_calendar(txn, room_id=ROOM_ID, calendar_id="cal", now_ms=123) + sql = txn.execute.call_args.args[0] + params = txn.execute.call_args.args[1] + assert "ON CONFLICT" in sql + assert params == (ROOM_ID, "cal", 123, 123) + + +def test_db_get_calendar_query(): + txn = _make_txn(fetchone=("cal",)) + db.get_calendar(txn, room_id=ROOM_ID) + params = txn.execute.call_args.args[1] + assert params == (ROOM_ID,) \ No newline at end of file diff --git a/modules/vod_service/README b/modules/vod_service/README new file mode 100644 index 00000000000..0605ed67e24 --- /dev/null +++ b/modules/vod_service/README @@ -0,0 +1,201 @@ +# VOD Service Module for Matrix Synapse + +This module exposes **HTTP endpoints** to create, list and retrieve recorded +streams (VODs), returning ready-to-play playback URLs. + +It replaces the AdonisJS `Stream` model + `/api/streams/vods` route. +The URL building logic that used to live in the model's `@computed()` getters now lives in `VodService`, pointing to **Oracle Object Storage** instead of AWS CloudFront. + +--- + +## Features + +- expose custom REST endpoints for VOD listing and retrieval +- register new VODs from a recording pipeline (service-token protected) +- paginate recorded streams by room +- build HLS master playlist URLs (Oracle Object Storage native URLs) +- build thumbnail URLs (including the "middle of the recording" thumbnail) +- serve VOD metadata scoped to a room + +--- + +## How it works + +1. Synapse loads the module on startup +2. The module initializes a `VodService` +3. Custom HTTP endpoints are registered under: + /\_synapse/vod_service/\* +4. Callers use these endpoints: + - the frontend reads VODs for a room (list/get) + - a recording pipeline registers finished recordings (create) +5. The module: + - reads/writes stream metadata in the database + - builds Object Storage URLs from `recording_path` + - returns JSON + +The module **does not** record or store video. Files are produced elsewhere (streaming / recording pipeline) and uploaded to Object Storage. This module only stores the metadata and serves the resulting URLs. + +--- + +## Registered endpoints + +- /\_synapse/vod_service/list +- /\_synapse/vod_service/get +- /\_synapse/vod_service/create + +--- + +### List VODs + +POST /list + +Returns a paginated list of recorded streams (ended_at IS NOT NULL) for a room. + +**Body:** + +- `room_id` (required) – Matrix room ID +- `page` (default 1) +- `limit` (default 10, max 100) + +**Response fields (per item):** + +- `id` +- `streamId` +- `roomId` +- `title` +- `categoryId` +- `recordingPath` +- `recordingDurationMs` +- `startedAt` +- `endedAt` +- `masterPlaylistUrl` +- `thumbnailBaseUrl` +- `latestThumbnail` +- `isLive` +- `isVod` + +**Meta fields:** + +- `total` +- `page` +- `perPage` +- `lastPage` + +**Errors:** + +- `400` – room_id is required / invalid page or limit + +--- + +### Get VOD + +POST /get + +Returns a single VOD by id. + +**Body:** + +- `stream_id` (required) + +**Errors:** + +- `400` – stream_id is required +- `404` – VOD not found + +--- + +### Create VOD + +POST /create + +Registers a finished recording as a new VOD. Intended to be called by a +recording pipeline (Jibri `finalize.sh`, an OBS post-hook, a streaming server, +etc.) — **not** by a logged-in user. The caller authenticates with a shared +**service token**, not a Matrix access token. + +**Auth:** + +- Header `Authorization: Bearer ` (see Configuration) + +**Body:** + +- `room_id` (required) – the room the VOD belongs to +- `recording_path` (required) – the object key prefix in the bucket +- `title` (optional) +- `started_at` (optional, epoch ms) +- `ended_at` (optional, epoch ms) +- `recording_duration_ms` (optional) + +The `stream_id` is generated server-side; callers don't provide it. + +**Response:** the created VOD, in the same shape as list/get items. + +**Errors:** + +- `400` – room_id or recording_path missing +- `403` – missing or invalid service token +- `503` – VOD creation not configured (no service token set) + +--- + +## URL format + +Playback URLs follow the Oracle Object Storage native format: + +``` +https://objectstorage..oraclecloud.com/n//b//o//media/hls/master.m3u8 +``` + +The `/media/hls/` + `/media/thumbnails/` layout follows the HLS auto-record convention and is kept unchanged. + +The old AdonisJS format was: + +``` +https:////media/hls/master.m3u8 +``` +There is no CDN in front of Object Storage for now, so objects must be publicly readable. If access needs to be restricted, `VodService` must be changed to generate **pre-signed URLs** instead of plain string URLs. This is a real gap for VODs in paid/private rooms: a public bucket serves anyone who has the URL, regardless of room membership. + +--- + +## Notes + +- VODs are scoped by `room_id`: each recording belongs to a Matrix room, and the `streams` table has a `room_id` column (FK to `room_business`). +- Whether the VODs tab is shown for a room is controlled separately, by the `room_features` table (feature `vods`) — not by this module. +- `create` only registers metadata. Uploading the actual files to Object Storage is the pipeline's responsibility, and must happen before (or alongside) the create call, or playback URLs will 404. + +--- + +## Configuration + +Example `homeserver.yaml` configuration: + +```yaml +modules: + - module: modules.vod_service.module.VodServiceModule + config: + object_storage_base_url: "https://objectstorage..oraclecloud.com" + namespace: "" + bucket: "vod-teste" + service_token_file: "/path/to/vod_service_token.txt" +``` + +`service_token_file` (or `service_token` inline) enables the `create` +endpoint. If neither is set, `create` returns `503` and the module still +serves list/get normally. + +Generate a token with: + +``` +openssl rand -hex 32 > /path/to/vod_service_token.txt +``` + +--- + +## Database + +The `streams` table needs read **and write** privileges for the Synapse DB user: + +```sql +GRANT SELECT, INSERT, UPDATE, DELETE ON streams TO ; +GRANT USAGE, SELECT ON SEQUENCE streams_id_seq TO ; +``` \ No newline at end of file diff --git a/modules/vod_service/__init__.py b/modules/vod_service/__init__.py new file mode 100644 index 00000000000..e69de29bb2d diff --git a/modules/vod_service/db.py b/modules/vod_service/db.py new file mode 100644 index 00000000000..81999c9ed84 --- /dev/null +++ b/modules/vod_service/db.py @@ -0,0 +1,86 @@ +def get_vods(txn, room_id: str, limit: int, offset: int): + txn.execute( + """ + SELECT id, stream_id, room_id, title, category_id, + recording_path, recording_duration_ms, + started_at, ended_at + FROM streams + WHERE room_id = ? + AND ended_at IS NOT NULL + ORDER BY started_at DESC + LIMIT ? OFFSET ? + """, + (room_id, limit, offset), + ) + return txn.fetchall() + + +def count_vods(txn, room_id: str): + txn.execute( + """ + SELECT COUNT(*) + FROM streams + WHERE room_id = ? + AND ended_at IS NOT NULL + """, + (room_id,), + ) + row = txn.fetchone() + return row[0] if row else 0 + + +def get_vod_by_id(txn, stream_id: int): + txn.execute( + """ + SELECT id, stream_id, room_id, title, category_id, + recording_path, recording_duration_ms, + started_at, ended_at + FROM streams + WHERE id = ? + """, + (stream_id,), + ) + return txn.fetchone() + + +def insert_vod( + txn, + stream_id: str, + room_id: str, + title, + recording_path: str, + recording_duration_ms, + started_at, + ended_at, + created_at: int, + updated_at: int, +): + """ + Insere um novo VOD (gravacao finalizada) e devolve a linha criada + no mesmo formato que os SELECTs, para o service serializar. + """ + txn.execute( + """ + INSERT INTO streams ( + stream_id, room_id, title, recording_path, + recording_duration_ms, started_at, ended_at, + created_at, updated_at + ) + VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?) + RETURNING id, stream_id, room_id, title, category_id, + recording_path, recording_duration_ms, + started_at, ended_at + """, + ( + stream_id, + room_id, + title, + recording_path, + recording_duration_ms, + started_at, + ended_at, + created_at, + updated_at, + ), + ) + return txn.fetchone() \ No newline at end of file diff --git a/modules/vod_service/module.py b/modules/vod_service/module.py new file mode 100644 index 00000000000..7189f8c1ae5 --- /dev/null +++ b/modules/vod_service/module.py @@ -0,0 +1,67 @@ +# module.py +import logging +import os + +from synapse.module_api import ModuleApi + +from .service import VodService +from .resources.list_vods import ListVodsResource +from .resources.get_vod import GetVodResource +from .resources.create_vod import CreateVodResource + +logger = logging.getLogger(__name__) + + +def _read_service_token(config: dict): + """ + Le o token de servico usado para criar VODs. + Aceita service_token_file (preferido) ou service_token direto. + Se nenhum for configurado, a criacao de VOD fica desabilitada (503). + """ + token_file = config.get("service_token_file") + if token_file and os.path.exists(token_file): + with open(token_file, "r") as f: + return f.read().strip() + return config.get("service_token") + + +class VodServiceModule: + def __init__(self, config, api: ModuleApi): + self.hs = api._hs + + object_storage_base_url = config["object_storage_base_url"] + namespace = config["namespace"] + bucket = config["bucket"] + service_token = _read_service_token(config) + + service = VodService( + api=api, + object_storage_base_url=object_storage_base_url, + namespace=namespace, + bucket=bucket, + service_token=service_token, + ) + + self.hs.vod_service = service + + api.register_web_resource( + "/_synapse/vod_service/list", + ListVodsResource(api, service), + ) + + api.register_web_resource( + "/_synapse/vod_service/get", + GetVodResource(api, service), + ) + + api.register_web_resource( + "/_synapse/vod_service/create", + CreateVodResource(api, service), + ) + + logger.info( + "VodServiceModule carregado (namespace=%s bucket=%s create=%s)", + namespace, + bucket, + "on" if service_token else "off", + ) \ No newline at end of file diff --git a/modules/vod_service/resources/__init__.py b/modules/vod_service/resources/__init__.py new file mode 100644 index 00000000000..e69de29bb2d diff --git a/modules/vod_service/resources/create_vod.py b/modules/vod_service/resources/create_vod.py new file mode 100644 index 00000000000..0fcbc33c305 --- /dev/null +++ b/modules/vod_service/resources/create_vod.py @@ -0,0 +1,52 @@ +# create_vod.py +from synapse.api.errors import SynapseError +from synapse.http.server import DirectServeJsonResource +from synapse.http.servlet import parse_json_object_from_request + + +def _extract_service_token(request) -> str: + """ + Le o token de servico do header Authorization: Bearer . + Nao usa get_user_by_req porque quem chama e uma maquina (Jibri/OBS), + nao um usuario logado. + """ + auth = request.getHeader(b"Authorization") + if not auth: + return "" + auth = auth.decode("utf-8", "ignore") + if auth.lower().startswith("bearer "): + return auth[7:].strip() + return "" + + +class CreateVodResource(DirectServeJsonResource): + isLeaf = True + + def __init__(self, api, service): + super().__init__() + self.api = api + self.service = service + + async def _async_render_POST(self, request): + service_token = _extract_service_token(request) + + content = parse_json_object_from_request(request) + + room_id = content.get("room_id") + recording_path = content.get("recording_path") + + if not room_id: + raise SynapseError(400, "room_id is required") + if not recording_path: + raise SynapseError(400, "recording_path is required") + + result = await self.service.create_vod( + service_token=service_token, + room_id=room_id, + recording_path=recording_path, + title=content.get("title"), + started_at=content.get("started_at"), + ended_at=content.get("ended_at"), + recording_duration_ms=content.get("recording_duration_ms"), + ) + return 200, result \ No newline at end of file diff --git a/modules/vod_service/resources/get_vod.py b/modules/vod_service/resources/get_vod.py new file mode 100644 index 00000000000..c17b30c2afc --- /dev/null +++ b/modules/vod_service/resources/get_vod.py @@ -0,0 +1,23 @@ +# get_vod.py +from synapse.api.errors import SynapseError +from synapse.http.server import DirectServeJsonResource +from synapse.http.servlet import parse_json_object_from_request + + +class GetVodResource(DirectServeJsonResource): + isLeaf = True + + def __init__(self, api, service): + super().__init__() + self.api = api + self.service = service + + async def _async_render_POST(self, request): + content = parse_json_object_from_request(request) + stream_id = content.get("stream_id") + + if not stream_id: + raise SynapseError(400, "stream_id is required") + + result = await self.service.get_vod(stream_id=stream_id) + return 200, result diff --git a/modules/vod_service/resources/list_vods.py b/modules/vod_service/resources/list_vods.py new file mode 100644 index 00000000000..cd2fc1a9b04 --- /dev/null +++ b/modules/vod_service/resources/list_vods.py @@ -0,0 +1,23 @@ +from synapse.http.server import DirectServeJsonResource +from synapse.http.servlet import parse_integer, parse_string + + +class ListVodsResource(DirectServeJsonResource): + isLeaf = True + + def __init__(self, api, service): + super().__init__() + self.api = api + self.service = service + + async def _async_render_GET(self, request): + room_id = parse_string(request, "room_id", required=True) + page = parse_integer(request, "page", default=1) + limit = parse_integer(request, "limit", default=10) + + result = await self.service.list_vods( + room_id=room_id, + page=page, + limit=limit, + ) + return 200, result \ No newline at end of file diff --git a/modules/vod_service/service.py b/modules/vod_service/service.py new file mode 100644 index 00000000000..244234cbc22 --- /dev/null +++ b/modules/vod_service/service.py @@ -0,0 +1,233 @@ +# service.py +import logging +import time +import uuid + +from synapse.module_api import ModuleApi +from synapse.api.errors import SynapseError + +from . import db + +logger = logging.getLogger(__name__) + + +class VodService: + def __init__( + self, + api: ModuleApi, + object_storage_base_url: str, + namespace: str, + bucket: str, + service_token: str = None, + ): + self.api = api + self.hs = api._hs + self.object_storage_base_url = object_storage_base_url.rstrip("/") + self.namespace = namespace + self.bucket = bucket + # segredo compartilhado com quem cria VODs (finalize.sh do Jibri, OBS, etc). + # Nao e um usuario logado - e maquina chamando maquina. + self.service_token = service_token + + self.store = self.hs.get_datastores().main + + # ---------------- URL BUILDERS ---------------- + def _object_prefix(self, recording_path: str) -> str: + """ + Monta o prefixo da URL nativa do Oracle Object Storage. + + Formato: + https://objectstorage..oraclecloud.com + /n//b//o/ + """ + return ( + f"{self.object_storage_base_url}" + f"/n/{self.namespace}" + f"/b/{self.bucket}" + f"/o/{recording_path}" + ) + + def master_playlist_url(self, recording_path: str) -> str: + return f"{self._object_prefix(recording_path)}/media/hls/master.m3u8" + + def thumbnail_base_url(self, recording_path: str) -> str: + return f"{self._object_prefix(recording_path)}/media/latest_thumbnail" + + def thumbnail_url(self, recording_path: str, thumb_num: int) -> str: + return ( + f"{self._object_prefix(recording_path)}" + f"/media/thumbnails/thumb{thumb_num}.jpg" + ) + + def latest_thumbnail( + self, + recording_path: str, + started_at, + ended_at, + ) -> str: + """ + Pega o thumbnail do meio da gravacao (4 thumbs por minuto). + """ + if started_at is None or ended_at is None: + return self.thumbnail_url(recording_path, 0) + + try: + total_minutes = (ended_at - started_at) / 60000.0 + middle_minute = int(total_minutes // 2) + thumb_num = middle_minute * 4 + + if thumb_num < 0: + return self.thumbnail_url(recording_path, 0) + + return self.thumbnail_url(recording_path, thumb_num) + except Exception: + logger.exception( + "latest_thumbnail: falha ao calcular para recording_path=%s", + recording_path, + ) + return self.thumbnail_url(recording_path, 0) + + # ---------------- SERIALIZER ---------------- + def _serialize(self, row) -> dict: + ( + row_id, + stream_id, + room_id, + title, + category_id, + recording_path, + recording_duration_ms, + started_at, + ended_at, + ) = row + + return { + "id": row_id, + "streamId": stream_id, + "roomId": room_id, + "title": title, + "categoryId": category_id, + "recordingPath": recording_path, + "recordingDurationMs": recording_duration_ms, + "startedAt": started_at, + "endedAt": ended_at, + "masterPlaylistUrl": self.master_playlist_url(recording_path), + "thumbnailBaseUrl": self.thumbnail_base_url(recording_path), + "latestThumbnail": self.latest_thumbnail( + recording_path, started_at, ended_at + ), + "isLive": ended_at is None, + "isVod": ended_at is not None, + } + + # ---------------- LIST ---------------- + async def list_vods(self, room_id: str, page: int = 1, limit: int = 10): + logger.info( + "list_vods: room_id=%s page=%s limit=%s", room_id, page, limit + ) + + if not room_id: + raise SynapseError(400, "room_id is required") + if page < 1: + raise SynapseError(400, "page must be >= 1") + if limit < 1 or limit > 100: + raise SynapseError(400, "limit must be between 1 and 100") + + offset = (page - 1) * limit + + rows = await self.store.db_pool.runInteraction( + "get_vods", + db.get_vods, + room_id, + limit, + offset, + ) + + total = await self.store.db_pool.runInteraction( + "count_vods", + db.count_vods, + room_id, + ) + + logger.info("list_vods: %d vods retornados (total=%d)", len(rows), total) + + return { + "data": [self._serialize(row) for row in rows], + "meta": { + "total": total, + "page": page, + "perPage": limit, + "lastPage": (total + limit - 1) // limit if limit else 1, + }, + } + + # ---------------- GET ---------------- + async def get_vod(self, stream_id: int): + logger.info("get_vod: stream_id=%s", stream_id) + + row = await self.store.db_pool.runInteraction( + "get_vod_by_id", + db.get_vod_by_id, + stream_id, + ) + + if not row: + raise SynapseError(404, "VOD not found") + + return self._serialize(row) + + # ---------------- CREATE (escrita) ---------------- + def assert_service_token(self, token: str) -> None: + """ + Autentica quem cria VODs. Nao e usuario logado - e um servico + (finalize.sh do Jibri, OBS, pipeline) mandando um segredo compartilhado. + """ + if not self.service_token: + raise SynapseError(503, "VOD creation is not configured (no service token)") + if not token or token != self.service_token: + raise SynapseError(403, "Invalid service token") + + async def create_vod( + self, + *, + service_token: str, + room_id: str, + recording_path: str, + title=None, + started_at=None, + ended_at=None, + recording_duration_ms=None, + ) -> dict: + self.assert_service_token(service_token) + + if not room_id: + raise SynapseError(400, "room_id is required") + if not recording_path: + raise SynapseError(400, "recording_path is required") + + now_ms = int(time.time() * 1000) + # o backend gera o stream_id, o caller nao precisa se preocupar + stream_id = f"vod_{uuid.uuid4().hex}" + + logger.info( + "create_vod: room_id=%s recording_path=%s stream_id=%s", + room_id, + recording_path, + stream_id, + ) + + row = await self.store.db_pool.runInteraction( + "insert_vod", + db.insert_vod, + stream_id, + room_id, + title, + recording_path, + recording_duration_ms, + started_at, + ended_at, + now_ms, + now_ms, + ) + + return self._serialize(row) \ No newline at end of file diff --git a/modules/vod_service/tests/__init__.py b/modules/vod_service/tests/__init__.py new file mode 100644 index 00000000000..e69de29bb2d diff --git a/modules/vod_service/tests/test_list_vods.py b/modules/vod_service/tests/test_list_vods.py new file mode 100644 index 00000000000..c09f0af48c5 --- /dev/null +++ b/modules/vod_service/tests/test_list_vods.py @@ -0,0 +1,279 @@ +import json +import pytest +from unittest.mock import AsyncMock, MagicMock +from io import BytesIO + +from twisted.web.test.requesthelper import DummyRequest +from twisted.web.http_headers import Headers + +from synapse.api.errors import SynapseError +from modules.vod_service.resources.list_vods import ListVodsResource +from modules.vod_service.resources.get_vod import GetVodResource +from modules.vod_service.module import VodServiceModule + +# testa +# 1. endpoints respondem 200 no formato esperado +# 2. query params sao parseados e repassados ao service +# 3. defaults aplicados quando os params nao vem +# 4. stream_id ausente retorna 400 +# 5. rotas registradas e config obrigatoria validada no boot + +# para testar: PYTHONPATH=. pytest -vv modules/vod_service/tests/test_list_vods.py + + +# ---------------- HELPERS ---------------- + +def _make_list_resource(return_value=None): + api = MagicMock() + service = MagicMock() + service.list_vods = AsyncMock( + return_value=return_value + if return_value is not None + else {"data": [], "meta": {"total": 0, "page": 1, "perPage": 10, "lastPage": 0}} + ) + return ListVodsResource(api, service), service + + +def _make_list_request(args=None): + request = DummyRequest([b"_synapse", b"vod_service", b"list"]) + request.args = args or {} + return request + + +def _make_get_resource(return_value=None): + api = MagicMock() + service = MagicMock() + service.get_vod = AsyncMock(return_value=return_value or {}) + return GetVodResource(api, service), service + + +def _make_get_request(body: dict): + request = DummyRequest([b"_synapse", b"vod_service", b"get"]) + request.method = b"POST" + request.requestHeaders = Headers({ + b"Content-Type": [b"application/json"], + }) + request.content = BytesIO(json.dumps(body).encode()) + return request + + +def _make_config(**overrides): + config = { + "object_storage_base_url": "https://objectstorage.sa-saopaulo-1.oraclecloud.com", + "namespace": "grmsqxilk3cb", + "bucket": "vod-teste", + } + config.update(overrides) + return config + + +def _make_module_api(): + api = MagicMock() + api._hs = MagicMock() + return api + + +# ---------------- LIST VODS RESOURCE ---------------- + +@pytest.mark.asyncio +async def test_list_vods_success(): + """ + Deve retornar 200 e a lista de vods com as URLs montadas. + """ + resource, service = _make_list_resource({ + "data": [ + { + "id": 1, + "title": "VOD de teste", + "recordingPath": "abc123", + "masterPlaylistUrl": ( + "https://objectstorage.sa-saopaulo-1.oraclecloud.com" + "/n/grmsqxilk3cb/b/vod-teste/o/abc123/media/hls/master.m3u8" + ), + "isVod": True, + } + ], + "meta": {"total": 1, "page": 1, "perPage": 10, "lastPage": 1}, + }) + + request = _make_list_request({ + b"room_id": [b"!room:localhost"], + }) + + code, body = await resource._async_render_GET(request) + + assert code == 200 + assert "data" in body + assert len(body["data"]) == 1 + + vod = body["data"][0] + assert vod["recordingPath"] == "abc123" + assert vod["masterPlaylistUrl"].endswith("/media/hls/master.m3u8") + assert "oraclecloud.com" in vod["masterPlaylistUrl"] + + service.list_vods.assert_called_once() + + +@pytest.mark.asyncio +async def test_list_vods_sem_room_id(): + """ + Sem room_id deve retornar 400 (room_id é obrigatório, não tem default). + """ + resource, service = _make_list_resource() + + request = _make_list_request() + + with pytest.raises(SynapseError) as err: + await resource._async_render_GET(request) + + assert err.value.code == 400 + service.list_vods.assert_not_awaited() + + +@pytest.mark.asyncio +async def test_list_vods_repassa_query_params(): + """ + Os query params devem chegar no service (room_id como string). + """ + resource, service = _make_list_resource() + + request = _make_list_request({ + b"room_id": [b"!abc:localhost"], + b"page": [b"3"], + b"limit": [b"25"], + }) + + await resource._async_render_GET(request) + + service.list_vods.assert_awaited_once_with( + room_id="!abc:localhost", + page=3, + limit=25, + ) + + +@pytest.mark.asyncio +async def test_list_vods_retorna_meta_da_paginacao(): + """ + O meta da paginacao deve ser repassado intacto para o front. + """ + resource, _ = _make_list_resource({ + "data": [], + "meta": {"total": 25, "page": 2, "perPage": 10, "lastPage": 3}, + }) + + request = _make_list_request({ + b"room_id": [b"!room:localhost"], + }) + + code, body = await resource._async_render_GET(request) + + assert code == 200 + assert body["meta"]["total"] == 25 + assert body["meta"]["lastPage"] == 3 + + +# ---------------- GET VOD RESOURCE ---------------- + +@pytest.mark.asyncio +async def test_get_vod_success(): + """ + Deve retornar 200 e o vod serializado. + """ + resource, service = _make_get_resource({ + "id": 1, + "title": "VOD de teste", + "recordingPath": "abc123", + "masterPlaylistUrl": ( + "https://objectstorage.sa-saopaulo-1.oraclecloud.com" + "/n/grmsqxilk3cb/b/vod-teste/o/abc123/media/hls/master.m3u8" + ), + "isVod": True, + }) + + request = _make_get_request({"stream_id": 1}) + + code, body = await resource._async_render_POST(request) + + assert code == 200 + assert body["id"] == 1 + assert body["recordingPath"] == "abc123" + + service.get_vod.assert_awaited_once_with(stream_id=1) + + +@pytest.mark.asyncio +async def test_get_vod_sem_stream_id(): + """ + Deve retornar 400 quando stream_id nao vem no body. + """ + resource, service = _make_get_resource() + + request = _make_get_request({}) + + with pytest.raises(SynapseError) as err: + await resource._async_render_POST(request) + + assert err.value.code == 400 + service.get_vod.assert_not_awaited() + + +@pytest.mark.asyncio +async def test_get_vod_stream_id_null(): + """ + stream_id explicitamente null tambem deve retornar 400. + """ + resource, _ = _make_get_resource() + + request = _make_get_request({"stream_id": None}) + + with pytest.raises(SynapseError) as err: + await resource._async_render_POST(request) + + assert err.value.code == 400 + + +# ---------------- MODULE ---------------- + +def test_registra_as_duas_rotas(): + """ + Deve registrar /list e /get sob o namespace do modulo. + """ + api = _make_module_api() + + VodServiceModule(config=_make_config(), api=api) + + paths = [call.args[0] for call in api.register_web_resource.call_args_list] + + assert "/_synapse/vod_service/list" in paths + assert "/_synapse/vod_service/get" in paths + + +def test_expoe_service_no_homeserver(): + """ + O service deve ficar acessivel no hs, igual ao room_service. + """ + api = _make_module_api() + + VodServiceModule(config=_make_config(), api=api) + + assert api._hs.vod_service is not None + assert api._hs.vod_service.bucket == "vod-teste" + assert api._hs.vod_service.namespace == "grmsqxilk3cb" + + +@pytest.mark.parametrize( + "missing_key", + ["object_storage_base_url", "namespace", "bucket"], +) +def test_config_obrigatoria_ausente_falha_no_boot(missing_key): + """ + Config incompleta deve quebrar no startup, nao na primeira request. + """ + api = _make_module_api() + + config = _make_config() + del config[missing_key] + + with pytest.raises(KeyError): + VodServiceModule(config=config, api=api) \ No newline at end of file diff --git a/modules/vod_service/tests/test_vod_service_urls.py b/modules/vod_service/tests/test_vod_service_urls.py new file mode 100644 index 00000000000..7fa4d98ca0a --- /dev/null +++ b/modules/vod_service/tests/test_vod_service_urls.py @@ -0,0 +1,586 @@ +import pytest +from unittest.mock import MagicMock, AsyncMock + +from synapse.api.errors import SynapseError +from modules.vod_service.service import VodService +from modules.vod_service import db + +# testa +# 1. formato da URL nativa do Oracle Object Storage +# 2. porte do @computed() latestThumbnail do model AdonisJS +# 3. validacao de page/limit e calculo da paginacao +# 4. serializacao das linhas do banco nos campos que o front espera +# 5. vod inexistente retorna 404 explicito +# 6. SQL filtra so gravacoes encerradas (ended_at IS NOT NULL) + +# para testar: PYTHONPATH=. pytest -vv modules/vod_service/tests/test_vod_service_urls.py + + +ROOM_ID = "!room:localhost" + + +# ---------------- HELPERS ---------------- + +def _make_service(base_url="https://objectstorage.sa-saopaulo-1.oraclecloud.com"): + api = MagicMock() + return VodService( + api=api, + object_storage_base_url=base_url, + namespace="grmsqxilk3cb", + bucket="vod-teste", + ) + + +def _make_service_with_db(rows=None, total=0): + """Service com o db_pool mockado, para testar list_vods.""" + api = MagicMock() + hs = MagicMock() + api._hs = hs + + store = MagicMock() + + async def fake_run_interaction(desc, func, *args): + if desc == "get_vods": + return rows if rows is not None else [] + if desc == "count_vods": + return total + raise AssertionError(f"runInteraction inesperado: {desc}") + + store.db_pool.runInteraction = AsyncMock(side_effect=fake_run_interaction) + hs.get_datastores.return_value.main = store + + service = VodService( + api=api, + object_storage_base_url="https://objectstorage.sa-saopaulo-1.oraclecloud.com", + namespace="grmsqxilk3cb", + bucket="vod-teste", + ) + service._store_mock = store + return service + + +def _make_service_with_row(row=None): + """Service com o db_pool mockado, para testar get_vod.""" + api = MagicMock() + hs = MagicMock() + api._hs = hs + + store = MagicMock() + store.db_pool.runInteraction = AsyncMock(return_value=row) + hs.get_datastores.return_value.main = store + + service = VodService( + api=api, + object_storage_base_url="https://objectstorage.sa-saopaulo-1.oraclecloud.com", + namespace="grmsqxilk3cb", + bucket="vod-teste", + ) + service._store_mock = store + return service + + +def _row(row_id=1, recording_path="abc123", started_at=0, ended_at=600_000): + return ( + row_id, + f"teste-{row_id}", + ROOM_ID, + "VOD de teste", + None, + recording_path, + 25900, + started_at, + ended_at, + ) + + +def _make_txn(fetchall=None, fetchone=None): + txn = MagicMock() + txn.fetchall.return_value = fetchall if fetchall is not None else [] + txn.fetchone.return_value = fetchone + return txn + + +# ---------------- URL: PLAYLIST ---------------- + +def test_master_playlist_url(): + """ + Deve montar a URL nativa do Oracle Object Storage. + """ + service = _make_service() + url = service.master_playlist_url("abc123") + + assert url == ( + "https://objectstorage.sa-saopaulo-1.oraclecloud.com" + "/n/grmsqxilk3cb/b/vod-teste/o/abc123/media/hls/master.m3u8" + ) + + +def test_base_url_com_barra_no_fim_nao_duplica(): + """ + base_url configurada com barra no fim nao deve gerar '//' na URL. + """ + service = _make_service( + base_url="https://objectstorage.sa-saopaulo-1.oraclecloud.com/" + ) + url = service.master_playlist_url("abc123") + + assert "//n/" not in url + assert url.startswith("https://objectstorage.sa-saopaulo-1.oraclecloud.com/n/") + + +def test_recording_path_aninhado(): + """ + recording_path com barras (formato do IVS) deve ser preservado. + """ + service = _make_service() + url = service.master_playlist_url("canal/2026/07/abc123") + + assert url.endswith("/o/canal/2026/07/abc123/media/hls/master.m3u8") + + +# ---------------- URL: THUMBNAILS ---------------- + +def test_thumbnail_base_url(): + """ + Deve apontar para o latest_thumbnail do IVS. + """ + service = _make_service() + url = service.thumbnail_base_url("abc123") + + assert url.endswith("/o/abc123/media/latest_thumbnail") + + +def test_thumbnail_url_numerada(): + """ + Deve montar a URL de um thumbnail especifico. + """ + service = _make_service() + url = service.thumbnail_url("abc123", 12) + + assert url.endswith("/media/thumbnails/thumb12.jpg") + + +def test_latest_thumbnail_sem_datas(): + """ + Sem started_at/ended_at deve cair no thumb0. + """ + service = _make_service() + url = service.latest_thumbnail("abc123", None, None) + + assert url.endswith("/media/thumbnails/thumb0.jpg") + + +def test_latest_thumbnail_sem_ended_at(): + """ + Stream ainda ao vivo (sem ended_at) deve cair no thumb0. + """ + service = _make_service() + url = service.latest_thumbnail("abc123", 0, None) + + assert url.endswith("/media/thumbnails/thumb0.jpg") + + +def test_latest_thumbnail_meio_da_gravacao(): + """ + Gravacao de 10 min -> minuto do meio = 5 -> thumb20 (4 thumbs por minuto). + """ + service = _make_service() + + started_at = 0 + ended_at = 10 * 60 * 1000 # 10 minutos em ms + + url = service.latest_thumbnail("abc123", started_at, ended_at) + + assert url.endswith("/media/thumbnails/thumb20.jpg") + + +def test_latest_thumbnail_gravacao_curta(): + """ + Gravacao de menos de 2 min -> minuto do meio = 0 -> thumb0. + """ + service = _make_service() + + url = service.latest_thumbnail("abc123", 0, 30_000) + + assert url.endswith("/media/thumbnails/thumb0.jpg") + + +def test_latest_thumbnail_datas_invertidas(): + """ + ended_at antes de started_at nao deve gerar thumb negativo. + """ + service = _make_service() + + url = service.latest_thumbnail("abc123", 600_000, 0) + + assert url.endswith("/media/thumbnails/thumb0.jpg") + + +def test_latest_thumbnail_gravacao_longa(): + """ + Gravacao de 2h -> minuto do meio = 60 -> thumb240. + """ + service = _make_service() + + started_at = 0 + ended_at = 120 * 60 * 1000 # 2 horas em ms + + url = service.latest_thumbnail("abc123", started_at, ended_at) + + assert url.endswith("/media/thumbnails/thumb240.jpg") + + +# ---------------- LIST VODS: VALIDACAO ---------------- + +@pytest.mark.asyncio +async def test_list_vods_rejeita_room_id_vazio(): + """ + room_id vazio deve retornar 400. + """ + service = _make_service_with_db() + + with pytest.raises(SynapseError) as err: + await service.list_vods(room_id="", page=1, limit=10) + + assert err.value.code == 400 + + +@pytest.mark.asyncio +async def test_list_vods_rejeita_page_zero(): + """ + Deve retornar 400 quando page < 1. + """ + service = _make_service_with_db() + + with pytest.raises(SynapseError) as err: + await service.list_vods(room_id=ROOM_ID, page=0, limit=10) + + assert err.value.code == 400 + + +@pytest.mark.asyncio +async def test_list_vods_rejeita_limit_zero(): + """ + Deve retornar 400 quando limit < 1. + """ + service = _make_service_with_db() + + with pytest.raises(SynapseError) as err: + await service.list_vods(room_id=ROOM_ID, page=1, limit=0) + + assert err.value.code == 400 + + +@pytest.mark.asyncio +async def test_list_vods_rejeita_limit_acima_do_maximo(): + """ + Deve retornar 400 quando limit > 100, evitando dump da tabela inteira. + """ + service = _make_service_with_db() + + with pytest.raises(SynapseError) as err: + await service.list_vods(room_id=ROOM_ID, page=1, limit=101) + + assert err.value.code == 400 + + +@pytest.mark.asyncio +async def test_list_vods_aceita_limit_no_limite(): + """ + Deve aceitar limit = 100 (fronteira valida). + """ + service = _make_service_with_db(rows=[], total=0) + + result = await service.list_vods(room_id=ROOM_ID, page=1, limit=100) + + assert result["meta"]["perPage"] == 100 + + +# ---------------- LIST VODS: PAGINACAO ---------------- + +@pytest.mark.asyncio +async def test_list_vods_calcula_offset_da_primeira_pagina(): + """ + Pagina 1 deve consultar com offset 0. + """ + service = _make_service_with_db(rows=[], total=0) + + await service.list_vods(room_id=ROOM_ID, page=1, limit=10) + + call = service._store_mock.db_pool.runInteraction.await_args_list[0] + # (desc, func, room_id, limit, offset) + assert call.args[2] == ROOM_ID + assert call.args[3] == 10 + assert call.args[4] == 0 + + +@pytest.mark.asyncio +async def test_list_vods_calcula_offset_da_terceira_pagina(): + """ + Pagina 3 com limit 10 deve consultar com offset 20. + """ + service = _make_service_with_db(rows=[], total=0) + + await service.list_vods(room_id=ROOM_ID, page=3, limit=10) + + call = service._store_mock.db_pool.runInteraction.await_args_list[0] + assert call.args[4] == 20 + + +@pytest.mark.asyncio +async def test_list_vods_calcula_last_page_com_resto(): + """ + 25 vods com limit 10 devem resultar em 3 paginas (arredonda pra cima). + """ + service = _make_service_with_db(rows=[], total=25) + + result = await service.list_vods(room_id=ROOM_ID, page=1, limit=10) + + assert result["meta"]["total"] == 25 + assert result["meta"]["lastPage"] == 3 + + +@pytest.mark.asyncio +async def test_list_vods_calcula_last_page_exata(): + """ + 20 vods com limit 10 devem resultar em exatamente 2 paginas. + """ + service = _make_service_with_db(rows=[], total=20) + + result = await service.list_vods(room_id=ROOM_ID, page=1, limit=10) + + assert result["meta"]["lastPage"] == 2 + + +@pytest.mark.asyncio +async def test_list_vods_sem_resultados(): + """ + Deve retornar lista vazia sem quebrar quando nao ha vods. + """ + service = _make_service_with_db(rows=[], total=0) + + result = await service.list_vods(room_id=ROOM_ID, page=1, limit=10) + + assert result["data"] == [] + assert result["meta"]["total"] == 0 + + +# ---------------- LIST VODS: SERIALIZACAO ---------------- + +@pytest.mark.asyncio +async def test_list_vods_serializa_campos(): + """ + Deve devolver os campos que o front espera, incluindo roomId. + """ + service = _make_service_with_db(rows=[_row()], total=1) + + result = await service.list_vods(room_id=ROOM_ID, page=1, limit=10) + vod = result["data"][0] + + assert vod["id"] == 1 + assert vod["streamId"] == "teste-1" + assert vod["roomId"] == ROOM_ID + assert vod["title"] == "VOD de teste" + assert vod["recordingPath"] == "abc123" + assert vod["recordingDurationMs"] == 25900 + assert vod["isVod"] is True + assert vod["isLive"] is False + + +@pytest.mark.asyncio +async def test_list_vods_monta_master_playlist_url(): + """ + A URL de playback deve apontar para o Object Storage da Oracle. + """ + service = _make_service_with_db(rows=[_row(recording_path="xyz789")], total=1) + + result = await service.list_vods(room_id=ROOM_ID, page=1, limit=10) + vod = result["data"][0] + + assert vod["masterPlaylistUrl"] == ( + "https://objectstorage.sa-saopaulo-1.oraclecloud.com" + "/n/grmsqxilk3cb/b/vod-teste/o/xyz789/media/hls/master.m3u8" + ) + + +@pytest.mark.asyncio +async def test_list_vods_marca_live_quando_ended_at_null(): + """ + Stream sem ended_at deve ser marcada como live, nao como vod. + """ + service = _make_service_with_db(rows=[_row(ended_at=None)], total=1) + + result = await service.list_vods(room_id=ROOM_ID, page=1, limit=10) + vod = result["data"][0] + + assert vod["isLive"] is True + assert vod["isVod"] is False + + +@pytest.mark.asyncio +async def test_list_vods_serializa_varios_vods(): + """ + Deve serializar todas as linhas retornadas pelo banco. + """ + rows = [_row(row_id=1), _row(row_id=2), _row(row_id=3)] + service = _make_service_with_db(rows=rows, total=3) + + result = await service.list_vods(room_id=ROOM_ID, page=1, limit=10) + + assert len(result["data"]) == 3 + assert [v["id"] for v in result["data"]] == [1, 2, 3] + + +# ---------------- GET VOD ---------------- + +@pytest.mark.asyncio +async def test_get_vod_not_found(): + """ + Deve retornar 404 quando o stream_id nao existe. + """ + service = _make_service_with_row(row=None) + + with pytest.raises(SynapseError) as err: + await service.get_vod(stream_id=999) + + assert err.value.code == 404 + + +@pytest.mark.asyncio +async def test_get_vod_success(): + """ + Deve retornar o vod serializado com a URL de playback montada. + """ + service = _make_service_with_row(row=_row()) + + result = await service.get_vod(stream_id=1) + + assert result["id"] == 1 + assert result["recordingPath"] == "abc123" + assert result["isVod"] is True + assert result["masterPlaylistUrl"].endswith("/abc123/media/hls/master.m3u8") + + +@pytest.mark.asyncio +async def test_get_vod_passa_id_para_o_banco(): + """ + O stream_id recebido deve ser o mesmo consultado no banco. + """ + service = _make_service_with_row(row=_row(row_id=7)) + + await service.get_vod(stream_id=7) + + call = service._store_mock.db_pool.runInteraction.await_args + assert call.args[0] == "get_vod_by_id" + assert call.args[2] == 7 + + +# ---------------- DB: GET VODS ---------------- + +def test_get_vods_filtra_apenas_gravacoes_encerradas(): + """ + Live em andamento (ended_at NULL) nao deve aparecer na listagem de vods. + """ + txn = _make_txn(fetchall=[]) + + db.get_vods(txn, room_id=ROOM_ID, limit=10, offset=0) + + sql = txn.execute.call_args.args[0] + assert "ended_at IS NOT NULL" in sql + + +def test_get_vods_ordena_do_mais_recente(): + """ + Vods devem vir do mais recente para o mais antigo. + """ + txn = _make_txn(fetchall=[]) + + db.get_vods(txn, room_id=ROOM_ID, limit=10, offset=0) + + sql = txn.execute.call_args.args[0] + assert "ORDER BY started_at DESC" in sql + + +def test_get_vods_passa_parametros_na_ordem(): + """ + Os parametros devem ser (room_id, limit, offset). + """ + txn = _make_txn(fetchall=[]) + + db.get_vods(txn, room_id=ROOM_ID, limit=25, offset=50) + + params = txn.execute.call_args.args[1] + assert params == (ROOM_ID, 25, 50) + + +def test_get_vods_retorna_linhas_do_cursor(): + """ + Deve devolver exatamente o que o cursor retornou. + """ + rows = [_row()] + txn = _make_txn(fetchall=rows) + + result = db.get_vods(txn, room_id=ROOM_ID, limit=10, offset=0) + + assert result == rows + + +# ---------------- DB: COUNT VODS ---------------- + +def test_count_vods_usa_mesmo_filtro_do_list(): + """ + O total precisa contar so gravacoes encerradas, senao a paginacao mente. + """ + txn = _make_txn(fetchone=(3,)) + + db.count_vods(txn, room_id=ROOM_ID) + + sql = txn.execute.call_args.args[0] + assert "ended_at IS NOT NULL" in sql + + +def test_count_vods_retorna_total(): + """ + Deve extrair o total da primeira coluna. + """ + txn = _make_txn(fetchone=(42,)) + + result = db.count_vods(txn, room_id=ROOM_ID) + + assert result == 42 + + +def test_count_vods_sem_linhas(): + """ + Deve retornar 0 quando o cursor nao devolve nada. + """ + txn = _make_txn(fetchone=None) + + result = db.count_vods(txn, room_id=ROOM_ID) + + assert result == 0 + + +# ---------------- DB: GET VOD BY ID ---------------- + +def test_get_vod_by_id_passa_id(): + """ + Deve consultar pelo id recebido. + """ + txn = _make_txn(fetchone=None) + + db.get_vod_by_id(txn, stream_id=7) + + params = txn.execute.call_args.args[1] + assert params == (7,) + + +def test_get_vod_by_id_retorna_none_quando_nao_existe(): + """ + Deve devolver None para o service transformar em 404. + """ + txn = _make_txn(fetchone=None) + + result = db.get_vod_by_id(txn, stream_id=999) + + assert result is None \ No newline at end of file