| Server IP : 85.155.190.233 / Your IP : 216.73.216.103 Web Server : nginx/1.24.0 System : Linux antigravity-cli 6.8.0-31-generic #31-Ubuntu SMP PREEMPT_DYNAMIC Sat Apr 20 00:40:06 UTC 2024 x86_64 User : wp-moonbloom ( 1001) PHP Version : 8.3.6 Disable Function : NONE MySQL : OFF | cURL : ON | WGET : ON | Perl : ON | Python : OFF | Sudo : ON | Pkexec : OFF Directory : /opt/moonbloom-dashboard/collector/ |
Upload File : |
"""
Orchestrator CLI for the MoonBloom analytics collector.
Runs every source (Google Ads, Pinterest Ads, Pinterest organic, Pinterest
audience demographics, GA4, Etsy orders CSV, Etsy Ads CSV, Etsy Offsite Ads
manual snapshot) for a computed
date window and upserts the results into Postgres via db.upsert_rows. Each
source is isolated in its own try/except so one failing source (e.g. no
Etsy Ads CSV dropped in yet, or a transient Pinterest per-ad breakdown
error) never stops the others from running — partial success still
persists whatever worked. The run exits 1 only after every source has been
attempted, so cron/monitoring can tell "something failed" without losing
the data that did come in.
Both ad platforms also write campaign/ad_group/ad hierarchy + status rows
into ad_structure_daily (google_ads.fetch()['structure'] /
pinterest_ads.fetch()['structure']), upserted with the same
try/except-per-source isolation as everything else — a structure-fetch
failure never loses that source's spend data, and vice versa, since both
are upserted inside the same source-level try block right after their
metrics rows.
This runs as __main__ inside the container via the Dockerfile's
ENTRYPOINT, not via stdin, so a plain load_dotenv() call is fine here (it
is a no-op if no .env file is present, which is expected inside the
container — real config there comes from docker-compose's env_file).
Usage (via docker-compose, from a later deploy phase):
docker compose run --rm collector --backfill 90
docker compose run --rm collector --days 3 # default if no flag given
"""
from __future__ import annotations
import argparse
import logging
import os
import sys
from datetime import date, timedelta
from dotenv import load_dotenv
load_dotenv()
import db
from sources import etsy_csv, ga4, google_ads, pinterest_ads, pinterest_audience, pinterest_organic
logging.basicConfig(level=logging.INFO, format="%(asctime)s %(levelname)s %(name)s: %(message)s")
log = logging.getLogger("collector")
def _parse_args(argv: list[str] | None = None) -> argparse.Namespace:
parser = argparse.ArgumentParser(description="MoonBloom analytics collector")
group = parser.add_mutually_exclusive_group()
group.add_argument("--backfill", type=int, metavar="N", help="days back from today to backfill (e.g. 90)")
group.add_argument("--days", type=int, metavar="N", help="incremental lookback window in days (default 3)")
args = parser.parse_args(argv)
if args.backfill is None and args.days is None:
args.days = 3
return args
def _date_window(args: argparse.Namespace) -> tuple[str, str]:
n = args.backfill if args.backfill is not None else args.days
end = date.today()
start = end - timedelta(days=n)
return start.isoformat(), end.isoformat()
def _safe_rollback(conn) -> None:
"""A failed statement leaves a psycopg2 connection in an aborted
transaction state until rolled back — without this, one source's DB
error would poison every subsequent source's upserts on the shared
connection, defeating the whole point of isolating sources."""
try:
conn.rollback()
except Exception: # noqa: BLE001 - best-effort cleanup only
pass
def main(argv: list[str] | None = None) -> int:
args = _parse_args(argv)
start_date, end_date = _date_window(args)
log.info("collector run starting: window %s .. %s", start_date, end_date)
conn = db.get_connection()
summary: list[tuple[str, int]] = []
failures: list[str] = []
# --- Google Ads: campaign / keyword / search-term daily -------------
try:
credentials = {
"developer_token": os.environ.get("GOOGLE_ADS_DEVELOPER_TOKEN"),
"client_id": os.environ.get("GOOGLE_ADS_CLIENT_ID"),
"client_secret": os.environ.get("GOOGLE_ADS_CLIENT_SECRET"),
"refresh_token": os.environ.get("GOOGLE_ADS_REFRESH_TOKEN"),
}
customer_id = os.environ.get("GOOGLE_ADS_CUSTOMER_ID")
login_customer_id = os.environ.get("GOOGLE_ADS_LOGIN_CUSTOMER_ID")
data = google_ads.fetch(customer_id, login_customer_id, credentials, start_date, end_date)
n = 0
n += db.upsert_rows(conn, "ad_metrics_daily", data["campaign_daily"], ["date", "source", "campaign_id"])
n += db.upsert_rows(
conn,
"google_keywords_daily",
data["keywords_daily"],
["date", "campaign", "ad_group", "keyword", "match_type"],
)
n += db.upsert_rows(
conn, "google_search_terms_daily", data["search_terms_daily"], ["date", "search_term", "campaign", "ad_group"]
)
n += db.upsert_rows(
conn, "ad_structure_daily", data["structure"], ["date", "source", "level", "entity_id"]
)
summary.append(("google_ads", n))
log.info("[google_ads] OK: %d rows", n)
except Exception as exc: # noqa: BLE001
log.exception("[google_ads] FAILED: %s", exc)
failures.append("google_ads")
_safe_rollback(conn)
# --- Pinterest Ads: account-level + per-ad daily ---------------------
try:
access_token = os.environ.get("PINTEREST_ACCESS_TOKEN")
ad_account_id = os.environ.get("PINTEREST_AD_ACCOUNT_ID")
data = pinterest_ads.fetch(access_token, ad_account_id, start_date, end_date)
n = 0
n += db.upsert_rows(conn, "ad_metrics_daily", data["account_daily"], ["date", "source", "campaign_id"])
n += db.upsert_rows(conn, "pinterest_ads_pins_daily", data["pins_daily"], ["date", "ad_id"])
n += db.upsert_rows(
conn, "ad_structure_daily", data["structure"], ["date", "source", "level", "entity_id"]
)
summary.append(("pinterest_ads", n))
log.info("[pinterest_ads] OK: %d rows", n)
except Exception as exc: # noqa: BLE001
log.exception("[pinterest_ads] FAILED: %s", exc)
failures.append("pinterest_ads")
_safe_rollback(conn)
# --- Pinterest organic: lifetime-metrics snapshot under today -------
try:
access_token = os.environ.get("PINTEREST_ACCESS_TOKEN")
rows = pinterest_organic.fetch(access_token)
n = db.upsert_rows(conn, "pinterest_organic_daily", rows, ["date", "pin_id"])
summary.append(("pinterest_organic", n))
log.info("[pinterest_organic] OK: %d rows", n)
except Exception as exc: # noqa: BLE001
log.exception("[pinterest_organic] FAILED: %s", exc)
failures.append("pinterest_organic")
_safe_rollback(conn)
# --- Pinterest audience demographics: age/gender/device/country/metro --
try:
access_token = os.environ.get("PINTEREST_ACCESS_TOKEN")
ad_account_id = os.environ.get("PINTEREST_AD_ACCOUNT_ID")
rows = pinterest_audience.fetch(access_token, ad_account_id)
n = db.upsert_rows(
conn, "pinterest_audience_daily", rows, ["date", "audience_type", "dimension", "key"]
)
summary.append(("pinterest_audience", n))
log.info("[pinterest_audience] OK: %d rows", n)
except Exception as exc: # noqa: BLE001
log.exception("[pinterest_audience] FAILED: %s", exc)
failures.append("pinterest_audience")
_safe_rollback(conn)
# --- GA4: traffic by source/medium + landing-page engagement --------
try:
property_id = os.environ.get("GA4_PROPERTY_ID")
data = ga4.fetch(property_id=property_id, start_date=start_date, end_date=end_date)
n = 0
n += db.upsert_rows(conn, "traffic_daily", data["traffic_daily"], ["date", "source_medium"])
n += db.upsert_rows(conn, "ga4_landing_daily", data["landing_daily"], ["date", "channel", "landing_page"])
summary.append(("ga4", n))
log.info("[ga4] OK: %d rows", n)
except Exception as exc: # noqa: BLE001
log.exception("[ga4] FAILED: %s", exc)
failures.append("ga4")
_safe_rollback(conn)
# --- Etsy orders CSV ---------------------------------------------------
# Where this file lives inside the container is a deploy-phase wiring
# detail (mount / copy it in), not something to hard-code here. Default
# matches a plausible future docker-compose volume mount; a missing env
# var or missing file is logged and skipped, never a hard failure.
try:
orders_csv_path = os.environ.get("ETSY_ORDERS_CSV", "/data/orders.csv")
if not os.path.exists(orders_csv_path):
log.info(
"[etsy_orders] SKIPPED: ETSY_ORDERS_CSV (%s) not found — "
"set ETSY_ORDERS_CSV and mount the file once deploy wiring is in place",
orders_csv_path,
)
summary.append(("etsy_orders", 0))
else:
rows = etsy_csv.fetch_orders(orders_csv_path)
n = db.upsert_rows(conn, "etsy_orders", rows, ["order_id"])
summary.append(("etsy_orders", n))
log.info("[etsy_orders] OK: %d rows", n)
except Exception as exc: # noqa: BLE001
log.exception("[etsy_orders] FAILED: %s", exc)
failures.append("etsy_orders")
_safe_rollback(conn)
# --- Etsy Ads CSV (manual export, real column mapping confirmed) -----
try:
ads_csv_dir = os.environ.get("ETSY_ADS_CSV_DIR", "/data/etsy_ads_imports")
rows = etsy_csv.fetch_ads_csv(ads_csv_dir)
n = db.upsert_rows(conn, "ad_metrics_daily", rows, ["date", "source", "campaign_id"])
summary.append(("etsy_ads_csv", n))
log.info("[etsy_ads_csv] OK: %d rows", n)
except Exception as exc: # noqa: BLE001
log.exception("[etsy_ads_csv] FAILED: %s", exc)
failures.append("etsy_ads_csv")
_safe_rollback(conn)
# --- Etsy Offsite Ads (manual monthly snapshot, no export exists) ----
try:
offsite_csv_path = os.environ.get(
"ETSY_OFFSITE_ADS_CSV", "/data/etsy_offsite_ads_imports/offsite_ads_monthly.csv"
)
rows = etsy_csv.fetch_offsite_ads_snapshots(offsite_csv_path)
n = db.upsert_rows(conn, "etsy_offsite_ads_monthly", rows, ["month"])
summary.append(("etsy_offsite_ads", n))
log.info("[etsy_offsite_ads] OK: %d rows", n)
except Exception as exc: # noqa: BLE001
log.exception("[etsy_offsite_ads] FAILED: %s", exc)
failures.append("etsy_offsite_ads")
_safe_rollback(conn)
conn.close()
print()
print(f"Collector run summary ({start_date} .. {end_date})")
print("-" * 42)
for label, n in summary:
print(f" {label:<20} {n:>6} rows")
if failures:
print(f"\nFAILED sources: {', '.join(failures)}")
print()
return 1 if failures else 0
if __name__ == "__main__":
sys.exit(main())