backend env updates and brand json files
This commit is contained in:
@@ -140,6 +140,18 @@ DATABASE_URL = os.getenv(
|
||||
f"postgresql://{DB_USER}:{DB_PASSWORD}@{DB_HOST}:{DB_PORT}/{DB_NAME}",
|
||||
)
|
||||
|
||||
# How often to re-run the brand-table <-> seed-catalog reconcile after boot.
|
||||
#
|
||||
# The startup run alone only catches what existed at boot. A brand table created
|
||||
# directly in the database while the server is up - by hand, by a script, or by
|
||||
# another machine sharing this database - is not mirrored into the seed catalogs
|
||||
# until the next restart. This interval is what closes that window.
|
||||
#
|
||||
# reconcile_brand_catalogs() is idempotent and non-destructive, so a sweep that
|
||||
# finds nothing to do costs one COUNT(*) per brand table and writes nothing.
|
||||
# 0 disables the loop, leaving the startup run and POST /api/system/brand-sync.
|
||||
BRAND_SYNC_INTERVAL_SECONDS = int(os.getenv("BRAND_SYNC_INTERVAL_SECONDS", "300"))
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# S3 / DigitalOcean Spaces (product image storage) - optional
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
46
app/main.py
46
app/main.py
@@ -20,7 +20,7 @@ from fastapi.staticfiles import StaticFiles
|
||||
from fastapi.responses import FileResponse
|
||||
|
||||
from app.infrastructure.persistence import restore_bundled_assets
|
||||
from app.infrastructure.settings import API_CORS_ORIGINS
|
||||
from app.infrastructure.settings import API_CORS_ORIGINS, BRAND_SYNC_INTERVAL_SECONDS
|
||||
from app.api.routers import health, brands, search, chat, catalog, system
|
||||
from app.api.routers import stores, discounts, analytics as store_analytics, trending, recommendations, store_admin
|
||||
from app.api.routers import nutrition, nutrition_admin, upload
|
||||
@@ -46,6 +46,12 @@ except Exception as e: # pragma: no cover - depends on an optional dependency
|
||||
MCP_PATH = "/mcp"
|
||||
|
||||
|
||||
# Signals the background thread below to stop. Doubles as its sleep: waiting on
|
||||
# an Event rather than time.sleep() means shutdown is immediate instead of
|
||||
# blocking until the current BRAND_SYNC_INTERVAL_SECONDS elapses.
|
||||
_brand_sync_stop = threading.Event()
|
||||
|
||||
|
||||
@asynccontextmanager
|
||||
async def lifespan(_app: FastAPI):
|
||||
"""
|
||||
@@ -55,6 +61,10 @@ async def lifespan(_app: FastAPI):
|
||||
Postgres, which may be slow or briefly unreachable on a cold boot, and the
|
||||
server answering /api/health in under a second is what lets the container
|
||||
healthcheck pass while that settles.
|
||||
|
||||
That same thread then stays alive to re-run the brand reconcile on an
|
||||
interval, so a brand table that appears in the database after boot reaches
|
||||
the catalog without a restart.
|
||||
"""
|
||||
|
||||
# Runs before the thread below, and synchronously: the seed catalogs and
|
||||
@@ -80,6 +90,26 @@ async def lifespan(_app: FastAPI):
|
||||
except Exception as e:
|
||||
logger.warning("Startup background init warning: %s", e)
|
||||
|
||||
# Everything above is a one-shot boot sweep. Keep going on an interval so
|
||||
# a brand table created directly in the database - bypassing the app, and
|
||||
# therefore bypassing the cache invalidation every write path does - still
|
||||
# reaches the catalog. Without this it waits for the next restart.
|
||||
if BRAND_SYNC_INTERVAL_SECONDS <= 0:
|
||||
logger.info("Periodic brand sync disabled (BRAND_SYNC_INTERVAL_SECONDS=0)")
|
||||
return
|
||||
|
||||
logger.info("Periodic brand sync every %ss", BRAND_SYNC_INTERVAL_SECONDS)
|
||||
while not _brand_sync_stop.wait(BRAND_SYNC_INTERVAL_SECONDS):
|
||||
try:
|
||||
# Imported per iteration for the same reason as above: a failure
|
||||
# to import must not be what kills the loop.
|
||||
from app.api.routers.system import _run_background_brand_sync
|
||||
_run_background_brand_sync()
|
||||
except Exception as e:
|
||||
# One bad sweep - an unreachable database, a malformed seed file -
|
||||
# must not end the loop, or recovery needs a restart again.
|
||||
logger.warning("Periodic brand sync warning: %s", e)
|
||||
|
||||
threading.Thread(target=_async_init, daemon=True).start()
|
||||
|
||||
# The mounted MCP app carries its own lifespan, which starts the session
|
||||
@@ -87,11 +117,17 @@ async def lifespan(_app: FastAPI):
|
||||
# that lifespan - only the outermost app's is executed - so without this
|
||||
# the endpoint exists, accepts a connection, and then fails on the first
|
||||
# message with a session manager that was never started.
|
||||
if _mcp_app is not None:
|
||||
async with _mcp_app.lifespan(_mcp_app):
|
||||
try:
|
||||
if _mcp_app is not None:
|
||||
async with _mcp_app.lifespan(_mcp_app):
|
||||
yield
|
||||
else:
|
||||
yield
|
||||
else:
|
||||
yield
|
||||
finally:
|
||||
# Releases the interval wait immediately. The thread is a daemon, so this
|
||||
# is not what lets the process exit - it is what stops a reconcile from
|
||||
# starting against a database the shutdown is already tearing down.
|
||||
_brand_sync_stop.set()
|
||||
|
||||
|
||||
app = FastAPI(
|
||||
|
||||
@@ -579,10 +579,16 @@ def get_brand_overview(force_refresh: bool = False) -> List[Dict[str, Any]]:
|
||||
|
||||
img = None
|
||||
image_cols = [c for c in ("image_id", "image_url", "image_urls") if c in columns]
|
||||
if "image_url" in columns or "image_urls" in columns:
|
||||
where = " OR ".join(
|
||||
f"({c} IS NOT NULL)" for c in ("image_url", "image_urls") if c in columns
|
||||
)
|
||||
url_cols = [c for c in ("image_url", "image_urls") if c in columns]
|
||||
# Prefer a row that carries a stored URL. When the table has
|
||||
# no URL column at all, still pick a row with an image_id:
|
||||
# the router resolves that against S3 (brands.py), so the
|
||||
# card gets an image instead of falling back to initials.
|
||||
# Brand tables written by a different DDL really can lack
|
||||
# both URL columns while their objects sit in the bucket.
|
||||
filter_cols = url_cols or [c for c in image_cols if c == "image_id"]
|
||||
if filter_cols:
|
||||
where = " OR ".join(f"({c} IS NOT NULL)" for c in filter_cols)
|
||||
order = " ORDER BY updated_at DESC" if "updated_at" in columns else ""
|
||||
cur.execute(
|
||||
f"SELECT {', '.join(image_cols)} FROM {table_name} "
|
||||
|
||||
Reference in New Issue
Block a user