"""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("""""") def _right(html: str) -> str: """Wrap pill/chip HTML for card_title(right_html=...) so it keeps its casing.""" return f'{html}' 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('
' 'Unlock admin in the sidebar to change settings — everything below is read-only.
') # --------------------------------------------------------------------------- # 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()