#!/usr/bin/env python3
"""PumpGTM experimental Events polling demo. Python 3.9+, standard library only.
Creates local review tasks, never contacts prospects or writes to a CRM.
Keep this SQLite file private; it contains your selected job observations.
"""
import argparse
import hashlib
import json
import os
import sqlite3
import time
from urllib.request import Request, urlopen
from urllib.error import HTTPError


def request(url, key, body):
    req = Request(url, data=json.dumps(body).encode(), headers={
        "Authorization": "Bearer " + key, "Content-Type": "application/json",
        "Accept": "application/json, text/event-stream", "MCP-Protocol-Version": "2025-03-26",
    })
    with urlopen(req, timeout=60) as response:
        raw = response.read().decode()
    if raw.lstrip().startswith("{"):
        return json.loads(raw)
    for line in raw.splitlines():
        if line.startswith("data:"):
            parsed = json.loads(line[5:])
            if "result" in parsed or "error" in parsed:
                return parsed
    raise RuntimeError("No MCP result in response")


def main():
    parser = argparse.ArgumentParser(description=__doc__)
    parser.add_argument("--url", default="https://app.pumpgtm.com/api/v1/mcp")
    parser.add_argument("--roles", default="Sales Development Representative", help="Comma-separated role phrases")
    parser.add_argument("--state", default="hiring-events.sqlite3")
    parser.add_argument("--once", action="store_true")
    parser.add_argument("--recent", action="store_true", help="Backfill on first run using the ordinary MCP tool")
    parser.add_argument("--replay", action="store_true", help="Repeat the previous request once; do not advance the saved checkpoint")
    args = parser.parse_args()
    key = os.environ.get("PUMPGTM_KEY")
    if not key:
        parser.error("Set PUMPGTM_KEY to your workspace key")
    if not args.url.startswith("https://") and not args.url.startswith("http://localhost:"):
        parser.error("Use HTTPS, or localhost for development")
    filters = {"roles": [role.strip() for role in args.roles.split(",") if role.strip()]}
    identity = hashlib.sha256(json.dumps([args.url, filters, hashlib.sha256(key.encode()).hexdigest()], sort_keys=True).encode()).hexdigest()
    os.umask(0o077)
    db = sqlite3.connect(args.state, timeout=60)
    db.execute("create table if not exists state (id integer primary key check(id=1), identity text not null, cursor text, last_request text)")
    db.execute("create table if not exists review_task (event_id text primary key, payload text not null)")
    db.commit()
    while True:
        # Hold a SQLite write lock through the read and commit. Two processes
        # sharing this state cannot move the checkpoint past each other.
        db.execute("begin immediate")
        try:
            saved = db.execute("select identity,cursor,last_request from state where id=1").fetchone()
            if saved and saved[0] != identity:
                raise RuntimeError("Endpoint, filters, or key changed. Use a separate --state file.")
            if args.replay:
                if not saved or not saved[2]:
                    raise RuntimeError("No previous request to replay")
                body = json.loads(saved[2])
            elif not saved and args.recent:
                body = {"jsonrpc": "2.0", "id": 1, "method": "tools/call", "params": {"name": "list_signal_events", "arguments": {"name": "hiring.job_observed", "filters": filters, "start": "recent", "limit": 50}}}
            else:
                body = {"jsonrpc": "2.0", "id": 1, "method": "events/poll", "params": {"name": "hiring.job_observed", "arguments": filters, "cursor": saved[1] if saved else None, "maxEvents": 50}}
            envelope = request(args.url, key, body)
            if "error" in envelope:
                raise RuntimeError(json.dumps(envelope["error"]))
            result = envelope["result"]
            if body["method"] == "tools/call":
                if result.get("isError"):
                    raise RuntimeError(json.dumps(result.get("structuredContent", result)))
                result = result["structuredContent"]
            if result.get("truncated"):
                raise RuntimeError("Replay gap reported; inspect before advancing the cursor")
            inserted = 0
            for event in result["events"]:
                inserted += db.execute("insert or ignore into review_task values (?,?)", (event["eventId"], json.dumps(event))).rowcount
            if not args.replay:
                db.execute("insert into state values (1,?,?,?) on conflict(id) do update set cursor=excluded.cursor,last_request=excluded.last_request", (identity, result["cursor"], json.dumps(body)))
            total = db.execute("select count(*) from review_task").fetchone()[0]
            db.commit()
            print(f"Received {len(result['events'])}; new review tasks {inserted}; total tasks {total}; Energy 0")
        except Exception:
            db.rollback()
            raise
        if args.once or args.replay:
            break
        time.sleep(1 if result.get("hasMore") else max(60, result.get("nextPollMs", 60000) / 1000))


if __name__ == "__main__":
    try:
        main()
    except HTTPError as error:
        raise SystemExit(f"HTTP {error.code}: {error.read().decode()[:500]}")
    except (RuntimeError, OSError, ValueError) as error:
        raise SystemExit(str(error))
