Health score updates in backend

This commit is contained in:
sriram
2026-09-04 12:06:40 +05:30
parent 1300d4c678
commit 0c3e23fac4
33 changed files with 59703 additions and 94 deletions

View File

@@ -245,6 +245,14 @@ class BatchOut(BaseModel):
# sits queued until the orchestrator picks it up, which the UI has to be
# able to say out loud rather than showing a run that looks stuck.
runner: str = batch_ingest.RUNNER_INPROCESS
# The nutrition-enrichment job queued when this batch finished, pollable at
# GET /api/admin/nutrition-intelligence/jobs/{job_id}. Surfaced here so the
# poll a client is already doing for the ingestion also reveals the scoring
# that follows it, instead of leaving the client to guess that it happened.
#
# None while the batch is still running, when AUTO_ENRICH_ON_UPLOAD is off,
# and for every batch that finished before this existed.
nutrition_job_id: Optional[str] = None
# The 11 stage names, in order, so a client can draw the whole pipeline
# before a file has entered any of it. Served rather than duplicated in the
# frontend so the two cannot drift when a stage is added.

View File

@@ -128,3 +128,56 @@ class NutritionEnrichmentJobOut(BaseModel):
unavailable: Optional[int] = None
duration_seconds: Optional[float] = None
error: Optional[str] = None
# ---------------------------------------------------------------------------
# Health-score listing (GET /api/nutrition/health-scores)
# ---------------------------------------------------------------------------
class HealthScoreItemOut(BaseModel):
"""One scored consumable product.
`health_score` is REQUIRED here, against this module's all-Optional house
style, and that is the point: the query filters on `health_score IS NOT
NULL`, so a null arriving in this model means the query lost its guarantee.
Declaring it required is what turns that into a loud failure instead of a
blank cell in whatever dashboard is reading this.
"""
brand: str
image_id: str
product_name: Optional[str] = None
category: Optional[str] = None
health_score: float
nutrition_score: Optional[float] = None
health_band: str # derived from health_score; see _health_band()
scoring_version: Optional[str] = None
# Provenance travels with the number. A score is only as trustworthy as the
# source it was computed from, and the consumer should be able to show that.
data_status: Optional[str] = None
data_source: Optional[str] = None
source_url: Optional[str] = None
calories_kcal: Optional[float] = None
protein_g: Optional[float] = None
dietary_fiber_g: Optional[float] = None
total_sugar_g: Optional[float] = None
sodium_mg: Optional[float] = None
diet_tags: Optional[List[str]] = None
allergens: Optional[List[str]] = None
model_config = {"extra": "ignore"}
class HealthScoreListOut(BaseModel):
"""Envelope matching the catalogue convention in `schemas.ProductListOut`
rather than the bare-array convention of the older nutrition list
endpoints - `total` is the whole reason this endpoint exists, since without
it a client paging a thousand products cannot tell when it has finished."""
total: int
limit: int
offset: int
generated_at: str
items: List[HealthScoreItemOut] = []

View File

