"""Google Sheets connector.
Ingests marketing spend data from a shared Google Sheet (e.g., Windsor.ai rollup)
into the AdMetric model. Uses OAuth 2.0 Authorization Code flow with PKCE so the
user simply clicks "Connect", authorises in a Google popup, and picks a spreadsheet.
Expected sheet columns (auto-detected, case-insensitive):
date, platform/channel/source, spend, impressions, clicks, conversions,
ctr, cpc, roas
"""
from __future__ import annotations
import hashlib
import logging
import time
from base64 import urlsafe_b64encode
from datetime import datetime, date, timezone, timedelta
from typing import Any, Dict, List, Optional, Tuple
from . import OAuthConnector, register_connector, _REGISTRY
logger = logging.getLogger(__name__)
# -- Registration metadata --------------------------------------------------
register_connector(
"google_sheets",
{
"service": "google_sheets",
"name": "Google Sheets",
"category": "data_source",
"description": "Ingest marketing spend data from a Google Sheet (e.g., Windsor.ai rollup) via OAuth.",
"auth_type": "oauth2",
"auth_fields": [], # OAuth flow — no manual fields needed
"capabilities": ["spreadsheet_sync", "marketing_spend", "auto_detect_columns"],
"rate_limit": "100 requests/minute; spreadsheet read is typically 1-2 requests",
"docs_url": "https://developers.google.com/sheets/api",
},
)
# -- OAuth 2.0 constants -----------------------------------------------------
_GOOGLE_AUTH_URL = "https://accounts.google.com/o/oauth2/v2/auth"
_GOOGLE_TOKEN_URL = "https://oauth2.googleapis.com/token"
_GOOGLE_SCOPES = [
"https://www.googleapis.com/auth/spreadsheets.readonly",
"https://www.googleapis.com/auth/drive.readonly",
]
# -- Column name detection map ----------------------------------------------
# Maps various column header names to AdMetric fields.
# Keys are normalised (lower-cased, stripped of spaces/hyphens).
_COLUMN_ALIASES: Dict[str, str] = {
# date
"date": "metric_date",
"daterange": "metric_date",
"campaign_date": "metric_date",
"report_date": "metric_date",
# platform / channel (stored in metadata, not a direct AdMetric field)
"platform": "platform",
"channel": "platform",
"source": "platform",
"ad_source": "platform",
"adplatform": "platform",
# spend
"spend": "spend",
"amount": "spend",
"cost": "spend",
"adspend": "spend",
"total_spend": "spend",
# impressions
"impressions": "impressions",
"imps": "impressions",
"total_impressions": "impressions",
# clicks
"clicks": "clicks",
"total_clicks": "clicks",
# conversions
"conversions": "conversions",
"conversion": "conversions",
"total_conversions": "conversions",
# ctr
"ctr": "ctr",
"clickthroughrate": "ctr",
"click_through_rate": "ctr",
# cpc
"cpc": "cpc",
"costperclick": "cpc",
"cost_per_click": "cpc",
# roas
"roas": "roas",
"returnonadspend": "roas",
"returnonadsspend": "roas",
"return_on_ad_spend": "roas",
"return_on_ads_spend": "roas",
}
# Fields that are numeric and should default to 0 when empty
_NUMERIC_ZERO_DEFAULTS = {"spend", "impressions", "clicks", "conversions"}
# Fields that are numeric but nullable (optional)
_NUMERIC_NULL_DEFAULTS = {"ctr", "cpc", "cpv", "roas"}
# -- Helpers -----------------------------------------------------------------
def _normalise_header(header: str) -> str:
"""Lower-case, strip, remove spaces/hyphens for flexible matching."""
import re
return re.sub(r"[\s\-_]+", "", header.strip().lower())
def _parse_headers(headers: List[str]) -> Dict[int, str]:
"""Return a mapping of column index -> target AdMetric field name.
Platform columns map to ``"platform"`` so they are stored in ``metadata_json``.
"""
mapping: Dict[int, str] = {}
for idx, header in enumerate(headers):
norm = _normalise_header(header)
if norm in _COLUMN_ALIASES:
mapping[idx] = _COLUMN_ALIASES[norm]
return mapping
def _parse_date(value: Any) -> Optional[date]:
"""Parse date from multiple formats: YYYY-MM-DD, MM/DD/YYYY, etc."""
if isinstance(value, date):
return value
if value is None:
return None
raw = str(value).strip()
if not raw:
return None
# ISO format (YYYY-MM-DD)
try:
return datetime.strptime(raw, "%Y-%m-%d").date()
except (ValueError, TypeError):
pass
# US format (MM/DD/YYYY)
try:
return datetime.strptime(raw, "%m/%d/%Y").date()
except (ValueError, TypeError):
pass
# US format with 2-digit year (MM/DD/YY)
try:
return datetime.strptime(raw, "%m/%d/%y").date()
except (ValueError, TypeError):
pass
# European format (DD.MM.YYYY)
try:
return datetime.strptime(raw, "%d.%m.%Y").date()
except (ValueError, TypeError):
pass
# Try parsing as a full datetime string
try:
return datetime.fromisoformat(raw.replace("Z", "+00:00")).date()
except (ValueError, TypeError):
pass
# Try Python's fromisoformat (handles most ISO-like formats)
try:
return datetime.fromisoformat(raw).date()
except (ValueError, TypeError):
pass
logger.warning("Could not parse date value: %r", value)
return None
def _parse_currency(value: Any) -> Optional[float]:
"""Strip currency symbols and commas, return float or None."""
import re
if value is None:
return None
raw = str(value).strip()
if not raw:
return None
# Remove $, €, £, commas, spaces, and percent signs
cleaned = re.sub(r"[$€£,\\s%]", "", raw)
try:
return float(cleaned)
except (ValueError, TypeError):
logger.warning("Could not parse currency value: %r", value)
return None
def _parse_int(value: Any, default: int = 0) -> int:
"""Parse an integer value, stripping commas. Returns *default* on failure."""
if value is None:
return default
raw = str(value).strip()
if not raw:
return default
cleaned = raw.replace(",", "")
try:
return int(float(cleaned))
except (ValueError, TypeError):
return default
def _parse_float(value: Any) -> Optional[float]:
"""Parse a float value or return None."""
if value is None:
return None
raw = str(value).strip()
if not raw:
return None
cleaned = raw.replace(",", "")
try:
return float(cleaned)
except (ValueError, TypeError):
return None
def _build_external_campaign_id(platform: str, metric_date: date) -> str:
"""Build a deterministic external_campaign_id from platform + date."""
return f"google_sheets_{platform}_{metric_date.isoformat()}"
# -- PKCE helpers ------------------------------------------------------------
def _pkce_code_verifier() -> str:
"""Generate a PKCE code verifier (43-128 chars, base64url)."""
import secrets
return urlsafe_b64encode(secrets.token_bytes(64)).decode("ascii").rstrip("=")
def _pkce_code_challenge(verifier: str) -> str:
"""S256 transform of a code verifier."""
digest = hashlib.sha256(verifier.encode("ascii")).digest()
return urlsafe_b64encode(digest).decode("ascii").rstrip("=")
# -- Google Sheets API wrapper ----------------------------------------------
def _build_oauth_credentials(config: Dict[str, Any]) -> Any:
"""Build Google OAuth 2.0 credentials from stored tokens.
Auto-refreshes if the access token is expired.
Returns a credentials object suitable for googleapiclient.build().
"""
from google.oauth2 import credentials
client_id = config.get("client_id", "")
client_secret = config.get("client_secret", "")
access_token = config.get("access_token", "")
refresh_token = config.get("refresh_token", "")
token_expiry_str = config.get("token_expiry")
# Check if token is expired and needs refresh
if token_expiry_str:
try:
expiry = datetime.fromisoformat(token_expiry_str)
if expiry <= datetime.now(timezone.utc):
logger.info("Google Sheets: token expired, refreshing...")
new_config = _refresh_tokens(client_id, client_secret, refresh_token)
if new_config:
config.update(new_config)
access_token = config.get("access_token", access_token)
refresh_token = config.get("refresh_token", refresh_token)
except (ValueError, TypeError):
logger.warning("Google Sheets: invalid token_expiry format, attempting refresh")
creds = credentials.Credentials(
token=access_token,
refresh_token=refresh_token,
token_uri=_GOOGLE_TOKEN_URL,
client_id=client_id,
client_secret=client_secret,
scopes=_GOOGLE_SCOPES,
)
return creds
def _refresh_tokens(client_id: str, client_secret: str, refresh_token: str) -> Optional[Dict[str, str]]:
"""Refresh Google OAuth tokens via the token endpoint."""
import requests
if not refresh_token:
logger.error("Google Sheets: no refresh_token available")
return None
resp = requests.post(
_GOOGLE_TOKEN_URL,
data={
"grant_type": "refresh_token",
"refresh_token": refresh_token,
"client_id": client_id,
"client_secret": client_secret,
},
headers={"Content-Type": "application/x-www-form-urlencoded"},
timeout=15,
)
resp.raise_for_status()
data = resp.json()
expires_in = int(data.get("expires_in", 3600))
token_expiry = (datetime.now(timezone.utc) + timedelta(seconds=expires_in)).isoformat()
result = {
"access_token": data.get("access_token", ""),
"refresh_token": data.get("refresh_token", refresh_token),
"token_expiry": token_expiry,
"expires_in": expires_in,
}
logger.info("Google Sheets: tokens refreshed, new expiry=%s", token_expiry)
return result
def _fetch_spreadsheet_data(
spreadsheet_id: str,
credentials: Any,
range_name: Optional[str] = None,
) -> Dict[str, Any]:
"""Fetch data from a Google Sheet using the Sheets API v4.
Returns a dict with ``values`` (list of rows) and ``column_metadata``.
"""
from googleapiclient.discovery import build
service = build("sheets", "v4", credentials=credentials)
sheet = service.spreadsheets()
if range_name is None:
# Default to the first sheet, all used range
spreadsheet = sheet.get(spreadsheetId=spreadsheet_id).execute()
sheet_properties = spreadsheet.get("sheets", [{}])[0].get("properties", {})
sheet_id = sheet_properties.get("sheetId", 0)
# Try to get sheet name
sheet_name = "Sheet1"
for sp in spreadsheet.get("sheets", []):
if sp.get("properties", {}).get("sheetId") == sheet_id:
sheet_name = sp["properties"].get("title", "Sheet1")
break
range_name = f"'{sheet_name}'!A:Z"
result = sheet.values().get(spreadsheetId=spreadsheet_id, range=range_name).execute()
return result
def _list_spreadsheets(credentials: Any) -> List[Dict[str, Any]]:
"""List user's spreadsheets via Drive API for the spreadsheet picker."""
from googleapiclient.discovery import build
drive_service = build("drive", "v3", credentials=credentials)
results = drive_service.files().list(
q="mimeType='application/vnd.google-apps.spreadsheet'",
pageSize=50,
fields="files(id, name, modifiedTime, webViewLink)",
).execute()
files = results.get("files", [])
return [
{
"id": f["id"],
"name": f["name"],
"modifiedTime": f.get("modifiedTime", ""),
"webViewLink": f.get("webViewLink", ""),
}
for f in files
]
# -- Connector class --------------------------------------------------------
class GoogleSheetsConnector(OAuthConnector):
"""Google Sheets integration — OAuth 2.0 with PKCE, ingests marketing spend data into AdMetric."""
_SERVICE = "google_sheets"
# -- OAuth 2.0 ------------------------------------------------------------
OAUTH_AUTHORIZE_URL = _GOOGLE_AUTH_URL
OAUTH_TOKEN_URL = _GOOGLE_TOKEN_URL
OAUTH_SCOPES = _GOOGLE_SCOPES
# -- OAuth abstract methods -----------------------------------------------
def oauth_authorize_url(self, state: str) -> str:
"""Build Google OAuth authorization URL with PKCE."""
from urllib.parse import urlencode
from flask import current_app
# Generate PKCE verifier and challenge
verifier = _pkce_code_verifier()
challenge = _pkce_code_challenge(verifier)
# Store the PKCE verifier in Redis state by updating the existing
# state entry generated by generate_oauth_state.
from app.connectors import _update_oauth_state
_update_oauth_state(service=self._SERVICE, state=state, pkce_verifier=verifier)
redirect_uri = self._get_redirect_uri()
params = {
"client_id": self._get_client_id(),
"response_type": "code",
"scope": " ".join(self.OAUTH_SCOPES),
"redirect_uri": redirect_uri,
"state": state,
"access_type": "offline",
"prompt": "consent",
"code_challenge": challenge,
"code_challenge_method": "S256",
}
return f"{self.OAUTH_AUTHORIZE_URL}?{urlencode(params)}"
def exchange_code_for_tokens(self, code: str, pkce_verifier: str = "") -> Dict[str, Any]:
"""Exchange authorization code for access/refresh tokens via Google."""
import requests
client_id = self._get_client_id()
client_secret = self._get_client_secret()
if not client_id or not client_secret:
raise ValueError(
"Google Sheets OAuth credentials not configured. "
"Set GOOGLE_SHEETS_CLIENT_ID and GOOGLE_SHEETS_CLIENT_SECRET."
)
redirect_uri = self._get_redirect_uri()
payload = {
"grant_type": "authorization_code",
"code": code,
"redirect_uri": redirect_uri,
"client_id": client_id,
"client_secret": client_secret,
}
if pkce_verifier:
payload["code_verifier"] = pkce_verifier
logger.debug("Google Sheets token exchange (PKCE=%s)", "yes" if pkce_verifier else "no")
response = requests.post(
self.OAUTH_TOKEN_URL,
data=payload,
headers={"Content-Type": "application/x-www-form-urlencoded"},
timeout=15,
)
# Log full response for debugging 400 errors
if response.status_code != 200:
logger.error(
"Google Sheets token exchange failed: %d %s",
response.status_code, response.text[:500],
)
response.raise_for_status()
data = response.json()
expires_in = int(data.get("expires_in", 3600))
token_expiry = (datetime.now(timezone.utc) + timedelta(seconds=expires_in)).isoformat()
return {
"access_token": data.get("access_token", ""),
"refresh_token": data.get("refresh_token", ""),
"expires_in": expires_in,
"token_expiry": token_expiry,
"token_type": data.get("token_type", "Bearer"),
"client_id": client_id,
"client_secret": client_secret,
}
def refresh_access_token(self, refresh_token: str) -> Dict[str, str]:
"""Refresh expired access token via Google token endpoint."""
client_id = self._get_client_id()
client_secret = self._get_client_secret()
if not client_id or not client_secret:
raise ValueError("Google Sheets OAuth credentials not configured.")
new_tokens = _refresh_tokens(client_id, client_secret, refresh_token)
if new_tokens is None:
raise ValueError("Token refresh failed — no refresh token available")
return {
"access_token": new_tokens["access_token"],
"refresh_token": new_tokens["refresh_token"],
"expires_in": new_tokens.get("expires_in", 3600),
}
# -- connect / disconnect -----------------------------------------------
def connect(self) -> Dict[str, Any]:
"""Validate OAuth tokens by attempting to read spreadsheet metadata."""
self._log(event_type="connect_attempt", status="pending")
start = time.monotonic()
try:
spreadsheet_id = self.config.get("spreadsheet_id", "")
if not spreadsheet_id:
# No spreadsheet selected yet — that's ok, tokens are valid
# Validate tokens work by checking token info
creds = _build_oauth_credentials(self.config)
from googleapiclient.discovery import build
drive = build("drive", "v3", credentials=creds)
# Just verify we can make an API call
drive.about().get(fields="user").execute()
self._connected = True
duration_ms = int((time.monotonic() - start) * 1000)
self._log(
event_type="connect_success",
status="success",
duration_ms=duration_ms,
details={"spreadsheet_selected": False},
)
return {
"status": "connected",
"service": "google_sheets",
"spreadsheet_selected": False,
"duration_ms": duration_ms,
}
credentials = _build_oauth_credentials(self.config)
# Validate by trying to read the spreadsheet metadata
from googleapiclient.discovery import build
service = build("sheets", "v4", credentials=credentials)
sheet = service.spreadsheets()
spreadsheet = sheet.get(spreadsheetId=spreadsheet_id).execute()
sheet_name = spreadsheet.get("properties", {}).get("title", "Unknown")
self._connected = True
# Save updated tokens back to config (in case refresh happened)
if credentials.refresh_token:
self.config["access_token"] = credentials.token
self.config["refresh_token"] = credentials.refresh_token
duration_ms = int((time.monotonic() - start) * 1000)
self._log(
event_type="connect_success",
status="success",
duration_ms=duration_ms,
details={
"spreadsheet_id": spreadsheet_id,
"spreadsheet_title": sheet_name,
},
)
return {
"status": "connected",
"service": "google_sheets",
"spreadsheet_id": spreadsheet_id,
"spreadsheet_title": sheet_name,
"duration_ms": duration_ms,
}
except Exception as exc:
self._log(event_type="connect_error", status="error", error_message=str(exc))
return {"status": "error", "error": str(exc)}
def disconnect(self) -> Dict[str, Any]:
"""Reset connection state. Revoke tokens if possible."""
refresh_token = self.config.get("refresh_token", "")
if refresh_token:
try:
import requests
requests.post(
"https://oauth2.googleapis.com/revoke",
params={"token": refresh_token},
timeout=10,
)
except Exception:
logger.warning("Google Sheets: token revocation failed (non-critical)")
self._connected = False
self._log(event_type="disconnect", status="success")
return {"status": "disconnected", "service": "google_sheets"}
# -- sync ----------------------------------------------------------------
def sync(self) -> Dict[str, Any]:
"""Read the shared spreadsheet and upsert AdMetric records."""
self._log(event_type="sync_start", status="pending")
start = time.monotonic()
try:
spreadsheet_id = self.config.get("spreadsheet_id", "")
range_name = self.config.get("range_name") # Optional override
if not spreadsheet_id:
raise ValueError("spreadsheet_id is required — please select a spreadsheet first")
# Build credentials (auto-refreshes if expired)
credentials = _build_oauth_credentials(self.config)
# Fetch data from Google Sheets
data = self._retry(
_fetch_spreadsheet_data,
spreadsheet_id=spreadsheet_id,
credentials=credentials,
range_name=range_name,
)
# Save updated tokens after potential refresh
if credentials.refresh_token:
self.config["access_token"] = credentials.token
self.config["refresh_token"] = credentials.refresh_token
values: List[List[Any]] = data.get("values", [])
if not values:
raise ValueError("Spreadsheet returned no data rows")
# First row = headers; rest = data
headers = [str(h) for h in values[0]]
rows = values[1:]
logger.info(
"Google Sheets: detected %d columns, %d data rows",
len(headers), len(rows),
)
# Auto-detect column mapping
col_map = _parse_headers(headers)
logger.info("Google Sheets: column mapping = %s", col_map)
if not col_map:
raise ValueError(
"No recognisable columns found. Expected: date, platform, spend, "
"impressions, clicks, conversions, ctr, cpc, roas"
)
# Process rows
from app.models import db, AdMetric
records_created = 0
records_updated = 0
rows_skipped = 0
for row in rows:
try:
metric_date = None
platform = "Unknown"
metric_fields: Dict[str, Any] = {}
for idx, field_name in col_map.items():
value = row[idx] if idx < len(row) else None
if field_name == "metric_date":
metric_date = _parse_date(value)
elif field_name == "platform":
platform = str(value).strip() if value else "Unknown"
elif field_name == "spend":
metric_fields["spend"] = _parse_currency(value) or 0.0
elif field_name == "impressions":
metric_fields["impressions"] = _parse_int(value, 0)
elif field_name == "clicks":
metric_fields["clicks"] = _parse_int(value, 0)
elif field_name == "conversions":
metric_fields["conversions"] = _parse_float(value) or 0.0
elif field_name == "ctr":
metric_fields["ctr"] = _parse_float(value)
elif field_name == "cpc":
metric_fields["cpc"] = _parse_float(value)
elif field_name == "roas":
metric_fields["roas"] = _parse_float(value)
# Date is required
if metric_date is None:
rows_skipped += 1
continue
# Build external_campaign_id for dedup
external_campaign_id = _build_external_campaign_id(platform, metric_date)
# Build metadata with platform and raw row
metadata: Dict[str, Any] = {
"platform": platform,
"source": "google_sheets",
"spreadsheet_id": spreadsheet_id,
}
# Merge/upsert
existing = (
AdMetric.query.filter_by(
company_id=self.company_id,
source_service="google_sheets",
metric_date=metric_date,
external_campaign_id=external_campaign_id,
)
.first()
)
if existing:
for key, value in metric_fields.items():
setattr(existing, key, value)
existing.metadata_json = metadata
existing.updated_at = datetime.now(timezone.utc)
records_updated += 1
else:
record = AdMetric(
company_id=self.company_id,
external_campaign_id=external_campaign_id,
source_service="google_sheets",
metric_date=metric_date,
metadata_json=metadata,
**metric_fields,
)
db.session.add(record)
records_created += 1
except Exception as row_exc:
logger.warning(
"Google Sheets: skipped row %d — %s",
rows.index(row) + 2, # +2 because 1-indexed and header row
row_exc,
)
rows_skipped += 1
continue
db.session.commit()
total_records = records_created + records_updated
duration_ms = int((time.monotonic() - start) * 1000)
self._log(
event_type="sync_complete",
status="success",
record_count=total_records,
duration_ms=duration_ms,
details={
"rows_processed": len(rows),
"records_created": records_created,
"records_updated": records_updated,
"rows_skipped": rows_skipped,
"column_mapping": {str(k): v for k, v in col_map.items()},
},
)
self.config["last_sync_at"] = datetime.now(timezone.utc).isoformat()
return {
"status": "success",
"record_count": total_records,
"duration_ms": duration_ms,
"details": {
"rows_processed": len(rows),
"records_created": records_created,
"records_updated": records_updated,
"rows_skipped": rows_skipped,
},
"last_sync_at": self.config.get("last_sync_at"),
}
except Exception as exc:
duration_ms = int((time.monotonic() - start) * 1000)
self._log(
event_type="sync_error",
status="error",
duration_ms=duration_ms,
error_message=str(exc),
)
return {
"status": "error",
"error": str(exc),
}
# -- status --------------------------------------------------------------
def status(self) -> Dict[str, Any]:
"""Check connection health by attempting to read spreadsheet metadata."""
if not self._connected:
return {
"service": "google_sheets",
"connected": False,
"config_present": bool(
self.config.get("spreadsheet_id")
or self.config.get("access_token")
),
}
try:
spreadsheet_id = self.config.get("spreadsheet_id", "")
credentials = _build_oauth_credentials(self.config)
from googleapiclient.discovery import build
sheets_service = build("sheets", "v4", credentials=credentials)
spreadsheet = (
sheets_service.spreadsheets()
.get(spreadsheetId=spreadsheet_id)
.execute()
)
return {
"service": "google_sheets",
"connected": True,
"spreadsheet_title": spreadsheet.get("properties", {}).get("title"),
"spreadsheet_id": spreadsheet_id,
"last_sync_at": self.config.get("last_sync_at"),
}
except Exception as exc:
return {
"service": "google_sheets",
"connected": False,
"error": str(exc),
}
# -- Register in the framework registry -------------------------------------
_REGISTRY["google_sheets"] = GoogleSheetsConnector