175 lines
5.3 KiB
Python
175 lines
5.3 KiB
Python
"""Task endpoints."""
|
||
|
||
from datetime import datetime, timezone
|
||
|
||
from fastapi import APIRouter, Depends, HTTPException, Query
|
||
from pydantic import BaseModel
|
||
from sqlalchemy import select, func
|
||
|
||
from ..deps import get_persistence
|
||
from ...storage.db import PersistenceService
|
||
from ...storage.models import ScrapeTask, SyncRun
|
||
from ...worker.tasks import sync_vehicle_task, sync_listing_task
|
||
|
||
router = APIRouter()
|
||
|
||
|
||
# --- Схемы.
|
||
|
||
class SyncVehicleRequest(BaseModel):
|
||
vehicle_url: str
|
||
lane: str = "iaai"
|
||
|
||
|
||
class SyncListingRequest(BaseModel):
|
||
make: str | None = None
|
||
model: str | None = None
|
||
lane: str = "iaai_cars"
|
||
limit: int | None = None
|
||
only_new: bool | None = None
|
||
|
||
|
||
# --- Эндпоинты.
|
||
|
||
@router.post("/tasks/sync-vehicle")
|
||
def start_sync_vehicle(
|
||
body: SyncVehicleRequest,
|
||
persistence: PersistenceService = Depends(get_persistence),
|
||
):
|
||
"""Запустить задачу скрапинга одного автомобиля через Celery."""
|
||
result = sync_vehicle_task.delay(body.vehicle_url, lane=body.lane)
|
||
|
||
# Регистрируем задачу в БД.
|
||
persistence.create_scrape_task(
|
||
celery_task_id=result.id,
|
||
task_type="sync_vehicle",
|
||
vehicle_url=body.vehicle_url,
|
||
)
|
||
|
||
return {
|
||
"task_id": result.id,
|
||
"status": "queued",
|
||
"vehicle_url": body.vehicle_url,
|
||
}
|
||
|
||
|
||
@router.post("/tasks/sync-listing")
|
||
def start_sync_listing(
|
||
body: SyncListingRequest,
|
||
persistence: PersistenceService = Depends(get_persistence),
|
||
):
|
||
"""Запустить задачу полного цикла листинга через Celery."""
|
||
result = sync_listing_task.delay(
|
||
make=body.make,
|
||
model=body.model,
|
||
lane=body.lane,
|
||
limit=body.limit,
|
||
only_new=body.only_new,
|
||
)
|
||
|
||
persistence.create_scrape_task(
|
||
celery_task_id=result.id,
|
||
task_type="sync_listing",
|
||
)
|
||
|
||
return {
|
||
"task_id": result.id,
|
||
"status": "queued",
|
||
}
|
||
|
||
|
||
@router.get("/tasks/{task_id}")
|
||
def get_task_status(
|
||
task_id: str,
|
||
persistence: PersistenceService = Depends(get_persistence),
|
||
):
|
||
"""Статус Celery-задачи."""
|
||
task_info = persistence.get_scrape_task(task_id)
|
||
if task_info is None:
|
||
raise HTTPException(status_code=404, detail="Task not found")
|
||
return task_info
|
||
|
||
|
||
@router.get("/tasks")
|
||
def list_tasks(
|
||
page: int = Query(1, ge=1),
|
||
per_page: int = Query(20, ge=1, le=100),
|
||
status: str | None = None,
|
||
task_type: str | None = None,
|
||
persistence: PersistenceService = Depends(get_persistence),
|
||
):
|
||
"""Список задач с пагинацией."""
|
||
with persistence.session_scope() as session:
|
||
query = select(ScrapeTask)
|
||
if status:
|
||
query = query.where(ScrapeTask.status == status)
|
||
if task_type:
|
||
query = query.where(ScrapeTask.task_type == task_type)
|
||
|
||
total = session.execute(
|
||
select(func.count()).select_from(query.subquery())
|
||
).scalar() or 0
|
||
|
||
offset = (page - 1) * per_page
|
||
tasks = session.execute(
|
||
query.order_by(ScrapeTask.created_at.desc()).offset(offset).limit(per_page)
|
||
).scalars().all()
|
||
|
||
return {
|
||
"total": total,
|
||
"page": page,
|
||
"per_page": per_page,
|
||
"items": [
|
||
{
|
||
"id": t.id,
|
||
"celery_task_id": t.celery_task_id,
|
||
"task_type": t.task_type,
|
||
"status": t.status,
|
||
"vehicle_url": t.vehicle_url,
|
||
"created_at": t.created_at.isoformat() if t.created_at else None,
|
||
"started_at": t.started_at.isoformat() if t.started_at else None,
|
||
"finished_at": t.finished_at.isoformat() if t.finished_at else None,
|
||
"result_summary": t.result_summary,
|
||
"error_message": t.error_message,
|
||
}
|
||
for t in tasks
|
||
],
|
||
}
|
||
|
||
|
||
@router.get("/sync-runs")
|
||
def list_sync_runs(
|
||
page: int = Query(1, ge=1),
|
||
per_page: int = Query(20, ge=1, le=100),
|
||
persistence: PersistenceService = Depends(get_persistence),
|
||
):
|
||
"""История запусков синхронизации."""
|
||
with persistence.session_scope() as session:
|
||
total = session.execute(select(func.count(SyncRun.id))).scalar() or 0
|
||
|
||
offset = (page - 1) * per_page
|
||
runs = session.execute(
|
||
select(SyncRun).order_by(SyncRun.started_at.desc()).offset(offset).limit(per_page)
|
||
).scalars().all()
|
||
|
||
return {
|
||
"total": total,
|
||
"page": page,
|
||
"per_page": per_page,
|
||
"items": [
|
||
{
|
||
"id": r.id,
|
||
"started_at": r.started_at.isoformat() if r.started_at else None,
|
||
"finished_at": r.finished_at.isoformat() if r.finished_at else None,
|
||
"status": r.status,
|
||
"lane": r.lane,
|
||
"ids_fetched": r.ids_fetched,
|
||
"cars_upserted": r.cars_upserted,
|
||
"cars_failed": r.cars_failed,
|
||
"images_upserted": r.images_upserted,
|
||
"error_summary": r.error_summary,
|
||
}
|
||
for r in runs
|
||
],
|
||
}
|