diff --git a/src/youtube_extension/backend/api/v1/router.py b/src/youtube_extension/backend/api/v1/router.py index 7bc80726b..0fa7e5637 100644 --- a/src/youtube_extension/backend/api/v1/router.py +++ b/src/youtube_extension/backend/api/v1/router.py @@ -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: @@ -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. + + async def _emit_event(event_type: str, data: dict, subject: str | None = None) -> None: """Emit a CloudEvent if the publisher is available.""" @@ -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]: @@ -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") @@ -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") @@ -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)} @@ -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 = { @@ -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)}") 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 = ( @@ -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 @@ -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 @@ -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}, @@ -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)}, @@ -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 @@ -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") @@ -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( @@ -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") @@ -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") @@ -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") @@ -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") @@ -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") @@ -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") @@ -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") @@ -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") @@ -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") @@ -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") @@ -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") @@ -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") @@ -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") @@ -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") @@ -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") @@ -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 @@ -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: @@ -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( @@ -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 @@ -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( @@ -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") @@ -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") @@ -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") @@ -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: @@ -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 @@ -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") @@ -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( 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/src/youtube_extension/backend/services/video_processing_service.py b/src/youtube_extension/backend/services/video_processing_service.py index 8ea026c60..a0919895d 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]]: @@ -490,7 +491,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 @@ -517,14 +518,14 @@ async def _try_langextract_fallback(self, video_url: str) -> Optional[dict[str, return None if proc.returncode != 0: - logger.warning(f"LangExtract MCP call failed: {stderr.decode()[:200]}") + logger.warning(f"LangExtract MCP call failed: {_safe_log(stderr.decode()[:200])}") return None out = 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", "") @@ -543,7 +544,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): @@ -555,4 +556,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_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 diff --git a/tests/unit/test_v1_router_extended.py b/tests/unit/test_v1_router_extended.py index bb206e505..497b7b19a 100644 --- a/tests/unit/test_v1_router_extended.py +++ b/tests/unit/test_v1_router_extended.py @@ -2373,3 +2373,42 @@ 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" + + 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"