mirror of
https://github.com/suitenumerique/meet.git
synced 2026-08-05 16:37:43 +00:00
✨(summary) more precise analytics events
* Specific transcript and summary events * Improve observability on summary tasks
This commit is contained in:
committed by
aleb_the_flash
parent
75ba2ff146
commit
e7f15b50ff
@@ -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
|
||||
|
||||
|
||||
@@ -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,
|
||||
},
|
||||
|
||||
@@ -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):
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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
|
||||
|
||||
Reference in New Issue
Block a user