Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
25 commits
Select commit Hold shift + click to select a range
b236310
chore: reload nginx after deployment
Muizzyranking Jun 12, 2026
da82be9
feat: handlers factory
Muizzyranking Jun 12, 2026
87793c3
fix: correct import path
Muizzyranking Jun 12, 2026
2260538
chore: pagination query params
Muizzyranking Jun 12, 2026
e1cfb92
feat: description for swagger
Muizzyranking Jun 12, 2026
8d73742
feat: benchmark runner for the scheduling algorithms
Muizzyranking Jun 12, 2026
29dd1c6
feat: sceheduler that clean up jobs and runs aging loop
Muizzyranking Jun 12, 2026
b08aab7
feat: the aging worker, ensures aging a job does starve
Muizzyranking Jun 12, 2026
c9c07d1
feat: job processor
Muizzyranking Jun 12, 2026
4c6ec64
feat: main worker that manages the process
Muizzyranking Jun 12, 2026
7501027
feat: create more docker service for app
Muizzyranking Jun 12, 2026
d27fc86
feat: benchmark route
Muizzyranking Jun 12, 2026
5670e6a
feat: bin route to manage soft deleted jobs
Muizzyranking Jun 12, 2026
f6d8921
feat: dashabord route
Muizzyranking Jun 12, 2026
208283f
refactor: disable the api key header
Muizzyranking Jun 12, 2026
b9eb3ca
feat: dlq route, manages dead letter queues
Muizzyranking Jun 12, 2026
d50f9f6
feat: job route endpoints to run crud operations on jobs
Muizzyranking Jun 12, 2026
0c20a61
feat: logs endpoint
Muizzyranking Jun 12, 2026
b68a480
feat: settings endpoint to change settings on the fly
Muizzyranking Jun 12, 2026
30b4aeb
feat: sse endpoint, for live events
Muizzyranking Jun 12, 2026
d175805
feat(WIP): workers manager
Muizzyranking Jun 12, 2026
bd7e2b4
feat: api routes
Muizzyranking Jun 12, 2026
f5f2962
feat: pydantic schemas
Muizzyranking Jun 12, 2026
8aee15c
chore: add more packages
Muizzyranking Jun 12, 2026
99ca97a
chore: ruff format
Muizzyranking Jun 12, 2026
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions .github/workflows/cd.yml
Original file line number Diff line number Diff line change
Expand Up @@ -28,3 +28,4 @@ jobs:
docker-compose down
docker-compose up -d --build
docker system prune -f
sudo nginx -t && sudo systemctl reload nginx
4 changes: 2 additions & 2 deletions Dockerfile
Original file line number Diff line number Diff line change
@@ -1,4 +1,4 @@
FROM python:3.11-slim
FROM python:3.13-slim

COPY --from=ghcr.io/astral-sh/uv:latest /uv /bin/uv

Expand All @@ -12,4 +12,4 @@ COPY . .

EXPOSE 8000

CMD ["sh", "-c", "uv run alembic upgrade head && uv run gunicorn app.main:app -w 4 -k uvicorn.workers.UvicornWorker -b 0.0.0.0:8000"]
CMD ["sh", "-c", "uv run alembic upgrade head && uv run gunicorn app.main:app -w 2 -k uvicorn.workers.UvicornWorker -b 0.0.0.0:8000"]
28 changes: 28 additions & 0 deletions app/api/router.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,28 @@
from fastapi import APIRouter, Depends

from app.api.v1 import (
benchmark,
bin,
dashboard,
dlq,
jobs,
logs,
settings,
sse,
workers,
)
from app.core.security import verify_api_key

router = APIRouter()
# NOTE: this is to require api key heeader. will add later
api_router = APIRouter(dependencies=[Depends(verify_api_key)])

router.include_router(jobs.router, prefix="/jobs", tags=["Jobs"])
router.include_router(dlq.router, prefix="/dlq", tags=["DLQ"])
router.include_router(bin.router, prefix="/bin", tags=["Bin"])
router.include_router(settings.router, prefix="/settings", tags=["Settings"])
router.include_router(workers.router, prefix="/workers", tags=["Workers"])
router.include_router(logs.router, prefix="/logs", tags=["Logs"])
router.include_router(benchmark.router, prefix="/benchmark", tags=["Benchmark"])
router.include_router(dashboard.router, prefix="/dashboard", tags=["Dashboard"])
router.include_router(sse.router, prefix="/sse", tags=["SSE"])
28 changes: 28 additions & 0 deletions app/api/v1/benchmark.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,28 @@
from fastapi import APIRouter

