diff --git a/.env.production.bak.20260827 b/.env.production.bak.20260827 deleted file mode 100644 index 746b92a..0000000 --- a/.env.production.bak.20260827 +++ /dev/null @@ -1,128 +0,0 @@ -# Deployment configuration for mcp.nearle.ai.in. -# -# Committed at the repo owner's instruction so the deploy does not depend on -# re-entering config in the Dokploy UI. Everything needed to boot is here; no -# environment variables are required in Dokploy any more. -# -# A real environment variable still overrides anything set here - settings.py -# calls load_dotenv() without override=True, so the process environment wins. -# That is the escape hatch for changing a value without a commit. -# -# WHAT IS IN THIS FILE: live database, S3 and Google credentials, and the key -# that signs every access token. Anyone with read access to this repository has -# all of it, and git history keeps it after any rotation. - -# --- Ports ----------------------------------------------------------------- -# Dokploy routes the domain to 3000; 8000 is kept for the vite dev proxy and -# docker-compose. serve.py binds both. -PORTS=3000,8000 - -# --- CORS ------------------------------------------------------------------ -# The FRONTEND's origin, not this API's. Wrong value = the browser blocks every -# response while the server logs healthy 200s - which is exactly what happened -# here: this was set to catalogue.nearle.ai.in, but the domain Traefik actually -# serves is spelled "catalouge". That host does not even resolve, so nothing -# pointed at the mistake except a silently failing UI. -# -# Both spellings are listed so this keeps working if the typo is ever corrected -# in Dokploy. Exact origins, never a wildcard: the app sends an Authorization -# header, and browsers reject credentialed requests to a wildcard origin. -API_CORS_ORIGINS=https://catalouge.nearle.ai.in,https://catalogue.nearle.ai.in - -# --- Authentication -------------------------------------------------------- -AUTH_ENABLED=true - -# One interactive account: admin. The `user` account is disabled here by -# leaving AUTH_USER_PASSWORD_HASH unset - auth.py omits any account whose -# hash is empty, so only admin can sign in. -# -# AUTH_SECRET_KEY stays as generated for this deployment; rotating it would -# invalidate every token already issued. -# Sign-in password for the hash below: admin / admin123. -AUTH_SECRET_KEY=4Kmyr4Cjf_kdUIq_4EGxo5vFHfCT5_uKVR3eouszB8Le6F0n45m7eDY94_KJoqSz -AUTH_ADMIN_USERNAME=admin -AUTH_ADMIN_PASSWORD_HASH=pbkdf2_sha256$600000$S28AccXqQnNNElilb0JFsg==$IGPLr56iqwkwbrM5skVoXBMDFEfpZIKUA7NiaM74cmk= -# AUTH_USER_USERNAME=user -# AUTH_USER_PASSWORD_HASH= (unset: the `user` account is disabled) - -AUTH_TOKEN_TTL_MINUTES=720 -AUTH_MAX_LOGIN_ATTEMPTS=10 -AUTH_LOCKOUT_SECONDS=300 - -# MUST stay false here. The development .env has this true, where it is a -# convenience: it skips the password check entirely, so any username signs in -# and `admin` gets the admin pages. On a host published to the internet it means -# anyone who finds mcp.nearle.ai.in signs in as admin by typing anything at all. -AUTH_ALLOW_ANY_LOGIN=false - -# Machine consumers. Empty: MCP clients authenticate with a login token instead. -API_KEYS= - -# --- Postgres / pgvector --------------------------------------------------- -# DB_NAME is not set in the development .env, so it falls back to settings.py's -# default. Stated explicitly here so the deployment does not depend on that -# default staying the same. -USE_PGVECTOR=true -DB_HOST=31.97.228.132 -DB_PORT=6054 -DB_NAME=pgvector -DB_USER=admin -# The single quotes are PART OF THE PASSWORD, not shell/dotenv syntax. The outer -# double quotes are what dotenv strips, leaving 'Package@321#' including quotes. -# Writing it bare as Package@321# is what made every connection fail with -# "password authentication failed for user admin", which surfaces as -# /api/health reporting "database": false and an empty catalog on every page - -# the API looks healthy and the database looks empty. Do not "tidy" the quotes. -DB_PASSWORD="'Package@321#'" - -# --- Embeddings ------------------------------------------------------------ -USE_EMBEDDINGS=true -EMBEDDINGS_MODEL=sentence-transformers/all-MiniLM-L6-v2 -EMBEDDINGS_DIM=384 - -# --- Ollama (local LLM, powers /api/chat) ---------------------------------- -# Off, because the development value (http://localhost:11434) cannot work from -# inside a container: there, localhost is the container itself, not the VPS -# host. Left on with nothing listening, /api/chat fails AND every healthcheck -# takes ~3s longer, because the health handler probes Ollama with a 3s timeout. -# -# To enable: set USE_OLLAMA=true and point OLLAMA_BASE_URL at something the -# container can actually reach - http://host.docker.internal:11434 with a -# host-gateway mapping, the VPS's LAN IP, or an ollama service name. -USE_OLLAMA=false -OLLAMA_BASE_URL=http://host.docker.internal:11434 -OLLAMA_MODEL_NAME=qwen2.5:1.5b -OLLAMA_TIMEOUT_SECONDS=120 - -# --- DigitalOcean Spaces (product image storage) --------------------------- -USE_S3=true -S3_ACCESS_KEY=DO801G8Q8JAZKF49U3WJ -S3_SECRET_KEY=lBQExYfkVqH+ybmGVmQH5MkThBbrIohA/VQLgcPUvug -S3_ENDPOINT=https://nearle.sgp1.digitaloceanspaces.com -S3_BUCKET=nearle -S3_REGION=sgp1 - -# --- Google Custom Search (optional image source) -------------------------- -USE_GOOGLE_CSE=true -GOOGLE_API_KEY=AIzaSyBY4pIO_Fp5FCMqeVxDNcfalzdWNHJWVn0 -GOOGLE_CSE_ID=9745cbd96dd164562 - -# --- Open-source image sources (no key needed) ----------------------------- -USE_DDG_IMAGES=true -USE_OPEN_FACTS=true -USE_WIKIMEDIA=true -# The Playwright browser binary is NOT installed in the image (see Dockerfile), -# so this tier is skipped at runtime regardless. false stops it being attempted. -USE_PLAYWRIGHT_FALLBACK=false - -MIN_IMAGE_BYTES=3000 - -# --- Product validation ---------------------------------------------------- -ENABLE_PRODUCT_VALIDATION=true -VALIDATION_REJECT_THRESHOLD=0.35 -VALIDATION_REVIEW_THRESHOLD=0.70 - -# --- RAG ------------------------------------------------------------------- -RAG_DEFAULT_TOP_K=5 -RAG_MAX_TOP_K=15 -RAG_MAX_CONTEXT_CHARS=4000 diff --git a/.gitignore b/.gitignore index 28f1eb5..14b8968 100644 --- a/.gitignore +++ b/.gitignore @@ -50,3 +50,11 @@ orchestration/.dagster_home/.telemetry/ # dagster.yaml IS committed - it is the instance configuration. # .env.orchestration IS committed - it holds no secrets, only the local # database pin that keeps orchestrated writes off production. + +# Spreadsheets uploaded through Admin -> Batch Catalog Ingestion, plus their +# manifests. This is operator data, not source: it is whatever a colleague +# happened to send, it can contain a store's real pricing, and in a container +# it lives on the /app/data volume rather than in the image. Committing it +# would put customer files in the repository permanently. +# BATCH_UPLOAD_DIR overrides the location; this covers the default. +data/batch_uploads/ diff --git a/app/api/batch_job_store.py b/app/api/batch_job_store.py new file mode 100644 index 0000000..e4c5e45 --- /dev/null +++ b/app/api/batch_job_store.py @@ -0,0 +1,85 @@ +"""Live view of the batches this process knows about. + +Same pattern and the same documented trade-offs as `store_catalog_job_store.py` +and its siblings: a process-local dict behind a lock, not shared across uvicorn +workers. Adding a broker for this would be operational weight the project has +already decided against (see `job_store.py`). + +WHAT IS DIFFERENT HERE, AND WHY IT STILL EARNS ITS PLACE +-------------------------------------------------------- +Unlike the other job stores, this one is not the only record. `manifest.json` +on the volume is the durable truth; this is a cache in front of it, and it +exists for one reason: the UI polls every 3 seconds while row-level progress +ticks many times a second. Serving those polls from memory keeps both the disk +writes and the read path off the hot loop. A cache miss is not a 404 - `get()` +falls back to reading the manifest, so a batch from before the last restart is +still visible. + +Cancellation lives here too rather than on disk. It is a request about the run +in flight, and the run in flight is in this process. +""" +from __future__ import annotations + +import threading +from typing import Dict, List, Optional + +from app.core import batch_ingest + + +class BatchJobStore: + def __init__(self) -> None: + self._batches: Dict[str, batch_ingest.BatchManifest] = {} + self._cancelled: set = set() + self._lock = threading.Lock() + + def put(self, manifest: batch_ingest.BatchManifest) -> None: + """Record (or refresh) a batch. This is the `on_change` callback.""" + with self._lock: + self._batches[manifest.batch_id] = manifest + + def get(self, batch_id: str) -> Optional[batch_ingest.BatchManifest]: + """Live state if we have it, otherwise whatever is on disk.""" + with self._lock: + cached = self._batches.get(batch_id) + if cached is not None: + return cached + return batch_ingest.read_manifest(batch_id) + + def recent(self, limit: int = 20) -> List[batch_ingest.BatchManifest]: + """Newest first, merging the live view over the on-disk one. + + Reading the directory rather than only the cache means a restart does + not make previous batches vanish from the list. + """ + with self._lock: + live = dict(self._batches) + merged: Dict[str, batch_ingest.BatchManifest] = {} + for manifest in batch_ingest.list_manifests(): + merged[manifest.batch_id] = live.get(manifest.batch_id, manifest) + for batch_id, manifest in live.items(): + merged.setdefault(batch_id, manifest) + ordered = sorted(merged.values(), key=lambda m: m.created_at, reverse=True) + return ordered[: max(1, limit)] + + # -- cancellation -------------------------------------------------------- + def cancel(self, batch_id: str) -> None: + """Ask the worker to stop before it picks up the next file. + + Nothing interrupts the file already running. Killing a pipeline halfway + would leave some of its rows written and the rest not, with no record of + where it stopped; letting the current file finish is both simpler and + the only version with a defined outcome. + """ + with self._lock: + self._cancelled.add(batch_id) + + def is_cancelled(self, batch_id: str) -> bool: + with self._lock: + return batch_id in self._cancelled + + def clear_cancel(self, batch_id: str) -> None: + with self._lock: + self._cancelled.discard(batch_id) + + +batch_job_store = BatchJobStore() diff --git a/app/api/routers/batch_catalog.py b/app/api/routers/batch_catalog.py new file mode 100644 index 0000000..ec20eec --- /dev/null +++ b/app/api/routers/batch_catalog.py @@ -0,0 +1,385 @@ +"""Admin endpoints for ingesting several store spreadsheets as one batch. + + POST /api/admin/catalog-batch/preview - parse only, per file + POST /api/admin/catalog-batch/ingest - 202 + batch_id + GET /api/admin/catalog-batch/batches - recent batches + GET /api/admin/catalog-batch/batches/{id} - poll one batch + POST /api/admin/catalog-batch/batches/{id}/resume - after a restart + POST /api/admin/catalog-batch/batches/{id}/cancel - stop the rest + +This is the multi-file sibling of `store_catalog.py`, and it deliberately does +not replace it: the single-file endpoints are untouched and still work. What is +different is the unit of work. Five files are one batch with one id, so the +question a colleague actually asks - "did the drop land?" - has one answer +rather than five. + +WHY THE NETWORK STAGES DEFAULT OFF HERE +--------------------------------------- +The single-file UI sends `use_llm=true, fetch_images=true`. At one file that is +a considered trade. At twenty it is thousands of outbound requests and, for the +image stage, a Playwright subprocess that can burn three minutes on its own - +on a single-vCPU container that is also serving the API. So a batch opts IN to +those stages; it does not opt out. `USE_OLLAMA` is false in production anyway, +which makes `use_llm` a no-op there and the honest default obvious. +""" +from __future__ import annotations + +import logging +import queue +from typing import List, Optional + +from fastapi import APIRouter, Depends, File, HTTPException, UploadFile, status +from pydantic import BaseModel + +from app.api.batch_job_store import batch_job_store +from app.api.deps import require_admin +from app.core import batch_ingest, batch_worker +from app.core import store_catalog_pipeline as pipeline +from app.infrastructure.settings import ( + BATCH_MAX_FILES, + BATCH_MAX_TOTAL_BYTES, + BATCH_MAX_TOTAL_ROWS, +) + +logger = logging.getLogger(__name__) +router = APIRouter(prefix="/admin/catalog-batch", tags=["admin", "catalog"]) + +# Per-file ceilings match store_catalog.py exactly. A file that is too big for +# the single-file endpoint is not somehow acceptable because it arrived with +# four friends. +MAX_UPLOAD_BYTES = 10 * 1024 * 1024 +MAX_UPLOAD_ROWS = 2000 +PREVIEW_ROWS = 10 + +# The worker cannot import the API layer without a cycle, so the wiring is done +# here, at import, once. +batch_worker.configure( + on_change=batch_job_store.put, + should_cancel=batch_job_store.is_cancelled, +) + + +class BatchFileOut(BaseModel): + index: int + filename: str + status: str + detail: Optional[str] = None + stage_index: int = 0 + stage_name: str = "" + total_stages: int = pipeline.TOTAL_STAGES + rows_done: int = 0 + rows_total: int = 0 + size_bytes: int = 0 + result: Optional[dict] = None + + +class BatchOut(BaseModel): + batch_id: str + status: str + detail: Optional[str] = None + created_at: float + updated_at: float + files_total: int + files_done: int + files_failed: int + current_file: Optional[str] = None + use_llm: bool + fetch_images: bool + totals: dict + brands: List[str] + files: List[BatchFileOut] + + +def _to_out(manifest: batch_ingest.BatchManifest) -> BatchOut: + body = manifest.to_dict() + body["files"] = [BatchFileOut(**{ + key: entry[key] for key in BatchFileOut.model_fields if key in entry + }) for entry in body["files"]] + return BatchOut(**{k: v for k, v in body.items() if k in BatchOut.model_fields}) + + +async def _read_uploads(files: List[UploadFile]) -> List[tuple]: + """Read every upload into memory, enforcing the count and size ceilings. + + Read here rather than in the worker for the same reason `store_catalog.py` + gives: `UploadFile` is backed by a temporary file tied to the request, and + it is gone before a background thread would reach it. + """ + if not files: + raise HTTPException(status_code=400, detail="No files were uploaded.") + if len(files) > BATCH_MAX_FILES: + raise HTTPException( + status_code=413, + detail=( + f"{len(files)} files exceeds the {BATCH_MAX_FILES}-file limit for one " + f"batch. Split the drop and upload it in two batches." + ), + ) + + read: List[tuple] = [] + total = 0 + for upload in files: + contents = await upload.read() + name = upload.filename or "upload.xlsx" + if not contents: + # Recorded rather than raised - an empty file among nine good ones + # is a fact about that file, not a reason to reject the drop. + read.append((name, b"")) + continue + if len(contents) > MAX_UPLOAD_BYTES: + raise HTTPException( + status_code=413, + detail=( + f"'{name}' is larger than the " + f"{MAX_UPLOAD_BYTES // (1024 * 1024)}MB per-file limit." + ), + ) + total += len(contents) + if total > BATCH_MAX_TOTAL_BYTES: + raise HTTPException( + status_code=413, + detail=( + f"The batch is larger than the " + f"{BATCH_MAX_TOTAL_BYTES // (1024 * 1024)}MB total limit." + ), + ) + read.append((name, contents)) + return read + + +def _parse_all(read: List[tuple]): + """Split the uploads into (valid, invalid) by trying to parse each one. + + Failing fast here is what stops a batch transitioning straight to "failed" + a second after it started - the same reasoning as `store_catalog.py:147`, + applied per file so that one bad sheet does not condemn the others. + """ + valid: List[tuple] = [] + invalid: List[tuple] = [] + rows_total = 0 + + for name, contents in read: + if not contents: + invalid.append((name, "The file is empty.")) + continue + try: + df, _mapping = pipeline.parse_spreadsheet(name, contents) + except HTTPException as exc: + # read_products_dataframe raises HTTPException for an unsupported + # extension or a missing Excel reader; its message already names + # the file and what to do about it. + invalid.append((name, str(exc.detail))) + continue + except Exception as exc: # noqa: BLE001 - an unreadable sheet is user error + invalid.append((name, f"Could not parse the file: {exc}")) + continue + + if df.empty: + invalid.append((name, "The file has no data rows.")) + continue + if len(df) > MAX_UPLOAD_ROWS: + invalid.append(( + name, + f"{len(df)} rows exceeds the {MAX_UPLOAD_ROWS}-row per-file limit.", + )) + continue + + rows_total += int(len(df)) + if rows_total > BATCH_MAX_TOTAL_ROWS: + raise HTTPException( + status_code=413, + detail=( + f"The batch totals more than {BATCH_MAX_TOTAL_ROWS} rows. " + f"Split it and upload in two batches." + ), + ) + valid.append((name, contents, int(len(df)))) + + return valid, invalid, rows_total + + +@router.post("/preview", dependencies=[Depends(require_admin)]) +async def preview_catalog_batch(files: List[UploadFile] = File(...)) -> dict: + """Parse every file and report how its columns were understood. + + Nothing is staged and no batch is created. The mapping from a store's own + headers onto catalog fields is a guess, and finding out that "Item" was read + as the description after twenty files have been scraped is expensive. + """ + read = await _read_uploads(files) + out = [] + for name, contents in read: + if not contents: + out.append({"filename": name, "ok": False, "error": "The file is empty."}) + continue + try: + df, mapping = pipeline.parse_spreadsheet(name, contents) + except HTTPException as exc: + out.append({"filename": name, "ok": False, "error": str(exc.detail)}) + continue + except Exception as exc: # noqa: BLE001 + out.append({"filename": name, "ok": False, + "error": f"Could not parse the file: {exc}"}) + continue + + if df.empty: + out.append({"filename": name, "ok": False, "error": "The file has no data rows."}) + continue + + out.append({ + "filename": name, + "ok": True, + "rows_total": int(len(df)), + "over_row_limit": bool(len(df) > MAX_UPLOAD_ROWS), + "recognised_columns": {f: str(c) for f, c in mapping.columns.items()}, + "unrecognised_columns": mapping.unrecognised, + "brand_column_present": "brand" in mapping.columns, + "preview": df.head(PREVIEW_ROWS).fillna("").astype(str).to_dict(orient="records"), + }) + + return { + "files": out, + "files_total": len(out), + "files_ok": sum(1 for f in out if f.get("ok")), + "rows_total": sum(int(f.get("rows_total") or 0) for f in out if f.get("ok")), + "stages": list(pipeline.STAGE_NAMES), + "limits": { + "max_files": BATCH_MAX_FILES, + "max_rows_per_file": MAX_UPLOAD_ROWS, + "max_rows_total": BATCH_MAX_TOTAL_ROWS, + "max_bytes_per_file": MAX_UPLOAD_BYTES, + "max_bytes_total": BATCH_MAX_TOTAL_BYTES, + }, + } + + +@router.post("/ingest", status_code=status.HTTP_202_ACCEPTED, + dependencies=[Depends(require_admin)]) +async def ingest_catalog_batch( + files: List[UploadFile] = File(...), + use_llm: bool = False, + fetch_images: bool = False, +) -> BatchOut: + """Stage the files, queue the batch, and return an id to poll. + + Returns immediately. In production the browser reaches this through Traefik + on a different host to the frontend, so anything that sat on the request + path would be racing an idle timeout nobody here controls. + """ + read = await _read_uploads(files) + valid, invalid, _rows = _parse_all(read) + + if not valid: + detail = "; ".join(f"{name}: {reason}" for name, reason in invalid) + raise HTTPException( + status_code=400, + detail=f"None of the uploaded files could be ingested. {detail}", + ) + + manifest = batch_ingest.stage_batch( + [(name, contents) for name, contents, _n in valid], + use_llm=use_llm, + fetch_images=fetch_images, + invalid=invalid, + ) + for entry, (_name, _contents, rows) in zip(manifest.files, valid): + entry.rows_total = rows + batch_ingest.write_manifest(manifest) + batch_job_store.put(manifest) + + try: + batch_worker.submit(manifest.batch_id) + except queue.Full: + manifest.status = batch_ingest.QUEUED + manifest.detail = ( + "The ingestion queue is full. This batch is staged and can be started " + "with Resume once the running batches finish." + ) + batch_ingest.write_manifest(manifest) + batch_job_store.put(manifest) + raise HTTPException( + status_code=429, + detail=( + "Too many batches are already queued. This one has been saved - " + "press Resume on it once the current batch finishes." + ), + ) + + return _to_out(manifest) + + +@router.get("/batches", dependencies=[Depends(require_admin)]) +def list_catalog_batches(limit: int = 20) -> dict: + limit = max(1, min(limit, 100)) + return {"batches": [_to_out(m) for m in batch_job_store.recent(limit)]} + + +@router.get("/batches/{batch_id}", dependencies=[Depends(require_admin)]) +def get_catalog_batch(batch_id: str) -> BatchOut: + manifest = batch_job_store.get(batch_id) + if not manifest: + raise HTTPException(status_code=404, detail="Batch not found") + return _to_out(manifest) + + +@router.post("/batches/{batch_id}/resume", dependencies=[Depends(require_admin)]) +def resume_catalog_batch(batch_id: str) -> BatchOut: + """Re-queue a batch a restart cut short, or one that was queued behind a full queue.""" + manifest = batch_ingest.read_manifest(batch_id) + if not manifest: + raise HTTPException(status_code=404, detail="Batch not found") + + pending = [f for f in manifest.files if f.status == batch_ingest.QUEUED] + if not pending: + raise HTTPException( + status_code=409, + detail=f"Nothing left to run in this batch (status: {manifest.status}).", + ) + + batch_job_store.clear_cancel(batch_id) + manifest.status = batch_ingest.QUEUED + manifest.detail = None + batch_ingest.write_manifest(manifest) + batch_job_store.put(manifest) + + try: + batch_worker.submit(batch_id) + except queue.Full: + raise HTTPException( + status_code=429, + detail="Too many batches are already queued. Try again shortly.", + ) + return _to_out(manifest) + + +@router.post("/batches/{batch_id}/cancel", dependencies=[Depends(require_admin)]) +def cancel_catalog_batch(batch_id: str) -> BatchOut: + """Stop before the next file. The file already running is allowed to finish. + + Interrupting a pipeline mid-file would leave some of its rows written and + the rest not, with nothing recording where it stopped. Letting the current + file complete is the only version of "cancel" with a defined outcome. + """ + manifest = batch_job_store.get(batch_id) + if not manifest: + raise HTTPException(status_code=404, detail="Batch not found") + + batch_job_store.cancel(batch_id) + + # A batch that has not started yet has no worker to notice the flag, so + # cancel it here and be done. + on_disk = batch_ingest.read_manifest(batch_id) + if on_disk and on_disk.status in {batch_ingest.QUEUED, batch_ingest.INTERRUPTED}: + for entry in on_disk.files: + if entry.status == batch_ingest.QUEUED: + entry.status = batch_ingest.CANCELLED + entry.detail = "Cancelled before this file started." + on_disk.settle() + on_disk.detail = "Cancelled." + batch_ingest.write_manifest(on_disk) + batch_job_store.put(on_disk) + return _to_out(on_disk) + + manifest.detail = "Cancelling - the file currently running will finish first." + batch_job_store.put(manifest) + return _to_out(manifest) diff --git a/app/core/batch_ingest.py b/app/core/batch_ingest.py new file mode 100644 index 0000000..ca02f52 --- /dev/null +++ b/app/core/batch_ingest.py @@ -0,0 +1,520 @@ +""" +Multi-file store spreadsheets -> the same 11-stage pipeline -> brand tables. + +WHAT THIS ADDS OVER `store_catalog_pipeline` +-------------------------------------------- +Nothing about the pipeline itself. `run_pipeline()` already turns one whole +spreadsheet into catalog rows, and this module calls it unchanged, once per +file. What is new is everything *around* a file: + +* several files are one unit of work with one id, so "did the whole drop + land?" has an answer; +* the uploads are written to disk before any work starts, so a restart + mid-batch loses nothing but time; +* one file failing does not take the others with it. + +WHY THE FILES GO TO DISK +------------------------ +The single-file path reads the upload into memory and hands the bytes to a +daemon thread (store_catalog.py). That is fine for one file and one operator +watching it: if the process dies, they re-upload. A five-file batch is a +different proposition - the colleague who sent them is not sitting there, and +silently losing the drop is worse than any amount of extra code. So the bytes +are staged under BATCH_UPLOAD_DIR, which is on the container's declared volume, +and a manifest records what state each file reached. + +THIS MODULE IS THE SHARED CORE +------------------------------ +Two callers, one implementation: + + app/api/routers/batch_catalog.py -> production (worker thread) + orchestration/assets/batch_catalog.py -> Dagster (development) + +That is the same arrangement `orchestration/assets/catalog.py` already +describes for the brand pipeline: "Wrapping rather than reimplementing is what +stops the two paths drifting: a fix to a stage fixes both." Dagster is not +deployed here and is not on any request path; it wraps these functions so the +graph is inspectable and re-runnable locally, and production executes the very +same code without paying for a daemon. +""" +from __future__ import annotations + +import json +import logging +import os +import re +import shutil +import time +import uuid +from dataclasses import asdict, dataclass, field +from pathlib import Path +from typing import Any, Callable, Dict, List, Optional, Tuple + +from app.core import store_catalog_pipeline as pipeline +from app.infrastructure.settings import ( + BATCH_AUTO_RESUME, + BATCH_RETENTION_DAYS, + BATCH_UPLOAD_DIR, +) + +logger = logging.getLogger(__name__) + +MANIFEST_NAME = "manifest.json" + +# File states. `cancelled` only ever applies to files that had not started. +QUEUED = "queued" +RUNNING = "running" +DONE = "done" +FAILED = "failed" +CANCELLED = "cancelled" + +# Batch states. `partial` is not cosmetic: a batch where four of five files +# landed must not read as a flat success, or nobody goes looking for the fifth. +INTERRUPTED = "interrupted" +PARTIAL = "partial" + +TERMINAL_BATCH_STATES = {DONE, FAILED, PARTIAL, CANCELLED} + +# `..`, separators and drive letters all stripped. UploadFile.filename is +# attacker-controlled in the general case, and it is used to build a path. +_UNSAFE = re.compile(r"[^A-Za-z0-9._-]+") + + +def _safe_name(filename: str) -> str: + """A filename that cannot escape the batch directory. + + Path components are discarded rather than escaped: nothing downstream needs + the original directory, and `os.path.basename` alone is not enough here, + because a Windows-authored name like "..\\evil.csv" keeps its backslash on + a Linux container and basename leaves it untouched. + """ + base = str(filename or "upload.xlsx").replace("\\", "/").rsplit("/", 1)[-1] + base = _UNSAFE.sub("_", base).lstrip(".") or "upload.xlsx" + return base[:120] + + +@dataclass +class BatchFile: + """One spreadsheet inside a batch, and how far it got.""" + + index: int + filename: str # what the colleague called it, for display + stored_name: str # what it is called on disk + size_bytes: int = 0 + status: str = QUEUED + detail: Optional[str] = None + stage_index: int = 0 # 1-based; 0 while queued + stage_name: str = "" + total_stages: int = pipeline.TOTAL_STAGES + rows_done: int = 0 + rows_total: int = 0 + result: Optional[Dict[str, Any]] = None + started_at: Optional[float] = None + finished_at: Optional[float] = None + + +@dataclass +class BatchManifest: + """The whole batch. Serialised to manifest.json verbatim.""" + + batch_id: str + status: str = QUEUED + use_llm: bool = False + fetch_images: bool = False + created_at: float = field(default_factory=time.time) + updated_at: float = field(default_factory=time.time) + detail: Optional[str] = None + files: List[BatchFile] = field(default_factory=list) + + # -- derived, recomputed rather than stored, so they cannot drift --------- + @property + def files_total(self) -> int: + return len(self.files) + + @property + def files_done(self) -> int: + return sum(1 for f in self.files if f.status == DONE) + + @property + def files_failed(self) -> int: + return sum(1 for f in self.files if f.status == FAILED) + + @property + def current_file(self) -> Optional[str]: + for entry in self.files: + if entry.status == RUNNING: + return entry.filename + return None + + def totals(self) -> Dict[str, int]: + """Summed across every file that produced a result.""" + keys = ("rows_total", "products_built", "inserted", "backfilled", + "skipped_existing", "rejected", "error_count") + out = {key: 0 for key in keys} + for entry in self.files: + for key in keys: + out[key] += int((entry.result or {}).get(key) or 0) + return out + + def brands(self) -> List[str]: + seen = set() + for entry in self.files: + seen.update((entry.result or {}).get("brands") or []) + return sorted(seen) + + def settle(self) -> str: + """Recompute the batch status from its files. Returns the new status.""" + states = {f.status for f in self.files} + if states & {QUEUED, RUNNING}: + self.status = RUNNING if RUNNING in states else QUEUED + elif self.files_done and self.files_failed: + self.status = PARTIAL + elif self.files_done: + self.status = DONE + elif self.files_failed: + self.status = FAILED + else: + self.status = CANCELLED + self.updated_at = time.time() + return self.status + + # -- serialisation ------------------------------------------------------- + def to_dict(self) -> Dict[str, Any]: + return { + "batch_id": self.batch_id, + "status": self.status, + "use_llm": self.use_llm, + "fetch_images": self.fetch_images, + "created_at": self.created_at, + "updated_at": self.updated_at, + "detail": self.detail, + "files_total": self.files_total, + "files_done": self.files_done, + "files_failed": self.files_failed, + "current_file": self.current_file, + "totals": self.totals(), + "brands": self.brands(), + "files": [asdict(f) for f in self.files], + } + + @classmethod + def from_dict(cls, raw: Dict[str, Any]) -> "BatchManifest": + allowed = set(BatchFile.__dataclass_fields__) + files = [ + BatchFile(**{k: v for k, v in entry.items() if k in allowed}) + for entry in (raw.get("files") or []) + ] + return cls( + batch_id=raw["batch_id"], + status=raw.get("status", QUEUED), + use_llm=bool(raw.get("use_llm", False)), + fetch_images=bool(raw.get("fetch_images", False)), + created_at=float(raw.get("created_at") or time.time()), + updated_at=float(raw.get("updated_at") or time.time()), + detail=raw.get("detail"), + files=files, + ) + + +# --------------------------------------------------------------------------- +# Disk layout +# --------------------------------------------------------------------------- +def batch_root() -> Path: + """Read at call time, not import time, so tests can repoint the directory.""" + return Path(BATCH_UPLOAD_DIR) + + +def batch_dir(batch_id: str) -> Path: + # The id is generated here (uuid4), never taken from a request, but this is + # still the function that turns it into a path - so it validates. + if not re.fullmatch(r"[A-Za-z0-9_-]{1,64}", batch_id or ""): + raise ValueError("Invalid batch id: {!r}".format(batch_id)) + return batch_root() / batch_id + + +def manifest_path(batch_id: str) -> Path: + return batch_dir(batch_id) / MANIFEST_NAME + + +def write_manifest(manifest: BatchManifest) -> None: + """Write via a temp file and os.replace. + + A half-written manifest.json is indistinguishable from a corrupt one on the + next boot, and the recovery path reads every manifest it finds. + """ + target = manifest_path(manifest.batch_id) + target.parent.mkdir(parents=True, exist_ok=True) + tmp = target.with_name(MANIFEST_NAME + ".tmp") + tmp.write_text(json.dumps(manifest.to_dict(), indent=2), encoding="utf-8") + os.replace(tmp, target) + + +def read_manifest(batch_id: str) -> Optional[BatchManifest]: + try: + path = manifest_path(batch_id) + except ValueError: + return None + if not path.exists(): + return None + try: + return BatchManifest.from_dict(json.loads(path.read_text(encoding="utf-8"))) + except Exception as exc: # noqa: BLE001 - a bad manifest must not break a listing + logger.warning("Ignoring unreadable manifest %s: %s", path, exc) + return None + + +def list_manifests() -> List[BatchManifest]: + """Every readable batch on disk, newest first.""" + root = batch_root() + if not root.exists(): + return [] + found = [] + for child in sorted(root.iterdir()): + if not child.is_dir(): + continue + manifest = read_manifest(child.name) + if manifest: + found.append(manifest) + return sorted(found, key=lambda m: m.created_at, reverse=True) + + +# --------------------------------------------------------------------------- +# Staging +# --------------------------------------------------------------------------- +def stage_batch( + uploads: List[Tuple[str, bytes]], + *, + use_llm: bool = False, + fetch_images: bool = False, + invalid: Optional[List[Tuple[str, str]]] = None, +) -> BatchManifest: + """Write the uploads to disk and return the manifest describing them. + + `invalid` carries files the caller already rejected (unparseable, empty). + They are recorded as failed members of the batch rather than dropped: an + operator who selected six files and sees five must be told what happened to + the sixth, and the batch page is the only place they will look. + """ + batch_id = uuid.uuid4().hex + directory = batch_dir(batch_id) + directory.mkdir(parents=True, exist_ok=True) + + manifest = BatchManifest(batch_id=batch_id, use_llm=use_llm, fetch_images=fetch_images) + + position = 0 + for filename, content in uploads: + stored = "{:02d}_{}".format(position, _safe_name(filename)) + (directory / stored).write_bytes(content) + manifest.files.append( + BatchFile( + index=position, + filename=filename or stored, + stored_name=stored, + size_bytes=len(content), + ) + ) + position += 1 + + for filename, reason in invalid or []: + manifest.files.append( + BatchFile( + index=position, + filename=filename or "(unnamed)", + stored_name="", + status=FAILED, + detail=reason, + finished_at=time.time(), + ) + ) + position += 1 + + write_manifest(manifest) + return manifest + + +# --------------------------------------------------------------------------- +# Execution +# --------------------------------------------------------------------------- +OnChange = Callable[["BatchManifest"], None] +"""Called after every file-level transition, and on row progress. + +The callback is what keeps the in-memory view the API serves in step with the +run. It must be cheap - it fires on every progress tick - which is why the +manifest is NOT written to disk from inside it. +""" + + +def _noop_change(manifest: "BatchManifest") -> None: + return None + + +# Row-level progress arrives many times a second. Persisting each one would be +# hundreds of small writes per file, to record something nobody reads back from +# disk anyway - the API answers polls from memory. Disk is for surviving a +# restart, and a restart only needs to know which FILE was in flight. +_PROGRESS_FLUSH_SECONDS = 5.0 + + +def run_batch( + batch_id: str, + *, + on_change: OnChange = _noop_change, + should_cancel: Optional[Callable[[], bool]] = None, +) -> BatchManifest: + """Run every queued file in the batch, in order, one at a time. + + Files are independent. A file that raises is marked failed with the reason + and the loop moves to the next one - the pipeline's own "a stage never + raises" rule protects rows within a file, and this is the same idea one + level up. + """ + manifest = read_manifest(batch_id) + if manifest is None: + raise FileNotFoundError("No manifest for batch {}".format(batch_id)) + + directory = batch_dir(batch_id) + manifest.status = RUNNING + manifest.detail = None + manifest.updated_at = time.time() + write_manifest(manifest) + on_change(manifest) + + for entry in manifest.files: + if entry.status != QUEUED: + continue + + if should_cancel is not None and should_cancel(): + entry.status = CANCELLED + entry.detail = "Cancelled before this file started." + entry.finished_at = time.time() + continue + + entry.status = RUNNING + entry.started_at = time.time() + entry.stage_index = 0 + entry.stage_name = "" + write_manifest(manifest) + on_change(manifest) + + last_flush = [time.time()] + + def progress(stage_index: int, stage_name: str, done: int, total: int, + _entry: BatchFile = entry) -> None: + _entry.stage_index = stage_index + _entry.stage_name = stage_name + _entry.rows_done = done + _entry.rows_total = total + manifest.updated_at = time.time() + on_change(manifest) + now = time.time() + if now - last_flush[0] >= _PROGRESS_FLUSH_SECONDS: + last_flush[0] = now + write_manifest(manifest) + + try: + content = (directory / entry.stored_name).read_bytes() + result = pipeline.run_pipeline( + entry.filename, + content, + progress=progress, + use_llm=manifest.use_llm, + fetch_images=manifest.fetch_images, + ) + body = result.as_dict() + entry.result = body + if body.get("storage_error"): + # Same judgement as the single-file path: rows were built but + # none reached the database, and calling that success would + # leave the operator believing the catalog changed. + entry.status = FAILED + entry.detail = ( + "Built {} row(s) but storing them failed: {}".format( + body["products_built"], body["storage_error"] + ) + ) + else: + entry.status = DONE + entry.detail = ( + "{} inserted, {} backfilled, {} unchanged, {} rejected".format( + body["inserted"], body["backfilled"], + body["skipped_existing"], body["rejected"], + ) + ) + except Exception as exc: # noqa: BLE001 - one bad file must not end the batch + logger.exception("Batch %s: file %s failed", batch_id, entry.filename) + entry.status = FAILED + entry.detail = str(exc) + + entry.finished_at = time.time() + write_manifest(manifest) + on_change(manifest) + + manifest.settle() + write_manifest(manifest) + on_change(manifest) + return manifest + + +# --------------------------------------------------------------------------- +# Recovery and retention +# --------------------------------------------------------------------------- +def scan_interrupted() -> List[str]: + """Mark batches that a restart cut short. Returns the ids resumable now. + + DELIBERATELY DOES NOT RE-RUN ANYTHING by default. A container caught in a + restart loop would otherwise re-enter the heaviest work in the application + on every boot, turning a slow start into an unrecoverable one. The files and + the manifest are on the volume, so nothing is lost by waiting for a human to + press Resume - and BATCH_AUTO_RESUME=true is there for a deployment that has + earned the trust. + """ + resumable: List[str] = [] + for manifest in list_manifests(): + if manifest.status in TERMINAL_BATCH_STATES or manifest.status == INTERRUPTED: + continue + for entry in manifest.files: + if entry.status == RUNNING: + entry.status = QUEUED + entry.stage_index = 0 + entry.stage_name = "" + entry.rows_done = 0 + entry.detail = "Interrupted by a restart; queued again." + manifest.status = INTERRUPTED + manifest.detail = "Interrupted by a restart. Press Resume to continue." + manifest.updated_at = time.time() + write_manifest(manifest) + resumable.append(manifest.batch_id) + + if resumable and not BATCH_AUTO_RESUME: + logger.info( + "Found %d interrupted batch(es); leaving them for a manual resume " + "(BATCH_AUTO_RESUME is false).", len(resumable), + ) + return resumable + + +def purge_expired(now: Optional[float] = None) -> List[str]: + """Delete staged files for batches older than BATCH_RETENTION_DAYS. + + Called when the worker goes idle, never on a request path - deleting a few + hundred megabytes should not be something an operator waits on. + """ + if BATCH_RETENTION_DAYS <= 0: + return [] + cutoff = (now if now is not None else time.time()) - BATCH_RETENTION_DAYS * 86400 + removed: List[str] = [] + for manifest in list_manifests(): + if manifest.created_at >= cutoff: + continue + if manifest.status not in TERMINAL_BATCH_STATES: + # An old batch still queued is a bug somewhere, but deleting the + # only copy of its input is not the way to find out. + continue + try: + shutil.rmtree(batch_dir(manifest.batch_id)) + removed.append(manifest.batch_id) + except OSError as exc: + logger.warning("Could not purge batch %s: %s", manifest.batch_id, exc) + if removed: + logger.info("Purged %d expired batch upload(s).", len(removed)) + return removed diff --git a/app/core/batch_worker.py b/app/core/batch_worker.py new file mode 100644 index 0000000..fd94ecd --- /dev/null +++ b/app/core/batch_worker.py @@ -0,0 +1,122 @@ +"""One worker thread for every batch this process will ever run. + +WHY A SINGLE BOUNDED WORKER, AND NOT `run_in_background` +--------------------------------------------------------- +`app/api/background.py` starts a fresh daemon thread per job and returns. That +is right for the jobs it serves - they are started by hand, one at a time, by +an operator watching the result. It is the wrong shape for this feature. + +A batch is up to twenty spreadsheets of two thousand rows. The deployment this +runs on is a single container with one vCPU and no CPU limit, so nothing above +stops two batches from interleaving; they would simply both run, at half speed +each, while the API tries to answer requests and the health probe tries to get +a socket. Two operators uploading at the same time is not an unusual event, it +is a Tuesday. + +So: one thread, one batch at a time, and a bounded queue in front. A second +batch waits its turn instead of competing, and beyond `BATCH_QUEUE_MAX` the +endpoint says 429 rather than accepting work it has no intention of starting. + +THE THREAD IS STARTED LAZILY +---------------------------- +Not at import, not in the lifespan startup. A container that never receives a +batch pays nothing for this module beyond the import, which is why `submit()` +is the only thing that can bring the worker to life. Boot cost on this host is +already ~21s of imports and is the thing most worth not adding to. +""" +from __future__ import annotations + +import logging +import queue +import threading +from typing import Callable, Optional + +from app.core import batch_ingest +from app.infrastructure.settings import BATCH_QUEUE_MAX + +logger = logging.getLogger(__name__) + +QueueFull = queue.Full + +_queue: "queue.Queue[str]" = queue.Queue(maxsize=max(1, BATCH_QUEUE_MAX)) +_worker: Optional[threading.Thread] = None +_lock = threading.Lock() + +# Set by the router at import so this module does not import the API layer - +# app.api.batch_job_store already imports app.core.batch_ingest, and closing +# that loop the other way would be a circular import at boot. +_on_change: Optional[Callable[[batch_ingest.BatchManifest], None]] = None +_should_cancel: Optional[Callable[[str], bool]] = None + + +def configure( + *, + on_change: Callable[[batch_ingest.BatchManifest], None], + should_cancel: Callable[[str], bool], +) -> None: + """Wire the worker to the job store. Called once, by the router module.""" + global _on_change, _should_cancel + _on_change = on_change + _should_cancel = should_cancel + + +def queue_depth() -> int: + return _queue.qsize() + + +def is_running() -> bool: + return _worker is not None and _worker.is_alive() + + +def submit(batch_id: str) -> None: + """Enqueue a batch and make sure the worker exists. Raises `queue.Full`. + + `put_nowait` rather than `put`: blocking here would block the request + handler, which is the one thing an endpoint that returns 202 must never do. + """ + _queue.put_nowait(batch_id) + _ensure_worker() + + +def _ensure_worker() -> None: + global _worker + with _lock: + if _worker is not None and _worker.is_alive(): + return + _worker = threading.Thread(target=_loop, name="catalog-batch-worker", daemon=True) + _worker.start() + + +def _loop() -> None: + """Drain the queue forever. + + Every iteration is wrapped, because a worker that dies on one bad batch + would leave every future batch queued behind a thread that is not there - + a failure that looks, from the UI, exactly like a batch that is merely slow. + """ + while True: + batch_id = _queue.get() + try: + _run_one(batch_id) + except Exception: # noqa: BLE001 - see docstring + logger.exception("Batch worker: unhandled error on batch %s", batch_id) + finally: + _queue.task_done() + + if _queue.empty(): + # Retention runs when there is nothing waiting, so deleting old + # uploads never delays a batch and never sits on a request path. + try: + batch_ingest.purge_expired() + except Exception: # noqa: BLE001 - housekeeping must not kill the worker + logger.exception("Batch worker: purge failed") + + +def _run_one(batch_id: str) -> None: + on_change = _on_change or (lambda manifest: None) + cancelled = _should_cancel or (lambda _id: False) + batch_ingest.run_batch( + batch_id, + on_change=on_change, + should_cancel=lambda: cancelled(batch_id), + ) diff --git a/app/infrastructure/settings.py b/app/infrastructure/settings.py index fd8a0f6..ddacfa8 100644 --- a/app/infrastructure/settings.py +++ b/app/infrastructure/settings.py @@ -134,6 +134,41 @@ MODEL_ARTIFACTS_DIR = _dir( "MODEL_ARTIFACTS_DIR", _BACKEND_ROOT / "app" / "intelligence" / "artifacts" ) +# --------------------------------------------------------------------------- +# Batch catalog ingestion (multi-file upload -> the 11-stage pipeline) +# --------------------------------------------------------------------------- +# Staged uploads live under DATA_DIR because that path is already a declared +# volume (backend/Dockerfile). A batch that survives a container restart is the +# whole point of writing the files down instead of holding them in the worker +# thread the way the single-file path does. +# +# Every ceiling below is enforced in the application, not at the proxy. In +# production the browser calls mcp.nearle.ai.in directly, so neither nginx's +# client_max_body_size nor Caddy's request_body cap is in front of these +# endpoints - whatever Traefik defaults to is, and it is not ours to rely on. +BATCH_UPLOAD_DIR = _dir("BATCH_UPLOAD_DIR", DATA_DIR / "batch_uploads") + +# Per-file limits stay at the single-upload values (10MB / 2000 rows, see +# app/api/routers/store_catalog.py); these bound the BATCH on top of that. +BATCH_MAX_FILES = int(os.getenv("BATCH_MAX_FILES", "20")) +BATCH_MAX_TOTAL_BYTES = int(os.getenv("BATCH_MAX_TOTAL_BYTES", str(50 * 1024 * 1024))) +BATCH_MAX_TOTAL_ROWS = int(os.getenv("BATCH_MAX_TOTAL_ROWS", "20000")) + +# Batches waiting behind the one running. Past this the endpoint returns 429 +# rather than accepting work it has no intention of starting soon. +BATCH_QUEUE_MAX = int(os.getenv("BATCH_QUEUE_MAX", "4")) + +# Staged files are deleted this many days after the batch was created. Without +# this the upload directory only grows, on a host whose disk is the scarcest +# resource it has. +BATCH_RETENTION_DAYS = int(os.getenv("BATCH_RETENTION_DAYS", "7")) + +# Deliberately false. A batch interrupted by a restart is marked "interrupted" +# and waits for someone to press Resume. Auto-resuming would mean a container +# stuck in a restart loop re-runs the heaviest work in the app on every boot, +# which is precisely how a slow start turns into an unrecoverable spiral. +BATCH_AUTO_RESUME = _bool("BATCH_AUTO_RESUME", "false") + # Pristine copies of the bundled seed catalogs and pre-trained models, placed # here by the Dockerfile at a path that is never itself mounted over. # diff --git a/app/main.py b/app/main.py index 0db6465..98c3573 100644 --- a/app/main.py +++ b/app/main.py @@ -22,6 +22,7 @@ from fastapi.responses import FileResponse from app.infrastructure.persistence import restore_bundled_assets from app.infrastructure.settings import ( API_CORS_ORIGINS, + BATCH_AUTO_RESUME, BRAND_SYNC_INTERVAL_SECONDS, cleaned_env_names, ) @@ -30,7 +31,7 @@ from app.api.routers import health, brands, search, suggest, chat, catalog, syst from app.api.routers import stores, discounts, analytics as store_analytics, trending, recommendations, store_admin from app.api.routers import nutrition, nutrition_admin, upload from app.api.routers import auth, user_products, admin_train, mcp_info -from app.api.routers import store_catalog +from app.api.routers import store_catalog, batch_catalog from app.services.store_db import ensure_store_intelligence_schema from app.services.nutrition_db import ensure_nutrition_schema @@ -82,6 +83,31 @@ async def lifespan(_app: FastAPI): except Exception as e: logger.warning("Could not restore bundled assets: %s", e) + # Cheap, and it must happen before anyone can press Resume: a batch that a + # restart cut short is still marked "running" on disk, and until it is + # reconciled the UI shows it as in flight with nothing behind it. This only + # rewrites manifests - it deliberately starts no work. See + # batch_ingest.scan_interrupted() for why auto-resume is not the default. + try: + from app.core.batch_ingest import scan_interrupted + + interrupted = scan_interrupted() + if interrupted: + logger.info("Marked %d interrupted catalog batch(es).", len(interrupted)) + if interrupted and BATCH_AUTO_RESUME: + # Opt-in only. Importing the worker here rather than at module + # scope keeps the queue and its thread out of a boot that never + # needs them. + from app.core import batch_worker + + for batch_id in interrupted: + try: + batch_worker.submit(batch_id) + except Exception as exc: # noqa: BLE001 - a full queue is not fatal + logger.warning("Could not auto-resume batch %s: %s", batch_id, exc) + except Exception as e: + logger.warning("Could not reconcile interrupted catalog batches: %s", e) + def _async_init(): try: ensure_store_intelligence_schema() @@ -266,6 +292,7 @@ app.include_router(nutrition.router, prefix="/api") app.include_router(nutrition_admin.router, prefix="/api") app.include_router(upload.router, prefix="/api") app.include_router(store_catalog.router, prefix="/api") +app.include_router(batch_catalog.router, prefix="/api") app.include_router(mcp_info.router, prefix="/api") # MCP lives outside /api on purpose: it is a protocol endpoint for AI clients, diff --git a/orchestration/README.md b/orchestration/README.md index f7fbd3e..3af744c 100644 --- a/orchestration/README.md +++ b/orchestration/README.md @@ -144,10 +144,21 @@ If you find yourself writing business logic in this package, it belongs in | `embedding_refresh_job` | embed rows lacking a vector, verify the index | | `nutrition_enrichment_job` | fetch nutrition, score it, refit the two models | | `ml_training_job` | rebuild the store dataset, fit the served models, report | +| `batch_ingestion_job` | a batch of uploaded spreadsheets -> the 11 stages -> Postgres | -Four jobs rather than one because they differ in cost and cadence: ingestion is +Five jobs rather than one because they differ in cost and cadence: ingestion is offline and cheap, embedding drives the transformer, nutrition leaves the -machine, and ML training depends on order history rather than catalog freshness. +machine, ML training depends on order history rather than catalog freshness, and +batch ingestion is event-driven - it exists only when somebody uploads files. + +`batch_ingestion_job` shares every stage with the upload path in production; +both call `app/core/batch_ingest.py`. Pick a batch with run config: + +```json +{"ops": {"batch_manifest": {"config": {"batch_id": ""}}}} +``` + +Left empty it takes the oldest batch still waiting under `BATCH_UPLOAD_DIR`. From the command line: @@ -170,6 +181,7 @@ dagster asset materialize -m orchestration.definitions \ | `daily_nutrition_refresh` | 04:00 | slowest, and the only one making outbound calls | | `weekly_ml_retrain` | Sun 05:00 | training data is a deterministic simulation, so nightly would repeat itself | | `seed_catalog_sensor` | every 60s | re-ingest a brand when its seed file changes | +| `batch_upload_sensor` | every 60s | run an uploaded batch that is staged and waiting | Shipping them stopped is deliberate. This machine also runs Postgres, Ollama, the API and Vite; a schedule that started itself the moment `dagster dev` @@ -239,6 +251,17 @@ allocates the whole 8GB VPS (ollama 3G, backend 2560M, postgres 1G), so there is no headroom for a webserver plus daemon. Orchestration is a development concern here. +**This includes `batch_ingestion_job`.** The Batch Catalog Ingestion screen on +the site does not talk to Dagster and does not need it running. It stages the +uploaded files, then executes the same functions this job wraps +(`app/core/batch_ingest.py`) on a single bounded worker thread inside the API +process - see `app/core/batch_worker.py` for why one worker rather than a +thread per upload. The job here exists so the graph is inspectable and a batch +can be re-run from the Launchpad on a development machine. + +Do not turn `batch_upload_sensor` on next to a running API: both would pick up +the same staged batch. That is why it ships STOPPED like the rest. + Port 3030, not Dagster's default 3000 - `serve.py` binds `PORTS=3000,8000` and the frontend nginx also listens on 3000. diff --git a/orchestration/assets/__init__.py b/orchestration/assets/__init__.py index 1e1f453..c2069d3 100644 --- a/orchestration/assets/__init__.py +++ b/orchestration/assets/__init__.py @@ -1,3 +1,3 @@ -from orchestration.assets import catalog, ml, nutrition +from orchestration.assets import batch_catalog, catalog, ml, nutrition -__all__ = ["catalog", "nutrition", "ml"] +__all__ = ["catalog", "nutrition", "ml", "batch_catalog"] diff --git a/orchestration/assets/batch_catalog.py b/orchestration/assets/batch_catalog.py new file mode 100644 index 0000000..55d63a3 --- /dev/null +++ b/orchestration/assets/batch_catalog.py @@ -0,0 +1,315 @@ +"""Batch ingestion assets: a staged upload -> parsed -> ingested -> reported. + +Every asset delegates to `app.core.batch_ingest`, which is the same module the +production endpoint calls. What Dagster adds here is the lineage between the +steps, per-run metadata that survives a restart, retries on the step that +reaches the network, and a launchpad for re-running a batch that went wrong - +none of which the worker thread and its in-memory store can give. + +WHY THESE ASSETS ARE NOT PARTITIONED +------------------------------------ +The brand assets in `catalog.py` are partitioned per brand, because a brand is +a stable, enumerable thing and one partition per brand buys per-brand retry and +bounded memory. A batch is neither stable nor enumerable in advance - it comes +into existence when somebody uploads files. Modelling it as a partition would +mean adding a dynamic partition per upload and never removing it, so the +partition set grows forever and the UI fills with dead keys. The batch id lives +in run config instead, which is what run config is for. + +WHERE THIS RUNS +--------------- +`dagster dev` on a developer's machine. Dagster is not deployed: it is absent +from docker-compose.prod.yml, its dependencies are in a separate requirements +file the API image does not install, and it is not on any request path. The +production path executes the identical functions on a bounded worker thread +(app/core/batch_worker.py). That is deliberate - a webserver plus daemon costs +400-600MB resident, and the host this ships to does not have it to give. +""" +# NOTE: deliberately no `from __future__ import annotations` here. +# Dagster resolves the decorated function signatures at definition time to +# validate the `context` parameter and to infer asset input types. Under +# PEP 563/649 the annotations arrive as strings and that validation fails +# with "Cannot annotate `context` parameter with type AssetExecutionContext". +# Local Python is 3.14, which defers annotations by default, so this is not +# hypothetical. + +import time +from typing import Any, Dict, List + +from dagster import ( + AssetExecutionContext, + Backoff, + Failure, + Jitter, + MetadataValue, + RetryPolicy, + asset, +) + +from orchestration.config import BatchConfig, database_target, require_local_database + +# Same policy and same reasoning as catalog.py: only the step that leaves the +# machine is retried. A file that fails to parse will fail identically on a +# second attempt, and retrying it only delays the red run somebody has to read. +NETWORK_RETRY = RetryPolicy( + max_retries=2, delay=5, backoff=Backoff.EXPONENTIAL, jitter=Jitter.PLUS_MINUS +) + +GROUP = "batch_catalog" + + +def _pick_batch_id(config: BatchConfig) -> str: + """Run config wins; otherwise the oldest batch still waiting to run.""" + from app.core import batch_ingest + + if config.batch_id: + return config.batch_id + + waiting = [ + m for m in batch_ingest.list_manifests() + if m.status in (batch_ingest.QUEUED, batch_ingest.INTERRUPTED) + ] + if not waiting: + raise Failure( + description=( + "No batch_id given and no staged batch is waiting to run.\n\n" + "Upload one through Admin -> Batch Catalog Ingestion on the site, " + "or pass {\"batch_id\": \"...\"} in the Launchpad. Staged batches " + "live under BATCH_UPLOAD_DIR ({}).".format(batch_ingest.batch_root()) + ), + metadata={"staged_batches": len(batch_ingest.list_manifests())}, + ) + return sorted(waiting, key=lambda m: m.created_at)[0].batch_id + + +@asset( + group_name=GROUP, + description="The staged upload batch this run will ingest.", +) +def batch_manifest(context: AssetExecutionContext, config: BatchConfig) -> Dict[str, Any]: + """Root of the lineage: resolve which batch, and prove it is readable. + + Resolving here rather than inside each downstream asset means the run page + shows which files a run actually touched, up front, instead of it being + discovered halfway through. + """ + from app.core import batch_ingest + + batch_id = _pick_batch_id(config) + manifest = batch_ingest.read_manifest(batch_id) + if manifest is None: + raise Failure( + description="No manifest for batch '{}' under {}.".format( + batch_id, batch_ingest.batch_root() + ), + metadata={"batch_id": batch_id}, + ) + + pending = [f for f in manifest.files if f.status == batch_ingest.QUEUED] + context.add_output_metadata( + { + "batch_id": batch_id, + "status": manifest.status, + "files_total": manifest.files_total, + "files_pending": len(pending), + "use_llm": config.use_llm or manifest.use_llm, + "fetch_images": config.fetch_images or manifest.fetch_images, + "database": database_target(), + "files": MetadataValue.md( + "\n".join( + "- `{}` - {}".format(f.filename, f.status) for f in manifest.files + ) or "_none_" + ), + } + ) + return {"batch_id": batch_id, "files_pending": len(pending)} + + +@asset( + group_name=GROUP, + description="Column mapping and row count for each staged file.", +) +def batch_parsed_files( + context: AssetExecutionContext, batch_manifest: Dict[str, Any] +) -> List[Dict[str, Any]]: + """Parse every staged file without ingesting anything. + + This is the Dagster equivalent of the UI's Preview step, and it exists for + the same reason: the mapping from a store's headers onto catalog fields is + a guess, and it is much cheaper to see it now than after eleven stages have + run over two thousand rows. + + A file that will not parse is reported, not raised. One unreadable sheet + among ten is a data problem worth seeing, not a reason to block the nine. + """ + from app.core import batch_ingest + from app.core import store_catalog_pipeline as pipeline + + batch_id = batch_manifest["batch_id"] + manifest = batch_ingest.read_manifest(batch_id) + directory = batch_ingest.batch_dir(batch_id) + + parsed: List[Dict[str, Any]] = [] + lines = [] + for entry in manifest.files: + if not entry.stored_name: + parsed.append({"filename": entry.filename, "ok": False, + "error": entry.detail or "not staged"}) + lines.append("- `{}` - **not staged**".format(entry.filename)) + continue + try: + content = (directory / entry.stored_name).read_bytes() + frame, mapping = pipeline.parse_spreadsheet(entry.filename, content) + except Exception as exc: # noqa: BLE001 - see docstring + parsed.append({"filename": entry.filename, "ok": False, "error": str(exc)}) + lines.append("- `{}` - **unreadable**: {}".format(entry.filename, exc)) + continue + + recognised = {f: str(c) for f, c in mapping.columns.items()} + parsed.append({ + "filename": entry.filename, + "ok": True, + "rows": int(len(frame)), + "recognised_columns": recognised, + "unrecognised_columns": list(mapping.unrecognised), + }) + lines.append("- `{}` - {} rows, {} mapped column(s)".format( + entry.filename, len(frame), len(recognised) + )) + + context.add_output_metadata( + { + "batch_id": batch_id, + "files": len(parsed), + "files_readable": sum(1 for p in parsed if p["ok"]), + "rows_total": sum(int(p.get("rows") or 0) for p in parsed), + "detail": MetadataValue.md("\n".join(lines) or "_none_"), + } + ) + return parsed + + +@asset( + group_name=GROUP, + retry_policy=NETWORK_RETRY, + description="The 11-stage pipeline run over every pending file in the batch.", +) +def batch_catalog_rows( + context: AssetExecutionContext, + config: BatchConfig, + batch_manifest: Dict[str, Any], + batch_parsed_files: List[Dict[str, Any]], +) -> Dict[str, Any]: + """Stages 1-11, per file, via the shared `batch_ingest.run_batch`. + + THE WRITE GUARD IS CALLED HERE, not only in the reporting asset downstream. + `run_batch` reaches stage 11 and upserts, so by the time the next asset runs + the write has already happened - checking there would be checking after the + fact. backend/.env points at production; see orchestration/config.py. + + Per-file progress is logged rather than pushed anywhere. In production the + same callback feeds the polling endpoint; in Dagster the run log IS the + progress view, and writing to both would be two sources of truth. + """ + from app.core import batch_ingest + + target = require_local_database("batch_catalog_rows") + batch_id = batch_manifest["batch_id"] + started = time.time() + + # Run config overrides what the upload asked for, so a batch uploaded with + # images off can be re-run with them on without re-uploading. + manifest = batch_ingest.read_manifest(batch_id) + manifest.use_llm = bool(config.use_llm or manifest.use_llm) + manifest.fetch_images = bool(config.fetch_images or manifest.fetch_images) + batch_ingest.write_manifest(manifest) + + seen = {"file": None} + + def on_change(current): + name = current.current_file + if name and name != seen["file"]: + seen["file"] = name + context.log.info("Ingesting %s", name) + + result = batch_ingest.run_batch(batch_id, on_change=on_change) + totals = result.totals() + + context.add_output_metadata( + { + "batch_id": batch_id, + "status": result.status, + "files_total": result.files_total, + "files_done": result.files_done, + "files_failed": result.files_failed, + "rows_ingested": totals["rows_total"], + "products_built": totals["products_built"], + "inserted": totals["inserted"], + "backfilled": totals["backfilled"], + "unchanged": totals["skipped_existing"], + "rejected": totals["rejected"], + "brands": ", ".join(result.brands()) or "none", + "database": target, + "cleanup": "False (never deletes rows outside this batch)", + "duration_s": round(time.time() - started, 2), + } + ) + return result.to_dict() + + +@asset( + group_name=GROUP, + description="Per-file outcome of the batch, and a red run if none of it landed.", +) +def batch_ingest_report( + context: AssetExecutionContext, batch_catalog_rows: Dict[str, Any] +) -> Dict[str, Any]: + """Turn the batch result into a readable report, and decide run colour. + + A batch where every file failed materializes GREEN without this: nothing + raised, so from Dagster's point of view the work completed. That is the + worst outcome available here - the run looks fine and the catalog did not + change - so it is made loud, in the same spirit as `raw_products` failing on + an empty brand. + + A partial batch stays green. Four files landing out of five is a data + problem to read about, not a pipeline failure, and failing the run would + make the four look like they had not happened. + """ + files = batch_catalog_rows.get("files") or [] + done = [f for f in files if f.get("status") == "done"] + failed = [f for f in files if f.get("status") == "failed"] + + lines = [] + for entry in files: + lines.append("- `{}` - **{}** - {}".format( + entry.get("filename"), entry.get("status"), entry.get("detail") or "" + )) + + context.add_output_metadata( + { + "batch_id": batch_catalog_rows.get("batch_id"), + "status": batch_catalog_rows.get("status"), + "files_done": len(done), + "files_failed": len(failed), + "totals": MetadataValue.json(batch_catalog_rows.get("totals") or {}), + "report": MetadataValue.md("\n".join(lines) or "_no files_"), + } + ) + + if files and not done: + raise Failure( + description=( + "Every file in batch '{}' failed; nothing reached the catalog.".format( + batch_catalog_rows.get("batch_id") + ) + ), + metadata={"files_failed": len(failed)}, + ) + + return { + "batch_id": batch_catalog_rows.get("batch_id"), + "status": batch_catalog_rows.get("status"), + "files_done": len(done), + "files_failed": len(failed), + } diff --git a/orchestration/config.py b/orchestration/config.py index 5d0350c..cc5354c 100644 --- a/orchestration/config.py +++ b/orchestration/config.py @@ -94,6 +94,25 @@ class BrandConfig(Config): max_products_per_brand: Optional[int] = None +class BatchConfig(Config): + """Per-run config for the uploaded-spreadsheet batch job. + + `batch_id` selects a batch staged under BATCH_UPLOAD_DIR by the API. Left + empty, the job picks the oldest batch still queued, which is what the + sensor wants and what makes a manual "run it now" from the Launchpad + convenient. + + The two network flags default OFF for the same reason they do on + BrandConfig, with more force: a batch is up to twenty files, so leaving + image search on would fire tens of thousands of outbound requests from one + run. + """ + + batch_id: Optional[str] = None + use_llm: bool = False + fetch_images: bool = False + + def resolve_brands(cfg_brands: Optional[List[str]]) -> List[str]: """Run config wins; otherwise ACTIVE_BRANDS; otherwise whatever the DB has. diff --git a/orchestration/definitions.py b/orchestration/definitions.py index ab7df34..d5b7d7f 100644 --- a/orchestration/definitions.py +++ b/orchestration/definitions.py @@ -61,11 +61,11 @@ _ENV_SOURCE = _load_orchestration_env() from dagster import Definitions, load_assets_from_modules, multiprocess_executor # noqa: E402 from orchestration import config as orchestration_config # noqa: E402 -from orchestration.assets import catalog, ml, nutrition # noqa: E402 +from orchestration.assets import batch_catalog, catalog, ml, nutrition # noqa: E402 from orchestration.jobs import ALL_JOBS # noqa: E402 from orchestration.schedules import ALL_SCHEDULES, ALL_SENSORS # noqa: E402 -_all_assets = load_assets_from_modules([catalog, nutrition, ml]) +_all_assets = load_assets_from_modules([catalog, nutrition, ml, batch_catalog]) defs = Definitions( assets=_all_assets, diff --git a/orchestration/jobs.py b/orchestration/jobs.py index d9ea729..363ca22 100644 --- a/orchestration/jobs.py +++ b/orchestration/jobs.py @@ -76,9 +76,31 @@ ml_training_job = define_asset_job( ), ) +# Uploaded spreadsheets -> the same 11 stages -> brand tables. +# +# Separate from catalog_ingestion_job because the source is different, not +# because the work is: that job reads seed JSON for a brand, this one reads a +# batch of files somebody uploaded. They share every stage in between, and +# keeping them apart means a failure here does not read as the brand pipeline +# being broken. +batch_ingestion_job = define_asset_job( + name="batch_ingestion_job", + description=( + "Ingest a batch of uploaded store spreadsheets: resolve the batch, " + "parse each file, run the 11 stages, and report what landed." + ), + selection=AssetSelection.assets( + "batch_manifest", + "batch_parsed_files", + "batch_catalog_rows", + "batch_ingest_report", + ), +) + ALL_JOBS = [ catalog_ingestion_job, embedding_refresh_job, nutrition_enrichment_job, ml_training_job, + batch_ingestion_job, ] diff --git a/orchestration/schedules.py b/orchestration/schedules.py index ee3780f..3eb5fe4 100644 --- a/orchestration/schedules.py +++ b/orchestration/schedules.py @@ -31,6 +31,7 @@ from dagster import ( ) from orchestration.jobs import ( + batch_ingestion_job, catalog_ingestion_job, embedding_refresh_job, ml_training_job, @@ -149,4 +150,66 @@ def seed_catalog_sensor(context: SensorEvaluationContext): ] -ALL_SENSORS = [seed_catalog_sensor] +@sensor( + job=batch_ingestion_job, + minimum_interval_seconds=60, + default_status=DefaultSensorStatus.STOPPED, + description="Run a batch of uploaded spreadsheets once one is staged and waiting.", +) +def batch_upload_sensor(context: SensorEvaluationContext): + """Watch BATCH_UPLOAD_DIR for batches the API staged but did not run. + + There is NO schedule for batch ingestion, deliberately. A batch exists + because a person uploaded files; there is nothing to do on a cron, and a + schedule that woke up hourly to find nothing would be pure cost on a + machine that is already sharing itself with Postgres, the API and Vite. + + Like seed_catalog_sensor this is as simple as it can be - a directory + listing once a minute, no watcher process, no broker. The cursor is the set + of batch ids already requested, so a batch is launched once even though it + stays on disk afterwards. + + NOTE ON THE NORMAL PRODUCTION PATH: nothing here is involved. The API runs + a staged batch itself, on a bounded worker thread, and this sensor exists + for the development machine where Dagster is the thing driving the work. + Turning it on alongside a running API would mean both trying to ingest the + same batch, which is why it ships STOPPED. + """ + from app.core import batch_ingest + + root = batch_ingest.batch_root() + if not root.exists(): + return SkipReason("Batch upload directory {} does not exist.".format(root)) + + waiting = [ + m for m in batch_ingest.list_manifests() + if m.status in (batch_ingest.QUEUED, batch_ingest.INTERRUPTED) + ] + if not waiting: + return SkipReason("No staged batch is waiting to run.") + + already = set(filter(None, (context.cursor or "").split(","))) + fresh = [m for m in waiting if m.batch_id not in already] + if not fresh: + return SkipReason( + "{} staged batch(es), all already requested.".format(len(waiting)) + ) + + # Keep the cursor bounded - it is a string in the Dagster instance, not a + # log, and only the recent tail is ever consulted. + context.update_cursor(",".join(list(already | {m.batch_id for m in fresh})[-200:])) + + context.log.info( + "Requesting runs for staged batch(es): %s", + ", ".join(m.batch_id for m in fresh), + ) + return [ + RunRequest( + run_key=m.batch_id, + run_config={"ops": {"batch_manifest": {"config": {"batch_id": m.batch_id}}}}, + ) + for m in fresh + ] + + +ALL_SENSORS = [seed_catalog_sensor, batch_upload_sensor] diff --git a/tests/test_batch_catalog_ingest.py b/tests/test_batch_catalog_ingest.py new file mode 100644 index 0000000..d156171 --- /dev/null +++ b/tests/test_batch_catalog_ingest.py @@ -0,0 +1,640 @@ +"""Tests for multi-file batch ingestion. + +Follows test_store_catalog_pipeline.py: monkeypatch the storage and embedding +boundary *on the pipeline module object*, build real spreadsheets in memory, +and point the batch directory at tmp_path so nothing is written into the repo. + +Nothing here loads sentence-transformers or torch, and nothing reaches the +network: every batch runs with use_llm=False and fetch_images=False, which are +also the defaults the endpoint ships. +""" +from __future__ import annotations + +import io +import json +import queue +import threading +import time + +import pytest + +from app.core import batch_ingest +from app.core import store_catalog_pipeline as pipeline + +openpyxl = pytest.importorskip("openpyxl") + +HEADERS = ["Product Name", "Category", "Brand"] +ROWS = [ + ["Amul Butter 100g", "Butter", "Amul"], + ["Amul Cheese Slices 200g", "Cheese", "Amul"], +] + + +def _sheet(headers=HEADERS, rows=ROWS) -> bytes: + wb = openpyxl.Workbook() + ws = wb.active + ws.append(headers) + for row in rows: + ws.append(row) + buf = io.BytesIO() + wb.save(buf) + return buf.getvalue() + + +def _csv(headers=HEADERS, rows=ROWS) -> bytes: + lines = [",".join(headers)] + [",".join(r) for r in rows] + return ("\n".join(lines)).encode("utf-8") + + +@pytest.fixture(autouse=True) +def _isolate_sku_counter(tmp_path, monkeypatch): + """Same reason as test_store_catalog_pipeline: keep the SKU sequence file + out of the working tree and make numbering deterministic per test.""" + from app.services import sku_service + monkeypatch.setattr(sku_service, "_data_dir", tmp_path / "sku_sequences") + + +@pytest.fixture(autouse=True) +def batch_root(tmp_path, monkeypatch): + """Point BATCH_UPLOAD_DIR at tmp_path. + + `batch_ingest.batch_root()` reads the setting at call time precisely so + this is possible; patching it at import time would not reach the module. + """ + root = tmp_path / "batch_uploads" + monkeypatch.setattr(batch_ingest, "BATCH_UPLOAD_DIR", root) + return root + + +@pytest.fixture(autouse=True) +def no_background_worker(monkeypatch): + """Keep the real worker thread out of every test that does not want it. + + THIS FIXTURE EXISTS BECAUSE ITS ABSENCE CORRUPTED THE WORKING TREE. + + `batch_root()` resolves BATCH_UPLOAD_DIR at call time. That is correct in + production, where the setting never moves. In a test it moves twice: the + fixture points it at tmp_path, and monkeypatch puts it back when the test + ends. An endpoint test POSTs to /ingest, the worker starts on a background + thread, the test returns - and the worker, waking up after teardown, writes + its manifest into the repository's real `data/batch_uploads/`. Two such + directories were committed-adjacent before this was noticed. + + So `submit` is stubbed by default and endpoint tests assert on what they + are actually about. The two tests that exercise the worker itself take + `real_submit` and wait for it to drain before returning. + """ + from app.core import batch_worker + + real = batch_worker.submit + submitted: list = [] + monkeypatch.setattr(batch_worker, "submit", submitted.append) + + class Handle: + pass + + handle = Handle() + handle.submitted = submitted + handle.real_submit = real + return handle + + +@pytest.fixture +def store(monkeypatch): + """A fake brand table shared by every file in a batch.""" + table: dict = {} + + def fake_upsert(brand, rows, cleanup=False): + assert cleanup is False, "cleanup=True would delete the brand's existing catalog" + for row in rows: + table[row["image_id"]] = dict(row) + return len(rows) + + monkeypatch.setattr(pipeline, "upsert_brand_products", fake_upsert) + monkeypatch.setattr(pipeline, "get_products_by_brand", lambda b, **kw: list(table.values())) + monkeypatch.setattr(pipeline, "embed_texts", lambda texts: [[0.0] * 384 for _ in texts]) + return table + + +# --------------------------------------------------------------------------- +# Staging +# --------------------------------------------------------------------------- +def test_staging_writes_files_and_a_manifest(batch_root): + manifest = batch_ingest.stage_batch([("a.csv", _csv()), ("b.xlsx", _sheet())]) + + directory = batch_root / manifest.batch_id + assert (directory / "manifest.json").exists() + assert sorted(p.name for p in directory.iterdir()) == [ + "00_a.csv", "01_b.xlsx", "manifest.json", + ] + assert [f.status for f in manifest.files] == ["queued", "queued"] + assert manifest.files_total == 2 + + +def test_stored_names_cannot_escape_the_batch_directory(batch_root): + """UploadFile.filename is attacker-controlled and is used to build a path.""" + manifest = batch_ingest.stage_batch([ + ("../../evil.csv", _csv()), + ("..\\windows\\evil2.csv", _csv()), + ("/etc/passwd", _csv()), + ]) + + directory = batch_root / manifest.batch_id + for entry in manifest.files: + written = directory / entry.stored_name + assert written.resolve().parent == directory.resolve() + assert ".." not in entry.stored_name + assert "/" not in entry.stored_name and "\\" not in entry.stored_name + # And nothing landed outside. + assert not (batch_root.parent / "evil.csv").exists() + + +def test_invalid_files_are_recorded_not_dropped(batch_root): + """An operator who selected three files and sees two must be told why.""" + manifest = batch_ingest.stage_batch( + [("good.csv", _csv())], + invalid=[("broken.pdf", "Unsupported file type")], + ) + assert manifest.files_total == 2 + bad = [f for f in manifest.files if f.filename == "broken.pdf"][0] + assert bad.status == "failed" + assert bad.detail == "Unsupported file type" + assert bad.stored_name == "" + + +def test_manifest_round_trips_through_disk(batch_root): + manifest = batch_ingest.stage_batch([("a.csv", _csv())], fetch_images=True) + reloaded = batch_ingest.read_manifest(manifest.batch_id) + assert reloaded.batch_id == manifest.batch_id + assert reloaded.fetch_images is True + assert [f.filename for f in reloaded.files] == ["a.csv"] + + +def test_unreadable_manifest_is_skipped_not_raised(batch_root): + manifest = batch_ingest.stage_batch([("a.csv", _csv())]) + batch_ingest.manifest_path(manifest.batch_id).write_text("{ not json", encoding="utf-8") + assert batch_ingest.read_manifest(manifest.batch_id) is None + assert batch_ingest.list_manifests() == [] + + +def test_batch_dir_rejects_a_traversing_id(batch_root): + with pytest.raises(ValueError): + batch_ingest.batch_dir("../escape") + + +# --------------------------------------------------------------------------- +# Running a batch +# --------------------------------------------------------------------------- +def test_every_file_runs_and_totals_are_summed(store, batch_root): + manifest = batch_ingest.stage_batch([("a.csv", _csv()), ("b.xlsx", _sheet())]) + result = batch_ingest.run_batch(manifest.batch_id) + + assert result.status == "done" + assert result.files_done == 2 + assert result.files_failed == 0 + assert [f.status for f in result.files] == ["done", "done"] + # Both files carry the same two products, so the second file backfills or + # skips rather than inserting again - but every row was seen. + assert result.totals()["rows_total"] == 4 + assert result.totals()["products_built"] >= 4 + assert "Amul" in result.brands() + + +def test_one_unparseable_file_does_not_stop_the_others(store, batch_root): + """The whole point of a batch: file 2 failing must not lose files 1 and 3.""" + manifest = batch_ingest.stage_batch([ + ("good1.csv", _csv()), + ("truncated.xlsx", _sheet()), + ]) + # Corrupt the second staged file after staging, so run_batch meets it cold. + # It has to be the .xlsx: pandas' CSV reader is tolerant enough that binary + # noise in a .csv still parses into *something*, which is not the failure + # this test is about. + directory = batch_ingest.batch_dir(manifest.batch_id) + (directory / manifest.files[1].stored_name).write_bytes(b"\x00\x01\x02not a sheet") + + result = batch_ingest.run_batch(manifest.batch_id) + + assert result.files_done == 1 + assert result.files_failed == 1 + assert result.status == "partial", "a partial batch must not read as success" + assert result.files[1].detail + + +def test_a_batch_with_a_pre_failed_file_is_partial(store, batch_root): + manifest = batch_ingest.stage_batch( + [("good.csv", _csv())], + invalid=[("broken.pdf", "Unsupported file type")], + ) + result = batch_ingest.run_batch(manifest.batch_id) + assert result.status == "partial" + assert result.files_done == 1 + assert result.files_failed == 1 + + +def test_storage_error_marks_the_file_failed_not_done(monkeypatch, store, batch_root): + """Rows built but nothing written is a failure, exactly as in the single-file path.""" + def exploding_upsert(brand, rows, cleanup=False): + raise RuntimeError("connection refused") + + monkeypatch.setattr(pipeline, "upsert_brand_products", exploding_upsert) + + manifest = batch_ingest.stage_batch([("a.csv", _csv())]) + result = batch_ingest.run_batch(manifest.batch_id) + + assert result.files[0].status == "failed" + assert "connection refused" in result.files[0].detail + assert result.status == "failed" + + +def test_progress_reaches_the_callback(store, batch_root): + manifest = batch_ingest.stage_batch([("a.csv", _csv())]) + seen = [] + + def on_change(current): + entry = current.files[0] + if entry.stage_index: + seen.append((entry.stage_index, entry.stage_name)) + + batch_ingest.run_batch(manifest.batch_id, on_change=on_change) + + assert seen, "expected stage progress" + assert max(s for s, _ in seen) == pipeline.TOTAL_STAGES + assert seen[-1][1] == pipeline.STAGE_NAMES[-1] + + +def test_cancel_stops_before_the_next_file(store, batch_root): + manifest = batch_ingest.stage_batch([ + ("a.csv", _csv()), ("b.csv", _csv()), ("c.csv", _csv()), + ]) + calls = {"n": 0} + + def should_cancel(): + # Let the first file through, then cancel. + calls["n"] += 1 + return calls["n"] > 1 + + result = batch_ingest.run_batch(manifest.batch_id, should_cancel=should_cancel) + + assert result.files[0].status == "done" + assert [f.status for f in result.files[1:]] == ["cancelled", "cancelled"] + + +def test_rerunning_a_finished_batch_does_nothing(store, batch_root): + manifest = batch_ingest.stage_batch([("a.csv", _csv())]) + batch_ingest.run_batch(manifest.batch_id) + before = batch_ingest.read_manifest(manifest.batch_id).files[0].finished_at + + again = batch_ingest.run_batch(manifest.batch_id) + assert again.files[0].finished_at == before, "a done file must not run twice" + + +# --------------------------------------------------------------------------- +# Recovery and retention +# --------------------------------------------------------------------------- +def test_scan_interrupted_marks_but_does_not_run(monkeypatch, batch_root): + """A crash loop must not re-enter the heaviest work in the app on every boot.""" + monkeypatch.setattr(batch_ingest, "BATCH_AUTO_RESUME", False) + + manifest = batch_ingest.stage_batch([("a.csv", _csv()), ("b.csv", _csv())]) + manifest.status = "running" + manifest.files[0].status = "done" + manifest.files[1].status = "running" + batch_ingest.write_manifest(manifest) + + resumable = batch_ingest.scan_interrupted() + + assert resumable == [manifest.batch_id] + reloaded = batch_ingest.read_manifest(manifest.batch_id) + assert reloaded.status == "interrupted" + assert reloaded.files[0].status == "done", "a finished file is not re-run" + assert reloaded.files[1].status == "queued", "the in-flight file is queued again" + + +def test_scan_interrupted_leaves_finished_batches_alone(store, batch_root): + manifest = batch_ingest.stage_batch([("a.csv", _csv())]) + batch_ingest.run_batch(manifest.batch_id) + + assert batch_ingest.scan_interrupted() == [] + assert batch_ingest.read_manifest(manifest.batch_id).status == "done" + + +def test_purge_removes_old_finished_batches_only(store, batch_root, monkeypatch): + monkeypatch.setattr(batch_ingest, "BATCH_RETENTION_DAYS", 7) + + old = batch_ingest.stage_batch([("a.csv", _csv())]) + batch_ingest.run_batch(old.batch_id) + stale = batch_ingest.read_manifest(old.batch_id) + stale.created_at = time.time() - 30 * 86400 + batch_ingest.write_manifest(stale) + + recent = batch_ingest.stage_batch([("b.csv", _csv())]) + batch_ingest.run_batch(recent.batch_id) + + removed = batch_ingest.purge_expired() + + assert removed == [old.batch_id] + assert not batch_ingest.batch_dir(old.batch_id).exists() + assert batch_ingest.batch_dir(recent.batch_id).exists() + + +def test_purge_spares_an_old_batch_that_never_ran(batch_root, monkeypatch): + """Deleting the only copy of an unrun batch's input is not how to debug it.""" + monkeypatch.setattr(batch_ingest, "BATCH_RETENTION_DAYS", 7) + + manifest = batch_ingest.stage_batch([("a.csv", _csv())]) + manifest.created_at = time.time() - 30 * 86400 + batch_ingest.write_manifest(manifest) + + assert batch_ingest.purge_expired() == [] + assert batch_ingest.batch_dir(manifest.batch_id).exists() + + +def test_purge_is_disabled_when_retention_is_zero(batch_root, monkeypatch): + monkeypatch.setattr(batch_ingest, "BATCH_RETENTION_DAYS", 0) + manifest = batch_ingest.stage_batch([("a.csv", _csv())]) + manifest.created_at = time.time() - 999 * 86400 + manifest.status = "done" + batch_ingest.write_manifest(manifest) + assert batch_ingest.purge_expired() == [] + + +# --------------------------------------------------------------------------- +# The worker +# --------------------------------------------------------------------------- +def test_worker_runs_one_batch_at_a_time(store, batch_root, monkeypatch, + no_background_worker): + """Two batches must serialise, not compete for the single vCPU.""" + from app.core import batch_worker + + # Idle housekeeping deletes directories. It runs after the queue drains, + # which can be after teardown has moved BATCH_UPLOAD_DIR back to the real + # one - so it is disabled here rather than raced with. + monkeypatch.setattr(batch_worker.batch_ingest, "purge_expired", lambda *a, **k: []) + + overlapping = [] + active = {"n": 0} + guard = threading.Lock() + finished = threading.Event() + done_count = {"n": 0} + + real_run = batch_ingest.run_batch + + def instrumented(batch_id, **kwargs): + with guard: + active["n"] += 1 + if active["n"] > 1: + overlapping.append(batch_id) + try: + time.sleep(0.05) + return real_run(batch_id, **kwargs) + finally: + with guard: + active["n"] -= 1 + done_count["n"] += 1 + if done_count["n"] == 2: + finished.set() + + monkeypatch.setattr(batch_worker.batch_ingest, "run_batch", instrumented) + + first = batch_ingest.stage_batch([("a.csv", _csv())]) + second = batch_ingest.stage_batch([("b.csv", _csv())]) + no_background_worker.real_submit(first.batch_id) + no_background_worker.real_submit(second.batch_id) + + assert finished.wait(timeout=60), "batches did not finish" + # Drain before teardown moves BATCH_UPLOAD_DIR back. + batch_worker._queue.join() + assert overlapping == [], "two batches ran concurrently" + + +def test_worker_survives_a_batch_that_raises(store, batch_root, monkeypatch, + no_background_worker): + """A worker that dies would leave every later batch queued behind nothing.""" + from app.core import batch_worker + + monkeypatch.setattr(batch_worker.batch_ingest, "purge_expired", lambda *a, **k: []) + + done = threading.Event() + real_run = batch_ingest.run_batch + seen = [] + + def flaky(batch_id, **kwargs): + seen.append(batch_id) + if len(seen) == 1: + raise RuntimeError("boom") + try: + return real_run(batch_id, **kwargs) + finally: + done.set() + + monkeypatch.setattr(batch_worker.batch_ingest, "run_batch", flaky) + + bad = batch_ingest.stage_batch([("a.csv", _csv())]) + good = batch_ingest.stage_batch([("b.csv", _csv())]) + no_background_worker.real_submit(bad.batch_id) + no_background_worker.real_submit(good.batch_id) + + assert done.wait(timeout=60), "the worker died on the first batch" + batch_worker._queue.join() + assert seen == [bad.batch_id, good.batch_id] + + +# --------------------------------------------------------------------------- +# Endpoints +# --------------------------------------------------------------------------- +INGEST = "/api/admin/catalog-batch/ingest" +PREVIEW = "/api/admin/catalog-batch/preview" + + +def _files(*pairs): + return [("files", (name, io.BytesIO(content), "application/octet-stream")) + for name, content in pairs] + + +def test_endpoints_require_an_admin(client, user_headers): + assert client.post(PREVIEW, files=_files(("a.csv", _csv()))).status_code == 401 + assert client.post( + PREVIEW, files=_files(("a.csv", _csv())), headers=user_headers + ).status_code == 403 + assert client.get("/api/admin/catalog-batch/batches").status_code == 401 + + +def test_preview_reports_each_file_separately(client, admin_headers): + response = client.post( + PREVIEW, + files=_files(("a.csv", _csv()), ("broken.pdf", b"%PDF-1.4")), + headers=admin_headers, + ) + assert response.status_code == 200 + body = response.json() + assert body["files_total"] == 2 + assert body["files_ok"] == 1 + good = [f for f in body["files"] if f["filename"] == "a.csv"][0] + assert good["ok"] is True + assert good["rows_total"] == 2 + assert good["brand_column_present"] is True + bad = [f for f in body["files"] if f["filename"] == "broken.pdf"][0] + assert bad["ok"] is False and bad["error"] + assert body["stages"] == list(pipeline.STAGE_NAMES) + + +def test_ingest_returns_202_and_a_batch_id(client, admin_headers, store, batch_root): + response = client.post( + INGEST, files=_files(("a.csv", _csv()), ("b.xlsx", _sheet())), headers=admin_headers + ) + assert response.status_code == 202 + body = response.json() + assert body["files_total"] == 2 + assert body["use_llm"] is False, "the LLM stage must be opt-in for a batch" + assert body["fetch_images"] is False, "image search must be opt-in for a batch" + assert len(body["files"]) == 2 + assert body["files"][0]["total_stages"] == pipeline.TOTAL_STAGES + + polled = client.get( + f"/api/admin/catalog-batch/batches/{body['batch_id']}", headers=admin_headers + ) + assert polled.status_code == 200 + assert polled.json()["batch_id"] == body["batch_id"] + + +def test_ingest_rejects_a_batch_where_nothing_is_usable(client, admin_headers): + response = client.post( + INGEST, files=_files(("broken.pdf", b"%PDF-1.4")), headers=admin_headers + ) + assert response.status_code == 400 + assert "broken.pdf" in response.json()["detail"] + + +def test_ingest_keeps_the_good_files_when_one_is_bad(client, admin_headers, store, batch_root): + response = client.post( + INGEST, + files=_files(("a.csv", _csv()), ("broken.pdf", b"%PDF-1.4")), + headers=admin_headers, + ) + assert response.status_code == 202 + body = response.json() + assert body["files_total"] == 2 + statuses = {f["filename"]: f["status"] for f in body["files"]} + assert statuses["broken.pdf"] == "failed" + + +def test_too_many_files_is_413(client, admin_headers, monkeypatch): + from app.api.routers import batch_catalog + + monkeypatch.setattr(batch_catalog, "BATCH_MAX_FILES", 2) + response = client.post( + INGEST, + files=_files(("a.csv", _csv()), ("b.csv", _csv()), ("c.csv", _csv())), + headers=admin_headers, + ) + assert response.status_code == 413 + assert "2-file limit" in response.json()["detail"] + + +def test_a_file_over_the_size_limit_is_413(client, admin_headers, monkeypatch): + from app.api.routers import batch_catalog + + monkeypatch.setattr(batch_catalog, "MAX_UPLOAD_BYTES", 64) + response = client.post( + INGEST, files=_files(("big.csv", b"x" * 500)), headers=admin_headers + ) + assert response.status_code == 413 + + +def test_too_many_rows_across_the_batch_is_413(client, admin_headers, monkeypatch): + from app.api.routers import batch_catalog + + monkeypatch.setattr(batch_catalog, "BATCH_MAX_TOTAL_ROWS", 3) + response = client.post( + INGEST, + files=_files(("a.csv", _csv()), ("b.csv", _csv())), + headers=admin_headers, + ) + assert response.status_code == 413 + assert "3 rows" in response.json()["detail"] + + +def test_a_full_queue_is_429_and_the_batch_is_kept(client, admin_headers, monkeypatch, batch_root): + """Refusing work is fine; losing the upload that was refused is not.""" + from app.api.routers import batch_catalog + + def full(_batch_id): + raise queue.Full() + + monkeypatch.setattr(batch_catalog.batch_worker, "submit", full) + + response = client.post(INGEST, files=_files(("a.csv", _csv())), headers=admin_headers) + assert response.status_code == 429 + + staged = batch_ingest.list_manifests() + assert len(staged) == 1, "the refused batch must still be on disk" + assert staged[0].status == "queued" + + +def test_unknown_batch_is_404(client, admin_headers): + response = client.get( + "/api/admin/catalog-batch/batches/deadbeef", headers=admin_headers + ) + assert response.status_code == 404 + + +def test_resume_requeues_an_interrupted_batch(client, admin_headers, monkeypatch, batch_root): + from app.api.routers import batch_catalog + + submitted = [] + monkeypatch.setattr(batch_catalog.batch_worker, "submit", submitted.append) + + manifest = batch_ingest.stage_batch([("a.csv", _csv())]) + manifest.status = "interrupted" + batch_ingest.write_manifest(manifest) + + response = client.post( + f"/api/admin/catalog-batch/batches/{manifest.batch_id}/resume", headers=admin_headers + ) + assert response.status_code == 200 + assert submitted == [manifest.batch_id] + + +def test_resume_on_a_finished_batch_is_409(client, admin_headers, store, batch_root): + manifest = batch_ingest.stage_batch([("a.csv", _csv())]) + batch_ingest.run_batch(manifest.batch_id) + + response = client.post( + f"/api/admin/catalog-batch/batches/{manifest.batch_id}/resume", headers=admin_headers + ) + assert response.status_code == 409 + + +def test_cancel_on_a_queued_batch_cancels_its_files(client, admin_headers, batch_root): + manifest = batch_ingest.stage_batch([("a.csv", _csv()), ("b.csv", _csv())]) + + response = client.post( + f"/api/admin/catalog-batch/batches/{manifest.batch_id}/cancel", headers=admin_headers + ) + assert response.status_code == 200 + body = response.json() + assert body["status"] == "cancelled" + assert {f["status"] for f in body["files"]} == {"cancelled"} + + +def test_listing_shows_batches_from_disk(client, admin_headers, batch_root): + first = batch_ingest.stage_batch([("a.csv", _csv())]) + second = batch_ingest.stage_batch([("b.csv", _csv())]) + + response = client.get("/api/admin/catalog-batch/batches", headers=admin_headers) + assert response.status_code == 200 + ids = {b["batch_id"] for b in response.json()["batches"]} + assert {first.batch_id, second.batch_id} <= ids + + +def test_manifest_on_disk_is_valid_json_after_a_run(store, batch_root): + """The recovery path reads every manifest it finds; a torn write breaks boot.""" + manifest = batch_ingest.stage_batch([("a.csv", _csv())]) + batch_ingest.run_batch(manifest.batch_id) + + raw = json.loads(batch_ingest.manifest_path(manifest.batch_id).read_text(encoding="utf-8")) + assert raw["status"] == "done" + assert raw["files"][0]["status"] == "done" + # The temp file must not survive the write. + assert not (batch_ingest.batch_dir(manifest.batch_id) / "manifest.json.tmp").exists() diff --git a/tests/test_orchestration_defs.py b/tests/test_orchestration_defs.py index e93d49a..670c0b9 100644 --- a/tests/test_orchestration_defs.py +++ b/tests/test_orchestration_defs.py @@ -54,6 +54,11 @@ def test_every_expected_asset_exists(asset_graph): "training_dataset", "trained_models", "model_evaluation", + # batch (uploaded spreadsheets) + "batch_manifest", + "batch_parsed_files", + "batch_catalog_rows", + "batch_ingest_report", } @@ -70,6 +75,9 @@ def test_every_expected_asset_exists(asset_graph): ("nutrition_models", {"nutrition_data"}), ("trained_models", {"training_dataset"}), ("model_evaluation", {"trained_models"}), + ("batch_parsed_files", {"batch_manifest"}), + ("batch_catalog_rows", {"batch_manifest", "batch_parsed_files"}), + ("batch_ingest_report", {"batch_catalog_rows"}), ], ) def test_lineage_edges(asset_graph, asset_key, expected_parents): @@ -113,12 +121,32 @@ def test_cross_brand_assets_are_not_partitioned(asset_graph): assert not asset_graph.get(AssetKey(key)).is_partitioned, key +def test_batch_assets_are_not_partitioned(asset_graph): + """A batch is created by an upload, not enumerated in advance. + + Partitioning these would mean a dynamic partition per upload that is never + removed, so the partition set grows without bound and the UI fills with + dead keys. The batch id belongs in run config, which is what run config is + for. + """ + from dagster import AssetKey + + for key in ( + "batch_manifest", + "batch_parsed_files", + "batch_catalog_rows", + "batch_ingest_report", + ): + assert not asset_graph.get(AssetKey(key)).is_partitioned, key + + def test_all_four_jobs_resolve(defs): assert {job.name for job in defs.jobs} == { "catalog_ingestion_job", "embedding_refresh_job", "nutrition_enrichment_job", "ml_training_job", + "batch_ingestion_job", } @@ -138,7 +166,10 @@ def test_every_schedule_ships_stopped(defs): def test_sensor_ships_stopped_and_is_not_hot(defs): from dagster import DefaultSensorStatus - assert defs.sensors, "expected the seed-catalog sensor" + assert {s.name for s in defs.sensors} == { + "seed_catalog_sensor", + "batch_upload_sensor", + } for sensor in defs.sensors: assert sensor.default_status == DefaultSensorStatus.STOPPED, sensor.name assert sensor.minimum_interval_seconds >= 60, sensor.name @@ -148,7 +179,8 @@ def test_network_assets_retry_and_are_bounded(asset_graph): """Retries must exist on the flaky steps and must never be unbounded.""" from dagster import AssetKey - for key in ("raw_products", "enriched_products", "product_embeddings"): + for key in ("raw_products", "enriched_products", "product_embeddings", + "batch_catalog_rows"): policy = asset_graph.get(AssetKey(key)).assets_def.op.retry_policy assert policy is not None, key assert 0 < policy.max_retries <= 3, key