From 3df2dc5991705541e9187ff27c3b66cb10d08cd4 Mon Sep 17 00:00:00 2001 From: sriram Date: Mon, 31 Aug 2026 14:59:14 +0530 Subject: [PATCH] Backend upload-automation file --- .env.example | 38 ++++++ app/api/routers/uploads.py | 176 +++++++++++++++++--------- app/infrastructure/settings.py | 38 +++++- docs/INGESTION_API.md | 150 ++++++++++++++++------ tests/conftest.py | 19 +++ tests/test_drop_lifecycle.py | 8 ++ tests/test_review_inbox.py | 10 ++ tests/test_uploads_api.py | 10 ++ tests/test_uploads_autorun.py | 223 +++++++++++++++++++++++++++++++++ 9 files changed, 569 insertions(+), 103 deletions(-) create mode 100644 tests/test_uploads_autorun.py diff --git a/.env.example b/.env.example index 551aff3..33a22c4 100644 --- a/.env.example +++ b/.env.example @@ -179,8 +179,46 @@ BATCH_RETENTION_DAYS=7 # waits for someone to press Resume. Auto-resuming means a container stuck in a # restart loop re-runs the heaviest work in the app on every boot, which is how # a slow start becomes an unrecoverable spiral. +# +# Worth knowing alongside UPLOAD_AUTORUN below: with uploads running unattended, +# a redeploy in the middle of one leaves that batch "interrupted" and waiting +# for a human. It is the one place manual intervention comes back. BATCH_AUTO_RESUME=false +# --------------------------------------------------------------------------- +# Unattended ingestion +# --------------------------------------------------------------------------- +# Whether POST /api/uploads/catalog runs the pipeline on arrival (true) or parks +# the files in the admin review inbox for someone to start by hand (false). +# +# READ THIS BEFORE CHANGING IT. That endpoint takes NO credential - it was +# opened on purpose so colleagues could send spreadsheets without one being +# issued to them. With autorun on, "anyone who can reach this host" and "anyone +# who can write to the live catalogue" are the same set of people, and an ingest +# is an upsert with no undo. +# +# What still bounds it is throughput, not identity: the per-request ceilings +# above, and BATCH_QUEUE_MAX behind a single worker. A sender can occupy the +# ingestion worker; they cannot multiply it. +# +# Set false and the review inbox comes back with no code change - the INBOX_* +# ceilings below apply only on that path. +UPLOAD_AUTORUN=true + +# How an auto-started run behaves. Deliberately NOT accepted from the request: +# the sender is anonymous, and letting an anonymous caller switch on the +# expensive outbound stages is the one thing this endpoint must not allow. +# +# Images on, because a product landing without one is the failure this endpoint +# exists to avoid. Stage 6 is the slowest stage and reaches the network, but +# only one batch runs at a time, so nothing else competes with it. +# +# LLM off, because use_llm gates only description generation in stage 2, and +# production runs USE_OLLAMA=false - turning it on there buys nothing and costs +# a connection timeout per row. +UPLOAD_AUTORUN_FETCH_IMAGES=true +UPLOAD_AUTORUN_USE_LLM=false + USE_OLLAMA=true OLLAMA_BASE_URL=http://localhost:11434 OLLAMA_MODEL_NAME=qwen2.5:1.5b diff --git a/app/api/routers/uploads.py b/app/api/routers/uploads.py index 2d3de3f..12f93e4 100644 --- a/app/api/routers/uploads.py +++ b/app/api/routers/uploads.py @@ -4,12 +4,10 @@ GET /api/uploads/catalog - the batches this caller has sent GET /api/uploads/catalog/{batch_id} - progress and result of one of them -This is deliberately its own router, with its own prefix and its own guard, so -that the difference between it and everything else in the app is visible in one -screen rather than inferred from a decorator halfway down a 400-line admin -module. Everything else that touches catalog data is `require_admin`; this is -`require_permission("upload_catalog")`, and that single line is the whole -security boundary of the feature. +This is deliberately its own router, with its own prefix, so that the difference +between it and everything else in the app is visible in one screen rather than +inferred from a decorator halfway down a 400-line admin module. Everything else +that touches catalog data is `require_admin`; the POST here has no guard at all. WHAT HAPPENS WHEN A FILE ARRIVES -------------------------------- @@ -23,19 +21,27 @@ The response is a `batch_id`. Ingestion is far too slow to finish inside a request - it is thousands of rows through eleven stages - so the caller polls GET /api/uploads/catalog/{batch_id} until `status` leaves `queued`/`running`. -THIS CREDENTIAL NOW COSTS CPU, AND THAT IS THE POINT ----------------------------------------------------- -An earlier version of this endpoint parked files in a review inbox and started -nothing, so that a leaked key could cost only disk. That is not the product: -an API user sends a file in order for it to be ingested, and a queue that needs -an admin to press a button is not an API. +AN UNAUTHENTICATED POST NOW STARTS REAL WORK +-------------------------------------------- +Be clear-eyed about what that means. This endpoint takes no credential, and +`UPLOAD_AUTORUN` (default true) runs the pipeline the moment a file lands. So +"anyone who can reach this host" and "anyone who can write to the live catalog" +are the same set of people, and an ingest is an upsert with no undo. -So the bound is no longer "this role cannot start work" but "all work, from -every source, goes through one worker". `batch_worker` runs a single batch at a -time behind a queue of `BATCH_QUEUE_MAX`; past that this endpoint answers 429. -An uploader key can therefore occupy the ingestion worker, which is what it is -for - it cannot multiply it, which is what matters on a one-vCPU host that is -also serving the API and its healthcheck. +That was chosen deliberately, over the alternative of issuing the sender an +`uploader` API key and auto-running only credentialed requests. The requirement +was uploads that run without manual intervention, and a review queue that needs +an admin to press a button is not that. + +What still bounds it is throughput, not identity: the per-request ceilings in +`_limits()`, and `batch_worker` running a single batch at a time behind a queue +of `BATCH_QUEUE_MAX`, past which this endpoint answers 429. A sender can occupy +the ingestion worker - that is what it is for - but cannot multiply it, which is +what matters on a one-vCPU host also serving the API and its healthcheck. + +Setting `UPLOAD_AUTORUN=false` restores the review inbox, where files wait for +an admin and the cost of an unwanted drop is disk rather than products. Both +paths are live and both are tested; see `stage_pending` in batch_common. WHAT THIS ENDPOINT STILL CANNOT DO ---------------------------------- @@ -62,6 +68,9 @@ from app.infrastructure.settings import ( BATCH_MAX_TOTAL_ROWS, INBOX_MAX_PENDING_BYTES, INBOX_MAX_PENDING_FILES, + UPLOAD_AUTORUN, + UPLOAD_AUTORUN_FETCH_IMAGES, + UPLOAD_AUTORUN_USE_LLM, ) logger = logging.getLogger(__name__) @@ -145,29 +154,35 @@ async def ingest_catalog_files( sender: Optional[str] = Form(None), principal: Optional[Principal] = Depends(get_optional_principal), ) -> CatalogUploadOut: - """Accept spreadsheets into the review inbox. No credential required. + """Accept spreadsheets and start the pipeline over them. No credential. Returns 202 and a `batch_id` to poll. A drop where some files parse and some - do not is a partial success, not a failure: the good ones are stored and the - bad ones come back in `files` as `status: "failed"` with the reason, so the - caller knows exactly which sheet to fix and resend. + do not is a partial success, not a failure: the good ones are accepted and + the bad ones come back in `files` as `status: "failed"` with the reason, so + the caller knows exactly which sheet to fix and resend. - NOTHING SENT HERE RUNS ON ARRIVAL. - The files are staged and the batch is left `pending`. An admin sees it in the - review inbox, ticks the sheets they want and presses Start; only then does - anything reach the pipeline or the catalog. That gate is what makes an - endpoint anybody can post to acceptable: the cost of an unwanted drop is - disk until someone declines it, not products in the live catalog. + WHAT HAPPENS ON ARRIVAL depends on one setting, `UPLOAD_AUTORUN`: - `use_llm` and `fetch_images` are NOT accepted here, though they used to be. - They decide how a run behaves, and the person who decides that is now the - admin pressing Start - not the sender. Leaving them on this endpoint would - let an anonymous caller commit the host to Playwright image search. + true (default) the batch is queued and the 11 stages run immediately. The + id returned IS the run id - poll it and watch `stages[]`. + false the batch is left `pending` in the admin review inbox and + nothing runs until someone presses Start. The id returned + is a DROP id; the run gets a different one, reachable + through the file's `released_to`. + + Clients should not care which is configured: both return 202 and an id that + `GET /api/uploads/catalog/{id}` understands. Only the number of hops differs. + + `use_llm` and `fetch_images` are still NOT accepted from the request, and + that has not changed with autorun - if anything it matters more. The caller + is anonymous, and letting an anonymous caller switch on the expensive + outbound stages is the one thing this endpoint must not allow. They come + from `UPLOAD_AUTORUN_FETCH_IMAGES` / `UPLOAD_AUTORUN_USE_LLM` instead. `sender` is a free-text label, not identity - it is whatever the caller - typed. It exists because the inbox groups drops by who sent them, and three - colleagues all showing as "anonymous" is an inbox nobody can triage. A real - credential, if one is presented, wins over it. + typed. It exists so a run can be attributed to a person, and three + colleagues all showing as "anonymous" is a batch list nobody can triage. A + real credential, if one is presented, wins over it. """ limits = _limits() read = await batch_common.read_uploads(files, limits) @@ -180,10 +195,15 @@ async def ingest_catalog_files( detail=f"None of the uploaded files could be ingested. {detail}", ) - _inbox_capacity_or_429( - incoming_files=len(valid), - incoming_bytes=sum(len(contents) for _n, contents, _r in valid), - ) + # Only meaningful on the review-inbox path. Under autorun nothing ever + # awaits review, so the count it guards is permanently zero and the check + # could never fire - and a guard that cannot guard anything reads, to the + # next person, like protection that is actually there. + if not UPLOAD_AUTORUN: + _inbox_capacity_or_429( + incoming_files=len(valid), + incoming_bytes=sum(len(contents) for _n, contents, _r in valid), + ) # A presented credential still names the sender - `get_optional_principal` # returns None only when NO credential was sent, and still raises on one @@ -194,30 +214,68 @@ async def ingest_catalog_files( else ((sender or "").strip()[:60] or "anonymous") ) - manifest = batch_common.stage_pending( - valid, - invalid, - submitted_by=submitted_by, - ) - - logger.info( - "Catalog drop %s received for review: %d file(s), %d rejected, from %s%s", - manifest.batch_id, len(valid), len(invalid), submitted_by, - "" if principal else " (no credential)", - ) + if UPLOAD_AUTORUN: + manifest, started = batch_common.stage_and_queue( + valid, + invalid, + use_llm=UPLOAD_AUTORUN_USE_LLM, + fetch_images=UPLOAD_AUTORUN_FETCH_IMAGES, + submitted_by=submitted_by, + ) + if not started: + # Staged and durable, but not running, and this caller has no Resume + # button - that lives on the admin router. So the honest instruction + # is to send it again shortly. The id is named so an admin can find + # and resume THIS batch instead, if the caller reports it. + raise HTTPException( + status_code=429, + detail=( + f"Too many batches are already queued. Batch " + f"{manifest.batch_id} has been saved but not started; retry " + f"this upload shortly." + ), + ) + logger.info( + "Catalog batch %s queued: %d file(s), %d rejected, from %s%s", + manifest.batch_id, len(valid), len(invalid), submitted_by, + "" if principal else " (no credential)", + ) + else: + manifest = batch_common.stage_pending( + valid, + invalid, + submitted_by=submitted_by, + ) + logger.info( + "Catalog drop %s received for review: %d file(s), %d rejected, from %s%s", + manifest.batch_id, len(valid), len(invalid), submitted_by, + "" if principal else " (no credential)", + ) body = batch_common.to_out(manifest).model_dump() - message = ( - f"{len(valid)} file(s) received and waiting for review. Nothing runs " - f"until an admin starts them. " - f"Poll GET /api/uploads/catalog/{manifest.batch_id} for status." - ) - if invalid: + if UPLOAD_AUTORUN: message = ( - f"{len(valid)} file(s) received and waiting for review. " - f"{len(invalid)} could not be read - see 'files' for the reason on " - f"each, and resend those." + f"{len(valid)} file(s) accepted and queued for ingestion. " + f"Poll GET /api/uploads/catalog/{manifest.batch_id} for progress." ) + if invalid: + message = ( + f"{len(valid)} file(s) accepted and queued for ingestion. " + f"{len(invalid)} could not be read - see 'files' for the reason " + f"on each, and resend those." + ) + else: + message = ( + f"{len(valid)} file(s) received and waiting for review. Nothing runs " + f"until an admin starts them. " + f"Poll GET /api/uploads/catalog/{manifest.batch_id} for status." + ) + if invalid: + message = ( + f"{len(valid)} file(s) received and waiting for review. " + f"{len(invalid)} could not be read - see 'files' for the reason on " + f"each, and resend those." + ) return CatalogUploadOut(**body, message=message) diff --git a/app/infrastructure/settings.py b/app/infrastructure/settings.py index 3386092..17e3fbd 100644 --- a/app/infrastructure/settings.py +++ b/app/infrastructure/settings.py @@ -163,10 +163,42 @@ BATCH_QUEUE_MAX = int(os.getenv("BATCH_QUEUE_MAX", "4")) # resource it has. BATCH_RETENTION_DAYS = int(os.getenv("BATCH_RETENTION_DAYS", "7")) +# --- Unattended ingestion -------------------------------------------------- +# Whether POST /api/uploads/catalog runs the pipeline on arrival, or parks the +# files in the admin review inbox for someone to start by hand. +# +# READ THIS BEFORE CHANGING IT. That endpoint takes NO credential - it was +# opened deliberately so colleagues could send spreadsheets without one being +# issued to them. With autorun on, "anyone who can reach this host" and "anyone +# who can write to the live catalogue" become the same set of people, and an +# ingest is an upsert with no undo. That trade was made knowingly: the ask was +# for uploads to run without manual intervention, and a review queue that needs +# an admin to press a button is not that. +# +# What still bounds it: the per-request ceilings above (20 files / 50MB / 20k +# rows), and BATCH_QUEUE_MAX behind a single worker thread - so a sender can +# occupy the ingestion worker but cannot multiply it. Those cap throughput, not +# who. If that stops being an acceptable trade, set this to false and the review +# inbox comes back with no code change; everything it needs is still here. +UPLOAD_AUTORUN = _bool("UPLOAD_AUTORUN", "true") + +# How an auto-started run behaves. Not accepted from the request: the sender is +# anonymous, and letting an anonymous caller turn on the expensive stages is the +# one thing the open endpoint must not allow. +# +# Images ON, because a product landing without one is the failure this endpoint +# exists to avoid - stage 6 is the slowest stage and reaches the network, but +# only one batch runs at a time so nothing else is competing with it. +# +# LLM OFF, because `use_llm` gates only description generation in +# stage_2_row_intake, and production runs USE_OLLAMA=false: turning it on there +# buys nothing and costs a connection timeout per row. +UPLOAD_AUTORUN_FETCH_IMAGES = _bool("UPLOAD_AUTORUN_FETCH_IMAGES", "true") +UPLOAD_AUTORUN_USE_LLM = _bool("UPLOAD_AUTORUN_USE_LLM", "false") + # --- Review inbox ---------------------------------------------------------- -# POST /api/uploads/catalog accepts files with NO credential, so that colleagues -# can send spreadsheets without one being issued to them. Nothing it accepts is -# queued - files wait in the admin review inbox - which removes BATCH_QUEUE_MAX +# The bound that applies only when UPLOAD_AUTORUN is false. Files then wait in +# the admin review inbox rather than being queued, which removes BATCH_QUEUE_MAX # as the bound on that endpoint and leaves the volume as the only thing an # anonymous sender can exhaust. These are that bound; past either, the endpoint # answers 429 and stages nothing. diff --git a/docs/INGESTION_API.md b/docs/INGESTION_API.md index a16d983..aa03852 100644 --- a/docs/INGESTION_API.md +++ b/docs/INGESTION_API.md @@ -10,12 +10,17 @@ If you only need to **send sheets and watch them run**, §1 is the whole story a need no credential. §2 covers the admin side — the review inbox, starting a batch, and picking who executes it. -> **Deployment status.** Everything in §1 and §2 is **live on -> `https://mcp.nearle.ai.in`** — verified 2026-08-31 by probing the deployment -> itself, not by reading the source. The open drop takes no credential, and the -> stage timeline in [Stage-by-stage progress](#stage-by-stage-progress) is -> present in the served schema. If you want to re-confirm before integrating, -> [§8 Checking what is deployed](#8-checking-what-is-deployed) is a few curls. +> **Deployment status.** The endpoints, the open drop and the stage timeline in +> [Stage-by-stage progress](#stage-by-stage-progress) are **live on +> `https://mcp.nearle.ai.in`** — verified 2026-08-31 by probing the deployment, +> not by reading the source. +> +> **`UPLOAD_AUTORUN` is the exception: committed, not yet deployed.** Until it +> ships, an upload still lands in the review inbox and answers `pending`. Both +> modes are documented below and both are supported; if you are integrating +> right now, write the client so it does not care which is running — +> [§8 Checking what is deployed](#8-checking-what-is-deployed) shows how to tell +> in one curl. --- @@ -34,9 +39,18 @@ A blank allergens cell records *"we were not told"* — which is not the claim * ## Authentication **The catalogue drop (§1) needs no credential.** Anyone who can reach the host can -send it a spreadsheet. That is only safe because nothing sent there runs on arrival: -files wait in an admin review inbox, so the cost of an unwanted drop is disk until -somebody declines it — never products in the live catalogue. +send it a spreadsheet — and with `UPLOAD_AUTORUN` on, that spreadsheet runs through +the pipeline immediately and its products land in the live catalogue. There is no +undo; ingestion is an upsert. + +That is a deliberate trade, not an oversight. The requirement was uploads that run +without manual intervention, and a review queue that needs an admin to press a button +is not that. What bounds the endpoint is throughput rather than identity: the limits +in [Limits](#limits), and a single worker thread behind a queue of `BATCH_QUEUE_MAX`. +A sender can occupy the ingestion worker; they cannot multiply it. + +With `UPLOAD_AUTORUN=false` the older behaviour returns — files wait in an admin +review inbox and the cost of an unwanted drop is disk until somebody declines it. The orchestration routes (§2) and the three operator imports (§3) write straight to the database and stay credentialed. Two credential types are accepted there: @@ -75,9 +89,18 @@ anonymous submission the sender then cannot find. This is the one to give a colleague. No credential, up to 20 sheets per request, answers immediately with a `batch_id`. -**Nothing you send here runs on arrival.** The files are stored and an admin sees them -in the review inbox; only when they select the sheets and press Start does anything -reach the 11-stage pipeline or the catalogue. +**What happens on arrival depends on one setting, `UPLOAD_AUTORUN`:** + +| | `true` (the intended production mode) | `false` | +| --- | --- | --- | +| On arrival | The batch is queued and the 11 stages run | Files wait in the admin review inbox | +| Batch `status` | `queued`, then `running` | `pending` | +| The id you get back | **is the run** — poll it and watch `stages[]` | is a *drop* id; the run gets a different one, reached via the file's `released_to` | +| Anything to do | No | An admin must tick the files and press Start | + +**Write your client so it does not care which is running.** Both answer `202` with an +id that `GET /api/uploads/catalog/{id}` understands; only the number of hops differs. +Poll the id you were given, and follow `released_to` if a file ever grows one. ### `POST /api/uploads/catalog` → `202` @@ -131,8 +154,8 @@ The bad file is kept as a failed member rather than dropped, so a sender who sub ```jsonc { "batch_id": "49a82536866a483a9189954d3c749243", - "status": "pending", - "detail": "Waiting for review. Nothing runs until an admin starts it.", + "status": "queued", // "pending" when UPLOAD_AUTORUN=false + "detail": null, "submitted_by": "priya", "created_at": 1756370000.0, "updated_at": 1756370000.0, @@ -140,8 +163,9 @@ The bad file is kept as a failed member rather than dropped, so a sender who sub "files_done": 0, "files_failed": 1, "current_file": null, - "use_llm": false, - "fetch_images": false, + "use_llm": false, // UPLOAD_AUTORUN_USE_LLM + "fetch_images": true, // UPLOAD_AUTORUN_FETCH_IMAGES + "runner": "inprocess", "totals": { "rows_total": 0, "products_built": 0, "inserted": 0, "backfilled": 0, "skipped_existing": 0, "rejected": 0 }, "brands": [], @@ -151,13 +175,18 @@ The bad file is kept as a failed member rather than dropped, so a sender who sub { "index": 1, "filename": "notes.txt", "status": "failed", "detail": "The file has no data rows." } ], - "message": "1 file(s) received and waiting for review. 1 could not be read - see 'files' for the reason on each, and resend those." + "message": "1 file(s) accepted and queued for ingestion. 1 could not be read - see 'files' for the reason on each, and resend those." } ``` -Note the two levels of status. The **batch** is `pending` — waiting on a person. An -individual **file** inside it reads `queued`, meaning it is intact and eligible to be -run; it is waiting behind a decision, not behind the worker. +`use_llm` and `fetch_images` are reported, never accepted. They decide how much +outbound work a run commits the host to, and this endpoint's caller is anonymous, so +they come from settings — sending them in the request has no effect. + +There are two levels of status, and they are not the same question. `status: "queued"` +on the **batch** means it is behind the worker; on a **file** it means intact and not +started yet. On the review-inbox path a `pending` batch containing `queued` files means +something different again: those files are waiting behind a *decision*, not the worker. ### Polling @@ -165,6 +194,11 @@ 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`. +Under autorun that id is already the run, so the rest of this subsection does not +apply — go straight to [Stage-by-stage progress](#stage-by-stage-progress). The +release hop below is the `UPLOAD_AUTORUN=false` path, and is kept because that mode +is still supported and a client that handles both needs no branch. + **Your drop id stays valid for the whole lifecycle.** Poll it and read the per-file `status`: @@ -273,9 +307,9 @@ below. | Status | Meaning | | --- | --- | -| `pending` | In the review inbox. Nothing has run. | +| `pending` | In the review inbox. Nothing has run. `UPLOAD_AUTORUN=false` only | | `retired` | Every file in this drop has been released or dismissed. Read the per-file `status` | -| `queued` | Approved and waiting for the worker | +| `queued` | Waiting for the worker. Under autorun this is the first status an upload gets | | `running` | In the pipeline now | | `done` | Every file completed | | `partial` | Some landed, some failed. Deliberately not `done` — four of five succeeding must not read as flat success | @@ -368,12 +402,19 @@ But that stability is per `image_id`: change the product name and you get a new | Rows per file | 2,000 | — | file marked `failed`, others accepted | | Bytes per request | 50 MB | `BATCH_MAX_TOTAL_BYTES` | `413` | | Rows per request | 20,000 | `BATCH_MAX_TOTAL_ROWS` | `413` | -| Files awaiting review | 200 | `INBOX_MAX_PENDING_FILES` | `429` — nothing stored | -| Bytes awaiting review | 200 MB | `INBOX_MAX_PENDING_BYTES` | `429` — nothing stored | +| Batches queued behind the running one | 4 | `BATCH_QUEUE_MAX` | `429` — **autorun only** | +| Files awaiting review | 200 | `INBOX_MAX_PENDING_FILES` | `429` — `UPLOAD_AUTORUN=false` only | +| Bytes awaiting review | 200 MB | `INBOX_MAX_PENDING_BYTES` | `429` — `UPLOAD_AUTORUN=false` only | -The last two are the ceiling on the inbox as a whole, across every sender. They count -only what is still *awaiting review*, so starting or dismissing a drop frees its share -immediately. A `429` here stores nothing — resend once an admin has cleared space. +**Under autorun, expect a `429` occasionally and retry.** One batch runs at a time and +only four may wait behind it, so a handful of uploads in quick succession will hit the +ceiling. The files *are* staged when this happens and the message names the batch id, +so an admin can resume that one instead of you resending — but a client that treats +429 as a hard failure will report a working system as broken. Back off and retry. + +The two inbox ceilings apply only when `UPLOAD_AUTORUN=false`. They count what is still +*awaiting review*, so starting or dismissing a drop frees its share immediately, and a +`429` from them stores nothing — resend once an admin has cleared space. Drops nobody acts on are deleted after `BATCH_RETENTION_DAYS` (7). @@ -595,7 +636,8 @@ Upserts use `COALESCE` throughout, so a later, thinner sheet can never blank a v | `403` | Valid credential, wrong permission | Check the role table | | `413` | Over a file-count, byte or row ceiling | Split the drop | | `422` | **Operator imports only:** not one row could be imported | `detail.errors` lists row numbers and what each was missing | -| `429` | **Catalogue drop:** the review inbox is full | Nothing was stored. Ask an admin to clear it, then resend | +| `429` | **Catalogue drop, autorun:** four batches already queued | Normal under load. Back off and retry; `detail` names the staged batch id | +| `429` | **Catalogue drop, `UPLOAD_AUTORUN=false`:** the review inbox is full | Nothing was stored. Ask an admin to clear it, then resend | ```jsonc // 400 — wrong columns @@ -619,12 +661,12 @@ Sample sheets with the correct headers: `GET /api/upload/template/stores`, `/ana | Was | Is now | | --- | --- | | `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` | +| Files waited for review (28 Aug – 31 Aug 2026) | **They run on arrival again.** Expect `queued`, not `pending`, and no `released_to` hop — the id you get back is the run. `UPLOAD_AUTORUN=false` restores the review inbox | | 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` | | Progress was one `stage_index` scalar | Each file also carries a `stages[]` timeline, and the batch carries `stage_names` and `runner` — see [Stage-by-stage progress](#stage-by-stage-progress) | -| `?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 | +| `?use_llm` / `?fetch_images` on the drop | Still ignored. Under autorun they come from `UPLOAD_AUTORUN_USE_LLM` / `UPLOAD_AUTORUN_FETCH_IMAGES`; the response reports what was used | +| `429` meant the review inbox was full | Under autorun it means the worker queue is full again — retryable, and the files were staged | | `200` with `rows_imported: 0` | `422` with per-row reasons. Handle as a client error, not a server one | | Missing `customer_id` became `cust_imported` | Row is skipped. Supply a real customer id | | Missing `cost_price` invented as `mrp × 0.7` | Prices skipped for that row; inventory still lands. Send all three | @@ -669,18 +711,32 @@ Constraints enforced at boot, before any request is served: ## 7. What an open drop costs -Disk, and nothing else, until somebody looks at it. The endpoint queues nothing, so it -cannot occupy the ingestion worker and cannot reach the catalogue on its own — the -review gate is what makes accepting files from anyone acceptable. +**With `UPLOAD_AUTORUN=true`, it costs CPU and it reaches the catalogue.** This is the +part to be clear-eyed about: the endpoint takes no credential, so anyone who can reach +the host can cause products to be written to the live catalogue, and an ingest is an +upsert with no undo. That was chosen knowingly — the requirement was uploads that run +without manual intervention — but it should never be a surprise to whoever operates +this next. -`INBOX_MAX_PENDING_FILES` and `INBOX_MAX_PENDING_BYTES` bound what unreviewed -submissions can occupy; past either the endpoint answers `429` and stores nothing. -Dismissing a drop deletes its bytes immediately, and retention reclaims anything nobody -ever looks at. +What still bounds it is throughput, not identity: -If the host is reachable from the open internet, consider an IP allow-list at the proxy -as defence in depth. The controls above bound the damage; they do not stop a stranger -from filling the inbox with sheets an admin then has to decline. +- the per-request ceilings in [Limits](#limits) — 20 files, 50 MB, 20,000 rows; +- one worker thread running a single batch at a time, with `BATCH_QUEUE_MAX` waiting + behind it and a `429` past that. A sender can occupy the ingestion worker — that is + what it is for — but cannot multiply it, and cannot touch the request path the + healthcheck reads; +- `BATCH_RETENTION_DAYS`, which reclaims staged bytes either way. + +Nothing here bounds *who*. **If the host is reachable from the open internet, put an IP +allow-list on this route at the proxy** — with autorun on, that is no longer defence in +depth, it is the only control over who may write. + +**With `UPLOAD_AUTORUN=false`** the cost is disk and nothing else until somebody looks: +the endpoint queues nothing, so it cannot occupy the worker or reach the catalogue on +its own. `INBOX_MAX_PENDING_FILES` and `INBOX_MAX_PENDING_BYTES` bound what unreviewed +submissions occupy, past either it answers `429` and stores nothing, and dismissing a +drop deletes its bytes immediately. The worst a stranger can then do is fill an inbox +an admin has to decline. --- @@ -697,6 +753,18 @@ curl -s https://mcp.nearle.ai.in/api/health old, credentialed build had been rolled back. - **`GET /api/uploads/catalog/<32 random hex chars>` returns `404`, not `401`** → the anonymous read is live and the id really is the credential. +- **Which mode is running.** There is no settings endpoint, so send one small sheet and + read the batch `status` it comes back with: + ```bash + printf 'Product Name\nAmul Butter 100g\n' > /tmp/probe.csv + curl -s -X POST https://mcp.nearle.ai.in/api/uploads/catalog \ + -F 'files=@/tmp/probe.csv' -F 'sender=deployment-probe' \ + | python -c "import json,sys; print(json.load(sys.stdin)['status'])" + # queued -> UPLOAD_AUTORUN is on; that row is being ingested now + # pending -> the review inbox is on; nothing runs until an admin starts it + ``` + Note what the first answer means: the probe row **is ingested**. Use a product name + you are willing to see in the catalogue, or run this against a staging host. - **The served schema carries the stage timeline** → `stages[]` and `stage_names` are there, so a client can render the eleven stages: ```bash diff --git a/tests/conftest.py b/tests/conftest.py index 06bc8b9..6cd5fcd 100644 --- a/tests/conftest.py +++ b/tests/conftest.py @@ -109,6 +109,25 @@ def _reset_login_throttle(): auth_router._failures.clear() +@pytest.fixture +def review_inbox_mode(monkeypatch): + """Pin POST /api/uploads/catalog to the review-inbox path. + + `UPLOAD_AUTORUN` ships true, so an upload now runs the pipeline on arrival. + The inbox is still a supported mode - it is exactly what setting the flag + false turns back on - and the modules covering it declare this fixture + autouse, so their subject is stated rather than inherited from whatever the + default happens to be on the day. + + Patched on the ROUTER module, not on settings: the handler reads the name + out of its own globals at call time, which is the idiom batch_common's + docstring describes for the BATCH_* limits. + """ + from app.api.routers import uploads + + monkeypatch.setattr(uploads, "UPLOAD_AUTORUN", False) + + def _token(client: TestClient, username: str, password: str) -> str: resp = client.post("/api/auth/login", json={"username": username, "password": password}) assert resp.status_code == 200, resp.text diff --git a/tests/test_drop_lifecycle.py b/tests/test_drop_lifecycle.py index fae5664..cba3b96 100644 --- a/tests/test_drop_lifecycle.py +++ b/tests/test_drop_lifecycle.py @@ -33,6 +33,14 @@ HEADERS = ["Product Name", "Category", "Brand"] ROWS = [["Amul Butter 100g", "Butter", "Amul"]] +@pytest.fixture(autouse=True) +def _inbox_path(review_inbox_mode): + """This file is the drop -> release -> run lifecycle, which only exists on + the review-inbox path. Under UPLOAD_AUTORUN there is no release step: the id + the sender is handed is already the run. That shorter path is covered in + test_uploads_autorun.py.""" + + def _csv(headers=HEADERS, rows=ROWS) -> bytes: return ("\n".join([",".join(headers)] + [",".join(r) for r in rows])).encode() diff --git a/tests/test_review_inbox.py b/tests/test_review_inbox.py index c6df11a..6e376bf 100644 --- a/tests/test_review_inbox.py +++ b/tests/test_review_inbox.py @@ -35,6 +35,16 @@ HEADERS = ["Product Name", "Category", "Brand"] ROWS = [["Amul Butter 100g", "Butter", "Amul"]] +@pytest.fixture(autouse=True) +def _inbox_path(review_inbox_mode): + """Every test in this module is about the review-inbox path. + + Uploads run on arrival by default now (UPLOAD_AUTORUN), which would leave + this file asserting against an inbox nothing ever reaches. The gate itself + is still real and still shipped; this pins the mode that exercises it. + """ + + def _csv(rows=ROWS) -> bytes: return ("\n".join([",".join(HEADERS)] + [",".join(r) for r in rows])).encode() diff --git a/tests/test_uploads_api.py b/tests/test_uploads_api.py index 7e89e65..5e7222e 100644 --- a/tests/test_uploads_api.py +++ b/tests/test_uploads_api.py @@ -40,6 +40,16 @@ OTHER_KEY = "z" * 43 UPLOAD = "/api/uploads/catalog" +@pytest.fixture(autouse=True) +def _inbox_path(review_inbox_mode): + """The sender half of the review inbox, so pin that mode. + + The autorun path this endpoint now defaults to is covered separately in + test_uploads_autorun.py. Splitting them keeps each file asserting one + behaviour instead of branching inside every test. + """ + + def _csv(rows=ROWS) -> bytes: return ("\n".join([",".join(HEADERS)] + [",".join(r) for r in rows])).encode() diff --git a/tests/test_uploads_autorun.py b/tests/test_uploads_autorun.py new file mode 100644 index 0000000..c7deb85 --- /dev/null +++ b/tests/test_uploads_autorun.py @@ -0,0 +1,223 @@ +"""POST /api/uploads/catalog when it runs the pipeline on arrival. + +This is the shipped default (`UPLOAD_AUTORUN=true`). The review-inbox path that +the same endpoint takes when the flag is false is covered by +test_uploads_api.py and test_review_inbox.py; keeping them in separate files +means neither has to branch inside every test. + +WHAT THESE TESTS ARE REALLY GUARDING +------------------------------------ +The endpoint takes no credential, so with autorun on, an unauthenticated POST +starts real work against the live catalogue. That trade was made deliberately. +What must not happen is it being made *accidentally* - by the flag inverting in +a refactor, by the anonymous caller gaining a way to switch on the expensive +outbound stages, or by the worker queue quietly ceasing to bound anything. Each +of those has a test below. + +The other half is the sender's contract. Under autorun the id they are handed +IS the run id - there is no drop-then-release hop - and a client that polls it +must see stages, not a batch parked forever waiting for a human. +""" +from __future__ import annotations + +import io +import queue + +import pytest + +from app.api.batch_job_store import batch_job_store +from app.core import batch_ingest + +UPLOAD = "/api/uploads/catalog" +INBOX = "/api/admin/catalog-batch/inbox" + +HEADERS = ["Product Name", "Category", "Brand"] +ROWS = [["Amul Butter 100g", "Butter", "Amul"]] + + +def _csv(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): + """BATCH_UPLOAD_DIR under tmp_path, so nothing lands in the real data dir.""" + monkeypatch.setattr(batch_ingest, "BATCH_UPLOAD_DIR", tmp_path / "batch_uploads") + + +@pytest.fixture(autouse=True) +def submitted(monkeypatch): + """Stub the worker and record what reached it. + + Two jobs, as in test_uploads_api: a real worker still running after teardown + resolves BATCH_UPLOAD_DIR again and writes into the repository's data/ + directory, and stubbing turns "did this actually start?" into something a + test can assert rather than infer from a status string. + """ + 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() + + +# --------------------------------------------------------------------------- +# The feature +# --------------------------------------------------------------------------- +def test_an_anonymous_upload_starts_the_pipeline_immediately(client, submitted): + """The whole point. No credential, no admin, no button - it runs.""" + response = client.post(UPLOAD, files=_files(("a.csv", _csv()))) + assert response.status_code == 202, response.text + + body = response.json() + assert body["status"] != batch_ingest.PENDING + assert submitted == [body["batch_id"]], "the batch never reached the worker" + + +def test_the_id_the_sender_gets_is_the_run_itself(client): + """Under autorun there is no drop-then-release hop, and a client written + against the inbox contract would otherwise wait for a `released_to` that is + never coming.""" + body = client.post(UPLOAD, files=_files(("a.csv", _csv()))).json() + + polled = client.get(f"{UPLOAD}/{body['batch_id']}") + assert polled.status_code == 200 + assert polled.json()["batch_id"] == body["batch_id"] + assert all(f.get("released_to") is None for f in polled.json()["files"]) + + +def test_the_response_says_it_started_rather_than_that_it_is_waiting(client): + """The message is the only part of this a human reads. Telling a sender + their file is 'waiting for review' when it is already running is how a + working integration gets reported as broken.""" + message = client.post(UPLOAD, files=_files(("a.csv", _csv()))).json()["message"] + assert "queued for ingestion" in message + assert "review" not in message.lower() + + +def test_the_run_is_attributed_to_the_sender_label(client): + """`sender` is free text, not identity - but without it every auto-started + run in the batch list reads 'anonymous' and nobody can tell them apart.""" + body = client.post( + UPLOAD, files=_files(("a.csv", _csv())), data={"sender": "priya"} + ).json() + assert body["submitted_by"] == "priya" + + +# --------------------------------------------------------------------------- +# What the anonymous caller still cannot do +# --------------------------------------------------------------------------- +def test_the_sender_cannot_switch_on_the_expensive_stages(client): + """Image search and the LLM reach the network and are chosen by settings, + never by the request. An anonymous caller who could flip these could commit + a one-vCPU host to outbound work at will.""" + body = client.post( + UPLOAD, + files=_files(("a.csv", _csv())), + data={"use_llm": "true", "fetch_images": "false"}, + ).json() + + # The configured values, not the ones asked for. + assert body["use_llm"] is False + assert body["fetch_images"] is True + + +def test_the_configured_defaults_are_what_reach_the_run(client, monkeypatch): + """Both directions, so a wrong default cannot hide behind a matching one.""" + from app.api.routers import uploads + + monkeypatch.setattr(uploads, "UPLOAD_AUTORUN_FETCH_IMAGES", False) + monkeypatch.setattr(uploads, "UPLOAD_AUTORUN_USE_LLM", True) + + body = client.post(UPLOAD, files=_files(("a.csv", _csv()))).json() + assert body["fetch_images"] is False + assert body["use_llm"] is True + + +# --------------------------------------------------------------------------- +# The bound that replaces the inbox ceiling +# --------------------------------------------------------------------------- +def test_a_full_worker_queue_is_a_429_that_names_the_batch(client, monkeypatch): + """BATCH_QUEUE_MAX is the only thing bounding what an anonymous sender can + cost in CPU, so the refusal has to be real - and it has to name the batch, + because the files ARE staged and an admin can resume that id instead of the + caller resending the whole drop. + """ + from app.core import batch_worker + + def full(_batch_id): + raise queue.Full() + + monkeypatch.setattr(batch_worker, "submit", full) + + response = client.post(UPLOAD, files=_files(("a.csv", _csv()))) + 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].batch_id in response.json()["detail"], ( + "the sender cannot resume, so the id has to be in the message for " + "whoever can" + ) + + +def test_the_inbox_ceiling_no_longer_gates_an_autorun_upload(client, monkeypatch): + """INBOX_MAX_PENDING_* counts files AWAITING REVIEW. Under autorun nothing + ever awaits review, so that count is permanently zero: the check is not a + bound here and must not be consulted. Left in place it would read like + protection that is not actually doing anything. + """ + from app.api.routers import uploads + + monkeypatch.setattr(uploads, "INBOX_MAX_PENDING_FILES", 0) + + response = client.post(UPLOAD, files=_files(("a.csv", _csv()))) + assert response.status_code == 202, ( + "a zero inbox ceiling must not block a run that never touches the inbox" + ) + + +def test_nothing_accumulates_in_the_review_inbox(client, admin_headers): + """The visible consequence of the change, asserted so it is a decision and + not a surprise: the admin inbox is empty because uploads no longer stop + there. Recent runs, not pending files, is where they now appear.""" + client.post(UPLOAD, files=_files(("a.csv", _csv()))) + + inbox = client.get(INBOX, headers=admin_headers) + assert inbox.status_code == 200 + assert inbox.json()["pending_count"] == 0 + + +# --------------------------------------------------------------------------- +# The switch itself +# --------------------------------------------------------------------------- +def test_turning_the_flag_off_restores_the_review_inbox( + client, admin_headers, submitted, review_inbox_mode +): + """The escape hatch has to work, or "set it to false" is not a real answer + to somebody who decides the open endpoint was a mistake.""" + body = client.post(UPLOAD, files=_files(("a.csv", _csv()))).json() + + assert body["status"] == batch_ingest.PENDING + assert submitted == [], "nothing may reach the worker on the inbox path" + assert client.get(INBOX, headers=admin_headers).json()["pending_count"] == 1