upload-catalog-integration

This commit is contained in:
sriram
2026-08-29 11:21:11 +05:30
parent 5aa2669f7d
commit 27d53fa957
8 changed files with 584 additions and 48 deletions

View File

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

View File

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

View File

@@ -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}")

View File

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

View File

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