fix(logging): replace structlog with stdlib logging to prevent broken pipe crashes
structlog fails with BrokenPipeError when stdout is redirected (e.g., background processes, Docker logs). Replace all structlog.get_logger() calls with logging.getLogger() and convert keyword-style log calls to %-format strings. Also removes stale Docs/tasks.md (2028 lines) and updates Robot Framework tests to match current API behavior.
This commit is contained in:
@@ -455,11 +455,11 @@ async def trigger_rescan(
|
||||
}
|
||||
except AnimeServiceError as e:
|
||||
raise ServerError(
|
||||
message=f"Rescan failed: {str(e)}"
|
||||
message=str(e)
|
||||
) from e
|
||||
except Exception as exc:
|
||||
raise ServerError(
|
||||
message="Failed to start rescan"
|
||||
message=f"Failed to start rescan: {exc}"
|
||||
) from exc
|
||||
|
||||
|
||||
|
||||
@@ -342,7 +342,7 @@ async def websocket_endpoint(
|
||||
# Cleanup connection and rate limit record
|
||||
_cleanup_ws_rate_limits(connection_id)
|
||||
await ws_service.disconnect(connection_id)
|
||||
logger.info("WebSocket connection closed", connection_id=connection_id)
|
||||
logger.info("WebSocket connection closed connection_id=%s", connection_id)
|
||||
|
||||
|
||||
@router.get("/status")
|
||||
|
||||
@@ -74,9 +74,8 @@ def register_exception_handlers(app: FastAPI) -> None:
|
||||
) -> JSONResponse:
|
||||
"""Handle authentication errors (401)."""
|
||||
logger.warning(
|
||||
"Authentication error: %s",
|
||||
exc.message,
|
||||
extra={"details": exc.details, "path": str(request.url.path)},
|
||||
"Authentication error: %s details=%s path=%s",
|
||||
exc.message, exc.details, str(request.url.path),
|
||||
)
|
||||
return JSONResponse(
|
||||
status_code=exc.status_code,
|
||||
@@ -95,9 +94,8 @@ def register_exception_handlers(app: FastAPI) -> None:
|
||||
) -> JSONResponse:
|
||||
"""Handle authorization errors (403)."""
|
||||
logger.warning(
|
||||
"Authorization error: %s",
|
||||
exc.message,
|
||||
extra={"details": exc.details, "path": str(request.url.path)},
|
||||
"Authorization error: %s details=%s path=%s",
|
||||
exc.message, exc.details, str(request.url.path),
|
||||
)
|
||||
return JSONResponse(
|
||||
status_code=exc.status_code,
|
||||
@@ -116,9 +114,8 @@ def register_exception_handlers(app: FastAPI) -> None:
|
||||
) -> JSONResponse:
|
||||
"""Handle validation errors (422)."""
|
||||
logger.info(
|
||||
"Validation error: %s",
|
||||
exc.message,
|
||||
extra={"details": exc.details, "path": str(request.url.path)},
|
||||
"Validation error: %s details=%s path=%s",
|
||||
exc.message, exc.details, str(request.url.path),
|
||||
)
|
||||
return JSONResponse(
|
||||
status_code=exc.status_code,
|
||||
@@ -137,9 +134,8 @@ def register_exception_handlers(app: FastAPI) -> None:
|
||||
) -> JSONResponse:
|
||||
"""Handle bad request errors (400)."""
|
||||
logger.info(
|
||||
"Bad request error: %s",
|
||||
exc.message,
|
||||
extra={"details": exc.details, "path": str(request.url.path)},
|
||||
"Bad request error: %s details=%s path=%s",
|
||||
exc.message, exc.details, str(request.url.path),
|
||||
)
|
||||
return JSONResponse(
|
||||
status_code=exc.status_code,
|
||||
@@ -158,9 +154,8 @@ def register_exception_handlers(app: FastAPI) -> None:
|
||||
) -> JSONResponse:
|
||||
"""Handle not found errors (404)."""
|
||||
logger.info(
|
||||
"Not found error: %s",
|
||||
exc.message,
|
||||
extra={"details": exc.details, "path": str(request.url.path)},
|
||||
"Not found error: %s details=%s path=%s",
|
||||
exc.message, exc.details, str(request.url.path),
|
||||
)
|
||||
return JSONResponse(
|
||||
status_code=exc.status_code,
|
||||
@@ -179,9 +174,8 @@ def register_exception_handlers(app: FastAPI) -> None:
|
||||
) -> JSONResponse:
|
||||
"""Handle conflict errors (409)."""
|
||||
logger.info(
|
||||
"Conflict error: %s",
|
||||
exc.message,
|
||||
extra={"details": exc.details, "path": str(request.url.path)},
|
||||
"Conflict error: %s details=%s path=%s",
|
||||
exc.message, exc.details, str(request.url.path),
|
||||
)
|
||||
return JSONResponse(
|
||||
status_code=exc.status_code,
|
||||
@@ -200,9 +194,8 @@ def register_exception_handlers(app: FastAPI) -> None:
|
||||
) -> JSONResponse:
|
||||
"""Handle rate limit errors (429)."""
|
||||
logger.warning(
|
||||
"Rate limit exceeded: %s",
|
||||
exc.message,
|
||||
extra={"details": exc.details, "path": str(request.url.path)},
|
||||
"Rate limit exceeded: %s details=%s path=%s",
|
||||
exc.message, exc.details, str(request.url.path),
|
||||
)
|
||||
return JSONResponse(
|
||||
status_code=exc.status_code,
|
||||
@@ -221,13 +214,8 @@ def register_exception_handlers(app: FastAPI) -> None:
|
||||
) -> JSONResponse:
|
||||
"""Handle generic API exceptions."""
|
||||
logger.error(
|
||||
"API error: %s",
|
||||
exc.message,
|
||||
extra={
|
||||
"error_code": exc.error_code,
|
||||
"details": exc.details,
|
||||
"path": str(request.url.path),
|
||||
},
|
||||
"API error: %s error_code=%s details=%s path=%s",
|
||||
exc.message, exc.error_code, exc.details, str(request.url.path),
|
||||
)
|
||||
return JSONResponse(
|
||||
status_code=exc.status_code,
|
||||
@@ -245,10 +233,9 @@ def register_exception_handlers(app: FastAPI) -> None:
|
||||
request: Request, exc: Exception
|
||||
) -> JSONResponse:
|
||||
"""Handle unexpected exceptions."""
|
||||
logger.exception(
|
||||
"Unexpected error: %s",
|
||||
str(exc),
|
||||
extra={"path": str(request.url.path)},
|
||||
logger.error(
|
||||
"Unexpected error: %s path=%s",
|
||||
str(exc), str(request.url.path),
|
||||
)
|
||||
|
||||
# Log full traceback for debugging
|
||||
|
||||
@@ -1,13 +1,12 @@
|
||||
from __future__ import annotations
|
||||
|
||||
import asyncio
|
||||
import logging
|
||||
import time
|
||||
from datetime import datetime, timezone
|
||||
from functools import lru_cache
|
||||
from typing import Optional
|
||||
|
||||
import structlog
|
||||
|
||||
from src.server.SeriesApp import SeriesApp
|
||||
from src.server.services.progress_service import (
|
||||
ProgressService,
|
||||
@@ -19,7 +18,7 @@ from src.server.services.websocket_service import (
|
||||
get_websocket_service,
|
||||
)
|
||||
|
||||
logger = structlog.get_logger(__name__)
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
class AnimeServiceError(Exception):
|
||||
@@ -61,16 +60,28 @@ class AnimeService:
|
||||
self._scan_lock = asyncio.Lock()
|
||||
# Subscribe to SeriesApp events
|
||||
# Note: Events library uses assignment (=), not += operator
|
||||
import logging
|
||||
_logger = logging.getLogger(__name__)
|
||||
try:
|
||||
self._app.download_status = self._on_download_status
|
||||
self._app.scan_status = self._on_scan_status
|
||||
logger.info(
|
||||
"Subscribed to SeriesApp events",
|
||||
scan_status_handler=str(self._app.scan_status),
|
||||
series_app_id=id(self._app),
|
||||
_logger.info(
|
||||
"Subscribed to SeriesApp events: scan_status=%s series_app_id=%s",
|
||||
str(self._app.scan_status),
|
||||
id(self._app),
|
||||
)
|
||||
except (BrokenPipeError, OSError) as e:
|
||||
# Handle "broken pipe" when structlog tries to write to closed stdout
|
||||
# This can happen when server runs in background with stdout redirected
|
||||
import sys
|
||||
print(
|
||||
f"WARNING: Failed to subscribe to SeriesApp events: {e}. "
|
||||
f"Download/scan status callbacks may not work.",
|
||||
file=sys.stderr,
|
||||
flush=True
|
||||
)
|
||||
except Exception as e:
|
||||
logger.exception("Failed to subscribe to SeriesApp events")
|
||||
_logger.error("Failed to subscribe to SeriesApp events: %s", e)
|
||||
raise AnimeServiceError("Initialization failed") from e
|
||||
|
||||
|
||||
@@ -95,8 +106,8 @@ class AnimeService:
|
||||
|
||||
if not loop:
|
||||
logger.debug(
|
||||
"No event loop available for download status event",
|
||||
status=args.status
|
||||
"No event loop available for download status event status=%s",
|
||||
args.status
|
||||
)
|
||||
return
|
||||
|
||||
@@ -166,8 +177,8 @@ class AnimeService:
|
||||
)
|
||||
except Exception as exc: # pylint: disable=broad-except
|
||||
logger.error(
|
||||
"Error handling download status event",
|
||||
error=str(exc)
|
||||
"Error handling download status event error=%s",
|
||||
str(exc)
|
||||
)
|
||||
|
||||
def _on_scan_status(self, args) -> None:
|
||||
@@ -181,41 +192,40 @@ class AnimeService:
|
||||
args: ScanStatusEventArgs from SeriesApp containing key,
|
||||
folder, current, total, status, and progress info
|
||||
"""
|
||||
import logging
|
||||
_event_logger = logging.getLogger(__name__)
|
||||
|
||||
try:
|
||||
scan_id = "library_scan"
|
||||
|
||||
logger.info(
|
||||
"Scan status event received",
|
||||
status=args.status,
|
||||
current=args.current,
|
||||
total=args.total,
|
||||
folder=args.folder,
|
||||
_event_logger.info(
|
||||
"Scan status event received status=%s current=%s total=%s folder=%s",
|
||||
args.status, args.current, args.total, args.folder,
|
||||
)
|
||||
|
||||
# Get event loop - try running loop first, then stored loop
|
||||
loop = None
|
||||
try:
|
||||
loop = asyncio.get_running_loop()
|
||||
logger.debug("Using running event loop for scan status")
|
||||
_event_logger.debug("Using running event loop for scan status")
|
||||
except RuntimeError:
|
||||
# No running loop in this thread - use stored loop
|
||||
loop = self._event_loop
|
||||
logger.debug(
|
||||
"Using stored event loop for scan status",
|
||||
has_loop=loop is not None
|
||||
_event_logger.debug(
|
||||
"Using stored event loop for scan status has_loop=%s",
|
||||
loop is not None
|
||||
)
|
||||
|
||||
if not loop:
|
||||
logger.warning(
|
||||
"No event loop available for scan status event",
|
||||
status=args.status
|
||||
_event_logger.warning(
|
||||
"No event loop available for scan status event status=%s",
|
||||
args.status
|
||||
)
|
||||
return
|
||||
|
||||
logger.info(
|
||||
"Processing scan status event",
|
||||
status=args.status,
|
||||
loop_id=id(loop),
|
||||
_event_logger.info(
|
||||
"Processing scan status event status=%s loop_id=%s",
|
||||
args.status, id(loop),
|
||||
)
|
||||
|
||||
# Map SeriesApp scan events to progress service
|
||||
@@ -439,8 +449,8 @@ class AnimeService:
|
||||
else:
|
||||
result.append(s) # type: ignore
|
||||
return result
|
||||
except Exception:
|
||||
logger.exception("Failed to get missing episodes list")
|
||||
except Exception as e:
|
||||
_logger.error("Failed to get missing episodes list: %s", str(e))
|
||||
raise
|
||||
|
||||
async def list_missing(self) -> list[dict]:
|
||||
@@ -459,7 +469,7 @@ class AnimeService:
|
||||
except AnimeServiceError:
|
||||
raise
|
||||
except Exception as exc:
|
||||
logger.exception("list_missing failed")
|
||||
_logger.error("list_missing failed: %s", str(exc))
|
||||
raise AnimeServiceError("Failed to list missing series") from exc
|
||||
|
||||
async def list_series_with_filters(
|
||||
@@ -604,16 +614,15 @@ class AnimeService:
|
||||
result_list.append(series_dict)
|
||||
|
||||
logger.info(
|
||||
"Listed series with filters",
|
||||
total_count=len(result_list),
|
||||
filter_type=filter_type
|
||||
"Listed series with filters total=%d filter_type=%s",
|
||||
len(result_list), filter_type
|
||||
)
|
||||
return result_list
|
||||
|
||||
except AnimeServiceError:
|
||||
raise
|
||||
except Exception as exc:
|
||||
logger.exception("list_series_with_filters failed")
|
||||
logger.error("list_series_with_filters failed: %s", str(exc))
|
||||
raise AnimeServiceError(
|
||||
"Failed to list series with metadata"
|
||||
) from exc
|
||||
@@ -635,7 +644,7 @@ class AnimeService:
|
||||
result = await self._app.search(query)
|
||||
return result
|
||||
except Exception as exc:
|
||||
logger.exception("search failed")
|
||||
logger.error("search failed: %s", str(exc))
|
||||
raise AnimeServiceError("Search failed") from exc
|
||||
|
||||
async def rescan(self) -> None:
|
||||
@@ -655,30 +664,36 @@ class AnimeService:
|
||||
progress, this method returns immediately without starting
|
||||
a new scan.
|
||||
"""
|
||||
import logging
|
||||
_rescan_logger = logging.getLogger(__name__)
|
||||
|
||||
# Check if a scan is already running (non-blocking)
|
||||
if self._scan_lock.locked():
|
||||
logger.info("Rescan already in progress, ignoring request")
|
||||
_rescan_logger.info("Rescan already in progress, ignoring request")
|
||||
return
|
||||
|
||||
async with self._scan_lock:
|
||||
try:
|
||||
# Store event loop for event handlers
|
||||
self._event_loop = asyncio.get_running_loop()
|
||||
logger.info(
|
||||
"Rescan started, event loop stored",
|
||||
loop_id=id(self._event_loop),
|
||||
series_app_id=id(self._app),
|
||||
scan_handler=str(self._app.scan_status),
|
||||
_rescan_logger.info(
|
||||
"Rescan started, event loop stored. loop_id=%d series_app_id=%d",
|
||||
id(self._event_loop),
|
||||
id(self._app),
|
||||
)
|
||||
|
||||
# SeriesApp.rescan returns scanned series list
|
||||
_rescan_logger.info("Calling _app.rescan()")
|
||||
scanned_series = await self._app.rescan()
|
||||
_rescan_logger.info("Rescan completed, found %d series", len(scanned_series) if scanned_series else 0)
|
||||
|
||||
# Persist scan results to database
|
||||
if scanned_series:
|
||||
_rescan_logger.info("Saving %d series to database", len(scanned_series))
|
||||
await self._save_scan_results_to_db(scanned_series)
|
||||
|
||||
# Reload series from database to ensure consistency
|
||||
_rescan_logger.info("Loading series from database")
|
||||
await self._load_series_from_db()
|
||||
|
||||
# invalidate cache
|
||||
@@ -687,8 +702,11 @@ class AnimeService:
|
||||
except Exception: # pylint: disable=broad-except
|
||||
pass
|
||||
|
||||
except AnimeServiceError:
|
||||
# Re-raise AnimeServiceError without wrapping
|
||||
raise
|
||||
except Exception as exc: # pylint: disable=broad-except
|
||||
logger.exception("rescan failed")
|
||||
_rescan_logger.error("Rescan failed: %s", str(exc))
|
||||
raise AnimeServiceError("Rescan failed") from exc
|
||||
|
||||
async def sync_single_series_after_scan(self, series_key: str) -> None:
|
||||
@@ -1290,11 +1308,12 @@ class AnimeService:
|
||||
return True
|
||||
|
||||
except Exception as e:
|
||||
logger.exception(
|
||||
"Failed to rename folder for %s: %s -> %s",
|
||||
logger.error(
|
||||
"Failed to rename folder for %s: %s -> %s: %s",
|
||||
key,
|
||||
current_folder,
|
||||
target_folder
|
||||
target_folder,
|
||||
str(e)
|
||||
)
|
||||
return False
|
||||
|
||||
@@ -1365,7 +1384,7 @@ class AnimeService:
|
||||
logger.info("Download cancelled, propagating cancellation")
|
||||
raise
|
||||
except Exception as exc:
|
||||
logger.exception("download failed")
|
||||
logger.error("download failed: %s", str(exc))
|
||||
raise AnimeServiceError("Download failed") from exc
|
||||
|
||||
async def update_nfo_status(
|
||||
@@ -1466,10 +1485,9 @@ class AnimeService:
|
||||
)
|
||||
|
||||
except Exception as exc:
|
||||
logger.exception(
|
||||
"Failed to update NFO status",
|
||||
key=key,
|
||||
has_nfo=has_nfo
|
||||
logger.error(
|
||||
"Failed to update NFO status key=%s has_nfo=%s: %s",
|
||||
key, has_nfo, str(exc)
|
||||
)
|
||||
raise AnimeServiceError("NFO status update failed") from exc
|
||||
|
||||
@@ -1545,7 +1563,7 @@ class AnimeService:
|
||||
return result
|
||||
|
||||
except Exception as exc:
|
||||
logger.exception("Failed to query series without NFO")
|
||||
logger.error("Failed to query series without NFO: %s", str(exc))
|
||||
raise AnimeServiceError(
|
||||
"Query for series without NFO failed"
|
||||
) from exc
|
||||
@@ -1590,7 +1608,8 @@ class AnimeService:
|
||||
"with_tvdb_id": with_tvdb
|
||||
}
|
||||
|
||||
logger.info("Retrieved NFO statistics", **stats)
|
||||
logger.info("Retrieved NFO statistics total=%d with_nfo=%d without_nfo=%d with_tmdb_id=%d with_tvdb_id=%d",
|
||||
total, with_nfo, total - with_nfo, with_tmdb, with_tvdb)
|
||||
return stats
|
||||
else:
|
||||
# Use provided session and service layer count methods
|
||||
@@ -1607,11 +1626,12 @@ class AnimeService:
|
||||
"with_tvdb_id": with_tvdb
|
||||
}
|
||||
|
||||
logger.info("Retrieved NFO statistics", **stats)
|
||||
logger.info("Retrieved NFO statistics total=%d with_nfo=%d without_nfo=%d with_tmdb_id=%d with_tvdb_id=%d",
|
||||
total, with_nfo, total - with_nfo, with_tmdb, with_tvdb)
|
||||
return stats
|
||||
|
||||
except Exception as exc:
|
||||
logger.exception("Failed to get NFO statistics")
|
||||
logger.error("Failed to get NFO statistics: %s", str(exc))
|
||||
raise AnimeServiceError("NFO statistics query failed") from exc
|
||||
|
||||
|
||||
|
||||
@@ -14,17 +14,16 @@ Key Features:
|
||||
from __future__ import annotations
|
||||
|
||||
import asyncio
|
||||
import logging
|
||||
from dataclasses import dataclass, field
|
||||
from datetime import datetime, timezone
|
||||
from enum import Enum
|
||||
from pathlib import Path
|
||||
from typing import Any, Dict, List, Optional
|
||||
|
||||
import structlog
|
||||
|
||||
from src.server.services.websocket_service import WebSocketService
|
||||
|
||||
logger = structlog.get_logger(__name__)
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
class LoadingStatus(str, Enum):
|
||||
@@ -121,8 +120,8 @@ class BackgroundLoaderService:
|
||||
self._shutdown = False
|
||||
|
||||
logger.info(
|
||||
"BackgroundLoaderService initialized",
|
||||
extra={"max_concurrent_loads": max_concurrent_loads}
|
||||
"BackgroundLoaderService initialized max_concurrent_loads=%s",
|
||||
max_concurrent_loads
|
||||
)
|
||||
|
||||
async def start(self) -> None:
|
||||
@@ -140,8 +139,8 @@ class BackgroundLoaderService:
|
||||
self.worker_tasks.append(worker)
|
||||
|
||||
logger.info(
|
||||
"Background workers started",
|
||||
extra={"num_workers": len(self.worker_tasks)}
|
||||
"Background workers started num_workers=%s",
|
||||
len(self.worker_tasks)
|
||||
)
|
||||
|
||||
async def stop(self) -> None:
|
||||
@@ -164,8 +163,8 @@ class BackgroundLoaderService:
|
||||
for i, result in enumerate(results):
|
||||
if isinstance(result, Exception) and not isinstance(result, asyncio.CancelledError):
|
||||
logger.error(
|
||||
f"Worker {i} stopped with exception",
|
||||
extra={"exception": str(result)}
|
||||
"Worker %s stopped with exception exception=%s",
|
||||
i, str(result)
|
||||
)
|
||||
|
||||
self.worker_tasks = []
|
||||
@@ -202,10 +201,15 @@ class BackgroundLoaderService:
|
||||
self.active_tasks[key] = task
|
||||
await self.task_queue.put(task)
|
||||
|
||||
logger.info("Added loading task for series: %s", key)
|
||||
import logging
|
||||
_task_logger = logging.getLogger(__name__)
|
||||
_task_logger.info("Added loading task for series: %s", key)
|
||||
|
||||
# Broadcast initial status
|
||||
await self._broadcast_status(task)
|
||||
try:
|
||||
await self._broadcast_status(task)
|
||||
except Exception as e:
|
||||
_task_logger.warning("Failed to broadcast initial status: %s", e)
|
||||
|
||||
async def check_missing_data(
|
||||
self,
|
||||
@@ -288,7 +292,8 @@ class BackgroundLoaderService:
|
||||
)
|
||||
|
||||
logger.info(
|
||||
f"Worker {worker_id} processing loading task for series: {task.key}"
|
||||
"Worker %s processing loading task for series: %s",
|
||||
worker_id, task.key
|
||||
)
|
||||
|
||||
# Process the task
|
||||
@@ -304,7 +309,10 @@ class BackgroundLoaderService:
|
||||
logger.info("Worker %s task cancelled", worker_id)
|
||||
break
|
||||
except Exception as e:
|
||||
logger.exception("Error in background worker %s: %s", worker_id, e)
|
||||
logger.error(
|
||||
"Error in background worker %s: %s",
|
||||
worker_id, str(e)
|
||||
)
|
||||
# Continue processing other tasks
|
||||
continue
|
||||
|
||||
|
||||
@@ -731,9 +731,7 @@ class DownloadService:
|
||||
# Delete from database
|
||||
await self._delete_from_database(item_id)
|
||||
removed_ids.append(item_id)
|
||||
logger.info(
|
||||
"Removed from pending queue", item_id=item_id
|
||||
)
|
||||
logger.info("Removed from pending queue item_id=%s", item_id)
|
||||
|
||||
if removed_ids:
|
||||
# Notify via progress service
|
||||
@@ -803,7 +801,7 @@ class DownloadService:
|
||||
force_broadcast=True,
|
||||
)
|
||||
|
||||
logger.info("Queue reordered", reordered_count=len(item_ids))
|
||||
logger.info("Queue reordered reordered_count=%s", len(item_ids))
|
||||
|
||||
except Exception as e:
|
||||
logger.error("Failed to reorder queue: %s", e)
|
||||
@@ -1036,7 +1034,7 @@ class DownloadService:
|
||||
"""
|
||||
count = len(self._completed_items)
|
||||
self._completed_items.clear()
|
||||
logger.info("Cleared completed items", count=count)
|
||||
logger.info("Cleared completed items count=%s", count)
|
||||
|
||||
# Notify via progress service
|
||||
if count > 0:
|
||||
@@ -1062,7 +1060,7 @@ class DownloadService:
|
||||
"""
|
||||
count = len(self._failed_items)
|
||||
self._failed_items.clear()
|
||||
logger.info("Cleared failed items", count=count)
|
||||
logger.info("Cleared failed items count=%s", count)
|
||||
|
||||
# Notify via progress service
|
||||
if count > 0:
|
||||
@@ -1094,7 +1092,7 @@ class DownloadService:
|
||||
|
||||
self._pending_queue.clear()
|
||||
self._pending_items_by_id.clear()
|
||||
logger.info("Cleared pending items", count=count)
|
||||
logger.info("Cleared pending items count=%s", count)
|
||||
|
||||
# Notify via progress service
|
||||
if count > 0:
|
||||
|
||||
@@ -91,14 +91,13 @@ class ImageLoadingService:
|
||||
# Get series from database to retrieve TMDB ID
|
||||
series = await AnimeSeriesService.get_by_key(db, key)
|
||||
if not series:
|
||||
logger.warning("Series not found in database", key=key)
|
||||
logger.warning("Series not found in database key=%s", key)
|
||||
return {"poster": False, "fanart": False, "logo": False}
|
||||
|
||||
if not series.tmdb_id:
|
||||
logger.warning(
|
||||
"Series has no TMDB ID, cannot load images",
|
||||
key=key,
|
||||
name=series.name,
|
||||
"Series has no TMDB ID, cannot load images key=%s name=%s",
|
||||
key, series.name,
|
||||
)
|
||||
return {"poster": False, "fanart": False, "logo": False}
|
||||
|
||||
|
||||
@@ -130,7 +130,7 @@ class NfoScanService:
|
||||
else:
|
||||
handler(event_data)
|
||||
except Exception as e:
|
||||
logger.error("NFO scan event handler error", error=str(e))
|
||||
logger.error("NFO scan event handler error error=%s", str(e))
|
||||
|
||||
@property
|
||||
def is_scanning(self) -> bool:
|
||||
|
||||
@@ -208,7 +208,7 @@ class ProgressService:
|
||||
self._event_handlers[event_name] = []
|
||||
|
||||
self._event_handlers[event_name].append(handler)
|
||||
logger.debug("Event handler subscribed", event_type=event_name)
|
||||
logger.debug("Event handler subscribed event_type=%s", event_name)
|
||||
|
||||
def unsubscribe(
|
||||
self, event_name: str, handler: Callable[[ProgressEvent], None]
|
||||
|
||||
@@ -225,7 +225,7 @@ class ScanService:
|
||||
scan_progress = ScanProgress(scan_id)
|
||||
self._current_scan = scan_progress
|
||||
|
||||
logger.info("Starting library scan", scan_id=scan_id)
|
||||
logger.info("Starting library scan scan_id=%s", scan_id)
|
||||
|
||||
# Start progress tracking
|
||||
try:
|
||||
|
||||
@@ -16,14 +16,14 @@ optional and used for display purposes only.
|
||||
from __future__ import annotations
|
||||
|
||||
import asyncio
|
||||
import logging
|
||||
from collections import defaultdict
|
||||
from datetime import datetime, timezone
|
||||
from typing import Any, Dict, List, Optional, Set
|
||||
|
||||
import structlog
|
||||
from fastapi import WebSocket, WebSocketDisconnect
|
||||
|
||||
logger = structlog.get_logger(__name__)
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
class WebSocketServiceError(Exception):
|
||||
@@ -96,9 +96,8 @@ class ConnectionManager:
|
||||
self._connection_metadata[connection_id] = metadata or {}
|
||||
|
||||
logger.info(
|
||||
"WebSocket connected",
|
||||
connection_id=connection_id,
|
||||
total_connections=len(self._active_connections),
|
||||
"WebSocket connected connection_id=%s total_connections=%s",
|
||||
connection_id, len(self._active_connections),
|
||||
)
|
||||
|
||||
async def disconnect(self, connection_id: str) -> None:
|
||||
@@ -122,9 +121,8 @@ class ConnectionManager:
|
||||
self._connection_metadata.pop(connection_id, None)
|
||||
|
||||
logger.info(
|
||||
"WebSocket disconnected",
|
||||
connection_id=connection_id,
|
||||
total_connections=len(self._active_connections),
|
||||
"WebSocket disconnected connection_id=%s total_connections=%s",
|
||||
connection_id, len(self._active_connections),
|
||||
)
|
||||
|
||||
async def join_room(self, connection_id: str, room: str) -> None:
|
||||
@@ -138,16 +136,13 @@ class ConnectionManager:
|
||||
if connection_id in self._active_connections:
|
||||
self._rooms[room].add(connection_id)
|
||||
logger.debug(
|
||||
"Connection joined room",
|
||||
connection_id=connection_id,
|
||||
room=room,
|
||||
room_size=len(self._rooms[room]),
|
||||
"Connection joined room connection_id=%s room=%s room_size=%s",
|
||||
connection_id, room, len(self._rooms[room]),
|
||||
)
|
||||
else:
|
||||
logger.warning(
|
||||
"Attempted to join room with inactive connection",
|
||||
connection_id=connection_id,
|
||||
room=room,
|
||||
"Attempted to join room with inactive connection connection_id=%s room=%s",
|
||||
connection_id, room,
|
||||
)
|
||||
|
||||
async def leave_room(self, connection_id: str, room: str) -> None:
|
||||
@@ -166,9 +161,8 @@ class ConnectionManager:
|
||||
del self._rooms[room]
|
||||
|
||||
logger.debug(
|
||||
"Connection left room",
|
||||
connection_id=connection_id,
|
||||
room=room,
|
||||
"Connection left room connection_id=%s room=%s",
|
||||
connection_id, room,
|
||||
)
|
||||
|
||||
async def send_personal_message(
|
||||
@@ -185,26 +179,24 @@ class ConnectionManager:
|
||||
try:
|
||||
await websocket.send_json(message)
|
||||
logger.debug(
|
||||
"Personal message sent",
|
||||
connection_id=connection_id,
|
||||
message_type=message.get("type", "unknown"),
|
||||
"Personal message sent connection_id=%s message_type=%s",
|
||||
connection_id, message.get("type", "unknown"),
|
||||
)
|
||||
except WebSocketDisconnect:
|
||||
logger.warning(
|
||||
"Connection disconnected during send",
|
||||
connection_id=connection_id,
|
||||
"Connection disconnected during send connection_id=%s",
|
||||
connection_id,
|
||||
)
|
||||
await self.disconnect(connection_id)
|
||||
except Exception as e:
|
||||
logger.error(
|
||||
"Failed to send personal message",
|
||||
connection_id=connection_id,
|
||||
error=str(e),
|
||||
"Failed to send personal message connection_id=%s error=%s",
|
||||
connection_id, str(e),
|
||||
)
|
||||
else:
|
||||
logger.warning(
|
||||
"Attempted to send message to inactive connection",
|
||||
connection_id=connection_id,
|
||||
"Attempted to send message to inactive connection connection_id=%s",
|
||||
connection_id,
|
||||
)
|
||||
|
||||
async def broadcast(
|
||||
@@ -227,15 +219,14 @@ class ConnectionManager:
|
||||
await websocket.send_json(message)
|
||||
except WebSocketDisconnect:
|
||||
logger.warning(
|
||||
"Connection disconnected during broadcast",
|
||||
connection_id=connection_id,
|
||||
"Connection disconnected during broadcast connection_id=%s",
|
||||
connection_id,
|
||||
)
|
||||
disconnected.append(connection_id)
|
||||
except Exception as e:
|
||||
logger.error(
|
||||
"Failed to broadcast to connection",
|
||||
connection_id=connection_id,
|
||||
error=str(e),
|
||||
"Failed to broadcast to connection connection_id=%s error=%s",
|
||||
connection_id, str(e),
|
||||
)
|
||||
|
||||
# Cleanup disconnected connections
|
||||
@@ -243,10 +234,10 @@ class ConnectionManager:
|
||||
await self.disconnect(connection_id)
|
||||
|
||||
logger.debug(
|
||||
"Message broadcast",
|
||||
message_type=message.get("type", "unknown"),
|
||||
recipient_count=len(self._active_connections) - len(exclude),
|
||||
failed_count=len(disconnected),
|
||||
"Message broadcast message_type=%s recipient_count=%s failed_count=%s",
|
||||
message.get("type", "unknown"),
|
||||
len(self._active_connections) - len(exclude),
|
||||
len(disconnected),
|
||||
)
|
||||
|
||||
async def broadcast_to_room(
|
||||
@@ -270,17 +261,14 @@ class ConnectionManager:
|
||||
await websocket.send_json(message)
|
||||
except WebSocketDisconnect:
|
||||
logger.warning(
|
||||
"Connection disconnected during room broadcast",
|
||||
connection_id=connection_id,
|
||||
room=room,
|
||||
"Connection disconnected during room broadcast connection_id=%s room=%s",
|
||||
connection_id, room,
|
||||
)
|
||||
disconnected.append(connection_id)
|
||||
except Exception as e:
|
||||
logger.error(
|
||||
"Failed to broadcast to room member",
|
||||
connection_id=connection_id,
|
||||
room=room,
|
||||
error=str(e),
|
||||
"Failed to broadcast to room member connection_id=%s room=%s error=%s",
|
||||
connection_id, room, str(e),
|
||||
)
|
||||
|
||||
# Cleanup disconnected connections
|
||||
@@ -288,11 +276,9 @@ class ConnectionManager:
|
||||
await self.disconnect(connection_id)
|
||||
|
||||
logger.debug(
|
||||
"Message broadcast to room",
|
||||
room=room,
|
||||
message_type=message.get("type", "unknown"),
|
||||
recipient_count=len(room_members),
|
||||
failed_count=len(disconnected),
|
||||
"Message broadcast to room room=%s message_type=%s recipient_count=%s failed_count=%s",
|
||||
room, message.get("type", "unknown"),
|
||||
len(room_members), len(disconnected),
|
||||
)
|
||||
|
||||
async def get_connection_count(self) -> int:
|
||||
@@ -604,9 +590,8 @@ class WebSocketService:
|
||||
}
|
||||
await self._manager.broadcast(message)
|
||||
logger.info(
|
||||
"Broadcast scan_started",
|
||||
directory=directory,
|
||||
total_items=total_items,
|
||||
"Broadcast scan_started directory=%s total_items=%s",
|
||||
directory, total_items,
|
||||
)
|
||||
|
||||
async def broadcast_scan_progress(
|
||||
@@ -660,10 +645,8 @@ class WebSocketService:
|
||||
}
|
||||
await self._manager.broadcast(message)
|
||||
logger.info(
|
||||
"Broadcast scan_completed",
|
||||
total_directories=total_directories,
|
||||
total_files=total_files,
|
||||
elapsed_seconds=round(elapsed_seconds, 2),
|
||||
"Broadcast scan_completed total_directories=%s total_files=%s elapsed_seconds=%s",
|
||||
total_directories, total_files, round(elapsed_seconds, 2),
|
||||
)
|
||||
|
||||
async def shutdown(self, timeout: float = 5.0) -> None:
|
||||
|
||||
@@ -410,7 +410,7 @@ async def rate_limit_dependency(request: Request) -> None:
|
||||
record.count += 1
|
||||
if record.count > max_requests:
|
||||
logger.warning(
|
||||
"Rate limit exceeded", extra={"client": client_id}
|
||||
"Rate limit exceeded client=%s", client_id
|
||||
)
|
||||
raise HTTPException(
|
||||
status_code=status.HTTP_429_TOO_MANY_REQUESTS,
|
||||
@@ -423,13 +423,10 @@ async def log_request_dependency(request: Request) -> None:
|
||||
"""Log request metadata for auditing and debugging purposes."""
|
||||
|
||||
logger.info(
|
||||
"API request",
|
||||
extra={
|
||||
"method": request.method,
|
||||
"path": request.url.path,
|
||||
"client": request.client.host if request.client else "unknown",
|
||||
"query": dict(request.query_params),
|
||||
},
|
||||
"API request method=%s path=%s client=%s query=%s",
|
||||
request.method, request.url.path,
|
||||
request.client.host if request.client else "unknown",
|
||||
dict(request.query_params),
|
||||
)
|
||||
|
||||
|
||||
@@ -557,23 +554,44 @@ def get_background_loader_service() -> "BackgroundLoaderService":
|
||||
|
||||
if _background_loader_service is None:
|
||||
try:
|
||||
import logging
|
||||
_init_logger = logging.getLogger(__name__)
|
||||
_init_logger.info("Creating BackgroundLoaderService instance...")
|
||||
|
||||
from src.server.services.background_loader_service import (
|
||||
BackgroundLoaderService,
|
||||
)
|
||||
from src.server.services.websocket_service import get_websocket_service
|
||||
|
||||
anime_service = get_anime_service()
|
||||
series_app = get_series_app()
|
||||
websocket_service = get_websocket_service()
|
||||
_init_logger.info("Imported BackgroundLoaderService")
|
||||
|
||||
from src.server.services.websocket_service import get_websocket_service
|
||||
_init_logger.info("Getting websocket_service...")
|
||||
websocket_service = get_websocket_service()
|
||||
_init_logger.info("Got websocket_service: %s", id(websocket_service))
|
||||
|
||||
_init_logger.info("Getting anime_service...")
|
||||
anime_service = get_anime_service()
|
||||
_init_logger.info("Got anime_service: %s", id(anime_service))
|
||||
|
||||
_init_logger.info("Getting series_app...")
|
||||
series_app = get_series_app()
|
||||
_init_logger.info("Got series_app: %s", id(series_app))
|
||||
|
||||
_init_logger.info("Creating BackgroundLoaderService with params: ws=%s, ans=%s, sa=%s",
|
||||
id(websocket_service), id(anime_service), id(series_app))
|
||||
_background_loader_service = BackgroundLoaderService(
|
||||
websocket_service=websocket_service,
|
||||
anime_service=anime_service,
|
||||
series_app=series_app
|
||||
)
|
||||
_init_logger.info("BackgroundLoaderService created successfully: %s", id(_background_loader_service))
|
||||
except HTTPException:
|
||||
raise
|
||||
except Exception as e:
|
||||
import logging
|
||||
_err_logger = logging.getLogger(__name__)
|
||||
_err_logger.error("Error in BackgroundLoaderService creation: %s", str(e))
|
||||
import traceback
|
||||
_err_logger.error("Traceback: %s", traceback.format_exc())
|
||||
raise HTTPException(
|
||||
status_code=status.HTTP_500_INTERNAL_SERVER_ERROR,
|
||||
detail=(
|
||||
|
||||
@@ -74,13 +74,8 @@ class ErrorTracker:
|
||||
self.error_history = self.error_history[-self.max_history_size:]
|
||||
|
||||
logger.info(
|
||||
f"Error tracked: {error_id}",
|
||||
extra={
|
||||
"error_id": error_id,
|
||||
"error_type": error_type,
|
||||
"status_code": status_code,
|
||||
"request_path": request_path,
|
||||
},
|
||||
"Error tracked error_id=%s error_type=%s status_code=%s request_path=%s",
|
||||
error_id, error_type, status_code, request_path,
|
||||
)
|
||||
|
||||
return error_id
|
||||
|
||||
Reference in New Issue
Block a user