Skip to content

OTEL Collector

otel

OTLP HTTP/JSON receiver and token/cost summary.

Lightweight collector that accepts OTLP exports for metrics and logs, tracks token usage over a sliding window, and prints a summary.

start_collector(run_dir, bind_addr='127.0.0.1')

Start the OTEL collector as a subprocess. Returns (proc, port).

Source code in src/agentic_ci/otel.py
def start_collector(run_dir, bind_addr="127.0.0.1"):
    """Start the OTEL collector as a subprocess. Returns (proc, port)."""
    otel_log = os.path.join(run_dir, "claude-otel.jsonl")
    otel_rate = os.path.join(run_dir, "claude-otel-rate.json")
    port_file = os.path.join(run_dir, "otel-port")

    for f in [otel_log, port_file]:
        try:
            os.unlink(f)
        except FileNotFoundError:
            pass

    env = {
        **os.environ,
        "OTEL_LOG_FILE": otel_log,
        "OTEL_RATE_FILE": otel_rate,
        "OTEL_COLLECTOR_PORT": "0",
        "OTEL_PORT_FILE": port_file,
        "OTEL_BIND_ADDR": bind_addr,
    }
    proc = subprocess.Popen(
        [sys.executable, "-m", "agentic_ci.otel"],
        env=env,
        stderr=subprocess.DEVNULL,
    )

    for _ in range(50):
        if os.path.exists(port_file):
            break
        time.sleep(0.1)
    else:
        proc.kill()
        raise RuntimeError("OTEL collector did not write port file")

    with open(port_file) as f:
        port = int(f.read().strip())

    return proc, port, otel_log, otel_rate

stop_collector(proc)

Stop the OTEL collector subprocess.

Source code in src/agentic_ci/otel.py
def stop_collector(proc):
    """Stop the OTEL collector subprocess."""
    proc.terminate()
    try:
        proc.wait(timeout=5)
    except subprocess.TimeoutExpired:
        proc.kill()
        proc.wait()

parse_metrics(records)

Parse OTLP JSONL records into structured token/cost data.

