Files

94 lines
3.4 KiB
Python
Raw Permalink Normal View History

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)