#!/usr/bin/env python3
"""Forward Claude Code session logs → Adapterly Ingest as OTLP traces.

Claude Code writes one JSONL file per session at
`~/.claude/projects/<encoded-cwd>/<session-uuid>.jsonl`. Each line is a
record of type assistant, user, progress, system, or file-history-snapshot.

We care only about `assistant` records — those carry model + usage. We
convert each to an OTel span, batch into one OTLP/JSON payload per source
file, and POST to the ingest endpoint.

Idempotency: the server uses (tenant_id, trace_id, span_id) as a primary
key with INSERT OR REPLACE, so re-running this script is a no-op for already-
sent events.

Usage:
    export INGEST_API_KEY=adk_...
    export INGEST_URL=https://adapterly.ai/ingest/v1/traces  (default)
    python3 claude_code_forwarder.py              # forward everything
    python3 claude_code_forwarder.py --dry-run    # print what would be sent
"""
from __future__ import annotations

import argparse
import hashlib
import json
import os
import sys
import urllib.error
import urllib.request
from datetime import datetime
from pathlib import Path


DEFAULT_URL = os.environ.get("INGEST_URL", "https://adapterly.ai/ingest/v1/traces")
DEFAULT_LOG_DIR = Path(os.environ.get("CLAUDE_LOG_DIR", str(Path.home() / ".claude" / "projects")))


def _hex16(s: str) -> str:
    """16-byte hex (32 chars) — valid OTel trace_id."""
    return hashlib.sha256(s.encode()).hexdigest()[:32]


def _hex8(s: str) -> str:
    """8-byte hex (16 chars) — valid OTel span_id."""
    return hashlib.sha256(s.encode()).hexdigest()[:16]


def _iso_to_nanos(iso: str) -> int | None:
    """Parse Claude Code's ISO timestamp to unix nanoseconds."""
    if not iso:
        return None
    try:
        s = iso.replace("Z", "+00:00")
        dt = datetime.fromisoformat(s)
        return int(dt.timestamp() * 1_000_000_000)
    except (ValueError, TypeError):
        return None


def _span_from_record(obj: dict) -> dict | None:
    """Convert a Claude Code `assistant` record to an OTel span dict."""
    if obj.get("type") != "assistant":
        return None
    msg = obj.get("message") or {}
    usage = msg.get("usage") or {}
    if not usage:
        return None

    session_id = obj.get("sessionId") or ""
    uuid = obj.get("uuid") or ""
    request_id = obj.get("requestId") or ""
    model = msg.get("model") or "unknown"
    ts = _iso_to_nanos(obj.get("timestamp"))
    if not ts:
        return None

    cc = usage.get("cache_creation") or {}
    attrs = [
        {"key": "gen_ai.system", "value": {"stringValue": "anthropic"}},
        {"key": "gen_ai.request.model", "value": {"stringValue": model}},
        {"key": "gen_ai.response.model", "value": {"stringValue": model}},
        {"key": "gen_ai.usage.input_tokens", "value": {"intValue": str(usage.get("input_tokens") or 0)}},
        {"key": "gen_ai.usage.output_tokens", "value": {"intValue": str(usage.get("output_tokens") or 0)}},
        {"key": "gen_ai.usage.cache_read_input_tokens", "value": {"intValue": str(usage.get("cache_read_input_tokens") or 0)}},
        {"key": "gen_ai.usage.cache_creation_5m_input_tokens", "value": {"intValue": str(cc.get("ephemeral_5m_input_tokens") or 0)}},
        {"key": "gen_ai.usage.cache_creation_1h_input_tokens", "value": {"intValue": str(cc.get("ephemeral_1h_input_tokens") or 0)}},
    ]
    if obj.get("isSidechain"):
        attrs.append({"key": "adapterly.sidechain", "value": {"boolValue": True}})

    return {
        "traceId": _hex16(session_id or uuid),
        "spanId": _hex8(request_id or uuid),
        "name": "anthropic.messages",
        "startTimeUnixNano": str(ts),
        "endTimeUnixNano": str(ts + 1),  # we don't know duration; +1ns keeps the span valid
        "attributes": attrs,
    }


