Files
foodlinkk-command-center/tools-api/app/export_intel.py
T
Aissa c5ea8cbb63 Ship War Room UX, CEO inbox, mobile chat, and Agents NAS workspace.
Make A2A togglable, keep reply context while reading, route agent data to //10.4.7.11/share/Agents, and polish themes plus export-intel sync cancel.
2026-08-16 22:40:32 +00:00

1460 lines
51 KiB
Python

"""Export Intel API — Foodlinkk B2B global export intelligence."""
from __future__ import annotations
import csv
import io
import hashlib
import json
import time
from collections import defaultdict
from typing import Any, Optional
from fastapi import APIRouter, HTTPException, Query
from fastapi.responses import StreamingResponse
from pydantic import BaseModel, Field
from app.db import execute, execute_returning, fetch_all, fetch_one
from app.export_halal_scoring import COUNTRY_HALAL_MARKET, halal_pin_tier, score_entity
from app.export_intel_crm import push_entities_to_crm
HALAL_REASON_NL: dict[str, str] = {
"sterke_halal_markt": "Sterke halal-markt in dit land",
"halal_in_naam": "Halal expliciet in bedrijfsnaam",
"halal_signaal": "Kebab/döner/vlees-signaal in naam",
"direct_vlees_kanaal": "Direct vlees-verkoopkanaal",
"B2B_kanaal": "B2B inkoop / distributie",
"heeft_contact": "Contactpersoon in registry",
}
def _halal_reason_labels(reasons: list[str]) -> list[str]:
out: list[str] = []
for r in reasons:
if r in HALAL_REASON_NL:
out.append(HALAL_REASON_NL[r])
elif r.startswith("type:"):
out.append("Type: " + r.split(":", 1)[1].replace("_", " "))
else:
out.append(r.replace("_", " "))
return out
def _normalize_url(url: Any) -> Optional[str]:
if url is None:
return None
s = str(url).strip()
if not s or s.lower() in ("null", "none", "-", "—"):
return None
if s.startswith("mailto:") or s.startswith("tel:"):
return s
if not s.startswith(("http://", "https://")):
s = "https://" + s.lstrip("/")
return s
def _display_url(url: str) -> str:
return url.replace("https://", "").replace("http://", "").rstrip("/")[:96]
def _collect_websites(
row: dict[str, Any],
contacts: Optional[list[dict[str, Any]]] = None,
brand_website: Optional[str] = None,
) -> list[dict[str, str]]:
seen: set[str] = set()
out: list[dict[str, str]] = []
def add(label: str, raw: Any) -> None:
href = _normalize_url(raw)
if not href:
return
key = href.lower().rstrip("/")
if key in seen:
return
seen.add(key)
out.append({
"label": (label or "Website")[:64],
"url": href,
"display": _display_url(href),
})
add("Organisatie", row.get("website"))
add("Brand", brand_website)
for c in contacts or []:
role = (c.get("role") or c.get("name") or "Contact").strip()
add(role, c.get("website_contact_url"))
if c.get("linkedin_url"):
add("LinkedIn", c.get("linkedin_url"))
src = c.get("source_url")
if src and "linkedin.com" not in str(src).lower():
add("Bron", src)
return out
def _clean_str(val: Any) -> Optional[str]:
if val is None:
return None
s = str(val).strip()
if not s or s.lower() in ("null", "none", "-", "—", "n/a"):
return None
return s
def _dedupe_items(items: list[dict[str, str]], key: str = "value") -> list[dict[str, str]]:
seen: set[str] = set()
out: list[dict[str, str]] = []
for item in items:
val = _clean_str(item.get(key))
if not val:
continue
k = val.lower()
if k in seen:
continue
seen.add(k)
out.append({**item, key: val})
return out
def _collect_emails(row: dict[str, Any], contacts: Optional[list[dict[str, Any]]] = None) -> list[dict[str, str]]:
items: list[dict[str, str]] = []
def add(label: str, raw: Any) -> None:
val = _clean_str(raw)
if not val or "@" not in val:
return
items.append({"label": (label or "E-mail")[:64], "value": val, "href": f"mailto:{val}"})
add("Organisatie", row.get("email"))
for c in contacts or []:
who = _clean_str(c.get("name")) or _clean_str(c.get("role")) or "Contact"
add(who, c.get("email"))
alt = _clean_str(c.get("email_alt"))
if alt and alt != _clean_str(c.get("email")):
add(f"{who} (alt)", alt)
return _dedupe_items(items)
def _collect_phones(row: dict[str, Any], contacts: Optional[list[dict[str, Any]]] = None) -> list[dict[str, str]]:
items: list[dict[str, str]] = []
def add(label: str, raw: Any, kind: str = "tel") -> None:
val = _clean_str(raw)
if not val:
return
digits = "".join(ch for ch in val if ch.isdigit() or ch == "+")
href = f"https://wa.me/{digits.lstrip('+')}" if kind == "wa" and digits else f"tel:{val}"
items.append({"label": (label or "Telefoon")[:64], "value": val, "href": href, "kind": kind})
add("Organisatie", row.get("phone"))
for c in contacts or []:
who = _clean_str(c.get("name")) or _clean_str(c.get("role")) or "Contact"
add(who, c.get("phone"))
add(f"{who} mobiel", c.get("mobile"))
add(f"{who} WhatsApp", c.get("whatsapp"), "wa")
add(f"{who} fax", c.get("fax"))
return _dedupe_items(items)
def _collect_addresses(row: dict[str, Any], contacts: Optional[list[dict[str, Any]]] = None) -> list[dict[str, str]]:
items: list[dict[str, str]] = []
def add(label: str, line: Any, city: Any = None, postal: Any = None) -> None:
parts = [_clean_str(line), _clean_str(postal), _clean_str(city)]
parts = [p for p in parts if p]
if not parts:
return
val = ", ".join(parts)
items.append({"label": (label or "Adres")[:64], "value": val})
add("Locatie", row.get("address_line"), row.get("city"), row.get("postal_code"))
for c in contacts or []:
who = _clean_str(c.get("name")) or _clean_str(c.get("role")) or "Contact"
add(who, c.get("address_line"), c.get("city"), c.get("postal_code"))
return _dedupe_items(items)
def _collect_social(row: dict[str, Any], contacts: Optional[list[dict[str, Any]]] = None) -> list[dict[str, str]]:
items: list[dict[str, str]] = []
def add(label: str, raw: Any) -> None:
href = _normalize_url(raw)
if not href:
return
items.append({"label": label[:64], "value": _display_url(href), "href": href})
for c in contacts or []:
who = _clean_str(c.get("name")) or _clean_str(c.get("role")) or "Contact"
if c.get("linkedin_url"):
add(f"LinkedIn · {who}", c.get("linkedin_url"))
return _dedupe_items(items, key="href")
def _public_contacts(contacts: Optional[list[dict[str, Any]]], limit: int = 8) -> list[dict[str, Any]]:
out: list[dict[str, Any]] = []
for c in (contacts or [])[:limit]:
out.append({
"name": _clean_str(c.get("name")),
"role": _clean_str(c.get("role")),
"title": _clean_str(c.get("title")),
"department": _clean_str(c.get("department")),
"email": _clean_str(c.get("email")),
"email_alt": _clean_str(c.get("email_alt")),
"phone": _clean_str(c.get("phone")),
"mobile": _clean_str(c.get("mobile")),
"whatsapp": _clean_str(c.get("whatsapp")),
"fax": _clean_str(c.get("fax")),
"linkedin_url": _normalize_url(c.get("linkedin_url")),
"website": _normalize_url(c.get("website_contact_url")),
"is_primary": bool(c.get("is_primary")),
"source": _clean_str(c.get("source")),
})
return out
def _build_reach(
row: dict[str, Any],
contacts: Optional[list[dict[str, Any]]] = None,
brand_website: Optional[str] = None,
) -> dict[str, Any]:
lat, lon = row.get("lat"), row.get("lon")
maps_url = row.get("maps_url")
if not maps_url and lat is not None and lon is not None:
maps_url = f"https://www.google.com/maps?q={lat},{lon}"
place_id = _clean_str(row.get("google_place_id"))
if place_id and not maps_url:
maps_url = f"https://www.google.com/maps/search/?api=1&query_place_id={place_id}"
reach = {
"emails": _collect_emails(row, contacts),
"phones": _collect_phones(row, contacts),
"websites": _collect_websites(row, contacts, brand_website),
"addresses": _collect_addresses(row, contacts),
"social": _collect_social(row, contacts),
"maps_url": maps_url,
"google_place_id": place_id,
"postal_code": _clean_str(row.get("postal_code")),
"contact_people": _public_contacts(contacts),
}
reach["has_any"] = bool(
reach["emails"] or reach["phones"] or reach["websites"]
or reach["addresses"] or reach["social"] or reach["maps_url"]
)
return reach
def _map_contacts_batch(entity_ids: list[int]) -> dict[int, list[dict[str, Any]]]:
if not entity_ids:
return {}
rows = fetch_all(
"""
SELECT entity_id, role, name, title, department, email, email_alt,
phone, mobile, whatsapp, fax, linkedin_url, website_contact_url,
source_url, address_line, city, postal_code, is_primary, confidence, source
FROM export_entity_contacts
WHERE entity_id = ANY(%s)
ORDER BY is_primary DESC NULLS LAST, confidence DESC NULLS LAST, id
""",
(entity_ids,),
)
grouped: dict[int, list[dict[str, Any]]] = defaultdict(list)
for row in rows:
grouped[int(row["entity_id"])].append(row)
return grouped
def _brand_websites(codes: list[str]) -> dict[str, Optional[str]]:
if not codes:
return {}
rows = fetch_all(
"SELECT code, website FROM export_caterer_brands WHERE code = ANY(%s)",
(codes,),
)
return {r["code"]: r.get("website") for r in rows}
def _enrich_entity(row: dict[str, Any], contacts: list[dict[str, Any]]) -> dict[str, Any]:
row["contacts"] = contacts
row["contact_count"] = len(contacts)
for c in contacts:
if not row.get("email") and c.get("email"):
row["email"] = c["email"]
if not row.get("phone"):
phone = c.get("phone") or c.get("mobile")
if phone:
row["phone"] = phone
primary = contacts[0] if contacts else None
row["primary_contact"] = primary
brand_website = row.pop("brand_website", None)
row["reach"] = _build_reach(row, contacts, brand_website)
row["websites"] = row["reach"]["websites"]
if row["websites"] and not row.get("website"):
row["website"] = row["websites"][0]["url"]
if row["reach"]["emails"] and not row.get("email"):
row["email"] = row["reach"]["emails"][0]["value"]
if row["reach"]["phones"] and not row.get("phone"):
row["phone"] = row["reach"]["phones"][0]["value"]
h_score, h_reasons = score_entity(row)
row["halal_score"] = h_score
row["halal_reasons"] = h_reasons
row["halal_reason_labels"] = _halal_reason_labels(h_reasons)
if row.get("lat") and row.get("lon"):
row["maps_url"] = f"https://www.google.com/maps?q={row['lat']},{row['lon']}"
return row
from app.export_intel_sync import (
enrich_contacts_country,
materialize_caterer_presence,
sync_all_priority,
sync_caterers,
sync_contacts,
sync_customers,
sync_distributors,
sync_osm_country,
sync_ted_tenders,
sync_world_regions,
)
_MAP_CACHE: dict[str, tuple[float, dict[str, Any]]] = {}
_MAP_CACHE_TTL = 120.0
def _map_entity_limit(
country: Optional[str],
region: Optional[str],
q: Optional[str],
entity_type: Optional[str],
entity_types: Optional[str],
halal_min: Optional[float],
favorite_only: bool,
crm_linked: Optional[bool],
) -> int:
if q or entity_type or entity_types or halal_min is not None or favorite_only or crm_linked is not None:
return 8000
if country:
return 12000
if region:
return 8000
return 2500
def _map_cache_get(key: str) -> Optional[dict[str, Any]]:
hit = _MAP_CACHE.get(key)
if not hit:
return None
ts, payload = hit
if time.time() - ts > _MAP_CACHE_TTL:
_MAP_CACHE.pop(key, None)
return None
return payload
def _map_cache_set(key: str, payload: dict[str, Any]) -> None:
if len(_MAP_CACHE) > 48:
oldest = min(_MAP_CACHE.items(), key=lambda x: x[1][0])[0]
_MAP_CACHE.pop(oldest, None)
_MAP_CACHE[key] = (time.time(), payload)
router = APIRouter(prefix="/export-intel", tags=["export-intel"])
DISTRIBUTOR_TYPES = ("distributor", "wholesaler", "importer", "logistics", "cold_storage", "port_agent")
CUSTOMER_TYPES = ("restaurant", "doner_shoarma", "butcher", "foodservice")
CATERER_TYPES = ("contract_caterer")
TYPE_FILTER_ALIASES: dict[str, tuple[str, ...]] = {
"b2b": DISTRIBUTOR_TYPES,
"distributor": DISTRIBUTOR_TYPES,
"importer": DISTRIBUTOR_TYPES,
"logistics": DISTRIBUTOR_TYPES,
"horeca": CUSTOMER_TYPES,
"restaurant": CUSTOMER_TYPES,
"doner_shoarma": CUSTOMER_TYPES,
"butcher": CUSTOMER_TYPES,
"foodservice": CUSTOMER_TYPES,
"caterer": CATERER_TYPES,
"contract_caterer": CATERER_TYPES,
}
def _expand_entity_types(
entity_type: Optional[str],
entity_types: Optional[str],
) -> tuple[Optional[str], Optional[str]]:
expanded: list[str] = []
if entity_types:
for raw in entity_types.split(","):
key = raw.strip()
if not key:
continue
alias = TYPE_FILTER_ALIASES.get(key)
expanded.extend(alias if alias else (key,))
if entity_type:
key = entity_type.strip()
alias = TYPE_FILTER_ALIASES.get(key)
expanded.extend(alias if alias else (key,))
if not expanded:
return entity_type, entity_types
deduped: list[str] = []
seen: set[str] = set()
for key in expanded:
if key not in seen:
seen.add(key)
deduped.append(key)
if len(deduped) == 1:
return deduped[0], None
return None, ",".join(deduped)
class ContactCreateIn(BaseModel):
role: str = "general"
name: Optional[str] = None
email: Optional[str] = None
phone: Optional[str] = None
mobile: Optional[str] = None
is_primary: bool = False
class SyncRequestIn(BaseModel):
country_iso2: Optional[str] = None
region_code: Optional[str] = None
max_priority: int = 2
entity_types: Optional[list[str]] = None
class EntityIdsIn(BaseModel):
entity_ids: list[int]
class EntityUpdateIn(BaseModel):
name: Optional[str] = None
entity_type: Optional[str] = None
country_iso2: Optional[str] = None
city: Optional[str] = None
address_line: Optional[str] = None
postal_code: Optional[str] = None
phone: Optional[str] = None
email: Optional[str] = None
website: Optional[str] = None
volume_band: Optional[str] = None
pipeline_stage: Optional[str] = None
confidence: Optional[int] = Field(None, ge=0, le=100)
halal_cert_notes: Optional[str] = None
cold_chain: Optional[bool] = None
class FavoritesIn(BaseModel):
entity_ids: list[int]
favorite: bool = True
class CrmPushIn(BaseModel):
entity_ids: list[int]
create_deals: bool = True
client_stage: str = "intake"
deal_stage: str = "lead"
def _entity_filters(
country: Optional[str],
region: Optional[str],
entity_type: Optional[str],
entity_types: Optional[str],
favorite_only: bool = False,
crm_linked: Optional[bool] = None,
) -> tuple[str, list[Any]]:
clauses: list[str] = ["1=1"]
params: list[Any] = []
if country:
clauses.append("e.country_iso2 = %s")
params.append(country.upper())
if region:
clauses.append("t.region_code = %s")
params.append(region)
clauses.append("e.territory_code = t.code")
entity_type, entity_types = _expand_entity_types(entity_type, entity_types)
if entity_type:
clauses.append("e.entity_type = %s")
params.append(entity_type)
if entity_types:
types = [t.strip() for t in entity_types.split(",") if t.strip()]
if types:
clauses.append("e.entity_type = ANY(%s)")
params.append(types)
if favorite_only:
clauses.append("e.is_favorite = TRUE")
if crm_linked is True:
clauses.append("e.client_id IS NOT NULL")
elif crm_linked is False:
clauses.append("e.client_id IS NULL")
return " AND ".join(clauses), params
@router.get("/stats")
def export_stats(
country: Optional[str] = Query(None),
region: Optional[str] = Query(None),
) -> dict[str, Any]:
entity_where = ""
params: list[Any] = []
if country:
entity_where = " WHERE country_iso2 = %s"
params = [country.upper()]
elif region:
entity_where = " WHERE territory_code IN (SELECT code FROM export_territories WHERE region_code = %s)"
params = [region]
def count(sql_suffix: str, extra_params: tuple[Any, ...] = ()) -> int:
w = entity_where
if sql_suffix:
w = f"{entity_where} AND {sql_suffix}" if entity_where else f" WHERE {sql_suffix}"
row = fetch_one(
f"SELECT COUNT(*) AS n FROM export_market_entities{w}",
tuple(params) + extra_params if params or extra_params else None,
)
return int(row["n"] or 0)
entities_n = count("")
type_geo_rows = fetch_all(
"""
SELECT entity_type, COUNT(*) AS n
FROM export_market_entities
WHERE lat IS NOT NULL AND lon IS NOT NULL
GROUP BY entity_type
"""
)
entity_types_geo = {r["entity_type"]: int(r["n"] or 0) for r in type_geo_rows}
country_count_rows = fetch_all(
"""
SELECT country_iso2, COUNT(*) AS n
FROM export_market_entities
GROUP BY country_iso2
"""
)
country_entity_counts = {r["country_iso2"]: int(r["n"] or 0) for r in country_count_rows}
distributors_n = count("entity_type IN ('distributor','wholesaler','importer','logistics','cold_storage','port_agent')")
caterers_n = count("entity_type = 'contract_caterer'")
contacts = fetch_one("SELECT COUNT(*) AS n FROM export_entity_contacts")
with_email = fetch_one(
"SELECT COUNT(DISTINCT entity_id) AS n FROM export_entity_contacts WHERE email IS NOT NULL AND email <> ''"
)
territories = fetch_one("SELECT COUNT(*) AS n FROM export_territories WHERE is_active = TRUE")
presence = fetch_one("SELECT COUNT(*) AS n FROM export_caterer_presence WHERE is_active = TRUE")
tenders_open = 0
try:
row = fetch_one("SELECT COUNT(*) AS n FROM export_tenders WHERE status = 'open'")
tenders_open = int(row["n"] or 0) if row else 0
except Exception:
pass
fav_row = fetch_one("SELECT COUNT(*) AS n FROM export_market_entities WHERE is_favorite = TRUE")
crm_row = fetch_one("SELECT COUNT(*) AS n FROM export_market_entities WHERE client_id IS NOT NULL")
return {
"entities": entities_n,
"distributors": distributors_n,
"caterers": caterers_n,
"contacts": int(contacts["n"] or 0),
"entities_with_email": int(with_email["n"] or 0),
"territories": int(territories["n"] or 0),
"caterer_presence": int(presence["n"] or 0),
"tenders_open": tenders_open,
"favorites": int(fav_row["n"] or 0) if fav_row else 0,
"crm_linked": int(crm_row["n"] or 0) if crm_row else 0,
"phase": "B",
"enrichment_status": "active",
"entity_types_geo": entity_types_geo,
"country_entity_counts": country_entity_counts,
}
@router.get("/regions")
def list_regions() -> list[dict[str, Any]]:
return fetch_all("SELECT * FROM export_regions ORDER BY sort_order, code")
def _dedupe_territories(rows: list[dict[str, Any]]) -> list[dict[str, Any]]:
seen: set[str] = set()
out: list[dict[str, Any]] = []
for row in rows:
iso = row.get("country_iso2")
if iso and iso in seen:
continue
if iso:
seen.add(iso)
out.append(row)
return out
@router.get("/territories")
def list_territories(region: Optional[str] = Query(None)) -> list[dict[str, Any]]:
if region:
rows = fetch_all(
"SELECT * FROM export_territories WHERE is_active = TRUE AND region_code = %s ORDER BY sync_priority, name_nl",
(region,),
)
else:
rows = fetch_all(
"SELECT * FROM export_territories WHERE is_active = TRUE ORDER BY sync_priority, name_nl"
)
return _dedupe_territories(rows)
@router.get("/entities")
def list_entities(
country: Optional[str] = Query(None),
region: Optional[str] = Query(None),
entity_type: Optional[str] = Query(None),
entity_types: Optional[str] = Query(None),
q: Optional[str] = Query(None),
favorite_only: bool = Query(False),
crm_linked: Optional[bool] = Query(None),
limit: int = Query(200, ge=1, le=500),
offset: int = Query(0, ge=0),
) -> dict[str, Any]:
join = ""
if region:
join = "JOIN export_territories t ON e.territory_code = t.code"
where_sql, params = _entity_filters(
country, region, entity_type, entity_types, favorite_only, crm_linked
)
if q:
where_sql += " AND (e.name ILIKE %s OR e.city ILIKE %s)"
params.extend([f"%{q}%", f"%{q}%"])
rows = fetch_all(
f"""
SELECT e.*,
cl.name AS crm_client_name,
(SELECT COUNT(*) FROM export_entity_contacts ec WHERE ec.entity_id = e.id) AS contact_count,
(SELECT ec.email FROM export_entity_contacts ec WHERE ec.entity_id = e.id AND ec.is_primary = TRUE LIMIT 1) AS primary_email,
(SELECT ec.phone FROM export_entity_contacts ec WHERE ec.entity_id = e.id AND ec.is_primary = TRUE LIMIT 1) AS primary_phone
FROM export_market_entities e
LEFT JOIN clients cl ON cl.id = e.client_id
{join}
WHERE {where_sql}
ORDER BY e.confidence DESC, e.name
LIMIT %s OFFSET %s
""",
tuple(params) + (limit, offset),
)
total = fetch_one(
f"SELECT COUNT(*) AS n FROM export_market_entities e LEFT JOIN clients cl ON cl.id = e.client_id {join} WHERE {where_sql}",
tuple(params),
)
return {"items": rows, "total": int(total["n"] or 0), "limit": limit, "offset": offset}
@router.get("/entities/{entity_id}")
def get_entity(entity_id: int) -> dict[str, Any]:
row = fetch_one(
"""
SELECT e.*, cl.name AS crm_client_name, d.title AS crm_deal_title, d.stage AS crm_deal_stage,
b.website AS brand_website
FROM export_market_entities e
LEFT JOIN clients cl ON cl.id = e.client_id
LEFT JOIN deals d ON d.id = e.deal_id
LEFT JOIN export_caterer_brands b ON b.code = e.brand_code
WHERE e.id = %s
""",
(entity_id,),
)
if not row:
raise HTTPException(404, "Entity not found")
contacts = fetch_all(
"SELECT * FROM export_entity_contacts WHERE entity_id = %s ORDER BY is_primary DESC, confidence DESC",
(entity_id,),
)
return _enrich_entity(row, contacts)
@router.get("/contacts")
def list_contacts(
country: Optional[str] = Query(None),
region: Optional[str] = Query(None),
entity_type: Optional[str] = Query(None),
has_email: Optional[bool] = Query(None),
crm_linked: Optional[bool] = Query(None),
q: Optional[str] = Query(None),
limit: int = Query(200, ge=1, le=500),
offset: int = Query(0, ge=0),
) -> dict[str, Any]:
join = ""
if region:
join = "JOIN export_territories t ON e.territory_code = t.code"
where_sql, params = _entity_filters(country, region, entity_type, None, False, crm_linked)
if has_email is True:
where_sql += " AND c.email IS NOT NULL AND c.email <> ''"
elif has_email is False:
where_sql += " AND (c.email IS NULL OR c.email = '')"
if q:
like = f"%{q}%"
where_sql += (
" AND (e.name ILIKE %s OR e.city ILIKE %s OR c.email ILIKE %s"
" OR c.name ILIKE %s OR c.phone ILIKE %s OR c.mobile ILIKE %s)"
)
params.extend([like, like, like, like, like, like])
rows = fetch_all(
f"""
SELECT c.*, e.name AS entity_name, e.entity_type, e.country_iso2,
e.city AS entity_city, e.client_id, e.deal_id,
cl.name AS crm_client_name
FROM export_entity_contacts c
JOIN export_market_entities e ON e.id = c.entity_id
LEFT JOIN clients cl ON cl.id = e.client_id
{join}
WHERE {where_sql}
ORDER BY e.name, e.city NULLS LAST, c.is_primary DESC, c.email NULLS LAST
LIMIT %s OFFSET %s
""",
tuple(params) + (limit, offset),
)
total = fetch_one(
f"""
SELECT COUNT(*) AS n FROM export_entity_contacts c
JOIN export_market_entities e ON e.id = c.entity_id
{join}
WHERE {where_sql}
""",
tuple(params),
)
return {"items": rows, "total": int(total["n"] or 0), "limit": limit, "offset": offset}
@router.get("/contacts/export.csv")
def export_contacts_csv(
country: Optional[str] = Query(None),
region: Optional[str] = Query(None),
entity_type: Optional[str] = Query(None),
has_email: Optional[bool] = Query(None),
crm_linked: Optional[bool] = Query(None),
q: Optional[str] = Query(None),
) -> StreamingResponse:
data = list_contacts(
country=country,
region=region,
entity_type=entity_type,
has_email=has_email,
crm_linked=crm_linked,
q=q,
limit=5000,
offset=0,
)
buf = io.StringIO()
writer = csv.writer(buf)
writer.writerow(
["entity_name", "entity_type", "country", "role", "name", "email", "phone", "mobile", "source", "confidence"]
)
for r in data["items"]:
writer.writerow(
[
r.get("entity_name"),
r.get("entity_type"),
r.get("country_iso2"),
r.get("role"),
r.get("name"),
r.get("email"),
r.get("phone"),
r.get("mobile"),
r.get("source"),
r.get("confidence"),
]
)
buf.seek(0)
return StreamingResponse(
iter([buf.getvalue()]),
media_type="text/csv",
headers={"Content-Disposition": "attachment; filename=export-intel-contacts.csv"},
)
@router.patch("/entities/{entity_id}")
def update_entity(entity_id: int, payload: EntityUpdateIn) -> dict[str, Any]:
row = fetch_one("SELECT id FROM export_market_entities WHERE id = %s", (entity_id,))
if not row:
raise HTTPException(status_code=404, detail="Entity not found")
data = payload.model_dump(exclude_unset=True)
if not data:
raise HTTPException(status_code=400, detail="Geen velden om bij te werken")
if "country_iso2" in data and data["country_iso2"]:
data["country_iso2"] = data["country_iso2"].upper()[:2]
sets = [f"{k} = %s" for k in data]
sets.append("updated_at = NOW()")
params = list(data.values()) + [entity_id]
execute(
f"UPDATE export_market_entities SET {', '.join(sets)} WHERE id = %s",
tuple(params),
)
_MAP_CACHE.clear()
return get_entity(entity_id)
@router.post("/entities/{entity_id}/contacts")
def add_contact(entity_id: int, payload: ContactCreateIn) -> dict[str, Any]:
ent = fetch_one("SELECT id FROM export_market_entities WHERE id = %s", (entity_id,))
if not ent:
raise HTTPException(404, "Entity not found")
row = execute_returning(
"""
INSERT INTO export_entity_contacts (entity_id, role, name, email, phone, mobile, is_primary, source, confidence)
VALUES (%s, %s, %s, %s, %s, %s, %s, 'manual', 80)
RETURNING *
""",
(
entity_id,
payload.role,
payload.name,
payload.email,
payload.phone,
payload.mobile,
payload.is_primary,
),
)
return dict(row)
@router.get("/caterers/brands")
def list_caterer_brands() -> list[dict[str, Any]]:
return fetch_all(
"""
SELECT b.*,
(SELECT COUNT(*) FROM export_caterer_presence p WHERE p.brand_code = b.code AND p.is_active) AS country_count
FROM export_caterer_brands b
ORDER BY b.tier, b.name
"""
)
@router.get("/caterers/presence")
def list_caterer_presence(country: Optional[str] = Query(None)) -> list[dict[str, Any]]:
if country:
return fetch_all(
"""
SELECT p.*, b.name AS brand_name, b.tier, b.website AS brand_website
FROM export_caterer_presence p
JOIN export_caterer_brands b ON b.code = p.brand_code
WHERE p.country_iso2 = %s AND p.is_active = TRUE
ORDER BY b.tier, b.name
""",
(country.upper(),),
)
return fetch_all(
"""
SELECT p.*, b.name AS brand_name, b.tier, b.website AS brand_website
FROM export_caterer_presence p
JOIN export_caterer_brands b ON b.code = p.brand_code
WHERE p.is_active = TRUE
ORDER BY p.country_iso2, b.tier, b.name
"""
)
@router.get("/gov-sources")
def list_gov_sources(country: Optional[str] = Query(None)) -> list[dict[str, Any]]:
if country:
return fetch_all(
"SELECT * FROM export_gov_data_sources WHERE is_active = TRUE AND country_iso2 = %s ORDER BY category, name",
(country.upper(),),
)
return fetch_all(
"SELECT * FROM export_gov_data_sources WHERE is_active = TRUE ORDER BY country_iso2, category, name"
)
@router.get("/tenders")
def list_tenders(
country: Optional[str] = Query(None),
status: Optional[str] = Query("open"),
limit: int = Query(100, ge=1, le=200),
) -> dict[str, Any]:
clauses = ["1=1"]
params: list[Any] = []
if country:
clauses.append("country_iso2 = %s")
params.append(country.upper())
if status:
clauses.append("status = %s")
params.append(status)
where = " AND ".join(clauses)
rows = fetch_all(
f"SELECT * FROM export_tenders WHERE {where} ORDER BY created_at DESC LIMIT %s",
tuple(params) + (limit,),
)
total = fetch_one(f"SELECT COUNT(*) AS n FROM export_tenders WHERE {where}", tuple(params))
return {"items": rows, "total": int(total["n"] or 0)}
@router.get("/map/bundle")
def map_bundle(
country: Optional[str] = Query(None),
region: Optional[str] = Query(None),
entity_type: Optional[str] = Query(None),
entity_types: Optional[str] = Query(None),
q: Optional[str] = Query(None),
halal_min: Optional[float] = Query(None, ge=0, le=100),
favorite_only: bool = Query(False),
crm_linked: Optional[bool] = Query(None),
) -> dict[str, Any]:
join = ""
if region:
join = "JOIN export_territories t ON e.territory_code = t.code"
filter_sql, params = _entity_filters(
country, region, entity_type, entity_types, favorite_only, crm_linked
)
where = f"e.lat IS NOT NULL AND e.lon IS NOT NULL AND {filter_sql}"
if q:
where += " AND (e.name ILIKE %s OR e.city ILIKE %s)"
params.extend([f"%{q}%", f"%{q}%"])
row_limit = _map_entity_limit(
country, region, q, entity_type, entity_types, halal_min, favorite_only, crm_linked
)
cache_key = hashlib.md5(
json.dumps(
{
"country": country,
"region": region,
"entity_type": entity_type,
"entity_types": entity_types,
"q": q,
"halal_min": halal_min,
"favorite_only": favorite_only,
"crm_linked": crm_linked,
"limit": row_limit,
},
sort_keys=True,
default=str,
).encode()
).hexdigest()
cached = _map_cache_get(cache_key)
if cached is not None:
return cached
entities = fetch_all(
f"""
SELECT e.id, e.name, e.entity_type, e.country_iso2, e.lat, e.lon,
e.confidence, e.pipeline_stage, e.city, e.address_line, e.postal_code,
e.email, e.phone, e.website, e.volume_band, e.source,
e.brand_code, e.is_favorite, e.client_id, e.deal_id,
cp.contact_email, cp.contact_phone,
COALESCE(cc.contact_count, 0) AS contact_count
FROM export_market_entities e
{join}
LEFT JOIN LATERAL (
SELECT c.email AS contact_email,
COALESCE(NULLIF(c.phone, ''), NULLIF(c.mobile, '')) AS contact_phone
FROM export_entity_contacts c
WHERE c.entity_id = e.id
ORDER BY c.is_primary DESC NULLS LAST, c.id
LIMIT 1
) cp ON TRUE
LEFT JOIN (
SELECT entity_id, COUNT(*)::int AS contact_count
FROM export_entity_contacts
GROUP BY entity_id
) cc ON cc.entity_id = e.id
WHERE {where}
ORDER BY e.confidence DESC NULLS LAST, e.id
LIMIT %s
""",
(tuple(params) + (row_limit,)) if params else (row_limit,),
)
entity_ids = [int(r["id"]) for r in entities]
contacts_by_entity = _map_contacts_batch(entity_ids)
brand_codes = list({r["brand_code"] for r in entities if r.get("brand_code")})
brands_by_code = _brand_websites(brand_codes)
scored: list[tuple[dict[str, Any], float, list[str]]] = []
halal_by_country: dict[str, dict[str, Any]] = {}
for row in entities:
if row.get("lat") is None or row.get("lon") is None:
continue
h_score, h_reasons = score_entity(row)
iso = row.get("country_iso2") or ""
bucket = halal_by_country.setdefault(
iso,
{"country_iso2": iso, "market_score": COUNTRY_HALAL_MARKET.get(iso, 40), "entity_count": 0, "avg_halal_score": 0.0, "hot_count": 0},
)
bucket["entity_count"] += 1
bucket["avg_halal_score"] += h_score
if h_score >= 75:
bucket["hot_count"] += 1
scored.append((row, h_score, h_reasons))
for bucket in halal_by_country.values():
n = bucket["entity_count"] or 1
bucket["avg_halal_score"] = round(bucket["avg_halal_score"] / n, 1)
bucket["combined_score"] = round(bucket["market_score"] * 0.4 + bucket["avg_halal_score"] * 0.6, 1)
if halal_min is not None:
scored = [s for s in scored if s[1] >= halal_min]
scored.sort(key=lambda x: x[1], reverse=True)
features = [
{
"type": "Feature",
"geometry": {"type": "Point", "coordinates": [row["lon"], row["lat"]]},
"properties": {
"id": row["id"],
"name": row["name"],
"entity_type": row["entity_type"],
"country_iso2": row["country_iso2"],
"city": row.get("city"),
"address_line": row.get("address_line"),
"email": row.get("email") or row.get("contact_email"),
"phone": row.get("phone") or row.get("contact_phone"),
"website": row.get("website"),
"reach": _build_reach(
row,
contacts_by_entity.get(int(row["id"]), []),
brands_by_code.get(row.get("brand_code")),
),
"websites": _build_reach(
row,
contacts_by_entity.get(int(row["id"]), []),
brands_by_code.get(row.get("brand_code")),
)["websites"],
"volume_band": row.get("volume_band"),
"confidence": row.get("confidence"),
"pipeline_stage": row.get("pipeline_stage"),
"contact_count": int(row.get("contact_count") or 0),
"halal_score": h_score,
"halal_tier": halal_pin_tier(h_score),
},
}
for row, h_score, h_reasons in scored
]
top_markets = sorted(halal_by_country.values(), key=lambda x: x["combined_score"], reverse=True)[:15]
top_entities = [
{
"id": row["id"],
"name": row["name"],
"entity_type": row["entity_type"],
"country_iso2": row["country_iso2"],
"city": row.get("city"),
"address_line": row.get("address_line"),
"email": row.get("email") or row.get("contact_email"),
"phone": row.get("phone") or row.get("contact_phone"),
"website": row.get("website"),
"halal_score": h_score,
"halal_reason_labels": _halal_reason_labels(h_reasons),
"is_favorite": bool(row.get("is_favorite")),
"client_id": row.get("client_id"),
"has_contact": bool(
row.get("email") or row.get("contact_email")
or row.get("phone") or row.get("contact_phone")
or int(row.get("contact_count") or 0) > 0
),
}
for row, h_score, h_reasons in scored[:25]
]
density_rows = fetch_all(
"""
SELECT country_iso2, COUNT(*) AS n
FROM export_market_entities
WHERE lat IS NOT NULL AND lon IS NOT NULL
GROUP BY country_iso2
"""
)
trade_choropleth = {r["country_iso2"]: r["n"] for r in density_rows}
sidebar_entities = [
{
"id": row["id"],
"name": row["name"],
"entity_type": row["entity_type"],
"country_iso2": row["country_iso2"],
"city": row.get("city"),
"email": row.get("email") or row.get("contact_email"),
"phone": row.get("phone") or row.get("contact_phone"),
"halal_score": h_score,
"client_id": row.get("client_id"),
"is_favorite": bool(row.get("is_favorite")),
"has_contact": bool(
row.get("email") or row.get("contact_email")
or row.get("phone") or row.get("contact_phone")
or int(row.get("contact_count") or 0) > 0
),
}
for row, h_score, h_reasons in scored[:200]
]
payload = {
"entities": {"type": "FeatureCollection", "features": features},
"trade_choropleth": trade_choropleth,
"meta": {
"entity_count": len(features),
"map_limit": row_limit,
"truncated": len(scored) >= row_limit,
"countries_with_trade": len(trade_choropleth),
"halal_min": halal_min,
"top_markets": top_markets,
"top_halal_entities": top_entities,
"sidebar_entities": sidebar_entities,
"phase": "B",
"note": "Halal-score: type + land + naam/signaal + contact",
},
}
_map_cache_set(cache_key, payload)
return payload
@router.get("/halal/markets")
def halal_markets(
region: Optional[str] = Query(None),
country: Optional[str] = Query(None),
limit: int = Query(30, ge=1, le=100),
) -> dict[str, Any]:
join = ""
where = "e.lat IS NOT NULL"
params: list[Any] = []
if region:
join = "JOIN export_territories t ON e.territory_code = t.code"
where += " AND t.region_code = %s"
params.append(region)
if country:
where += " AND e.country_iso2 = %s"
params.append(country.upper())
rows = fetch_all(
f"""
SELECT e.id, e.name, e.entity_type, e.country_iso2, e.city,
e.confidence, e.product_interest, e.halal_cert_notes, e.email,
(SELECT COUNT(*) FROM export_entity_contacts c WHERE c.entity_id = e.id) AS contact_count
FROM export_market_entities e
{join}
WHERE {where}
LIMIT 15000
""",
tuple(params) if params else None,
)
by_country: dict[str, dict[str, Any]] = {}
top_entities: list[dict[str, Any]] = []
for row in rows:
h_score, reasons = score_entity(row)
iso = row.get("country_iso2") or ""
bucket = by_country.setdefault(
iso,
{"country_iso2": iso, "market_score": COUNTRY_HALAL_MARKET.get(iso, 40), "entities": 0, "score_sum": 0.0, "hot": 0},
)
bucket["entities"] += 1
bucket["score_sum"] += h_score
if h_score >= 75:
bucket["hot"] += 1
top_entities.append({
"id": row["id"],
"name": row["name"],
"entity_type": row["entity_type"],
"country_iso2": iso,
"city": row.get("city"),
"halal_score": h_score,
"halal_reasons": reasons[:3],
})
markets = []
for iso, b in by_country.items():
n = b["entities"] or 1
avg = round(b["score_sum"] / n, 1)
markets.append({
"country_iso2": iso,
"market_score": b["market_score"],
"entity_count": b["entities"],
"avg_halal_score": avg,
"hot_entities": b["hot"],
"combined_score": round(b["market_score"] * 0.4 + avg * 0.6, 1),
})
markets.sort(key=lambda x: x["combined_score"], reverse=True)
top_entities.sort(key=lambda x: x["halal_score"], reverse=True)
return {
"markets": markets[:limit],
"top_entities": top_entities[:50],
"total_entities": len(rows),
}
@router.post("/entities/favorites")
def set_favorites(payload: FavoritesIn) -> dict[str, Any]:
if not payload.entity_ids:
raise HTTPException(400, "entity_ids required")
fav_val = payload.favorite
execute(
"""
UPDATE export_market_entities
SET is_favorite = %s,
favorited_at = CASE WHEN %s THEN NOW() ELSE NULL END,
updated_at = NOW()
WHERE id = ANY(%s)
""",
(fav_val, fav_val, payload.entity_ids),
)
return {"ok": True, "updated": len(payload.entity_ids), "favorite": fav_val}
@router.get("/entities/favorites")
def list_favorite_entities(
limit: int = Query(200, ge=1, le=500),
country: Optional[str] = Query(None),
region: Optional[str] = Query(None),
) -> dict[str, Any]:
return list_entities(
country=country,
region=region,
favorite_only=True,
limit=limit,
offset=0,
)
@router.post("/crm/push")
def crm_push(payload: CrmPushIn) -> dict[str, Any]:
if not payload.entity_ids:
raise HTTPException(400, "entity_ids required")
if len(payload.entity_ids) > 100:
raise HTTPException(400, "max 100 entities per request")
return push_entities_to_crm(
payload.entity_ids,
create_deals=payload.create_deals,
client_stage=payload.client_stage,
deal_stage=payload.deal_stage,
)
@router.get("/crm/pipeline")
def crm_pipeline(
limit: int = Query(100, ge=1, le=500),
country: Optional[str] = Query(None),
region: Optional[str] = Query(None),
) -> dict[str, Any]:
join = ""
if region:
join = "JOIN export_territories t ON e.territory_code = t.code"
where_sql, params = _entity_filters(country, region, None, None, False, True)
rows = fetch_all(
f"""
SELECT e.id, e.name, e.entity_type, e.country_iso2, e.city,
e.pipeline_stage, e.crm_pushed_at, e.is_favorite,
e.client_id, e.deal_id, e.email, e.phone,
cl.name AS crm_client_name, cl.stage AS crm_client_stage,
d.title AS crm_deal_title, d.stage AS crm_deal_stage
FROM export_market_entities e
LEFT JOIN clients cl ON cl.id = e.client_id
LEFT JOIN deals d ON d.id = e.deal_id
{join}
WHERE {where_sql}
ORDER BY e.crm_pushed_at DESC NULLS LAST, e.updated_at DESC
LIMIT %s
""",
tuple(params) + (limit,),
)
return {"items": rows, "total": len(rows)}
@router.post("/sync/contacts")
def api_sync_contacts(payload: SyncRequestIn) -> dict[str, Any]:
try:
result = sync_contacts(payload.country_iso2, payload.region_code)
return {"status": "completed", "phase": "B", "result": result}
except Exception as exc:
raise HTTPException(500, str(exc)) from exc
@router.post("/sync/caterers")
def api_sync_caterers(payload: SyncRequestIn) -> dict[str, Any]:
try:
result = sync_caterers(payload.country_iso2, payload.region_code)
return {"status": "completed", "phase": "B", "result": result}
except Exception as exc:
raise HTTPException(500, str(exc)) from exc
@router.post("/sync/distributors")
def api_sync_distributors(payload: SyncRequestIn) -> dict[str, Any]:
try:
result = sync_distributors(payload.country_iso2, payload.region_code)
return {"status": "completed", "phase": "B", "result": result}
except Exception as exc:
raise HTTPException(500, str(exc)) from exc
@router.post("/sync/customers")
def api_sync_customers(payload: SyncRequestIn) -> dict[str, Any]:
if not payload.country_iso2 and not payload.region_code:
raise HTTPException(400, "country_iso2 or region_code required")
try:
result = sync_customers(payload.country_iso2, payload.region_code)
return {"status": "completed", "phase": "B", "result": result}
except Exception as exc:
raise HTTPException(500, str(exc)) from exc
@router.post("/sync/tenders")
def api_sync_tenders() -> dict[str, Any]:
try:
result = sync_ted_tenders()
return {"status": "completed", "phase": "B", "result": result}
except Exception as exc:
raise HTTPException(500, str(exc)) from exc
@router.post("/sync/all")
def api_sync_all(payload: SyncRequestIn = SyncRequestIn()) -> dict[str, Any]:
try:
result = sync_all_priority(max_priority=payload.max_priority, region_code=payload.region_code)
return {"status": "completed", "phase": "B", "result": result}
except Exception as exc:
raise HTTPException(500, str(exc)) from exc
@router.post("/sync/missing")
def api_sync_missing(payload: SyncRequestIn = SyncRequestIn()) -> dict[str, Any]:
try:
from app.export_intel_sync import sync_missing_countries
result = sync_missing_countries(max_priority=payload.max_priority or 3)
return {"status": "completed", "phase": "B", "result": result}
except Exception as exc:
raise HTTPException(500, str(exc)) from exc
class ExportIntelSettingsIn(BaseModel):
google_maps_api_key: Optional[str] = None
@router.get("/settings")
def export_intel_settings_get() -> dict[str, Any]:
from app.export_intel_places import google_api_key_status
return google_api_key_status()
@router.put("/settings")
def export_intel_settings_put(payload: ExportIntelSettingsIn) -> dict[str, Any]:
from app.export_intel_places import google_api_key_status, set_google_api_key
if payload.google_maps_api_key is not None:
set_google_api_key(payload.google_maps_api_key)
return google_api_key_status()
@router.post("/sync/places")
def api_sync_places(payload: SyncRequestIn) -> dict[str, Any]:
from app.export_intel_places import enrich_entities_places_country, sync_places_country
from app.export_intel_sync import _countries_for_sync
from app.export_intel_sync_progress import finish_job, job_step, start_job, set_active_job
if payload.country_iso2:
job_id = start_job("places", f"Google Places · {payload.country_iso2.upper()}", total=1)
try:
job_step(job_id, message=f"Google Places {payload.country_iso2.upper()}", current=1, total=1)
result = sync_places_country(payload.country_iso2.upper())
finish_job(job_id, result=result)
return {"status": "completed", "phase": "B", "result": result}
except Exception as exc:
finish_job(job_id, error=str(exc))
raise HTTPException(500, str(exc)) from exc
finally:
set_active_job(None)
countries = _countries_for_sync(region_code=payload.region_code, max_priority=payload.max_priority or 3)
if not countries:
raise HTTPException(400, "country_iso2 or region_code required")
job_id = start_job("places", f"Google Places · {len(countries)} landen", total=len(countries))
results: dict[str, Any] = {}
try:
from app.export_intel_sync_progress import SyncCancelled, raise_if_cancelled
for i, c in enumerate(countries, start=1):
raise_if_cancelled()
job_step(job_id, message=f"Google Places {i}/{len(countries)}: {c}", country=c, current=i, total=len(countries))
results[c] = sync_places_country(c)
time.sleep(1)
finish_job(job_id, result={"countries": len(countries)})
except Exception as exc:
finish_job(job_id, error=str(exc))
raise HTTPException(500, str(exc)) from exc
finally:
set_active_job(None)
return {"status": "completed", "phase": "B", "result": results}
@router.post("/sync/cancel")
def api_sync_cancel() -> dict[str, Any]:
"""Request cancellation of the currently running sync job(s)."""
from app.export_intel_sync_progress import request_cancel, get_monitor_snapshot
out = request_cancel("user")
snap = get_monitor_snapshot()
out["any_running"] = snap.get("any_running")
out["active_jobs"] = [{"id": j.get("id"), "label": j.get("label"), "status": j.get("status")} for j in snap.get("active_jobs", [])]
return out
@router.get("/sync/status")
def sync_status() -> dict[str, Any]:
from app.export_intel_sync_progress import get_monitor_snapshot
total_row = fetch_one("SELECT COUNT(*) AS n FROM export_market_entities")
terr_row = fetch_one("SELECT COUNT(*) AS n FROM export_territories WHERE is_active = TRUE")
with_data = fetch_one(
"""
SELECT COUNT(DISTINCT country_iso2) AS n FROM export_market_entities
WHERE country_iso2 IS NOT NULL
"""
)
zero_rows = fetch_all(
"""
SELECT t.country_iso2, t.name_nl, t.region_code
FROM export_territories t
WHERE t.is_active
AND NOT EXISTS (
SELECT 1 FROM export_market_entities e WHERE e.country_iso2 = t.country_iso2
)
ORDER BY t.region_code, t.name_nl
"""
)
type_rows = fetch_all(
"""
SELECT entity_type, COUNT(*) AS n FROM export_market_entities GROUP BY entity_type ORDER BY n DESC
"""
)
coverage = {
"entities_total": int(total_row["n"] or 0) if total_row else 0,
"territories_total": int(terr_row["n"] or 0) if terr_row else 0,
"countries_with_data": int(with_data["n"] or 0) if with_data else 0,
"countries_empty": len(zero_rows),
"countries_empty_list": zero_rows[:30],
"entity_types": {r["entity_type"]: int(r["n"] or 0) for r in type_rows},
}
snap = get_monitor_snapshot(coverage)
return snap
@router.post("/sync/world")
def api_sync_world(payload: SyncRequestIn = SyncRequestIn()) -> dict[str, Any]:
"""Sync Europe, Middle East, Africa, Americas — long running."""
try:
result = sync_world_regions(max_priority=payload.max_priority)
return {"status": "completed", "phase": "B", "result": result}
except Exception as exc:
raise HTTPException(500, str(exc)) from exc
@router.post("/sync/region")
def api_sync_region(payload: SyncRequestIn) -> dict[str, Any]:
if not payload.region_code:
raise HTTPException(400, "region_code required")
try:
result = sync_all_priority(max_priority=payload.max_priority, region_code=payload.region_code)
return {"status": "completed", "phase": "B", "result": result}
except Exception as exc:
raise HTTPException(500, str(exc)) from exc