backend updates for bulk product uploads from user
This commit is contained in:
@@ -19,7 +19,20 @@ async def _run_job(job_id: str, brand: str, max_products: int) -> None:
|
||||
job_store.update(job_id, "running")
|
||||
try:
|
||||
summary = await ingest_brand(brand, max_products=max_products)
|
||||
job_store.update(job_id, "done", detail=f"{summary['total_products']} products ingested")
|
||||
# "Ingested" has to mean "in the database". The generation stages can all
|
||||
# succeed while the pgvector write fails, and reporting that as done is
|
||||
# how a run that stored nothing ends up looking successful in the UI.
|
||||
if summary.get("storage_error"):
|
||||
job_store.update(
|
||||
job_id,
|
||||
"failed",
|
||||
detail=(
|
||||
f"Generated {summary['total_products']} product(s) but storing them "
|
||||
f"failed, so none are in the catalog: {summary['storage_error']}"
|
||||
),
|
||||
)
|
||||
else:
|
||||
job_store.update(job_id, "done", detail=f"{summary['total_products']} products ingested")
|
||||
except Exception as e: # noqa: BLE001 - surface any failure to the UI
|
||||
logger.exception("Catalog ingestion job %s failed", job_id)
|
||||
job_store.update(job_id, "failed", detail=str(e))
|
||||
|
||||
File diff suppressed because it is too large
Load Diff
@@ -940,11 +940,17 @@ class ProductCatalogEngine:
|
||||
if 'image_id' not in p:
|
||||
id_source = p.get('product_name') or p.get('title', 'unknown_product')
|
||||
p['image_id'] = s3_service.generate_image_id(id_source)
|
||||
upsert_brand_products(brand, enhanced_products, cleanup=True)
|
||||
logger.info("🧠 Stored embeddings to pgvector")
|
||||
stored = upsert_brand_products(brand, enhanced_products, cleanup=True)
|
||||
logger.info("🧠 Stored %s product(s) with embeddings to pgvector", stored)
|
||||
except Exception as e:
|
||||
logger.warning(f"Vector storage skipped/failed: {e}")
|
||||
|
||||
# Stage 4 is where the catalog becomes readable by the app - every
|
||||
# API read goes to pgvector, not to the dict returned here. So a
|
||||
# failure at this stage means the run produced nothing the user can
|
||||
# see, and it has to travel back to the job status rather than being
|
||||
# logged and forgotten.
|
||||
logger.error("❌ Stage 4 (pgvector storage) failed for '%s': %s", brand, e)
|
||||
catalog['storage_error'] = str(e)
|
||||
|
||||
return catalog
|
||||
|
||||
def save_catalog(self, catalog: Dict[str, Any], filename: str = None) -> str:
|
||||
|
||||
@@ -40,11 +40,17 @@ async def ingest_brand(brand: str, max_products: int = 50) -> Dict[str, Any]:
|
||||
# Placed here rather than in the API router because this function is also
|
||||
# the CLI's entry point (cli/ingest_brand.py), and a failure to write the
|
||||
# file must not turn a successful ingest into a failed job.
|
||||
try:
|
||||
from app.services.brand_sync import export_brand_to_seed_file
|
||||
export_brand_to_seed_file(brand)
|
||||
except Exception: # noqa: BLE001 - DB rows are already committed
|
||||
logger.warning("Seed-catalog export failed for %s (DB rows intact)", brand, exc_info=True)
|
||||
storage_error = catalog.get("storage_error")
|
||||
|
||||
# Nothing reached the database, so there is nothing to mirror out of it -
|
||||
# and running the export anyway would either write an empty file or leave a
|
||||
# stale one looking current.
|
||||
if not storage_error:
|
||||
try:
|
||||
from app.services.brand_sync import export_brand_to_seed_file
|
||||
export_brand_to_seed_file(brand)
|
||||
except Exception: # noqa: BLE001 - DB rows are already committed
|
||||
logger.warning("Seed-catalog export failed for %s (DB rows intact)", brand, exc_info=True)
|
||||
|
||||
summary = {
|
||||
"brand": brand,
|
||||
@@ -52,6 +58,7 @@ async def ingest_brand(brand: str, max_products: int = 50) -> Dict[str, Any]:
|
||||
"total_images": catalog.get("total_images", 0),
|
||||
"duration_seconds": round(duration, 2),
|
||||
"engine_info": catalog.get("engine_info", {}),
|
||||
"storage_error": storage_error,
|
||||
}
|
||||
logger.info("Finished ingestion for brand=%s in %.2fs: %s products",
|
||||
brand, duration, summary["total_products"])
|
||||
|
||||
@@ -22,6 +22,28 @@ from app.services.s3_service import s3_service
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
class VectorStoreUnavailable(RuntimeError):
|
||||
"""The pgvector database could not be reached, so a write did not happen.
|
||||
|
||||
Exists because the silent alternative was a production bug that took a long
|
||||
time to see: upsert_brand_products() used to `return` when _connect() gave
|
||||
back None - unreachable host, wrong DB_PASSWORD, USE_PGVECTOR=false - and
|
||||
every caller read that as a successful write. The upload endpoints then
|
||||
answered "success", the seed JSON was updated, and not one row existed in
|
||||
the database. A write that cannot happen has to raise.
|
||||
"""
|
||||
|
||||
|
||||
class VectorStoreWriteFailed(RuntimeError):
|
||||
"""The INSERT ran without error but the rows are not in the table.
|
||||
|
||||
Guards against the failure modes an exception cannot catch: a statement
|
||||
silently rolled back, a trigger swallowing the row, or an ON CONFLICT
|
||||
target that quietly matched nothing. The only trustworthy proof of a write
|
||||
is reading it back.
|
||||
"""
|
||||
|
||||
|
||||
def _sanitize_name(name: str) -> str:
|
||||
"""Sanitize a brand name for use as a PostgreSQL table name suffix.
|
||||
|
||||
@@ -193,22 +215,44 @@ def ensure_brand_schema(brand: str) -> str:
|
||||
return table_name
|
||||
|
||||
|
||||
def upsert_brand_products(brand: str, products: List[Dict[str, Any]], cleanup: bool = False) -> None:
|
||||
def upsert_brand_products(brand: str, products: List[Dict[str, Any]], cleanup: bool = False) -> int:
|
||||
"""Insert products into brand-specific table - simplified with only essential fields
|
||||
|
||||
When `cleanup=True`, any products in the table whose image_id is NOT in the
|
||||
provided `products` list are deleted after the upsert. This ensures the
|
||||
database exactly reflects the source data. The caller is responsible for
|
||||
providing the complete set of products for the brand when using cleanup.
|
||||
|
||||
Returns the number of distinct image_ids confirmed present in the table
|
||||
afterwards - the rows are read back, so a non-raising call is proof of
|
||||
persistence rather than proof that a statement was merely sent.
|
||||
|
||||
Raises:
|
||||
VectorStoreUnavailable: the database is unreachable; nothing was written.
|
||||
VectorStoreWriteFailed: the statements ran but the rows are not there.
|
||||
ValueError: a product carries no image_id, which is the primary key
|
||||
every other product is deduplicated on.
|
||||
"""
|
||||
if not products:
|
||||
return 0
|
||||
|
||||
conn = _connect()
|
||||
if not conn:
|
||||
return
|
||||
|
||||
raise VectorStoreUnavailable(
|
||||
f"Cannot save products for '{brand}': the product database is unreachable "
|
||||
f"(USE_PGVECTOR={USE_PGVECTOR}, host={DB_HOST}:{DB_PORT}, db={DB_NAME}). "
|
||||
f"Nothing was saved. Check DB_HOST/DB_USER/DB_PASSWORD in backend/.env "
|
||||
f"and that Postgres is accepting connections."
|
||||
)
|
||||
|
||||
table_name = ensure_brand_schema(brand)
|
||||
if not table_name:
|
||||
return
|
||||
|
||||
conn.close()
|
||||
raise VectorStoreUnavailable(
|
||||
f"Cannot save products for '{brand}': the brand table could not be created "
|
||||
f"or verified in database '{DB_NAME}'. Nothing was saved."
|
||||
)
|
||||
|
||||
rows = []
|
||||
for p in products:
|
||||
# Extract only essential fields
|
||||
@@ -333,55 +377,89 @@ def upsert_brand_products(brand: str, products: List[Dict[str, Any]], cleanup: b
|
||||
|
||||
image_ids = [r[4] for r in rows if r[4]]
|
||||
|
||||
with conn, conn.cursor() as cur:
|
||||
cur.executemany(
|
||||
f"""
|
||||
INSERT INTO {table_name}
|
||||
(product_name, title, description, category, image_id, image_url, image_urls, price_range, size_variants, providers,
|
||||
fssai_license, product_sku, sku_source, hsn_code, final_selling_price, selling_price, barcode, barcode_type, highlights, nutrients, search_query, embedding)
|
||||
VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s)
|
||||
ON CONFLICT (image_id) DO UPDATE SET
|
||||
product_name = EXCLUDED.product_name,
|
||||
title = EXCLUDED.title,
|
||||
description = EXCLUDED.description,
|
||||
category = EXCLUDED.category,
|
||||
image_url = EXCLUDED.image_url,
|
||||
image_urls = EXCLUDED.image_urls,
|
||||
price_range = EXCLUDED.price_range,
|
||||
size_variants = EXCLUDED.size_variants,
|
||||
providers = EXCLUDED.providers,
|
||||
fssai_license = EXCLUDED.fssai_license,
|
||||
product_sku = EXCLUDED.product_sku,
|
||||
sku_source = EXCLUDED.sku_source,
|
||||
hsn_code = EXCLUDED.hsn_code,
|
||||
final_selling_price = EXCLUDED.final_selling_price,
|
||||
selling_price = EXCLUDED.selling_price,
|
||||
barcode = EXCLUDED.barcode,
|
||||
barcode_type = EXCLUDED.barcode_type,
|
||||
highlights = EXCLUDED.highlights,
|
||||
nutrients = EXCLUDED.nutrients,
|
||||
search_query = EXCLUDED.search_query,
|
||||
embedding = EXCLUDED.embedding,
|
||||
updated_at = CURRENT_TIMESTAMP
|
||||
""",
|
||||
rows,
|
||||
# An empty image_id is not a harmless blank: it is the conflict target, so
|
||||
# two such products would overwrite each other and the second would replace
|
||||
# the first instead of being added. Refuse the batch and name the rows.
|
||||
if len(image_ids) != len(rows):
|
||||
unnamed = [r[0] or "<no product_name>" for r in rows if not r[4]]
|
||||
conn.close()
|
||||
raise ValueError(
|
||||
f"{len(unnamed)} product(s) for '{brand}' have no image_id and cannot be "
|
||||
f"stored (a product needs a name with at least one letter or digit): "
|
||||
f"{', '.join(unnamed[:5])}"
|
||||
)
|
||||
logger.info(f"✅ Upserted {len(rows)} products into {table_name}")
|
||||
|
||||
# Remove stale products that were deleted from the source data.
|
||||
# Only runs when cleanup=True so that callers processing partial
|
||||
# product sets (e.g. multiple seed files contributing to the same
|
||||
# brand table) don't accidentally orphan each other's data.
|
||||
if cleanup and image_ids:
|
||||
cur.execute(
|
||||
f"DELETE FROM {table_name} WHERE image_id != ALL(%s::text[])",
|
||||
(image_ids,),
|
||||
expected_ids = sorted(set(image_ids))
|
||||
|
||||
try:
|
||||
with conn, conn.cursor() as cur:
|
||||
cur.executemany(
|
||||
f"""
|
||||
INSERT INTO {table_name}
|
||||
(product_name, title, description, category, image_id, image_url, image_urls, price_range, size_variants, providers,
|
||||
fssai_license, product_sku, sku_source, hsn_code, final_selling_price, selling_price, barcode, barcode_type, highlights, nutrients, search_query, embedding)
|
||||
VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s)
|
||||
ON CONFLICT (image_id) DO UPDATE SET
|
||||
product_name = EXCLUDED.product_name,
|
||||
title = EXCLUDED.title,
|
||||
description = EXCLUDED.description,
|
||||
category = EXCLUDED.category,
|
||||
image_url = EXCLUDED.image_url,
|
||||
image_urls = EXCLUDED.image_urls,
|
||||
price_range = EXCLUDED.price_range,
|
||||
size_variants = EXCLUDED.size_variants,
|
||||
providers = EXCLUDED.providers,
|
||||
fssai_license = EXCLUDED.fssai_license,
|
||||
product_sku = EXCLUDED.product_sku,
|
||||
sku_source = EXCLUDED.sku_source,
|
||||
hsn_code = EXCLUDED.hsn_code,
|
||||
final_selling_price = EXCLUDED.final_selling_price,
|
||||
selling_price = EXCLUDED.selling_price,
|
||||
barcode = EXCLUDED.barcode,
|
||||
barcode_type = EXCLUDED.barcode_type,
|
||||
highlights = EXCLUDED.highlights,
|
||||
nutrients = EXCLUDED.nutrients,
|
||||
search_query = EXCLUDED.search_query,
|
||||
embedding = EXCLUDED.embedding,
|
||||
updated_at = CURRENT_TIMESTAMP
|
||||
""",
|
||||
rows,
|
||||
)
|
||||
deleted = cur.rowcount
|
||||
if deleted:
|
||||
logger.info(f"🗑️ Removed {deleted} stale product(s) from {table_name}")
|
||||
|
||||
conn.close()
|
||||
# Remove stale products that were deleted from the source data.
|
||||
# Only runs when cleanup=True so that callers processing partial
|
||||
# product sets (e.g. multiple seed files contributing to the same
|
||||
# brand table) don't accidentally orphan each other's data.
|
||||
if cleanup and image_ids:
|
||||
cur.execute(
|
||||
f"DELETE FROM {table_name} WHERE image_id != ALL(%s::text[])",
|
||||
(image_ids,),
|
||||
)
|
||||
deleted = cur.rowcount
|
||||
if deleted:
|
||||
logger.info(f"🗑️ Removed {deleted} stale product(s) from {table_name}")
|
||||
|
||||
# Read the rows back. This is the line that turns "we sent an
|
||||
# INSERT" into "the data is in the table", and it is the only
|
||||
# signal the API layer is allowed to report success on.
|
||||
cur.execute(
|
||||
f"SELECT COUNT(DISTINCT image_id) FROM {table_name} WHERE image_id = ANY(%s::text[])",
|
||||
(expected_ids,),
|
||||
)
|
||||
row = cur.fetchone()
|
||||
persisted = int(row[0]) if row else 0
|
||||
|
||||
if persisted < len(expected_ids):
|
||||
raise VectorStoreWriteFailed(
|
||||
f"Wrote {len(rows)} product(s) for '{brand}' to {table_name} but only "
|
||||
f"{persisted} of {len(expected_ids)} are readable back afterwards. "
|
||||
f"The data was not saved - treat this as a failed import."
|
||||
)
|
||||
|
||||
logger.info("✅ Upserted %d product(s) into %s (%d verified in table)",
|
||||
len(rows), table_name, persisted)
|
||||
finally:
|
||||
conn.close()
|
||||
|
||||
# Every write path into the catalog funnels through here, so this is the
|
||||
# one place that has to invalidate the derived views: the brand cards'
|
||||
@@ -393,6 +471,8 @@ def upsert_brand_products(brand: str, products: List[Dict[str, Any]], cleanup: b
|
||||
except Exception: # noqa: BLE001 - cache invalidation must never fail a write
|
||||
pass
|
||||
|
||||
return persisted
|
||||
|
||||
|
||||
def get_existing_product_image_id(brand: str, product_name: str) -> Optional[str]:
|
||||
"""Check if a product with this name exists in the brand table and return its image_id"""
|
||||
|
||||
Reference in New Issue
Block a user