Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions .gitignore
Original file line number Diff line number Diff line change
Expand Up @@ -36,3 +36,4 @@ saved-config.yaml
*.db
*.db-wal
*.db-shm
uv.lock
2 changes: 1 addition & 1 deletion .pre-commit-config.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -23,7 +23,7 @@ repos:
hooks:
- id: mypy
name: mypy
entry: python -m mypy
entry: .venv/bin/python -m mypy --python-version 3.12
language: system
pass_filenames: false
types: [python]
72 changes: 70 additions & 2 deletions src/tokenops/control/http.py
Original file line number Diff line number Diff line change
Expand Up @@ -6,14 +6,18 @@

from __future__ import annotations

import csv
import io
import json
import os
from collections.abc import Awaitable, Callable, Mapping
from dataclasses import asdict
from typing import Any

import httpx
from chronicle.session import reset_session
from fastapi import FastAPI, Request
from fastapi.responses import JSONResponse, Response
from fastapi import FastAPI, Query, Request
from fastapi.responses import JSONResponse, Response, StreamingResponse

from tokenops.control.core import Halt
from tokenops.control.engine import Throttled
Expand Down Expand Up @@ -65,6 +69,70 @@ async def register_run(request: Request) -> JSONResponse:
)


_EXPORT_CSV_COLUMNS = [
"run_id",
"agent",
"status",
"cost_micros",
"steps",
"started_at",
"ended_at",
"duration_s",
"dims",
"halt_reason",
"detector",
"governance_events",
]


def _run_to_csv_row(rec): # type: ignore[no-untyped-def]
d = asdict(rec)
d["duration_s"] = round(rec.ended_at - rec.started_at, 2) if rec.ended_at else ""
d["dims"] = json.dumps(rec.dims) if rec.dims else ""
d["governance_events"] = json.dumps(rec.governance_events) if rec.governance_events else ""
return [d.get(col, "") for col in _EXPORT_CSV_COLUMNS]


def mount_export(app: FastAPI, store: Store) -> None:
"""Mount ``GET /v1/export`` — on-demand run-record export (CSV / JSON)."""

@app.get("/v1/export")
def export_runs(
from_at: float | None = Query(None, description="Start timestamp (epoch seconds)"),
to_at: float | None = Query(None, description="End timestamp (epoch seconds)"),
agent: str | None = Query(None, description="Filter by agent name"),
status: str | None = Query(None, description="Filter by run status"),
tenant: str | None = Query(None, description="Filter by tenant (from dims)"),
format: str = Query("json", description="Output format: json or csv"),
limit: int = Query(5000, ge=1, le=10_000, description="Max rows to return"),
) -> Response:
runs = store.export_runs(
from_at=from_at,
to_at=to_at,
agent=agent,
status=status,
tenant=tenant,
limit=limit,
)

if format == "csv":
buf = io.StringIO()
writer = csv.writer(buf)
writer.writerow(_EXPORT_CSV_COLUMNS)
for rec in runs:
writer.writerow(_run_to_csv_row(rec))
buf.seek(0)
return StreamingResponse(
iter([buf.getvalue()]),
media_type="text/csv",
headers={"Content-Disposition": 'attachment; filename="export.csv"'},
)

# JSON (default)
rows = [asdict(rec) for rec in runs]
return JSONResponse(rows, headers={"X-Total-Count": str(len(rows))})


def with_governance_errors(handler: Handler) -> Handler:
"""Wrap a task handler so Halt → 200 halted and Throttled → 429."""

Expand Down
64 changes: 64 additions & 0 deletions src/tokenops/control/store.py
Original file line number Diff line number Diff line change
Expand Up @@ -30,6 +30,7 @@
import time
import uuid
from collections.abc import Callable
from datetime import date, datetime
from typing import Any, TypeVar

from tokenops.control.ledger import LIFETIME, RUN_TOTAL_BUDGET
Expand Down Expand Up @@ -64,6 +65,26 @@ def _known_policy_templates() -> frozenset[str]:
return frozenset({*_TEMPLATES, "trajectory_hint"})


def _coerce_epoch(value: float | str | date | None) -> float | None:
"""Coerce a timestamp value to epoch float.

Accepts epoch floats (pass-through), ISO date strings (``"2026-09-19"``),
``datetime.date`` objects, or ``None``. Dates are converted to start-of-day
epoch seconds so they compare correctly against the ``started_at`` REAL column.
"""
if value is None:
return None
if isinstance(value, int | float):
return float(value)
if isinstance(value, date) and not isinstance(value, datetime):
return datetime.combine(value, datetime.min.time()).timestamp()
if isinstance(value, str):
return datetime.fromisoformat(value).timestamp()
if isinstance(value, datetime):
return value.timestamp()
return float(value)


_SCHEMA = """
CREATE TABLE IF NOT EXISTS segments (
id TEXT PRIMARY KEY, name TEXT NOT NULL, dimension TEXT NOT NULL,
Expand Down Expand Up @@ -519,6 +540,49 @@ def list_runs(self, *, problematic_only: bool = False, limit: int = 200) -> list
sql += " ORDER BY started_at DESC LIMIT ?"
return [self._run_with_ledger_cost(r) for r in self._db.execute(sql, (limit,))]

@_locked
def export_runs(
self,
*,
from_at: float | None = None,
to_at: float | None = None,
agent: str | None = None,
status: str | None = None,
tenant: str | None = None,
limit: int = 5000,
) -> list[RunRecord]:
"""Export run-records filtered by time range, agent, status, or tenant.

