Add cleanup and pending prepare logic in local_rag.py
This commit is contained in:
@@ -180,6 +180,22 @@ def _collect_library_paths(slug: str) -> Dict[str, Path]:
|
|||||||
}
|
}
|
||||||
|
|
||||||
|
|
||||||
|
def _cleanup_generated_artifacts(slug: str) -> None:
|
||||||
|
paths = _collect_library_paths(slug)
|
||||||
|
for key in (
|
||||||
|
"corpus",
|
||||||
|
"enhanced",
|
||||||
|
"shadow",
|
||||||
|
"shadow_index",
|
||||||
|
"shadow_store",
|
||||||
|
"content_index",
|
||||||
|
"content_store",
|
||||||
|
):
|
||||||
|
target = paths[key]
|
||||||
|
if target.exists():
|
||||||
|
target.unlink()
|
||||||
|
|
||||||
|
|
||||||
def library_payload(data: Dict[str, Any]) -> Dict[str, Any]:
|
def library_payload(data: Dict[str, Any]) -> Dict[str, Any]:
|
||||||
paths = _collect_library_paths(data["slug"])
|
paths = _collect_library_paths(data["slug"])
|
||||||
files = list(data.get("files", []))
|
files = list(data.get("files", []))
|
||||||
@@ -315,6 +331,80 @@ def _mark_pipeline_stage(slug: str, stage: str, source_signature: Optional[str])
|
|||||||
write_library(slug, data)
|
write_library(slug, data)
|
||||||
|
|
||||||
|
|
||||||
|
def _set_pending_prepare_signature(data: Dict[str, Any], source_signature: Optional[str]) -> None:
|
||||||
|
pipeline = data.get("pipeline")
|
||||||
|
if not isinstance(pipeline, dict):
|
||||||
|
pipeline = {}
|
||||||
|
data["pipeline"] = pipeline
|
||||||
|
|
||||||
|
if source_signature:
|
||||||
|
pipeline["pending_prepare_signature"] = source_signature
|
||||||
|
pipeline["pending_prepare_updated_at"] = now_iso()
|
||||||
|
else:
|
||||||
|
pipeline.pop("pending_prepare_signature", None)
|
||||||
|
pipeline.pop("pending_prepare_updated_at", None)
|
||||||
|
|
||||||
|
|
||||||
|
def _clear_pending_prepare(slug: str) -> None:
|
||||||
|
path = lib_json(slug)
|
||||||
|
if not path.exists():
|
||||||
|
return
|
||||||
|
|
||||||
|
data = _read_json(path)
|
||||||
|
_set_pending_prepare_signature(data, None)
|
||||||
|
write_library(slug, data)
|
||||||
|
|
||||||
|
|
||||||
|
async def _ensure_prepare_job(slug: str) -> Optional[str]:
|
||||||
|
path = lib_json(slug)
|
||||||
|
if not path.exists():
|
||||||
|
return None
|
||||||
|
|
||||||
|
lock = LIB_LOCKS.setdefault(slug, asyncio.Lock())
|
||||||
|
async with lock:
|
||||||
|
if _has_active_job(slug):
|
||||||
|
return None
|
||||||
|
|
||||||
|
data = read_library(slug)
|
||||||
|
payload = library_payload(data)
|
||||||
|
pipeline = _pipeline_meta(data)
|
||||||
|
pending_signature = pipeline.get("pending_prepare_signature")
|
||||||
|
|
||||||
|
if not payload["states"].get("has_files") or not pending_signature:
|
||||||
|
return None
|
||||||
|
if payload["states"].get("is_indexed") and payload.get("source_signature") == pending_signature:
|
||||||
|
_clear_pending_prepare(slug)
|
||||||
|
return None
|
||||||
|
|
||||||
|
return _start_job(slug, "prepare")
|
||||||
|
|
||||||
|
|
||||||
|
async def _handle_post_job_state(slug: str, job_type: str, status: str) -> None:
|
||||||
|
path = lib_json(slug)
|
||||||
|
if not path.exists():
|
||||||
|
return
|
||||||
|
|
||||||
|
data = read_library(slug)
|
||||||
|
payload = library_payload(data)
|
||||||
|
pipeline = _pipeline_meta(data)
|
||||||
|
pending_signature = pipeline.get("pending_prepare_signature")
|
||||||
|
|
||||||
|
if not payload["states"].get("has_files"):
|
||||||
|
_clear_pending_prepare(slug)
|
||||||
|
_cleanup_generated_artifacts(slug)
|
||||||
|
return
|
||||||
|
|
||||||
|
if pending_signature and payload["states"].get("is_indexed") and payload.get("source_signature") == pending_signature:
|
||||||
|
_clear_pending_prepare(slug)
|
||||||
|
return
|
||||||
|
|
||||||
|
if status != "succeeded":
|
||||||
|
return
|
||||||
|
|
||||||
|
if pending_signature and not payload["states"].get("is_indexed"):
|
||||||
|
await _ensure_prepare_job(slug)
|
||||||
|
|
||||||
|
|
||||||
def _scaled_progress(on_progress, start: float, end: float, prefix: str):
|
def _scaled_progress(on_progress, start: float, end: float, prefix: str):
|
||||||
def wrapped(phase: str, pct: float, detail: str):
|
def wrapped(phase: str, pct: float, detail: str):
|
||||||
clamped = min(1.0, max(0.0, float(pct)))
|
clamped = min(1.0, max(0.0, float(pct)))
|
||||||
@@ -431,6 +521,10 @@ async def _run_job(job_id: str, fn_name: str, **kwargs):
|
|||||||
job["error"] = f"{type(exc).__name__}: {exc}"
|
job["error"] = f"{type(exc).__name__}: {exc}"
|
||||||
finally:
|
finally:
|
||||||
job["finished_at"] = now_iso()
|
job["finished_at"] = now_iso()
|
||||||
|
try:
|
||||||
|
await _handle_post_job_state(job["slug"], fn_name, job["status"])
|
||||||
|
except Exception:
|
||||||
|
pass
|
||||||
|
|
||||||
|
|
||||||
def _start_job(slug: str, job_type: str, **kwargs) -> str:
|
def _start_job(slug: str, job_type: str, **kwargs) -> str:
|
||||||
@@ -537,7 +631,7 @@ def delete_library(slug: str):
|
|||||||
|
|
||||||
|
|
||||||
@router.post("/libraries/{slug}/files/register")
|
@router.post("/libraries/{slug}/files/register")
|
||||||
def register_paths(slug: str, req: RegisterPathsRequest):
|
async def register_paths(slug: str, req: RegisterPathsRequest):
|
||||||
data = read_library(slug)
|
data = read_library(slug)
|
||||||
stage = stage_dir(slug)
|
stage = stage_dir(slug)
|
||||||
existing = {entry.get("sha256"): entry for entry in data.get("files", [])}
|
existing = {entry.get("sha256"): entry for entry in data.get("files", [])}
|
||||||
@@ -566,26 +660,44 @@ def register_paths(slug: str, req: RegisterPathsRequest):
|
|||||||
added.append(entry)
|
added.append(entry)
|
||||||
existing[sha] = entry
|
existing[sha] = entry
|
||||||
|
|
||||||
|
job_id = None
|
||||||
|
if added:
|
||||||
|
_set_pending_prepare_signature(data, _source_signature(data.get("files", [])))
|
||||||
write_library(slug, data)
|
write_library(slug, data)
|
||||||
|
if added:
|
||||||
|
job_id = await _ensure_prepare_job(slug)
|
||||||
return {
|
return {
|
||||||
"added": added,
|
"added": added,
|
||||||
"skipped": skipped,
|
"skipped": skipped,
|
||||||
|
"job_id": job_id,
|
||||||
"library": library_payload(data),
|
"library": library_payload(data),
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
||||||
@router.delete("/libraries/{slug}/files")
|
@router.delete("/libraries/{slug}/files")
|
||||||
def remove_file(slug: str, req: RemoveFileRequest):
|
async def remove_file(slug: str, req: RemoveFileRequest):
|
||||||
data = read_library(slug)
|
data = read_library(slug)
|
||||||
before = len(data.get("files", []))
|
files = list(data.get("files", []))
|
||||||
data["files"] = [entry for entry in data.get("files", []) if entry.get("rel") != req.rel]
|
removed = next((entry for entry in files if entry.get("rel") == req.rel), None)
|
||||||
|
if not removed:
|
||||||
|
raise HTTPException(status_code=404, detail="File not found")
|
||||||
|
|
||||||
|
data["files"] = [entry for entry in files if entry.get("rel") != req.rel]
|
||||||
symlink_path = stage_dir(slug) / req.rel
|
symlink_path = stage_dir(slug) / req.rel
|
||||||
if symlink_path.exists():
|
if symlink_path.exists():
|
||||||
symlink_path.unlink()
|
symlink_path.unlink()
|
||||||
|
|
||||||
|
source_signature = _source_signature(data.get("files", []))
|
||||||
|
_set_pending_prepare_signature(data, source_signature)
|
||||||
write_library(slug, data)
|
write_library(slug, data)
|
||||||
if len(data["files"]) == before:
|
|
||||||
raise HTTPException(status_code=404, detail="File not found")
|
job_id = None
|
||||||
return {"ok": True, "library": library_payload(data)}
|
if source_signature:
|
||||||
|
job_id = await _ensure_prepare_job(slug)
|
||||||
|
else:
|
||||||
|
_cleanup_generated_artifacts(slug)
|
||||||
|
|
||||||
|
return {"ok": True, "job_id": job_id, "library": library_payload(data)}
|
||||||
|
|
||||||
|
|
||||||
@router.post("/libraries/{slug}/jobs/build")
|
@router.post("/libraries/{slug}/jobs/build")
|
||||||
|
|||||||
Reference in New Issue
Block a user