Files

512 lines
18 KiB
Python
Raw Permalink Normal View History

"""Retail 360 workspace API — notes, media, milestones, RSS, wholesalers."""
from __future__ import annotations
import json
import urllib.request
from datetime import datetime
from typing import Any, Optional
from fastapi import APIRouter, HTTPException, Query
from pydantic import BaseModel, Field
from app.db import fetch_all, fetch_one
from app.middleware import log_agent_event
from app import retail_360
from app import wholesaler_scrapers
from app.connectors import market_stocks, rss_feeds
from app.connectors import food_trends
router = APIRouter(prefix="/retail", tags=["retail-360"])
class NoteIn(BaseModel):
body: str = Field(..., min_length=1)
title: Optional[str] = None
note_type: str = "general"
class MilestoneIn(BaseModel):
title: str
milestone_type: str = "custom"
client_id: Optional[int] = None
deal_id: Optional[int] = None
target_date: Optional[str] = None
value_eur: Optional[float] = None
notes: Optional[str] = None
class OwnershipIn(BaseModel):
new_owner: str
previous_owner: Optional[str] = None
change_type: str = "acquisition"
effective_date: Optional[str] = None
source: Optional[str] = None
notes: Optional[str] = None
class CalendarIn(BaseModel):
title: str
starts_at: str
description: Optional[str] = None
ends_at: Optional[str] = None
client_id: Optional[int] = None
deal_id: Optional[int] = None
location: Optional[str] = None
class MediaIn(BaseModel):
filename: str
storage_path: str
content_type: str = "image/jpeg"
caption: Optional[str] = None
def _row(row: dict | None) -> dict[str, Any]:
if not row:
raise HTTPException(404, "Not found")
out: dict[str, Any] = {}
for k, v in row.items():
if hasattr(v, "isoformat"):
out[k] = v.isoformat()
elif v is not None and hasattr(v, "__float__") and type(v).__name__ == "Decimal":
out[k] = float(v)
else:
out[k] = v
return out
def _fetch_weather_forecast(lat: float, lon: float) -> list[dict[str, Any]]:
url = (
f"https://api.open-meteo.com/v1/forecast?latitude={lat}&longitude={lon}"
f"&daily=temperature_2m_max,precipitation_sum,weathercode"
f"&timezone=Europe%2FAmsterdam&forecast_days=7"
)
try:
with urllib.request.urlopen(url, timeout=15) as resp:
data = json.loads(resp.read().decode())
days = data.get("daily", {}).get("time", [])
temps = data.get("daily", {}).get("temperature_2m_max", [])
prec = data.get("daily", {}).get("precipitation_sum", [])
return [
{"date": days[i], "temperature_c": temps[i] if i < len(temps) else None,
"precipitation_mm": prec[i] if i < len(prec) else None, "source": "open-meteo-live"}
for i in range(len(days))
]
except Exception:
return []
@router.get("/360/{store_id}")
def get_360_view(store_id: int) -> dict[str, Any]:
try:
data = retail_360.get_store_360(store_id)
except ValueError as exc:
raise HTTPException(404, str(exc)) from exc
store = data["store"]
if store.get("latitude") and store.get("longitude"):
live = _fetch_weather_forecast(float(store["latitude"]), float(store["longitude"]))
if live:
data["weather_forecast"] = live
for key in ("notes", "media", "milestones", "ownership_changes", "calendar", "weather"):
data[key] = [_row(x) for x in data.get(key, [])]
if data.get("area_analysis"):
data["area_analysis"] = _row(data["area_analysis"])
return data
@router.post("/360/{store_id}/notes")
def add_store_note(store_id: int, payload: NoteIn) -> dict[str, Any]:
note = retail_360.add_note("supermarket", store_id, payload.body, payload.title, payload.note_type)
log_agent_event(agent_name="retail_360", event_type="note", title=f"Note on store {store_id}")
return {"note": _row(note)}
@router.post("/360/{store_id}/milestones")
def add_store_milestone(store_id: int, payload: MilestoneIn) -> dict[str, Any]:
ms = retail_360.add_milestone(store_id, payload.title, payload.milestone_type, **payload.model_dump(exclude={"title", "milestone_type"}))
return {"milestone": _row(ms)}
@router.post("/360/{store_id}/ownership")
def add_store_ownership(store_id: int, payload: OwnershipIn) -> dict[str, Any]:
store = fetch_one("SELECT chain FROM supermarkets WHERE id = %s", (store_id,))
row = retail_360.add_ownership(entity_id=store_id, chain=store.get("chain") if store else None, **payload.model_dump())
return {"ownership": _row(row)}
@router.post("/360/{store_id}/calendar")
def add_store_calendar(store_id: int, payload: CalendarIn) -> dict[str, Any]:
ev = retail_360.add_calendar_event(store_id, payload.title, payload.starts_at, **payload.model_dump(exclude={"title", "starts_at"}))
return {"event": _row(ev)}
@router.post("/360/{store_id}/media")
def register_store_media(store_id: int, payload: MediaIn) -> dict[str, Any]:
media = retail_360.register_media("supermarket", store_id, payload.filename, payload.storage_path, payload.content_type, payload.caption)
return {"media": _row(media)}
@router.get("/cities")
def list_cities(
limit: int = Query(200, ge=1, le=1000),
q: Optional[str] = None,
min_population: Optional[int] = None,
min_muslim_pct: Optional[float] = None,
sort: str = Query("population", pattern="^(population|muslim|stores|city)$"),
) -> dict[str, Any]:
clauses, params = [], []
if q:
clauses.append("c.city ILIKE %s")
params.append(f"%{q}%")
if min_population:
clauses.append("c.population >= %s")
params.append(min_population)
if min_muslim_pct:
clauses.append("c.muslim_proxy_pct >= %s")
params.append(min_muslim_pct)
where = (" WHERE " + " AND ".join(clauses)) if clauses else ""
order_map = {
"population": "c.population DESC NULLS LAST",
"muslim": "c.muslim_proxy_pct DESC NULLS LAST",
"stores": "store_count DESC",
"city": "c.city ASC",
}
order = order_map.get(sort, order_map["population"])
rows = fetch_all(
f"""SELECT c.*, (SELECT COUNT(*) FROM supermarkets s WHERE s.city ILIKE c.city) AS store_count
FROM city_demographics c{where} ORDER BY {order} LIMIT %s""",
tuple(params + [limit]),
)
return {"items": [_row(r) for r in rows], "count": len(rows)}
@router.post("/cities/sync")
def sync_cities(limit: int = Query(50, ge=1, le=200)) -> dict[str, Any]:
return retail_360.sync_city_demographics(limit)
@router.get("/wholesalers")
def list_wholesalers(
limit: int = Query(500, ge=1, le=2000),
q: Optional[str] = None,
province: Optional[str] = None,
city: Optional[str] = None,
halal_certified: Optional[bool] = None,
has_phone: Optional[bool] = None,
has_email: Optional[bool] = None,
sort: str = Query("name", pattern="^(name|city|province)$"),
) -> dict[str, Any]:
clauses, params = [], []
if q:
clauses.append("(name ILIKE %s OR city ILIKE %s OR address ILIKE %s OR email ILIKE %s)")
like = f"%{q}%"
params.extend([like, like, like, like])
if province:
clauses.append("province ILIKE %s")
params.append(province)
if city:
clauses.append("city ILIKE %s")
params.append(f"%{city}%")
if halal_certified is True:
clauses.append("halal_certified = TRUE")
if has_phone is True:
clauses.append("phone IS NOT NULL AND phone <> ''")
if has_email is True:
clauses.append("email IS NOT NULL AND email <> ''")
where = (" WHERE " + " AND ".join(clauses)) if clauses else ""
order = {"name": "name", "city": "city", "province": "province"}.get(sort, "name")
rows = fetch_all(f"SELECT * FROM wholesalers{where} ORDER BY {order} LIMIT %s", tuple(params + [limit]))
total = fetch_one(f"SELECT COUNT(*) AS n FROM wholesalers{where}", tuple(params) if params else None)
return {"items": [_row(r) for r in rows], "count": len(rows), "total": int(total["n"]) if total else len(rows)}
@router.get("/wholesalers/meta")
def wholesalers_meta() -> dict[str, Any]:
provinces = fetch_all(
"SELECT province, COUNT(*) AS n FROM wholesalers WHERE province IS NOT NULL GROUP BY province ORDER BY n DESC"
)
return {
"provinces": [_row(p) for p in provinces],
"total": _safe_count_wh("wholesalers"),
}
def _safe_count_wh(table: str) -> int:
row = fetch_one(f"SELECT COUNT(*) AS n FROM {table}")
return int(row["n"]) if row else 0
@router.get("/wholesalers/{wh_id}/contacts")
def wholesaler_contacts(wh_id: int) -> dict[str, Any]:
rows = fetch_all(
"SELECT * FROM wholesaler_contacts WHERE wholesaler_id = %s ORDER BY confidence DESC, full_name",
(wh_id,),
)
wh = fetch_one("SELECT id, name, phone, email, address, city, province, website, linkedin_url FROM wholesalers WHERE id = %s", (wh_id,))
if not wh:
raise HTTPException(404, "Wholesaler not found")
return {"wholesaler": _row(wh), "contacts": [_row(r) for r in rows]}
class WholesalerContactIn(BaseModel):
full_name: str
role: str = "contact"
phone: Optional[str] = None
email: Optional[str] = None
linkedin_url: Optional[str] = None
@router.post("/wholesalers/{wh_id}/contacts")
def add_wholesaler_contact(wh_id: int, payload: WholesalerContactIn) -> dict[str, Any]:
row = fetch_one(
"""INSERT INTO wholesaler_contacts (wholesaler_id, full_name, role, phone, email, linkedin_url, source)
VALUES (%s, %s, %s, %s, %s, %s, 'manual') RETURNING *""",
(wh_id, payload.full_name, payload.role, payload.phone, payload.email, payload.linkedin_url),
)
return {"contact": _row(row)}
class RssBookmarkIn(BaseModel):
rss_item_id: int
title: Optional[str] = None
link: Optional[str] = None
feed_name: Optional[str] = None
notes: Optional[str] = None
@router.get("/rss/bookmarks")
def list_rss_bookmarks(limit: int = Query(50, ge=1, le=200)) -> dict[str, Any]:
rows = fetch_all(
"""SELECT b.*, i.title AS item_title, i.link AS item_link, f.name AS feed_name
FROM rss_bookmarks b
LEFT JOIN rss_items i ON i.id = b.rss_item_id
LEFT JOIN rss_feeds f ON f.id = i.feed_id
ORDER BY b.created_at DESC LIMIT %s""",
(limit,),
)
out = []
for r in rows:
row = _row(r)
row["title"] = row.get("title") or row.get("item_title")
row["link"] = row.get("link") or row.get("item_link")
out.append(row)
return {"items": out, "count": len(out)}
@router.post("/rss/bookmarks")
def add_rss_bookmark(payload: RssBookmarkIn) -> dict[str, Any]:
from app.db import execute
item = fetch_one("SELECT id, title, link FROM rss_items WHERE id = %s", (payload.rss_item_id,))
if not item:
raise HTTPException(404, "RSS item not found")
execute(
"""INSERT INTO rss_bookmarks (rss_item_id, title, link, feed_name, notes)
VALUES (%s, %s, %s, %s, %s)
ON CONFLICT (rss_item_id) DO UPDATE SET title=EXCLUDED.title, link=EXCLUDED.link, feed_name=EXCLUDED.feed_name, notes=EXCLUDED.notes""",
(
payload.rss_item_id,
payload.title or item.get("title"),
payload.link or item.get("link"),
payload.feed_name,
payload.notes,
),
)
row = fetch_one("SELECT * FROM rss_bookmarks WHERE rss_item_id = %s", (payload.rss_item_id,))
return {"bookmark": _row(row)}
@router.delete("/rss/bookmarks/{rss_item_id}")
def delete_rss_bookmark(rss_item_id: int) -> dict[str, Any]:
from app.db import execute
execute("DELETE FROM rss_bookmarks WHERE rss_item_id = %s", (rss_item_id,))
return {"ok": True}
@router.get("/promo-campaigns")
def list_promo_campaigns(
chain: Optional[str] = None,
status: str = Query("active"),
q: Optional[str] = None,
folder_type: Optional[str] = None,
valid_days: Optional[int] = Query(None, ge=1, le=365),
limit: int = Query(100, ge=1, le=500),
) -> dict[str, Any]:
clauses, params = ["p.status = %s"], [status]
if chain:
clauses.append("p.chain ILIKE %s")
params.append(f"%{chain}%")
if q:
clauses.append("(p.title ILIKE %s OR p.chain ILIKE %s OR p.description ILIKE %s)")
params.extend([f"%{q}%"] * 3)
if folder_type:
clauses.append("(p.promo_type ILIKE %s OR p.metadata->>'folder_type' ILIKE %s)")
params.extend([f"%{folder_type}%", f"%{folder_type}%"])
if valid_days:
clauses.append("p.valid_to IS NOT NULL AND p.valid_to <= CURRENT_DATE + %s * INTERVAL '1 day'")
params.append(valid_days)
where = " WHERE " + " AND ".join(clauses)
rows = fetch_all(
f"""SELECT p.*, s.name AS store_name FROM promo_campaigns p
LEFT JOIN supermarkets s ON s.id = p.supermarket_id
{where} ORDER BY p.valid_to ASC NULLS LAST, p.chain ASC, p.created_at DESC LIMIT %s""",
tuple(params + [limit]),
)
return {"items": [_row(r) for r in rows], "count": len(rows)}
@router.get("/reclamefolder/chains")
def reclamefolder_chains() -> dict[str, Any]:
from app.connectors import reclamefolder
chains = reclamefolder.list_chains()
return {"chains": chains, "count": len(chains)}
class PromoCampaignIn(BaseModel):
chain: Optional[str] = None
title: str
folder_path: Optional[str] = None
folder_label: Optional[str] = None
description: Optional[str] = None
valid_from: Optional[str] = None
valid_to: Optional[str] = None
status: str = "active"
promo_type: str = "folder"
@router.post("/promo-campaigns")
def add_promo_campaign(payload: PromoCampaignIn) -> dict[str, Any]:
row = fetch_one(
"""INSERT INTO promo_campaigns (chain, title, folder_path, folder_label, description, valid_from, valid_to, status, promo_type)
VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s) RETURNING *""",
(
payload.chain, payload.title, payload.folder_path, payload.folder_label,
payload.description, payload.valid_from, payload.valid_to, payload.status, payload.promo_type,
),
)
return {"campaign": _row(row)}
@router.post("/reclamefolder/refresh")
def refresh_reclamefolders() -> dict[str, Any]:
from app.connectors import reclamefolder
from app.middleware import log_agent_event
log_agent_event(
agent_name="reclamefolder",
event_type="refresh",
title="Reclamefolder.nl sync",
)
return reclamefolder.sync_to_db()
@router.get("/reclamefolder/live")
def live_reclamefolders(limit: int = Query(50, ge=1, le=200)) -> dict[str, Any]:
from app.connectors import reclamefolder
try:
result = reclamefolder.sync_to_db()
items = result.get("items") or []
except Exception as exc:
items = reclamefolder.list_cached(limit)
return {"items": items, "count": len(items), "cached": True, "error": str(exc)}
return {"items": items[:limit], "count": len(items), "synced": True, "source": "reclamefolder.nl"}
@router.post("/wholesalers/import")
def import_wholesalers(background: bool = Query(False)) -> dict[str, Any]:
log_agent_event(agent_name="wholesale_scraper", event_type="import", title="OSM wholesalers import")
if background:
import threading
threading.Thread(target=wholesaler_scrapers.import_wholesalers, daemon=True).start()
return {"status": "started", "message": "Wholesaler import running in background"}
return wholesaler_scrapers.import_wholesalers()
@router.get("/rss/live")
def rss_live(limit: int = Query(30, ge=1, le=100), category: Optional[str] = None) -> dict[str, Any]:
rows = rss_feeds.list_live_feed(limit, category)
return {"items": [_row(r) for r in rows], "count": len(rows)}
@router.post("/rss/refresh")
def rss_refresh() -> dict[str, Any]:
log_agent_event(agent_name="rss_feeds", event_type="refresh", title="RSS feeds refresh")
return rss_feeds.refresh_all_feeds()
@router.get("/market/stocks")
def retail_market_stocks() -> dict[str, Any]:
quotes = market_stocks.fetch_retail_quotes()
return {
"items": quotes,
"summary": market_stocks.market_summary(quotes),
"updated_at": datetime.utcnow().isoformat(),
}
@router.get("/market/supermarkets")
def supermarket_market_board() -> dict[str, Any]:
quotes = market_stocks.fetch_supermarket_quotes()
listed = [q for q in quotes if q.get("listed")]
return {
"items": [_row(q) for q in quotes],
"listed": [_row(q) for q in listed],
"unlisted_nl": [_row(q) for q in quotes if not q.get("listed")],
"summary": market_stocks.market_summary(listed),
"data_source": market_stocks.DATA_SOURCE,
"updated_at": datetime.utcnow().isoformat(),
}
@router.get("/market/food-trends")
def market_food_trends() -> dict[str, Any]:
return food_trends.food_trends_dashboard()
@router.get("/market/concepts")
def market_concepts(limit: int = Query(6, ge=1, le=12)) -> dict[str, Any]:
listed = market_stocks.fetch_retail_quotes()
summary = market_stocks.market_summary(listed)
best = summary.get("best_performer")
concepts = food_trends.generate_concepts(market_best=best, limit=limit)
return {
"concepts": concepts,
"summary": summary,
"updated_at": datetime.utcnow().isoformat(),
}
@router.get("/regulations")
def retail_regulations(limit: int = Query(30, ge=1, le=100)) -> dict[str, Any]:
reg = rss_feeds.list_live_feed(limit, "regelgeving")
cbs = rss_feeds.list_live_feed(limit, "cbs")
markt = rss_feeds.list_live_feed(min(limit, 15), "markt")
return {
"regelgeving": [_row(r) for r in reg],
"cbs": [_row(r) for r in cbs],
"markt": [_row(r) for r in markt],
"updated_at": datetime.utcnow().isoformat(),
}
@router.get("/live-dashboard")
def live_dashboard() -> dict[str, Any]:
trends = fetch_all(
"SELECT * FROM market_trends ORDER BY updated_at DESC NULLS LAST LIMIT 8"
)
rss = rss_feeds.list_live_feed(12)
opportunities = fetch_all(
"""SELECT s.name, s.chain, s.city, ros.halal_opportunity_score
FROM retail_opportunity_scores ros JOIN supermarkets s ON s.id = ros.supermarket_id
ORDER BY ros.halal_opportunity_score DESC LIMIT 5"""
)
quotes = market_stocks.fetch_retail_quotes()
return {
"trends": [_row(t) for t in trends],
"rss": [_row(r) for r in rss],
"top_opportunities": [_row(o) for o in opportunities],
"market_stocks": quotes,
"market_summary": market_stocks.market_summary(quotes),
"updated_at": datetime.utcnow().isoformat(),
}