Uname:Linux antigravity-cli 6.8.0-31-generic #31-Ubuntu SMP PREEMPT_DYNAMIC Sat Apr 20 00:40:06 UTC 2024 x86_64

Base Dir : /var/www/moonbloom

User : wp-moonbloom


403WebShell
403Webshell
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 :
current_dir [ Writeable ] document_root [ Writeable ]

 

Command :


[ Back ]     

Current File : /opt/moonbloom-dashboard/collector/collect.py
"""
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())

Youez - 2016 - github.com/yon3zu
LinuXploit