#!/usr/bin/env python3
"""Runs the SEO cohort through every available route via the treg CLI and stores one raw
response per call with wall-clock latency. Deterministic; no LLM.
Usage: scripts/run-seo-cohort.py evidence/fixtures/seo-cohort-v1.json evidence/seo-runs/YYYY-MM-DD [--tasks kw,serp,bl]
Cost on 2026-09-15 list prices for v1: about $2.50 for the full run."""
import json, sys, os, subprocess, time, concurrent.futures as cf
cohort = json.load(open(sys.argv[1])); out = sys.argv[2]; os.makedirs(out, exist_ok=True)
tasks = (sys.argv[4] if len(sys.argv) > 4 and sys.argv[3] == '--tasks' else 'kw,serp,bl').split(',')
LOC = cohort['market']['location_code']
def call(label, argv, meta):
    f = os.path.join(out, label + '.json')
    if os.path.exists(f): return
    t0 = time.time()
    p = subprocess.run(['treg', 'call', *argv], capture_output=True, text=True)
    ms = int((time.time() - t0) * 1000)
    raw = p.stdout if p.stdout.strip() else p.stderr
    try: body = json.loads(raw)
    except Exception: body = {'_nonjson': raw[:2000]}
    rec = {'_run': {**meta, 'label': label, 'latency_ms': ms, 'exit': p.returncode, 'ts': time.strftime('%Y-%m-%dT%H:%M:%SZ', time.gmtime())}, 'response': body}
    json.dump(rec, open(f, 'w'))
    print(label, ms, 'ms', 'exit', p.returncode, flush=True)
jobs = []
Q = cohort['queries']; D = cohort['domains']
if 'kw' in tasks:
    # bulk providers: one call for all keywords (that is how buyers use them)
    jobs.append(('kw-dataforseo-all', ['dataforseo.google.keywords.volume', '--method', 'POST', '--data', json.dumps([{'keywords': Q, 'location_code': LOC, 'language_code': 'en'}])], {'task': 'kw', 'provider': 'dataforseo', 'n': len(Q)}))
    jobs.append(('kw-serpstat-all', ['serpstat.google.keywords.volume', '--method', 'POST', '--data', json.dumps({'id': '1', 'method': 'SerpstatKeywordProcedure.getKeywordsInfo', 'params': {'keywords': Q, 'se': 'g_us'}})], {'task': 'kw', 'provider': 'serpstat', 'n': len(Q)}))
    jobs.append(('kw-seranking-all', ['seranking.google.keywords.volume', '--method', 'POST', '--query', 'source=us', '--data', json.dumps({'keywords': Q})], {'task': 'kw', 'provider': 'seranking', 'n': len(Q)}))
    jobs.append(('kw-googleads-all', ['google-ads', 'v25/customers/2277522568:generateKeywordHistoricalMetrics', '--method', 'POST', '--data', json.dumps({'keywords': Q, 'geoTargetConstants': [f'geoTargetConstants/{LOC}'], 'language': 'languageConstants/1000', 'keywordPlanNetwork': 'GOOGLE_SEARCH'})], {'task': 'kw', 'provider': 'google-ads', 'n': len(Q)}))
if 'serp' in tasks:
    for i, q in enumerate(Q):
        k = f'q{i+1:02d}'
        jobs.append((f'serp-cloro-{k}', ['cloro.google.serp.organic', '--method', 'POST', '--data', json.dumps({'query': q, 'gl': 'US', 'pages': 1})], {'task': 'serp', 'provider': 'cloro', 'q': q}))
        jobs.append((f'serp-dataforseo-{k}', ['dataforseo.x.serp-google-organic-live-advanced', '--method', 'POST', '--data', json.dumps([{'language_code': 'en', 'location_code': LOC, 'keyword': q, 'depth': 10}])], {'task': 'serp', 'provider': 'dataforseo', 'q': q}))
        jobs.append((f'serp-serpstat-{k}', ['serpstat.google.serp.organic', '--method', 'POST', '--data', json.dumps({'id': '1', 'method': 'SerpstatKeywordProcedure.getKeywordTop', 'params': {'keyword': q, 'se': 'g_us', 'size': 10}})], {'task': 'serp', 'provider': 'serpstat', 'q': q}))
        jobs.append((f'serp-serpapi-{k}', ['serpapi.google.serp.organic', '--query', 'engine=google', '--query', f'q={q}', '--query', 'gl=us', '--query', 'hl=en', '--query', 'num=10'], {'task': 'serp', 'provider': 'serpapi', 'q': q}))
        jobs.append((f'serp-treg-{k}', ['treg.google.serp.organic', '--method', 'POST', '--data', json.dumps({'q': q})], {'task': 'serp', 'provider': 'treg-routed', 'q': q}))
if 'bl' in tasks:
    for i, d in enumerate(D):
        k = f'd{i+1:02d}'
        jobs.append((f'bl-serpstat-{k}', ['serpstat.web.backlinks.summary', '--method', 'POST', '--data', json.dumps({'id': '1', 'method': 'SerpstatBacklinksProcedure.getSummaryV2', 'params': {'query': d}})], {'task': 'bl', 'provider': 'serpstat', 'd': d}))
        jobs.append((f'bl-seranking-{k}', ['seranking.web.backlinks.summary', '--query', f'target={d}', '--query', 'mode=domain'], {'task': 'bl', 'provider': 'seranking', 'd': d}))
        jobs.append((f'bl-moz-{k}', ['moz.web.backlinks.summary', '--method', 'POST', '--data', json.dumps({'targets': [d]})], {'task': 'bl', 'provider': 'moz', 'd': d}))
# run: one worker per provider so latency is not distorted by our own concurrency
by_prov = {}
for j in jobs: by_prov.setdefault(j[2]['provider'], []).append(j)
def run_prov(lst):
    for label, argv, meta in lst: call(label, argv, meta)
with cf.ThreadPoolExecutor(max_workers=len(by_prov)) as ex:
    list(ex.map(run_prov, by_prov.values()))
print('done', len(jobs), 'jobs')
