squid-main-space / routes /ocean_batch.py
Zin299's picture
Deploy v4.3.0 persistent task center
e5e9d87
Raw
History Blame
4.74 kB
from __future__ import annotations
import json
from fastapi import HTTPException, Request
from fastapi.responses import JSONResponse
def _inject(context):globals().update(context)
async def create_ocean_batch(request:Request):
uid,_=await resolve_request_user(request)
try:body=await request.json()
except Exception as exc:raise HTTPException(400,"invalid json") from exc
prompt=str((body or {}).get("prompt") or "").strip()
spec=parse_ocean_batch_request(prompt)
if not spec:raise HTTPException(400,"该请求不是可自动拆分的 Ocean 多日批量导出;请提供日期范围/月、经纬度范围和变量。")
return await OCEAN_BATCH_MANAGER.create(uid,prompt,spec,str((body or {}).get("thread_id") or ""))
async def get_ocean_batch(job_id:str,request:Request):
uid,_=await resolve_request_user(request)
data=OCEAN_BATCH_MANAGER.get(job_id)
if not data:raise HTTPException(404,"batch job not found")
raw=OCEAN_BATCH_MANAGER._read(job_id)
if raw.get("user_id") and uid and raw.get("user_id")!=uid:raise HTTPException(403,"forbidden")
return data
async def cancel_ocean_batch(job_id:str,request:Request):
uid,_=await resolve_request_user(request)
raw=OCEAN_BATCH_MANAGER._read(job_id)
if not raw:raise HTTPException(404,"batch job not found")
if raw.get("user_id") and uid and raw.get("user_id")!=uid:raise HTTPException(403,"forbidden")
return await OCEAN_BATCH_MANAGER.cancel(job_id)
async def resume_ocean_batch(job_id:str,request:Request):
uid,_=await resolve_request_user(request)
raw=OCEAN_BATCH_MANAGER._read(job_id)
if not raw:raise HTTPException(404,"batch job not found")
if raw.get("user_id") and uid and raw.get("user_id")!=uid:raise HTTPException(403,"forbidden")
await OCEAN_BATCH_MANAGER.start(job_id)
return OCEAN_BATCH_MANAGER.get(job_id)
async def retry_ocean_batch(job_id:str,request:Request):
uid,_=await resolve_request_user(request)
raw=OCEAN_BATCH_MANAGER._read(job_id)
if not raw:raise HTTPException(404,"batch job not found")
if raw.get("user_id") and uid and raw.get("user_id")!=uid:raise HTTPException(403,"forbidden")
return await OCEAN_BATCH_MANAGER.retry_failed(job_id)
async def retry_ocean_batch_subtask(job_id:str,subtask_id:str,request:Request):
uid,_=await resolve_request_user(request)
raw=OCEAN_BATCH_MANAGER._read(job_id)
if not raw:raise HTTPException(404,"batch job not found")
if raw.get("user_id") and uid and raw.get("user_id")!=uid:raise HTTPException(403,"forbidden")
if not any(x.get("subtask_id")==subtask_id for x in raw.get("subtasks",[])):
raise HTTPException(404,"subtask not found")
return await OCEAN_BATCH_MANAGER.retry_subtask(job_id,subtask_id)
async def refresh_ocean_batch_download(job_id:str,subtask_id:str,request:Request):
uid,_=await resolve_request_user(request)
raw=OCEAN_BATCH_MANAGER._read(job_id)
if not raw:raise HTTPException(404,"batch job not found")
if raw.get("user_id") and uid and raw.get("user_id")!=uid:raise HTTPException(403,"forbidden")
if not any(x.get("subtask_id")==subtask_id for x in raw.get("subtasks",[])):
raise HTTPException(404,"subtask not found")
try:return await OCEAN_BATCH_MANAGER.refresh_download_url(job_id,subtask_id)
except Exception as exc:raise HTTPException(502,f"刷新下载链接失败:{str(exc)[:300]}") from exc
async def ocean_batch_manifest(job_id:str,request:Request):
uid,_=await resolve_request_user(request)
raw=OCEAN_BATCH_MANAGER._read(job_id)
if not raw:raise HTTPException(404,"batch job not found")
if raw.get("user_id") and uid and raw.get("user_id")!=uid:raise HTTPException(403,"forbidden")
return JSONResponse(OCEAN_BATCH_MANAGER.manifest(job_id),headers={"Content-Disposition":f'attachment; filename="ocean_batch_{job_id}_manifest.json"'})
def register_ocean_batch_routes(app,context):
_inject(context)
app.add_api_route("/api/ocean/batch",create_ocean_batch,methods=["POST"])
app.add_api_route("/api/ocean/batch/{job_id}",get_ocean_batch,methods=["GET"])
app.add_api_route("/api/ocean/batch/{job_id}",cancel_ocean_batch,methods=["DELETE"])
app.add_api_route("/api/ocean/batch/{job_id}/resume",resume_ocean_batch,methods=["POST"])
app.add_api_route("/api/ocean/batch/{job_id}/retry",retry_ocean_batch,methods=["POST"])
app.add_api_route("/api/ocean/batch/{job_id}/subtasks/{subtask_id}/retry",retry_ocean_batch_subtask,methods=["POST"])
app.add_api_route("/api/ocean/batch/{job_id}/files/{subtask_id}/refresh",refresh_ocean_batch_download,methods=["POST"])
app.add_api_route("/api/ocean/batch/{job_id}/manifest",ocean_batch_manifest,methods=["GET"])