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"])