Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
92 changes: 49 additions & 43 deletions src/youtube_extension/backend/api/v1/router.py
Original file line number Diff line number Diff line change
Expand Up @@ -49,6 +49,7 @@
from youtube_extension.services.workflows.transcript_action_workflow import (
TranscriptActionWorkflow,
)
from youtube_extension.utils.logsafe import safe_log as _safe_log

# CloudEvents integration (optional — falls back to file sink)
try:
Expand Down Expand Up @@ -119,6 +120,11 @@

logger = logging.getLogger(__name__)

# ``_safe_log`` (imported at the top of this module from
# ``youtube_extension.utils.logsafe``) strips log-forging characters from any
# user-controlled value — including exception text — before it is logged.
Comment thread
groupthinking marked this conversation as resolved.



async def _emit_event(event_type: str, data: dict, subject: str | None = None) -> None:
"""Emit a CloudEvent if the publisher is available."""
Expand All @@ -131,7 +137,7 @@ async def _emit_event(event_type: str, data: dict, subject: str | None = None) -
subject=subject,
)
except Exception as exc:
logger.debug("CloudEvent publish failed: %s", exc)
logger.debug("CloudEvent publish failed: %s", _safe_log(exc))


def _normalize_tag_list(raw_tags: Any) -> list[str]:
Expand Down Expand Up @@ -257,7 +263,7 @@ async def health_check_v1(
)
return HealthResponse(**health_status)
except Exception as e:
logger.error(f"Health check failed: {e}", exc_info=True)
logger.error(f"Health check failed: {_safe_log(e)}", exc_info=True)
raise HTTPException(status_code=500, detail="Internal server error")


Expand Down Expand Up @@ -286,7 +292,7 @@ async def detailed_health_check_v1(
"timestamp": datetime.now().isoformat(),
}
except Exception as e:
logger.error(f"Detailed health check failed: {e}", exc_info=True)
logger.error(f"Detailed health check failed: {_safe_log(e)}", exc_info=True)
raise HTTPException(status_code=500, detail="Internal server error")


Expand Down Expand Up @@ -351,7 +357,7 @@ async def get_capabilities_v1(
],
}
except Exception as e:
logger.error(f"Capabilities check failed: {e}")
logger.error(f"Capabilities check failed: {_safe_log(e)}")
return {"status": "error", "error": str(e)}


