Skip to content
Merged
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
44 changes: 43 additions & 1 deletion fintick/aggregate.py
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,7 @@
POST_AGGREGATION_MAX_ATTEMPTS,
V2Event,
open_database,
record_inference_usage,
set_post_aggregation_decision,
upsert_event,
)
Expand All @@ -36,7 +37,21 @@
# must leave room for both. 4096 gets fully consumed by reasoning on hard batches -> empty content;
# 16384 leaves ample headroom (observed ~5-7k total). Tunable per-model via the env var.
LLM_MAX_TOKENS = int(os.environ.get("FINTICK_LLM_MAX_TOKENS", "16384"))
# USD per 1M tokens (input, output). Cloud models are priced here; local/unknown
# models are free. Keys are matched case-insensitively. Extend as models are added.
LLM_PRICES = {"gpt-5.6-luna": (0.20, 1.20)}
WINDOW = timedelta(hours=6)


def inference_cost_usd(
model: str | None, prompt_tokens: int, completion_tokens: int
) -> float:
"""Dollar cost of one call's tokens for the given model (0 for local/unknown)."""
price = LLM_PRICES.get((model or "").lower())
if not price:
return 0.0
price_in, price_out = price
return (prompt_tokens / 1_000_000) * price_in + (completion_tokens / 1_000_000) * price_out
MAX_POSTS = 200
DEFAULT_BATCH = 50