Source code in src/agentic_ci/otel.py
def parse_metrics(records):
    """Parse OTLP JSONL records into structured token/cost data."""
    token_totals = defaultdict(float)
    cost_totals = defaultdict(float)
    api_requests = []
    codex_api_requests = []
    codex_response_requests = []
    active_time = defaultdict(float)

    for rec in records:
        path = rec.get("path", "")
        payload = rec.get("payload", {})

        if "/v1/metrics" in path:
            for rm in payload.get("resourceMetrics", []):
                for sm in rm.get("scopeMetrics", []):
                    for metric in sm.get("metrics", []):
                        name = metric.get("name", "")
                        data = metric.get("sum", metric.get("gauge", metric.get("histogram", {})))
                        for dp in data.get("dataPoints", []):
                            attrs = {
                                a.get("key"): _otel_attribute_value(a.get("value"))
                                for a in dp.get("attributes", [])
                                if a.get("key")
                            }
                            value = dp.get("asDouble", dp.get("asInt", 0))

                            if name == "claude_code.token.usage":
                                model = attrs.get("model", "unknown")
                                token_type = attrs.get("type", "unknown")
                                token_totals[(model, token_type)] += value
                            elif name == "claude_code.cost.usage":
                                model = attrs.get("model", "unknown")
                                cost_totals[model] += value
                            elif name == "claude_code.active_time.total":
                                time_type = attrs.get("type", "unknown")
                                active_time[time_type] += value

        elif "/v1/logs" in path:
            for rl in payload.get("resourceLogs", []):
                for sl in rl.get("scopeLogs", []):
                    for lr in sl.get("logRecords", []):
                        event_name = ""
                        event_attrs = {}
                        for a in lr.get("attributes", []):
                            key = a.get("key")
                            if not key:
                                continue
                            v = _otel_attribute_value(a.get("value"))
                            event_attrs[key] = v
                            if key == "event.name":
                                event_name = v
                        if event_name == "claude_code.api_request":
                            api_requests.append(event_attrs)
                        elif event_name == "codex.api_request":
                            api_requests.append(event_attrs)
                            codex_api_requests.append(event_attrs)
                        elif (
                            event_name == "codex.sse_event"
                            and event_attrs.get("event.kind") == "response.completed"
                        ):
                            codex_response_requests.append(event_attrs)
                            model = str(event_attrs.get("model") or "unknown")
                            input_tokens = _numeric_attribute(event_attrs, "input_token_count")
                            cached_tokens = _numeric_attribute(event_attrs, "cached_token_count")
                            cache_write_tokens = _numeric_attribute(
                                event_attrs, "cache_write_token_count"
                            )
                            # Codex reports input_token_count inclusive of the
                            # cached and cache-write tokens, so subtract them to
                            # get the fresh prompt tokens. The max(..., 0) guards
                            # against exporters that report them separately (in
                            # which case fresh_input clamps to 0 rather than
                            # going negative).
                            fresh_input = max(
                                input_tokens - cached_tokens - cache_write_tokens,
                                0,
                            )
                            token_totals[(model, "input")] += fresh_input
                            token_totals[(model, "cacheRead")] += cached_tokens
                            token_totals[(model, "cacheCreation")] += cache_write_tokens
                            token_totals[(model, "output")] += _numeric_attribute(
                                event_attrs, "output_token_count"
                            )

                            direct_cost = _optional_numeric_attribute(
                                event_attrs,
                                "cost_usd",
                                "total_cost_usd",
                                "cost",
                                "total_cost",
                            )
                            if direct_cost is not None and direct_cost >= 0:
                                cost_totals[model] += direct_cost
                            else:
                                service_tier = event_attrs.get("service_tier")
                                if not isinstance(service_tier, str):
                                    service_tier = None
                                estimated_cost = estimate_cost(
                                    model,
                                    input_tokens=fresh_input,
                                    output_tokens=_numeric_attribute(
                                        event_attrs, "output_token_count"
                                    ),
                                    cached_input_tokens=cached_tokens,
                                    cache_creation_input_tokens=cache_write_tokens,
                                    service_tier=service_tier,
                                )
                                if estimated_cost is not None:
                                    cost_totals[model] += estimated_cost

    # Recent Codex versions emit response.completed events but no separate
    # codex.api_request log event. Keep the explicit event when available and
    # use completed responses as a conservative request-count fallback.
    if not codex_api_requests:
        api_requests.extend(codex_response_requests)

    return token_totals, cost_totals, api_requests, active_time

print_summary(log_file)

Print a human-readable token/cost summary from an OTEL JSONL log.

Source code in src/agentic_ci/otel.py
def print_summary(log_file):
    """Print a human-readable token/cost summary from an OTEL JSONL log."""
    records = []
    try:
        with open(log_file) as f:
            for line in f:
                line = line.strip()
                if line:
                    records.append(json.loads(line))
    except FileNotFoundError:
        print("No OTEL data collected (log file not found).")
        return
    except json.JSONDecodeError as e:
        print(f"Error parsing OTEL log: {e}")
        return

    if not records:
        print("No OTEL data collected.")
        return

    token_totals, cost_totals, api_requests, active_time = parse_metrics(records)

    if token_totals:
        models = sorted(set(m for m, _ in token_totals.keys()))
        for model in models:
            print(f"\n  Model: {model}")
            print(f"  {'Token Type':<20} {'Count':>12}")
            print(f"  {'-' * 20} {'-' * 12}")
            model_tokens = {t: c for (m, t), c in token_totals.items() if m == model}
            for token_type in ["input", "cacheRead", "cacheCreation", "output"]:
                if token_type in model_tokens:
                    print(f"  {token_type:<20} {model_tokens[token_type]:>12,.0f}")
            total = sum(model_tokens.values())
            print(f"  {'TOTAL':<20} {total:>12,.0f}")

    if cost_totals:
        print(f"\n  {'Model':<30} {'Cost (USD)':>12}")
        print(f"  {'-' * 30} {'-' * 12}")
        grand_total = 0.0
        for model in sorted(cost_totals.keys()):
            cost = cost_totals[model]
            grand_total += cost
            print(f"  {model:<30} ${cost:>11.4f}")
        if len(cost_totals) > 1:
            print(f"  {'TOTAL':<30} ${grand_total:>11.4f}")

    if active_time:
        print("\n  Active Time:")
        for time_type, seconds in sorted(active_time.items()):
            mins, secs = divmod(int(seconds), 60)
            print(f"    {time_type}: {mins}m {secs}s")

    if api_requests:
        print(f"\n  API Requests: {len(api_requests)}")
        total_duration = sum(float(r.get("duration_ms", 0)) for r in api_requests)
        if total_duration:
            print(f"  Total API time: {total_duration / 1000:.1f}s")

