Ingest OpenTelemetry spans
(request: Request, supabase: AsyncSupabaseClient)
| 279 | |
| 280 | @router.post("/traces") |
| 281 | async def traces(request: Request, supabase: AsyncSupabaseClient): |
| 282 | """Ingest OpenTelemetry spans""" |
| 283 | try: |
| 284 | api_key = request.headers.get("X-Agentops-Auth") |
| 285 | |
| 286 | sessions, data = await asyncio.gather( |
| 287 | supabase.table("sessions").select("id").eq("api_key", api_key).limit(1).single().execute(), |
| 288 | request.json(), |
| 289 | ) |
| 290 | |
| 291 | session_id = data.get("session_id") |
| 292 | session_ids_for_project = [session["id"] for session in sessions] |
| 293 | if session_id not in session_ids_for_project: |
| 294 | raise RuntimeError("Invalid API Key for session") |
| 295 | |
| 296 | logger.debug(data) |
| 297 | |
| 298 | except RuntimeError as e: |
| 299 | message = {"message": f"/traces: Error processing spans: {e}"} |
| 300 | logger.error(message) |
| 301 | return JSONResponse(message, status_code=401) |
| 302 | |
| 303 | try: |
| 304 | spans_data = [] |
| 305 | for span in data.get("spans", []): |
| 306 | # Classify the span |
| 307 | span_type = await span_handlers.classify_span(span) |
| 308 | |
| 309 | # Route to the appropriate handler |
| 310 | if span_type == span_handlers.SESSION_UPDATE_SPAN: |
| 311 | span_data = await span_handlers.handle_session_update_span(span, session_id) |
| 312 | elif span_type == span_handlers.GEN_AI_SPAN: |
| 313 | span_data = await span_handlers.handle_gen_ai_span(span, session_id) |
| 314 | elif span_type == span_handlers.LOG_SPAN: |
| 315 | span_data = await span_handlers.handle_log_span(span, session_id) |
| 316 | else: |
| 317 | # Default to session update handler |
| 318 | span_data = await span_handlers.handle_session_update_span(span, session_id) |
| 319 | |
| 320 | spans_data.append(span_data) |
| 321 | |
| 322 | # Insert spans into the database |
| 323 | if spans_data: |
| 324 | await supabase.table("spans").upsert(spans_data).execute() |
| 325 | |
| 326 | logger.info(f"/traces: Completed POST request for {api_key}") |
| 327 | return JSONResponse({"status": "success"}) |
| 328 | except RuntimeError as e: |
| 329 | message = {"message": f"/traces: Error processing spans: {e}"} |
| 330 | logger.error(message) |
| 331 | return JSONResponse(message, status_code=400) |
| 332 | |
| 333 | |
| 334 | # @router.get("/openapi.yaml") |