Expand Down Expand Up @@ -419,6 +434,7 @@ def call_inference(
api_key: str | None = None,
model: str | None = None,
timeout: float = 300.0,
usage_sink: Callable[[dict[str, Any]], None] | None = None,
) -> str:
"""Make one forced-JSON call to an OpenAI-compatible chat endpoint.

Expand Down Expand Up @@ -449,10 +465,19 @@ def call_inference(
},
method="POST",
)
usage_payload: dict[str, Any] | None = None
try:
with urllib.request.urlopen(request, timeout=timeout) as response:
data = json.load(response)
content = data["choices"][0]["message"]["content"]
if usage_sink is not None:
usage = data.get("usage") if isinstance(data, dict) else None
usage = usage if isinstance(usage, dict) else {}
usage_payload = {
"model": model,
"prompt_tokens": int(usage.get("prompt_tokens") or 0),
"completion_tokens": int(usage.get("completion_tokens") or 0),
}
except (
OSError,
urllib.error.URLError,
Expand All @@ -462,6 +487,13 @@ def call_inference(
IndexError,
) as error:
raise RuntimeError(f"aggregation request failed: {error}") from error
# Record usage even when the content is empty — those tokens were still spent.
# A logging failure must never fail the aggregation call.
if usage_sink is not None and usage_payload is not None:
try:
usage_sink(usage_payload)
except Exception:
pass
if not isinstance(content, str) or not content.strip():
raise RuntimeError("aggregation model returned empty or non-text content")
return content
Expand Down Expand Up @@ -495,8 +527,18 @@ def aggregate_once(

try:
if call_model is None:
def _record_usage(usage: dict[str, Any]) -> None:
try:
with open_database(database) as connection:
record_inference_usage(
connection, usage["model"],
usage["prompt_tokens"], usage["completion_tokens"],
)
except Exception:
pass
call_model = lambda value: call_inference(
value, base_url=base_url, api_key=api_key, model=model
value, base_url=base_url, api_key=api_key, model=model,
usage_sink=_record_usage,
)
raw_content = call_model(prompt)
raw_object = _json_object(raw_content)
Expand Down
29 changes: 26 additions & 3 deletions fintick/dashboard.py
Original file line number Diff line number Diff line change
Expand Up @@ -12,8 +12,14 @@
from typing import Any, cast
from urllib.parse import parse_qs, urlsplit

from fintick.aggregate import inference_cost_usd
from fintick.service_handoff import database_identity
from fintick.storage import load_events, load_pipeline_health, open_database
from fintick.storage import (
load_events,
load_inference_usage,
load_pipeline_health,
open_database,
)

DEFAULT_LIMIT = 100
MAX_LIMIT = 250
Expand Down Expand Up @@ -91,7 +97,20 @@ def read_feed(database: str | Path, *, limit: int = DEFAULT_LIMIT) -> dict[str,
with open_database(database) as connection:
events = load_events(connection, limit=None)
pipeline = load_pipeline_health(connection)
usage = load_inference_usage(connection)
pipeline["database_identity"] = database_identity(database)
# Operator cost tracker: price each window's per-model token sums (see ?ops).
pipeline["cost"] = {
label: {
"usd": round(sum(
inference_cost_usd(
row["model"], row["prompt_tokens"], row["completion_tokens"]
) for row in rows
), 4),
"calls": sum(row["calls"] for row in rows),
}
for label, rows in usage.items()
}
now = datetime.now(UTC)
for event in events:
event["validations"] = _safe_validations(event.get("validations"))
Expand Down Expand Up @@ -169,7 +188,9 @@ def read_feed(database: str | Path, *, limit: int = DEFAULT_LIMIT) -> dict[str,
.theme-toggle{appearance:none;width:34px;height:34px;flex:0 0 auto;border:1px solid var(--line);background:var(--pill);color:var(--text);border-radius:8px;font-size:15px;line-height:1;cursor:pointer;display:flex;align-items:center;justify-content:center;transition:border-color .12s,background-color .12s}.theme-toggle:hover{border-color:var(--amber)}.theme-toggle:focus-visible{outline:2px solid var(--amber);outline-offset:2px}
.strip{display:flex;align-items:stretch;gap:12px;margin-top:14px}
/* Operator-only telemetry: hidden unless ?ops is set (see head script). */
:root:not([data-ops="1"]) .pipeline-health,:root:not([data-ops="1"]) .connection{display:none}
:root:not([data-ops="1"]) .pipeline-health,:root:not([data-ops="1"]) .connection,:root:not([data-ops="1"]) .cost{display:none}
.cost{display:flex;flex-wrap:wrap;align-items:center;gap:6px 16px;margin-top:10px;padding:8px 14px;border:1px solid var(--line);background:var(--pill);color:var(--muted);font-size:9px;letter-spacing:.1em;text-transform:uppercase}
.cost .lead{color:var(--amber)}.cost b{color:var(--text);font-weight:650;font-variant-numeric:tabular-nums}.cost .calls{color:var(--dim);text-transform:none;letter-spacing:0}
.pipeline-health{flex:0 0 auto;display:flex;flex-wrap:wrap;align-items:center;gap:6px 18px;padding:9px 14px;border:1px solid var(--line);background:var(--pill);color:var(--muted);font-size:9px;letter-spacing:.1em;text-transform:uppercase}.pipeline-health b{color:var(--text);font-weight:650}.pipeline-health .good b{color:var(--confirmed)}.pipeline-health .warn b{color:var(--developing)}.pipeline-health .bad b{color:var(--breaking)}
.metrics{flex:1 1 auto;display:flex;flex-wrap:wrap;justify-content:center;gap:10px;margin:0}
.metric{appearance:none;margin:0;padding:8px 14px;background:var(--pill);border:1px solid var(--line);border-radius:999px;font:inherit;font-size:10px;letter-spacing:.09em;text-transform:uppercase;color:var(--muted);cursor:pointer;display:inline-flex;align-items:center;gap:8px;white-space:nowrap;transition:background-color .12s,border-color .12s,color .12s}
Expand Down Expand Up @@ -208,6 +229,7 @@ def read_feed(database: str | Path, *, limit: int = DEFAULT_LIMIT) -> dict[str,
<div class="pipeline-health" id="pipeline-health" aria-label="Pipeline accounting">Awaiting pipeline health…</div>
<div class="tape" aria-label="Latest event ticker"><div class="tape-track" id="tape"></div></div>
</div>
<div class="cost" id="cost" aria-label="Inference cost tracker"></div>
</header>
<main>
<div class="board-head"><div><h2>The Edge Board</h2><p>What the stream caught—and whether the news has caught up.</p></div><span class="updated" id="updated">Awaiting events…</span></div>
Expand All @@ -222,7 +244,8 @@ def read_feed(database: str | Path, *, limit: int = DEFAULT_LIMIT) -> dict[str,
function safeStatus(value){return['breaking','confirmed','contradicted','developing','unconfirmed'].includes(value)?value:'developing'}
function badgeText(item){const n=Array.isArray(item.validations)?item.validations.length:0;switch(safeStatus(item.status)){case'breaking':return'BREAKING — no corroboration yet';case'unconfirmed':return'UNCONFIRMED — wire still silent';case'confirmed':return'CONFIRMED — '+n+' source'+(n===1?'':'s');case'contradicted':return'CONTRADICTED';default:return'DEVELOPING'}}
function lagText(seconds){if(!Number.isFinite(seconds))return'';const abs=Math.abs(seconds),value=abs<3600?Math.round(abs/60)+' min':(abs/3600).toFixed(1)+' hr';return seconds>=0?'news +'+value+' after the stream':'news '+value+' before the stream'}
function renderPipeline(value){const p=value&&typeof value==='object'?value:{},node=$('pipeline-health'),backlog=Number(p.backlog)||0,errors=Number(p.terminal_errors)||0,accounted=Number(p.accounted)||0,posts=Number(p.posts)||0;node.replaceChildren();let state='CAUGHT UP';if(errors>0){state='TERMINAL ERRORS'}else if(backlog>0){state='CATCHING UP'}const coverage=element('span','');coverage.append(document.createTextNode('accounted '),element('b','',accounted+' / '+posts));node.append(coverage);const queue=element('span',backlog?'warn':'good');queue.append(document.createTextNode('backlog '),element('b','',String(backlog)));node.append(queue);if(backlog&&p.oldest_pending_at){const oldest=element('span','');oldest.append(document.createTextNode('oldest '),element('b','',relativeTime(p.oldest_pending_at)));node.append(oldest)}if(errors){const terminal=element('span','bad');terminal.append(document.createTextNode('errors '),element('b','',String(errors)));node.append(terminal)}$('connection').textContent=state;$('pulse').classList.toggle('catchup',backlog>0&&!errors);$('pulse').classList.toggle('error',errors>0)}
function renderPipeline(value){const p=value&&typeof value==='object'?value:{},node=$('pipeline-health'),backlog=Number(p.backlog)||0,errors=Number(p.terminal_errors)||0,accounted=Number(p.accounted)||0,posts=Number(p.posts)||0;node.replaceChildren();let state='CAUGHT UP';if(errors>0){state='TERMINAL ERRORS'}else if(backlog>0){state='CATCHING UP'}const coverage=element('span','');coverage.append(document.createTextNode('accounted '),element('b','',accounted+' / '+posts));node.append(coverage);const queue=element('span',backlog?'warn':'good');queue.append(document.createTextNode('backlog '),element('b','',String(backlog)));node.append(queue);if(backlog&&p.oldest_pending_at){const oldest=element('span','');oldest.append(document.createTextNode('oldest '),element('b','',relativeTime(p.oldest_pending_at)));node.append(oldest)}if(errors){const terminal=element('span','bad');terminal.append(document.createTextNode('errors '),element('b','',String(errors)));node.append(terminal)}$('connection').textContent=state;$('pulse').classList.toggle('catchup',backlog>0&&!errors);$('pulse').classList.toggle('error',errors>0);renderCost(p.cost)}
function renderCost(cost){const node=$('cost');if(!node)return;node.replaceChildren();const c=cost&&typeof cost==='object'?cost:{},labels={hour:'1H',day:'24H',week:'7D',month:'30D'};node.append(element('span','lead','inference cost'));for(const key of ['hour','day','week','month']){const w=c[key]||{},usd=Number(w.usd)||0,calls=Number(w.calls)||0,span=element('span','');span.append(document.createTextNode(labels[key]+' '),element('b','','$'+usd.toFixed(usd<1?4:2)));if(calls)span.append(element('span','calls',' ('+calls+' call'+(calls===1?'':'s')+')'));node.append(span)}}
const FILTERS=[['all','all'],['breaking','breaking'],['unconfirmed','unconfirmed'],['developing','developing'],['confirmed','confirmed'],['contradicted','contradicted']];
const activeFilters=new Set();
function statusCount(items,status){return status==='all'?items.length:items.filter(x=>x.status===status).length}
Expand Down
61 changes: 61 additions & 0 deletions fintick/storage.py
Original file line number Diff line number Diff line change
Expand Up @@ -442,6 +442,16 @@ def set_state(connection: sqlite3.Connection, key: str, value: str) -> None:
""",
"CREATE INDEX IF NOT EXISTS post_aggregation_decisions_state_idx "
"ON post_aggregation_decisions(state, updated_at)",
"""
CREATE TABLE IF NOT EXISTS inference_usage (
id INTEGER PRIMARY KEY AUTOINCREMENT,
at TEXT NOT NULL,
model TEXT NOT NULL,
prompt_tokens INTEGER NOT NULL DEFAULT 0,
completion_tokens INTEGER NOT NULL DEFAULT 0
)
""",
"CREATE INDEX IF NOT EXISTS inference_usage_at_idx ON inference_usage(at)",
)