@@ -1,12 +1,17 @@
from __future__ import annotations
from typing import List, Optional
import csv
import io
from datetime import datetime, timezone
from typing import Any, Dict, Iterator, List, Optional
from fastapi import APIRouter, HTTPException, Query
from fastapi.responses import StreamingResponse
from app.api.nutrition_schemas import (
FullNutritionOut, HealthyAlternativeOut, NutritionInsightsOut,
PersonalizedRecommendationOut, ProductListItemOut, SimilarProductOut,
FullNutritionOut, HealthScoreItemOut, HealthScoreListOut,
HealthyAlternativeOut, NutritionInsightsOut, PersonalizedRecommendationOut,
ProductListItemOut, SimilarProductOut,
)
from app.intelligence import nutrition_recommendation, nutrition_similarity
from app.services import nutrition_alternatives_service, nutrition_analytics_service, nutrition_db
@@ -100,6 +105,139 @@ def filter_products(
return [ProductListItemOut(**r) for r in results]
# ---------------------------------------------------------------------------
# GET Health Scores - every scored consumable product
# ---------------------------------------------------------------------------
#
# Distinct from `/products` above, which exists to answer "show me the ten
# highest-protein snacks". This one exists to answer "give me the health score
# of every consumable product in the catalogue", and that needs three things
# `/products` does not offer: a total so a client knows when it has finished
# paging, a brand filter, and a one-request export.
# Derived from health_score alone, never stored - a stored band would be one
# more thing that can fall out of step with the score beside it.
#
# The 60 boundary is not a fresh invention: `nutrition_db.store_healthy_
# distribution()` already counts `health_score >= 60` as a "healthy product".
# Picking a different number here would leave two parts of the same API
# disagreeing about what healthy means.
_HEALTH_BANDS = ((80.0, "excellent"), (60.0, "good"), (40.0, "fair"))
def _health_band(score: float) -> str:
for threshold, label in _HEALTH_BANDS:
if score >= threshold:
return label
return "poor"
def _to_item(row: Dict[str, Any]) -> HealthScoreItemOut:
return HealthScoreItemOut(**row, health_band=_health_band(row["health_score"]))
_CSV_COLUMNS = (
"brand", "image_id", "product_name", "category",
"health_score", "health_band", "nutrition_score", "scoring_version",
"data_status", "data_source", "source_url",
"calories_kcal", "protein_g", "dietary_fiber_g", "total_sugar_g", "sodium_mg",
"diet_tags", "allergens",
)
def _csv_rows(rows: Iterator[Dict[str, Any]]) -> Iterator[str]:
"""Yields the export a row at a time so neither the full result set nor the
full response body is ever held in memory at once."""
buffer = io.StringIO()
writer = csv.writer(buffer, lineterminator=chr(10))
def flush() -> str:
value = buffer.getvalue()
buffer.seek(0)
buffer.truncate(0)
return value
writer.writerow(_CSV_COLUMNS)
yield flush()
for row in rows:
row = dict(row, health_band=_health_band(row["health_score"]))
writer.writerow([
# A Postgres TEXT[] would render as "['Vegan', 'Gluten Free']" via
# str(); pipe-joining keeps the cell readable in a spreadsheet.
"|".join(row[c] or []) if c in ("diet_tags", "allergens") else
("" if row.get(c) is None else row[c])
for c in _CSV_COLUMNS
])
yield flush()
@router.get(
"/health-scores",
response_model=None,
responses={200: {
"model": HealthScoreListOut,
"content": {"application/json": {}, "text/csv": {}},
"description": "Paged JSON envelope, or the whole list as CSV with format=csv.",
}},
)
def list_health_scores(
brand: Optional[str] = Query(None, description="Exact brand match, e.g. 'Own Products'"),
category: Optional[str] = Query(None, description="Substring match, case-insensitive"),
min_score: Optional[float] = Query(None, ge=0, le=100),
max_score: Optional[float] = Query(None, ge=0, le=100),
include_unknown: bool = Query(
False,
description="Also return products whose edibility was never confirmed. "
"Non-consumables are never returned either way.",
),
sort_by: str = Query("health_score", description="health_score|nutrition_score|protein|fiber|sugar|sodium|calcium|iron|vitamin_c|calories|product_name|brand|category"),
order: str = Query("desc", pattern="^(asc|desc)$"),
limit: int = Query(100, ge=1, le=500),
offset: int = Query(0, ge=0),
fmt: str = Query("json", alias="format", pattern="^(json|csv)$"),
):
"""Every consumable product that has a health score.
Only scored products are returned: a consumable whose nutrition could not be
matched to a verified source has no score, and is absent rather than present
with a null - `compute_scores` refuses to invent one, and this endpoint
refuses to imply one.
Non-consumables (soap, shampoo, mosquito repellent) are excluded by the
stored `edibility` verdict, as are products the classifier could not decide
on. `include_unknown=true` opts the undecided ones back in; nothing opts a
confirmed non-consumable in.
`format=csv` streams the entire matching set in one response and ignores
`limit`/`offset`, which is the point of it.
"""
filters = {
"brand": brand, "category": category,
"min_score": min_score, "max_score": max_score,
"include_unknown": include_unknown,
}
if fmt == "csv":
stamp = datetime.now(timezone.utc).strftime("%Y%m%d")
return StreamingResponse(
_csv_rows(nutrition_db.iter_health_scores(sort_by=sort_by, order=order, **filters)),
media_type="text/csv; charset=utf-8",
headers={"Content-Disposition": f'attachment; filename="health_scores_{stamp}.csv"'},
)
rows = nutrition_db.query_health_scores(
sort_by=sort_by, order=order, limit=limit, offset=offset, **filters)
return HealthScoreListOut(
# Counted with the SAME filters that produced `rows`, so the two cannot
# describe different populations.
total=nutrition_db.count_health_scores(**filters),
limit=limit, offset=offset,
generated_at=datetime.now(timezone.utc).isoformat(timespec="seconds"),
items=[_to_item(r) for r in rows],
)
# ---------------------------------------------------------------------------
# GET Nutrition Analytics (Feature 9)
# ---------------------------------------------------------------------------

View File

@@ -9,6 +9,7 @@ from app.api.background import run_in_background
from app.api.deps import require_admin
from app.api.nutrition_job_store import nutrition_job_store
from app.services import nutrition_enrichment_service
from app.services.nutrition_autoenrich import run_enrich_job
logger = logging.getLogger(__name__)
router = APIRouter(prefix="/admin/nutrition-intelligence", tags=["admin", "nutrition"])
@@ -19,30 +20,13 @@ class EnrichRequest(BaseModel):
generate_narrative: bool = True
max_products: int | None = None
def _run_enrich_job(job_id: str, skip_if_verified: bool, generate_narrative: bool, max_products: int | None) -> None:
nutrition_job_store.update(job_id, status="running")
def progress_cb(done: int, total: int) -> None:
nutrition_job_store.update(job_id, processed=done, total=total)
try:
result = nutrition_enrichment_service.enrich_all_products(
skip_if_verified=skip_if_verified, generate_narrative=generate_narrative,
progress_cb=progress_cb, max_products=max_products,
)
nutrition_job_store.update(
job_id, status="done", detail="Enrichment complete",
result={
"total_products": result.total_products, "verified": result.verified,
"partial": result.partial, "unavailable": result.unavailable,
"duration_seconds": result.duration_seconds, "error_count": len(result.errors),
"errors": result.errors[:20],
},
)
except Exception as e: # noqa: BLE001
logger.exception("Nutrition enrichment job %s failed", job_id)
nutrition_job_store.update(job_id, status="failed", detail=str(e))
# `enrich_all_products` has accepted these three for a while; the API had no
# way to pass them, so the endpoint could not do what `scripts/
# enrich_nutrition.py` could. In particular there was no way to widen past
# ACTIVE_BRANDS, which is what a catalogue-wide run needs.
brands: list[str] | None = None
include_inactive: bool = False
categories: list[str] | None = None
def _run_train_job(job_id: str) -> None:
@@ -64,7 +48,15 @@ def enrich_nutrition(payload: EnrichRequest) -> dict:
Equivalent to `python scripts/enrich_nutrition.py`."""
job = nutrition_job_store.create("enrich")
run_in_background(
lambda: _run_enrich_job(job.job_id, payload.skip_if_verified, payload.generate_narrative, payload.max_products),
lambda: run_enrich_job(
job.job_id,
skip_if_verified=payload.skip_if_verified,
generate_narrative=payload.generate_narrative,
max_products=payload.max_products,
brands=payload.brands,
include_inactive=payload.include_inactive,
categories=payload.categories,
),
name=f"nutrition-enrich-{job.job_id[:8]}",
)
return {"job_id": job.job_id, "status": job.status}

View File

@@ -49,7 +49,8 @@ from fastapi.responses import PlainTextResponse
from app.api.deps import require_permission
from app.services.vector_store import _connect
from app.services import store_db, nutrition_db
from app.services import store_db, nutrition_db, nutrition_scoring
from app.services.consumability import Edibility, classify_edibility
logger = logging.getLogger(__name__)
router = APIRouter(prefix="/upload", tags=["upload"])
@@ -566,6 +567,7 @@ async def upload_nutrition_file(file: UploadFile = File(...)) -> Dict[str, Any]:
imported_count = 0
errors: List[Dict[str, Any]] = []
skipped = 0
skipped_non_consumable = 0
status_counts = {"verified": 0, "partial": 0, "unavailable": 0}
try:
@@ -592,6 +594,20 @@ async def upload_nutrition_file(file: UploadFile = File(...)) -> Dict[str, Any]:
image_id = image_id or _slug_id(brand, product_name)
category = _opt_str(row, ['category', 'cat'])
# THE GATE THIS ENDPOINT NEVER HAD. Every other write path into
# nutrition_facts refuses non-food; a spreadsheet could walk
# straight past them and put a health score on a shampoo, which
# then appeared in the public listing endpoints and the
# analytics leaderboards. A refused row is reported, not
# silently dropped, so the uploader can see it was not imported.
verdict = classify_edibility(category or "", product_name or "")
if verdict.edibility is Edibility.NON_CONSUMABLE:
skipped += 1
skipped_non_consumable += 1
_record(errors, index,
"not a food or drink, so no nutrition was stored ({})".format(verdict.reason))
continue
values = {
"calories_kcal": _opt_float(row, ['calories', 'calories_kcal', 'energy']),
"protein_g": _opt_float(row, ['protein', 'protein_g']),
@@ -617,14 +633,15 @@ async def upload_nutrition_file(file: UploadFile = File(...)) -> Dict[str, Any]:
cur.execute(
"""
INSERT INTO nutrition_facts
(brand, image_id, product_name, category, data_status, data_source,
(brand, image_id, product_name, category, edibility, data_status, data_source,
calories_kcal, protein_g, carbohydrates_g, total_sugar_g, dietary_fiber_g,
total_fat_g, sodium_mg, calcium_mg, iron_mg, vitamin_c_mg, updated_at)
VALUES (%s, %s, %s, %s, %s, 'excel_upload',
VALUES (%s, %s, %s, %s, %s, %s, 'excel_upload',
%s, %s, %s, %s, %s, %s, %s, %s, %s, %s, CURRENT_TIMESTAMP)
ON CONFLICT (brand, image_id) DO UPDATE SET
product_name = COALESCE(EXCLUDED.product_name, nutrition_facts.product_name),
category = COALESCE(EXCLUDED.category, nutrition_facts.category),
edibility = EXCLUDED.edibility,
data_source = 'excel_upload',
-- Never downgrade a row a trusted source already verified.
data_status = CASE WHEN nutrition_facts.data_status = 'verified'
@@ -641,7 +658,7 @@ async def upload_nutrition_file(file: UploadFile = File(...)) -> Dict[str, Any]:
vitamin_c_mg = COALESCE(EXCLUDED.vitamin_c_mg, nutrition_facts.vitamin_c_mg),
updated_at = CURRENT_TIMESTAMP
""",
(brand, image_id, product_name, category, data_status,
(brand, image_id, product_name, category, verdict.edibility.value, data_status,
values["calories_kcal"], values["protein_g"], values["carbohydrates_g"],
values["total_sugar_g"], values["dietary_fiber_g"], values["total_fat_g"],
values["sodium_mg"], values["calcium_mg"], values["iron_mg"],
@@ -652,7 +669,7 @@ async def upload_nutrition_file(file: UploadFile = File(...)) -> Dict[str, Any]:
# Only ever built from values this row actually carried. A row
# that carried none produces no insight row at all, rather than
# a confident-looking one full of defaults.
health_score = _opt_float(row, ['health_score', 'nutrition_score', 'score'])
sheet_score = _opt_float(row, ['health_score', 'nutrition_score', 'score'])
diet_tags_raw = _opt_str(row, ['diet_tags', 'tags', 'diet'])
allergens_raw = _opt_str(row, ['allergens', 'allergen'])
@@ -664,14 +681,47 @@ async def upload_nutrition_file(file: UploadFile = File(...)) -> Dict[str, Any]:
# told", and it is what stops an empty list reading as "none".
allergen_source = "upload" if allergens is not None else "unavailable"
positives = []
if values["protein_g"] is not None:
positives.append(f"Contains {values['protein_g']}g protein per 100g")
if values["dietary_fiber_g"] is not None:
positives.append(f"Provides {values['dietary_fiber_g']}g dietary fiber")
cautions = []
if values["total_sugar_g"] is not None:
cautions.append(f"{values['total_sugar_g']}g sugar per 100g")
# THE SCORE IS COMPUTED HERE, NOT READ OFF THE SHEET.
#
# This endpoint used to store whatever was in a `health_score`
# column. That number shares a column - and every endpoint that
# sorts, ranks and filters on it - with scores produced by
# `nutrition_scoring.compute_scores` from USDA and Open Food
# Facts data. Two scales in one column means the ordering means
# nothing: an uploaded 95 outranks a computed 80 without either
# number having been measured the same way.
#
# So the same function scores every row in the system, and the
# same rules generate every insight. The uploaded nutrients are
# the input; the score is a conclusion drawn from them.
facts = dict(values, data_status=data_status,
allergens=allergens or [])
scores = nutrition_scoring.compute_scores(facts)
if scores:
nutrition_score = scores["nutrition_score"]
health_score = scores["health_score"]
scoring_version = scores["scoring_version"]
else:
# Too few nutrients to score fairly. The operator's own
# number is kept rather than discarded - they may have it
# from a source this sheet does not carry - but
# `scoring_version` says where it came from, so nobody
# later mistakes it for one of ours.
nutrition_score = None
health_score = sheet_score
scoring_version = "excel_upload"
positives = nutrition_scoring.generate_positive_insights(facts)
cautions = nutrition_scoring.generate_cautions(facts)
# Union, not replacement: the rules derive what the numbers
# support ("High Protein"), while the sheet may carry what they
# cannot ("Vegan" is an ingredient fact, not a nutrient one).
derived_tags = nutrition_scoring.classify_diet_tags(facts)
if derived_tags or diet_tags:
diet_tags = list(dict.fromkeys((diet_tags or []) + derived_tags))
if allergens:
allergens = nutrition_scoring.normalize_allergens(facts)
if health_score is not None or diet_tags or allergens or positives or cautions:
cur.execute(
@@ -680,13 +730,13 @@ async def upload_nutrition_file(file: UploadFile = File(...)) -> Dict[str, Any]:
(brand, image_id, data_status, nutrition_score, health_score,
scoring_version, positive_insights, nutritional_cautions,
diet_tags, allergens, allergen_source, generated_at)
VALUES (%s, %s, %s, %s, %s, 'excel_upload', %s, %s, %s, %s, %s, CURRENT_TIMESTAMP)
VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, CURRENT_TIMESTAMP)
ON CONFLICT (brand, image_id) DO UPDATE SET
data_status = CASE WHEN nutrition_insights.data_status = 'verified'
THEN 'verified' ELSE EXCLUDED.data_status END,
nutrition_score = COALESCE(EXCLUDED.nutrition_score, nutrition_insights.nutrition_score),
health_score = COALESCE(EXCLUDED.health_score, nutrition_insights.health_score),
scoring_version = 'excel_upload',
scoring_version = EXCLUDED.scoring_version,
positive_insights = COALESCE(EXCLUDED.positive_insights, nutrition_insights.positive_insights),
nutritional_cautions = COALESCE(EXCLUDED.nutritional_cautions, nutrition_insights.nutritional_cautions),
diet_tags = COALESCE(EXCLUDED.diet_tags, nutrition_insights.diet_tags),
@@ -696,8 +746,8 @@ async def upload_nutrition_file(file: UploadFile = File(...)) -> Dict[str, Any]:
ELSE EXCLUDED.allergen_source END,
generated_at = CURRENT_TIMESTAMP
""",
(brand, image_id, data_status, health_score, health_score,
positives or None, cautions or None,
(brand, image_id, data_status, nutrition_score, health_score,
scoring_version, positives or None, cautions or None,
diet_tags, allergens, allergen_source)
)
@@ -715,6 +765,11 @@ async def upload_nutrition_file(file: UploadFile = File(...)) -> Dict[str, Any]:
file.filename, len(df), imported_count, errors, skipped,
"nutritional intelligence items",
data_status_counts=status_counts,
# Reported separately from `rows_skipped`, which otherwise reads as "the
# sheet was malformed". These rows were well-formed and deliberately not
# imported. Named to match `EnrichmentResult.skipped_non_consumable`, so
# both write paths report a refusal the same way.
skipped_non_consumable=skipped_non_consumable,
)

View File

@@ -918,6 +918,16 @@ def _process_upload(filename: str, content: bytes) -> Dict[str, Any]:
logger.info("Upload '%s': saved %d, failed %d, skipped %d blank row(s)",
filename, body["added_count"], body["error_count"], skipped_blank_rows)
# Score the consumables that were just stored. `outcome.brands` is already
# the set of brands actually written, so nothing new is tracked to answer
# this. Returns None (and logs) rather than raising if scoring cannot be
# queued - the products are saved either way, and this response has already
# been shaped to say so.
from app.services.nutrition_autoenrich import submit_enrichment_for_brands
body["nutrition_job_id"] = submit_enrichment_for_brands(
outcome.brands, source="upload of '{}'".format(filename))
return body