Used by the ``GET /v1/export`` route for FinOps / chargeback CSV and JSON.
Returns at most *limit* rows (capped at 10 000) ordered by ``started_at`` DESC.
"""
limit = min(limit, 10_000)
from_at = _coerce_epoch(from_at)
to_at = _coerce_epoch(to_at)
where: list[str] = []
params: list[object] = []
if from_at is not None:
where.append("started_at >= ?")
params.append(from_at)
if to_at is not None:
where.append("started_at <= ?")
params.append(to_at)
if agent is not None:
where.append("agent = ?")
params.append(agent)
if status is not None:
where.append("status = ?")
params.append(status)
if tenant is not None:
where.append("json_extract(dims, '$.tenant') = ?")
params.append(tenant)
sql = "SELECT * FROM runs"
if where:
sql += " WHERE " + " AND ".join(where)
sql += " ORDER BY started_at DESC LIMIT ?"
params.append(limit)
return [self._run_with_ledger_cost(r) for r in self._db.execute(sql, tuple(params))]

def _run_with_ledger_cost(self, row: sqlite3.Row) -> RunRecord:
"""Build a RunRecord; prefer ``__run_total__`` ledger spend when present."""
rec = _run(row)
Expand Down
3 changes: 2 additions & 1 deletion src/tokenops/server/app.py
Original file line number Diff line number Diff line change
Expand Up @@ -7,7 +7,7 @@
from fastapi import FastAPI

from tokenops import __version__
from tokenops.control.http import mount_run_registration
from tokenops.control.http import mount_export, mount_run_registration
from tokenops.control.store import Store


Expand All @@ -32,6 +32,7 @@ def health() -> dict[str, str]:
return {"status": "ok", "service": "tokenops-control-plane"}

mount_run_registration(app, store)
mount_export(app, store)

# Placeholder for future plane APIs (observe, governance admin over HTTP, etc.).
# Agents keep using ControlPlaneClient; expand the plane surface here.
Expand Down
87 changes: 87 additions & 0 deletions src/tokenops/ui/views/dashboard.py
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,12 @@

from __future__ import annotations

import csv
import io
import json
from dataclasses import asdict
from datetime import datetime, time

import altair as alt
import pandas as pd
import streamlit as st
Expand Down Expand Up @@ -169,6 +175,87 @@ def _seg(r) -> str:
)
st.dataframe(table, use_container_width=True, hide_index=True)

# ---- export (CSV / JSON download) ---------------------------------------- #
_EXPORT_COLUMNS = [
"run_id",
"agent",
"status",
"cost_micros",
"steps",
"started_at",
"ended_at",
"duration_s",
"dims",
"halt_reason",
"detector",
"governance_events",
]


def _to_csv(runs_list):
buf = io.StringIO()
writer = csv.writer(buf)
writer.writerow(_EXPORT_COLUMNS)
for rec in runs_list:
d = asdict(rec)
d["duration_s"] = round(rec.ended_at - rec.started_at, 2) if rec.ended_at else ""
d["dims"] = json.dumps(rec.dims) if rec.dims else ""
d["governance_events"] = json.dumps(rec.governance_events) if rec.governance_events else ""
writer.writerow([d.get(c, "") for c in _EXPORT_COLUMNS])
return buf.getvalue()


def _date_to_epoch(d, end_of_day: bool = False) -> float:
"""Convert a date to epoch seconds. ``end_of_day`` returns 23:59:59.999."""
t = time.max if end_of_day else time.min
return datetime.combine(d, t).timestamp()


with st.expander("Export run data", expanded=False):
agents = sorted({r.agent for r in runs})
ecol1, ecol2, ecol3, ecol4 = st.columns(4)
with ecol1:
export_from = st.date_input("From", value=None, key="export_from")
with ecol2:
export_to = st.date_input("To", value=None, key="export_to")
with ecol3:
export_agent = st.selectbox("Agent", ["All"] + agents, key="export_agent")
with ecol4:
export_status = st.selectbox(
"Status",
["All", "completed", "halted", "error", "running"],
key="export_status",
)

export_runs = store.export_runs(
from_at=_date_to_epoch(export_from) if export_from else None,
to_at=_date_to_epoch(export_to, end_of_day=True) if export_to else None,
agent=export_agent if export_agent != "All" else None,
status=export_status if export_status != "All" else None,
limit=10_000,
)

st.caption(f"{len(export_runs)} runs matched")
st.markdown(
"<style>div[data-testid='stHorizontalBlock']{gap:0.5rem}</style>",
unsafe_allow_html=True,
)
dl_col1, dl_col2, _ = st.columns([1, 1, 6])
with dl_col1:
st.download_button(
"Download CSV",
data=_to_csv(export_runs),
file_name="export.csv",
mime="text/csv",
)
with dl_col2:
st.download_button(
"Download JSON",
data=json.dumps([asdict(r) for r in export_runs], default=str, indent=2),
file_name="export.json",
mime="application/json",
)

# ---- run detail picker (when not already focused) ------------------------ #
if not focus_run:
st.subheader("Run detail")
Expand Down
Loading
Loading