From e7f15b50ff1c721e2ecf4d6b5046215a6d8eda83 Mon Sep 17 00:00:00 2001 From: Florent Chehab Date: Tue, 7 Jul 2026 14:45:25 +0200 Subject: [PATCH] =?UTF-8?q?=E2=9C=A8(summary)=20more=20precise=20analytics?= =?UTF-8?q?=20events?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit * Specific transcript and summary events * Improve observability on summary tasks --- CHANGELOG.md | 1 + src/summary/summary/api/route/tasks_v2.py | 6 ++--- src/summary/summary/core/analytics.py | 17 ++++++++----- src/summary/summary/core/celery_worker.py | 31 +++++++++++++++++++---- src/summary/summary/core/config.py | 9 ++++--- 5 files changed, 46 insertions(+), 18 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 38e87612..30772d3d 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -28,6 +28,7 @@ and this project adheres to - ⬆️(frontend) upgrade @tanstack/react-query from 5.100.14 to 5.101.0 - ⬆️(frontend) update the frontend build image to Node 22 - 🔒️(frontend) update docker image to nginx-unprivileged:1.30.3-alpine3.23 +- ✨(summary) more precise analytics events ### Fixed diff --git a/src/summary/summary/api/route/tasks_v2.py b/src/summary/summary/api/route/tasks_v2.py index eebe88e8..ec175160 100644 --- a/src/summary/summary/api/route/tasks_v2.py +++ b/src/summary/summary/api/route/tasks_v2.py @@ -62,10 +62,9 @@ async def create_transcribe_task_v2( # We track the request, this also properly initializes the user in the # analytics system, so that later feature flags work properly analytics.capture( - settings.posthog_event_request, + settings.posthog_transcript_request, request.user_sub, properties={ - "kind": "transcribe", "$set": { "email": request.user_email, }, @@ -94,10 +93,9 @@ async def create_summarize_task_v2( # We track the request, this also properly initializes the user in the # analytics system, so that later feature flags work properly analytics.capture( - settings.posthog_event_request, + settings.posthog_summary_request, request.user_sub, properties={ - "kind": "summarize", "$set": { "email": request.user_email, }, diff --git a/src/summary/summary/core/analytics.py b/src/summary/summary/core/analytics.py index fb52ef8a..2fbcff88 100644 --- a/src/summary/summary/core/analytics.py +++ b/src/summary/summary/core/analytics.py @@ -11,7 +11,7 @@ from celery.utils.log import get_task_logger from posthog import Posthog from summary.core.config import get_settings -from summary.core.models import TranscribeTaskJob +from summary.core.models import SummarizeTaskJob, TranscribeTaskJob logger = get_task_logger(__name__) settings = get_settings() @@ -109,19 +109,15 @@ 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: str, task_payload: TranscribeTaskJob): + def create(self, task_id: str, task_payload: TranscribeTaskJob | SummarizeTaskJob): """Create initial metadata entry for a new task.""" if self._is_disabled or self.has_task_id(task_id): return 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": clean_url, "sub": task_payload.user_sub, # avoid None in redis, it shouldn't happen anyway in prod "email": task_payload.user_email or "", @@ -129,6 +125,15 @@ class MetadataManager: "queuing_time": round(start_time - task_payload.received_at.timestamp(), 2), } + if isinstance(task_payload, TranscribeTaskJob): + parts = urlsplit(task_payload.cloud_storage_url) + clean_url = urlunsplit((parts.scheme, parts.netloc, parts.path, "", "")) + initial_metadata["source_url"] = clean_url + initial_metadata["asr_model"] = settings.whisperx_asr_model + elif isinstance(task_payload, SummarizeTaskJob): + initial_metadata["content_length"] = len(task_payload.content) + initial_metadata["llm_model"] = settings.llm_model + self._save_metadata(task_id, initial_metadata) def retry(self, task_id): diff --git a/src/summary/summary/core/celery_worker.py b/src/summary/summary/core/celery_worker.py index 867cbcf0..cbfa0bcf 100644 --- a/src/summary/summary/core/celery_worker.py +++ b/src/summary/summary/core/celery_worker.py @@ -542,28 +542,28 @@ 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) + metadata_manager.capture(job_id, settings.posthog_transcript_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): +def task_started_transcript(task_id=None, task=None, args=None, **kwargs): """Signal handler called before task execution begins.""" if args: metadata_manager.create(task_id, TranscribeTaskJob.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): +def task_retry_handler_transcript(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): +def task_failure_handler_transcript(task_id, exception=None, **kwargs): """Signal handler called when task execution fails permanently.""" - metadata_manager.capture(task_id, settings.posthog_event_failure) + metadata_manager.capture(task_id, settings.posthog_transcript_failure) @signals.task_failure.connect(sender=process_audio_transcribe_v2_task) @@ -648,9 +648,30 @@ def summarize_v2_task( call_webhook_v2_task.apply_async( args=[success_payload.model_dump(), payload.tenant_id] ) + metadata_manager.capture(job_id, settings.posthog_summary_success) + return success_payload.model_dump() +@signals.task_prerun.connect(sender=summarize_v2_task) +def task_started_summary(task_id=None, task=None, args=None, **kwargs): + """Signal handler called before task execution begins.""" + if args: + metadata_manager.create(task_id, SummarizeTaskJob.model_validate(args[0])) + + +@signals.task_retry.connect(sender=summarize_v2_task) +def task_retry_handler_summary(request=None, reason=None, einfo=None, **kwargs): + """Signal handler called when task execution retries.""" + metadata_manager.retry(request.id) + + +@signals.task_failure.connect(sender=summarize_v2_task) +def task_failure_handler_summary(task_id, exception=None, **kwargs): + """Signal handler called when task execution fails permanently.""" + metadata_manager.capture(task_id, settings.posthog_summary_failure) + + @signals.task_failure.connect(sender=summarize_v2_task) def handle_summarize_v2_failed( sender, diff --git a/src/summary/summary/core/config.py b/src/summary/summary/core/config.py index 8ee4eabf..4afb377e 100644 --- a/src/summary/summary/core/config.py +++ b/src/summary/summary/core/config.py @@ -130,9 +130,12 @@ class Settings(BaseSettings): posthog_enabled: bool = False posthog_api_key: Optional[str] = None posthog_api_host: Optional[str] = "https://eu.i.posthog.com" - posthog_event_failure: str = "transcript-failure" - posthog_event_success: str = "transcript-success" - posthog_event_request: str = "transcript-request" + posthog_transcript_request: str = "transcript-request" + posthog_transcript_failure: str = "transcript-failure" + posthog_transcript_success: str = "transcript-success" + posthog_summary_request: str = "summary-request" + posthog_summary_failure: str = "summary-failure" + posthog_summary_success: str = "summary-success" # Langfuse (LLM Observability) langfuse_enabled: bool = False