diff --git a/fintick/aggregate.py b/fintick/aggregate.py index 573b694..2d62826 100644 --- a/fintick/aggregate.py +++ b/fintick/aggregate.py @@ -20,6 +20,7 @@ POST_AGGREGATION_MAX_ATTEMPTS, V2Event, open_database, + record_inference_usage, set_post_aggregation_decision, upsert_event, ) @@ -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 @@ -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. @@ -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, @@ -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 @@ -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) diff --git a/fintick/dashboard.py b/fintick/dashboard.py index 78985fe..3eff9d1 100644 --- a/fintick/dashboard.py +++ b/fintick/dashboard.py @@ -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 @@ -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")) @@ -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} @@ -208,6 +229,7 @@ def read_feed(database: str | Path, *, limit: int = DEFAULT_LIMIT) -> dict[str,
Awaiting pipeline health…
+

The Edge Board

What the stream caught—and whether the news has caught up.

Awaiting events…
@@ -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} diff --git a/fintick/storage.py b/fintick/storage.py index accef90..5ef6136 100644 --- a/fintick/storage.py +++ b/fintick/storage.py @@ -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)", ) @@ -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, *, diff --git a/tests/test_aggregate.py b/tests/test_aggregate.py index d106da3..74a9f9d 100644 --- a/tests/test_aggregate.py +++ b/tests/test_aggregate.py @@ -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: