diff --git a/public/admin/app.html b/public/admin/app.html new file mode 100644 index 0000000..532c726 --- /dev/null +++ b/public/admin/app.html @@ -0,0 +1,27 @@ + + + + + +Duriin Ops + + + + +
+
+
DURIIN
+
loading console…
+
+
+ + + + + + + + diff --git a/public/admin/assets/css/ops.css b/public/admin/assets/css/ops.css new file mode 100644 index 0000000..50f1331 --- /dev/null +++ b/public/admin/assets/css/ops.css @@ -0,0 +1,197 @@ +/* Ops console. Dark by default because this is a thing you stare at during an + incident, but the tokens are defined so a light theme is a one line change. */ +:root { + --bg: #0b0d10; + --panel: #12161b; + --panel-2: #171c22; + --line: #232a33; + --ink: #e6edf3; + --ink-dim: #9aa7b4; + --ink-faint: #6b7885; + --accent: #4c8dff; + --ok: #3fb950; + --warn: #d29922; + --bad: #f85149; + --radius: 10px; + --mono: ui-monospace, SFMono-Regular, "SF Mono", Menlo, Consolas, monospace; + --sans: -apple-system, BlinkMacSystemFont, "Segoe UI", Inter, Roboto, sans-serif; +} + +* { box-sizing: border-box; } + +body { + margin: 0; + background: var(--bg); + color: var(--ink); + font: 14px/1.5 var(--sans); + -webkit-font-smoothing: antialiased; +} + +.boot { display: grid; place-content: center; gap: 8px; height: 100vh; text-align: center; } +.boot-mark { font: 600 20px var(--mono); letter-spacing: .3em; } +.boot-note { color: var(--ink-faint); font-size: 13px; } + +/* ---------- shell ---------- */ +.shell { display: grid; grid-template-columns: 208px 1fr; min-height: 100vh; } + +.side { + border-right: 1px solid var(--line); + background: var(--panel); + padding: 18px 12px; + position: sticky; top: 0; height: 100vh; + display: flex; flex-direction: column; gap: 4px; +} +.brand { font: 600 14px var(--mono); letter-spacing: .28em; padding: 6px 10px 16px; } +.nav-item { + display: flex; align-items: center; justify-content: space-between; gap: 8px; + padding: 8px 10px; border-radius: 8px; cursor: pointer; + color: var(--ink-dim); text-decoration: none; font-size: 13.5px; + border: 1px solid transparent; +} +.nav-item:hover { background: var(--panel-2); color: var(--ink); } +.nav-item.active { background: var(--panel-2); color: var(--ink); border-color: var(--line); } +.nav-spacer { flex: 1; } + +.main { padding: 20px 24px 60px; min-width: 0; } + +.topbar { + display: flex; align-items: center; gap: 12px; flex-wrap: wrap; + margin-bottom: 18px; +} +.topbar h1 { font-size: 17px; margin: 0; font-weight: 600; } +.topbar .grow { flex: 1; } + +/* ---------- primitives ---------- */ +.grid { display: grid; gap: 14px; } +.cols-2 { grid-template-columns: repeat(auto-fit, minmax(340px, 1fr)); } +.cols-3 { grid-template-columns: repeat(auto-fit, minmax(230px, 1fr)); } + +.card { + background: var(--panel); + border: 1px solid var(--line); + border-radius: var(--radius); + padding: 14px 16px; + min-width: 0; +} +.card h2 { + font-size: 11px; text-transform: uppercase; letter-spacing: .12em; + color: var(--ink-faint); margin: 0 0 12px; font-weight: 600; +} + +.stat { font: 600 24px/1.15 var(--mono); } +.stat-sub { color: var(--ink-faint); font-size: 12px; margin-top: 2px; } + +.dot { width: 8px; height: 8px; border-radius: 50%; display: inline-block; flex: none; } +.dot.ok { background: var(--ok); } +.dot.warn { background: var(--warn); } +.dot.bad { background: var(--bad); } +.dot.idle { background: var(--ink-faint); } +.dot.live { box-shadow: 0 0 0 0 rgba(63,185,80,.6); animation: pulse 2.4s infinite; } +@keyframes pulse { + 70% { box-shadow: 0 0 0 7px rgba(63,185,80,0); } + 100% { box-shadow: 0 0 0 0 rgba(63,185,80,0); } +} + +.row { display: flex; align-items: center; gap: 10px; } +.row + .row { margin-top: 8px; } +.spread { justify-content: space-between; } +.muted { color: var(--ink-faint); } +.mono { font-family: var(--mono); } +.small { font-size: 12px; } +.nowrap { white-space: nowrap; } + +table { width: 100%; border-collapse: collapse; font-size: 13px; } +th { + text-align: left; font-weight: 500; color: var(--ink-faint); + font-size: 11px; text-transform: uppercase; letter-spacing: .08em; + padding: 0 10px 8px 0; border-bottom: 1px solid var(--line); +} +td { padding: 8px 10px 8px 0; border-bottom: 1px solid var(--line); vertical-align: top; } +tr:last-child td { border-bottom: 0; } +.num { text-align: right; font-family: var(--mono); } +.table-scroll { overflow-x: auto; } + +.pill { + display: inline-flex; align-items: center; gap: 6px; + padding: 2px 8px; border-radius: 999px; font-size: 11.5px; + border: 1px solid var(--line); background: var(--panel-2); color: var(--ink-dim); + font-family: var(--mono); +} +.pill.ok { color: var(--ok); border-color: rgba(63,185,80,.35); } +.pill.warn { color: var(--warn); border-color: rgba(210,153,34,.35); } +.pill.bad { color: var(--bad); border-color: rgba(248,81,73,.35); } + +button { + font: inherit; color: var(--ink); background: var(--panel-2); + border: 1px solid var(--line); border-radius: 8px; + padding: 7px 12px; cursor: pointer; +} +button:hover:not(:disabled) { border-color: var(--accent); } +button:disabled { opacity: .45; cursor: default; } +button.primary { background: var(--accent); border-color: var(--accent); color: #06101f; font-weight: 600; } +button.danger { color: var(--bad); border-color: rgba(248,81,73,.4); } +button.danger:hover:not(:disabled) { background: rgba(248,81,73,.12); border-color: var(--bad); } +button.sm { padding: 4px 9px; font-size: 12px; } + +input, select, textarea { + font: inherit; color: var(--ink); background: var(--bg); + border: 1px solid var(--line); border-radius: 8px; padding: 7px 10px; +} +textarea { font-family: var(--mono); font-size: 12.5px; width: 100%; resize: vertical; } +input:focus, select:focus, textarea:focus { outline: none; border-color: var(--accent); } + +pre { + margin: 0; padding: 12px; background: var(--bg); + border: 1px solid var(--line); border-radius: 8px; + overflow: auto; max-height: 460px; + font: 12px/1.5 var(--mono); color: var(--ink-dim); + white-space: pre; +} + +/* ---------- bars ---------- */ +.bar { height: 6px; border-radius: 3px; background: var(--panel-2); overflow: hidden; } +.bar > span { display: block; height: 100%; background: var(--accent); } +.bar.ok > span { background: var(--ok); } +.bar.warn > span { background: var(--warn); } +.bar.bad > span { background: var(--bad); } + +.gate { display: grid; gap: 9px; } +.gate-row { display: grid; grid-template-columns: 1fr auto; gap: 4px 10px; align-items: center; } +.gate-label { font-size: 12.5px; color: var(--ink-dim); } + +.banner { + display: flex; align-items: center; gap: 12px; flex-wrap: wrap; + border: 1px solid rgba(248,81,73,.4); background: rgba(248,81,73,.08); + border-radius: var(--radius); padding: 12px 14px; margin-bottom: 14px; +} +.banner.warn { border-color: rgba(210,153,34,.4); background: rgba(210,153,34,.08); } + +.toast { + position: fixed; right: 18px; bottom: 18px; z-index: 50; + background: var(--panel); border: 1px solid var(--line); + border-radius: 10px; padding: 11px 14px; max-width: 380px; + box-shadow: 0 10px 30px rgba(0,0,0,.45); +} +.toast.bad { border-color: rgba(248,81,73,.5); } +.toast.ok { border-color: rgba(63,185,80,.45); } + +.empty { color: var(--ink-faint); font-size: 13px; padding: 8px 0; } + +.skel { + background: linear-gradient(90deg, var(--panel-2) 25%, #1d232b 37%, var(--panel-2) 63%); + background-size: 400% 100%; + animation: shimmer 1.3s ease infinite; + border-radius: 6px; height: 14px; +} +@keyframes shimmer { 0% { background-position: 100% 0; } 100% { background-position: -100% 0; } } + +@media (max-width: 820px) { + .shell { grid-template-columns: 1fr; } + .side { + position: static; height: auto; flex-direction: row; overflow-x: auto; + border-right: 0; border-bottom: 1px solid var(--line); + } + .brand { padding: 6px 10px; } + .nav-spacer { display: none; } + .main { padding: 16px; } +} diff --git a/public/admin/assets/js/ops/controls.js b/public/admin/assets/js/ops/controls.js new file mode 100644 index 0000000..202f845 --- /dev/null +++ b/public/admin/assets/js/ops/controls.js @@ -0,0 +1,154 @@ +import { html, usePoll, api, useState, num } from './core.js'; +import { Card, Dot, Pill, Confirm, Empty, Skeleton } from './ui.js'; + +export function Controls({ notify }) { + const { data, error, loading, refresh } = usePoll('/admin/api/ops/overview', 8000); + const [busy, setBusy] = useState(null); + const [analysis, setAnalysis] = useState(null); + + const act = async (key, run, describe) => { + setBusy(key); + try { + const result = await run(); + notify({ tone: 'ok', message: describe(result) }); + await refresh(); + return result; + } catch (err) { + console.error('[ops] control failed:', key, err.message); + notify({ tone: 'bad', message: `${key} failed: ${err.message}`, sticky: true }); + return null; + } finally { + setBusy(null); + } + }; + + if (loading && !data) return html`
<${Skeleton} rows=${4} />
`; + + const controls = (data && data.controls) || {}; + const dead = (data && data.deadLetters) || []; + const deadTotal = dead.reduce((sum, row) => sum + Number(row.n || 0), 0); + const killed = Boolean(controls.killSwitch); + + return html` +
+ ${error && html``} + + <${Card} title="Execution" + right=${html`<${Pill} tone=${killed ? 'bad' : (controls.mode === 'paper' ? 'warn' : 'ok')}> + ${killed ? 'halted' : (controls.mode || 'unknown')} + `}> +
+
+
+
+ <${Dot} tone=${killed ? 'bad' : 'ok'} /> + Kill switch ${killed ? 'engaged' : 'released'} +
+
+ While engaged no order intents are created. Broker reconciliation keeps + running so the account stays visible. Takes effect within one poll, + no redeploy. +
+
+ <${Confirm} + label=${killed ? 'release' : 'engage kill switch'} + confirmLabel=${killed ? 'yes, release' : 'yes, halt trading'} + danger=${!killed} + busy=${busy === 'kill'} + onConfirm=${() => act('kill', + () => api('/admin/api/ops/settings', { method: 'POST', body: { killSwitch: !killed } }), + (r) => `Kill switch ${r.controls.killSwitch ? 'engaged' : 'released'}`)} /> +
+ +
+
+ Mode +
+ shadow records intents without touching the broker. + paper submits them to Alpaca paper. + Currently from ${controls.modeSource === 'settings' ? 'this dashboard' : 'the environment'}. +
+
+ <${Confirm} + label=${controls.mode === 'paper' ? 'switch to shadow' : 'switch to paper'} + confirmLabel=${controls.mode === 'paper' ? 'yes, shadow' : 'yes, submit to broker'} + danger=${controls.mode !== 'paper'} + busy=${busy === 'mode'} + onConfirm=${() => act('mode', + () => api('/admin/api/ops/settings', { + method: 'POST', + body: { mode: controls.mode === 'paper' ? 'shadow' : 'paper' }, + }), + (r) => `Execution mode is now ${r.controls.mode}`)} /> +
+
+ + + <${Card} title="Dead letters" + right=${deadTotal > 0 && html`<${Pill} tone="bad">${num(deadTotal)} stuck`}> + ${!deadTotal + ? html`<${Empty}>Nothing dead-lettered. Good.` + : html` +
+
+ Requeue moves rows from dead_letter back to + pending. It never deletes and never edits a payload, + so the worst case is repeated work. +
+
+ + + + ${dead.map((row) => html` + + + + + + + `)} + +
joblanenlast error
${row.job_type}<${Pill}>${row.lane}${num(row.n)}${row.sample_error || '—'} + <${Confirm} label="requeue" + busy=${busy === `dl:${row.job_type}:${row.lane}`} + confirmLabel=${`requeue ${num(row.n)}`} + onConfirm=${() => act(`dl:${row.job_type}:${row.lane}`, + () => api('/admin/api/ops/dead-letters/requeue', { + method: 'POST', body: { jobType: row.job_type, lane: row.lane }, + }), + (r) => `Requeued ${num(r.requeued)} ${row.job_type} jobs`)} /> +
+
+
`} + + + <${Card} title="Analyses"> +
+
+
+ Reaction conditioning +
+ Tests whether the initial market reaction predicts anything, on every + matured outcome. Read only, takes a few minutes, caches price history. +
+
+ +
+ ${analysis && html` +
+ ${analysis.error && html`
${analysis.error}
`} +
${analysis.output || '(no output)'}
+
`} +
+ +
`; +} diff --git a/public/admin/assets/js/ops/core.js b/public/admin/assets/js/ops/core.js new file mode 100644 index 0000000..9ef9321 --- /dev/null +++ b/public/admin/assets/js/ops/core.js @@ -0,0 +1,111 @@ +// Shared plumbing: htm binding, hash router, polling fetch, formatters. +const { createElement, useState, useEffect, useRef, useCallback } = React; +export const html = htm.bind(createElement); +export { useState, useEffect, useRef, useCallback }; + +export async function api(path, options = {}) { + const response = await fetch(path, { + ...options, + headers: { 'Content-Type': 'application/json', ...(options.headers || {}) }, + body: options.body ? JSON.stringify(options.body) : undefined, + }); + const text = await response.text(); + let parsed = null; + try { parsed = text ? JSON.parse(text) : null; } catch (error) { + // A non JSON body from an API that always speaks JSON means something upstream + // failed, so surface the raw text rather than a parse error nobody can act on. + console.error('[ops] non-JSON response from', path, error.message); + throw new Error(text.slice(0, 200) || `HTTP ${response.status}`); + } + if (!response.ok) throw new Error((parsed && parsed.error) || `HTTP ${response.status}`); + return parsed; +} + +// Poll without the spinner flash: keep showing the previous payload while the next +// one is in flight, and never let a slow response overwrite a newer one. +export function usePoll(path, intervalMs = 5000) { + const [state, setState] = useState({ data: null, error: null, loading: true, at: null }); + const seq = useRef(0); + const alive = useRef(true); + + const refresh = useCallback(async () => { + const mine = ++seq.current; + try { + const data = await api(path); + if (!alive.current || mine !== seq.current) return; + setState({ data, error: null, loading: false, at: Date.now() }); + } catch (error) { + if (!alive.current || mine !== seq.current) return; + console.error('[ops] poll failed for', path, error.message); + setState((prev) => ({ ...prev, error: error.message, loading: false })); + } + }, [path]); + + useEffect(() => { + alive.current = true; + refresh(); + const timer = setInterval(refresh, intervalMs); + const onVisible = () => { if (!document.hidden) refresh(); }; + document.addEventListener('visibilitychange', onVisible); + return () => { + alive.current = false; + clearInterval(timer); + document.removeEventListener('visibilitychange', onVisible); + }; + }, [refresh, intervalMs]); + + return { ...state, refresh }; +} + +export function useHashRoute(fallback) { + const read = () => (location.hash || '').replace(/^#\/?/, '') || fallback; + const [route, setRoute] = useState(read); + useEffect(() => { + const onHash = () => setRoute(read()); + addEventListener('hashchange', onHash); + return () => removeEventListener('hashchange', onHash); + }, []); + return route; +} + +/* ---------- formatting ---------- */ +export const num = (value) => (value === null || value === undefined || Number.isNaN(Number(value)) + ? '—' : Number(value).toLocaleString()); + +export function compact(value) { + const n = Number(value); + if (!Number.isFinite(n)) return '—'; + if (Math.abs(n) >= 1e9) return (n / 1e9).toFixed(2) + 'B'; + if (Math.abs(n) >= 1e6) return (n / 1e6).toFixed(2) + 'M'; + if (Math.abs(n) >= 1e3) return (n / 1e3).toFixed(1) + 'k'; + return String(n); +} + +export const pct = (value, digits = 1) => (Number.isFinite(Number(value)) + ? `${(Number(value) * 100).toFixed(digits)}%` : '—'); + +export function ago(minutes) { + if (minutes === null || minutes === undefined || !Number.isFinite(Number(minutes))) return 'unknown'; + const m = Number(minutes); + if (m < 1) return 'just now'; + if (m < 60) return `${Math.round(m)}m ago`; + if (m < 60 * 24) return `${(m / 60).toFixed(1)}h ago`; + return `${(m / 1440).toFixed(1)}d ago`; +} + +export function whenDate(value) { + if (!value) return '—'; + const stamp = String(value).slice(0, 10); + const days = Math.round((Date.parse(`${stamp}T00:00:00Z`) - Date.now()) / 86400000); + if (!Number.isFinite(days)) return stamp; + if (days === 0) return `${stamp} (today)`; + return days > 0 ? `${stamp} (in ${days}d)` : `${stamp} (${-days}d ago)`; +} + +export function health(minutes, threshold) { + if (minutes === null || minutes === undefined) return 'idle'; + if (!Number.isFinite(Number(threshold))) return 'ok'; + const m = Number(minutes); + if (m <= threshold) return 'ok'; + return m <= threshold * 3 ? 'warn' : 'bad'; +} diff --git a/public/admin/assets/js/ops/explore.js b/public/admin/assets/js/ops/explore.js new file mode 100644 index 0000000..12693ed --- /dev/null +++ b/public/admin/assets/js/ops/explore.js @@ -0,0 +1,156 @@ +import { html, usePoll, useState, num, compact, ago } from './core.js'; +import { Card, Stat, Pill, Skeleton, Empty } from './ui.js'; + +const PAGE = 50; + +function shortUrl(url) { + try { return new URL(url).hostname.replace(/^www\./, ''); } + catch (error) { return String(url || '').slice(0, 40); } +} + +function statusTone(status) { + if (status === 'ready') return 'ok'; + if (status === 'failed') return 'bad'; + if (status === 'pending') return 'warn'; + return ''; +} + +// One list component behind both Articles and Events. They differ only in columns +// and endpoint, and having two near-identical files is how they drift apart. +function ListView({ title, path, columns, searchable }) { + const [offset, setOffset] = useState(0); + const [query, setQuery] = useState(''); + const [applied, setApplied] = useState(''); + + const url = `${path}?limit=${PAGE}&offset=${offset}${applied ? `&q=${encodeURIComponent(applied)}` : ''}`; + const { data, error, loading } = usePoll(url, 20000); + + const rows = (data && data.rows) || []; + const total = data && data.total; + + return html` +
+ <${Card} title=${title} + right=${html`${total === undefined ? '' : `${num(total)} total`}`}> + ${searchable && html` +
{ e.preventDefault(); setOffset(0); setApplied(query.trim()); }}> + setQuery(e.target.value)} /> + + ${applied && html``} +
`} + + ${error && html`
${error}
`} + + ${loading && !data + ? html`<${Skeleton} rows=${6} />` + : !rows.length + ? html`<${Empty}>Nothing here.` + : html` +
+ + ${columns.map((c) => html` + `)} + + ${rows.map((row) => html` + + ${columns.map((c) => html` + `)} + `)} + +
${c.label}
${c.render(row)}
+
`} + +
+ + ${rows.length ? `${num(offset + 1)}–${num(offset + rows.length)}` : '—'} + + + + + +
+ +
`; +} + +export const Articles = () => html`<${ListView} + title="Articles" path="/admin/api/articles" searchable + columns=${[ + { key: 'title', label: 'title', render: (r) => html` + ${r.title || '(untitled)'} +
${shortUrl(r.url)}
` }, + { key: 'source', label: 'source', render: (r) => html`${r.source}` }, + { key: 'status', label: 'content', render: (r) => html` + <${Pill} tone=${statusTone(r.content_status)}>${r.content_status || 'unfetched'}` }, + { key: 'pub', label: 'published', render: (r) => html` + ${String(r.pub_date_effective || r.pub_date || '').slice(0, 16)}` }, + ]} />`; + +export const Events = () => html`<${ListView} + title="Events" path="/admin/api/events" + columns=${[ + { key: 'title', label: 'title', render: (r) => r.title || '(untitled)' }, + { key: 'articles', label: 'articles', num: true, render: (r) => num(r.article_count) }, + { key: 'created', label: 'created', render: (r) => html` + ${String(r.created_at || '').slice(0, 16)}` }, + ]} />`; + +export function Intelligence() { + const summary = usePoll('/admin/api/stats/summary', 15000); + const intel = usePoll('/admin/api/intelligence/stats', 15000); + + const s = summary.data || {}; + const i = intel.data || {}; + + return html` +
+
+ <${Stat} label="Articles" value=${compact(s.total)} + sub=${`${compact(s.withContent)} with content · ${compact(s.withEmbedding)} embedded`} /> + <${Stat} label="Events" value=${compact(s.eventCount)} sub="clustered" /> + <${Stat} label="Knowledge rows" value=${compact(i.knowledge)} + sub=${`${compact(i.companies)} companies tracked`} /> +
+ +
+ <${Card} title="Worker rates" right=${html`rows written`}> + ${!(i.workerRates || []).length + ? html`<${Empty}>No worker activity recorded.` + : html` + + + + ${i.workerRates.map((w) => html` + + + + + `)} + +
workerlast 1mlast 5m
${w.worker}${num(w.n1m)}${num(w.n5m)}
`} + + + <${Card} title="Article queue"> + ${!(i.queue || []).length + ? html`<${Empty}>Queue is empty.` + : html` + + + + ${i.queue.map((q) => html` + + + + `)} + +
statusn
<${Pill} tone=${statusTone(q.status)}>${q.status}${num(q.n)}
`} + +
+
`; +} diff --git a/public/admin/assets/js/ops/main.js b/public/admin/assets/js/ops/main.js new file mode 100644 index 0000000..0d24218 --- /dev/null +++ b/public/admin/assets/js/ops/main.js @@ -0,0 +1,87 @@ +import { html, useHashRoute, useState, usePoll, num, ago } from './core.js'; +import { Toast, Dot } from './ui.js'; +import { Overview } from './overview.js'; +import { Controls } from './controls.js'; +import { Articles, Events, Intelligence } from './explore.js'; +import { Sql } from './sql.js'; + +// The d3 force graph is a big specialised visualisation that already works. It is +// mounted in a frame rather than rewritten, because porting it would risk breaking +// something valuable to gain nothing the operator can see. +const Graph = () => html` +
+ +
`; + +const ROUTES = [ + { id: 'overview', label: 'Overview', view: Overview }, + { id: 'controls', label: 'Controls', view: Controls }, + { id: 'articles', label: 'Articles', view: Articles }, + { id: 'events', label: 'Events', view: Events }, + { id: 'intelligence', label: 'Intelligence', view: Intelligence }, + { id: 'graph', label: 'Graph', view: Graph }, + { id: 'sql', label: 'SQL', view: Sql }, +]; + +function Nav({ route, health }) { + return html` + `; +} + +function App() { + const route = useHashRoute('overview'); + const [toast, setToast] = useState(null); + + // A single cheap poll drives the sidebar so every view does not need its own. + const { data } = usePoll('/admin/api/ops/overview', 10000); + const controls = (data && data.controls) || {}; + const deadTotal = ((data && data.deadLetters) || []) + .reduce((sum, row) => sum + Number(row.n || 0), 0); + const ingestMinutes = data && data.freshness ? data.freshness.ingestMinutes : null; + + const halted = Boolean(controls.killSwitch); + const health = { + deadTotal, + halted, + ingestMinutes, + tone: halted ? 'bad' : (deadTotal > 0 ? 'warn' : 'ok'), + label: halted ? 'trading halted' : (controls.mode ? `${controls.mode} mode` : 'connecting…'), + }; + + const active = ROUTES.find((r) => r.id === route) || ROUTES[0]; + const View = active.view; + + return html` +
+ <${Nav} route=${active.id} health=${health} /> +
+
+

${active.label}

+ + ${halted && html`kill switch engaged`} +
+ <${View} notify=${setToast} /> +
+ <${Toast} toast=${toast} /> +
`; +} + +ReactDOM.createRoot(document.getElementById('root')).render(html`<${App} />`); diff --git a/public/admin/assets/js/ops/overview.js b/public/admin/assets/js/ops/overview.js new file mode 100644 index 0000000..b86a710 --- /dev/null +++ b/public/admin/assets/js/ops/overview.js @@ -0,0 +1,216 @@ +import { html, usePoll, num, compact, pct, ago, whenDate, health } from './core.js'; +import { Card, Stat, Dot, Pill, Bar, Skeleton, Empty } from './ui.js'; + +const GATE = { minSample: 30, minInstruments: 5, maxConcentration: 0.5 }; + +const ORIGIN_NOTE = { + live: 'genuine real-time work, the only thing that can authorise an order', + historical: 'coordinator backfill, training evidence only', + replay: 'walk-forward replay, training evidence only', +}; + +function PipelineRow({ label, tone, detail, note }) { + return html` +
+ + <${Dot} tone=${tone} live /> + ${label} + + + ${note && html`${note}`} + ${detail} + +
`; +} + +function LiveEvidence({ maturity, liveOpen, liveResolved }) { + const total = (maturity || []).reduce((sum, row) => sum + Number(row.n || 0), 0); + const next = (maturity || []) + .map((row) => row.first_matures).filter(Boolean).sort()[0]; + + if (!total && !liveResolved) { + return html`<${Empty}>No live predictions yet.`; + } + + return html` +
+
+ ${num(liveResolved)} + matured of ${num(liveOpen + liveResolved)} live +
+ <${Bar} value=${liveResolved} max=${liveOpen + liveResolved} + tone=${liveResolved > 0 ? 'ok' : 'warn'} /> +
+ ${liveResolved > 0 + ? 'Live evidence is accumulating.' + : html`Nothing has matured yet. First outcome ${html`${whenDate(next)}`}.`} +
+
+ + + + ${(maturity || []).map((row) => html` + + + + + `)} + +
horizonopenfirst matures
${row.horizon_days}d${num(row.n)}${whenDate(row.first_matures)}
+
+
`; +} + +function GateCard({ cohorts }) { + const live = (cohorts || []).filter((c) => c.source === 'live'); + const best = live.slice().sort((a, b) => Number(b.sample_size) - Number(a.sample_size))[0]; + const qualifying = live.filter((c) => Number(c.sample_size) >= GATE.minSample + && Number(c.distinct_instruments) >= GATE.minInstruments + && Number(c.top_instrument_share) <= GATE.maxConcentration); + + const rows = [ + { label: 'samples', value: best ? Number(best.sample_size) : 0, need: GATE.minSample, + fmt: (v) => num(v) }, + { label: 'distinct tickers', value: best ? Number(best.distinct_instruments) : 0, need: GATE.minInstruments, + fmt: (v) => num(v) }, + { label: 'top ticker share', value: best ? Number(best.top_instrument_share) : 1, + need: GATE.maxConcentration, invert: true, fmt: (v) => pct(v, 0) }, + ]; + + return html` +
+
+ ${qualifying.length} + live cohorts clear the gate +
+ ${!live.length + ? html`<${Empty}>No live calibration exists yet, so every decision abstains. Offline cohorts cannot authorise orders.` + : html` +
+ ${rows.map((row) => { + const passed = row.invert ? row.value <= row.need : row.value >= row.need; + return html` +
+ ${row.label} + + ${row.fmt(row.value)} ${passed ? '✓' : `/ ${row.fmt(row.need)}`} + +
+ <${Bar} value=${row.invert ? Math.max(0, 1 - row.value) : row.value} + max=${row.invert ? 1 : row.need} tone=${passed ? 'ok' : 'warn'} /> +
+
`; + })} +
+
Best-populated live cohort shown.
`} +
`; +} + +export function Overview({ notify }) { + const { data, error, loading, at } = usePoll('/admin/api/ops/overview', 5000); + + if (loading && !data) { + return html`
+ ${[0, 1, 2].map((i) => html`
<${Skeleton} rows=${2} />
`)} +
`; + } + if (error && !data) return html`<${Card} title="Overview">
${error}
`; + + const p = data.predictions || {}; + const f = data.freshness || {}; + const thresholds = f.thresholds || {}; + const deadTotal = (data.deadLetters || []).reduce((s, r) => s + Number(r.n || 0), 0); + const abstain = (data.decisions || []).find((d) => d.action === 'ABSTAIN'); + const acting = (data.decisions || []).filter((d) => d.action === 'BUY' || d.action === 'SELL') + .reduce((s, d) => s + Number(d.n || 0), 0); + + return html` +
+ ${error && html``} + + ${deadTotal > 0 && html` + `} + +
+ <${Stat} label="Live predictions open" value=${num(p.live_open)} + sub=${`${num(p.live_resolved)} matured · ${num(p.total)} total`} /> + <${Stat} label="Order intents" value=${num(data.orderIntents)} + tone=${Number(data.orderIntents) > 0 ? 'warn' : null} + sub=${acting ? `${num(acting)} actionable decisions` : 'every decision has abstained'} /> + <${Stat} label="Archive" value=${compact(data.archive && data.archive.max_id)} + sub=${`last ingest ${ago(f.ingestMinutes)}`} /> +
+ +
+ <${Card} title="Pipeline"> +
+ <${PipelineRow} label="Ingest" note="articles" + tone=${health(f.ingestMinutes, thresholds.ingest)} + detail=${ago(f.ingestMinutes)} /> + <${PipelineRow} label="Coordinator" note="predictions" + tone=${health(f.predictionMinutes, thresholds.prediction)} + detail=${ago(f.predictionMinutes)} /> + <${PipelineRow} label="Outcomes" note="scored" + tone=${health(f.outcomeMinutes, thresholds.outcome)} + detail=${ago(f.outcomeMinutes)} /> + <${PipelineRow} label="Execution" + note=${data.controls ? (data.controls.killSwitch ? 'kill switch on' : data.controls.mode) : ''} + tone=${data.controls && data.controls.killSwitch ? 'bad' : 'ok'} + detail=${`${num(data.orderIntents)} intents`} /> +
+ + + <${Card} title="Can it trade yet?"> + <${GateCard} cohorts=${data.cohorts} /> + + + <${Card} title="Live evidence"> + <${LiveEvidence} maturity=${data.maturity} + liveOpen=${Number(p.live_open || 0)} liveResolved=${Number(p.live_resolved || 0)} /> + + + <${Card} title="Accuracy by origin" + right=${html`only live counts as a track record`}> + ${!(data.byOrigin || []).length + ? html`<${Empty}>No scored outcomes yet.` + : html` +
+ + + + ${data.byOrigin.map((row) => { + const acc = Number(row.total) ? Number(row.correct) / Number(row.total) : null; + return html` + + + + + + `; + })} + +
originnaccuracymean excess
+
+ <${Pill} tone=${row.origin === 'live' ? 'ok' : ''}>${row.origin} +
+
${ORIGIN_NOTE[row.origin] || ''}
+
${num(row.total)}${pct(acc)}= 0 ? 'ok' : 'bad'})`}> + ${pct(row.mean_excess, 2)} +
+
`} + +
+ +
+ ${abstain ? `${num(abstain.n)} decisions, all ABSTAIN · ` : ''} + updated ${at ? new Date(at).toLocaleTimeString() : '—'} +
+
`; +} diff --git a/public/admin/assets/js/ops/sql.js b/public/admin/assets/js/ops/sql.js new file mode 100644 index 0000000..06ef8f4 --- /dev/null +++ b/public/admin/assets/js/ops/sql.js @@ -0,0 +1,106 @@ +import { html, api, useState, num } from './core.js'; +import { Card, Pill, Empty } from './ui.js'; + +const SAMPLES = [ + { label: 'live prediction maturity', db: 'intelligence', sql: +`SELECT horizon_days, COUNT(*) AS n, + MIN(date(information_cutoff, '+' || horizon_days || ' days')) AS first_matures +FROM autonomy_predictions +WHERE origin='live' AND status='open' +GROUP BY horizon_days ORDER BY horizon_days` }, + { label: 'accuracy by origin', db: 'intelligence', sql: +`SELECT p.origin, COUNT(*) AS n, + ROUND(100.0*AVG(o.direction_correct),1) AS acc_pct, + ROUND(100.0*AVG(o.excess_return),2) AS mean_excess_pct +FROM autonomy_predictions p JOIN autonomy_outcomes o ON o.prediction_id=p.id +GROUP BY p.origin` }, + { label: 'dead letters', db: 'intelligence', sql: +`SELECT job_type, lane, COUNT(*) AS n, substr(MAX(last_error),1,80) AS err +FROM autonomy_jobs WHERE status='dead_letter' GROUP BY job_type, lane` }, + { label: 'content backlog', db: 'archive', sql: +`SELECT COALESCE(content_status,'unfetched') AS status, COUNT(*) AS n +FROM articles GROUP BY status ORDER BY n DESC` }, +]; + +export function Sql({ notify }) { + const [sql, setSql] = useState(SAMPLES[0].sql); + const [database, setDatabase] = useState('intelligence'); + const [result, setResult] = useState(null); + const [busy, setBusy] = useState(false); + + const run = async () => { + setBusy(true); + const started = performance.now(); + try { + const response = await api('/admin/api/sql', { method: 'POST', body: { sql, database } }); + setResult({ ...response, clientMs: Math.round(performance.now() - started) }); + } catch (error) { + console.error('[ops] sql failed:', error.message); + notify({ tone: 'bad', message: error.message, sticky: true }); + setResult(null); + } finally { + setBusy(false); + } + }; + + const first = result && result.results && result.results[0]; + const rows = (first && first.rows) || []; + const columns = rows.length ? Object.keys(rows[0]) : []; + + return html` +
+ <${Card} title="Query" + right=${html` +
+ + +
`}> + +
+ samples: + ${SAMPLES.map((s) => html` + `)} + + ⌘/ctrl + enter +
+ + + ${result && html` + <${Card} title="Result" + right=${html` + <${Pill}>${num(rows.length)} rows + <${Pill}>${num(result.elapsed)}ms server + <${Pill}>${num(result.clientMs)}ms total + `}> + ${first && first.error + ? html`
${first.error}
` + : !rows.length + ? html`<${Empty}>No rows.` + : html` +
+ + ${columns.map((c) => html``)} + + ${rows.slice(0, 500).map((row, i) => html` + + ${columns.map((c) => html` + `)} + `)} + +
${c}
${row[c] === null ? '—' : String(row[c])}
+
+ ${rows.length > 500 && html` +
+ Showing the first 500 of ${num(rows.length)} rows. +
`}`} + `} +
`; +} diff --git a/public/admin/assets/js/ops/ui.js b/public/admin/assets/js/ops/ui.js new file mode 100644 index 0000000..25b4ff4 --- /dev/null +++ b/public/admin/assets/js/ops/ui.js @@ -0,0 +1,71 @@ +import { html, useState, useEffect } from './core.js'; + +export const Card = ({ title, right, children }) => html` +
+ ${title && html` +
+