def _resource_attrs(cwd: str | None) -> list[dict]:
    attrs = [{"key": "service.name", "value": {"stringValue": "claude-code"}}]
    if cwd:
        # Short friendly project name = last two path components
        parts = cwd.rstrip("/").split("/")
        project = "/".join(parts[-2:]) if len(parts) >= 2 else parts[-1]
        attrs.append({"key": "adapterly.project", "value": {"stringValue": project}})
    return attrs


def _file_to_payload(path: Path) -> dict | None:
    """Build one OTLP TracesData payload from one JSONL file."""
    spans_by_cwd: dict[str, list[dict]] = {}
    cwd = None
    with path.open("r", encoding="utf-8", errors="replace") as f:
        for line in f:
            line = line.strip()
            if not line:
                continue
            try:
                obj = json.loads(line)
            except json.JSONDecodeError:
                continue
            if obj.get("cwd"):
                cwd = obj["cwd"]
            span = _span_from_record(obj)
            if span is None:
                continue
            spans_by_cwd.setdefault(cwd or "", []).append(span)

    resource_spans = []
    for file_cwd, spans in spans_by_cwd.items():
        if not spans:
            continue
        resource_spans.append({
            "resource": {"attributes": _resource_attrs(file_cwd)},
            "scopeSpans": [{"spans": spans}],
        })
    if not resource_spans:
        return None
    return {"resourceSpans": resource_spans}


def _post(payload: dict, api_key: str, url: str) -> int:
    data = json.dumps(payload).encode()
    req = urllib.request.Request(
        url, data=data, method="POST",
        headers={"Content-Type": "application/json", "X-API-Key": api_key},
    )
    try:
        with urllib.request.urlopen(req, timeout=30) as r:
            return r.status
    except urllib.error.HTTPError as e:
        sys.stderr.write(f"HTTP {e.code}: {e.read().decode()[:300]}\n")
        raise


def main() -> None:
    p = argparse.ArgumentParser(description=__doc__.splitlines()[0])
    p.add_argument("--dry-run", action="store_true", help="Print counts without sending")
    p.add_argument("--log-dir", type=Path, default=DEFAULT_LOG_DIR)
    p.add_argument("--url", default=DEFAULT_URL)
    args = p.parse_args()

    api_key = os.environ.get("INGEST_API_KEY", "").strip()
    key_file = Path.home() / ".adapterly-ingest-key"
    if not api_key and key_file.exists():
        api_key = key_file.read_text().strip()
    if not api_key and not args.dry_run:
        sys.stderr.write("INGEST_API_KEY not set (env or ~/.adapterly-ingest-key)\n")
        sys.exit(2)

    files = sorted(args.log_dir.rglob("*.jsonl"))
    if not files:
        print(f"No .jsonl files under {args.log_dir}")
        return

    total_spans = 0
    total_files = 0
    total_sent = 0
    for path in files:
        payload = _file_to_payload(path)
        if not payload:
            continue
        n_spans = sum(len(rs["scopeSpans"][0]["spans"]) for rs in payload["resourceSpans"])
        total_spans += n_spans
        total_files += 1
        if args.dry_run:
            print(f"would send {n_spans:>4} spans from {path.name}")
            continue
        try:
            status = _post(payload, api_key, args.url)
        except Exception as e:
            print(f"FAILED  {path.name}: {e}")
            continue
        if status == 204:
            total_sent += n_spans
            print(f"sent    {n_spans:>4} spans from {path.name}")
        else:
            print(f"status {status}  {path.name}")

    print()
    print(f"Files scanned:   {total_files}")
    print(f"Spans found:     {total_spans}")
    if not args.dry_run:
        print(f"Spans sent:      {total_sent}")
        print(f"Endpoint:        {args.url}")


if __name__ == "__main__":
    main()
