Padmanav's picture
fix: import AsyncResult directly so test patches on app.api.routes.AsyncResult resolve
4c8f7d3
Raw
History Blame Contribute Delete
5.37 kB
from typing import Any
from fastapi import APIRouter, Depends, HTTPException
from starlette.requests import Request
from app.api.dependencies import verify_api_key
from app.api.limiter import limiter
from app.core.logger import get_logger
from app.core.metrics import metrics
from app.models.report import EngineeringReport
from app.models.repository import RepositoryRequest, RepositoryResponse
from app.models.job import JobStatus, JobStatusResponse, JobSubmitResponse
from app.worker.tasks import analyze_repository_task
from app.services.analysis_service import run_full_analysis
from app.agents import repo_analysis_agent
from app.models.review import ReviewFeedback
from app.tools.github_tool import clone_repository, delete_repository
from app.api.limiter import LIMIT_SYNC_ANALYZE, LIMIT_ASYNC_SUBMIT, LIMIT_JOB_POLL
from celery.result import AsyncResult
router = APIRouter()
logger = get_logger(__name__)
@router.get("/health")
def health_check() -> dict[str, str]:
return {"status": "ok", "message": "AI Code Review Agent is running"}
# REPLACE the body of analyze_repository with:
@router.post(
"/analyze", response_model=EngineeringReport, dependencies=[Depends(verify_api_key)]
)
@limiter.limit(LIMIT_SYNC_ANALYZE)
async def analyze_repository(
request: Request, body: RepositoryRequest
) -> EngineeringReport:
"""
Full pipeline with Redis cache:
1. Check cache by HEAD SHA
2. If miss: clone, run agents, store result
3. Return engineering report
"""
try:
return await run_full_analysis(body.github_url, body.base_sha)
except HTTPException:
raise
except Exception as exc:
logger.exception(
"Pipeline failed", extra={"url": body.github_url, "error": str(exc)}
)
raise HTTPException(
status_code=500, detail="Internal server error. Check server logs."
)
@router.post(
"/analyze/async",
response_model=JobSubmitResponse,
status_code=202,
dependencies=[Depends(verify_api_key)],
)
@limiter.limit(LIMIT_ASYNC_SUBMIT)
async def analyze_repository_async(
request: Request, body: RepositoryRequest
) -> JobSubmitResponse:
"""
Enqueue a full analysis pipeline run on a Celery worker and return
immediately with a job ID. Poll GET /jobs/{job_id} for status/result.
"""
task = analyze_repository_task.delay(body.github_url, body.base_sha)
return JobSubmitResponse(job_id=task.id, status=JobStatus.PENDING)
@router.get(
"/jobs/{job_id}",
response_model=JobStatusResponse,
dependencies=[Depends(verify_api_key)],
)
@limiter.limit(LIMIT_JOB_POLL)
def get_job_status(request: Request, job_id: str) -> JobStatusResponse:
"""Poll the status/result of an async analysis job."""
result = AsyncResult(job_id)
status_map = {
"PENDING": JobStatus.PENDING,
"STARTED": JobStatus.STARTED,
"SUCCESS": JobStatus.SUCCESS,
"FAILURE": JobStatus.FAILURE,
"RETRY": JobStatus.STARTED,
}
status = status_map.get(result.status, JobStatus.PENDING)
if status == JobStatus.SUCCESS:
return JobStatusResponse(job_id=job_id, status=status, result=result.result)
if status == JobStatus.FAILURE:
return JobStatusResponse(job_id=job_id, status=status, error=str(result.result))
return JobStatusResponse(job_id=job_id, status=status)
@router.post(
"/quick-analyze",
response_model=RepositoryResponse,
dependencies=[Depends(verify_api_key)],
)
@limiter.limit("10/minute")
async def quick_analyze(
request: Request, body: RepositoryRequest
) -> RepositoryResponse:
"""
Only run repository analysis agent.
Faster, for testing purposes.
"""
local_path = None
try:
local_path = clone_repository(body.github_url)
metadata = await repo_analysis_agent.run(body.github_url, local_path)
return RepositoryResponse(
success=True, message="Repository analyzed successfully", data=metadata
)
except HTTPException:
raise
except Exception as exc:
logger.exception(
"Quick analyze failed", extra={"url": body.github_url, "error": str(exc)}
)
raise HTTPException(
status_code=500, detail="Internal server error. Check server logs."
)
finally:
if local_path:
delete_repository(local_path)
@router.get("/metrics", dependencies=[Depends(verify_api_key)])
def get_metrics() -> dict[str, Any]:
"""Observability endpoint: agent run counts, error rates, latencies."""
return metrics.snapshot()
_feedback_store: list[dict] = []
def _clear_feedback_store() -> None: # test helper only
_feedback_store.clear()
@router.post("/feedback", dependencies=[Depends(verify_api_key)])
async def submit_feedback(body: ReviewFeedback) -> dict[str, str]:
_feedback_store.append(body.model_dump())
logger.info(
"Feedback received", extra={"useful": body.useful, "finding": body.finding[:50]}
)
return {"status": "recorded"}
@router.get("/feedback/summary", dependencies=[Depends(verify_api_key)])
def feedback_summary() -> dict[str, int]:
useful = sum(1 for f in _feedback_store if f["useful"])
not_useful = len(_feedback_store) - useful
return {"useful": useful, "not_useful": not_useful, "total": len(_feedback_store)}