import asyncio import secrets from datetime import datetime, timezone from fastapi import APIRouter, Depends, HTTPException, Request from app.core.security import require_permission from app.services import publish_service router = APIRouter() PUBLISH_LOCK = asyncio.Lock() PUBLISH_JOBS: dict[str, dict] = {} def utc_now_iso() -> str: return datetime.now(timezone.utc).isoformat().replace("+00:00", "Z") def public_job(job: dict): return { "job_id": job["job_id"], "status": job["status"], "message": job["message"], "result": job.get("result"), "error": job.get("error"), "created_at": job["created_at"], "updated_at": job["updated_at"], } async def run_publish_job(job_id: str, payload: dict): job = PUBLISH_JOBS[job_id] job.update(status="running", message="正在校验并发布版本", updated_at=utc_now_iso()) try: result = await publish_service.publish_payload(payload) job.update( status="success", message=result.get("msg") or "发布完成", result=result, updated_at=utc_now_iso(), ) except HTTPException as err: detail = err.detail if isinstance(err.detail, str) else str(err.detail) job.update(status="error", message=detail, error=err.detail, updated_at=utc_now_iso()) except Exception as err: job.update(status="error", message=f"发布失败:{err}", error=str(err), updated_at=utc_now_iso()) finally: if PUBLISH_LOCK.locked(): PUBLISH_LOCK.release() @router.post("/admin/publish") async def admin_publish_version(request: Request, auth=Depends(require_permission("publish:manage"))): # 发布版本会写数据库、上传大量文件到 MinIO,并可能持续较长时间。 # 同一时间只允许一个发布任务,避免两个版本并发写入导致目录、版本号或文件清单互相污染。 if PUBLISH_LOCK.locked(): raise HTTPException(status_code=409, detail="已有版本发布任务正在进行,请等待完成后再发布") async with PUBLISH_LOCK: return await publish_service.publish_version(request) @router.post("/admin/publish/job") async def admin_publish_version_job(request: Request, auth=Depends(require_permission("publish:manage"))): if PUBLISH_LOCK.locked(): raise HTTPException(status_code=409, detail="已有版本发布任务正在进行,请等待完成后再发布") await PUBLISH_LOCK.acquire() try: payload = await publish_service.make_background_publish_payload(request) job_id = secrets.token_urlsafe(16) PUBLISH_JOBS[job_id] = { "job_id": job_id, "status": "queued", "message": "发布文件已接收,等待后台处理", "result": None, "error": None, "created_at": utc_now_iso(), "updated_at": utc_now_iso(), } asyncio.create_task(run_publish_job(job_id, payload)) return public_job(PUBLISH_JOBS[job_id]) except Exception: if PUBLISH_LOCK.locked(): PUBLISH_LOCK.release() raise @router.get("/admin/publish/job/{job_id}") async def admin_publish_job_status(job_id: str, auth=Depends(require_permission("publish:manage"))): job = PUBLISH_JOBS.get(job_id) if not job: raise HTTPException(status_code=404, detail="发布任务不存在或服务端已重启") return public_job(job)