from app.schemas.benchmark import BenchmarkRequest, BenchmarkResult
from app.schemas.response import ApiResponse

router = APIRouter()


@router.post(
"/run",
summary="Run queue algorithm benchmark",
response_model=ApiResponse[BenchmarkResult],
)
async def run_benchmark(body: BenchmarkRequest):
try:
from benchmark.runner import run_benchmark as _run

result = await _run(n=body.n, algorithm=body.algorithm)
return ApiResponse[BenchmarkResult](
message="Benchmark completed successfully.",
data=BenchmarkResult(**result),
)

except Exception as exc:
return ApiResponse[None](
message="Benchmark failed.",
errors=[{"message": str(exc)}],
), 500
65 changes: 65 additions & 0 deletions app/api/v1/bin.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,65 @@
from uuid import UUID

from fastapi import APIRouter

from app.core.exceptions import FlintException
from app.dependencies import DBSession, PaginationParams
from app.schemas.job import JobResponse
from app.schemas.response import ApiResponse, Meta, error_response
from app.services import job

router = APIRouter()


@router.get(
"",
summary="List bin (soft-deleted jobs)",
response_model=ApiResponse[list[JobResponse]],
)
async def list_bin(db: DBSession, page_params: PaginationParams):
page = page_params.page
limit = page_params.limit
jobs, total = await job.get_bin_jobs(page, limit, db)
return ApiResponse[list[JobResponse]](
message="Bin retrieved successfully.",
data=[JobResponse.model_validate(j) for j in jobs],
meta=Meta(page=page, limit=limit, total=total),
)


@router.patch(
"/{job_id}/restore",
summary="Restore a job from the bin",
response_model=ApiResponse[JobResponse],
)
async def restore_job(job_id: UUID, db: DBSession):
try:
job_result = await job.restore_job(job_id, db)
return ApiResponse[JobResponse](
message="Job restored successfully.",
data=JobResponse.model_validate(job_result),
)
except FlintException as exc:
return error_response(
message=exc.message,
errors=[{"message": exc.message}],
status_code=exc.status_code,
)


@router.delete(
"/{job_id}",
summary="Permanently delete a job",
description=("Hard-deletes a job from the database."),
response_model=ApiResponse[None],
)
async def hard_delete_job(job_id: UUID, db: DBSession):
try:
await job.hard_delete_job(job_id, db)
return ApiResponse[None](message="Job permanently deleted.")
except FlintException as exc:
return error_response(
message=exc.message,
errors=[{"message": exc.message}],
status_code=exc.status_code,
)
20 changes: 20 additions & 0 deletions app/api/v1/dashboard.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,20 @@
from fastapi import APIRouter

from app.dependencies import DBSession
from app.schemas.response import ApiResponse
from app.services.job import get_job_counts_by_status

router = APIRouter()


@router.get(
"/stats",
summary="Get dashboard job counts",
description="Returns job counts grouped by status for the dashboard.",
response_model=ApiResponse[dict],
)
async def get_dashboard_stats(db: DBSession):
counts = await get_job_counts_by_status(db)
return ApiResponse[dict](
message="Dashboard stats retrieved successfully.", data=counts
)
68 changes: 68 additions & 0 deletions app/api/v1/dlq.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,68 @@
from uuid import UUID

from fastapi import APIRouter

from app.core.exceptions import FlintException
from app.dependencies import DBSession, PaginationParams
from app.queues.heapq import HeapQueue
from app.schemas.job import JobResponse
from app.schemas.response import ApiResponse, Meta, error_response
from app.services import dlq as dlq_service

router = APIRouter()

_queue = HeapQueue()


@router.get("", summary="List DLQ jobs", response_model=ApiResponse[list[JobResponse]])
async def list_dlq(
page_params: PaginationParams,
db: DBSession,
):
page = page_params.page
limit = page_params.limit
jobs, total = await dlq_service.get_dlq_jobs(page, limit, db)
return ApiResponse[list[JobResponse]](
message="DLQ retrieved successfully.",
data=[JobResponse.model_validate(j) for j in jobs],
meta=Meta(page=page, limit=limit, total=total),
)


