From d0d4b9844c00e8a7b15c99a6a70cbda554a9d648 Mon Sep 17 00:00:00 2001 From: Claude Date: Tue, 21 Jul 2026 08:14:31 +0000 Subject: [PATCH 1/3] fix(security): sanitize user-controlled values in API logs (CWE-117 log injection) Rebuilt on current main. Adds a `_safe_log()` CR/LF scrubber and applies it to every user-controlled value interpolated into a v1-router log line, closing log-injection (CWE-117): a newline in a path param or request field could otherwise forge additional log entries. Sinks sanitized (14): chat `request.message`/`request.session_id`; video-context, video-not-found, and processing-complete `video_id`; video/markdown/ video-to-software `request.video_url`; action-retrieval `video_id`; action-update `action_id`; `job.job_id` (persist), `job_id` (cloud-tasks fallback, missing-at-runtime, job-failed); and `execution.agent_id`. An AST scan over the router confirms zero unsanitized user-controlled logger sinks remain. Adds `TestSafeLog` regression tests. Co-Authored-By: Claude Opus 4.8 (1M context) Claude-Session: https://claude.ai/code/session_01FZcDgrGTkknC2ya13Uy6bU --- .../backend/api/v1/router.py | 38 ++++++++++++------- tests/unit/test_v1_router_extended.py | 17 +++++++++ 2 files changed, 41 insertions(+), 14 deletions(-) diff --git a/src/youtube_extension/backend/api/v1/router.py b/src/youtube_extension/backend/api/v1/router.py index b3b4cbc65..d8a6d09e4 100644 --- a/src/youtube_extension/backend/api/v1/router.py +++ b/src/youtube_extension/backend/api/v1/router.py @@ -110,6 +110,16 @@ logger = logging.getLogger(__name__) +def _safe_log(value: object) -> str: + """Strip CR/LF from a value before it enters a log line. + + Guards against log-injection (CWE-117): user-controlled inputs (path + params, request fields) can smuggle newlines to forge additional log + entries. Sanitize such values wherever they are interpolated into logs. + """ + return str(value).replace("\r", "").replace("\n", "") + + async def _emit_event(event_type: str, data: dict, subject: str | None = None) -> None: """Emit a CloudEvent if the publisher is available.""" @@ -572,7 +582,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 = { @@ -593,13 +603,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)}") 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 = ( @@ -609,7 +619,7 @@ 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}") @@ -674,7 +684,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}, @@ -743,7 +753,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 @@ -779,7 +789,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( @@ -1023,7 +1033,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)}: {e}", exc_info=True) raise HTTPException(status_code=500, detail="Internal server error") @@ -1071,7 +1081,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)}: {e}", exc_info=True) raise HTTPException(status_code=500, detail="Internal server error") @@ -1327,7 +1337,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), exc) # If we are in an async loop, offload serialization and I/O to a thread try: @@ -1424,7 +1434,7 @@ async def _queue_transcript_action_job( except Exception as exc: logger.info( "Cloud Tasks unavailable for %s, using local background task: %s", - job_id, + _safe_log(job_id), exc, ) asyncio.create_task( @@ -1513,7 +1523,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 @@ -1559,7 +1569,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: {exc}") @router.post( @@ -2089,7 +2099,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: {exc}") @router.get( diff --git a/tests/unit/test_v1_router_extended.py b/tests/unit/test_v1_router_extended.py index cd484b1c3..7713691f1 100644 --- a/tests/unit/test_v1_router_extended.py +++ b/tests/unit/test_v1_router_extended.py @@ -2357,3 +2357,20 @@ def test_eviction_targets_least_recently_touched(self): d["d"] = 4 # overflow -> evict b assert "b" not in d assert {"a", "c", "d"} <= set(d.keys()) + + +class TestSafeLog: + """Regression guard for log-injection (CWE-117) sanitization.""" + + def test_safe_log_strips_crlf(self): + from youtube_extension.backend.api.v1 import router as router_module + + malicious = "vid123\r\nERROR forged-admin-login-success" + scrubbed = router_module._safe_log(malicious) + assert "\n" not in scrubbed and "\r" not in scrubbed + assert scrubbed == "vid123ERROR forged-admin-login-success" + + def test_safe_log_coerces_non_str(self): + from youtube_extension.backend.api.v1 import router as router_module + + assert router_module._safe_log(42) == "42" From 6efd6a0bd70e0f172b743d73d9575f6216d5c2f3 Mon Sep 17 00:00:00 2001 From: Claude Date: Sun, 2 Aug 2026 13:35:51 +0000 Subject: [PATCH 2/3] fix(security): sanitize exception text and sibling log sinks (CWE-117) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Addresses the owner's scope objection on #810: the previous pass wrapped identifiers but left the exception in the same statement raw, and attacker data reaches these logs primarily *through* exception strings (e.g. `ValueError(f"Invalid YouTube URL: {video_url}")`), so the sanitized sites stayed injectable. - Extract the sanitizer to `youtube_extension/utils/logsafe.py` (`safe_log`) and share it; router re-exports it as `_safe_log`. It now also strips ESC, vertical tab, form feed, NEL, and U+2028/U+2029 in addition to CR/LF. - Wrap every unsanitized dynamic value interpolated into a logger call in `router.py` (34 sites — all `{e}`/`{exc}` and sibling user-data sinks) and in `services/video_processing_service.py` (17 sites, incl. the `{video_url}` logs called directly by the patched handlers). An AST scan confirms zero unsanitized dynamic logger interpolations remain in either file. - Add regression tests for exception-text sanitization and the new control characters. Full unit suite: 7543 passed, 0 failed. Behavior-preserving. Co-Authored-By: Claude Opus 4.8 (1M context) Claude-Session: https://claude.ai/code/session_01FZcDgrGTkknC2ya13Uy6bU --- .../backend/api/v1/router.py | 104 ++++++++++-------- .../services/video_processing_service.py | 35 +++--- src/youtube_extension/utils/logsafe.py | 27 +++++ tests/unit/test_v1_router_extended.py | 22 ++++ 4 files changed, 124 insertions(+), 64 deletions(-) create mode 100644 src/youtube_extension/utils/logsafe.py diff --git a/src/youtube_extension/backend/api/v1/router.py b/src/youtube_extension/backend/api/v1/router.py index d8a6d09e4..a9aeef961 100644 --- a/src/youtube_extension/backend/api/v1/router.py +++ b/src/youtube_extension/backend/api/v1/router.py @@ -16,11 +16,21 @@ from datetime import datetime, timezone from typing import Any, Optional -from fastapi import APIRouter, Depends, Header, HTTPException, Query, Request, Response, status +from fastapi import ( + APIRouter, + Depends, + Header, + HTTPException, + Query, + Request, + Response, + status, +) from fastapi.responses import JSONResponse from shared.youtube import RobustYouTubeMetadata from uvai.ml.client import get_uvai_ml_client + try: from youtube_extension.services.agents import AgentOrchestrator from youtube_extension.services.agents.adapters.agent_orchestrator import ( @@ -40,6 +50,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: @@ -72,6 +83,7 @@ AgentStatus, AgentStatusResponse, ApiResponse, + BlueprintRequest, CacheStats, ChatRequest, ChatResponse, @@ -87,6 +99,7 @@ GeminiCacheResponse, GeminiTokenRequest, GeminiTokenResponse, + GenerateCodeRequest, HealthResponse, JobStatus, KnowledgeIngestRequest, @@ -96,28 +109,21 @@ TranscriptActionRequest, TranscriptActionResponse, VideoJobStatusResponse, + VideoPackRequest, VideoProcessingRequest, VideoProcessJobRequest, VideoProcessJobResponse, VideoToSoftwareRequest, VideoToSoftwareResponse, - VideoPackRequest, - BlueprintRequest, - GenerateCodeRequest, ) performance_monitor = PerformanceMonitor() logger = logging.getLogger(__name__) -def _safe_log(value: object) -> str: - """Strip CR/LF from a value before it enters a log line. - - Guards against log-injection (CWE-117): user-controlled inputs (path - params, request fields) can smuggle newlines to forge additional log - entries. Sanitize such values wherever they are interpolated into logs. - """ - return str(value).replace("\r", "").replace("\n", "") +# ``_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. @@ -132,7 +138,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]: @@ -258,7 +264,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") @@ -287,7 +293,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") @@ -352,7 +358,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)} @@ -621,7 +627,7 @@ async def chat_v1( detail = data_service.get_video_detail(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 @@ -666,7 +672,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 @@ -720,7 +726,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)}, @@ -770,7 +776,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") @@ -808,7 +814,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") @@ -832,7 +838,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") @@ -861,7 +867,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") @@ -884,7 +890,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") @@ -904,7 +910,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") @@ -942,7 +948,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") @@ -966,7 +972,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") @@ -982,7 +988,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") @@ -1017,7 +1023,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") @@ -1033,7 +1039,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 {_safe_log(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") @@ -1081,7 +1087,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 {_safe_log(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") @@ -1128,7 +1134,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") @@ -1148,7 +1154,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") @@ -1168,7 +1174,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") @@ -1186,7 +1192,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") @@ -1317,7 +1323,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 @@ -1337,7 +1343,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", _safe_log(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: @@ -1435,7 +1441,7 @@ async def _queue_transcript_action_job( logger.info( "Cloud Tasks unavailable for %s, using local background task: %s", _safe_log(job_id), - exc, + _safe_log(exc), ) asyncio.create_task( _run_video_job( @@ -1569,7 +1575,7 @@ async def _run_video_job( job.status = JobStatus.failed job.error = str(exc) _persist_video_job(job) - logger.error(f"Video job {_safe_log(job_id)} failed: {exc}") + logger.error(f"Video job {_safe_log(job_id)} failed: {_safe_log(exc)}") @router.post( @@ -1704,7 +1710,11 @@ async def get_or_create_videopack(request: VideoPackRequest): # In a real implementation, this would look up in a VideoPackStore. # For MVP, we return a synthesized pack from the job or a 404. try: - from youtube_extension.videopack.schema import Provenance, Transcript, VideoPackV0 + from youtube_extension.videopack.schema import ( + Provenance, + Transcript, + VideoPackV0, + ) # Check if we have a job with results job = None @@ -1726,7 +1736,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") @@ -1775,7 +1785,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") @@ -1801,7 +1811,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") @@ -1921,7 +1931,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: @@ -1934,7 +1944,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 @@ -1960,7 +1970,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") @@ -2099,7 +2109,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 {_safe_log(execution.agent_id)} failed: {exc}") + logger.error(f"Agent {_safe_log(execution.agent_id)} failed: {_safe_log(exc)}") @router.get( diff --git a/src/youtube_extension/backend/services/video_processing_service.py b/src/youtube_extension/backend/services/video_processing_service.py index 5039f42b5..5e9aea3ce 100644 --- a/src/youtube_extension/backend/services/video_processing_service.py +++ b/src/youtube_extension/backend/services/video_processing_service.py @@ -15,6 +15,7 @@ from pathlib import Path from typing import Any, Optional +from youtube_extension.utils.logsafe import safe_log as _safe_log from youtube_extension.utils.proxy import get_proxy_url DEPLOYMENT_TARGET_ALIASES: dict[str, str] = { @@ -87,7 +88,7 @@ def get_video_processor(self): return processor except Exception as e: - logger.error(f"Error initializing video processor: {e}") + logger.error(f"Error initializing video processor: {_safe_log(e)}") return None async def process_video_for_markdown(self, video_url: str, force_regenerate: bool = False) -> dict[str, Any]: @@ -102,13 +103,13 @@ async def process_video_for_markdown(self, video_url: str, force_regenerate: boo Dict containing processing results and metadata """ try: - logger.info(f"Processing video for markdown: {video_url}") + logger.info(f"Processing video for markdown: {_safe_log(video_url)}") # Check cache first unless force regenerate if not force_regenerate: cached_result = self.cache_service.get_cached_result(video_url) if cached_result: - logger.info(f"Returning cached result for {cached_result['video_id']}") + logger.info(f"Returning cached result for {_safe_log(cached_result['video_id'])}") # Process cached content content = cached_result['markdown_content'] @@ -145,7 +146,7 @@ async def process_video_for_markdown(self, video_url: str, force_regenerate: boo if fallback_result: return fallback_result - logger.warning(f"Video processing failed for {video_url}") + logger.warning(f"Video processing failed for {_safe_log(video_url)}") raise ValueError("Video processing failed") # Process successful result @@ -167,7 +168,7 @@ async def process_video_for_markdown(self, video_url: str, force_regenerate: boo } except Exception as e: - logger.error(f"Error in video processing service: {e}") + logger.error(f"Error in video processing service: {_safe_log(e)}") # Try fallback on error if self.use_langextract_fallback: @@ -189,7 +190,7 @@ async def process_video_basic(self, video_url: str, options: Optional[dict[str, Dict containing processing results """ try: - logger.info(f"Basic video processing: {video_url}") + logger.info(f"Basic video processing: {_safe_log(video_url)}") options = options or {} processor = self.get_video_processor() @@ -205,7 +206,7 @@ async def process_video_basic(self, video_url: str, options: Optional[dict[str, logger.info("Returning processor-level cached result") return self._normalize_result(video_url, cached) except Exception as cache_err: - logger.debug(f"Cache lookup skipped due to error: {cache_err}") + logger.debug(f"Cache lookup skipped due to error: {_safe_log(cache_err)}") # Run processing start_time = time.time() @@ -219,14 +220,14 @@ async def process_video_basic(self, video_url: str, options: Optional[dict[str, return normalized except asyncio.TimeoutError as e: - logger.error(f"Timeout during video processing: {e}") + logger.error(f"Timeout during video processing: {_safe_log(e)}") # Re-raise to allow API layer to map to 408 raise except ValueError: # Bubble up ValueError for API layer to map to 400 raise except Exception as e: - logger.error(f"Error in basic video processing: {e}") + logger.error(f"Error in basic video processing: {_safe_log(e)}") raise def _normalize_result(self, video_url: str, result: dict[str, Any]) -> dict[str, Any]: @@ -332,7 +333,7 @@ async def process_video_to_software(self, video_url: str, project_type: str = "w start_time = time.time() try: - logger.info(f"Processing video to software: {video_url}") + logger.info(f"Processing video to software: {_safe_log(video_url)}") # Import required components try: @@ -342,7 +343,7 @@ async def process_video_to_software(self, video_url: str, project_type: str = "w ) pipeline_available = True except ImportError as e: - logger.error(f"Software generation pipeline not available: {e}") + logger.error(f"Software generation pipeline not available: {_safe_log(e)}") pipeline_available = False if not pipeline_available: @@ -457,7 +458,7 @@ async def process_video_to_software(self, video_url: str, project_type: str = "w except Exception as e: processing_time = time.time() - start_time - logger.error(f"Video-to-software processing failed: {e}") + logger.error(f"Video-to-software processing failed: {_safe_log(e)}") raise e async def _try_langextract_fallback(self, video_url: str) -> Optional[dict[str, Any]]: @@ -491,7 +492,7 @@ async def _try_langextract_fallback(self, video_url: str) -> Optional[dict[str, ) if not Path(server_path).is_file(): logger.warning( - f"LangExtract MCP server not found at {server_path}; " + f"LangExtract MCP server not found at {_safe_log(server_path)}; " "set LANGEXTRACT_MCP_SERVER to override" ) return None @@ -504,14 +505,14 @@ async def _try_langextract_fallback(self, video_url: str) -> Optional[dict[str, ) if proc.returncode != 0: - logger.warning(f"LangExtract MCP call failed: {proc.stderr.decode()[:200]}") + logger.warning(f"LangExtract MCP call failed: {_safe_log(proc.stderr.decode()[:200])}") return None out = proc.stdout.decode().strip().splitlines()[-1] res = _json.loads(out).get("result", {}) if "error" in res: - logger.warning(f"LangExtract error: {res['error']}") + logger.warning(f"LangExtract error: {_safe_log(res['error'])}") return None text = res.get("text", "") @@ -530,7 +531,7 @@ async def _try_langextract_fallback(self, video_url: str) -> Optional[dict[str, } except Exception as e: - logger.warning(f"LangExtract fallback exception: {e}") + logger.warning(f"LangExtract fallback exception: {_safe_log(e)}") return None async def cleanup(self): @@ -542,4 +543,4 @@ async def cleanup(self): await close_coro() logger.info("✅ Processor session closed") except Exception as e: - logger.warning(f"Processor cleanup warning: {e}") + logger.warning(f"Processor cleanup warning: {_safe_log(e)}") diff --git a/src/youtube_extension/utils/logsafe.py b/src/youtube_extension/utils/logsafe.py new file mode 100644 index 000000000..951889fb1 --- /dev/null +++ b/src/youtube_extension/utils/logsafe.py @@ -0,0 +1,27 @@ +"""Log-injection (CWE-117) sanitization helpers.""" + +# Characters that can forge or corrupt log records: CR/LF and the other +# Unicode line boundaries, plus terminal control chars (ESC) usable for +# ANSI-escape spoofing. Mapped to ``None`` so ``str.translate`` drops them. +_UNSAFE_LOG_CHARS = { + ord("\r"): None, # carriage return + ord("\n"): None, # line feed + 0x0B: None, # vertical tab + 0x0C: None, # form feed + 0x1B: None, # ESC (ANSI escape / terminal injection) + 0x85: None, # NEL (next line) + 0x2028: None, # line separator + 0x2029: None, # paragraph separator +} + + +def safe_log(value: object) -> str: + """Return ``value`` as a string with log-forging characters removed. + + Guards against log injection (CWE-117): user-controlled inputs — path + params, request fields, URLs, and the text of exceptions raised from + them — can smuggle newlines or terminal control sequences to forge or + corrupt log records. Apply to every dynamic value (including exception + objects) interpolated into a log message. + """ + return str(value).translate(_UNSAFE_LOG_CHARS) diff --git a/tests/unit/test_v1_router_extended.py b/tests/unit/test_v1_router_extended.py index 7713691f1..ea5b6b39e 100644 --- a/tests/unit/test_v1_router_extended.py +++ b/tests/unit/test_v1_router_extended.py @@ -2374,3 +2374,25 @@ def test_safe_log_coerces_non_str(self): from youtube_extension.backend.api.v1 import router as router_module assert router_module._safe_log(42) == "42" + + def test_safe_log_strips_terminal_and_unicode_line_breaks(self): + from youtube_extension.backend.api.v1 import router as router_module + + # ESC (ANSI injection), vertical tab, form feed, NEL, and the Unicode + # line/paragraph separators must all be removed, not just CR/LF. + control = [chr(0x1B), chr(0x0B), chr(0x0C), chr(0x85), chr(0x2028), chr(0x2029)] + malicious = 'id' + ''.join(control) + 'forged' + scrubbed = router_module._safe_log(malicious) + assert scrubbed == 'idforged' + for ch in control: + assert ch not in scrubbed + + def test_safe_log_sanitizes_exception_text(self): + # The primary log-injection vector is attacker data carried *inside* + # an exception message (e.g. ValueError(f"Invalid URL: {video_url}")). + from youtube_extension.backend.api.v1 import router as router_module + + exc = ValueError("Invalid YouTube URL: bad\r\nADMIN forged log line") + scrubbed = router_module._safe_log(exc) + assert "\n" not in scrubbed and "\r" not in scrubbed + assert scrubbed == "Invalid YouTube URL: badADMIN forged log line" From eec771244ffb8a1353bcb9d1ff4af032f3210488 Mon Sep 17 00:00:00 2001 From: Claude Date: Sun, 2 Aug 2026 16:19:51 +0000 Subject: [PATCH 3/3] fix(security): sanitize final rendered log record in StructuredFormatter (CWE-117) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Per-call `_safe_log` scrubs only the interpolated message; `exc_info=True` appends the raw traceback (including `str(exc)`) and `logger.exception` renders structured `extra` fields, so attacker-controlled `\r\n` inside an exception message (e.g. `ValueError(f"Invalid YouTube URL: {video_url}")`) still forged a log line. Verified against rendered handler output. Central fix: `StructuredFormatter.format()` now runs the fully-rendered record through the shared `logsafe` translation table, neutralizing CR/LF and other line separators in the message, traceback, and any `extra` fields at once — collapsing each record to a single line — while `exc_info=True` diagnostics are preserved server-side. Adds `test_logging_config_crlf.py`, which asserts the rendered record (not just `_safe_log`'s return) cannot forge lines via tracebacks, `extra` fields, or any line-separator code point. Addresses the current-head Copilot review finding on #810. Co-Authored-By: Claude Opus 4.8 (1M context) Claude-Session: https://claude.ai/code/session_01FZcDgrGTkknC2ya13Uy6bU --- .../backend/config/logging_config.py | 14 +++- tests/unit/test_logging_config_crlf.py | 74 +++++++++++++++++++ 2 files changed, 86 insertions(+), 2 deletions(-) create mode 100644 tests/unit/test_logging_config_crlf.py diff --git a/src/youtube_extension/backend/config/logging_config.py b/src/youtube_extension/backend/config/logging_config.py index 98f28489e..6a4df9628 100644 --- a/src/youtube_extension/backend/config/logging_config.py +++ b/src/youtube_extension/backend/config/logging_config.py @@ -14,6 +14,8 @@ from datetime import datetime from pathlib import Path +from youtube_extension.utils.logsafe import _UNSAFE_LOG_CHARS + class StructuredFormatter(logging.Formatter): """ @@ -36,10 +38,18 @@ def format(self, record: logging.LogRecord) -> str: if hasattr(record, 'request_id'): record.correlation_id = record.request_id - # Format the base message + # Format the base message (includes any exc_info traceback appended + # by the base formatter). formatted_message = super().format(record) - return formatted_message + # CWE-117: neutralize CR/LF and other line separators in the FINAL + # rendered record — so an exc_info traceback or a structured ``extra`` + # field carrying attacker-controlled text (e.g. a URL inside an + # exception message) cannot forge or corrupt log lines even when the + # value was not passed through ``_safe_log`` at the call site. This + # collapses each record to a single line, matching the ``logsafe`` + # drop behavior. + return formatted_message.translate(_UNSAFE_LOG_CHARS) def formatException(self, ei) -> str: """Format exception with enhanced stack trace""" diff --git a/tests/unit/test_logging_config_crlf.py b/tests/unit/test_logging_config_crlf.py new file mode 100644 index 000000000..7f3ceb759 --- /dev/null +++ b/tests/unit/test_logging_config_crlf.py @@ -0,0 +1,74 @@ +"""CWE-117 regression: the log formatter must neutralize CR/LF in the FINAL +rendered record, including exc_info tracebacks and structured ``extra`` fields. + +This guards the vector that per-call ``_safe_log`` cannot reach: ``exc_info=True`` +appends the raw exception text (and ``logger.exception`` renders ``extra``), +which can carry attacker-controlled ``\\r\\n`` and forge a log line. +""" + +import io +import logging + +import pytest + +from youtube_extension.backend.config.logging_config import StructuredFormatter +from youtube_extension.utils.logsafe import safe_log + + +def _render(record_call) -> str: + buf = io.StringIO() + handler = logging.StreamHandler(buf) + handler.setFormatter(StructuredFormatter("%(levelname)s - %(message)s")) + logger = logging.getLogger("crlf-regression") + logger.handlers[:] = [handler] + logger.setLevel(logging.DEBUG) + logger.propagate = False + record_call(logger) + return buf.getvalue() + + +def test_exc_info_traceback_cannot_forge_log_lines(): + def call(logger): + try: + raise ValueError("boom\r\nCRITICAL - FORGED ADMIN LINE") + except ValueError as exc: + # Message is already sanitized at the call site; the traceback is + # the vector under test. + logger.error("Error in chat endpoint: %s", safe_log(exc), exc_info=True) + + # The StreamHandler appends one trailing "\n" terminator — the legitimate + # record boundary. The record *content* must contain no CR/LF, so the whole + # record is a single line and nothing can be forged. + out = _render(call) + assert out.count("\n") == 1 and out.endswith("\n") + body = out[:-1] + assert "\r" not in body and "\n" not in body + assert "FORGED ADMIN LINE" in body # present, but on the one safe line + + +def test_extra_field_cannot_forge_log_lines(): + def call(logger): + logger.info( + "processing %(video_url)s", + {"video_url": "http://x\r\nADMIN forged-from-extra"}, + ) + + out = _render(call) + body = out[:-1] if out.endswith("\n") else out + assert "\r" not in body and "\n" not in body + + +@pytest.mark.parametrize( + "codepoint", + [0x0D, 0x0A, 0x1B, 0x0B, 0x0C, 0x85, 0x2028, 0x2029], +) +def test_all_line_separators_stripped_from_rendered_record(codepoint): + sep = chr(codepoint) + + def call(logger): + logger.warning("value=%s", f"a{sep}FORGED") + + out = _render(call) + # Ignore the single trailing terminator the handler appends. + body = out[:-1] if out.endswith("\n") else out + assert sep not in body