Expand Down Expand Up @@ -807,6 +817,57 @@ def load_pipeline_health(connection: sqlite3.Connection) -> dict[str, Any]:
}


def record_inference_usage(
connection: sqlite3.Connection,
model: str,
prompt_tokens: int,
completion_tokens: int,
) -> None:
"""Append one inference call's token usage for later cost accounting."""
connection.execute(
"INSERT INTO inference_usage (at, model, prompt_tokens, completion_tokens) "
"VALUES (?, ?, ?, ?)",
(
datetime.now(UTC).isoformat(),
str(model or "unknown"),
int(prompt_tokens or 0),
int(completion_tokens or 0),
),
)


# Rolling windows for the operator cost tracker (label -> lookback).
INFERENCE_COST_WINDOWS = {
"hour": timedelta(hours=1),
"day": timedelta(days=1),
"week": timedelta(weeks=1),
"month": timedelta(days=30),
}


def load_inference_usage(connection: sqlite3.Connection) -> dict[str, list[dict[str, Any]]]:
"""Sum token usage per model over each rolling window (pricing applied upstream)."""
now = datetime.now(UTC)
windows: dict[str, list[dict[str, Any]]] = {}
for label, delta in INFERENCE_COST_WINDOWS.items():
cutoff = (now - delta).isoformat()
rows = connection.execute(
"SELECT model, COALESCE(SUM(prompt_tokens),0), COALESCE(SUM(completion_tokens),0), "
"COUNT(*) FROM inference_usage WHERE at >= ? GROUP BY model",
(cutoff,),
).fetchall()
windows[label] = [
{
"model": row[0],
"prompt_tokens": int(row[1]),
"completion_tokens": int(row[2]),
"calls": int(row[3]),
}
for row in rows
]
return windows


