""" Camera API — Register, manage, and monitor RTSP cameras. """ from fastapi import APIRouter, Depends, HTTPException, status from sqlalchemy.ext.asyncio import AsyncSession from sqlalchemy import select from pydantic import BaseModel from typing import Optional from uuid import UUID from app.core.database import get_db from app.models.camera import Camera, CameraStatus from app.services.stream_manager import StreamManager from app.workers.pipeline import DetectionPipeline router = APIRouter(prefix="/cameras", tags=["cameras"]) class CameraCreate(BaseModel): name: str rtsp_url: str location: Optional[str] = None description: Optional[str] = None auto_start: bool = True class CameraResponse(BaseModel): id: UUID name: str rtsp_url: str location: Optional[str] description: Optional[str] status: str is_enabled: bool class Config: from_attributes = True @router.get("/", response_model=list[CameraResponse]) async def list_cameras(db: AsyncSession = Depends(get_db)): result = await db.execute(select(Camera).order_by(Camera.name)) return result.scalars().all() @router.post("/", response_model=CameraResponse, status_code=status.HTTP_201_CREATED) async def add_camera(payload: CameraCreate, db: AsyncSession = Depends(get_db)): """Register a new RTSP camera and optionally start streaming.""" # Check duplicate name existing = await db.execute(select(Camera).where(Camera.name == payload.name)) if existing.scalar_one_or_none(): raise HTTPException(status_code=400, detail="Camera name already exists") camera = Camera( name=payload.name, rtsp_url=payload.rtsp_url, location=payload.location, description=payload.description, status=CameraStatus.inactive, ) db.add(camera) await db.commit() await db.refresh(camera) if payload.auto_start: await _start_camera(camera) return camera @router.post("/{camera_id}/start") async def start_camera(camera_id: UUID, db: AsyncSession = Depends(get_db)): camera = await _get_camera(db, camera_id) await _start_camera(camera) camera.status = CameraStatus.active await db.commit() return {"message": f"Camera '{camera.name}' started"} @router.post("/{camera_id}/stop") async def stop_camera(camera_id: UUID, db: AsyncSession = Depends(get_db)): camera = await _get_camera(db, camera_id) manager = StreamManager.get() await manager.stop_stream(camera_id) camera.status = CameraStatus.inactive await db.commit() return {"message": f"Camera '{camera.name}' stopped"} @router.delete("/{camera_id}", status_code=status.HTTP_204_NO_CONTENT) async def delete_camera(camera_id: UUID, db: AsyncSession = Depends(get_db)): camera = await _get_camera(db, camera_id) await StreamManager.get().stop_stream(camera_id) await db.delete(camera) await db.commit() @router.get("/streams/status") async def stream_status(): """Get real-time status of all active streams.""" return StreamManager.get().get_all_status() # --- Helpers --- async def _get_camera(db: AsyncSession, camera_id: UUID) -> Camera: result = await db.execute(select(Camera).where(Camera.id == camera_id)) camera = result.scalar_one_or_none() if not camera: raise HTTPException(status_code=404, detail="Camera not found") return camera async def _start_camera(camera: Camera): manager = StreamManager.get() worker = await manager.start_stream(camera.id, camera.name, camera.rtsp_url) pipeline = DetectionPipeline(camera.id, camera.name, worker) await pipeline.start()