From 27d53fa9575949a00fdd2410bd87f0fc360891b5 Mon Sep 17 00:00:00 2001 From: sriram Date: Sat, 29 Aug 2026 11:21:11 +0530 Subject: [PATCH] upload-catalog-integration --- app/api/batch_common.py | 27 ++- app/api/routers/batch_catalog.py | 38 +++- app/api/routers/uploads.py | 2 +- app/core/batch_ingest.py | 87 +++++--- app/core/store_catalog_pipeline.py | 59 ++++++ docs/INGESTION_API.md | 89 +++++++- tests/test_drop_lifecycle.py | 320 +++++++++++++++++++++++++++++ tests/test_review_inbox.py | 10 +- 8 files changed, 584 insertions(+), 48 deletions(-) create mode 100644 tests/test_drop_lifecycle.py diff --git a/app/api/batch_common.py b/app/api/batch_common.py index c8c8346..fecaa57 100644 --- a/app/api/batch_common.py +++ b/app/api/batch_common.py @@ -187,6 +187,10 @@ class BatchFileOut(BaseModel): filename: str status: str detail: Optional[str] = None + # Set only on a file that was released out of the review inbox: the id of + # the run that took it, which is how a sender gets from the drop id they + # hold to the batch that carries their results. + released_to: Optional[str] = None stage_index: int = 0 stage_name: str = "" total_stages: int = pipeline.TOTAL_STAGES @@ -214,11 +218,26 @@ class BatchOut(BaseModel): files: List[BatchFileOut] -def to_out(manifest: batch_ingest.BatchManifest) -> BatchOut: +def to_out(manifest: batch_ingest.BatchManifest, *, slim: bool = False) -> BatchOut: + """Render a manifest for the API. + + `slim=True` drops the per-file `products` manifest, which is for LIST + responses. That list can run to thousands of rows per file + (store_catalog_pipeline.MAX_REPORTED_PRODUCTS), so twenty batches rendered + in full is a multi-megabyte response to a request that only wanted to know + what ran lately. The single-batch read keeps it - that is where a caller + goes to reconcile a specific run. + """ 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"]] + files = [] + for entry in body["files"]: + rendered = {key: entry[key] for key in BatchFileOut.model_fields if key in entry} + if slim and isinstance(rendered.get("result"), dict): + rendered["result"] = { + k: v for k, v in rendered["result"].items() if k != "products" + } + files.append(BatchFileOut(**rendered)) + body["files"] = files return BatchOut(**{k: v for k, v in body.items() if k in BatchOut.model_fields}) diff --git a/app/api/routers/batch_catalog.py b/app/api/routers/batch_catalog.py index 25283b6..b304e7f 100644 --- a/app/api/routers/batch_catalog.py +++ b/app/api/routers/batch_catalog.py @@ -189,7 +189,7 @@ def list_catalog_batches(limit: int = 20) -> dict: m for m in batch_job_store.recent(limit * 10) if m.status != batch_ingest.PENDING ][:limit] - return {"batches": [_to_out(m) for m in runs]} + return {"batches": [_to_out(m, slim=True) for m in runs]} @router.get("/batches/{batch_id}", dependencies=[Depends(require_admin)]) @@ -320,6 +320,27 @@ class InboxDismissOut(BaseModel): dismissed: int +def _retire(batch_id: str, indices: List[int], *, state: str, + released_to: Optional[str] = None) -> int: + """Retire files and refresh the cached manifest. Returns how many. + + The refresh is the part that is easy to forget and impossible to see: + `batch_job_store.get` prefers its in-memory copy over the disk, so a drop + retired without this keeps reporting `queued` to the sender polling it - + the exact question this whole record exists to answer. batch_ingest does + not know about the store (the store imports IT), so the refresh belongs + here, where both are already in hand. + """ + retired = batch_ingest.retire_files( + batch_id, indices, state=state, released_to=released_to, + ) + if retired: + updated = batch_ingest.read_manifest(batch_id) + if updated: + batch_job_store.put(updated) + return retired + + def _parse_file_ids(file_ids: List[str]) -> Dict[str, List[int]]: """Group "{batch_id}:{index}" into {batch_id: [index, ...]}. @@ -444,10 +465,15 @@ def start_batch_from_inbox(request: InboxStartRequest) -> BatchOut: ), ) - # Only now, once the bytes are safely staged under a new id. Removing them + # Only now, once the bytes are safely staged under a new id. Retiring them # first would lose the files outright if staging then failed. + # + # `released_to` is the whole point: the sender polls the drop id they were + # given, and this is how they learn which run took their sheet and where to + # follow it. Without it a release is indistinguishable from a deletion. for batch_id, indices in grouped.items(): - batch_ingest.remove_files(batch_id, indices) + _retire(batch_id, indices, state=batch_ingest.RELEASED, + released_to=manifest.batch_id) return _to_out(manifest) @@ -459,6 +485,10 @@ def dismiss_inbox_files(request: InboxSelection) -> InboxDismissOut: Deliberately irreversible and deliberately unceremonious: this is the disposal path for a drop nobody wants, and on an endpoint anyone can post to it is the control that keeps the volume from filling with rejected sheets. + + The file entry survives as a record reading `dismissed`, so the sender who + polls their drop id is told they were declined rather than left staring at a + 404. Only the bytes are gone. """ grouped = _parse_file_ids(request.file_ids) if not grouped: @@ -471,6 +501,6 @@ def dismiss_inbox_files(request: InboxSelection) -> InboxDismissOut: # files out of a batch that was mid-run. if not manifest or manifest.status != batch_ingest.PENDING: continue - dismissed += batch_ingest.remove_files(batch_id, indices) + dismissed += _retire(batch_id, indices, state=batch_ingest.DISMISSED) return InboxDismissOut(dismissed=dismissed) diff --git a/app/api/routers/uploads.py b/app/api/routers/uploads.py index 0b2c2e5..2d3de3f 100644 --- a/app/api/routers/uploads.py +++ b/app/api/routers/uploads.py @@ -239,7 +239,7 @@ def list_my_catalog_batches( mine = [ m for m in batch_job_store.recent(limit * 10) if _owns(m, principal) ][:limit] - return {"batches": [batch_common.to_out(m) for m in mine]} + return {"batches": [batch_common.to_out(m, slim=True) for m in mine]} @router.get("/catalog/{batch_id}") diff --git a/app/core/batch_ingest.py b/app/core/batch_ingest.py index 724a941..c7d87e9 100644 --- a/app/core/batch_ingest.py +++ b/app/core/batch_ingest.py @@ -68,6 +68,13 @@ DONE = "done" FAILED = "failed" CANCELLED = "cancelled" +# What became of a file that was sitting in the review inbox. Recorded on the +# file rather than deleted with it, because the sender polls the id THEY were +# given and has no other way to learn what happened: `released` carries the id +# of the run that took it, `dismissed` says an admin declined it. +RELEASED = "released" +DISMISSED = "dismissed" + # 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" @@ -79,7 +86,11 @@ PARTIAL = "partial" # `queued`, which is why `settle()` has to leave this status alone - see there. PENDING = "pending" -TERMINAL_BATCH_STATES = {DONE, FAILED, PARTIAL, CANCELLED} +# A drop whose files have all been released or declined. It is kept as a record +# for the sender to poll, holds no bytes, and ages out on normal retention. +RETIRED = "retired" + +TERMINAL_BATCH_STATES = {DONE, FAILED, PARTIAL, CANCELLED, RETIRED} # `..`, separators and drive letters all stripped. UploadFile.filename is # attacker-controlled in the general case, and it is used to build a path. @@ -105,7 +116,11 @@ class BatchFile: index: int filename: str # what the colleague called it, for display - stored_name: str # what it is called on disk + stored_name: str # what it is called on disk; "" once retired + # The run that took this file out of the review inbox. On the FILE and not + # the manifest because one drop can be released a few sheets at a time, into + # different runs, and the sender needs to know which of theirs went where. + released_to: Optional[str] = None size_bytes: int = 0 status: str = QUEUED detail: Optional[str] = None @@ -298,55 +313,59 @@ def list_manifests() -> List[BatchManifest]: return sorted(found, key=lambda m: m.created_at, reverse=True) -def remove_files(batch_id: str, indices) -> int: - """Drop these files from a staged batch, bytes and all. Returns the count. +def retire_files(batch_id: str, indices, *, state: str, + released_to: Optional[str] = None) -> int: + """Take these files out of the review inbox. Returns how many. - Used by both halves of the review inbox: dismissing files, and taking them - into a run (`from-inbox` copies the bytes into a fresh batch first, so the - original submission must not keep serving them a second time). + Both halves of triage land here: `state=RELEASED` with the id of the run + that took them, or `state=DISMISSED` when an admin declined them. + + THE BYTES GO; THE RECORD STAYS. Deleting the manifest as well - which is + what the first version did - left the sender polling the only id they were + ever given and getting a 404, unable to tell "accepted and running" from + "declined" from "lost". The file entry is therefore kept, marked, and for a + release annotated with the run to follow. Disk, which is the thing an open + endpoint can actually exhaust, is still reclaimed immediately. Surviving files KEEP their original `index`. The inbox addresses a file as - "{batch_id}:{index}", and the admin UI holds those ids in a selection set - across polls, so renumbering here would silently repoint a tick from the - file someone chose to whichever file slid into its place. - - A submission with nothing left is deleted outright rather than left as an - empty drop the inbox would render as a heading with no rows. + "{batch_id}:{index}" and the admin UI holds those ids in a selection set + across polls, so renumbering would silently repoint a tick at another file. """ manifest = read_manifest(batch_id) if not manifest: return 0 wanted = {int(i) for i in indices} - keep, dropped = [], 0 + retired = 0 for entry in manifest.files: - if entry.index not in wanted: - keep.append(entry) + if entry.index not in wanted or not entry.stored_name: continue - if entry.stored_name: - try: - (batch_dir(batch_id) / entry.stored_name).unlink(missing_ok=True) - except OSError as exc: - logger.warning( - "Could not delete %s from batch %s: %s", - entry.stored_name, batch_id, exc, - ) - dropped += 1 + try: + (batch_dir(batch_id) / entry.stored_name).unlink(missing_ok=True) + except OSError as exc: + logger.warning( + "Could not delete %s from drop %s: %s", + entry.stored_name, batch_id, exc, + ) + entry.stored_name = "" + entry.status = state + entry.released_to = released_to + entry.finished_at = time.time() + retired += 1 - if not dropped: + if not retired: return 0 - if not keep: - try: - shutil.rmtree(batch_dir(batch_id)) - except OSError as exc: - logger.warning("Could not remove empty batch %s: %s", batch_id, exc) - return dropped + # Nothing left for anyone to decide on, so it leaves the inbox. RETIRED is + # terminal, which is also what makes purge_expired willing to reclaim the + # directory once retention is up. + if not any(f.status == QUEUED and f.stored_name for f in manifest.files): + manifest.status = RETIRED + manifest.detail = "Every file here has been started or dismissed." - manifest.files = keep manifest.updated_at = time.time() write_manifest(manifest) - return dropped + return retired # --------------------------------------------------------------------------- diff --git a/app/core/store_catalog_pipeline.py b/app/core/store_catalog_pipeline.py index 4777285..54b4a72 100644 --- a/app/core/store_catalog_pipeline.py +++ b/app/core/store_catalog_pipeline.py @@ -161,6 +161,11 @@ class PipelineResult: recognised_columns: Dict[str, str] = field(default_factory=dict) unrecognised_columns: List[str] = field(default_factory=list) storage_error: Optional[str] = None + # What the sheet became, row by row: the identifiers a sender needs to + # reconcile their file against the catalog without matching on product name. + # See _record_product for why the unchanged rows are in here too. + products: List[Dict[str, Any]] = field(default_factory=list) + products_truncated: bool = False def as_dict(self) -> Dict[str, Any]: return { @@ -178,9 +183,18 @@ class PipelineResult: "recognised_columns": self.recognised_columns, "unrecognised_columns": self.unrecognised_columns, "storage_error": self.storage_error, + "products": self.products, + "products_truncated": self.products_truncated, } +# Ceiling on the per-file `products` manifest. Pack-size explosion multiplies +# the 2000-row file limit, and this list is serialised verbatim into +# manifest.json on the upload volume - the scarcest resource on this host. +# Past it the list stops growing and `products_truncated` says so, rather than +# a batch quietly writing tens of megabytes of JSON per file. +MAX_REPORTED_PRODUCTS = 5000 + ProgressFn = Callable[[int, str, int, int], None] """Called as (stage_index_1_based, stage_name, rows_done, rows_total).""" @@ -385,7 +399,20 @@ def stage_6_images(row: Dict[str, Any], *, enabled: bool = True) -> Dict[str, An # Stage 7 - SKU # --------------------------------------------------------------------------- def stage_7_sku(row: Dict[str, Any]) -> Dict[str, Any]: + """Resolve a SKU, but never over one the sheet already supplied. + + A sheet's own `sku` column wins outright and the resolver is not consulted - + the sender's identifier is a fact about their catalog, not a gap to fill. + + `sku_source` is stamped so the field is TOTAL rather than merely usually + populated: "sheet" here, "Internal" for a minted one, a marketplace label + when the lookup finds a real product id. Left null, a preserved SKU was + indistinguishable from an unknown one, and a caller could only tell by + diffing against the file they sent. + """ if not _blank(row.get("product_sku")): + if _blank(row.get("sku_source")): + row["sku_source"] = "sheet" return row try: resolved = resolve_product_sku( @@ -484,6 +511,32 @@ def _to_storage_row(row: Dict[str, Any]) -> Dict[str, Any]: } +def _record_product(result: PipelineResult, brand: str, row: Dict[str, Any], + disposition: str) -> None: + """Note the identity of one resolved catalog row. + + `image_id` is the useful field here and the one to join on: it is the key + `upsert_brand_products` deduplicates every product by. Matching on + product_name instead is what quietly creates duplicates, because a name that + differs by a character is a different id and therefore a different product. + + UNCHANGED ROWS ARE RECORDED TOO. They wrote nothing, but they are the normal + outcome of re-sending a sheet, and a caller who got an empty list back would + read a completely successful no-op as a total failure. + """ + if len(result.products) >= MAX_REPORTED_PRODUCTS: + result.products_truncated = True + return + result.products.append({ + "image_id": row.get("image_id"), + "brand": brand, + "product_name": row.get("product_name"), + "product_sku": row.get("product_sku"), + "sku_source": row.get("sku_source"), + "disposition": disposition, + }) + + def _merge_with_existing(new: Dict[str, Any], existing: Dict[str, Any]) -> Tuple[Dict[str, Any], bool]: """Keep the stored row, filling only the columns it left empty. @@ -529,13 +582,19 @@ def stage_11_store(rows: List[Dict[str, Any]], brand: str, result: PipelineResul if prior is None: to_write.append(row) result.inserted += 1 + _record_product(result, brand, row, "inserted") continue merged, changed = _merge_with_existing(row, prior) if changed: to_write.append(merged) result.backfilled += 1 + # The MERGED row: a backfill keeps the stored product_sku rather + # than the one this run minted, so reporting `row` would hand back + # an identifier that is not the one in the catalog. + _record_product(result, brand, merged, "backfilled") else: result.skipped_existing += 1 + _record_product(result, brand, prior, "unchanged") if not to_write: return diff --git a/docs/INGESTION_API.md b/docs/INGESTION_API.md index 3e6934e..dd3375b 100644 --- a/docs/INGESTION_API.md +++ b/docs/INGESTION_API.md @@ -158,10 +158,26 @@ run; it is waiting behind a decision, not behind the worker. `batch_id` is a `uuid4` handed only to whoever sent the drop, so holding it is the proof of having sent it. An id nobody issued is a `404`. -Expect `pending` until an admin acts. After that the batch you sent stops existing -under that id — the files move into a run with an id of its own, and your drop -reports `404`. **A `404` after `pending` means your files were accepted and started,** -not that they were lost. Ask the admin for the run id if you need to follow it. +**Your drop id stays valid for the whole lifecycle.** Poll it and read the per-file +`status`: + +| File `status` | `released_to` | Meaning | +| --- | --- | --- | +| `queued` | `null` | Still in the inbox; nobody has looked yet | +| `released` | the run id | Accepted and started. **Follow that id for progress and results** | +| `dismissed` | `null` | An admin declined this file | + +```jsonc +// GET /api/uploads/catalog/{drop_id}, after an admin pressed Start +{ "status": "retired", + "files": [ { "filename": "catalog.csv", + "status": "released", + "released_to": "8dcef8a2ad94..." } ] } +``` + +Then `GET /api/uploads/catalog/{released_to}` — also with no credential — for the run +itself. A run an admin assembled from several drops lists every file in it, so you may +see filenames batched alongside your own. `GET /api/uploads/catalog?limit=20` — your recent submissions, newest first, under `{"batches": [...]}`. This one **does** need a credential: letting an anonymous caller @@ -173,6 +189,7 @@ hold. | Status | Meaning | | --- | --- | | `pending` | In the review inbox. Nothing has run. | +| `retired` | Every file in this drop has been released or dismissed. Read the per-file `status` | | `queued` | Approved and waiting for the worker | | `running` | In the pipeline now | | `done` | Every file completed | @@ -181,6 +198,68 @@ hold. | `interrupted` | A restart cut a run short. An admin resumes it; it never auto-restarts | | `cancelled` | Stopped by an admin before the remaining files began | +### What a finished run tells you + +Each file in a completed run carries a `result`. Alongside the counts it lists **every +row that resolved**, with the identifiers needed to reconcile the sheet against the +catalogue: + +```jsonc +"result": { + "rows_total": 2, "inserted": 1, "backfilled": 0, "skipped_existing": 1, + "rejected": 0, "brands": ["amul"], + "products": [ + { "image_id": "amul_amul_butter_100g", "brand": "amul", + "product_name": "Amul Butter 100g", + "product_sku": "ACME-BUT-100", "sku_source": "sheet", + "disposition": "inserted" }, + { "image_id": "amul_amul_ghee_1l", "brand": "amul", + "product_name": "Amul Ghee 1L", + "product_sku": "AMUL-GHE-1-001", "sku_source": "Internal", + "disposition": "unchanged" } + ], + "products_truncated": false +} +``` + +| `disposition` | What happened | +| --- | --- | +| `inserted` | New product, created by this run | +| `backfilled` | Existed already; this sheet filled columns it had left empty | +| `unchanged` | Existed already and was complete. Nothing was written | + +**Join on `image_id`.** It is the key the catalogue deduplicates every product by. Do +not match on `product_name` — a name differing by one character is a different +`image_id` and therefore a different product, so name matching silently creates +duplicates instead of updating. + +`unchanged` rows are listed too. Re-sending a sheet is the normal case and writes +nothing; without them you would get an empty list back for a completely successful +upload. `products_truncated` is `true` when a file resolved more than 5,000 rows +(`MAX_REPORTED_PRODUCTS`) — the counts stay accurate, only the listing is cut. + +`products` is returned on the **single-batch read** only. The list endpoints omit it, +since twenty runs of several thousand rows each is not a list payload. + +### Which SKU wins + +| Your sheet's `sku` column | Result | `sku_source` | +| --- | --- | --- | +| Has a value | **Preserved verbatim.** The resolver is never called | `sheet` | +| Blank | A deterministic internal SKU is minted, e.g. `AMUL-GHE-1-001` | `Internal` | +| Blank, and marketplace lookup enabled | A real marketplace product id | the marketplace name | + +The third row does not occur in production: `ENABLE_SKU_WEB_LOOKUP` is `false` there, +so a blank cell always yields an `Internal` SKU. + +If your sheet carries its own `sku_source` column, that is preserved too and not +overwritten with `sheet`. + +**A product's SKU is stable once stored.** Re-sending the same product does not +renumber it — a backfill keeps the stored value rather than the one that run minted. +But that stability is per `image_id`: change the product name and you get a new +`image_id`, a new product, and a new SKU. Another reason to join on `image_id`. + ### The eleven stages 1. Brand Resolution & FSSAI Licence Mapping @@ -344,6 +423,8 @@ Sample sheets with the correct headers: `GET /api/upload/template/stores`, `/ana | --- | --- | | `POST /api/uploads/catalog` needed an `X-API-Key` | No credential. Send the file; drop the header | | Files ran on arrival | Files wait for review. Expect `pending`, not `queued` | +| The drop id 404'd once an admin started it | It stays valid. The file reads `released` and carries `released_to` | +| The result was counts only | It also lists `products` with `image_id` / `product_sku` / `disposition` | | `?use_llm` / `?fetch_images` on the drop | Ignored. The admin chooses at Start | | `429` meant the worker queue was full | `429` now means the review inbox is full | | `200` with `rows_imported: 0` | `422` with per-row reasons. Handle as a client error, not a server one | diff --git a/tests/test_drop_lifecycle.py b/tests/test_drop_lifecycle.py new file mode 100644 index 0000000..fae5664 --- /dev/null +++ b/tests/test_drop_lifecycle.py @@ -0,0 +1,320 @@ +"""What a sender can learn about their own drop, and what the result tells them. + +Three integrator questions produced this file, and each is pinned here rather +than only answered in prose: + +1. How does a held batch get released, and how does the sender find out? The + drop id they were handed has to stay meaningful for the whole lifecycle - + the first version of the inbox deleted the drop on release, which left them + polling a 404 with no way to tell "running" from "declined" from "lost". +2. Does the result carry the identifiers needed to reconcile a sheet against + the catalog? `image_id` is the join key; matching on product_name is what + silently creates duplicates. +3. Is a SKU the sheet supplied preserved, or replaced by a minted one? The + sheet wins, and `sku_source` has to say so without the caller diffing + against the file they sent. +""" +from __future__ import annotations + +import io + +import pytest + +from app.api.batch_job_store import batch_job_store +from app.core import batch_ingest +from app.core import store_catalog_pipeline as pipeline + +UPLOAD = "/api/uploads/catalog" +INBOX = "/api/admin/catalog-batch/inbox" +FROM_INBOX = "/api/admin/catalog-batch/from-inbox" +DISMISS = "/api/admin/catalog-batch/inbox/dismiss" + +HEADERS = ["Product Name", "Category", "Brand"] +ROWS = [["Amul Butter 100g", "Butter", "Amul"]] + + +def _csv(headers=HEADERS, rows=ROWS) -> bytes: + return ("\n".join([",".join(headers)] + [",".join(r) for r in rows])).encode() + + +def _files(*pairs): + return [("files", (n, io.BytesIO(c), "application/octet-stream")) + for n, c in pairs] + + +@pytest.fixture(autouse=True) +def _isolate_sku_counter(tmp_path, monkeypatch): + 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): + root = tmp_path / "batch_uploads" + monkeypatch.setattr(batch_ingest, "BATCH_UPLOAD_DIR", root) + return root + + +@pytest.fixture(autouse=True) +def submitted(monkeypatch): + from app.core import batch_worker + + seen: list = [] + monkeypatch.setattr(batch_worker, "submit", seen.append) + return seen + + +@pytest.fixture(autouse=True) +def _clean_job_store(): + batch_job_store._batches.clear() + batch_job_store._cancelled.clear() + yield + batch_job_store._batches.clear() + batch_job_store._cancelled.clear() + + +@pytest.fixture +def store(monkeypatch): + """Stand in for the catalog table, keyed the way the real one is. + + Keyed on image_id because that is what upsert_brand_products deduplicates + on - which is the whole claim question 2 rests on. + """ + table: dict = {} + + def _get(brand): + return list(table.get((brand or "").lower(), {}).values()) + + def _upsert(brand, products, cleanup=False): + rows = table.setdefault((brand or "").lower(), {}) + for product in products: + rows[product["image_id"]] = dict(product) + return len(rows) + + monkeypatch.setattr(pipeline, "get_products_by_brand", _get) + monkeypatch.setattr(pipeline, "upsert_brand_products", _upsert) + monkeypatch.setattr(pipeline, "USE_EMBEDDINGS", False) + return table + + +def _drop(client, sender="priya", *names): + response = client.post( + UPLOAD, + files=_files(*[(n, _csv()) for n in (names or ("a.csv",))]), + data={"sender": sender}, + ) + assert response.status_code == 202, response.text + return response.json()["batch_id"] + + +# --------------------------------------------------------------------------- +# 1. Releasing a held drop, as the sender sees it +# --------------------------------------------------------------------------- +def test_a_released_drop_points_the_sender_at_the_run(client, admin_headers): + """THE answer to "how does a held batch get released". + + The sender polls the one id they were given. After an admin starts their + file, that id must still resolve AND carry the id of the run that took it - + otherwise there is no path from the drop to the result. + """ + drop_id = _drop(client, "priya", "catalog.csv") + + started = client.post(FROM_INBOX, json={"file_ids": [f"{drop_id}:0"]}, + headers=admin_headers) + run_id = started.json()["batch_id"] + + polled = client.get(f"{UPLOAD}/{drop_id}") + assert polled.status_code == 200, polled.text + entry = polled.json()["files"][0] + assert entry["status"] == "released" + assert entry["released_to"] == run_id + + # And that id is readable with no credential either, so the sender can + # follow it without one being issued to them. + assert client.get(f"{UPLOAD}/{run_id}").status_code == 200 + + +def test_a_dismissed_drop_says_so_rather_than_vanishing(client, admin_headers): + drop_id = _drop(client, "priya") + + client.post(DISMISS, json={"file_ids": [f"{drop_id}:0"]}, headers=admin_headers) + + entry = client.get(f"{UPLOAD}/{drop_id}").json()["files"][0] + assert entry["status"] == "dismissed" + assert entry["released_to"] is None + + +def test_a_half_released_drop_keeps_the_rest_waiting(client, admin_headers): + drop_id = _drop(client, "priya", "one.csv", "two.csv") + + client.post(FROM_INBOX, json={"file_ids": [f"{drop_id}:0"]}, headers=admin_headers) + + files = client.get(f"{UPLOAD}/{drop_id}").json()["files"] + by_name = {f["filename"]: f for f in files} + assert by_name["one.csv"]["status"] == "released" + assert by_name["two.csv"]["status"] == "queued" + # Still in the inbox, so the drop as a whole has not been settled. + assert client.get(INBOX, headers=admin_headers).json()["pending_count"] == 1 + + +def test_a_retired_drop_leaves_the_inbox_and_frees_its_disk(client, admin_headers, + batch_root): + drop_id = _drop(client, "priya") + + client.post(FROM_INBOX, json={"file_ids": [f"{drop_id}:0"]}, headers=admin_headers) + + assert client.get(INBOX, headers=admin_headers).json()["pending_count"] == 0 + # The record survives; the bytes do not. Keeping the manifest is what makes + # the sender's poll work - keeping the sheet would defeat the disk control. + assert not list((batch_root / drop_id).glob("*.csv")) + assert batch_ingest.read_manifest(drop_id).status == "retired" + + +def test_a_retired_drop_is_purged_on_retention(client, admin_headers): + """The tombstone is small, but it is not permanent.""" + drop_id = _drop(client, "priya") + client.post(DISMISS, json={"file_ids": [f"{drop_id}:0"]}, headers=admin_headers) + + import time + purged = batch_ingest.purge_expired( + now=time.time() + (batch_ingest.BATCH_RETENTION_DAYS + 1) * 86400 + ) + + assert drop_id in purged + + +def test_a_released_drop_is_not_resurrected_by_a_restart(client, admin_headers): + """`retired` is terminal, so the interrupt sweep must leave it alone.""" + drop_id = _drop(client, "priya") + client.post(FROM_INBOX, json={"file_ids": [f"{drop_id}:0"]}, headers=admin_headers) + + assert drop_id not in batch_ingest.scan_interrupted() + assert batch_ingest.read_manifest(drop_id).status == "retired" + + +# --------------------------------------------------------------------------- +# 2. Identifiers in the result +# --------------------------------------------------------------------------- +def test_the_result_names_every_product_it_created(client, admin_headers, store): + drop_id = _drop(client, "priya", "catalog.csv") + run_id = client.post(FROM_INBOX, json={"file_ids": [f"{drop_id}:0"]}, + headers=admin_headers).json()["batch_id"] + + manifest = batch_ingest.run_batch(run_id) + + products = manifest.files[0].result["products"] + assert len(products) == 1 + product = products[0] + assert product["brand"] == "amul" + assert product["image_id"] + assert product["product_name"] + assert product["product_sku"] + assert product["disposition"] == "inserted" + + +def test_resending_a_sheet_reports_unchanged_not_an_empty_list(client, admin_headers, + store): + """The case that decides whether a caller can reconcile at all. + + An unchanged row writes nothing. Reporting only what was written would hand + back an empty list for a completely successful re-send, which reads as total + failure - and the ids are exactly what the caller needs to confirm the rows + are already there. + """ + def _run_one(): + drop_id = _drop(client, "priya") + run_id = client.post(FROM_INBOX, json={"file_ids": [f"{drop_id}:0"]}, + headers=admin_headers).json()["batch_id"] + return batch_ingest.run_batch(run_id).files[0].result["products"] + + first = _run_one() + second = _run_one() + + assert [p["disposition"] for p in first] == ["inserted"] + assert [p["disposition"] for p in second] == ["unchanged"] + # The join key is stable across both, which is what makes it a join key. + assert first[0]["image_id"] == second[0]["image_id"] + assert first[0]["product_sku"] == second[0]["product_sku"] + + +def test_the_product_manifest_is_absent_from_list_responses(client, admin_headers, + store): + """Thousands of rows per file is a single-batch payload, not a list one.""" + drop_id = _drop(client, "priya") + run_id = client.post(FROM_INBOX, json={"file_ids": [f"{drop_id}:0"]}, + headers=admin_headers).json()["batch_id"] + # on_change is what batch_worker passes in production; without it the job + # store keeps the pre-run copy and the list below would read a stale result. + batch_ingest.run_batch(run_id, on_change=batch_job_store.put) + + listed = client.get("/api/admin/catalog-batch/batches", + headers=admin_headers).json()["batches"] + listed_result = [b for b in listed if b["batch_id"] == run_id][0]["files"][0]["result"] + assert "products" not in listed_result + # The counts a list view actually needs are still there. + assert listed_result["inserted"] == 1 + + single = client.get(f"{UPLOAD}/{run_id}").json()["files"][0]["result"] + assert len(single["products"]) == 1 + + +def test_the_product_manifest_is_capped(client, admin_headers, store, monkeypatch): + monkeypatch.setattr(pipeline, "MAX_REPORTED_PRODUCTS", 1) + rows = [["Amul Butter 100g", "Butter", "Amul"], ["Amul Ghee 1L", "Ghee", "Amul"]] + response = client.post(UPLOAD, files=_files(("a.csv", _csv(rows=rows)))) + drop_id = response.json()["batch_id"] + run_id = client.post(FROM_INBOX, json={"file_ids": [f"{drop_id}:0"]}, + headers=admin_headers).json()["batch_id"] + + result = batch_ingest.run_batch(run_id).files[0].result + + assert len(result["products"]) == 1 + assert result["products_truncated"] is True + # Truncating the REPORT must not truncate the work. + assert result["inserted"] == 2 + + +# --------------------------------------------------------------------------- +# 3. Whose SKU wins +# --------------------------------------------------------------------------- +def test_a_sheet_sku_is_preserved_and_labelled(): + row = {"product_sku": "MY-SKU-1", "brand": "amul", "product_name": "Butter"} + + out = pipeline.stage_7_sku(dict(row)) + + assert out["product_sku"] == "MY-SKU-1" + assert out["sku_source"] == "sheet" + + +def test_a_sheets_own_sku_source_is_not_overwritten(): + row = {"product_sku": "MY-SKU-1", "sku_source": "erp-export"} + + assert pipeline.stage_7_sku(dict(row))["sku_source"] == "erp-export" + + +def test_a_blank_sku_is_minted_internally(): + """ENABLE_SKU_WEB_LOOKUP is false by default and unset in production, so + the marketplace branch is dead there and this is the only other outcome.""" + row = {"product_sku": "", "brand": "Amul", "product_name": "Butter", "size": "100g"} + + out = pipeline.stage_7_sku(dict(row)) + + assert out["product_sku"] + assert out["sku_source"] == "Internal" + + +def test_a_reupload_does_not_renumber_an_existing_products_sku(client, admin_headers, + store): + """The sequence counter hands out a fresh number every run, so the stability + of a minted SKU rests entirely on _merge_with_existing keeping the stored + value. Pinned here because nothing else would notice it changing.""" + def _run_one(): + drop_id = _drop(client, "priya") + run_id = client.post(FROM_INBOX, json={"file_ids": [f"{drop_id}:0"]}, + headers=admin_headers).json()["batch_id"] + return batch_ingest.run_batch(run_id).files[0].result["products"][0] + + first = _run_one() + second = _run_one() + + assert first["product_sku"] == second["product_sku"] diff --git a/tests/test_review_inbox.py b/tests/test_review_inbox.py index 67191ca..a1d5ebc 100644 --- a/tests/test_review_inbox.py +++ b/tests/test_review_inbox.py @@ -283,9 +283,17 @@ def test_dismissing_deletes_the_file_and_its_bytes(client, admin_headers, batch_ assert result.status_code == 200, result.text assert result.json() == {"dismissed": 1} assert client.get(INBOX, headers=admin_headers).json()["pending_count"] == 0 + # The disposal path on an endpoint anyone can post to. If the bytes survived # a dismiss, the volume would fill with sheets that were already refused. - assert not (batch_root / batch_id).exists() + assert not list((batch_root / batch_id).glob("*.csv")) + + # The RECORD survives, though, and that is deliberate: the sender polls the + # only id they were ever given, and must be told they were declined rather + # than left unable to tell a refusal from a lost file. + polled = client.get(f"{UPLOAD}/{batch_id}") + assert polled.status_code == 200, polled.text + assert polled.json()["files"][0]["status"] == "dismissed" def test_dismissing_cannot_touch_a_running_batch(client, admin_headers, monkeypatch):