generate_trace_context()

Generate W3C Trace Context components for orchestrator-owned root spans.

Returns (trace_id, span_id, traceparent) where traceparent is a valid W3C traceparent header value. If the agent's OTEL SDK picks up the TRACEPARENT env var, its spans become children of this root span.

Source code in src/agentic_ci/otel.py
def generate_trace_context():
    """Generate W3C Trace Context components for orchestrator-owned root spans.

    Returns (trace_id, span_id, traceparent) where traceparent is a valid
    W3C traceparent header value. If the agent's OTEL SDK picks up the
    TRACEPARENT env var, its spans become children of this root span.
    """
    trace_id = uuid.uuid4().hex
    span_id = uuid.uuid4().hex[:16]
    traceparent = f"00-{trace_id}-{span_id}-01"
    return trace_id, span_id, traceparent

root_span_attributes(backend, harness, model, effort=None, subagent_effort=None)

Return the synthetic root span attributes for one agent run.

Always sets agent.backend, agent.harness and agent.model; adds agent.reasoning_effort and agent.subagent_reasoning_effort when those efforts are in effect, so MLflow and JSONL consumers can see the effective effort next to the model.

Source code in src/agentic_ci/otel.py
def root_span_attributes(
    backend: str,
    harness: str,
    model: str,
    effort: str | None = None,
    subagent_effort: str | None = None,
) -> dict[str, str]:
    """Return the synthetic root span attributes for one agent run.

    Always sets ``agent.backend``, ``agent.harness`` and ``agent.model``;
    adds ``agent.reasoning_effort`` and ``agent.subagent_reasoning_effort``
    when those efforts are in effect, so MLflow and JSONL consumers can see
    the effective effort next to the model.
    """
    attributes = {"agent.backend": backend, "agent.harness": harness, "agent.model": model}
    if effort is not None:
        attributes["agent.reasoning_effort"] = effort
    if subagent_effort is not None:
        attributes["agent.subagent_reasoning_effort"] = subagent_effort
    return attributes

inject_root_spans(log_file, start_ns, end_ns, exit_code, fallback_trace_id=None, fallback_span_id=None, attributes=None)

Stitch the run's traces together and append a synthetic root.

All traces in the JSONL belong to one agent run, so they are first merged into a single trace (see :func:_stitch_run_traces). Then, for each trace that has child spans but no root span (agent ignored TRACEPARENT, container crash, OOM, timeout), a synthetic root span is appended to the JSONL file. If no spans exist at all and fallback_trace_id is provided, a standalone root span is created so MLflow always has at least one complete trace. When fallback_span_id is set, the fallback root reuses that span ID so children that were parented via TRACEPARENT stay connected.

Returns the number of synthetic root spans injected.

