afa67's picture
deploy dry-run demo
257a670 verified
Raw
History Blame
7.23 kB
"""Windsor.ai pull β†’ metrics_daily upsert + creative back-mapping (SPEC_NEW Β§7.6) β€” Phase 3.
Ingest sources ad-level daily metrics from the Windsor.ai Data API (configured Meta
connector), not the Meta Graph insights endpoint. Auth via WINDSOR_API_KEY. Verify the
exact Windsor.ai endpoint + field names at build time. Publishing stays on Meta
(see meta_client.py).
Flow (SPEC_NEW Β§7.6, Phase 3 acceptance Β§8):
- DRY_RUN or no WINDSOR_API_KEY β†’ load fixtures/windsor_sample.json (no network), so the
whole path is testable without credentials. Live β†’ GET the Windsor.ai Data API.
- Map Windsor field names β†’ metrics_daily columns via FIELD_MAP (override with WINDSOR_FIELDS).
- Upsert each row ON CONFLICT (ad_id, date), storing the full source row as raw jsonb.
- Back-map: ad_ids present in publishes are KNOWN (belong to a video); unknown ads are still
ingested (raw kept). metrics_daily has no video FK β€” the link is publishes.ad_id at read time.
"""
from __future__ import annotations
import json
import httpx
from psycopg.types.json import Json
from . import db
from .config import FIXTURES_DIR, get_settings
# Verify-at-build-time (same convention as VEO_MODEL_ID / META_API_VERSION, SPEC_NEW Β§3):
# Windsor.ai Data API. The connector path segment is s.windsor_connector ("facebook" for Meta).
# GET https://connectors.windsor.ai/{connector}?api_key=...&fields=...&date_from=...&date_to=...
# Response shape: {"data": [ {<field>: <value>, ...}, ... ]}.
WINDSOR_BASE_URL = "https://connectors.windsor.ai" # verify-at-build-time
WINDSOR_TIMEOUT_S = 120
FIXTURE = FIXTURES_DIR / "windsor_sample.json"
# metrics_daily column β†’ Windsor.ai field name. Override per-deploy with WINDSOR_FIELDS
# ("col=windsor_field,col=windsor_field,..."). The Meta-connector field names below are
# best-effort β€” verify-at-build-time against the Windsor.ai field reference for `facebook`
# (esp. video_3s / thruplay, which Meta names drift on). The fixture uses these defaults.
FIELD_MAP: dict[str, str] = {
"ad_id": "ad_id",
"date": "date",
"impressions": "impressions",
"spend": "spend",
"clicks": "clicks",
"video_3s": "video_3s", # Meta "3-second video plays" β€” verify-at-build-time
"thruplay": "thruplay", # Meta "ThruPlay" β€” verify-at-build-time
"purchases": "purchases",
"revenue": "revenue",
}
# Numeric metric columns (everything except the (ad_id, date) key) β€” used to build the upsert.
METRIC_COLS = ("impressions", "spend", "clicks", "video_3s", "thruplay", "purchases", "revenue")
def _field_map() -> dict[str, str]:
"""FIELD_MAP, with WINDSOR_FIELDS overrides applied ("col=field,col=field")."""
fields = dict(FIELD_MAP)
override = (get_settings().windsor_fields or "").strip()
if override:
for pair in override.split(","):
pair = pair.strip()
if not pair or "=" not in pair:
continue
col, _, name = pair.partition("=")
col, name = col.strip(), name.strip()
if col in fields and name:
fields[col] = name
return fields
def _fetch_rows(days: int) -> list[dict]:
"""The list of source rows from Windsor.ai (or the fixture in dry-run / no key)."""
s = get_settings()
fields = _field_map()
if s.dry_run or not s.windsor_api_key:
why = "DRY_RUN" if s.dry_run else "no WINDSOR_API_KEY"
print(f"[ingest] {why}: loading fixture {FIXTURE.name} instead of Windsor.ai "
f"(connector={s.windsor_connector}, days={days})")
payload = json.loads(FIXTURE.read_text())
return list(payload.get("data") or [])
# Live: GET the Windsor.ai Data API (verify-at-build-time β€” endpoint/fields/date params).
from datetime import date, timedelta
s.require("WINDSOR_API_KEY")
date_to = date.today()
date_from = date_to - timedelta(days=days)
params = {
"api_key": s.windsor_api_key,
"fields": ",".join(fields[c] for c in FIELD_MAP), # request exactly the mapped fields
"date_from": date_from.isoformat(),
"date_to": date_to.isoformat(),
}
url = f"{WINDSOR_BASE_URL}/{s.windsor_connector}"
redacted = {**params, "api_key": "***"}
print(f"[ingest] GET {url} {redacted}")
r = httpx.get(url, params=params, timeout=WINDSOR_TIMEOUT_S)
if r.status_code >= 400:
raise RuntimeError(f"Windsor.ai {r.status_code}: {r.text[:400]}")
return list(r.json().get("data") or [])
def _coerce(row: dict, fields: dict[str, str]) -> dict | None:
"""Pull metrics_daily column values out of one Windsor row via the field map.
Returns None (skip) if the row has no ad_id/date β€” those are the only required keys.
Metric values pass through as-is (None when absent); the full source row is kept in raw.
"""
ad_id = row.get(fields["ad_id"])
date = row.get(fields["date"])
if ad_id is None or date is None:
return None
out: dict = {"ad_id": str(ad_id), "date": date}
for col in METRIC_COLS:
out[col] = row.get(fields[col])
return out
def run(days: int = 14) -> int:
"""Pull `days` of ad-level daily metrics via Windsor.ai and upsert into metrics_daily.
Returns the number of metrics_daily rows upserted. Logs known (in publishes) vs unknown
ad counts; unknown ads are still ingested with their raw row kept (SPEC_NEW Β§7.6).
"""
fields = _field_map()
rows = _fetch_rows(days)
# Known = ad_ids that round-trip back to a video via publishes.ad_id (no FK; joined at read).
known_ids = {
r["ad_id"]
for r in db.fetch_all("select distinct ad_id from publishes where ad_id is not null")
}
upserted = 0
seen_known: set[str] = set()
seen_unknown: set[str] = set()
with db.connect() as conn:
for src in rows:
vals = _coerce(src, fields)
if vals is None:
continue
conn.execute(
"""
insert into metrics_daily
(ad_id, date, impressions, spend, clicks, video_3s, thruplay, purchases, revenue, raw)
values
(%(ad_id)s, %(date)s, %(impressions)s, %(spend)s, %(clicks)s, %(video_3s)s,
%(thruplay)s, %(purchases)s, %(revenue)s, %(raw)s)
on conflict (ad_id, date) do update set
impressions = excluded.impressions,
spend = excluded.spend,
clicks = excluded.clicks,
video_3s = excluded.video_3s,
thruplay = excluded.thruplay,
purchases = excluded.purchases,
revenue = excluded.revenue,
raw = excluded.raw
""",
{**vals, "raw": Json(src)},
)
upserted += 1
(seen_known if vals["ad_id"] in known_ids else seen_unknown).add(vals["ad_id"])
print(f"[ingest] upserted {upserted} metrics_daily row(s) over {days}d β€” "
f"{len(seen_known)} known ad(s) (linked to videos via publishes), "
f"{len(seen_unknown)} unknown ad(s) ingested raw-only")
return upserted