diff --git a/src/backend/core/transcription/__init__.py b/src/backend/core/transcription/__init__.py new file mode 100644 index 00000000..e69de29b diff --git a/src/summary/summary/core/locales/__init__.py b/src/backend/core/transcription/locales/__init__.py similarity index 79% rename from src/summary/summary/core/locales/__init__.py rename to src/backend/core/transcription/locales/__init__.py index 89dfea46..0746775a 100644 --- a/src/summary/summary/core/locales/__init__.py +++ b/src/backend/core/transcription/locales/__init__.py @@ -2,9 +2,10 @@ from typing import Optional -from summary.core.config import get_settings -from summary.core.locales import de, en, fr, nl -from summary.core.locales.strings import LocaleStrings +from django.conf import settings + +from core.transcription.locales import de, en, fr, nl +from core.transcription.locales.strings import LocaleStrings _LOCALES = {"fr": fr, "en": en, "de": de, "nl": nl} @@ -27,4 +28,4 @@ def get_locale(*languages: Optional[str]) -> LocaleStrings: if base_lang in _LOCALES: return _LOCALES[base_lang].STRINGS - return _LOCALES[get_settings().default_context_language].STRINGS + return _LOCALES[settings.TRANSCRIPTION_DEFAULT_LANGUAGE].STRINGS diff --git a/src/summary/summary/core/locales/de.py b/src/backend/core/transcription/locales/de.py similarity index 93% rename from src/summary/summary/core/locales/de.py rename to src/backend/core/transcription/locales/de.py index 9d13f413..eb64cd62 100644 --- a/src/summary/summary/core/locales/de.py +++ b/src/backend/core/transcription/locales/de.py @@ -1,6 +1,6 @@ """German locale strings.""" -from summary.core.locales.strings import LocaleStrings +from core.transcription.locales.strings import LocaleStrings STRINGS = LocaleStrings( empty_transcription=""" diff --git a/src/summary/summary/core/locales/en.py b/src/backend/core/transcription/locales/en.py similarity index 92% rename from src/summary/summary/core/locales/en.py rename to src/backend/core/transcription/locales/en.py index a70ef1a9..a9ea3f15 100644 --- a/src/summary/summary/core/locales/en.py +++ b/src/backend/core/transcription/locales/en.py @@ -1,6 +1,6 @@ """English locale strings.""" -from summary.core.locales.strings import LocaleStrings +from core.transcription.locales.strings import LocaleStrings STRINGS = LocaleStrings( empty_transcription=""" diff --git a/src/summary/summary/core/locales/fr.py b/src/backend/core/transcription/locales/fr.py similarity index 93% rename from src/summary/summary/core/locales/fr.py rename to src/backend/core/transcription/locales/fr.py index 21214176..c170c5b8 100644 --- a/src/summary/summary/core/locales/fr.py +++ b/src/backend/core/transcription/locales/fr.py @@ -1,6 +1,6 @@ """French locale strings (default).""" -from summary.core.locales.strings import LocaleStrings +from core.transcription.locales.strings import LocaleStrings STRINGS = LocaleStrings( empty_transcription=""" diff --git a/src/summary/summary/core/locales/nl.py b/src/backend/core/transcription/locales/nl.py similarity index 93% rename from src/summary/summary/core/locales/nl.py rename to src/backend/core/transcription/locales/nl.py index d57e6cd5..667f31ad 100644 --- a/src/summary/summary/core/locales/nl.py +++ b/src/backend/core/transcription/locales/nl.py @@ -1,6 +1,6 @@ """Dutch locale strings.""" -from summary.core.locales.strings import LocaleStrings +from core.transcription.locales.strings import LocaleStrings STRINGS = LocaleStrings( empty_transcription=""" diff --git a/src/summary/summary/core/locales/strings.py b/src/backend/core/transcription/locales/strings.py similarity index 100% rename from src/summary/summary/core/locales/strings.py rename to src/backend/core/transcription/locales/strings.py diff --git a/src/summary/summary/core/transcript_formatter.py b/src/backend/core/transcription/transcript_formatter.py similarity index 95% rename from src/summary/summary/core/transcript_formatter.py rename to src/backend/core/transcription/transcript_formatter.py index 4dec9e33..2f88d9ea 100644 --- a/src/summary/summary/core/transcript_formatter.py +++ b/src/backend/core/transcription/transcript_formatter.py @@ -5,10 +5,9 @@ from datetime import datetime from typing import Tuple from zoneinfo import ZoneInfo -from summary.core.config import get_settings -from summary.core.locales import LocaleStrings +from django.conf import settings -settings = get_settings() +from core.transcription.locales.strings import LocaleStrings logger = logging.getLogger(__name__) @@ -38,7 +37,7 @@ class TranscriptFormatter: return None - def format( + def format( # pylint: disable=too-many-arguments,too-many-positional-arguments self, transcription, room: str | None = None, diff --git a/src/backend/meet/settings.py b/src/backend/meet/settings.py index a855a20d..9db07e76 100755 --- a/src/backend/meet/settings.py +++ b/src/backend/meet/settings.py @@ -744,6 +744,10 @@ class Base(Configuration): SUMMARY_SERVICE_API_TOKEN = SecretFileValue( None, environ_name="SUMMARY_SERVICE_API_TOKEN", environ_prefix=None ) + TRANSCRIPTION_DEFAULT_LANGUAGE = values.Value( + default="fr", environ_name="TRANSCRIPTION_DEFAULT_LANGUAGE", environ_prefix=None + ) + SCREEN_RECORDING_BASE_URL = values.Value( None, environ_name="SCREEN_RECORDING_BASE_URL", environ_prefix=None ) 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 19c56f35..dcdb5671 100644 --- a/src/summary/summary/core/celery_worker.py +++ b/src/summary/summary/core/celery_worker.py @@ -4,7 +4,6 @@ import json import time -from datetime import datetime import openai import sentry_sdk @@ -16,8 +15,8 @@ from summary.core.analytics import MetadataManager, get_analytics from summary.core.config import get_settings from summary.core.file_service import FileService, FileServiceException from summary.core.llm_service import LLMException, LLMObservability, LLMService -from summary.core.locales import get_locale from summary.core.models import ( + RecordingMetadata, SummarizeTaskV2Payload, TranscribeTaskV2Payload, ) @@ -39,11 +38,9 @@ from summary.core.shared_models import ( WhisperXResponse, webhook_payload_adapter, ) -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() @@ -79,23 +76,17 @@ file_service = FileService() def transcribe_audio( *, task_id: str, - recording_filename: str | None = None, language: str, - cloud_storage_url=None, + cloud_storage_url: str, raises: bool = False, ): """Transcribe an audio file using WhisperX. - Downloads the audio from MinIO or a cloud storage URL, sends it to + Downloads the audio from a cloud storage URL, sends it to WhisperX for transcription, and tracks metadata throughout the process. Returns the transcription object, or None if the file could not be retrieved. """ - if bool(recording_filename) == bool(cloud_storage_url): - raise ValueError( - "Either filename or cloud_storage_url must be provided, but not both." - ) - logger.info("Initiating WhisperX client") whisperx_client = openai.OpenAI( api_key=settings.whisperx_api_key.get_secret_value(), @@ -106,7 +97,6 @@ def transcribe_audio( # Transcription try: with file_service.prepare_audio_file( - remote_object_key=recording_filename, cloud_storage_url=cloud_storage_url, ) as (audio_file, metadata): metadata_manager.track(task_id, {"audio_length": metadata["duration"]}) @@ -150,10 +140,8 @@ def transcribe_audio( ) logger.exception( ( - "Unexpected error while preparing file | filename: %s " - "| cloud_storage_url: %s" + "Unexpected error while preparing file %s " ), - recording_filename, redacted_cloud_storage_url, ) return None @@ -163,41 +151,31 @@ def transcribe_audio( def resolve_speaker_identities_and_apply_to( - transcription, recording_start_at, recording_end_at, metadata_filename, task_id -): + *, transcription: WhisperXResponse, recording_metadata: RecordingMetadata, task_id +) -> WhisperXResponse: """Assign users to detected speakers and rewrite the transcriptions. Args: transcription: output of meet-whisperx after transcription and diarization - recording_start_at: sourced from LiveKit FileInfo via the egress_ended webhook - recording_end_at: sourced from LiveKit FileInfo via the egress_ended webhook - metadata_filename: name of metadata file containing VAD information in S3 + recording_metadata: Metadata of the recording task_id: current task id, for logging purposes """ - recording_start_dt = ( - datetime.fromisoformat(recording_start_at) if recording_start_at else None - ) - recording_end_dt = ( - datetime.fromisoformat(recording_end_at) if recording_end_at else None - ) - logger.debug( "recording_start_dt: %s ; recording_end_dt: %s", - recording_start_dt, - recording_end_dt, + recording_metadata.start_at, + recording_metadata.end_at, ) - if (recording_start_dt is None) or (recording_end_dt is None): - logger.debug("Skipping resolve_speaker_identities") - return transcription logger.debug("Running resolve_speaker_identities") try: - metadata = file_service.read_json(metadata_filename) + metadata = file_service.read_cloud_storage_json( + recording_metadata.cloud_storage_url + ) speaker_mapping = resolve_speaker_identities( metadata, transcription, - recording_start_dt, - recording_end_dt, + recording_metadata.start_at, + recording_metadata.end_at, ) new_transcription = speaker_mapping.apply_to(transcription.model_dump()) return new_transcription @@ -221,34 +199,6 @@ def resolve_speaker_identities_and_apply_to( return transcription -def format_transcript( - transcription, - context_language: str | None, - language: str, - room: str | None, - recording_datetime: str | None, - owner_timezone: str | None, - download_link: str | None, -) -> tuple[str, str]: - """Format a transcription into readable content with a title. - - Resolves the locale from context_language / language, then uses - TranscriptFormatter to produce markdown content and a title. - - Returns a (content, title) tuple. - """ - locale = get_locale(context_language, language) - formatter = TranscriptFormatter(locale) - - return formatter.format( - transcription, - room=room, - recording_datetime=recording_datetime, - owner_timezone=owner_timezone, - download_link=download_link, - ) - - def format_actions(llm_output: dict) -> str: """Format the actions from the LLM output into a markdown list. @@ -267,126 +217,23 @@ 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, - ) - - # Format output - content, title = format_transcript( - transcription, - context_language, - language, - room, - recording_start_at, - owner_timezone, - download_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) +# @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( @@ -468,29 +315,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 ################################################################################## @@ -549,6 +373,17 @@ def process_audio_transcribe_v2_task( ).model_dump() ) + # Assign speakers and rewrite transcription/diarization output + if settings.is_resolve_speaker_identities_enabled and payload.metadata is not None: + try: + transcription_res = resolve_speaker_identities_and_apply_to( + transcription=transcription_res, + recording_metadata=payload.metadata, + task_id=job_id, + ) + except BaseException as e: + logger.error(f"Failed to resolve speaker identities, skipping: {e}") + file_service.store_transcript( transcript=transcription_res, job_id=job_id, @@ -561,6 +396,8 @@ def process_audio_transcribe_v2_task( call_webhook_v2_task.apply_async( args=[success_payload.model_dump(), payload.tenant_id] ) + metadata_manager.capture(job_id, settings.posthog_event_success) + return success_payload.model_dump() diff --git a/src/summary/summary/core/config.py b/src/summary/summary/core/config.py index 6858a49d..872bedcd 100644 --- a/src/summary/summary/core/config.py +++ b/src/summary/summary/core/config.py @@ -1,9 +1,8 @@ """Application configuration and settings.""" import logging -import os from functools import cached_property, lru_cache -from typing import Annotated, Any, List, Literal, Mapping, Optional, Set +from typing import Annotated, List, Literal, Mapping, Optional, Set from fastapi import Depends from pydantic import ( @@ -36,8 +35,6 @@ class AuthorizedTenant(BaseModel): ) -V1_DEFAULT_TENANT_ID = "__deprecated_meet_tenant__" - class Settings(BaseSettings): """Configuration settings loaded from environment variables and .env file.""" @@ -45,14 +42,12 @@ class Settings(BaseSettings): 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 # 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" @@ -114,9 +107,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}" @@ -145,33 +135,6 @@ 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.""" @@ -190,16 +153,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]: """Return a frozenset of authorized tenant API api_keys.""" diff --git a/src/summary/summary/core/file_service.py b/src/summary/summary/core/file_service.py index 994b3d1d..51ecb8ae 100644 --- a/src/summary/summary/core/file_service.py +++ b/src/summary/summary/core/file_service.py @@ -13,7 +13,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 @@ -144,54 +143,6 @@ class FileService: self._allowed_extensions = settings.recording_allowed_extensions self._max_duration = 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() - - if extension not in self._allowed_extensions: - logger.warning("Invalid file extension '%s'", extension) - raise ValueError(f"Invalid file extension '{extension}'") - - 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( @@ -302,33 +253,19 @@ class FileService: os.remove(output_path) raise RuntimeError("Failed to extract audio.") from e - def read_json(self, object_name: str) -> dict: + def read_cloud_storage_json(self, cloud_storage_url: 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 + logger.info("Reading JSON: %s", cloud_storage_url) + local_path = self._download_from_cloud_storage_url(cloud_storage_url) 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 + return json.load(local_path.open("r")) except (json.JSONDecodeError, UnicodeDecodeError) as e: raise FileServiceException("Invalid JSON content.") from e - finally: - if response: - response.close() - response.release_conn() @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. @@ -341,20 +278,9 @@ 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 + ) duration = self._validate_duration(downloaded_path) diff --git a/src/summary/summary/core/models.py b/src/summary/summary/core/models.py index 11983d15..ed97fd82 100644 --- a/src/summary/summary/core/models.py +++ b/src/summary/summary/core/models.py @@ -1,6 +1,6 @@ """Models for the API & Celery tasks creation.""" -from pydantic import BaseModel, Field, field_validator +from pydantic import AwareDatetime, BaseModel, Field, field_validator from summary.core.config import get_settings from summary.core.types import Url @@ -14,6 +14,17 @@ class SharedV2TaskCreation(BaseModel): user_sub: str = Field(title="User Sub", description="The user's sub.") +class RecordingMetadata(BaseModel): + """Model for recording metadata.""" + + cloud_storage_url: Url = Field( + title="Cloud Storage URL", + description="The URL of the metadata file for speaker assignement.", + ) + start_at: AwareDatetime = Field(title="Start time of the recording to transcribe") + end_at: AwareDatetime = Field(title="End time of the recording to transcribe") + + class TranscribeTaskV2Request(SharedV2TaskCreation): """Model for creating a transcribe and summarize task (used for API request).""" @@ -27,7 +38,12 @@ class TranscribeTaskV2Request(SharedV2TaskCreation): description="The language of the context text.", ) language: str = Field( - title="Language", description="The language of the content to summarize." + title="Language", description="The language of the content to transcribe." + ) + metadata: RecordingMetadata | None = Field( + title="Metadata", + description="The metadata for the transcribe task.", + default=None, ) @field_validator("language") diff --git a/src/summary/summary/core/shared_models.py b/src/summary/summary/core/shared_models.py index 043ea759..2f3caf79 100644 --- a/src/summary/summary/core/shared_models.py +++ b/src/summary/summary/core/shared_models.py @@ -159,4 +159,5 @@ __all__ = [ "SummarizeWebhookPayloads", "WebhookPayloads", "WhisperXResponse", + "webhook_payload_adapter", ] diff --git a/src/summary/summary/core/webhook_service.py b/src/summary/summary/core/webhook_service.py index 0e196cb7..bd38e225 100644 --- a/src/summary/summary/core/webhook_service.py +++ b/src/summary/summary/core/webhook_service.py @@ -46,55 +46,6 @@ def _post_with_retries(*, url, data, api_key: str | None = None): session.close() -def call_webhook_v1(*, tenant_id: str, payload: dict) -> None: - """Call webhook with payload a payload and optional token.""" - tenant = settings.get_authorized_tenant(tenant_id=tenant_id) - - logger.debug("Submitting to %s", tenant.webhook_url) - logger.debug("Request payload: %s", json.dumps(payload, indent=2)) - - response = _post_with_retries( - url=tenant.webhook_url, - api_key=tenant.webhook_api_key.get_secret_value(), - data=payload, - ) - - try: - response_data = response.json() - document_id = response_data.get("id", "N/A") - except (json.JSONDecodeError, AttributeError): - document_id = "Unable to parse response" - response_data = response.text - - logger.info( - "Delivery success | Document %s submitted (HTTP %s)", - document_id, - response.status_code, - ) - logger.debug("Full response: %s", response_data) - - -def submit_content(content: str, title: str, email: str, sub: str) -> None: - """Submit content to the configured webhook destination. - - Builds the payload, sends it with retries, and logs the outcome. - - Notes: - Deprecated: Use call_webhook_v2 directly instead. - - Deprecated: - This will route content to the v1 default tenant - """ - data = { - "title": title, - "content": content, - "email": email, - "sub": sub, - } - - call_webhook_v1(payload=data, tenant_id=settings.v1_tenant_id) - - def call_webhook_v2( *, tenant_id: str, 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)