|
|
|
|
@@ -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
|