Expand Down Expand Up @@ -581,7 +587,7 @@ async def chat_v1(
"""Chat endpoint with AI processing via AgentOrchestrator"""
try:
logger.info(
f"Chat request received: {request.message[:50]}... session={request.session_id}"
f"Chat request received: {_safe_log(request.message[:50])}... session={_safe_log(request.session_id)}"
)

params = {
Expand All @@ -602,13 +608,13 @@ async def chat_v1(
video_id = match.group(1)

if video_id:
logger.info(f"Adding video context for video_id: {video_id}")
logger.info(f"Adding video context for video_id: {_safe_log(video_id)}")
Comment thread
groupthinking marked this conversation as resolved.
detail = data_service.get_video_detail(video_id)

# If video not found, trigger real-time processing
if not detail and request.video_url:
logger.info(
f"Video not found for {video_id}, triggering real-time processing"
f"Video not found for {_safe_log(video_id)}, triggering real-time processing"
)
try:
proc_result = (
Expand All @@ -618,9 +624,9 @@ async def chat_v1(
)
if proc_result and proc_result.get("status") == "success":
detail = data_service.get_video_detail(video_id)
logger.info(f"Real-time processing complete for {video_id}")
logger.info(f"Real-time processing complete for {_safe_log(video_id)}")
except Exception as e:
logger.error(f"Real-time video processing failed: {e}")
logger.error(f"Real-time video processing failed: {_safe_log(e)}")

if detail:
params["video_id"] = video_id
Expand Down Expand Up @@ -665,7 +671,7 @@ async def chat_v1(
return response

except Exception as e:
logger.error(f"Error in chat endpoint: {e}", exc_info=True)
logger.error(f"Error in chat endpoint: {_safe_log(e)}", exc_info=True)
raise HTTPException(status_code=500, detail="Internal server error") from e


Expand All @@ -683,7 +689,7 @@ async def process_video_v1(
):
"""Basic video processing endpoint"""
try:
logger.info(f"Video processing request: {request.video_url}")
logger.info(f"Video processing request: {_safe_log(request.video_url)}")
await _emit_event(
"com.eventrelay.video.received",
{"url": request.video_url},
Expand Down Expand Up @@ -719,7 +725,7 @@ async def process_video_v1(
return result

except Exception as e:
logger.error(f"Error in video processing: {e}")
logger.error(f"Error in video processing: {_safe_log(e)}")
await _emit_event(
"com.eventrelay.pipeline.failed",
{"url": request.video_url, "error": str(e)},
Expand Down Expand Up @@ -752,7 +758,7 @@ async def process_video_markdown_v1(
health_service.increment_metric("process_video_markdown_total")

try:
logger.info(f"Markdown processing request: {request.video_url}")
logger.info(f"Markdown processing request: {_safe_log(request.video_url)}")

result = await video_processing_service.process_video_for_markdown(
request.video_url, request.force_regenerate
Expand All @@ -769,7 +775,7 @@ async def process_video_markdown_v1(
health_service.increment_metric("error_total")
raise
except Exception as e:
logger.error(f"Error in markdown processing: {e}", exc_info=True)
logger.error(f"Error in markdown processing: {_safe_log(e)}", exc_info=True)
health_service.increment_metric("error_total")
raise HTTPException(status_code=500, detail="Internal server error")

Expand All @@ -788,7 +794,7 @@ async def video_to_software_v1(
):
"""Convert YouTube video to deployed software"""
try:
logger.info(f"Video-to-software request: {request.video_url}")
logger.info(f"Video-to-software request: {_safe_log(request.video_url)}")
target_info = resolve_deployment_target(request.deployment_target)

result = await video_processing_service.process_video_to_software(
Expand All @@ -807,7 +813,7 @@ async def video_to_software_v1(
return VideoToSoftwareResponse(**result)

except Exception as e:
logger.error(f"Video-to-software processing failed: {e}", exc_info=True)
logger.error(f"Video-to-software processing failed: {_safe_log(e)}", exc_info=True)
raise HTTPException(status_code=500, detail="Internal server error")


Expand All @@ -831,7 +837,7 @@ async def get_cache_stats_v1(cache_service: CacheService = Depends(get_cache_ser
_stats_cache_time = now
return CacheStats(**stats)
except Exception as e:
logger.error(f"Error getting cache stats: {e}", exc_info=True)
logger.error(f"Error getting cache stats: {_safe_log(e)}", exc_info=True)
raise HTTPException(status_code=500, detail="Internal server error")


Expand Down Expand Up @@ -860,7 +866,7 @@ async def get_cached_video_v1(
except HTTPException:
raise
except Exception as e:
logger.error(f"Error retrieving cached video: {e}", exc_info=True)
logger.error(f"Error retrieving cached video: {_safe_log(e)}", exc_info=True)
raise HTTPException(status_code=500, detail="Internal server error")


Expand All @@ -883,7 +889,7 @@ async def clear_video_cache_v1(
"timestamp": datetime.now().isoformat(),
}
except Exception as e:
logger.error(f"Error clearing video cache: {e}", exc_info=True)
logger.error(f"Error clearing video cache: {_safe_log(e)}", exc_info=True)
raise HTTPException(status_code=500, detail="Internal server error")


Expand All @@ -903,7 +909,7 @@ async def clear_all_cache_v1(cache_service: CacheService = Depends(get_cache_ser
"timestamp": datetime.now().isoformat(),
}
except Exception as e:
logger.error(f"Error clearing all cache: {e}", exc_info=True)
logger.error(f"Error clearing all cache: {_safe_log(e)}", exc_info=True)
raise HTTPException(status_code=500, detail="Internal server error")


Expand Down Expand Up @@ -941,7 +947,7 @@ async def list_videos_v1(
}

except Exception as e:
logger.error(f"Error listing videos: {e}", exc_info=True)
logger.error(f"Error listing videos: {_safe_log(e)}", exc_info=True)
raise HTTPException(status_code=500, detail="Internal server error")


Expand All @@ -965,7 +971,7 @@ async def get_video_detail_v1(
except HTTPException:
raise
except Exception as e:
logger.error(f"Error getting video detail: {e}", exc_info=True)
logger.error(f"Error getting video detail: {_safe_log(e)}", exc_info=True)
raise HTTPException(status_code=500, detail="Internal server error")


Expand All @@ -981,7 +987,7 @@ async def get_learning_log_v1(data_service: DataService = Depends(get_data_servi
learning_log = data_service.get_learning_log()
return learning_log
except Exception as e:
logger.error(f"Error getting learning log: {e}", exc_info=True)
logger.error(f"Error getting learning log: {_safe_log(e)}", exc_info=True)
raise HTTPException(status_code=500, detail="Internal server error")


Expand Down Expand Up @@ -1016,7 +1022,7 @@ async def ingest_knowledge_v1(
except HTTPException:
raise
except Exception as exc:
logger.error(f"Error ingesting knowledge entry: {exc}")
logger.error(f"Error ingesting knowledge entry: {_safe_log(exc)}")
raise HTTPException(status_code=500, detail="Failed to store insight")


Expand All @@ -1032,7 +1038,7 @@ async def get_actions_by_video_v1(video_id: str):
actions = repo.get_by_video_id(video_id)
return actions
except Exception as e:
logger.error(f"Error retrieving actions for {video_id}: {e}", exc_info=True)
logger.error(f"Error retrieving actions for {_safe_log(video_id)}: {_safe_log(e)}", exc_info=True)
raise HTTPException(status_code=500, detail="Internal server error")


Expand Down Expand Up @@ -1080,7 +1086,7 @@ async def update_action_v1(action_id: str, payload: dict[str, Any]):
logger.debug("Action feedback recording failed", exc_info=True)
return {"success": bool(success)}
except Exception as e:
logger.error(f"Error updating action {action_id}: {e}", exc_info=True)
logger.error(f"Error updating action {_safe_log(action_id)}: {_safe_log(e)}", exc_info=True)
raise HTTPException(status_code=500, detail="Internal server error")


Expand Down Expand Up @@ -1127,7 +1133,7 @@ async def submit_feedback_v1(
raise HTTPException(status_code=500, detail="Failed to save feedback")

except Exception as e:
logger.error(f"Error saving feedback: {e}", exc_info=True)
logger.error(f"Error saving feedback: {_safe_log(e)}", exc_info=True)
raise HTTPException(status_code=500, detail="Internal server error")


Expand All @@ -1147,7 +1153,7 @@ async def get_metrics_v1(
metrics_lines = health_service.get_metrics_prometheus_format()
return Response(content="\n".join(metrics_lines), media_type="text/plain")
except Exception as e:
logger.error(f"Metrics endpoint failed: {e}", exc_info=True)
logger.error(f"Metrics endpoint failed: {_safe_log(e)}", exc_info=True)
raise HTTPException(status_code=500, detail="Internal server error")


Expand All @@ -1167,7 +1173,7 @@ async def ingest_performance_alert_v1(payload: dict[str, Any]):
)
return {"status": "ok", "recorded": metric_name}
except Exception as e:
logger.error(f"Failed to ingest performance alert: {e}", exc_info=True)
logger.error(f"Failed to ingest performance alert: {_safe_log(e)}", exc_info=True)
raise HTTPException(status_code=500, detail="Internal server error")


Expand All @@ -1185,7 +1191,7 @@ async def ingest_performance_report_v1(report: dict[str, Any]):
)
return {"status": "ok", "metrics_recorded": len(metrics)}
except Exception as e:
logger.error(f"Failed to ingest performance report: {e}", exc_info=True)
logger.error(f"Failed to ingest performance report: {_safe_log(e)}", exc_info=True)
raise HTTPException(status_code=500, detail="Internal server error")


Expand Down Expand Up @@ -1316,7 +1322,7 @@ async def _periodic_cleanup():
_agent_executions.evict_expired()
_dispatches.evict_expired()
except Exception as exc:
logger.debug("Periodic cleanup failed: %s", exc)
logger.debug("Periodic cleanup failed: %s", _safe_log(exc))
await asyncio.sleep(300) # Sweep every 5 minutes


Expand All @@ -1336,7 +1342,7 @@ def _sync_persist():
data = job.model_dump(mode="json")
get_job_store().save(job.job_id, data)
except Exception as exc:
logger.warning("Job persist failed for %s: %s", job.job_id, exc)
logger.warning("Job persist failed for %s: %s", _safe_log(job.job_id), _safe_log(exc))

# If we are in an async loop, offload serialization and I/O to a thread
try:
Expand Down Expand Up @@ -1433,8 +1439,8 @@ async def _queue_transcript_action_job(
except Exception as exc:
logger.info(
"Cloud Tasks unavailable for %s, using local background task: %s",
job_id,
exc,
_safe_log(job_id),
_safe_log(exc),
)
asyncio.create_task(
_run_video_job(
Expand Down Expand Up @@ -1522,7 +1528,7 @@ async def _run_video_job(
"""Background coroutine that drives the transcript-action workflow."""
job = _load_video_job(job_id)
if job is None:
logger.error("Video job %s missing at run time", job_id)
logger.error("Video job %s missing at run time", _safe_log(job_id))
return
try:
job.status = JobStatus.downloading
Expand Down Expand Up @@ -1568,7 +1574,7 @@ async def _run_video_job(
job.status = JobStatus.failed
job.error = str(exc)
_persist_video_job(job)
logger.error(f"Video job {job_id} failed: {exc}")
logger.error(f"Video job {_safe_log(job_id)} failed: {_safe_log(exc)}")


@router.post(
Expand Down Expand Up @@ -1729,7 +1735,7 @@ async def get_or_create_videopack(request: VideoPackRequest):
)
return ApiResponse.success(pack.model_dump())
except Exception as e:
logger.error(f"Failed to create VideoPack: {e}", exc_info=True)
logger.error(f"Failed to create VideoPack: {_safe_log(e)}", exc_info=True)
raise HTTPException(status_code=500, detail="Internal server error")


Expand Down Expand Up @@ -1778,7 +1784,7 @@ async def generate_blueprint(request: BlueprintRequest):
)
return ApiResponse.success(blueprint)
except Exception as e:
logger.error(f"Failed to generate blueprint: {e}", exc_info=True)
logger.error(f"Failed to generate blueprint: {_safe_log(e)}", exc_info=True)
raise HTTPException(status_code=500, detail="Internal server error")


Expand All @@ -1804,7 +1810,7 @@ async def generate_project_code(request: GenerateCodeRequest):
)
return ApiResponse.success(result)
except Exception as e:
logger.error(f"Code generation failed: {e}", exc_info=True)
logger.error(f"Code generation failed: {_safe_log(e)}", exc_info=True)
raise HTTPException(status_code=500, detail="Internal server error")


Expand Down Expand Up @@ -1924,7 +1930,7 @@ async def _extract_chunk(chunk: str) -> list[ExtractedEvent]:
)
)
except Exception as exc:
logger.warning(f"Direct Gemini extraction unavailable for chunk: {exc}")
logger.warning(f"Direct Gemini extraction unavailable for chunk: {_safe_log(exc)}")
return chunk_events

try:
Expand All @@ -1937,7 +1943,7 @@ async def _extract_chunk(chunk: str) -> list[ExtractedEvent]:
seen_titles.add(ev.title)
events.append(ev)
except Exception as exc:
logger.warning(f"Chunked extraction failed: {exc}")
logger.warning(f"Chunked extraction failed: {_safe_log(exc)}")

# Real-AI fallback: if no events yet, try the Vercel AI Gateway (uses
# VERCEL_API_KEY, routes to Gemini/GPT/Claude). This keeps the AI path
Expand All @@ -1963,7 +1969,7 @@ async def _extract_chunk(chunk: str) -> list[ExtractedEvent]:
"Extracted %d events via Vercel AI Gateway", len(gw_events)
)
except Exception as gw_exc: # noqa: BLE001
logger.warning(f"Vercel AI Gateway extraction failed: {gw_exc}")
logger.warning(f"Vercel AI Gateway extraction failed: {_safe_log(gw_exc)}")

if not events:
logger.warning("Falling back to heuristic extraction")
Expand Down Expand Up @@ -2102,7 +2108,7 @@ async def _run_agent(execution: AgentExecution, events: list[dict[str, Any]]):
except Exception as exc:
execution.status = AgentStatus.failed
execution.error = str(exc)
logger.error(f"Agent {execution.agent_id} failed: {exc}")
logger.error(f"Agent {_safe_log(execution.agent_id)} failed: {_safe_log(exc)}")


@router.get(
Expand Down
Loading
Loading