mirror of
https://github.com/suitenumerique/meet.git
synced 2026-07-27 12:19:10 +00:00
💥(summary) removed summary v1 related code
We are in the process of moving meet to summary v2 routes and tasks. This commits removes code from summary and moves some code from summary to meet, which should have the responsability of this code.
This commit is contained in:
@@ -14,7 +14,8 @@ dependencies = [
|
||||
"posthog==7.9.12",
|
||||
"requests==2.33.0",
|
||||
"sentry-sdk[fastapi, celery]==2.54.0",
|
||||
"langfuse==4.0.0"
|
||||
"langfuse==4.0.0",
|
||||
"ruff>=0.15.6",
|
||||
]
|
||||
|
||||
[project.optional-dependencies]
|
||||
|
||||
@@ -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"])
|
||||
|
||||
@@ -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}
|
||||
@@ -1,7 +1,10 @@
|
||||
"""API routes related to application tasks (V2 / tenant friendly)."""
|
||||
|
||||
import logging
|
||||
from datetime import datetime, timezone
|
||||
|
||||
from celery.result import AsyncResult
|
||||
from fastapi import APIRouter, Depends, HTTPException, Request
|
||||
from fastapi import APIRouter, Depends, HTTPException, Request, status
|
||||
|
||||
from summary.core.celery_worker import (
|
||||
celery,
|
||||
@@ -20,6 +23,7 @@ from summary.core.shared_models import (
|
||||
TranscribeWebhookSuccessPayload,
|
||||
)
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
router_tasks_v2 = APIRouter()
|
||||
|
||||
|
||||
@@ -29,8 +33,20 @@ async def create_transcribe_task_v2(
|
||||
request_tenant: AuthorizedTenant = Depends(verify_tenant_api_key_v2),
|
||||
):
|
||||
"""Create a transcription task."""
|
||||
if (
|
||||
request.push_to_docs_config is not None
|
||||
and not request_tenant.allowed_push_to_docs
|
||||
):
|
||||
logger.error(
|
||||
f"Push to docs is not allowed for this tenant ({request_tenant.id})."
|
||||
)
|
||||
raise HTTPException(
|
||||
status_code=status.HTTP_400_BAD_REQUEST,
|
||||
detail="Push to docs is not allowed for this tenant.",
|
||||
)
|
||||
|
||||
task = process_audio_transcribe_v2_task.apply_async(
|
||||
args=[{**request.model_dump(), "tenant_id": request_tenant.id}]
|
||||
args=[{**request.model_dump(), "tenant_id": request_tenant.id, "received_at": datetime.now(timezone.utc)}]
|
||||
)
|
||||
|
||||
return TranscribeWebhookPendingPayload(job_id=task.id).model_dump()
|
||||
|
||||
@@ -4,12 +4,14 @@ import json
|
||||
import time
|
||||
from collections import Counter
|
||||
from functools import lru_cache
|
||||
from urllib.parse import urlsplit, urlunsplit
|
||||
|
||||
import redis
|
||||
from celery.utils.log import get_task_logger
|
||||
from posthog import Posthog
|
||||
|
||||
from summary.core.config import get_settings
|
||||
from summary.core.models import TranscribeTaskV2Payload
|
||||
|
||||
logger = get_task_logger(__name__)
|
||||
settings = get_settings()
|
||||
@@ -107,23 +109,23 @@ class MetadataManager:
|
||||
"""Check if task_id exists in tasks metadata cache."""
|
||||
return self._redis.exists(self._get_redis_key(task_id))
|
||||
|
||||
def create(self, task_id, task_args):
|
||||
def create(self, task_id: str, task_payload: TranscribeTaskV2Payload):
|
||||
"""Create initial metadata entry for a new task."""
|
||||
if self._is_disabled or self.has_task_id(task_id):
|
||||
return
|
||||
|
||||
# Positional args mirror process_audio_transcribe_summarize_v2 signature:
|
||||
# owner_id, recording_filename, metadata_filename, email, sub, received_at, ...
|
||||
_, filename, _, email, _, received_at, *_ = task_args
|
||||
|
||||
start_time = time.time()
|
||||
parts = urlsplit(task_payload.cloud_storage_url)
|
||||
clean_url = urlunsplit((parts.scheme, parts.netloc, parts.path, "", ""))
|
||||
initial_metadata = {
|
||||
"start_time": start_time,
|
||||
"asr_model": settings.whisperx_asr_model,
|
||||
"retries": 0,
|
||||
"filename": filename,
|
||||
"email": email,
|
||||
"queuing_time": round(start_time - received_at, 2),
|
||||
"filename": clean_url,
|
||||
"sub": task_payload.user_sub,
|
||||
"email": task_payload.user_email,
|
||||
"tenant_id": task_payload.tenant_id,
|
||||
"queuing_time": round(start_time - task_payload.received_at.timestamp(), 2),
|
||||
}
|
||||
|
||||
self._save_metadata(task_id, initial_metadata)
|
||||
|
||||
@@ -4,7 +4,6 @@
|
||||
|
||||
import json
|
||||
import time
|
||||
from datetime import datetime
|
||||
|
||||
import openai
|
||||
import sentry_sdk
|
||||
@@ -14,10 +13,12 @@ from requests import exceptions
|
||||
|
||||
from summary.core.analytics import MetadataManager, get_analytics
|
||||
from summary.core.config import get_settings
|
||||
from summary.core.docs_service import create_document_in_docs
|
||||
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,
|
||||
)
|
||||
@@ -43,7 +44,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()
|
||||
@@ -79,23 +79,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 +100,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"]})
|
||||
@@ -149,11 +142,7 @@ def transcribe_audio(
|
||||
cloud_storage_url.split("?", 1)[0] if cloud_storage_url else None
|
||||
)
|
||||
logger.exception(
|
||||
(
|
||||
"Unexpected error while preparing file | filename: %s "
|
||||
"| cloud_storage_url: %s"
|
||||
),
|
||||
recording_filename,
|
||||
("Unexpected error while preparing file %s "),
|
||||
redacted_cloud_storage_url,
|
||||
)
|
||||
return None
|
||||
@@ -163,41 +152,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.started_at,
|
||||
recording_metadata.ended_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.started_at,
|
||||
recording_metadata.ended_at,
|
||||
)
|
||||
new_transcription = speaker_mapping.apply_to(transcription.model_dump())
|
||||
return new_transcription
|
||||
@@ -225,11 +204,8 @@ 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]:
|
||||
) -> str:
|
||||
"""Format a transcription into readable content with a title.
|
||||
|
||||
Resolves the locale from context_language / language, then uses
|
||||
@@ -242,9 +218,6 @@ def format_transcript(
|
||||
|
||||
return formatter.format(
|
||||
transcription,
|
||||
room=room,
|
||||
recording_datetime=recording_datetime,
|
||||
owner_timezone=owner_timezone,
|
||||
download_link=download_link,
|
||||
)
|
||||
|
||||
@@ -267,130 +240,9 @@ 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)
|
||||
|
||||
|
||||
def summarize_transcription_internals(
|
||||
*, owner_id: str, transcript: str, session_id: str
|
||||
*, user_sub: str, transcript: str, session_id: str
|
||||
) -> str:
|
||||
"""Generate a summary from the provided transcription text.
|
||||
|
||||
@@ -401,11 +253,11 @@ def summarize_transcription_internals(
|
||||
"""
|
||||
logger.info(
|
||||
"Starting summarization task | Owner: %s",
|
||||
owner_id,
|
||||
user_sub,
|
||||
)
|
||||
|
||||
user_has_tracing_consent = analytics.is_feature_enabled(
|
||||
"summary-tracing-consent", distinct_id=owner_id
|
||||
"summary-tracing-consent", distinct_id=user_sub
|
||||
)
|
||||
|
||||
# NOTE: We must instantiate a new LLMObservability client for each task invocation
|
||||
@@ -416,7 +268,7 @@ def summarize_transcription_internals(
|
||||
llm_observability = LLMObservability(
|
||||
user_has_tracing_consent=user_has_tracing_consent,
|
||||
session_id=session_id,
|
||||
user_id=owner_id,
|
||||
user_id=user_sub,
|
||||
)
|
||||
llm_service = LLMService(llm_observability=llm_observability)
|
||||
|
||||
@@ -468,29 +320,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 +378,64 @@ 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}")
|
||||
|
||||
# We do it synchronously for now
|
||||
if (
|
||||
payload.push_to_docs_config
|
||||
and settings.is_docs_integration_enabled
|
||||
and settings.get_authorized_tenant(
|
||||
tenant_id=payload.tenant_id
|
||||
).allowed_push_to_docs
|
||||
):
|
||||
# Format output
|
||||
content = format_transcript(
|
||||
transcription_res,
|
||||
payload.context_language,
|
||||
payload.language,
|
||||
payload.push_to_docs_config.download_link,
|
||||
)
|
||||
|
||||
create_document_in_docs(
|
||||
content=content,
|
||||
title=payload.push_to_docs_config.title,
|
||||
email=payload.push_to_docs_config.user_email,
|
||||
sub=payload.user_sub,
|
||||
)
|
||||
|
||||
if (
|
||||
payload.push_to_docs_config.auto_create_summary
|
||||
and analytics.is_feature_enabled(
|
||||
"summary-enabled", distinct_id=payload.user_sub
|
||||
)
|
||||
and settings.is_summary_enabled
|
||||
):
|
||||
summary = summarize_transcription_internals(
|
||||
user_sub=payload.user_sub,
|
||||
transcript=content,
|
||||
session_id=self.request.id,
|
||||
)
|
||||
locale = get_locale(payload.context_language, payload.language)
|
||||
create_document_in_docs(
|
||||
content=summary,
|
||||
title=locale.summary_title_template.format(
|
||||
title=payload.push_to_docs_config.title
|
||||
),
|
||||
email=payload.push_to_docs_config.user_email,
|
||||
sub=payload.user_sub,
|
||||
)
|
||||
|
||||
metadata_manager.capture(job_id, settings.posthog_event_success)
|
||||
|
||||
file_service.store_transcript(
|
||||
transcript=transcription_res,
|
||||
job_id=job_id,
|
||||
@@ -561,9 +448,30 @@ 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()
|
||||
|
||||
|
||||
|
||||
@signals.task_prerun.connect(sender=process_audio_transcribe_v2_task)
|
||||
def task_started(task_id=None, task=None, args=None, **kwargs):
|
||||
"""Signal handler called before task execution begins."""
|
||||
if args:
|
||||
metadata_manager.create(task_id, TranscribeTaskV2Payload.model_validate(args[0]))
|
||||
|
||||
|
||||
@signals.task_retry.connect(sender=process_audio_transcribe_v2_task)
|
||||
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_v2_task)
|
||||
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_failure.connect(sender=process_audio_transcribe_v2_task)
|
||||
def handle_transcribe_v2_failed(
|
||||
sender,
|
||||
@@ -621,7 +529,7 @@ def summarize_v2_task(
|
||||
"""
|
||||
payload = SummarizeTaskV2Payload.model_validate(payload)
|
||||
summary = summarize_transcription_internals(
|
||||
owner_id=payload.user_sub,
|
||||
user_sub=payload.user_sub,
|
||||
transcript=payload.content,
|
||||
session_id=self.request.id,
|
||||
)
|
||||
|
||||
@@ -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, Mapping, Optional, Set
|
||||
|
||||
from fastapi import Depends
|
||||
from pydantic import (
|
||||
@@ -34,9 +33,12 @@ class AuthorizedTenant(BaseModel):
|
||||
title="Webhook API Key",
|
||||
description="The api_key to authenticate the webhook request.",
|
||||
)
|
||||
|
||||
|
||||
V1_DEFAULT_TENANT_ID = "__deprecated_meet_tenant__"
|
||||
allowed_push_to_docs: bool = Field(
|
||||
title="Allow Push to Docs",
|
||||
description="Whether to allow pushing transcript"
|
||||
" and summaries to docs for this tenant.",
|
||||
default=False,
|
||||
)
|
||||
|
||||
|
||||
class Settings(BaseSettings):
|
||||
@@ -45,14 +47,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 +72,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,15 +112,16 @@ 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
|
||||
|
||||
# Docs service configuration
|
||||
is_docs_integration_enabled: bool = True
|
||||
docs_base_url: str = "https://example.com"
|
||||
docs_server_to_server_api_key: SecretStr = Field(
|
||||
title="API key for using docs server to server api", default="NO_API_KEY"
|
||||
)
|
||||
|
||||
# Sentry
|
||||
sentry_is_enabled: bool = False
|
||||
sentry_dsn: Optional[str] = None
|
||||
@@ -145,33 +144,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 +162,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."""
|
||||
|
||||
@@ -0,0 +1,78 @@
|
||||
"""Service for delivering content to external destinations."""
|
||||
|
||||
import json
|
||||
import logging
|
||||
|
||||
from requests import Session
|
||||
from requests.adapters import HTTPAdapter
|
||||
from urllib3.util import Retry
|
||||
|
||||
from summary.core.config import get_settings
|
||||
|
||||
settings = get_settings()
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
def _create_retry_session(api_key: str | None = None):
|
||||
"""Create an HTTP session configured with retry logic."""
|
||||
session = Session()
|
||||
retries = Retry(
|
||||
total=settings.webhook_max_retries,
|
||||
backoff_factor=settings.webhook_backoff_factor,
|
||||
status_forcelist=settings.webhook_status_forcelist,
|
||||
allowed_methods={"POST"},
|
||||
)
|
||||
session.mount("https://", HTTPAdapter(max_retries=retries))
|
||||
if api_key:
|
||||
session.headers.update({"Authorization": f"Bearer {api_key}"})
|
||||
|
||||
return session
|
||||
|
||||
|
||||
def _post_with_retries(*, url, data, api_key: str | None = None):
|
||||
"""Send POST request with automatic retries."""
|
||||
session = _create_retry_session(api_key=api_key)
|
||||
|
||||
try:
|
||||
response = session.post(url, json=data, timeout=(20, 3 * 60))
|
||||
response.raise_for_status()
|
||||
return response
|
||||
finally:
|
||||
session.close()
|
||||
|
||||
|
||||
def create_document_in_docs(*, content: str, title: str, email: str, sub: str) -> None:
|
||||
"""Call the Docs API to create a document on behalf of the user there.
|
||||
|
||||
Builds the payload, sends it with retries, and logs the outcome.
|
||||
"""
|
||||
data = {
|
||||
"title": title,
|
||||
"content": content,
|
||||
"email": email,
|
||||
"sub": sub,
|
||||
}
|
||||
|
||||
logger.debug("Submitting to %s", settings.docs_base_url)
|
||||
logger.debug("Request payload: %s", json.dumps(data, indent=2))
|
||||
|
||||
response = _post_with_retries(
|
||||
url=settings.docs_base_url,
|
||||
api_key=settings.docs_server_to_server_api_key.get_secret_value(),
|
||||
data=data,
|
||||
)
|
||||
|
||||
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)
|
||||
@@ -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,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)
|
||||
|
||||
duration = self._validate_duration(downloaded_path)
|
||||
|
||||
|
||||
@@ -30,4 +30,5 @@ Einige Punkte, die wir Ihnen empfehlen zu überprüfen:
|
||||
document_title_template=(
|
||||
'Besprechung "{room}" am {room_recording_date} um {room_recording_time}'
|
||||
),
|
||||
summary_title_template="Zusammenfassung von {title}",
|
||||
)
|
||||
|
||||
@@ -30,4 +30,5 @@ A few things we recommend you check:
|
||||
document_title_template=(
|
||||
'Meeting "{room}" on {room_recording_date} at {room_recording_time}'
|
||||
),
|
||||
summary_title_template="Summary of {title}",
|
||||
)
|
||||
|
||||
@@ -30,4 +30,5 @@ Quelques points que nous vous conseillons de vérifier :
|
||||
document_title_template=(
|
||||
'Réunion "{room}" du {room_recording_date} à {room_recording_time}'
|
||||
),
|
||||
summary_title_template="Résumé de {title}",
|
||||
)
|
||||
|
||||
@@ -30,4 +30,5 @@ Een paar punten die wij u aanraden te controleren:
|
||||
document_title_template=(
|
||||
'Vergadering "{room}" op {room_recording_date} om {room_recording_time}'
|
||||
),
|
||||
summary_title_template="Samenvatting van {title}",
|
||||
)
|
||||
|
||||
@@ -13,3 +13,4 @@ class LocaleStrings:
|
||||
hallucination_replacement_text: str
|
||||
document_default_title: str
|
||||
document_title_template: str
|
||||
summary_title_template: str
|
||||
|
||||
@@ -1,6 +1,7 @@
|
||||
"""Models for the API & Celery tasks creation."""
|
||||
from datetime import datetime
|
||||
|
||||
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
|
||||
@@ -12,6 +13,47 @@ class SharedV2TaskCreation(BaseModel):
|
||||
"""Model that holds basic information for task creation."""
|
||||
|
||||
user_sub: str = Field(title="User Sub", description="The user's sub.")
|
||||
user_email: str | None = Field(
|
||||
title="User Email", description="The user's email for analytics purposes."
|
||||
)
|
||||
|
||||
|
||||
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.",
|
||||
)
|
||||
started_at: AwareDatetime = Field(title="Start time of the recording to transcribe")
|
||||
ended_at: AwareDatetime = Field(title="End time of the recording to transcribe")
|
||||
|
||||
|
||||
class PushToDocsBaseConfig(BaseModel):
|
||||
"""Model containing information for pushing transcript and summaries to docs."""
|
||||
|
||||
user_email: str = Field(
|
||||
title="User Email", description="The user's email, future owner of the docs."
|
||||
)
|
||||
title: str = Field(title="Title", description="The title for the created document.")
|
||||
|
||||
|
||||
class PushToDocsTranscriptConfig(PushToDocsBaseConfig):
|
||||
"""Model for push to docs information for transcripts."""
|
||||
|
||||
download_link: str | None = Field(
|
||||
title="Download Link", description="The link to download the recording."
|
||||
)
|
||||
auto_create_summary: bool = Field(
|
||||
title="Auto Create Summary Docs",
|
||||
description="Whether to automatically create a summary "
|
||||
"for the transcription task and push it to docs.",
|
||||
default=False,
|
||||
)
|
||||
|
||||
|
||||
class PushToDocsSummaryConfig(PushToDocsBaseConfig):
|
||||
"""Model for push to docs information for summaries."""
|
||||
|
||||
|
||||
class TranscribeTaskV2Request(SharedV2TaskCreation):
|
||||
@@ -27,7 +69,17 @@ 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,
|
||||
)
|
||||
push_to_docs_config: PushToDocsTranscriptConfig | None = Field(
|
||||
title="Push to Docs info",
|
||||
description="If set, configuration for pushing to docs",
|
||||
default=None,
|
||||
)
|
||||
|
||||
@field_validator("language")
|
||||
@@ -46,6 +98,7 @@ class TranscribeTaskV2Payload(TranscribeTaskV2Request):
|
||||
"""Model for creating a transcribe and summarize task (used for actual task creation).""" # noqa: E501
|
||||
|
||||
tenant_id: str = Field(title="Tenant ID", description="The ID of the tenant.")
|
||||
received_at: datetime = Field(title="Received At", description="The time the task was received.")
|
||||
|
||||
|
||||
class SummarizeTaskV2Request(SharedV2TaskCreation):
|
||||
@@ -58,3 +111,4 @@ class SummarizeTaskV2Payload(SummarizeTaskV2Request):
|
||||
"""Model for creating a summarize task (used for actual task creation)."""
|
||||
|
||||
tenant_id: str = Field(title="Tenant ID", description="The ID of the tenant.")
|
||||
received_at: datetime = Field(title="Received At", description="The time the task was received.")
|
||||
|
||||
@@ -159,4 +159,5 @@ __all__ = [
|
||||
"SummarizeWebhookPayloads",
|
||||
"WebhookPayloads",
|
||||
"WhisperXResponse",
|
||||
"webhook_payload_adapter",
|
||||
]
|
||||
|
||||
@@ -1,9 +1,6 @@
|
||||
"""Transcript formatting into readable conversation format with speaker labels."""
|
||||
|
||||
import logging
|
||||
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
|
||||
@@ -41,12 +38,9 @@ class TranscriptFormatter:
|
||||
def format(
|
||||
self,
|
||||
transcription,
|
||||
room: str | None = None,
|
||||
recording_datetime: str | None = None,
|
||||
owner_timezone: str | None = None,
|
||||
download_link: str | None = None,
|
||||
) -> Tuple[str, str]:
|
||||
"""Format transcription into the final document and its title."""
|
||||
) -> str:
|
||||
"""Format transcription into the final document."""
|
||||
segments = self._get_segments(transcription)
|
||||
|
||||
if not segments:
|
||||
@@ -56,9 +50,7 @@ class TranscriptFormatter:
|
||||
content = self._remove_hallucinations(content)
|
||||
content = self._add_header(content, download_link)
|
||||
|
||||
title = self._generate_title(room, recording_datetime, owner_timezone)
|
||||
|
||||
return content, title
|
||||
return content
|
||||
|
||||
def _remove_hallucinations(self, content: str) -> str:
|
||||
"""Remove hallucination patterns from content."""
|
||||
@@ -96,23 +88,3 @@ class TranscriptFormatter:
|
||||
content = header + content
|
||||
|
||||
return content
|
||||
|
||||
def _generate_title(
|
||||
self,
|
||||
room: str | None = None,
|
||||
recording_datetime: str | None = None,
|
||||
owner_timezone: str | None = None,
|
||||
) -> str:
|
||||
"""Generate title from context or return default."""
|
||||
if not room or not recording_datetime:
|
||||
return self._locale.document_default_title
|
||||
|
||||
dt = datetime.fromisoformat(recording_datetime)
|
||||
if owner_timezone:
|
||||
dt = dt.astimezone(ZoneInfo(owner_timezone))
|
||||
|
||||
return self._locale.document_title_template.format(
|
||||
room=room,
|
||||
room_recording_date=dt.strftime("%Y-%m-%d"),
|
||||
room_recording_time=dt.strftime("%H:%M"),
|
||||
)
|
||||
|
||||
@@ -4,9 +4,6 @@ import json
|
||||
import logging
|
||||
|
||||
import requests
|
||||
from requests import Session
|
||||
from requests.adapters import HTTPAdapter
|
||||
from urllib3.util import Retry
|
||||
|
||||
from summary.core.config import get_settings
|
||||
from summary.core.shared_models import (
|
||||
@@ -18,83 +15,6 @@ settings = get_settings()
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
def _create_retry_session(api_key: str | None = None):
|
||||
"""Create an HTTP session configured with retry logic."""
|
||||
session = Session()
|
||||
retries = Retry(
|
||||
total=settings.webhook_max_retries,
|
||||
backoff_factor=settings.webhook_backoff_factor,
|
||||
status_forcelist=settings.webhook_status_forcelist,
|
||||
allowed_methods={"POST"},
|
||||
)
|
||||
session.mount("https://", HTTPAdapter(max_retries=retries))
|
||||
if api_key:
|
||||
session.headers.update({"Authorization": f"Bearer {api_key}"})
|
||||
|
||||
return session
|
||||
|
||||
|
||||
def _post_with_retries(*, url, data, api_key: str | None = None):
|
||||
"""Send POST request with automatic retries."""
|
||||
session = _create_retry_session(api_key=api_key)
|
||||
|
||||
try:
|
||||
response = session.post(url, json=data)
|
||||
response.raise_for_status()
|
||||
return response
|
||||
finally:
|
||||
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,
|
||||
|
||||
@@ -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)
|
||||
|
||||
Generated
+1475
File diff suppressed because it is too large
Load Diff
Reference in New Issue
Block a user