${title}

+ ${right} +
`} + ${children} +
`; + +export const Stat = ({ label, value, sub, tone }) => html` +
+

${label}

+
${value}
+ ${sub && html`
${sub}
`} +
`; + +export const Dot = ({ tone = 'idle', live }) => html` + `; + +export const Pill = ({ tone, children }) => html`${children}`; + +export const Bar = ({ value, max, tone }) => { + const width = max > 0 ? Math.min(100, (Number(value) / Number(max)) * 100) : 0; + return html`
`; +}; + +export const Skeleton = ({ rows = 3 }) => html` +
+ ${Array.from({ length: rows }, (_, i) => html` +
`)} +
`; + +export const Empty = ({ children }) => html`
${children}
`; + +export function Toast({ toast }) { + const [shown, setShown] = useState(toast); + useEffect(() => { + setShown(toast); + if (!toast) return undefined; + const timer = setTimeout(() => setShown(null), toast.sticky ? 12000 : 4500); + return () => clearTimeout(timer); + }, [toast]); + if (!shown) return null; + return html`
${shown.message}
`; +} + +// Anything that changes production asks twice. The second click is a different +// button label so muscle memory cannot carry you through both. +export function Confirm({ label, confirmLabel, onConfirm, danger, disabled, busy }) { + const [armed, setArmed] = useState(false); + useEffect(() => { + if (!armed) return undefined; + const timer = setTimeout(() => setArmed(false), 5000); + return () => clearTimeout(timer); + }, [armed]); + + if (busy) return html``; + if (!armed) { + return html``; + } + return html` + + + + `; +} diff --git a/src/autonomy/schema.js b/src/autonomy/schema.js index 69704c1..cf912c5 100644 --- a/src/autonomy/schema.js +++ b/src/autonomy/schema.js @@ -221,6 +221,16 @@ function initAutonomySchema(db) { created_at TEXT NOT NULL DEFAULT (datetime('now')) ); + -- Runtime knobs the operator can change without a redeploy. Execution mode used + -- to live only in AUTONOMY_EXECUTION_MODE, which means flipping it needed a + -- container restart, which is the last thing you want during an incident. + CREATE TABLE IF NOT EXISTS autonomy_settings ( + key TEXT PRIMARY KEY, + value TEXT NOT NULL, + updated_at TEXT NOT NULL DEFAULT (datetime('now')), + updated_by TEXT + ); + INSERT OR IGNORE INTO autonomy_schema(version) VALUES (${AUTONOMY_SCHEMA_VERSION}); `); for (const statement of [ diff --git a/src/autonomy/settings.js b/src/autonomy/settings.js new file mode 100644 index 0000000..2ca37be --- /dev/null +++ b/src/autonomy/settings.js @@ -0,0 +1,67 @@ +// Runtime settings for the autonomy stack. The env vars stay the default, the +// table is the override, so nothing changes behaviour until somebody deliberately +// writes a row. Kept deliberately tiny -- this is a control plane, not a config +// system, and every key here can move real money or stop the pipeline. + +const EXECUTION_MODES = ['shadow', 'paper']; + +const KEYS = { + executionMode: 'execution_mode', + killSwitch: 'execution_kill_switch', +}; + +function readSetting(db, key) { + try { + const row = db.prepare('SELECT value FROM autonomy_settings WHERE key = ?').get(key); + return row ? row.value : null; + } catch (error) { + // A missing table means an older schema, which should behave like "no override" + // rather than taking the worker down with it. + console.error(`[settings] could not read ${key}:`, error.message); + return null; + } +} + +function writeSetting(db, key, value, updatedBy = 'admin') { + db.prepare(` + INSERT INTO autonomy_settings(key, value, updated_by) VALUES (?, ?, ?) + ON CONFLICT(key) DO UPDATE SET value = excluded.value, updated_by = excluded.updated_by, + updated_at = datetime('now') + `).run(key, String(value), updatedBy); +} + +// Truthy strings people actually type, rather than only accepting 'true' +function isOn(value) { + return ['1', 'true', 'on', 'yes', 'engaged'].includes(String(value || '').trim().toLowerCase()); +} + +function getExecutionControls(db, env = process.env) { + const stored = readSetting(db, KEYS.executionMode); + const fallback = env.AUTONOMY_EXECUTION_MODE || 'shadow'; + const mode = EXECUTION_MODES.includes(String(stored)) ? String(stored) : fallback; + + return { + mode: EXECUTION_MODES.includes(mode) ? mode : 'shadow', + modeSource: EXECUTION_MODES.includes(String(stored)) ? 'settings' : 'env', + killSwitch: isOn(readSetting(db, KEYS.killSwitch)), + }; +} + +function setExecutionMode(db, mode, updatedBy) { + if (!EXECUTION_MODES.includes(mode)) { + throw new Error(`unsupported execution mode: ${mode} (expected ${EXECUTION_MODES.join(' or ')})`); + } + writeSetting(db, KEYS.executionMode, mode, updatedBy); + return getExecutionControls(db); +} + +function setKillSwitch(db, engaged, updatedBy) { + writeSetting(db, KEYS.killSwitch, engaged ? 'true' : 'false', updatedBy); + return getExecutionControls(db); +} + +module.exports = { + EXECUTION_MODES, KEYS, isOn, + readSetting, writeSetting, + getExecutionControls, setExecutionMode, setKillSwitch, +}; diff --git a/src/routes/admin.js b/src/routes/admin.js index db385bd..e0c253b 100644 --- a/src/routes/admin.js +++ b/src/routes/admin.js @@ -7,6 +7,7 @@ const config = require('../config'); const Database = require('better-sqlite3'); const { openRuntimeDb, isPostgresEnabled } = require('../db/runtime'); const pg = require('../db/pgAsync'); +const opsRoutes = require('./ops'); let idb = null; let adb = null; @@ -265,6 +266,7 @@ const pagesDir = path.join(publicDir, 'pages'); // map pretty url → page html file. keep these close to the routes so its // obvious when a page gets added or renamed. const pageMap = { + '/admin/console': path.join(publicDir, 'app.html'), '/admin/autonomy': path.join(pagesDir, 'autonomy.html'), '/admin/ingest/articles': path.join(pagesDir, 'ingest', 'articles.html'), '/admin/ingest/events': path.join(pagesDir, 'ingest', 'events.html'), @@ -283,6 +285,10 @@ function sendPage(reply, filePath) { } async function adminRoutes(fastify) { + // Control plane for the ops dashboard. Lives in its own file, shares this one's + // auth and db handles so there is exactly one of each. + fastify.register(opsRoutes, { checkAuth, getIntelligenceDb, getArchiveDb }); + // gate every request under /admin/* behind basic auth (covers pages, api, and assets) fastify.addHook('onRequest', async (request, reply) => { @@ -303,9 +309,11 @@ async function adminRoutes(fastify) { }, }); - // top-level entry — the autonomy control room is the primary product surface + // top-level entry — the ops console is the primary surface now. The older + // per-page admin is still served underneath, the console frames the d3 graph + // from it rather than duplicating that visualisation. fastify.get('/admin', async (request, reply) => { - reply.redirect('/admin/autonomy'); + reply.redirect('/admin/console'); }); // ingest root — redirect to the articles subsection diff --git a/src/routes/ops.js b/src/routes/ops.js new file mode 100644 index 0000000..8a60d39 --- /dev/null +++ b/src/routes/ops.js @@ -0,0 +1,254 @@ +// Control plane for the operations dashboard. +// +// Kept apart from admin.js on purpose: everything in here either changes what the +// autonomy stack does or is read by an operator while something is on fire, so it +// wants to stay small enough to audit in one sitting. +const { execFile } = require('child_process'); +const path = require('path'); +const { getExecutionControls, setExecutionMode, setKillSwitch, EXECUTION_MODES } = require('../autonomy/settings'); + +// Which of the pipeline stages we consider "recent enough to be alive". These are +// generous, they exist to catch a stall not to police a few seconds of jitter. +const STALE_AFTER_MINUTES = { ingest: 90, prediction: 180, outcome: 24 * 60 }; + +const ANALYSES = { + reaction: { + label: 'Reaction conditioning', + script: 'scripts/analyze-reaction-conditioning.js', + detail: 'Does the initial market reaction predict anything. Read only, a few minutes.', + }, +}; + +function minutesSince(value) { + if (!value) return null; + const stamp = String(value).includes('T') ? String(value) : `${String(value).replace(' ', 'T')}Z`; + const then = Date.parse(stamp); + if (!Number.isFinite(then)) return null; + return Math.max(0, (Date.now() - then) / 60000); +} + +function one(db, sql, params = []) { + try { + return db.prepare(sql).get(...params) || {}; + } catch (error) { + console.error('[ops] query failed:', sql.trim().slice(0, 80), error.message); + return {}; + } +} + +function many(db, sql, params = []) { + try { + return db.prepare(sql).all(...params) || []; + } catch (error) { + console.error('[ops] query failed:', sql.trim().slice(0, 80), error.message); + return []; + } +} + +// One request for the whole dashboard. The old admin made the browser fire a +// handful of sequential calls and stitch them together, which is most of why it +// felt sluggish even though every individual endpoint was fast. +function buildOverview(intel) { + const predictions = one(intel, ` + SELECT COUNT(*) AS total, + SUM(CASE WHEN origin='live' AND status='open' THEN 1 ELSE 0 END) AS live_open, + SUM(CASE WHEN origin='live' AND status='resolved' THEN 1 ELSE 0 END) AS live_resolved, + SUM(CASE WHEN origin IN ('historical','replay') AND status='resolved' THEN 1 ELSE 0 END) AS offline_resolved, + SUM(CASE WHEN status='unresolvable' THEN 1 ELSE 0 END) AS unresolvable, + MAX(created_at) AS latest + FROM autonomy_predictions + `); + + const byOrigin = many(intel, ` + SELECT p.origin, COUNT(*) AS total, + SUM(o.direction_correct) AS correct, + AVG(o.excess_return) AS mean_excess + FROM autonomy_predictions p JOIN autonomy_outcomes o ON o.prediction_id = p.id + GROUP BY p.origin + `); + + // When live evidence actually arrives. This is the number that decides whether + // anything can ever qualify, and nothing in the old UI showed it. + const maturity = many(intel, ` + SELECT horizon_days, COUNT(*) AS n, + MIN(date(information_cutoff, '+' || horizon_days || ' days')) AS first_matures + FROM autonomy_predictions + WHERE origin='live' AND status='open' + GROUP BY horizon_days ORDER BY horizon_days + `); + + const decisions = many(intel, "SELECT action, COUNT(*) AS n FROM autonomy_decisions GROUP BY action"); + const intents = one(intel, 'SELECT COUNT(*) AS n FROM autonomy_order_intents'); + + const jobs = many(intel, ` + SELECT job_type, lane, status, COUNT(*) AS n, MAX(created_at) AS newest + FROM autonomy_jobs GROUP BY job_type, lane, status + `); + const deadLetters = many(intel, ` + SELECT job_type, lane, COUNT(*) AS n, substr(MAX(last_error), 1, 160) AS sample_error + FROM autonomy_jobs WHERE status='dead_letter' GROUP BY job_type, lane + `); + + let cohorts = []; + try { + cohorts = many(intel, ` + SELECT cohort_key, source, sample_size, distinct_instruments, top_instrument_share, + directional_probability, expected_excess_return + FROM autonomy_calibration_snapshots + WHERE cohort_key LIKE 'v2|%' + GROUP BY cohort_key, source + HAVING MAX(created_at) = created_at + ORDER BY sample_size DESC LIMIT 40 + `); + } catch (error) { + console.error('[ops] cohort snapshot read failed:', error.message); + } + + const outcomes = one(intel, 'SELECT COUNT(*) AS n, MAX(evaluated_at) AS latest FROM autonomy_outcomes'); + + return { + generatedAt: new Date().toISOString(), + predictions, + byOrigin, + maturity, + decisions, + orderIntents: intents.n || 0, + outcomes, + jobs, + deadLetters, + cohorts, + freshness: { + predictionMinutes: minutesSince(predictions.latest), + outcomeMinutes: minutesSince(outcomes.latest), + thresholds: STALE_AFTER_MINUTES, + }, + }; +} + +function opsRoutes(fastify, options) { + const { checkAuth, getIntelligenceDb, getArchiveDb } = options; + + const withIntel = (reply) => { + const intel = getIntelligenceDb(); + if (!intel) { + reply.code(503).send({ error: 'intelligence database unavailable' }); + return null; + } + return intel; + }; + + fastify.get('/admin/api/ops/overview', async (request, reply) => { + if (!checkAuth(request, reply)) return; + const intel = withIntel(reply); + if (!intel) return; + + const payload = buildOverview(intel); + try { + payload.controls = getExecutionControls(intel); + } catch (error) { + console.error('[ops] could not read execution controls:', error.message); + payload.controls = null; + } + + // Archive counts are the one genuinely expensive thing here, so they are + // cheap approximations rather than COUNT(*) over 2.2M rows on every poll. + try { + const archive = getArchiveDb ? getArchiveDb() : null; + if (archive) { + payload.archive = one(archive, ` + SELECT MAX(id) AS max_id, MAX(ingested_at) AS latest_ingest FROM articles + `); + payload.freshness.ingestMinutes = minutesSince(payload.archive.latest_ingest); + } + } catch (error) { + console.error('[ops] archive probe failed:', error.message); + payload.archive = null; + } + return payload; + }); + + fastify.get('/admin/api/ops/settings', async (request, reply) => { + if (!checkAuth(request, reply)) return; + const intel = withIntel(reply); + if (!intel) return; + return { controls: getExecutionControls(intel), modes: EXECUTION_MODES, analyses: ANALYSES }; + }); + + fastify.post('/admin/api/ops/settings', async (request, reply) => { + if (!checkAuth(request, reply)) return; + const intel = withIntel(reply); + if (!intel) return; + const { mode, killSwitch } = request.body || {}; + try { + if (mode !== undefined) setExecutionMode(intel, String(mode), 'admin-ui'); + if (killSwitch !== undefined) setKillSwitch(intel, Boolean(killSwitch), 'admin-ui'); + } catch (error) { + console.error('[ops] settings write rejected:', error.message); + reply.code(400).send({ error: error.message }); + return; + } + const controls = getExecutionControls(intel); + console.log(`[ops] execution controls now mode=${controls.mode} kill=${controls.killSwitch}`); + return { controls }; + }); + + // Requeue is deliberately narrow: it only ever moves dead_letter back to pending, + // it never deletes and never edits payloads, so the worst case is repeated work. + fastify.post('/admin/api/ops/dead-letters/requeue', async (request, reply) => { + if (!checkAuth(request, reply)) return; + const intel = withIntel(reply); + if (!intel) return; + const { jobType, lane } = request.body || {}; + + const filters = ["status = 'dead_letter'"]; + const params = []; + if (jobType) { filters.push('job_type = ?'); params.push(String(jobType)); } + if (lane) { filters.push('lane = ?'); params.push(String(lane)); } + const where = filters.join(' AND '); + + try { + const before = one(intel, `SELECT COUNT(*) AS n FROM autonomy_jobs WHERE ${where}`, params).n || 0; + if (!before) return { requeued: 0, remaining: 0 }; + const result = intel.prepare(` + UPDATE autonomy_jobs + SET status='pending', attempts=0, available_at=datetime('now'), + leased_by=NULL, lease_expires_at=NULL, + last_error='requeued from the ops dashboard' + WHERE ${where} + `).run(...params); + const remaining = one(intel, `SELECT COUNT(*) AS n FROM autonomy_jobs WHERE ${where}`, params).n || 0; + console.log(`[ops] requeued ${result.changes} dead letters (jobType=${jobType || 'any'} lane=${lane || 'any'})`); + return { requeued: result.changes, before, remaining }; + } catch (error) { + console.error('[ops] requeue failed:', error.message, error.stack); + reply.code(500).send({ error: error.message }); + } + }); + + fastify.post('/admin/api/ops/analysis/:name', async (request, reply) => { + if (!checkAuth(request, reply)) return; + const spec = ANALYSES[request.params.name]; + if (!spec) { + reply.code(404).send({ error: `unknown analysis: ${request.params.name}` }); + return; + } + const script = path.resolve(__dirname, '..', '..', spec.script); + return new Promise((resolve) => { + execFile('node', [script], { timeout: 15 * 60 * 1000, maxBuffer: 8 * 1024 * 1024 }, + (error, stdout, stderr) => { + if (error) console.error(`[ops] analysis ${request.params.name} failed:`, error.message); + resolve({ + analysis: request.params.name, + label: spec.label, + ok: !error, + output: String(stdout || '').slice(-20000), + error: error ? String(stderr || error.message).slice(-4000) : null, + }); + }); + }); + }); +} + +module.exports = opsRoutes; +module.exports.buildOverview = buildOverview; +module.exports.ANALYSES = ANALYSES; diff --git a/workers/executionWorker.js b/workers/executionWorker.js index da3740b..8a13e17 100644 --- a/workers/executionWorker.js +++ b/workers/executionWorker.js @@ -3,6 +3,7 @@ const { openRuntimeDb } = require('../src/db/runtime'); const { initAutonomySchema } = require('../src/autonomy/schema'); const { createOrderIntent } = require('../src/autonomy/orderIntents'); const { createAlpacaPaperClient } = require('../src/brokers/alpacaPaper'); +const { getExecutionControls } = require('../src/autonomy/settings'); function sleep(ms) { return new Promise((resolve) => setTimeout(resolve, ms)); } @@ -12,10 +13,27 @@ async function runExecutionWorker({ intelligencePath, pollMs = 10000, mode = 'sh db.pragma('journal_mode = WAL'); db.pragma('busy_timeout = 5000'); initAutonomySchema(db); - const paperClient = mode === 'paper' - ? createAlpacaPaperClient({ keyId: process.env.ALPACA_PAPER_KEY_ID, secretKey: process.env.ALPACA_PAPER_SECRET_KEY }) - : null; + + // mode is re-read every poll now. the env var is only the default, so the kill + // switch actually works during an incident instead of needing a redeploy first. + let paperClient = null; + let lastMode = null; + let lastKill = null; + while (true) { + const controls = getExecutionControls(db, { AUTONOMY_EXECUTION_MODE: mode }); + if (controls.mode !== lastMode) { + console.log(`[${workerId}] execution mode = ${controls.mode} (from ${controls.modeSource})`); + lastMode = controls.mode; + } + if (controls.killSwitch !== lastKill) { + console.log(`[${workerId}] kill switch ${controls.killSwitch ? 'ENGAGED, no orders will be placed' : 'released'}`); + lastKill = controls.killSwitch; + } + if (controls.mode === 'paper' && !paperClient) { + paperClient = createAlpacaPaperClient({ keyId: process.env.ALPACA_PAPER_KEY_ID, secretKey: process.env.ALPACA_PAPER_SECRET_KEY }); + } + if (controls.mode !== 'paper') paperClient = null; if (paperClient) { try { const [account, positions, orders] = await Promise.all([ @@ -49,7 +67,9 @@ async function runExecutionWorker({ intelligencePath, pollMs = 10000, mode = 'sh console.error(`[${workerId}] broker reconciliation:`, error.message); } } - const decisions = db.prepare(` + // Broker reconciliation above still runs while the switch is engaged, we want + // to keep seeing the account. It is order creation specifically that stops. + const decisions = controls.killSwitch ? [] : db.prepare(` SELECT d.id FROM autonomy_decisions d JOIN autonomy_predictions p ON p.id = d.prediction_id LEFT JOIN autonomy_order_intents oi ON oi.decision_id = d.id @@ -59,8 +79,8 @@ async function runExecutionWorker({ intelligencePath, pollMs = 10000, mode = 'sh for (const decision of decisions) { try { const intent = createOrderIntent(db, decision.id, notional, { tradable: true, maxNotional: notional }); - if (mode === 'paper') db.prepare("UPDATE autonomy_order_intents SET status='pending', updated_at=datetime('now') WHERE client_order_id=?").run(intent.clientOrderId); - console.log(`[${workerId}] ${mode} intent ${intent.clientOrderId}`); + if (controls.mode === 'paper') db.prepare("UPDATE autonomy_order_intents SET status='pending', updated_at=datetime('now') WHERE client_order_id=?").run(intent.clientOrderId); + console.log(`[${workerId}] ${controls.mode} intent ${intent.clientOrderId}`); } catch (error) { console.error(`[${workerId}] decision ${decision.id}:`, error.message); }