SysOps: manual — 2026-07-30 09:33 UTC

This commit is contained in:
sysops
2026-07-30 09:33:03 +00:00
parent f21d58483b
commit 5dd69a1138
25 changed files with 4148 additions and 734 deletions
+462 -9
View File
@@ -6,6 +6,7 @@ import io
import hashlib
import json
import time
from collections import defaultdict
from typing import Any, Optional
from fastapi import APIRouter, HTTPException, Query
@@ -38,6 +39,241 @@ def _halal_reason_labels(reasons: list[str]) -> list[str]:
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)
@@ -50,6 +286,15 @@ def _enrich_entity(row: dict[str, Any], contacts: list[dict[str, Any]]) -> dict[
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
@@ -87,7 +332,7 @@ def _map_entity_limit(
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 4000
return 8000
if country:
return 12000
if region:
@@ -118,6 +363,49 @@ DISTRIBUTOR_TYPES = ("distributor", "wholesaler", "importer", "logistics", "cold
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"
@@ -186,6 +474,7 @@ def _entity_filters(
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)
@@ -228,7 +517,24 @@ def export_stats(
return int(row["n"] or 0)
entities_n = count("")
distributors_n = count("entity_type IN ('distributor','wholesaler','importer','logistics')")
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(
@@ -259,6 +565,8 @@ def export_stats(
"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,
}
@@ -267,16 +575,31 @@ 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:
return fetch_all(
rows = fetch_all(
"SELECT * FROM export_territories WHERE is_active = TRUE AND region_code = %s ORDER BY sync_priority, name_nl",
(region,),
)
return fetch_all(
"SELECT * FROM export_territories WHERE is_active = TRUE ORDER BY sync_priority, name_nl"
)
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")
@@ -328,10 +651,12 @@ def list_entities(
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
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,),
@@ -614,9 +939,9 @@ def map_bundle(
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.confidence, e.pipeline_stage, e.city, e.address_line, e.postal_code,
e.email, e.phone, e.website, e.volume_band, e.source,
e.is_favorite, e.client_id, e.deal_id,
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
@@ -641,6 +966,11 @@ def map_bundle(
(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:
@@ -678,8 +1008,23 @@ def map_bundle(
"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),
@@ -971,6 +1316,114 @@ def api_sync_all(payload: SyncRequestIn = SyncRequestIn()) -> dict[str, Any]:
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:
for i, c in enumerate(countries, start=1):
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.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."""
+24
View File
@@ -0,0 +1,24 @@
#!/usr/bin/env python3
"""Background sync with live progress tracking."""
import time
from app.export_intel_sync import sync_customers, sync_missing_countries
from app.export_intel_sync_progress import finish_job, job_step, start_job
REGIONS = ["europe", "middle_east", "gcc", "africa", "americas", "caribbean", "asia", "oceania"]
def main():
job_id = start_job("background", "Volledige wereld-update", total=len(REGIONS) + 1)
try:
job_step(job_id, message="Lege landen synchroniseren…", step="missing", current=0, total=len(REGIONS) + 1)
sync_missing_countries(max_priority=3)
for i, region in enumerate(REGIONS, start=1):
job_step(job_id, message=f"Horeca regio {i}/{len(REGIONS)}: {region}", step="horeca", current=i, total=len(REGIONS) + 1)
sync_customers(region_code=region)
time.sleep(2)
finish_job(job_id, result={"regions": len(REGIONS)})
except Exception as exc:
finish_job(job_id, error=str(exc))
raise
if __name__ == "__main__":
main()
+216 -41
View File
@@ -1,20 +1,26 @@
"""Google Places sync for Export Intel (optional — requires GOOGLE_MAPS_API_KEY)."""
"""Google Places sync + entity enrichment for Export Intel."""
from __future__ import annotations
import json
import math
import os
import time
from typing import Any, Optional
import httpx
from app.db import fetch_one
from app.db import execute, fetch_all, fetch_one
from app.export_intel_sync import _get_territory, _upsert_contact
PLACES_URL = "https://maps.googleapis.com/maps/api/place/textsearch/json"
DETAILS_URL = "https://maps.googleapis.com/maps/api/place/details/json"
FIND_URL = "https://maps.googleapis.com/maps/api/place/findplacefromtext/json"
SETTINGS_KEY = "export_intel.google_maps_api_key"
QUERIES = [
"groothandel groente fruit",
"wholesale food",
"meat wholesaler",
"halal meat distributor",
"poultry distributor",
@@ -22,13 +28,58 @@ QUERIES = [
"contract catering",
"food importer",
"halal butcher",
"horeca groothandel",
]
LARGE_COUNTRIES = frozenset({"US", "CA", "BR", "MX", "AR", "AU", "RU", "CN", "IN", "CD", "DZ", "EG", "IR", "SA"})
def get_google_api_key() -> Optional[str]:
env = (os.getenv("GOOGLE_MAPS_API_KEY") or os.getenv("GOOGLE_PLACES_API_KEY") or "").strip()
if env:
return env
row = fetch_one("SELECT value FROM app_settings WHERE key = %s", (SETTINGS_KEY,))
if not row or not row.get("value"):
return None
val = row["value"]
if isinstance(val, dict):
key = (val.get("api_key") or val.get("key") or "").strip()
return key or None
if isinstance(val, str):
try:
parsed = json.loads(val)
if isinstance(parsed, dict):
return (parsed.get("api_key") or parsed.get("key") or "").strip() or None
except Exception:
return val.strip() or None
return None
def set_google_api_key(api_key: str) -> None:
key = (api_key or "").strip()
if not key:
execute("DELETE FROM app_settings WHERE key = %s", (SETTINGS_KEY,))
return
execute(
"""
INSERT INTO app_settings (key, value, updated_at)
VALUES (%s, %s::jsonb, NOW())
ON CONFLICT (key) DO UPDATE SET value = EXCLUDED.value, updated_at = NOW()
""",
(SETTINGS_KEY, json.dumps({"api_key": key})),
)
def google_api_key_status() -> dict[str, Any]:
key = get_google_api_key()
return {
"google_maps_api_key_set": bool(key),
"source": "env" if (os.getenv("GOOGLE_MAPS_API_KEY") or os.getenv("GOOGLE_PLACES_API_KEY")) else ("db" if key else None),
}
def _api_key() -> Optional[str]:
return os.getenv("GOOGLE_MAPS_API_KEY") or os.getenv("GOOGLE_PLACES_API_KEY")
return get_google_api_key()
def _places_search(query: str, lat: float, lon: float, radius_m: int = 80000) -> list[dict[str, Any]]:
@@ -51,6 +102,46 @@ def _places_search(query: str, lat: float, lon: float, radius_m: int = 80000) ->
return data.get("results", [])[:12]
def _find_place_id(name: str, lat: float, lon: float, radius_m: int = 3000) -> Optional[str]:
key = _api_key()
if not key or not name.strip():
return None
params = {
"input": name.strip()[:128],
"inputtype": "textquery",
"locationbias": f"circle:{radius_m}@{lat},{lon}",
"fields": "place_id,name,geometry",
"key": key,
}
with httpx.Client(timeout=20.0) as client:
resp = client.get(FIND_URL, params=params)
if resp.status_code >= 400:
return None
data = resp.json()
if data.get("status") != "OK":
return None
cands = data.get("candidates") or []
if not cands:
return None
best = cands[0]
geo = best.get("geometry", {}).get("location", {})
clat, clon = geo.get("lat"), geo.get("lng")
if clat is not None and clon is not None:
dist = _haversine_m(lat, lon, float(clat), float(clon))
if dist > max(radius_m * 2, 8000):
return None
return best.get("place_id")
def _haversine_m(lat1: float, lon1: float, lat2: float, lon2: float) -> float:
r = 6371000.0
p1, p2 = math.radians(lat1), math.radians(lat2)
dp = math.radians(lat2 - lat1)
dl = math.radians(lon2 - lon1)
a = math.sin(dp / 2) ** 2 + math.cos(p1) * math.cos(p2) * math.sin(dl / 2) ** 2
return 2 * r * math.asin(math.sqrt(a))
def _place_details(place_id: str) -> dict[str, Any]:
key = _api_key()
if not key:
@@ -60,7 +151,7 @@ def _place_details(place_id: str) -> dict[str, Any]:
DETAILS_URL,
params={
"place_id": place_id,
"fields": "name,formatted_phone_number,website,formatted_address,geometry,international_phone_number",
"fields": "name,formatted_phone_number,website,formatted_address,geometry,international_phone_number,place_id",
"key": key,
},
)
@@ -84,55 +175,132 @@ def _bbox_points(bbox: Any, grid: int = 3) -> list[tuple[float, float]]:
return points
def _upsert_place(place: dict[str, Any], details: dict[str, Any], country_iso2: str, territory_code: str, query: str) -> tuple[str, int]:
from app.db import execute, execute_returning
def _apply_place_to_entity(entity_id: int, place_id: str, details: dict[str, Any]) -> bool:
phone = details.get("formatted_phone_number") or details.get("international_phone_number")
website = details.get("website")
geo = details.get("geometry", {}).get("location", {})
plat, plon = geo.get("lat"), geo.get("lng")
execute(
"""
UPDATE export_market_entities SET
phone = COALESCE(%s, phone),
website = COALESCE(%s, website),
google_place_id = COALESCE(google_place_id, %s),
lat = COALESCE(%s, lat),
lon = COALESCE(%s, lon),
confidence = GREATEST(confidence, 68),
updated_at = NOW(),
metadata = COALESCE(metadata, '{}'::jsonb) || %s::jsonb
WHERE id = %s
""",
(
phone,
website,
place_id,
plat,
plon,
json.dumps({"google_place_id": place_id, "google_enriched": True}),
entity_id,
),
)
if phone:
_upsert_contact(entity_id, None, phone, "google_places", 72, role="sales", name="Google Places")
return bool(phone or website)
pid = place.get("place_id")
def _upsert_place(place: dict[str, Any], details: dict[str, Any], country_iso2: str, territory_code: str, query: str) -> tuple[str, int]:
from app.db import execute_returning
pid = place.get("place_id") or details.get("place_id")
if not pid:
return "skipped", 0
name = details.get("name") or place.get("name") or "Unknown"
geo = details.get("geometry", {}).get("location", place.get("geometry", {}).get("location", {}))
plat = geo.get("lat")
plon = geo.get("lng")
phone = details.get("formatted_phone_number") or details.get("international_phone_number")
website = details.get("website")
entity_type = "contract_caterer" if "catering" in query else "distributor"
entity_type = "contract_caterer" if "catering" in query else "wholesaler"
existing = fetch_one(
"SELECT id FROM export_market_entities WHERE metadata->>'google_place_id' = %s",
"SELECT id FROM export_market_entities WHERE google_place_id = %s",
(pid,),
)
if existing:
execute(
"""
UPDATE export_market_entities SET
phone = COALESCE(%s, phone), website = COALESCE(%s, website),
lat = COALESCE(%s, lat), lon = COALESCE(%s, lon),
confidence = GREATEST(confidence, 65), updated_at = NOW()
WHERE id = %s
""",
(phone, website, plat, plon, existing["id"]),
)
entity_id = existing["id"]
action = "updated"
else:
row = execute_returning(
"""
INSERT INTO export_market_entities (
name, entity_type, country_iso2, territory_code, lat, lon,
phone, website, google_place_id, source, sources, confidence, metadata
) VALUES (%s,%s,%s,%s,%s,%s,%s,%s,%s,'google_places','["google_places"]'::jsonb,65,%s)
RETURNING id
""",
(
name[:255], entity_type, country_iso2.upper(), territory_code,
plat, plon, phone, website, pid,
json.dumps({"google_place_id": pid}),
),
)
entity_id = row["id"]
action = "inserted"
_apply_place_to_entity(int(existing["id"]), pid, details)
return "updated", int(existing["id"])
row = execute_returning(
"""
INSERT INTO export_market_entities (
name, entity_type, country_iso2, territory_code, lat, lon,
phone, website, google_place_id, source, sources, confidence, metadata
) VALUES (%s,%s,%s,%s,%s,%s,%s,%s,%s,'google_places','["google_places"]'::jsonb,65,%s)
RETURNING id
""",
(
name[:255],
entity_type,
country_iso2.upper(),
territory_code,
plat,
plon,
phone,
website,
pid,
json.dumps({"google_place_id": pid}),
),
)
entity_id = int(row["id"])
if phone:
_upsert_contact(entity_id, None, phone, "google_places", 70, role="sales", name="Google Places")
return action, entity_id
return "inserted", entity_id
def enrich_entities_places_country(country_iso2: str, limit: int = 400) -> dict[str, Any]:
"""Match existing OSM entities to Google Places for phone/website."""
if not _api_key():
return {"skipped": True, "reason": "GOOGLE_MAPS_API_KEY not set"}
entities = fetch_all(
"""
SELECT id, name, lat, lon, city, phone, website
FROM export_market_entities
WHERE country_iso2 = %s
AND lat IS NOT NULL AND lon IS NOT NULL
AND (phone IS NULL OR phone = '')
AND (google_place_id IS NULL OR google_place_id = '')
ORDER BY confidence DESC NULLS LAST, id
LIMIT %s
""",
(country_iso2.upper(), limit),
)
enriched = 0
not_found = 0
errors = 0
for ent in entities:
query = ent["name"] or ""
if ent.get("city"):
query = f"{query} {ent['city']}"
try:
pid = _find_place_id(query, float(ent["lat"]), float(ent["lon"]))
if not pid:
not_found += 1
time.sleep(0.12)
continue
details = _place_details(pid)
if _apply_place_to_entity(int(ent["id"]), pid, details):
enriched += 1
except Exception:
errors += 1
time.sleep(0.15)
return {
"country": country_iso2.upper(),
"candidates": len(entities),
"enriched": enriched,
"not_found": not_found,
"errors": errors,
}
def sync_places_country(country_iso2: str) -> dict[str, Any]:
@@ -165,12 +333,19 @@ def sync_places_country(country_iso2: str) -> dict[str, Any]:
action, _ = _upsert_place(place, details, iso, territory_code, q)
if action == "inserted":
inserted += 1
else:
elif action == "updated":
updated += 1
time.sleep(0.15)
time.sleep(0.5)
return {"inserted": inserted, "updated": updated, "queries": len(QUERIES), "search_points": len(search_points)}
enrich = enrich_entities_places_country(iso, limit=500)
return {
"inserted": inserted,
"updated": updated,
"queries": len(QUERIES),
"search_points": len(search_points),
"enrich_existing": enrich,
}
def sync_places_countries(countries: list[str]) -> dict[str, Any]:
+210 -35
View File
@@ -25,6 +25,15 @@ TED_RSS_URLS = [
from app.export_territory_seeds import BBOX_ONLY_ISO2
from app.export_intel_sync_progress import (
finish_job,
get_active_job,
job_step,
set_active_job,
start_job,
update_job,
)
EMAIL_RE = re.compile(r"[a-zA-Z0-9._%+-]+@[a-zA-Z0-9.-]+\.[a-zA-Z]{2,}")
SKIP_EMAIL_SUFFIX = ("example.com", "sentry.io", "wixpress.com", "png", "jpg", "jpeg")
@@ -140,9 +149,22 @@ def _upsert_osm_entity(
hn = tags.get("addr:housenumber") or ""
address = f"{street} {hn}".strip() or None
city = tags.get("addr:city") or tags.get("addr:town") or tags.get("addr:village")
phone = (tags.get("phone") or tags.get("contact:phone") or "")[:64] or None
phone = (
tags.get("phone")
or tags.get("contact:phone")
or tags.get("contact:mobile")
or tags.get("mobile")
or tags.get("phone:mobile")
or ""
)[:64] or None
email = (tags.get("email") or tags.get("contact:email") or "")[:255] or None
website = (tags.get("website") or tags.get("contact:website") or "")[:512] or None
website = (
tags.get("website")
or tags.get("contact:website")
or tags.get("contact:url")
or tags.get("url")
or ""
)[:512] or None
existing = fetch_one(
"SELECT id FROM export_market_entities WHERE metadata->>'osm_external_id' = %s",
@@ -185,7 +207,7 @@ def _upsert_osm_entity(
entity_id = row["id"]
action = "inserted"
if email:
if email or phone:
_upsert_contact(entity_id, email, phone, "osm", 60)
return action, entity_id
@@ -285,6 +307,9 @@ def sync_osm_country(country_iso2: str, kinds: Optional[list[str]] = None) -> di
parts = mapping.get(kind, [])
if not parts:
continue
active = get_active_job()
if active:
job_step(active, message=f"OSM {kind} · {country_iso2}", country=country_iso2, step=kind)
query = _build_country_query(parts, country_iso2, bbox)
elements = _fetch_overpass(query)
ki = {"fetched": len(elements), "inserted": 0, "updated": 0}
@@ -352,6 +377,46 @@ def materialize_caterer_presence(country_iso2: str) -> int:
return created
PHONE_RE = re.compile(
r"(?:tel:|phone|telefoon)[^\d+]{0,12}(\+?[\d\s\(\)\-\.]{8,20})|(\+31[\s\-\.]?(?:\(?\d{2}\)?[\s\-\.]?)?[\d\s\-\.]{6,12})",
re.I,
)
def _normalize_phone(raw: str) -> Optional[str]:
s = re.sub(r"[^\d+]", "", (raw or "").strip())
if len(s) < 8:
return None
if s.startswith("00"):
s = "+" + s[2:]
if s.startswith("0") and not s.startswith("+"):
s = "+31" + s[1:]
return s[:32]
def _scrape_phones(url: str) -> list[str]:
if not url:
return []
if not url.startswith("http"):
url = "https://" + url
try:
with httpx.Client(timeout=15.0, follow_redirects=True) as client:
resp = client.get(url, headers={"User-Agent": "FoodlinkkExportIntel/1.0"})
if resp.status_code >= 400:
return []
text = resp.text[:500000]
except Exception:
return []
found: set[str] = set()
for m in PHONE_RE.finditer(text):
raw = (m.group(1) or m.group(2) or "").strip()
norm = _normalize_phone(raw)
if norm:
found.add(norm)
return list(found)[:3]
def _scrape_emails(url: str) -> list[str]:
if not url:
return []
@@ -386,8 +451,25 @@ def enrich_contacts_country(country_iso2: str) -> dict[str, Any]:
)
added = 0
scraped = 0
phones_added = 0
for ent in entities:
emails = _scrape_emails(ent["website"])
site = ent["website"]
emails = _scrape_emails(site)
phones = _scrape_phones(site) if not ent.get("phone") else []
# probe common contact pages (gratis, geen API)
if not emails or not phones:
base = site if site.startswith("http") else ("https://" + site)
base = base.rstrip("/")
for path in ("/contact", "/contact/", "/contact-ons", "/contactus", "/over-ons", "/about"):
try:
emails = list(dict.fromkeys(emails + _scrape_emails(base + path)))
if not ent.get("phone"):
phones = list(dict.fromkeys(phones + _scrape_phones(base + path)))
except Exception:
pass
if emails and (ent.get("phone") or phones):
break
time.sleep(0.15)
scraped += 1
for em in emails:
role = "procurement" if any(x in em for x in ("inkoop", "procurement", "sales", "info", "order")) else "general"
@@ -398,13 +480,27 @@ def enrich_contacts_country(country_iso2: str) -> dict[str, Any]:
if not before:
_upsert_contact(ent["id"], em, None, "website_scrape", 65, role=role, name="Website")
added += 1
for ph in phones:
before_p = fetch_one(
"SELECT id FROM export_entity_contacts WHERE entity_id = %s AND phone = %s",
(ent["id"], ph),
)
if not before_p:
_upsert_contact(ent["id"], None, ph, "website_scrape", 60, role="general", name="Website")
added += 1
phones_added += 1
if not ent.get("email") and emails:
execute(
"UPDATE export_market_entities SET email = %s, updated_at = NOW() WHERE id = %s",
(emails[0], ent["id"]),
)
if not ent.get("phone") and phones:
execute(
"UPDATE export_market_entities SET phone = %s, updated_at = NOW() WHERE id = %s",
(phones[0], ent["id"]),
)
time.sleep(0.3)
return {"entities_scraped": scraped, "contacts_added": added}
return {"entities_scraped": scraped, "contacts_added": added, "phones_added": phones_added}
def sync_ted_tenders(limit: int = 50) -> dict[str, Any]:
@@ -476,46 +572,105 @@ def sync_ted_tenders(limit: int = 50) -> dict[str, Any]:
def sync_all_priority(max_priority: int = 2, region_code: Optional[str] = None) -> dict[str, Any]:
"""Full OSM + contacts + Places sync for active territories."""
"""Full OSM + contacts sync for active territories (gratis)."""
iso_list = _countries_for_sync(region_code=region_code, max_priority=max_priority)
label = f"Sync {region_code or 'alle landen'} ({len(iso_list)} landen)"
job_id = start_job("all", label, total=len(iso_list), meta={"region": region_code, "max_priority": max_priority})
report: dict[str, Any] = {
"countries": iso_list,
"country_count": len(iso_list),
"region": region_code,
"max_priority": max_priority,
"steps": {},
"job_id": job_id,
}
try:
update_job(job_id, step="ted", message="TED tenders ophalen…")
report["steps"]["ted"] = sync_ted_tenders(40)
report["steps"]["ted"] = sync_ted_tenders(40)
dist: dict[str, Any] = {}
cater: dict[str, Any] = {}
cust: dict[str, Any] = {}
contacts: dict[str, Any] = {}
for i, c in enumerate(iso_list, start=1):
job_step(job_id, message=f"Land {i}/{len(iso_list)}: {c}", country=c, step="osm", current=i, total=len(iso_list))
try:
materialize_caterer_presence(c)
dist[c] = sync_osm_country(c, kinds=["distributors"])
cater[c] = sync_osm_country(c, kinds=["caterers"])
cust[c] = sync_osm_country(c, kinds=["customers"])
contacts[c] = enrich_contacts_country(c)
except Exception as exc:
dist[c] = {"error": str(exc)}
time.sleep(2)
report["steps"]["distributors"] = dist
report["steps"]["caterers"] = cater
report["steps"]["customers"] = cust
report["steps"]["contacts"] = contacts
report["steps"]["google_places"] = {"skipped": True, "reason": "gratis-modus: geen Google Places"}
finish_job(job_id, result={"country_count": len(iso_list), "region": region_code})
except Exception as exc:
finish_job(job_id, error=str(exc))
raise
finally:
set_active_job(None)
return report
def sync_missing_countries(max_priority: int = 3) -> dict[str, Any]:
"""OSM + contacts sync for active territories with zero entities."""
rows = fetch_all(
"""
SELECT t.country_iso2
FROM export_territories t
WHERE t.is_active
AND t.sync_priority <= %s
AND NOT EXISTS (
SELECT 1 FROM export_market_entities e WHERE e.country_iso2 = t.country_iso2
)
GROUP BY t.country_iso2
ORDER BY MIN(t.sync_priority), t.country_iso2
""",
(max_priority,),
)
iso_list = [r["country_iso2"] for r in rows]
job_id = start_job("missing", f"Lege landen ({len(iso_list)})", total=len(iso_list))
report: dict[str, Any] = {"countries": iso_list, "country_count": len(iso_list), "steps": {}, "job_id": job_id}
dist: dict[str, Any] = {}
cater: dict[str, Any] = {}
cust: dict[str, Any] = {}
contacts: dict[str, Any] = {}
for c in iso_list:
try:
materialize_caterer_presence(c)
dist[c] = sync_osm_country(c, kinds=["distributors"])
cater[c] = sync_osm_country(c, kinds=["caterers"])
cust[c] = sync_osm_country(c, kinds=["customers"])
contacts[c] = enrich_contacts_country(c)
except Exception as exc:
dist[c] = {"error": str(exc)}
time.sleep(2)
report["steps"]["distributors"] = dist
report["steps"]["caterers"] = cater
report["steps"]["customers"] = cust
report["steps"]["contacts"] = contacts
from app.export_intel_places import sync_places_countries
report["steps"]["google_places"] = sync_places_countries(iso_list)
try:
for i, c in enumerate(iso_list, start=1):
job_step(job_id, message=f"Leeg land {i}/{len(iso_list)}: {c}", country=c, current=i, total=len(iso_list))
try:
materialize_caterer_presence(c)
dist[c] = sync_osm_country(c, kinds=["distributors"])
cater[c] = sync_osm_country(c, kinds=["caterers"])
cust[c] = sync_osm_country(c, kinds=["customers"])
contacts[c] = enrich_contacts_country(c)
except Exception as exc:
dist[c] = {"error": str(exc)}
time.sleep(2)
report["steps"] = {
"distributors": dist,
"caterers": cater,
"customers": cust,
"contacts": contacts,
}
finish_job(job_id, result={"country_count": len(iso_list)})
except Exception as exc:
finish_job(job_id, error=str(exc))
raise
finally:
set_active_job(None)
return report
def sync_world_regions(max_priority: int = 2) -> dict[str, Any]:
"""Sync all target regions sequentially (Europe, MENA, Africa, Americas)."""
regions = ["europe", "middle_east", "gcc", "africa", "americas", "caribbean"]
regions = ["europe", "middle_east", "gcc", "africa", "americas", "caribbean", "asia", "oceania"]
out: dict[str, Any] = {"regions": {}, "totals": {"countries": 0}}
for reg in regions:
countries = _countries_for_sync(region_code=reg, max_priority=max_priority)
@@ -571,13 +726,33 @@ def sync_contacts(country_iso2: Optional[str] = None, region_code: Optional[str]
def sync_customers(country_iso2: Optional[str] = None, region_code: Optional[str] = None) -> dict[str, Any]:
if country_iso2:
return sync_osm_country(country_iso2, kinds=["customers"])
countries = _countries_for_sync(region_code=region_code)
results: dict[str, Any] = {}
for c in countries:
job_id = start_job("customers", f"Horeca {country_iso2.upper()}", total=1, meta={"country": country_iso2.upper()})
try:
results[c] = sync_osm_country(c, kinds=["customers"])
job_step(job_id, message=f"Horeca OSM · {country_iso2.upper()}", country=country_iso2.upper(), current=1, total=1)
result = sync_osm_country(country_iso2, kinds=["customers"])
finish_job(job_id, result=result)
return result
except Exception as exc:
results[c] = {"error": str(exc)}
time.sleep(3)
return {"osm": results}
finish_job(job_id, error=str(exc))
raise
finally:
set_active_job(None)
countries = _countries_for_sync(region_code=region_code)
label = f"Horeca {region_code or 'regio'} ({len(countries)} landen)"
job_id = start_job("customers", label, total=len(countries), meta={"region": region_code})
results: dict[str, Any] = {}
try:
for i, c in enumerate(countries, start=1):
job_step(job_id, message=f"Horeca {i}/{len(countries)}: {c}", country=c, current=i, total=len(countries))
try:
results[c] = sync_osm_country(c, kinds=["customers"])
except Exception as exc:
results[c] = {"error": str(exc)}
time.sleep(3)
finish_job(job_id, result={"countries": len(countries), "region": region_code})
except Exception as exc:
finish_job(job_id, error=str(exc))
raise
finally:
set_active_job(None)
return {"osm": results, "job_id": job_id}
+186
View File
@@ -0,0 +1,186 @@
"""File-backed sync progress for Export Intel background and UI jobs."""
from __future__ import annotations
import json
import os
import time
from pathlib import Path
from typing import Any, Optional
STATUS_PATH = Path("/tmp/export-intel-sync-status.json")
LOG_PATH = Path("/tmp/export-intel-sync-progress.jsonl")
LOCK_PATH = Path("/tmp/export-intel-sync.lock")
_active_job_id: Optional[str] = None
def _now() -> float:
return time.time()
def _read_json(path: Path, default: Any) -> Any:
try:
if path.exists():
return json.loads(path.read_text())
except Exception:
pass
return default
def _write_json(path: Path, data: Any) -> None:
path.write_text(json.dumps(data, indent=2, default=str))
def log_event(event: str, **data: Any) -> None:
row = {"ts": _now(), "event": event, **data}
with LOG_PATH.open("a") as fh:
fh.write(json.dumps(row, default=str) + "\n")
def set_active_job(job_id: Optional[str]) -> None:
global _active_job_id
_active_job_id = job_id
def get_active_job() -> Optional[str]:
return _active_job_id
def _load_status() -> dict[str, Any]:
data = _read_json(STATUS_PATH, {"jobs": {}, "history": []})
if "jobs" not in data:
data["jobs"] = {}
if "history" not in data:
data["history"] = []
return data
def _save_status(data: dict[str, Any]) -> None:
_write_json(STATUS_PATH, data)
def _touch_lock(job_id: str, label: str) -> None:
LOCK_PATH.write_text(json.dumps({"pid": os.getpid(), "job_id": job_id, "label": label, "ts": _now()}))
def _clear_lock() -> None:
LOCK_PATH.unlink(missing_ok=True)
def start_job(job_type: str, label: str, total: int = 0, meta: Optional[dict[str, Any]] = None) -> str:
job_id = f"{job_type}-{_now():.0f}"
status = _load_status()
job = {
"id": job_id,
"type": job_type,
"label": label,
"status": "running",
"started_at": _now(),
"updated_at": _now(),
"current": 0,
"total": total,
"percent": 0,
"country": None,
"step": None,
"message": label,
"meta": meta or {},
"error": None,
}
status["jobs"][job_id] = job
_save_status(status)
_touch_lock(job_id, label)
set_active_job(job_id)
log_event("job_start", job_id=job_id, type=job_type, label=label, total=total)
return job_id
def update_job(job_id: str, **fields: Any) -> None:
status = _load_status()
job = status["jobs"].get(job_id)
if not job:
return
job.update(fields)
job["updated_at"] = _now()
cur = int(job.get("current") or 0)
tot = int(job.get("total") or 0)
if tot > 0:
job["percent"] = min(100, round(cur * 100 / tot))
status["jobs"][job_id] = job
_save_status(status)
if LOCK_PATH.exists():
_touch_lock(job_id, job.get("label") or job_id)
def job_step(
job_id: str,
*,
message: str,
country: Optional[str] = None,
step: Optional[str] = None,
current: Optional[int] = None,
total: Optional[int] = None,
) -> None:
fields: dict[str, Any] = {"message": message}
if country is not None:
fields["country"] = country
if step is not None:
fields["step"] = step
if current is not None:
fields["current"] = current
if total is not None:
fields["total"] = total
update_job(job_id, **fields)
log_event("step", job_id=job_id, message=message, country=country, step=step, current=current, total=total)
def finish_job(job_id: str, result: Optional[dict[str, Any]] = None, error: Optional[str] = None) -> None:
status = _load_status()
job = status["jobs"].get(job_id)
if not job:
return
job["status"] = "error" if error else "completed"
job["updated_at"] = _now()
job["finished_at"] = _now()
job["percent"] = 100 if not error else job.get("percent", 0)
job["error"] = error
if result is not None:
job["result"] = result
status["jobs"][job_id] = job
status["history"] = ([job_id] + [h for h in status.get("history", []) if h != job_id])[:20]
_save_status(status)
log_event("job_finish", job_id=job_id, status=job["status"], error=error)
if get_active_job() == job_id:
set_active_job(None)
running = [j for j in status["jobs"].values() if j.get("status") == "running"]
if not running:
_clear_lock()
def tail_log(limit: int = 40) -> list[dict[str, Any]]:
if not LOG_PATH.exists():
return []
lines = LOG_PATH.read_text().strip().splitlines()
out: list[dict[str, Any]] = []
for line in lines[-limit:]:
try:
out.append(json.loads(line))
except Exception:
continue
return out
def get_monitor_snapshot(db_counts: Optional[dict[str, Any]] = None) -> dict[str, Any]:
status = _load_status()
jobs = status.get("jobs", {})
active = [j for j in jobs.values() if j.get("status") == "running"]
recent = [jobs[jid] for jid in status.get("history", [])[:8] if jid in jobs and jobs[jid].get("status") != "running"]
lock = _read_json(LOCK_PATH, None) if LOCK_PATH.exists() else None
return {
"active_jobs": active,
"recent_jobs": recent,
"events": tail_log(50),
"lock": lock,
"any_running": bool(active) or bool(lock),
"coverage": db_counts or {},
"updated_at": _now(),
}