@router.post(
"/{job_id}/retry",
summary="Retry a DLQ job",
)
async def retry_dlq_job(job_id: UUID, db: DBSession):
try:
job = await dlq_service.retry_dlq_job(job_id, db, _queue)
return ApiResponse[JobResponse](
message="Job re-queued from DLQ successfully.",
data=JobResponse.model_validate(job),
)
except FlintException as exc:
return error_response(
message=exc.message,
errors=[{"message": exc.message}],
status_code=exc.status_code,
)


@router.delete(
"/{job_id}",
summary="Remove a job from DLQ (soft-delete)",
response_model=ApiResponse[None],
)
async def remove_from_dlq(
job_id: UUID,
db: DBSession,
):
try:
await dlq_service.remove_from_dlq(job_id, db)
return ApiResponse[None](message="Job removed from DLQ and moved to bin.")
except FlintException as exc:
return error_response(
message=exc.message,
errors=[{"message": exc.message}],
status_code=exc.status_code,
)
134 changes: 134 additions & 0 deletions app/api/v1/jobs.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,134 @@
from typing import Annotated
from uuid import UUID

from fastapi import APIRouter, Depends

from app.core.exceptions import FlintException
from app.dependencies import DBSession
from app.schemas.job import (
JobCreate,
JobFilterParams,
JobLogResponse,
JobResponse,
)
from app.schemas.response import ApiResponse, Meta, error_response
from app.services import job as job_service

router = APIRouter()


@router.post(
"",
summary="Create a job",
response_model=ApiResponse[JobResponse],
status_code=201,
)
async def create_job(
body: JobCreate,
db: DBSession,
):
try:
job = await job_service.create_job(body, db)
return ApiResponse[JobResponse](
message="Job created successfully.",
data=JobResponse.model_validate(job),
)
except FlintException as exc:
return error_response(
message=exc.message,
errors=[{"message": exc.message}],
status_code=exc.status_code,
)
except ValueError as exc:
return error_response(
message="Validation error.",
errors=[{"message": str(exc)}],
status_code=422,
)


@router.get(
"",
summary="List jobs",
response_model=ApiResponse[list[JobResponse]],
)
async def list_jobs(
db: DBSession,
filters: Annotated[JobFilterParams, Depends(JobFilterParams)],
):
jobs, total = await job_service.get_jobs(filters, db)
return ApiResponse[list[JobResponse]](
message="Jobs retrieved successfully.",
data=[JobResponse.model_validate(j) for j in jobs],
meta=Meta(
page=filters.page,
limit=filters.limit,
total=total,
),
)


@router.get(
"/{job_id}",
summary="Get job detail",
response_model=ApiResponse[JobResponse],
)
async def get_job(
job_id: UUID,
db: DBSession,
):
try:
job, dependency_ids, logs = await job_service.get_job_with_details(job_id, db)
data = JobResponse.model_validate(job)
data.dependencies = dependency_ids
data.logs = [JobLogResponse.model_validate(log) for log in logs]
return ApiResponse[JobResponse](
message="Job retrieved successfully.",
data=data,
)
except FlintException as exc:
return error_response(
message=exc.message,
errors=[{"message": exc.message}],
status_code=exc.status_code,
)


@router.patch(
"/{job_id}/cancel",
summary="Cancel a job",
response_model=ApiResponse[JobResponse],
)
async def cancel_job(
job_id: UUID,
db: DBSession,
):
try:
job = await job_service.cancel_job(job_id, db)
return ApiResponse[JobResponse](
message="Cancellation requested successfully.",
data=JobResponse.model_validate(job),
)
except FlintException as exc:
return error_response(
message=exc.message,
errors=[{"message": exc.message}],
status_code=exc.status_code,
)


@router.delete(
"/{job_id}",
summary="Soft-delete a job (move to bin)",
response_model=ApiResponse[None],
)
async def soft_delete_job(job_id: UUID, db: DBSession):
try:
await job_service.soft_delete_job(job_id, db)
return ApiResponse[None](message="Job moved to bin.")
except FlintException as exc:
return error_response(
message=exc.message,
errors=[{"message": exc.message}],
status_code=exc.status_code,
)
Loading
Loading