diff --git a/.github/workflows/meet.yml b/.github/workflows/meet.yml index da5c14d9..e7960ac0 100644 --- a/.github/workflows/meet.yml +++ b/.github/workflows/meet.yml @@ -305,7 +305,6 @@ jobs: working-directory: src/summary env: - V1_TENANT_ID: 'test-tenant' AUTHORIZED_TENANTS: '[{"id": "test-tenant", "api_key": "test-api-token", "webhook_url": "https://example.com/webhook", "webhook_api_key": "test-webhook-api-key"}]' AWS_STORAGE_BUCKET_NAME: "http://meet-media-storage" AWS_S3_ENDPOINT_URL: "minio:9000" diff --git a/compose.yml b/compose.yml index 68d62ed6..659f145f 100644 --- a/compose.yml +++ b/compose.yml @@ -301,7 +301,7 @@ services: context: ./src/summary dockerfile: Dockerfile target: production - command: celery -A summary.core.celery_worker worker --pool=solo --loglevel=debug -Q transcribe-queue + command: celery -A summary.core.celery_worker worker --pool=solo --loglevel=debug -Q transcribe-queue-v2 env_file: - env.d/development/summary volumes: @@ -321,7 +321,7 @@ services: context: ./src/summary dockerfile: Dockerfile target: production - command: celery -A summary.core.celery_worker worker --pool=solo --loglevel=debug -Q summarize-queue + command: celery -A summary.core.celery_worker worker --pool=solo --loglevel=debug -Q summarize-queue-v2 env_file: - env.d/development/summary volumes: diff --git a/src/helm/env.d/common.yaml.gotmpl b/src/helm/env.d/common.yaml.gotmpl index ae0cf195..26fbb043 100644 --- a/src/helm/env.d/common.yaml.gotmpl +++ b/src/helm/env.d/common.yaml.gotmpl @@ -28,7 +28,25 @@ _summaryEnvVars: &summaryEnvVars AWS_S3_ACCESS_KEY_ID: meet AWS_S3_SECRET_ACCESS_KEY: password AWS_S3_SECURE_ACCESS: False - AUTHORIZED_TENANTS: '[{"id": "dictaphone", "api_key": "dictaphone_token", "webhook_url": "http://dictaphone-backend.dictaphone.svc.cluster.local/api/v1.0/ai-jobs/webhook/", "webhook_api_key": "token_summary"}]' + AUTHORIZED_TENANTS: > + [ + { + "id": "dictaphone", + "api_key": "dictaphone_token", + "webhook_url": "http://dictaphone-backend.dictaphone.svc.cluster.local/api/v1.0/ai-jobs/webhook/", + "webhook_api_key": "token_summary", + "allowed_push_to_docs": false + }, + { + "id": "visio", + "api_key": "password", + "webhook_url": "https://meet.127.0.0.1.nip.io/api/v1.0/recordings/external-process-hook/", + "webhook_api_key": "webhook-password", + "allowed_push_to_docs": false + } + ] + SSL_CERT_FILE: /usr/local/lib/python3.13/site-packages/certifi/cacert.pem + IS_DOCS_INTEGRATION_ENABLED: false WHISPERX_API_KEY: secretKeyRef: name: secret-dev @@ -169,8 +187,9 @@ backend: RECORDING_ENABLE: True RECORDING_STORAGE_EVENT_ENABLE: True RECORDING_STORAGE_EVENT_TOKEN: password - SUMMARY_SERVICE_ENDPOINT: http://meet-summary:80/api/v1/tasks/ + SUMMARY_SERVICE_ENDPOINT: http://meet-summary:80/api/v2/async-jobs/transcribe/ SUMMARY_SERVICE_API_TOKEN: password + SUMMARY_SERVICE_WEBHOOK_API_TOKEN: webhook-password RECORDING_DOWNLOAD_BASE_URL: https://meet.127.0.0.1.nip.io/recording migrate: @@ -291,7 +310,7 @@ celeryTranscribe: - "--pool=solo" - "--loglevel=info" - "-Q" - - "transcribe-queue,transcribe-queue-v2" + - "transcribe-queue-v2" celerySummarize: replicas: 1 @@ -308,7 +327,7 @@ celerySummarize: - "--pool=solo" - "--loglevel=info" - "-Q" - - "summarize-queue,summarize-queue-v2" + - "summarize-queue-v2" celerySummaryBackend: replicas: 1 diff --git a/src/summary/summary/api/main.py b/src/summary/summary/api/main.py index 65937827..e553473c 100644 --- a/src/summary/summary/api/main.py +++ b/src/summary/summary/api/main.py @@ -2,11 +2,8 @@ from fastapi import APIRouter, Depends -from summary.api.route import tasks, tasks_v2 +from summary.api.route import tasks_v2 from summary.core.security import verify_tenant_api_key -api_router_v1 = APIRouter(dependencies=[Depends(verify_tenant_api_key)]) -api_router_v1.include_router(tasks.router_tasks_v1, tags=["tasks"]) - api_router_v2 = APIRouter(dependencies=[Depends(verify_tenant_api_key)]) api_router_v2.include_router(tasks_v2.router_tasks_v2, tags=["tasks"]) diff --git a/src/summary/summary/api/route/tasks.py b/src/summary/summary/api/route/tasks.py deleted file mode 100644 index 91705f09..00000000 --- a/src/summary/summary/api/route/tasks.py +++ /dev/null @@ -1,79 +0,0 @@ -"""API routes related to application tasks.""" - -import time -from typing import Optional - -from celery.result import AsyncResult -from fastapi import APIRouter -from pydantic import BaseModel, field_validator - -from summary.core.celery_worker import ( - process_audio_transcribe_summarize_v2, -) -from summary.core.config import get_settings - -settings = get_settings() - - -class TranscribeSummarizeTaskCreation(BaseModel): - """Transcription and summarization parameters.""" - - owner_id: str - recording_filename: str - metadata_filename: Optional[str] = None - email: str - sub: str - version: Optional[int] = 2 - room: Optional[str] - owner_timezone: Optional[str] - language: Optional[str] - download_link: Optional[str] - context_language: Optional[str] = None - recording_start_at: Optional[str] = None - recording_end_at: Optional[str] = None - - @field_validator("language") - @classmethod - def validate_language(cls, v): - """Validate 'language' parameter.""" - if v is not None and v not in settings.whisperx_allowed_languages: - raise ValueError( - f"Language '{v}' is not allowed. " - f"Allowed languages: {', '.join(settings.whisperx_allowed_languages)}" - ) - return v - - -router_tasks_v1 = APIRouter(prefix="/tasks") - - -@router_tasks_v1.post("/") -async def create_transcribe_summarize_task(request: TranscribeSummarizeTaskCreation): - """Create a transcription and summarization task.""" - task = process_audio_transcribe_summarize_v2.apply_async( - args=[ - request.owner_id, - request.recording_filename, - request.metadata_filename, - request.email, - request.sub, - time.time(), - request.room, - request.owner_timezone, - request.language, - request.download_link, - request.context_language, - request.recording_start_at, - request.recording_end_at, - ], - queue=settings.transcribe_queue, - ) - - return {"id": task.id, "message": "Task created"} - - -@router_tasks_v1.get("/{task_id}") -async def get_task_status(task_id: str): - """Check task status by ID.""" - task = AsyncResult(task_id) - return {"id": task_id, "status": task.status} diff --git a/src/summary/summary/core/celery_worker.py b/src/summary/summary/core/celery_worker.py index ada448b7..ff803cac 100644 --- a/src/summary/summary/core/celery_worker.py +++ b/src/summary/summary/core/celery_worker.py @@ -50,7 +50,6 @@ from summary.core.transcript_formatter import TranscriptFormatter from summary.core.user_assign import resolve_speaker_identities from summary.core.webhook_service import ( call_webhook_v2, - submit_content, ) settings = get_settings() @@ -317,136 +316,6 @@ def format_actions(llm_output: dict) -> str: return "" -@celery.task( - bind=True, - autoretry_for=[exceptions.HTTPError], - max_retries=settings.celery_max_retries, - queue=settings.transcribe_queue, -) -def process_audio_transcribe_summarize_v2( - self, - owner_id: str, - recording_filename: str, - metadata_filename: str | None, - email: str, - sub: str, - received_at: float, - room: str | None, - owner_timezone: str | None, - language: str | None, - download_link: str | None, - context_language: str | None = None, - recording_start_at: str | None = None, - recording_end_at: str | None = None, -): - """Process an audio file by transcribing it and generating a summary. - - This Celery task orchestrates: - 1. Audio transcription via WhisperX - 2. Transcript formatting - 3. Webhook submission - 4. Conditional summarization queuing - - Args: - self: Celery task instance (passed on with bind=True) - owner_id: Unique identifier of the recording owner. - recording_filename: Name of the audio file in MinIO storage. - metadata_filename: Name of the audio file in MinIO storage. - email: Email address of the recording owner. - sub: OIDC subject identifier of the recording owner. - received_at: Unix timestamp when the recording was received. - room: room name where the recording took place. - owner_timezone: IANA timezone of the recording owner (e.g. "Europe/Paris"). - language: ISO 639-1 language code for transcription. - download_link: URL to download the original recording. - context_language: ISO 639-1 language code of the meeting summary context text. - recording_start_at: ISO 8601 timestamp of when file recording actually started - (from LiveKit FileInfo.started_at via the egress_ended webhook). - recording_end_at: ISO 8601 timestamp of when file recording ended - (from LiveKit FileInfo.ended_at via the egress_ended webhook). - """ - logger.info( - "Notification received | Owner: %s | Room: %s", - owner_id, - room, - ) - - task_id = self.request.id - - # Transcribe the audio - transcription = transcribe_audio( - task_id=task_id, recording_filename=recording_filename, language=language - ) - if transcription is None: - return - - # Assign speakers and rewrite transcription/diarization output - if settings.is_resolve_speaker_identities_enabled and ( - metadata_filename is not None - ): - transcription = resolve_speaker_identities_and_apply_to( - transcription, - recording_start_at, - recording_end_at, - metadata_filename, - task_id, - ) - - form_base_url = settings.transcription_satisfaction_form_base_url - form_link = ( - f"{form_base_url}?room_id={room}" - if (form_base_url and metadata_filename is not None) - else None - ) - - # Format output - content, title = format_transcript( - transcription, - context_language, - language, - room, - recording_start_at, - owner_timezone, - download_link, - form_link, - ) - - submit_content(content, title, email, sub) - metadata_manager.capture(task_id, settings.posthog_event_success) - - # LLM Summarization - if ( - analytics.is_feature_enabled("summary-enabled", distinct_id=owner_id) - and settings.is_summary_enabled - ): - logger.info("Queuing summary generation task.") - summarize_transcription.apply_async( - args=[owner_id, content, email, sub, title], - queue=settings.summarize_queue, - ) - else: - logger.info("Summary generation not enabled for this user. Skipping.") - - -@signals.task_prerun.connect(sender=process_audio_transcribe_summarize_v2) -def task_started(task_id=None, task=None, args=None, **kwargs): - """Signal handler called before task execution begins.""" - task_args = args or [] - metadata_manager.create(task_id, task_args) - - -@signals.task_retry.connect(sender=process_audio_transcribe_summarize_v2) -def task_retry_handler(request=None, reason=None, einfo=None, **kwargs): - """Signal handler called when task execution retries.""" - metadata_manager.retry(request.id) - - -@signals.task_failure.connect(sender=process_audio_transcribe_summarize_v2) -def task_failure_handler(task_id, exception=None, **kwargs): - """Signal handler called when task execution fails permanently.""" - metadata_manager.capture(task_id, settings.posthog_event_failure) - - def summarize_transcription_internals( *, owner_id: str, transcript: str, session_id: str ) -> str: @@ -526,29 +395,6 @@ def summarize_transcription_internals( return summary -@celery.task( - bind=True, - autoretry_for=[LLMException, Exception], - max_retries=settings.celery_max_retries, - queue=settings.summarize_queue, -) -def summarize_transcription( - self, owner_id: str, transcript: str, email: str, sub: str, title: str -): - """Generate a summary from the provided transcription text. - - This Celery task performs the following operations: - 1. Run summary internals - 2. Sends the final summary via webhook. - """ - summary = summarize_transcription_internals( - owner_id=owner_id, transcript=transcript, session_id=self.request.id - ) - summary_title = settings.summary_title_template.format(title=title) - - submit_content(summary, summary_title, email, sub) - - ################################################################################## # Tasks v2 ################################################################################## diff --git a/src/summary/summary/core/config.py b/src/summary/summary/core/config.py index b536713e..c3089cf8 100644 --- a/src/summary/summary/core/config.py +++ b/src/summary/summary/core/config.py @@ -36,23 +36,18 @@ class AuthorizedTenant(BaseModel): ) -V1_DEFAULT_TENANT_ID = "__deprecated_meet_tenant__" - - class Settings(BaseSettings): """Configuration settings loaded from environment variables and .env file.""" model_config = SettingsConfigDict(env_file=".env", frozen=True) app_name: str = "summary" - app_api_v1_str: str = "/api/v1" app_api_v2_str: str = "/api/v2" # Authorized Tenants # Using env variables to store authorized tenants for now # to avoid any other external dependency (DB) - authorized_tenants: tuple[AuthorizedTenant, ...] = Field(default_factory=tuple) - v1_tenant_id: str = V1_DEFAULT_TENANT_ID + authorized_tenants: tuple[AuthorizedTenant, ...] = Field(min_length=1) # Audio recordings recording_max_duration: Optional[int] = None @@ -72,8 +67,6 @@ class Settings(BaseSettings): celery_result_backend: str = "redis://redis/0" celery_max_retries: int = 1 - transcribe_queue: str = "transcribe-queue" - summarize_queue: str = "summarize-queue" # v2 tasks transcribe_queue_v2: str = "transcribe-queue-v2" summarize_queue_v2: str = "summarize-queue-v2" @@ -113,12 +106,6 @@ class Settings(BaseSettings): webhook_status_forcelist: List[int] = [502, 503, 504] webhook_backoff_factor: float = 0.1 - # Locale - default_context_language: Literal["de", "en", "fr", "nl"] = "fr" - - # Output related settings - summary_title_template: Optional[str] = "Résumé de {title}" - # Summary related settings is_summary_enabled: bool = True transcription_satisfaction_form_base_url: Optional[str] = None @@ -145,38 +132,9 @@ class Settings(BaseSettings): task_tracker_redis_url: str = "redis://redis/0" task_tracker_prefix: str = "task_metadata:" - @model_validator(mode="before") - @classmethod - def legacy_default_tenant_config(cls, data: Any) -> Any: - """Migrate the legacy default tenant configuration.""" - if isinstance(data, dict): - api_key = os.getenv("APP_API_TOKEN") - webhook_api_key = os.getenv("WEBHOOK_API_TOKEN") - webhook_url = os.getenv("WEBHOOK_URL") - if api_key and webhook_api_key and webhook_url: - logger.warning( - "Deprecated legacy app configuration detected, " - "please use only the new 'authorized_tenants' field instead." - ) - - authorized_tenants = list(data.get("authorized_tenants", [])) - authorized_tenants.append( - AuthorizedTenant( - id=V1_DEFAULT_TENANT_ID, - api_key=SecretStr(api_key), - webhook_url=webhook_url, - webhook_api_key=SecretStr(webhook_api_key), - ) - ) - data["authorized_tenants"] = tuple(authorized_tenants) - - return data - @model_validator(mode="after") def validate_authorized_tenants(self): """Validate authorized tenants configuration.""" - if len(self.authorized_tenants) == 0: - raise ValueError("No authorized tenants configured") tenant_ids = {tenant.id for tenant in self.authorized_tenants} if len(tenant_ids) != len(self.authorized_tenants): @@ -190,15 +148,6 @@ class Settings(BaseSettings): raise ValueError("Duplicate application API api_keys are not allowed") return self - @model_validator(mode="after") - def validate_default_v1_tenant(self): - """Validate default v1 tenant configuration.""" - if not any( - tenant.id == self.v1_tenant_id for tenant in self.authorized_tenants - ): - raise ValueError("v1 tenant is not configured in authorized tenants") - - return self @cached_property def authorized_tenant_api_keys(self) -> frozenset[str]: diff --git a/src/summary/summary/core/file_service.py b/src/summary/summary/core/file_service.py index 1b4437f4..920f577e 100644 --- a/src/summary/summary/core/file_service.py +++ b/src/summary/summary/core/file_service.py @@ -14,7 +14,6 @@ from urllib.parse import urlparse import requests from minio import Minio -from minio.error import MinioException, S3Error from summary.core.config import get_settings from summary.core.shared_models import WhisperXResponse @@ -284,50 +283,6 @@ class FileService: self._max_duration_seconds = settings.recording_max_duration - def _download_from_minio(self, remote_object_key) -> Path: - """Download file from MinIO to local temporary file. - - The file is downloaded to a temporary location for local manipulation - such as validation, conversion, or processing before being used. - """ - logger.info("Download recording | object_key: %s", remote_object_key) - - if not remote_object_key: - logger.warning("Invalid object_key '%s'", remote_object_key) - raise ValueError("Invalid object_key") - - extension = Path(remote_object_key).suffix.lower() - - response = None - - try: - response = self._minio_client.get_object( - self._bucket_name, remote_object_key - ) - - with tempfile.NamedTemporaryFile( - suffix=extension, delete=False, prefix="minio_download_" - ) as tmp: - for chunk in response.stream(self._stream_chunk_size): - tmp.write(chunk) - - tmp.flush() - local_path = Path(tmp.name) - - logger.info("Recording successfully downloaded") - logger.debug("Recording local file path: %s", local_path) - - return local_path - - except (MinioException, S3Error) as e: - raise FileServiceException( - "Unexpected error while downloading object." - ) from e - - finally: - if response: - response.close() - def _download_from_cloud_storage_url(self, cloud_storage_url: str) -> Path: """Download file from a cloud storage URL to local temporary file.""" logger.info( @@ -392,33 +347,17 @@ class FileService: logger.error(error_msg) raise MediaDurationTooLongError(error_msg) - def read_json(self, object_name: str) -> dict: - """Read and parse a JSON file from MinIO storage.""" - logger.info("Reading JSON: %s", object_name) - - if not object_name: - raise ValueError("Invalid object_name") - - response = None - try: - response = self._minio_client.get_object(self._bucket_name, object_name) - return json.loads(response.read()) - except (MinioException, S3Error) as e: - raise FileServiceException( - "Unexpected error while reading JSON object." - ) from e - except (json.JSONDecodeError, UnicodeDecodeError) as e: - raise FileServiceException("Invalid JSON content.") from e - finally: - if response: - response.close() - response.release_conn() + def read_cloud_storage_json(self, cloud_storage_url: str) -> dict: + """Read and parse a JSON file from a url.""" + logger.info("Reading JSON: %s", cloud_storage_url) + res = requests.get(cloud_storage_url, timeout=(10, 100)) + res.raise_for_status() + return res.json() @contextmanager def prepare_audio_file( self, - remote_object_key: str | None = None, - cloud_storage_url: str | None = None, + cloud_storage_url: str, ): """Download and prepare audio file for processing. @@ -431,20 +370,7 @@ class FileService: file_handle = None try: - if bool(remote_object_key) == bool(cloud_storage_url): - raise ValueError( - ( - "Exactly one of 'remote_object_key' or " - "'cloud_storage_url' must be provided." - ) - ) - - if cloud_storage_url: - downloaded_path = self._download_from_cloud_storage_url( - cloud_storage_url - ) - else: - downloaded_path = self._download_from_minio(remote_object_key) + downloaded_path = self._download_from_cloud_storage_url(cloud_storage_url) media_info = get_media_info(downloaded_path) diff --git a/src/summary/summary/main.py b/src/summary/summary/main.py index 6447671d..a63d2413 100644 --- a/src/summary/summary/main.py +++ b/src/summary/summary/main.py @@ -4,7 +4,7 @@ import sentry_sdk from fastapi import FastAPI from summary.api import health -from summary.api.main import api_router_v1, api_router_v2 +from summary.api.main import api_router_v2 from summary.core.config import get_settings settings = get_settings() @@ -17,6 +17,5 @@ app = FastAPI( title=settings.app_name, ) -app.include_router(api_router_v1, prefix=settings.app_api_v1_str) app.include_router(api_router_v2, prefix=settings.app_api_v2_str) app.include_router(health.router) diff --git a/src/summary/tests/api/test_api_tasks.py b/src/summary/tests/api/test_api_tasks.py deleted file mode 100644 index 98ef066b..00000000 --- a/src/summary/tests/api/test_api_tasks.py +++ /dev/null @@ -1,91 +0,0 @@ -"""Integration tests for the API tasks endpoints.""" - -# tests/unit/test_api_tasks.py -from unittest.mock import MagicMock, patch - - -class TestTasks: - """Tests for the /v1/tasks endpoint.""" - - @patch( - "summary.api.route.tasks.process_audio_transcribe_summarize_v2.apply_async", - return_value=MagicMock(id="task-id-abc"), - ) - @patch("summary.api.route.tasks.time.time", return_value=1735725600.0) - def test_create_task_returns_task_id(self, mock_time, mock_apply_async, client): - """POST /tasks/ with valid payload returns id and dispatches Celery task.""" - response = client.post( - "api/v1/tasks/", - headers={"Authorization": "Bearer test-api-token"}, - json={ - "owner_id": "owner-123", - "recording_filename": "recording.mp4", - "metadata_filename": "metadata.json", - "email": "user@example.com", - "sub": "sub-123", - "room": "room-abc", - "owner_timezone": "UTC", - "language": None, - "download_link": "https://example.com/file.mp4", - "form_link": "https://satisfaction.com?room_id=test", - }, - ) - - assert response.status_code == 200 - assert response.json() == {"id": "task-id-abc", "message": "Task created"} - - args = mock_apply_async.call_args.kwargs["args"] - assert args == [ - "owner-123", # owner_id - "recording.mp4", # recording_filename - "metadata.json", # metadata_filename - "user@example.com", # email - "sub-123", # sub - 1735725600.0, # received_at - "room-abc", # room - "UTC", # owner_timezone - None, # language - "https://example.com/file.mp4", # download_link - None, # context_language - None, # recording_start_at - None, # recording_end_at - ] - - def test_create_task_invalid_language(self, client): - """POST /tasks/ with an unsupported language returns 422.""" - payload = {"language": "klingon"} - response = client.post( - "/api/v1/tasks/", - headers={"Authorization": "Bearer test-api-token"}, - json=payload, - ) - - assert response.status_code == 422 - - @patch( - "summary.api.route.tasks.AsyncResult", - return_value=MagicMock(status="PENDING"), - ) - def test_get_task_status_pending(self, mock_result, client): - """GET /tasks/{id} returns PENDING status when the task has not started yet.""" - response = client.get( - "/api/v1/tasks/task-id-abc", - headers={"Authorization": "Bearer test-api-token"}, - ) - - assert response.status_code == 200 - assert response.json() == {"id": "task-id-abc", "status": "PENDING"} - - @patch( - "summary.api.route.tasks.AsyncResult", - return_value=MagicMock(status="SUCCESS"), - ) - def test_get_task_status_success(self, mock_result, client): - """GET /tasks/{id} returns SUCCESS status when the task has completed.""" - response = client.get( - "/api/v1/tasks/task-id-abc", - headers={"Authorization": "Bearer test-api-token"}, - ) - - assert response.status_code == 200 - assert response.json()["status"] == "SUCCESS" diff --git a/src/summary/tests/conftest.py b/src/summary/tests/conftest.py index 6971ed5c..f24f8db5 100644 --- a/src/summary/tests/conftest.py +++ b/src/summary/tests/conftest.py @@ -11,7 +11,6 @@ from summary.main import app def get_settings_override(): """Return settings for tests.""" return Settings( - v1_tenant_id="test-tenant", authorized_tenants=( AuthorizedTenant( webhook_url="https://example.com/webhook",