AutonomousAgent / page_files /Run_Control.py
Mathias Heider
Claude Opus 5.5
Cost "f": price the agent at the model it runs; thinking/vision token breakdown; model switch; near-duplicate gate
2cf2ccd
Raw History Blame Contribute Delete
28.7 kB
"""Run Control — autonomy switch, caps, query frontier. Admin-gated."""
import pandas as pd
import streamlit as st
from agent import agdb, scheduler
from agent import config as C
from agent import ui_theme as T
from agent.ui_common import admin_gate, status_banner, df_or_empty
_TITLE = "Run Control"
_SUB = "Autonomy switch, per-cycle budget, discovery sources and the query frontier."
# Page-local CSS: card_title() applies uppercase/letter-spacing to the whole
# title row; pills/chips passed as right_html must keep their own casing.
# (candidate to hoist into ui_theme.card_title)
st.html("""<style>
.aim-card-title .aim-rc-right { text-transform: none; letter-spacing: 0; display: inline-flex;
gap: 6px; align-items: center; flex-wrap: wrap; justify-content: flex-end; }
.aim-card-title .aim-rc-right .aim-chips { margin: 0; }
.aim-rc-note { font-size: 12.5px; color: var(--aim-ink2); background: var(--aim-tint);
border: 1px solid var(--aim-hairline); border-radius: 8px; padding: 6px 10px; margin: -.3rem 0 .6rem 0;
display: inline-flex; gap: 6px; align-items: center; }
</style>""")
def _right(html: str) -> str:
"""Wrap pill/chip HTML for card_title(right_html=...) so it keeps its casing."""
return f'<span class="aim-rc-right">{html}</span>' if html else ""
def _num_sum(df: pd.DataFrame, col: str) -> int:
"""NaN/NULL-proof integer column sum (0 when the column is missing)."""
if col not in df:
return 0
try:
return int(pd.to_numeric(df[col], errors="coerce").fillna(0).sum())
except Exception:
return 0
# Provisional header (no pills yet) so a DB-unreachable error still shows the
# page title; it is replaced with the pill version once the state is known.
_hdr = st.empty()
with _hdr:
T.page_header(_TITLE, _SUB)
status_banner()
admin = admin_gate()
conn = agdb.connect()
try:
cfg = agdb.get_config(conn)
lock = agdb.lock_is_active(conn) # fresh lock only — stale ≠ running
dlock = agdb.lock_is_active(conn, key="discovery_lock")
try:
qstats = agdb.queue_stats(conn)
except Exception:
qstats = {"by_status": {}, "queued_by_source": {}, "last_discovery": None, "oldest_queued": None}
finally:
conn.close()
_autonomy_on = bool(cfg.get("autonomy"))
_interval_h = float(cfg.get("interval_hours", 6))
_nxt = scheduler.next_run_time()
_nxt_txt = _nxt.strftime("%H:%M UTC") if _nxt else "—"
_autonomy_pill = (T.pill_html("ok", f"autonomy ON · every {_interval_h:g} h")
if _autonomy_on else T.pill_html("off", "autonomy OFF"))
_state_pill = (T.pill_html("running", "cycle running") if lock
else T.pill_html("idle", "idle"))
_pills = _autonomy_pill + _state_pill
if not admin:
_pills += T.pill_html("info", "read-only")
with _hdr:
T.page_header(_TITLE, _SUB, right=_pills)
if not admin:
st.html('<div class="aim-rc-note"><span aria-hidden="true">🔒</span>'
'Unlock admin in the sidebar to change settings — everything below is read-only.</div>')
# ---------------------------------------------------------------------------
# Card grid — Row 1: Autonomy | Sources & behavior; Row 2: Per-cycle budget |
# Manual run; then the Save row and the Query frontier. The column containers
# are created up front and filled later so the widget/button order in code
# stays exactly as before: autonomy → budget → sources → save → manual run →
# query frontier.
# ---------------------------------------------------------------------------
row1_left, row1_right = st.columns([5, 3], gap="medium")
row2_left, row2_right = st.columns([5, 3], gap="medium")
# --- autonomy ---------------------------------------------------------------
with row1_left:
with st.container(border=True, height="stretch"):
T.card_title("Autonomy", right_html=_right(_autonomy_pill))
col1, col2 = st.columns([1, 1], vertical_alignment="bottom")
with col1:
autonomy = st.toggle("Autonomous crawling", value=bool(cfg.get("autonomy")),
disabled=not admin, key="rc_autonomy")
with col2:
interval = st.number_input("Interval (hours)", min_value=0.5, max_value=168.0,
value=float(cfg.get("interval_hours", 6)),
step=0.5, disabled=not admin, key="rc_interval")
st.caption("While ON, the agent runs a full discover→dedupe(DOI)→download→"
"extract→publish cycle at this interval, viewers or not. "
"Restarts resume automatically.")
st.html(T.chips_html([
("next scheduled cycle", _nxt_txt, "muted"),
("interval", f"{_interval_h:g} h", "muted"),
]))
# --- per-cycle budget -------------------------------------------------------
with row2_left:
with st.container(border=True, height="stretch"):
T.card_title("Per-cycle budget", sub=f"hard cap {C.HARD_MAX_PDFS_PER_CYCLE} PDFs / cycle")
b1, b2, b3 = st.columns(3, vertical_alignment="bottom")
with b1:
max_pdfs = st.number_input("Max new PDFs per cycle", 1, 25,
int(cfg.get("max_new_pdfs_per_cycle", 5)),
disabled=not admin, key="rc_max_pdfs")
with b2:
per_query = st.number_input("Max results per query/source", 1, 50,
int(cfg.get("max_per_query", 12)),
disabled=not admin, key="rc_per_query",
help="Relevant candidates taken per query from each "
"source. OpenAlex/S2 calls cost the same at 200 "
"rows as at 20, so this is a download-effort knob, "
"not an API-cost knob.")
with b3:
n_queries = st.number_input("Queries per cycle", 1, 10,
int(cfg.get("queries_per_cycle", 3)),
disabled=not admin, key="rc_n_queries")
b4, b5, b6 = st.columns(3, vertical_alignment="bottom")
with b4:
target_pdfs = st.number_input("Target new PDFs per cycle", 1, 25,
int(cfg.get("min_new_pdfs_per_cycle", 10)),
disabled=not admin, key="rc_target_pdfs",
help="Below this the cycle keeps taking the next "
"frontier slice (up to Max rounds); at the cap "
"above it stops regardless. Each extra round is "
"one more search per query per lane.")
with b5:
max_rounds = st.number_input("Max rounds per cycle", 1, 10,
int(cfg.get("max_rounds_per_cycle", 5)),
disabled=not admin, key="rc_max_rounds")
with b6:
workers = st.number_input("Ingest workers", 1, C.HARD_MAX_INGEST_WORKERS,
int(cfg.get("ingest_workers", 3)),
disabled=not admin, key="rc_workers",
help="PDFs extracted concurrently (Gemini calls are "
"latency-bound). Each worker holds one Postgres "
"connection; 1 = one PDF at a time.")
st.caption("Each cycle picks this many queries, asks every enabled source for at most "
"this many results each, downloads at most this many new PDFs, and "
"extracts them with this many parallel workers.")
# --- sources & behavior -----------------------------------------------------
with row1_right:
with st.container(border=True, height="stretch"):
T.card_title("Sources & behavior")
use_oa = st.checkbox("OpenAlex + Unpaywall", value=bool(cfg.get("use_openalex", True)),
disabled=not admin, key="rc_use_openalex")
use_ax = st.checkbox("arXiv", value=bool(cfg.get("use_arxiv", True)),
disabled=not admin, key="rc_use_arxiv")
use_s2 = st.checkbox("Semantic Scholar (bulk, OA-PDF only)",
value=bool(cfg.get("use_semantic_scholar", False)),
disabled=not admin,
help="One bulk call per query returns up to 1,000 papers that "
"have an open-access PDF. Works without a key on a shared "
"rate bucket; an S2_API_KEY secret gives a dedicated quota.",
key="rc_use_s2")
use_ntrs = st.checkbox("NASA NTRS", value=bool(cfg.get("use_ntrs", True)),
disabled=not admin,
help="NASA Technical Reports Server: free JSON API, no key, "
"no rate limit, direct PDF links (Langley thermoplastic "
"composite reports). Slide decks are skipped.",
key="rc_use_ntrs")
use_topics = st.checkbox("OpenAlex topic walk",
value=bool(cfg.get("use_openalex_topics", True)),
disabled=not admin,
help="Each cycle also takes the next page of open-access works "
"in the configured OpenAlex topics (a 1-credit filter call "
"per topic vs 10 credits per full-text search).",
key="rc_use_topics")
use_ds = st.checkbox("Manufacturer datasheets",
value=bool(cfg.get("use_datasheets", True)),
disabled=not admin,
help="Crawl the datasheet seed pages (Victrex, Hexcel, Toray, RTP, …) "
"for new PDFs; each seed is revisited at most once a day.",
key="rc_use_datasheets")
expand = st.checkbox("Gemini query expansion", value=bool(cfg.get("expand_queries", False)),
disabled=not admin,
help="When a cycle finds nothing new, ask Gemini for fresh search intents.",
key="rc_expand")
st.caption("Discovery sources for each cycle; query expansion only kicks in "
"when a cycle yields nothing new. A source that answers 'budget "
"exhausted' is skipped until its reset time (see Runs & Logs).")
# --- candidate queue ---------------------------------------------------------
_qs = qstats.get("by_status", {})
_q_depth = int(_qs.get("queued", 0))
_q_pill = (T.pill_html("running", "discovery running") if dlock
else T.pill_html("ok", f"{_q_depth:,} queued") if _q_depth
else T.pill_html("warn", "queue empty"))
with st.container(border=True):
T.card_title("Candidate queue",
sub="a discovery pass fills it from every lane; cycles drain it",
right_html=_right(_q_pill))
qc1, qc2, qc3, qc4 = st.columns([1.2, 1, 1, 1], vertical_alignment="bottom")
with qc1:
use_queue = st.toggle("Use the candidate queue",
value=bool(cfg.get("use_candidate_queue", True)),
disabled=not admin, key="rc_use_queue",
help="ON: cycles pop queued candidates and only search inline "
"when the queue is empty; a discovery pass runs on the "
"interval below and whenever the queue falls under the "
"low-water mark. OFF: every cycle searches inline (pre-28 Sep).")
with qc2:
disc_hours = st.number_input("Discovery every (hours)", min_value=1.0, max_value=168.0,
value=float(cfg.get("discovery_interval_hours", 24)),
step=1.0, disabled=not admin, key="rc_disc_hours")
with qc3:
low_water = st.number_input("Low-water mark", 0, 5000,
int(cfg.get("queue_low_water", 150)),
disabled=not admin, key="rc_low_water",
help="A cycle that finds fewer queued candidates than this "
"asks for a discovery pass right after it.")
with qc4:
disc_queries = st.number_input("Queries per pass", 1, 500,
int(cfg.get("discovery_queries", 40)),
disabled=not admin, key="rc_disc_queries",
help="Frontier queries searched on every lane per pass; "
"the first N (below) also hit OpenAlex's 10-credit "
"phrase search.")
qd1, qd2, qd3 = st.columns(3, vertical_alignment="bottom")
with qd1:
disc_oa = st.number_input("…of which on OpenAlex phrase search", 0, 500,
int(cfg.get("discovery_openalex_queries", 30)),
disabled=not admin, key="rc_disc_oa")
with qd2:
disc_pages = st.number_input("Topic pages per topic per pass", 0, 200,
int(cfg.get("discovery_topic_pages", 10)),
disabled=not admin, key="rc_disc_pages",
help="OpenAlex topic walk: 50 works per page, 1 credit each.")
with qd3:
disc_per_query = st.number_input("Candidates per query per lane", 1, 200,
int(cfg.get("discovery_per_query", 40)),
disabled=not admin, key="rc_disc_per_query")
_last = qstats.get("last_discovery")
_nxt_d = scheduler.next_discovery_time()
_by_src = qstats.get("queued_by_source", {})
st.html(T.chips_html(
[("queued", f"{_q_depth:,}", "good" if _q_depth else "warn"),
("ingested", f"{int(_qs.get('ingested', 0)):,}", "muted"),
("downloaded", f"{int(_qs.get('downloaded', 0)):,}", "muted"),
("failed", f"{int(_qs.get('failed', 0)):,}", "error" if _qs.get("failed") else "muted"),
("seen", f"{int(_qs.get('seen', 0)):,}", "muted"),
("last pass", T.fmt_ts(_last) or "never", "muted"),
("next pass", _nxt_d.strftime("%b %d %H:%M UTC") if _nxt_d else "—", "muted")]
+ [(f"queued · {k}", f"{v:,}", "muted") for k, v in list(_by_src.items())[:8]]))
if dlock:
st.info("A discovery pass is running right now — follow it in Runs & Logs.")
if st.button("🔍 Run a discovery pass now", disabled=not admin or bool(dlock),
key="rc_disc_now"):
scheduler.run_discovery_now()
st.toast("Discovery pass started in the background.", icon="🔍")
st.success("Discovery pass started — it appears in Runs & Logs as a 'discovery' run.")
st.caption("UDSpace (UD theses) is on by default; CORE needs the free CORE_API_KEY secret "
"and is inactive without it. Semantic Scholar, arXiv, NTRS and OpenAlex follow "
"the Sources & behavior toggles.")
# --- Gemini cost ----------------------------------------------------------------
_THINK_OPTS = ["low", "minimal", "off", "dynamic", "high"]
_think_now = str(cfg.get("gemini_thinking", "low") or "dynamic").lower()
if _think_now not in _THINK_OPTS:
_THINK_OPTS.append(_think_now) # an explicit token budget set by env
_MODEL_OPTS = list(C.MODEL_PRICES)
_model_now = str(cfg.get("gemini_model") or "gemini-3.5-flash").strip()
if _model_now not in _MODEL_OPTS:
_MODEL_OPTS.insert(0, _model_now)
_VMODEL_OPTS = ["(same as text model)"] + list(C.MODEL_PRICES)
_vmodel_now = str(cfg.get("gemini_vision_model") or "").strip() or "(same as text model)"
if _vmodel_now not in _VMODEL_OPTS:
_VMODEL_OPTS.append(_vmodel_now)
def _price_label(m: str) -> str:
pin, pout = C.model_price(m)
return f"{m} (${pin:g} / ${pout:g} per M)"
with st.container(border=True):
T.card_title("Gemini cost", sub="thinking tokens are billed as output (6× the input price "
"on gemini-3.5-flash)")
k1, k2, k3 = st.columns(3, vertical_alignment="bottom")
with k1:
gemini_thinking = st.selectbox("Thinking per call", _THINK_OPTS,
index=_THINK_OPTS.index(_think_now),
disabled=not admin, key="rc_thinking",
help="dynamic = the API default, unbounded (the 28 Sep "
"cycles: ~30k output tokens per PDF for a ~4k-token "
"extraction). low / minimal are Gemini 3 thinking "
"levels; off = budget 0. Applied to every Gemini call "
"(text extraction, figure classification, mining). "
"If the model rejects the option the call is repeated "
"without it.")
with k2:
gemini_model = st.selectbox("Text extraction model", _MODEL_OPTS,
index=_MODEL_OPTS.index(_model_now),
format_func=_price_label,
disabled=not admin, key="rc_model",
help="Every row is stamped with this model. Change it only "
"after comparing on the gold set (cost_ab.py) — the "
"paper's numbers are per model.")
with k3:
gemini_vision_model = st.selectbox("Figure (vision) model", _VMODEL_OPTS,
index=_VMODEL_OPTS.index(_vmodel_now),
format_func=lambda m: m if m.startswith("(")
else _price_label(m),
disabled=not admin, key="rc_vmodel",
help="The one classify call per PDF (and mining, if "
"on). An easy task: the Lite model is ~5x cheaper.")
dedupe_sim = st.slider("Skip near-duplicate PDFs (text similarity ≥)", 0.0, 1.0,
float(cfg.get("dedupe_text_similarity", 0.9) or 0.0), 0.01,
disabled=not admin, key="rc_dedupe",
help="Before any Gemini call, a PDF whose text is this similar to one "
"already ingested (regional datasheet variants, preprint vs "
"published version) is skipped. 0 = off.")
st.caption("Real cost of the 28–29 Sep cycles on gemini-3.5-flash with dynamic thinking and "
"mining on: ~$0.41 per PDF (36k in / 40k out tokens). Watch Cost and $/PDF in "
"Runs & Logs after a change; rows per paper and the flagged share tell you if "
"recall moved.")
# --- figure mining ----------------------------------------------------------
with st.container(border=True):
T.card_title("Figure & graph mining",
sub=f"hard cap {C.HARD_MAX_FIGURES_PER_PDF} figures / PDF",
right_html=_right(T.pill_html("ok", "quarantined by design")))
g1, g2, g3 = st.columns(3, vertical_alignment="bottom")
with g1:
figures_on = st.toggle("Harvest figures each cycle",
value=bool(cfg.get("figures_enabled", True)),
disabled=not admin, key="rc_figures_on")
with g2:
figure_mine = st.toggle("Mine plot/table values",
value=bool(cfg.get("figure_mining", False)),
disabled=not admin, key="rc_figure_mine",
help="OFF (default since 28 Sep) = harvest + classify only: one "
"vision call per PDF, figure linking still works. ON also "
"reads values off property plots and table images — one "
"thinking vision call per figure, into rows that stay "
"quarantined until a human promotes them.")
with g3:
max_figs = st.number_input("Max figures per PDF", 1,
C.HARD_MAX_FIGURES_PER_PDF,
int(cfg.get("max_figures_per_pdf", 12)),
disabled=not admin, key="rc_max_figs")
st.caption("Each ingested PDF also passes through figures.py: plots and table "
"images are detected, captioned and classified (one vision call per "
"PDF), and mineable ones are read for values. A number read off a "
"graph is an estimate — figure rows land as status "
"'figure_estimate' and only a human promote in the Review Queue "
"can publish them.")
link_on = st.toggle("Link text rows to the figures they cite",
value=bool(cfg.get("link_figures", True)),
disabled=not admin, key="rc_link_figures",
help="figure_links.py: a row whose evidence says 'see Fig. 3' "
"gets that figure's id, the citation, a tier score, and "
"the PNG embedded in its image column. Local work, no "
"model call; the row's value and status are untouched.")
with st.container(horizontal=True, vertical_alignment="center", gap="medium"):
if admin and st.button("💾 Save settings", type="primary", key="rc_save"):
conn = agdb.connect()
try:
agdb.set_config(conn, {
"autonomy": autonomy, "interval_hours": float(interval),
"max_new_pdfs_per_cycle": int(max_pdfs),
"max_per_query": int(per_query), "queries_per_cycle": int(n_queries),
"min_new_pdfs_per_cycle": int(target_pdfs),
"max_rounds_per_cycle": int(max_rounds), "ingest_workers": int(workers),
"use_openalex": use_oa, "use_arxiv": use_ax,
"use_semantic_scholar": use_s2, "use_ntrs": use_ntrs,
"expand_queries": expand,
"use_openalex_topics": use_topics, "use_datasheets": use_ds,
"use_candidate_queue": use_queue, "discovery_interval_hours": float(disc_hours),
"queue_low_water": int(low_water), "discovery_queries": int(disc_queries),
"discovery_openalex_queries": int(disc_oa), "discovery_topic_pages": int(disc_pages),
"discovery_per_query": int(disc_per_query),
"figures_enabled": figures_on, "figure_mining": figure_mine,
"gemini_thinking": str(gemini_thinking),
"gemini_model": str(gemini_model),
"gemini_vision_model": ("" if str(gemini_vision_model).startswith("(")
else str(gemini_vision_model)),
"dedupe_text_similarity": float(dedupe_sim),
"max_figures_per_pdf": int(max_figs),
"link_figures": link_on,
})
finally:
conn.close()
scheduler.sync_from_config()
st.toast("Saved — schedule updated.", icon="💾")
st.success("Saved — schedule updated.")
st.rerun()
if admin:
st.caption("Saves autonomy, budget and sources together; the schedule is re-synced immediately.")
# --- manual cycle -----------------------------------------------------------
with row2_right:
with st.container(border=True, height="stretch"):
T.card_title("Manual run", right_html=_right(_state_pill))
if lock:
st.info("A cycle is running right now — watch Runs & Logs.")
if st.button("▶️ Run one cycle now", type="primary",
disabled=not admin or bool(lock), key="rc_run_now"):
scheduler.run_now()
st.toast("Cycle started in the background.", icon="▶️")
st.success("Cycle started in the background — follow it in Runs & Logs.")
st.caption("Fires one full cycle immediately in the background, independent of "
"the autonomy switch. Only one cycle runs at a time; counters are "
"written when it finishes.")
# --- query frontier ---------------------------------------------------------
qdf = df_or_empty(
"SELECT id, enabled, query, origin, times_used, pdfs_found, rows_yielded, "
"last_used_at FROM agent_queries ORDER BY enabled DESC, id")
if qdf.empty:
_q_sub = ""
_q_chips = ""
else:
_n_q = int(len(qdf))
_n_on = int(qdf["enabled"].fillna(False).astype(bool).sum()) if "enabled" in qdf else 0
_n_manual = int((qdf["origin"].astype(str) == "manual").sum()) if "origin" in qdf else 0
_pdfs = _num_sum(qdf, "pdfs_found")
_rows = _num_sum(qdf, "rows_yielded")
_q_sub = f"{_n_q:,} quer{'y' if _n_q == 1 else 'ies'} · {_n_on:,} enabled"
_q_chips = T.chips_html([
("manual", f"{_n_manual:,}", "muted"),
("PDFs found", f"{_pdfs:,}", "muted"),
("rows yielded", f"{_rows:,}", "muted"),
])
with st.container(border=True):
T.card_title("Query frontier", sub=_q_sub, right_html=_right(_q_chips))
if qdf.empty:
T.empty_state("🧭", "No queries yet",
"The frontier is seeded on first boot; add your own search intent below.")
else:
edited = st.data_editor(
qdf, hide_index=True, width='stretch',
disabled=[] if admin else list(qdf.columns),
column_config={
"id": st.column_config.NumberColumn("ID", width=50, format="%d"),
"enabled": st.column_config.CheckboxColumn("On", width=58),
"query": st.column_config.TextColumn("Query"),
"origin": st.column_config.TextColumn("Origin", width=85),
"times_used": st.column_config.NumberColumn("Times used", width=105, format="%d"),
"pdfs_found": st.column_config.NumberColumn("PDFs found", width=105, format="%d"),
"rows_yielded": st.column_config.NumberColumn("Rows yielded", width=115, format="%d"),
"last_used_at": st.column_config.DatetimeColumn(
"Last used", format="MMM D, HH:mm", width=110),
},
key="query_editor")
st.caption("Each cycle takes the least-recently-used enabled queries; edit the text or "
"toggle a query, then save. Only 'On' and 'Query' are stored.")
if admin and st.button("💾 Save query changes", key="rc_save_queries"):
conn = agdb.connect()
try:
with conn.cursor() as cur:
for _, row in edited.iterrows():
cur.execute(
"UPDATE agent_queries SET enabled = %s, query = %s WHERE id = %s",
(bool(row["enabled"]), str(row["query"]), int(row["id"])))
conn.commit()
finally:
conn.close()
st.toast("Queries updated.", icon="💾")
st.success("Queries updated.")
st.rerun()
with st.container(horizontal=True, vertical_alignment="bottom", gap="small"):
new_q = st.text_input("Add a query", placeholder="e.g. LCP glass fiber dielectric properties",
disabled=not admin, key="rc_new_query")
if admin and st.button("➕ Add query", key="rc_add_query") and new_q.strip():
conn = agdb.connect()
try:
agdb.seed_queries(conn, [new_q.strip()], origin="manual")
finally:
conn.close()
st.toast("Query added to the frontier.", icon="➕")
st.rerun()