71 lines
2.2 KiB
Python
71 lines
2.2 KiB
Python
from fastapi import FastAPI, Depends, HTTPException
|
|
from contextlib import asynccontextmanager
|
|
from sqlalchemy.ext.asyncio import AsyncSession
|
|
from app.db import Base, engine, get_session
|
|
from app.schemas import CreateJobRequest, JobInfo, TaskInfo
|
|
from app.schema_utils import job_to_schema
|
|
from app.services import create_job, get_job_with_tasks
|
|
from app.models import Job, Task
|
|
|
|
from app.routers.mail import router as mail_router
|
|
|
|
from typing import List
|
|
from uuid import UUID
|
|
import re, idna
|
|
|
|
@asynccontextmanager
|
|
async def lifespan(app: FastAPI):
|
|
# STARTUP Beispiele
|
|
# await init_db_pool
|
|
# await init_broker()
|
|
|
|
# Optionale Werte Beispiele
|
|
# app.state.db_pool = get_db_pool
|
|
# app.state.broker = get_broker()
|
|
|
|
# Beispiel: DB-Update beim Startup
|
|
#await conn.run_sync(Base.metadata.create_all)
|
|
|
|
try:
|
|
yield
|
|
finally:
|
|
# SHUTDOWN Beispiele
|
|
# await close_broker()
|
|
# await close_db_pool()
|
|
pass
|
|
|
|
app = FastAPI(
|
|
title="Task Queue API",
|
|
version="0.3.0",
|
|
lifespan=lifespan)
|
|
|
|
app.include_router(mail_router)
|
|
|
|
|
|
DOMAIN_RE = re.compile(r"^(?=.{1,253}$)(?!-)[A-Za-z0-9-]{1,63}(?<!-)(\.[A-Za-z0-9-]{1,63})+$")
|
|
|
|
def _validate_domain(domain: str) -> str:
|
|
try:
|
|
_ = idna.encode(domain).decode()
|
|
except Exception:
|
|
raise HTTPException(status_code=400, detail="Invalid domain (IDNA)")
|
|
if not DOMAIN_RE.match(domain):
|
|
raise HTTPException(status_code=400, detail="Invalid domain format")
|
|
return domain
|
|
|
|
@app.post("/jobs", response_model=JobInfo)
|
|
async def create_job_endpoint(req: CreateJobRequest, session: AsyncSession = Depends(get_session)):
|
|
if not req.job_type:
|
|
raise HTTPException(status_code=400, detail="job_type is required")
|
|
domain = _validate_domain(req.domain.strip().lower())
|
|
job = await create_job(session, req.job_type, domain, req.payload)
|
|
job = await get_job_with_tasks(session, job.id)
|
|
return job_to_schema(job)
|
|
|
|
@app.get("/jobs/{job_id}", response_model=JobInfo)
|
|
async def get_job_endpoint(job_id: UUID, session: AsyncSession = Depends(get_session)):
|
|
job = await get_job_with_tasks(session, job_id)
|
|
if not job:
|
|
raise HTTPException(status_code=404, detail="Job not found")
|
|
return job_to_schema(job)
|
|
|