def load_events(
connection: sqlite3.Connection,
*,
Expand Down
54 changes: 54 additions & 0 deletions tests/test_aggregate.py
Original file line number Diff line number Diff line change
Expand Up @@ -82,6 +82,60 @@ def test_empty_content_raises(self, urlopen: mock.Mock) -> None:
with self.assertRaises(RuntimeError):
call_inference("[]")

@mock.patch("fintick.aggregate.urllib.request.urlopen")
def test_usage_sink_receives_token_counts(self, urlopen: mock.Mock) -> None:
urlopen.return_value = BytesIO(json.dumps({
"choices": [{"message": {"content": "{\"events\":[]}"}}],
"usage": {"prompt_tokens": 1735, "completion_tokens": 3497},
}).encode())
captured: list[dict[str, object]] = []
call_inference("[]", model="gpt-5.6-luna", usage_sink=captured.append)
self.assertEqual(captured, [
{"model": "gpt-5.6-luna", "prompt_tokens": 1735, "completion_tokens": 3497},
])

@mock.patch("fintick.aggregate.urllib.request.urlopen")
def test_usage_recorded_even_when_content_is_empty(self, urlopen: mock.Mock) -> None:
# Empty output still burned tokens — the cost must be captured before raising.
urlopen.return_value = BytesIO(json.dumps({
"choices": [{"message": {"content": ""}}],
"usage": {"prompt_tokens": 1700, "completion_tokens": 16384},
}).encode())
captured: list[dict[str, object]] = []
with self.assertRaises(RuntimeError):
call_inference("[]", model="gpt-5.6-luna", usage_sink=captured.append)
self.assertEqual(len(captured), 1)
self.assertEqual(captured[0]["completion_tokens"], 16384)


class InferenceCostTests(unittest.TestCase):
def test_prices_cloud_model_and_frees_local(self) -> None:
from fintick.aggregate import inference_cost_usd
self.assertAlmostEqual(
inference_cost_usd("gpt-5.6-luna", 1_000_000, 1_000_000), 0.20 + 1.20
)
self.assertEqual(inference_cost_usd("qwen3.8:27b", 5000, 4000), 0.0)
self.assertEqual(inference_cost_usd(None, 5000, 4000), 0.0)

def test_usage_windows_sum_per_model(self) -> None:
import tempfile
from fintick.storage import (
open_database, record_inference_usage, load_inference_usage,
)
with tempfile.TemporaryDirectory() as tmp:
db = Path(tmp) / "usage.db"
with open_database(db) as connection:
record_inference_usage(connection, "gpt-5.6-luna", 1735, 3497)
record_inference_usage(connection, "gpt-5.6-luna", 1700, 5000)
with open_database(db) as connection:
windows = load_inference_usage(connection)
self.assertEqual(set(windows), {"hour", "day", "week", "month"})
hour = windows["hour"]
self.assertEqual(len(hour), 1)
self.assertEqual(hour[0]["calls"], 2)
self.assertEqual(hour[0]["prompt_tokens"], 3435)
self.assertEqual(hour[0]["completion_tokens"], 8497)


class AccountedAggregationTests(unittest.TestCase):
def test_short_ids_map_to_uris_and_every_post_gets_a_decision(self) -> None:
Expand Down