Source code in src/agentic_ci/otel.py
def inject_root_spans(
    log_file,
    start_ns,
    end_ns,
    exit_code,
    fallback_trace_id=None,
    fallback_span_id=None,
    attributes=None,
):
    """Stitch the run's traces together and append a synthetic root.

    All traces in the JSONL belong to one agent run, so they are first merged
    into a single trace (see :func:`_stitch_run_traces`). Then, for each trace
    that has child spans but no root span (agent ignored TRACEPARENT,
    container crash, OOM, timeout), a synthetic root span is appended to the
    JSONL file. If no spans exist at all and fallback_trace_id is provided, a
    standalone root span is created so MLflow always has at least one
    complete trace. When fallback_span_id is set, the fallback root reuses
    that span ID so children that were parented via TRACEPARENT stay
    connected.

    Returns the number of synthetic root spans injected.
    """
    records = []
    try:
        with open(log_file) as f:
            for line in f:
                line = line.strip()
                if not line:
                    continue
                try:
                    records.append(json.loads(line))
                except json.JSONDecodeError:
                    continue
    except FileNotFoundError:
        pass

    primary, _ = _stitch_run_traces(records, fallback_trace_id, fallback_span_id)
    if primary:
        _write_records(log_file, records)

    orphans = _find_orphan_traces(records)

    injected = 0
    root_ids = {}

    for trace_id, (child_start, child_end, dangling_parent) in orphans.items():
        if dangling_parent:
            span_id = dangling_parent
        elif fallback_span_id:
            span_id = fallback_span_id
        else:
            span_id = uuid.uuid4().hex[:16]
        root_ids[trace_id] = span_id

    # An agent can drop more than one parent span (Codex loses its exec root
    # and a session span). The synthetic root reuses the most common missing
    # parent; every other span whose parent was never exported is re-homed
    # under it so nothing renders detached. Traces that already have a real
    # root get their dangling spans re-homed under that root instead.
    for trace_id, root_span_id in _real_root_ids(records).items():
        root_ids.setdefault(trace_id, root_span_id)
    if _rehome_dangling_spans(records, root_ids):
        _write_records(log_file, records)

    for trace_id, (child_start, child_end, _dangling_parent) in orphans.items():
        span_id = root_ids[trace_id]
        span_start = min(start_ns, child_start) if child_start else start_ns
        span_end = max(end_ns, child_end) if child_end else end_ns
        record = _build_root_span_record(
            trace_id, span_id, span_start, span_end, exit_code, attributes
        )
        with open(log_file, "a") as f:
            f.write(json.dumps(record) + "\n")
        injected += 1

    if not orphans and not _has_any_traces(records) and fallback_trace_id:
        span_id = fallback_span_id if fallback_span_id else uuid.uuid4().hex[:16]
        record = _build_root_span_record(
            fallback_trace_id, span_id, start_ns, end_ns, exit_code, attributes
        )
        with open(log_file, "a") as f:
            f.write(json.dumps(record) + "\n")
        injected += 1

    return injected

main()

Run the OTEL collector server.

Source code in src/agentic_ci/otel.py
def main():
    """Run the OTEL collector server."""
    port = int(os.environ.get("OTEL_COLLECTOR_PORT", "4318"))
    bind_addr = os.environ.get("OTEL_BIND_ADDR", "127.0.0.1")
    server = HTTPServer((bind_addr, port), OTLPHandler)
    actual_port = server.server_address[1]
    port_file = os.environ.get("OTEL_PORT_FILE")
    if port_file:
        with open(port_file, "w") as f:
            f.write(str(actual_port))
    signal.signal(signal.SIGTERM, lambda *_: sys.exit(0))
    log_file = os.environ.get("OTEL_LOG_FILE", "/tmp/claude-otel.jsonl")
    print(
        f"OTLP collector listening on {bind_addr}:{actual_port}, writing to {log_file}",
        file=sys.stderr,
    )
    try:
        server.serve_forever()
    except KeyboardInterrupt:
        pass
    server.server_close()