From d4e6313b63b92a379d5a903beb091d12dd4a3f5d Mon Sep 17 00:00:00 2001 From: benjsmith Date: Fri, 18 Sep 2026 20:30:31 +0000 Subject: [PATCH 01/25] feat: Phase 4a same-origin embed reverse-proxy (no iframes) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Daemon proxies /embed/ce/* → 127.0.0.1:8766 and /embed/okstratr/* → :8767 with loopback-only upstream guard and X-CE-Host / X-Okstratr-Host: switchbay. Feature-flag Graph/Agents to proxied same-document panels (default off). ADR-004 + allowlist tests. --- docs/ADR-004-same-origin-embed-proxy.md | 59 ++++ frontend/src/center/builtinTabs.tsx | 16 +- frontend/src/index.css | 73 +++++ frontend/src/layout/SettingsModal.tsx | 66 ++++ .../src/widgets/embed/ProxiedSkillPanel.tsx | 183 +++++++++++ .../widgets/embed/useProxiedSkillEmbeds.ts | 38 +++ frontend/vite.config.ts | 2 + src/switchbay/app_settings.py | 22 ++ src/switchbay/daemon.py | 14 +- src/switchbay/embed_proxy.py | 307 ++++++++++++++++++ tests/unit/test_embed_proxy.py | 145 +++++++++ 11 files changed, 919 insertions(+), 6 deletions(-) create mode 100644 docs/ADR-004-same-origin-embed-proxy.md create mode 100644 frontend/src/widgets/embed/ProxiedSkillPanel.tsx create mode 100644 frontend/src/widgets/embed/useProxiedSkillEmbeds.ts create mode 100644 src/switchbay/embed_proxy.py create mode 100644 tests/unit/test_embed_proxy.py diff --git a/docs/ADR-004-same-origin-embed-proxy.md b/docs/ADR-004-same-origin-embed-proxy.md new file mode 100644 index 0000000..26a52f5 --- /dev/null +++ b/docs/ADR-004-same-origin-embed-proxy.md @@ -0,0 +1,59 @@ +# ADR-004: Same-origin embed reverse-proxy (no iframes) + +- **Status:** Accepted (Phase 4a) +- **Date:** 2026-09-18 +- **Deciders:** Ben / skill-shell rationalization charter + +## Context + +Switchbay hosts Graph (curiosity-engine atlas/wiki) and Agents (okstratr +observer/desk) surfaces. Cross-origin iframes pointed at +`127.0.0.1:8766` / `:8767` break hosted-mode control (okstratr HTML +settings must hide when `host=switchbay`) and create an opaque nested +browsing context. + +Charter locked decision #1: shells **same-origin reverse-proxy** skill +daemons; in-app panels load first-party proxied routes — **not nested +frames**. + +## Decision + +1. **Daemon reverse-proxy** (always on; independent of the UI flag): + - `/embed/ce/*` → `http://127.0.0.1:8766/*` (override: + `SWITCHBAY_CE_UPSTREAM`, must remain loopback) + - `/embed/okstratr/*` → `http://127.0.0.1:8767/*` (override: + `SWITCHBAY_OKSTRATR_UPSTREAM`, must remain loopback) + - Upstream allowlist is **loopback-only** (`127.0.0.1` / `::1` / + `localhost`). Non-loopback upstreams are rejected (502). + - Inject hosted-shell headers on every proxied request: + - `X-CE-Host: switchbay` + - `X-Okstratr-Host: switchbay` + +2. **No iframes** for Graph/Agents skill surfaces. Feature flag + `proxied_skill_embeds` (default **false**) switches Graph → CE panel + and Agents → okstratr panel that navigate `/embed/*` via same-origin + `fetch` + same-document rendering (script-stripped HTML / JSON). + Built-in GraphTab / AgentDashboardTab / filebrowser remain the + default and are **not deleted**. + +3. **Settings → okstratr registry (TODO).** Switchbay settings will + become a client that writes okstratr's harness/model registry + (`harnesses.toml` via API). Out of scope for 4a; tracked here so the + shell does not grow a second allowlist. + +## Consequences + +- Dev Vite must proxy `/embed` to the daemon (`vite.config.ts`). +- CE and okstratr should honor public-base + `X-*-Host` (okstratr + Phase 1a `public_base.py`; CE equivalent TBD). +- Full atlas/observer chrome parity is gated by the parity checklist + before any Switchbay duplicate deletion (Phase 4b+). +- WebSocket upgrade through the embed proxy is deferred; HTTP(S) first. + +## Alternatives considered + +| Option | Why not | +|--------|---------| +| Cross-origin iframe to :8766/:8767 | Opaque origin; hosted settings leak; charter forbid | +| Same-origin iframe to `/embed/*` | Still a nested frame; charter: no iframes | +| Delete built-in Graph/Agents now | Violates parity checklist / dual-stack gate | diff --git a/frontend/src/center/builtinTabs.tsx b/frontend/src/center/builtinTabs.tsx index 5ceb3b2..7452403 100644 --- a/frontend/src/center/builtinTabs.tsx +++ b/frontend/src/center/builtinTabs.tsx @@ -1,6 +1,8 @@ import { lazy } from "react"; import SketchErrorBoundary from "../widgets/sketch/ErrorBoundary"; import { registerTabKind, type TabComponent } from "./tabRegistry"; +import ProxiedSkillPanel from "../widgets/embed/ProxiedSkillPanel"; +import { useProxiedSkillEmbeds } from "../widgets/embed/useProxiedSkillEmbeds"; /** * Wire each built-in tab kind into the registry. Called once from @@ -40,9 +42,11 @@ const ReportDocTab = lazy(() => import("../widgets/library/ReportDocTab")); const ThrustersTab = lazy(() => import("../widgets/thrusters/ThrustersTab")); const OwidTab = lazy(() => import("../widgets/owid/OwidTab")); -const GraphAdapter: TabComponent = ({ graphData, graphError }) => ( - -); +const GraphAdapter: TabComponent = ({ graphData, graphError }) => { + const proxied = useProxiedSkillEmbeds(); + if (proxied) return ; + return ; +}; const EditorAdapter: TabComponent = ({ tab }) => ; const DuckDBAdapter: TabComponent = () => ; const SheetAdapter: TabComponent = () => ; @@ -50,7 +54,11 @@ const VegaAdapter: TabComponent = () => ; const SketchAdapter: TabComponent = () => ( ); -const AgentsAdapter: TabComponent = () => ; +const AgentsAdapter: TabComponent = () => { + const proxied = useProxiedSkillEmbeds(); + if (proxied) return ; + return ; +}; const ProjectsAdapter: TabComponent = () => ; const PackFileListAdapter: TabComponent = ({ tab }) => ; diff --git a/frontend/src/index.css b/frontend/src/index.css index f212a1b..da36056 100644 --- a/frontend/src/index.css +++ b/frontend/src/index.css @@ -8986,3 +8986,76 @@ a.sy-source-cite { color: var(--accent); } color: var(--danger, #c44); font-size: 10px; } + + +/* ── Phase 4a: same-origin proxied skill panel (no iframe) ───────── */ +.sy-proxied-skill { + display: flex; + flex-direction: column; + height: 100%; + min-height: 0; + background: var(--sy-bg, #0f1115); + color: var(--sy-fg, #e8eaed); +} +.sy-proxied-skill-bar { + display: flex; + flex-wrap: wrap; + align-items: center; + gap: 8px 12px; + padding: 8px 12px; + border-bottom: 1px solid var(--sy-border, #2a2f3a); + flex-shrink: 0; +} +.sy-proxied-skill-hint { + opacity: 0.65; + font-size: 0.85em; +} +.sy-proxied-skill-nav { + display: flex; + gap: 6px; + flex: 1; + min-width: 200px; +} +.sy-proxied-skill-nav input { + flex: 1; + font-family: ui-monospace, SFMono-Regular, Menlo, monospace; + font-size: 0.85em; + padding: 4px 8px; + border-radius: 4px; + border: 1px solid var(--sy-border, #2a2f3a); + background: var(--sy-input-bg, #1a1f2a); + color: inherit; +} +.sy-proxied-skill-nav button { + padding: 4px 10px; + border-radius: 4px; + border: 1px solid var(--sy-border, #2a2f3a); + background: var(--sy-btn-bg, #222833); + color: inherit; + cursor: pointer; +} +.sy-proxied-skill-body { + flex: 1; + min-height: 0; + overflow: auto; + padding: 12px; +} +.sy-proxied-skill-muted { opacity: 0.7; } +.sy-proxied-skill-error { + border: 1px solid #a44; + border-radius: 6px; + padding: 12px; + background: rgba(160, 40, 40, 0.12); +} +.sy-proxied-skill-error pre, +.sy-proxied-skill-json, +.sy-proxied-skill-text { + white-space: pre-wrap; + word-break: break-word; + font-family: ui-monospace, SFMono-Regular, Menlo, monospace; + font-size: 0.85em; +} +.sy-proxied-skill-html { + /* Hosted skill HTML body, scripts stripped — same-document panel. */ + max-width: 100%; +} diff --git a/frontend/src/layout/SettingsModal.tsx b/frontend/src/layout/SettingsModal.tsx index fca160e..63c42e8 100644 --- a/frontend/src/layout/SettingsModal.tsx +++ b/frontend/src/layout/SettingsModal.tsx @@ -1732,6 +1732,7 @@ type SettingsBody = { admin_ceiling?: number | null; chief_counted?: boolean; note?: string; + proxied_skill_embeds?: boolean; media?: { modalities?: Record; note?: string; @@ -3646,6 +3647,42 @@ function StoragePanel({ open }: { open: boolean }) { } }; + const toggleProxiedEmbeds = async () => { + if (!settings || busy) return; + const next = !settings.proxied_skill_embeds; + setBusy(true); + setStatus(null); + try { + const r = await fetch("/api/settings", { + method: "POST", + headers: { "Content-Type": "application/json" }, + body: JSON.stringify({ proxied_skill_embeds: next }), + }); + if (!r.ok) { + const body = await r.json().catch(() => ({} as Record)); + setStatus({ ok: false, msg: body.error || `HTTP ${r.status}` }); + return; + } + const body = (await r.json()) as SettingsBody; + setSettings(body); + window.dispatchEvent( + new CustomEvent("sy:proxied-skill-embeds", { + detail: { enabled: !!body.proxied_skill_embeds }, + }), + ); + setStatus({ + ok: true, + msg: next + ? "Graph/Agents now use same-origin /embed/* panels (no iframe)." + : "Graph/Agents restored to built-in tabs.", + }); + } catch (e) { + setStatus({ ok: false, msg: (e as Error).message }); + } finally { + setBusy(false); + } + }; + const setEmbeddingBackend = async (backend: string) => { if (!settings || busy) return; setBusy(true); @@ -3728,6 +3765,35 @@ function StoragePanel({ open }: { open: boolean }) { Machine-local avoids that at the cost of per-machine history.

+

+ Proxied skill embeds (Phase 4a) +

+

+ When on, Graph and Agents load CE / okstratr through the daemon + same-origin reverse proxy (/embed/ce/,{" "} + /embed/okstratr/) — in-app panels, no iframes. + Default off keeps the built-in tabs. Settings will later write the + okstratr harness registry (TODO). +

+
+ + Proxied embeds:{" "} + {settings.proxied_skill_embeds ? "on (Graph→CE, Agents→okstratr)" : "off (built-in)"} + + + +
+

Semantic recall (conversation memory)

diff --git a/frontend/src/widgets/embed/ProxiedSkillPanel.tsx b/frontend/src/widgets/embed/ProxiedSkillPanel.tsx new file mode 100644 index 0000000..36eda98 --- /dev/null +++ b/frontend/src/widgets/embed/ProxiedSkillPanel.tsx @@ -0,0 +1,183 @@ +import { useCallback, useEffect, useState } from "react"; + +/** + * Phase 4a same-origin skill panel (NO iframe). + * + * Loads CE or okstratr through the daemon reverse-proxy under + * `/embed/ce/*` or `/embed/okstratr/*`. In-app path navigation uses + * fetch + a same-document panel — never a nested frame. + * + * Full atlas/observer chrome lands as those skills grow hosted-mode + * fragment/API UIs; this panel proves the proxy path and stays usable + * when upstream is down. + */ + +export type SkillEmbedKind = "ce" | "okstratr"; + +const PREFIX: Record = { + ce: "/embed/ce", + okstratr: "/embed/okstratr", +}; + +const LABEL: Record = { + ce: "Curiosity Engine", + okstratr: "okstratr", +}; + +const DEFAULT_PATH: Record = { + ce: "/", + okstratr: "/observer/", +}; + +type Props = { + kind: SkillEmbedKind; +}; + +type LoadState = + | { status: "loading" } + | { status: "ok"; contentType: string; text: string; httpStatus: number } + | { status: "error"; message: string; httpStatus?: number }; + +function stripScripts(html: string): string { + return html + .replace(/)<[^<]*)*<\/script>/gi, "") + .replace(/\son\w+="[^"]*"/gi, "") + .replace(/\son\w+='[^']*'/gi, ""); +} + +function extractBody(html: string): string { + const m = html.match(/]*>([\s\S]*)<\/body>/i); + return m ? m[1] : html; +} + +export default function ProxiedSkillPanel({ kind }: Props) { + const prefix = PREFIX[kind]; + const [path, setPath] = useState(DEFAULT_PATH[kind]); + const [draft, setDraft] = useState(DEFAULT_PATH[kind]); + const [state, setState] = useState({ status: "loading" }); + + const load = useCallback( + async (p: string) => { + setState({ status: "loading" }); + const url = prefix + (p.startsWith("/") ? p : `/${p}`); + try { + const r = await fetch(url, { + headers: { Accept: "text/html, application/json;q=0.9, */*;q=0.8" }, + }); + const ct = r.headers.get("content-type") || ""; + const text = await r.text(); + if (!r.ok) { + setState({ + status: "error", + message: text.slice(0, 500) || r.statusText, + httpStatus: r.status, + }); + return; + } + setState({ + status: "ok", + contentType: ct, + text, + httpStatus: r.status, + }); + } catch (e) { + setState({ + status: "error", + message: (e as Error).message || "fetch failed", + }); + } + }, + [prefix], + ); + + useEffect(() => { + void load(path); + }, [load, path]); + + const onNavigate = (ev: React.FormEvent) => { + ev.preventDefault(); + const next = draft.trim() || "/"; + setPath(next.startsWith("/") ? next : `/${next}`); + }; + + return ( +
+
+ {LABEL[kind]} + + same-origin via {prefix} (no iframe) + +
+ setDraft(e.target.value)} + aria-label="Proxied path" + spellCheck={false} + /> + + +
+
+ +
+ {state.status === "loading" && ( +

Loading {prefix}{path}…

+ )} + {state.status === "error" && ( +
+

+ Upstream unreachable or proxy error + {state.httpStatus != null ? ` (${state.httpStatus})` : ""}. +

+
{state.message}
+

+ Start the skill daemon on loopback + {kind === "ce" ? " :8766" : " :8767"} (charter), or turn off + Settings → “Proxied skill embeds” to use the built-in tab. +

+
+ )} + {state.status === "ok" && ( + + )} +
+
+ ); +} + +function ProxiedContent({ + contentType, + text, +}: { + contentType: string; + text: string; +}) { + if (contentType.includes("application/json")) { + let pretty = text; + try { + pretty = JSON.stringify(JSON.parse(text), null, 2); + } catch { + /* keep raw */ + } + return
{pretty}
; + } + if (contentType.includes("text/html") || /^\s* + ); + } + return
{text}
; +} diff --git a/frontend/src/widgets/embed/useProxiedSkillEmbeds.ts b/frontend/src/widgets/embed/useProxiedSkillEmbeds.ts new file mode 100644 index 0000000..2d084ad --- /dev/null +++ b/frontend/src/widgets/embed/useProxiedSkillEmbeds.ts @@ -0,0 +1,38 @@ +import { useEffect, useState } from "react"; + +/** + * Phase 4a feature flag: when true, Graph/Agents use `/embed/*` panels. + * Default false (built-in tabs). Reads `/api/settings` once; listens for + * `sy:proxied-skill-embeds` so Settings can flip without reload. + */ +export function useProxiedSkillEmbeds(): boolean { + const [on, setOn] = useState(false); + + useEffect(() => { + let cancelled = false; + void fetch("/api/settings") + .then((r) => (r.ok ? r.json() : null)) + .then((j) => { + if (!cancelled && j && typeof j.proxied_skill_embeds === "boolean") { + setOn(j.proxied_skill_embeds); + } + }) + .catch(() => { + /* older daemon — leave off */ + }); + + const onEvt = (ev: Event) => { + const detail = (ev as CustomEvent<{ enabled?: boolean }>).detail; + if (detail && typeof detail.enabled === "boolean") { + setOn(detail.enabled); + } + }; + window.addEventListener("sy:proxied-skill-embeds", onEvt); + return () => { + cancelled = true; + window.removeEventListener("sy:proxied-skill-embeds", onEvt); + }; + }, []); + + return on; +} diff --git a/frontend/vite.config.ts b/frontend/vite.config.ts index 0266f6d..ba9bbe9 100644 --- a/frontend/vite.config.ts +++ b/frontend/vite.config.ts @@ -18,6 +18,8 @@ export default defineConfig({ // daemon so handle_figure_file can serve from /figures/ or // /wiki/figures/ (PDF page rasters, sketch PNGs, etc.). "/figures": DAEMON, + // Phase 4a: same-origin embed reverse-proxy (CE + okstratr). + "/embed": DAEMON, }, }, }); diff --git a/src/switchbay/app_settings.py b/src/switchbay/app_settings.py index a06ed57..0735a1f 100644 --- a/src/switchbay/app_settings.py +++ b/src/switchbay/app_settings.py @@ -13,6 +13,7 @@ "workspaces_home": "~/Workspaces", # see below "orchestration_preference": 0.5, # 0 Economy … 1 Maximum "orchestration_denied_models": [] # seed until a workspace file exists + "proxied_skill_embeds": false # Phase 4a: Graph/Agents via /embed/* } `rail_history_local` — where the per-workspace rail-history DB @@ -32,6 +33,7 @@ from __future__ import annotations import json +import os from typing import Any from . import workspaces @@ -158,3 +160,23 @@ def set_desk_max_live_workers(value: int) -> int: data["desk_max_live_workers"] = n save(data) return n + + +# `proxied_skill_embeds` — Phase 4a feature flag. When true, the Graph and +# Agents tabs use same-origin proxied panels under `/embed/ce/` and +# `/embed/okstratr/` instead of the built-in implementations. Default +# false so existing UX is unchanged. Env SWITCHBAY_PROXIED_SKILL_EMBEDS=1 +# overrides settings.json to true (useful for local E2E). +def get_proxied_skill_embeds() -> bool: + env = (os.environ.get("SWITCHBAY_PROXIED_SKILL_EMBEDS") or "").strip().lower() + if env in ("1", "true", "yes", "on"): + return True + if env in ("0", "false", "no", "off"): + return False + return bool(load().get("proxied_skill_embeds", False)) + + +def set_proxied_skill_embeds(value: bool) -> None: + data = load() + data["proxied_skill_embeds"] = bool(value) + save(data) diff --git a/src/switchbay/daemon.py b/src/switchbay/daemon.py index 1436a26..8c6da4f 100644 --- a/src/switchbay/daemon.py +++ b/src/switchbay/daemon.py @@ -37,7 +37,7 @@ ce_tools, command_palettes, commands, conversations, curation_history, dbintrospect, demo_workspace, - duckdb_starters, file_state, fileops, llm_config, llmgateway, + duckdb_starters, embed_proxy, file_state, fileops, llm_config, llmgateway, localllm, orchestrator_fs, schedules, mcpstore, merging, model_cache, modestore, owid, packstore, pasteboard, permissions, plots, html_decks, library, local_models, media_settings, micro_edits, projects, proposals, protocol, rail, report_html, report_packages, reports, secrets, selection, service, share, sheets, @@ -905,6 +905,7 @@ async def handle_settings_get(request: web.Request) -> web.Response: "orchestration_preference": orchestration_policy.get_preference(), "orchestration_denied_models": orchestration_policy.get_denied_models(workspace), **_desk_live_view(workspace), + "proxied_skill_embeds": app_settings.get_proxied_skill_embeds(), }) @@ -957,6 +958,8 @@ async def handle_settings_post(request: web.Request) -> web.Response: # Re-select the backend on next use; the drain's reconcile then # rebuilds the index if the vector space changed. conversations.reset_embedder() + if "proxied_skill_embeds" in body: + app_settings.set_proxied_skill_embeds(bool(body["proxied_skill_embeds"])) # Media: { media: { image?: {provider, model}|null, video?: …, voice?: … } } if "media" in body and isinstance(body["media"], dict): if not admin_policy.feature_enabled("media_generation"): @@ -17091,7 +17094,11 @@ async def handle_spa(request: web.Request) -> web.StreamResponse: # Never serve the SPA for API/WS paths — an unregistered /api GET # must 404 (so clients that probe GET-then-POST fall back correctly), # not silently receive index.html. - if rel == "api" or rel.startswith("api/") or rel == "ws" or rel.startswith("ws/"): + if ( + rel == "api" or rel.startswith("api/") + or rel == "ws" or rel.startswith("ws/") + or rel == "embed" or rel.startswith("embed/") + ): return web.Response(status=404) if rel: candidate = (dist / rel).resolve() @@ -17448,6 +17455,9 @@ def build_app(workspace: Path) -> web.Application: app.router.add_post("/api/tabs/vault-doc", handle_tab_vault_doc_add) app.router.add_post("/api/tabs/vault-doc/remove", handle_tab_vault_doc_remove) app.router.add_get("/ws", handle_ws) + # Phase 4a: same-origin reverse proxy for CE (:8766) and okstratr + # (:8767). Must register BEFORE the SPA catch-all. + embed_proxy.register_routes(app) # Catch-all LAST: serves the built SPA + PWA assets (manifest, icons) # for any non-API GET. aiohttp matches in registration order, so the # specific /api and /ws routes above always take precedence. diff --git a/src/switchbay/embed_proxy.py b/src/switchbay/embed_proxy.py new file mode 100644 index 0000000..e46a983 --- /dev/null +++ b/src/switchbay/embed_proxy.py @@ -0,0 +1,307 @@ +"""Same-origin reverse proxy for CE + okstratr (Phase 4a). + +Charter locked decision #1: Switchbay reverse-proxies loopback skill +daemons under ``/embed/ce/`` and ``/embed/okstratr/`` so in-app panels +can load first-party routes (no cross-origin iframes). + +Upstream targets are **loopback-only** (127.0.0.1 / ::1). Non-loopback +upstream configuration is rejected. Proxied requests carry hosted-shell +headers so CE/okstratr can enter hosted mode: + + X-CE-Host: switchbay + X-Okstratr-Host: switchbay +""" + +from __future__ import annotations + +import ipaddress +import logging +import os +from typing import Mapping +from urllib.parse import urlsplit, urlunsplit + +from aiohttp import ClientSession, ClientTimeout, web + +log = logging.getLogger("switchbay.embed_proxy") + +# Public path prefixes (no trailing slash). +PREFIX_CE = "/embed/ce" +PREFIX_OKSTRATR = "/embed/okstratr" + +HOST_HEADER_CE = "X-CE-Host" +HOST_HEADER_OKSTRATR = "X-Okstratr-Host" +HOSTED_SHELL = "switchbay" + +# Charter defaults: CE :8766, okstratr :8767. +_DEFAULT_CE = "http://127.0.0.1:8766" +_DEFAULT_OKSTRATR = "http://127.0.0.1:8767" + +_HOP_BY_HOP = frozenset({ + "connection", + "keep-alive", + "proxy-authenticate", + "proxy-authorization", + "te", + "trailers", + "transfer-encoding", + "upgrade", + "host", + "content-length", +}) + + +class UpstreamNotLoopback(ValueError): + """Raised when an embed upstream is not a loopback address.""" + + +def _env_upstream(name: str, default: str) -> str: + raw = (os.environ.get(name) or "").strip() + return raw or default + + +def ce_upstream() -> str: + return _normalize_base(_env_upstream("SWITCHBAY_CE_UPSTREAM", _DEFAULT_CE)) + + +def okstratr_upstream() -> str: + return _normalize_base( + _env_upstream("SWITCHBAY_OKSTRATR_UPSTREAM", _DEFAULT_OKSTRATR) + ) + + +def _normalize_base(url: str) -> str: + u = url.rstrip("/") + if "://" not in u: + u = "http://" + u + return u + + +def is_loopback_host(host: str) -> bool: + """True iff *host* (no port) is a loopback literal.""" + h = (host or "").strip().lower() + if h in ("localhost", "127.0.0.1", "::1"): + return True + # Bracketed IPv6 from urlsplit netloc + if h.startswith("[") and h.endswith("]"): + h = h[1:-1] + try: + return ipaddress.ip_address(h).is_loopback + except ValueError: + return False + + +def assert_loopback_upstream(url: str) -> str: + """Validate upstream URL is http(s) to a loopback host. Return normalized base.""" + base = _normalize_base(url) + parts = urlsplit(base) + if parts.scheme not in ("http", "https"): + raise UpstreamNotLoopback(f"embed upstream scheme must be http(s): {url!r}") + host = parts.hostname or "" + if not is_loopback_host(host): + raise UpstreamNotLoopback( + f"embed upstream must be loopback-only (got host {host!r} from {url!r})" + ) + # Rebuild without path/query/fragment so join is predictable + netloc = parts.netloc + return urlunsplit((parts.scheme, netloc, "", "", "")).rstrip("/") + + +def embed_targets() -> dict[str, dict[str, str]]: + """Prefix → {upstream, host_header, host_value} (validated loopback).""" + return { + PREFIX_CE: { + "upstream": assert_loopback_upstream(ce_upstream()), + "host_header": HOST_HEADER_CE, + "host_value": HOSTED_SHELL, + }, + PREFIX_OKSTRATR: { + "upstream": assert_loopback_upstream(okstratr_upstream()), + "host_header": HOST_HEADER_OKSTRATR, + "host_value": HOSTED_SHELL, + }, + } + + +def match_embed_prefix(path: str) -> tuple[str, str] | None: + """Return ``(prefix, rest)`` if *path* is under an embed prefix. + + *rest* is the upstream path including a leading ``/`` (or ``/`` when + the request targeted the prefix alone). + """ + p = path or "/" + for prefix in (PREFIX_CE, PREFIX_OKSTRATR): + if p == prefix or p == prefix + "/": + return prefix, "/" + if p.startswith(prefix + "/"): + rest = p[len(prefix) :] + return prefix, rest if rest else "/" + return None + + +def build_upstream_url(prefix: str, rest: str, query: str = "") -> str: + """Compose absolute upstream URL for a matched embed request.""" + targets = embed_targets() + if prefix not in targets: + raise KeyError(prefix) + base = targets[prefix]["upstream"] + path = rest if rest.startswith("/") else "/" + rest + q = f"?{query}" if query else "" + return f"{base}{path}{q}" + + +def filter_request_headers( + headers: Mapping[str, str], + *, + host_header: str, + host_value: str, +) -> dict[str, str]: + """Copy inbound headers minus hop-by-hop; inject hosted-shell header.""" + out: dict[str, str] = {} + for k, v in headers.items(): + if k.lower() in _HOP_BY_HOP: + continue + # Drop client-supplied hosted-shell headers; we set them. + if k.lower() in (HOST_HEADER_CE.lower(), HOST_HEADER_OKSTRATR.lower()): + continue + out[k] = v + out[host_header] = host_value + return out + + +def filter_response_headers(headers: Mapping[str, str]) -> dict[str, str]: + out: dict[str, str] = {} + for k, v in headers.items(): + if k.lower() in _HOP_BY_HOP: + continue + out[k] = v + return out + + +async def _proxy_http(request: web.Request) -> web.StreamResponse: + matched = match_embed_prefix(request.path) + if not matched: + return web.Response(status=404, text="not an embed path") + prefix, rest = matched + try: + targets = embed_targets() + except UpstreamNotLoopback as e: + log.error("embed upstream rejected: %s", e) + return web.json_response({"error": str(e)}, status=502) + + meta = targets[prefix] + upstream_url = build_upstream_url(prefix, rest, request.query_string) + # Defense in depth: re-check the concrete URL we are about to hit. + try: + assert_loopback_upstream( + urlunsplit(urlsplit(upstream_url)[:2] + ("", "", "")) + ) + except UpstreamNotLoopback as e: + return web.json_response({"error": str(e)}, status=502) + + req_headers = filter_request_headers( + request.headers, + host_header=meta["host_header"], + host_value=meta["host_value"], + ) + body = await request.read() + timeout = ClientTimeout(total=120, connect=5, sock_connect=5) + + session: ClientSession | None = request.app.get("embed_http") + owns_session = False + if session is None or session.closed: + session = ClientSession(timeout=timeout) + owns_session = True + + try: + async with session.request( + request.method, + upstream_url, + headers=req_headers, + data=body if body else None, + allow_redirects=False, + ) as upstream: + resp_headers = filter_response_headers(upstream.headers) + payload = await upstream.read() + return web.Response( + status=upstream.status, + body=payload, + headers=resp_headers, + ) + except OSError as e: + log.warning("embed upstream unreachable %s: %s", upstream_url, e) + return web.json_response( + { + "error": "embed upstream unreachable", + "upstream": meta["upstream"], + "prefix": prefix, + "detail": str(e), + }, + status=502, + ) + except Exception as e: # noqa: BLE001 + log.exception("embed proxy failed for %s", upstream_url) + return web.json_response( + {"error": "embed proxy failed", "detail": str(e)}, + status=502, + ) + finally: + if owns_session and session is not None: + await session.close() + + +async def handle_embed(request: web.Request) -> web.StreamResponse: + """HTTP reverse-proxy handler for ``/embed/ce/*`` and ``/embed/okstratr/*``.""" + if request.method == "OPTIONS": + # Let upstream decide CORS; we still inject host headers via proxy. + return await _proxy_http(request) + return await _proxy_http(request) + + +def register_routes(app: web.Application) -> None: + """Mount embed proxy routes (call before the SPA catch-all).""" + # Bare prefix + wildcard. Methods mirror what skill UIs typically need. + paths = ( + PREFIX_CE, + PREFIX_CE + "/", + PREFIX_CE + "/{tail:.*}", + PREFIX_OKSTRATR, + PREFIX_OKSTRATR + "/", + PREFIX_OKSTRATR + "/{tail:.*}", + ) + methods = ("GET", "POST", "PUT", "PATCH", "DELETE", "OPTIONS", "HEAD") + for path in paths: + for method in methods: + app.router.add_route(method, path, handle_embed) + + async def _open_session(application: web.Application) -> None: + application["embed_http"] = ClientSession( + timeout=ClientTimeout(total=120, connect=5, sock_connect=5), + ) + + async def _close_session(application: web.Application) -> None: + sess = application.get("embed_http") + if sess is not None and not sess.closed: + await sess.close() + + app.on_startup.append(_open_session) + app.on_cleanup.append(_close_session) + + +# ── allowlist helpers (unit-tested) ───────────────────────────────────── + +def allowed_upstream_bases() -> frozenset[str]: + """Current configured (validated) upstream bases.""" + t = embed_targets() + return frozenset(v["upstream"] for v in t.values()) + + +def upstream_allowed(url: str) -> bool: + """True if *url* targets a configured loopback embed upstream.""" + try: + parts = urlsplit(url) + base = assert_loopback_upstream( + urlunsplit((parts.scheme, parts.netloc, "", "", "")) + ) + except UpstreamNotLoopback: + return False + return base in allowed_upstream_bases() diff --git a/tests/unit/test_embed_proxy.py b/tests/unit/test_embed_proxy.py new file mode 100644 index 0000000..550f93a --- /dev/null +++ b/tests/unit/test_embed_proxy.py @@ -0,0 +1,145 @@ +"""Phase 4a: same-origin embed reverse-proxy allowlist + routing.""" + +from __future__ import annotations + +import pytest +from aiohttp import web +from aiohttp.test_utils import TestClient, TestServer + +from switchbay import embed_proxy + + +def test_loopback_hosts_accepted(): + assert embed_proxy.is_loopback_host("127.0.0.1") + assert embed_proxy.is_loopback_host("localhost") + assert embed_proxy.is_loopback_host("::1") + assert embed_proxy.is_loopback_host("[::1]") + + +def test_non_loopback_hosts_rejected(): + assert not embed_proxy.is_loopback_host("example.com") + assert not embed_proxy.is_loopback_host("8.8.8.8") + assert not embed_proxy.is_loopback_host("10.0.0.1") + assert not embed_proxy.is_loopback_host("") + + +def test_assert_loopback_upstream_ok(): + assert embed_proxy.assert_loopback_upstream("http://127.0.0.1:8766") == ( + "http://127.0.0.1:8766" + ) + assert embed_proxy.assert_loopback_upstream("http://localhost:8767/") == ( + "http://localhost:8767" + ) + + +def test_assert_loopback_upstream_rejects_remote(): + with pytest.raises(embed_proxy.UpstreamNotLoopback): + embed_proxy.assert_loopback_upstream("http://evil.example:8766") + with pytest.raises(embed_proxy.UpstreamNotLoopback): + embed_proxy.assert_loopback_upstream("http://192.168.1.5:8766") + + +def test_match_embed_prefix(): + assert embed_proxy.match_embed_prefix("/embed/ce") == ("/embed/ce", "/") + assert embed_proxy.match_embed_prefix("/embed/ce/") == ("/embed/ce", "/") + assert embed_proxy.match_embed_prefix("/embed/ce/api/health") == ( + "/embed/ce", + "/api/health", + ) + assert embed_proxy.match_embed_prefix("/embed/okstratr/observer/") == ( + "/embed/okstratr", + "/observer/", + ) + assert embed_proxy.match_embed_prefix("/api/health") is None + assert embed_proxy.match_embed_prefix("/embed/other") is None + + +def test_build_upstream_url_defaults(monkeypatch): + monkeypatch.delenv("SWITCHBAY_CE_UPSTREAM", raising=False) + monkeypatch.delenv("SWITCHBAY_OKSTRATR_UPSTREAM", raising=False) + assert embed_proxy.build_upstream_url("/embed/ce", "/api/x", "") == ( + "http://127.0.0.1:8766/api/x" + ) + assert embed_proxy.build_upstream_url( + "/embed/okstratr", "/observer/", "host=x" + ) == "http://127.0.0.1:8767/observer/?host=x" + + +def test_filter_request_headers_injects_host(): + h = embed_proxy.filter_request_headers( + {"Accept": "application/json", "Host": "127.0.0.1:8765", + "X-CE-Host": "forged"}, + host_header=embed_proxy.HOST_HEADER_CE, + host_value="switchbay", + ) + assert h["X-CE-Host"] == "switchbay" + assert "Host" not in h + assert h["Accept"] == "application/json" + + +def test_upstream_allowed_allowlist(monkeypatch): + monkeypatch.delenv("SWITCHBAY_CE_UPSTREAM", raising=False) + monkeypatch.delenv("SWITCHBAY_OKSTRATR_UPSTREAM", raising=False) + assert embed_proxy.upstream_allowed("http://127.0.0.1:8766/foo") + assert embed_proxy.upstream_allowed("http://127.0.0.1:8767/bar") + assert not embed_proxy.upstream_allowed("http://127.0.0.1:9999/x") + assert not embed_proxy.upstream_allowed("http://example.com/") + + +def test_env_override_must_still_be_loopback(monkeypatch): + monkeypatch.setenv("SWITCHBAY_CE_UPSTREAM", "http://evil.example:8766") + with pytest.raises(embed_proxy.UpstreamNotLoopback): + embed_proxy.embed_targets() + + +@pytest.mark.asyncio +async def test_proxy_forwards_with_host_header(monkeypatch): + """End-to-end: embed proxy hits a loopback stub and injects X-*-Host.""" + seen: dict[str, str] = {} + + async def upstream_handler(request: web.Request) -> web.Response: + seen["path"] = request.path + seen["x_ce"] = request.headers.get("X-CE-Host", "") + seen["x_ok"] = request.headers.get("X-Okstratr-Host", "") + return web.json_response({"ok": True, "path": request.path}) + + up_app = web.Application() + up_app.router.add_get("/api/health", upstream_handler) + up_app.router.add_get("/observer/", upstream_handler) + up_server = TestServer(up_app) + await up_server.start_server() + try: + # Force loopback literal — TestServer may report 0.0.0.0. + base = f"http://127.0.0.1:{up_server.port}" + monkeypatch.setenv("SWITCHBAY_CE_UPSTREAM", base) + monkeypatch.setenv("SWITCHBAY_OKSTRATR_UPSTREAM", base) + + app = web.Application() + embed_proxy.register_routes(app) + async with TestClient(TestServer(app)) as client: + r = await client.get("/embed/ce/api/health") + assert r.status == 200 + body = await r.json() + assert body["ok"] is True + assert seen["path"] == "/api/health" + assert seen["x_ce"] == "switchbay" + + seen.clear() + r2 = await client.get("/embed/okstratr/observer/") + assert r2.status == 200 + assert seen["path"] == "/observer/" + assert seen["x_ok"] == "switchbay" + finally: + await up_server.close() + + +@pytest.mark.asyncio +async def test_proxy_rejects_if_upstream_env_not_loopback(monkeypatch): + monkeypatch.setenv("SWITCHBAY_CE_UPSTREAM", "http://203.0.113.1:8766") + app = web.Application() + embed_proxy.register_routes(app) + async with TestClient(TestServer(app)) as client: + r = await client.get("/embed/ce/api/health") + assert r.status == 502 + data = await r.json() + assert "loopback" in data["error"].lower() From eb5faa578989db4339652af8a46205dd6344cb23 Mon Sep 17 00:00:00 2001 From: benjsmith Date: Sat, 19 Sep 2026 08:37:24 +0200 Subject: [PATCH 02/25] fix(embed): wait for proxied flag before mounting Agents/Graph Avoid flashing the lazy AgentDashboardTab while /api/settings loads when proxied embeds are on (that race showed as module script import failures). Clarify 502 upstream messaging. --- frontend/src/center/builtinTabs.tsx | 6 ++++++ frontend/src/widgets/embed/ProxiedSkillPanel.tsx | 9 ++++++--- .../src/widgets/embed/useProxiedSkillEmbeds.ts | 16 ++++++++++------ 3 files changed, 22 insertions(+), 9 deletions(-) diff --git a/frontend/src/center/builtinTabs.tsx b/frontend/src/center/builtinTabs.tsx index 7452403..f2d9d9d 100644 --- a/frontend/src/center/builtinTabs.tsx +++ b/frontend/src/center/builtinTabs.tsx @@ -44,6 +44,9 @@ const OwidTab = lazy(() => import("../widgets/owid/OwidTab")); const GraphAdapter: TabComponent = ({ graphData, graphError }) => { const proxied = useProxiedSkillEmbeds(); + if (proxied === null) { + return

Loading…

; + } if (proxied) return ; return ; }; @@ -56,6 +59,9 @@ const SketchAdapter: TabComponent = () => ( ); const AgentsAdapter: TabComponent = () => { const proxied = useProxiedSkillEmbeds(); + if (proxied === null) { + return

Loading…

; + } if (proxied) return ; return ; }; diff --git a/frontend/src/widgets/embed/ProxiedSkillPanel.tsx b/frontend/src/widgets/embed/ProxiedSkillPanel.tsx index 36eda98..cf0d102 100644 --- a/frontend/src/widgets/embed/ProxiedSkillPanel.tsx +++ b/frontend/src/widgets/embed/ProxiedSkillPanel.tsx @@ -139,9 +139,12 @@ export default function ProxiedSkillPanel({ kind }: Props) {

{state.message}

- Start the skill daemon on loopback - {kind === "ce" ? " :8766" : " :8767"} (charter), or turn off - Settings → “Proxied skill embeds” to use the built-in tab. + Start the skill upstream on loopback + {kind === "ce" + ? " (CE viewer — default :8766, or whatever SWITCHBAY_CE_UPSTREAM points at)" + : " (okstratr :8767)"} + , or turn off Settings → Storage → “Proxied skill embeds” for the built-in tab. + A 502 almost always means that upstream process exited.

)} diff --git a/frontend/src/widgets/embed/useProxiedSkillEmbeds.ts b/frontend/src/widgets/embed/useProxiedSkillEmbeds.ts index 2d084ad..d59181c 100644 --- a/frontend/src/widgets/embed/useProxiedSkillEmbeds.ts +++ b/frontend/src/widgets/embed/useProxiedSkillEmbeds.ts @@ -2,23 +2,27 @@ import { useEffect, useState } from "react"; /** * Phase 4a feature flag: when true, Graph/Agents use `/embed/*` panels. - * Default false (built-in tabs). Reads `/api/settings` once; listens for - * `sy:proxied-skill-embeds` so Settings can flip without reload. + * Returns `null` until `/api/settings` resolves so adapters do not flash the + * built-in (lazy) tabs — that flash caused "Importing a module script failed" + * when Agents opened while proxied was actually on. */ -export function useProxiedSkillEmbeds(): boolean { - const [on, setOn] = useState(false); +export function useProxiedSkillEmbeds(): boolean | null { + const [on, setOn] = useState(null); useEffect(() => { let cancelled = false; void fetch("/api/settings") .then((r) => (r.ok ? r.json() : null)) .then((j) => { - if (!cancelled && j && typeof j.proxied_skill_embeds === "boolean") { + if (cancelled) return; + if (j && typeof j.proxied_skill_embeds === "boolean") { setOn(j.proxied_skill_embeds); + } else { + setOn(false); } }) .catch(() => { - /* older daemon — leave off */ + if (!cancelled) setOn(false); }); const onEvt = (ev: Event) => { From 679c067c7361505480622952579dd7588870b7db Mon Sep 17 00:00:00 2001 From: benjsmith Date: Sat, 19 Sep 2026 09:53:21 +0200 Subject: [PATCH 03/25] feat(embed): keep CE viewer alive when proxied embeds are on Start and health-check curiosity-engine viewer.sh serve for /embed/ce, restart on death, and rebuild the static bundle when wiki mtime changes. Toggle via Settings proxied_skill_embeds; supervisor runs as a daemon background task so Graph does not 502 when the viewer process exits. --- src/switchbay/ce_viewer_supervisor.py | 309 ++++++++++++++++++++++++++ src/switchbay/daemon.py | 57 ++++- 2 files changed, 364 insertions(+), 2 deletions(-) create mode 100644 src/switchbay/ce_viewer_supervisor.py diff --git a/src/switchbay/ce_viewer_supervisor.py b/src/switchbay/ce_viewer_supervisor.py new file mode 100644 index 0000000..59068c4 --- /dev/null +++ b/src/switchbay/ce_viewer_supervisor.py @@ -0,0 +1,309 @@ +"""Keep the CE HTML viewer alive for Switchbay proxied Graph embeds. + +When ``proxied_skill_embeds`` is on, Graph loads CE through ``/embed/ce/*``. +That requires a long-lived CE ``viewer.sh serve`` on loopback. This module +starts it, health-checks it, and restarts it if it dies. Optional wiki +mtime polling triggers ``viewer.sh build`` so the served bundle tracks +wiki edits (reload still needed in the Phase 4a HTML panel). +""" + +from __future__ import annotations + +import asyncio +import logging +import os +import signal +import subprocess +import time +from pathlib import Path +from typing import Any +from urllib.error import URLError +from urllib.parse import urlparse +from urllib.request import urlopen + +from . import app_settings, cebridge, embed_proxy + +log = logging.getLogger("switchbay.ce_viewer_supervisor") + +_PID_NAME = "ce-viewer.pid" +_LOG_NAME = "ce-viewer.log" +_HEALTH_PATH = "/" +_POLL_SEC = 8.0 +_WIKI_POLL_SEC = 15.0 + + +def _state_dir() -> Path: + d = Path.home() / ".local" / "state" / "switchbay" + d.mkdir(parents=True, exist_ok=True) + return d + + +def pid_path() -> Path: + return _state_dir() / _PID_NAME + + +def log_path() -> Path: + return _state_dir() / _LOG_NAME + + +def upstream_base() -> str: + return embed_proxy.ce_upstream() + + +def upstream_port() -> int: + u = urlparse(upstream_base()) + if u.port: + return int(u.port) + return 443 if u.scheme == "https" else 80 + + +def health_url() -> str: + return upstream_base().rstrip("/") + _HEALTH_PATH + + +def is_healthy(*, timeout: float = 2.0) -> bool: + try: + with urlopen(health_url(), timeout=timeout) as resp: # noqa: S310 — loopback only + return 200 <= int(getattr(resp, "status", 200)) < 500 + except (URLError, OSError, TimeoutError, ValueError): + return False + + +def _read_pid() -> int | None: + p = pid_path() + if not p.is_file(): + return None + try: + raw = p.read_text(encoding="utf-8").strip() + return int(raw) if raw else None + except (OSError, ValueError): + return None + + +def _pid_alive(pid: int) -> bool: + if pid <= 0: + return False + try: + os.kill(pid, 0) + return True + except OSError: + return False + + +def _public_base() -> str: + # Match embed proxy prefix so asset URLs work under /embed/ce/. + return (os.environ.get("CE_PUBLIC_BASE") or "/embed/ce").rstrip("/") or "/embed/ce" + + +def start(workspace: Path) -> dict[str, Any]: + """Start CE viewer.sh serve if not healthy. Idempotent.""" + workspace = Path(workspace).expanduser().resolve() + if is_healthy(): + return { + "ok": True, + "already_running": True, + "healthy": True, + "url": upstream_base(), + "port": upstream_port(), + "pid": _read_pid(), + } + + if not cebridge.has_wiki(workspace): + return { + "ok": False, + "error": f"no wiki/ under workspace {workspace}", + "hint": "Open a CE wiki workspace or run curiosity-engine setup.", + } + + script = cebridge.ce_root() / "scripts" / "viewer.sh" + if not script.is_file(): + return { + "ok": False, + "error": f"viewer.sh not found under {cebridge.ce_root()}", + "hint": "Install CE skill or set SWITCHBAY_CE_ROOT.", + } + + port = upstream_port() + # Drop stale pid if process is gone. + old = _read_pid() + if old and not _pid_alive(old): + try: + pid_path().unlink(missing_ok=True) + except OSError: + pass + + env = os.environ.copy() + env["CE_PUBLIC_BASE"] = _public_base() + # Prefer the same CE root Switchbay already resolved. + env.setdefault("SWITCHBAY_CE_ROOT", str(cebridge.ce_root())) + + logf = log_path().open("a", encoding="utf-8") + try: + proc = subprocess.Popen( # noqa: S603 + ["bash", str(script), "serve", str(port)], + cwd=str(workspace), + env=env, + stdout=logf, + stderr=subprocess.STDOUT, + start_new_session=True, + ) + except OSError as e: + logf.close() + return {"ok": False, "error": f"failed to spawn viewer: {e}"} + + try: + pid_path().write_text(str(proc.pid) + "\n", encoding="utf-8") + except OSError as e: + log.warning("could not write ce-viewer pid: %s", e) + + # Wait briefly for bind. + deadline = time.time() + 12.0 + while time.time() < deadline: + if is_healthy(): + return { + "ok": True, + "already_running": False, + "healthy": True, + "url": upstream_base(), + "port": port, + "pid": proc.pid, + "workspace": str(workspace), + "log": str(log_path()), + } + if proc.poll() is not None: + return { + "ok": False, + "error": f"viewer exited early code={proc.returncode}", + "log": str(log_path()), + } + time.sleep(0.4) + + return { + "ok": False, + "error": "viewer started but health check timed out", + "pid": proc.pid, + "log": str(log_path()), + "url": upstream_base(), + } + + +def stop() -> dict[str, Any]: + pid = _read_pid() + if not pid: + return {"ok": True, "stopped": False, "note": "no pid file"} + try: + os.kill(pid, signal.SIGTERM) + except OSError as e: + pid_path().unlink(missing_ok=True) + return {"ok": True, "stopped": False, "note": f"pid gone: {e}"} + deadline = time.time() + 5.0 + while time.time() < deadline and _pid_alive(pid): + time.sleep(0.2) + if _pid_alive(pid): + try: + os.kill(pid, signal.SIGKILL) + except OSError: + pass + pid_path().unlink(missing_ok=True) + return {"ok": True, "stopped": True, "pid": pid} + + +def status(workspace: Path | None = None) -> dict[str, Any]: + pid = _read_pid() + return { + "ok": True, + "healthy": is_healthy(), + "url": upstream_base(), + "port": upstream_port(), + "pid": pid, + "pid_alive": _pid_alive(pid) if pid else False, + "public_base": _public_base(), + "ce_root": str(cebridge.ce_root()), + "workspace": str(workspace) if workspace else None, + "proxied_skill_embeds": app_settings.get_proxied_skill_embeds(), + "log": str(log_path()), + } + + +def rebuild_bundle(workspace: Path) -> dict[str, Any]: + """Run viewer.sh build so the served bundle picks up wiki edits.""" + workspace = Path(workspace).expanduser().resolve() + if not cebridge.has_wiki(workspace): + return {"ok": False, "error": "no wiki/"} + script = cebridge.ce_root() / "scripts" / "viewer.sh" + if not script.is_file(): + return {"ok": False, "error": f"viewer.sh missing under {cebridge.ce_root()}"} + env = os.environ.copy() + env["CE_PUBLIC_BASE"] = _public_base() + try: + proc = subprocess.run( # noqa: S603 + ["bash", str(script), "build"], + cwd=str(workspace), + env=env, + capture_output=True, + text=True, + timeout=600, + check=False, + ) + except (OSError, subprocess.TimeoutExpired) as e: + return {"ok": False, "error": str(e)} + if proc.returncode != 0: + tail = (proc.stderr or proc.stdout or "")[-800:] + return {"ok": False, "error": f"build exit {proc.returncode}", "detail": tail} + return {"ok": True, "rebuilt": True} + + +def _wiki_mtime(workspace: Path) -> float: + wiki = Path(workspace) / "wiki" + if not wiki.is_dir(): + return 0.0 + newest = 0.0 + try: + for p in wiki.rglob("*"): + if p.is_file(): + try: + newest = max(newest, p.stat().st_mtime) + except OSError: + continue + except OSError: + return newest + return newest + + +async def run_supervisor(app: Any) -> None: + """Background task: keep CE up while proxied embeds are enabled.""" + last_wiki_mtime = 0.0 + last_rebuild = 0.0 + while True: + try: + if not app_settings.get_proxied_skill_embeds(): + await asyncio.sleep(_POLL_SEC) + continue + ws = Path(app["workspace"]) + if not is_healthy(): + log.info("CE viewer unhealthy — starting for %s", ws) + out = await asyncio.to_thread(start, ws) + if not out.get("ok"): + log.warning("CE viewer start failed: %s", out.get("error")) + else: + # Rebuild static bundle when wiki changes (Phase 4a panel + # still needs a browser reload; keep-alive is the hard part). + mtime = await asyncio.to_thread(_wiki_mtime, ws) + now = time.time() + if ( + mtime > last_wiki_mtime + and last_wiki_mtime > 0.0 + and (now - last_rebuild) > _WIKI_POLL_SEC + ): + last_rebuild = now + reb = await asyncio.to_thread(rebuild_bundle, ws) + log.info("CE viewer rebuild after wiki change: %s", reb) + if last_wiki_mtime == 0.0: + last_wiki_mtime = mtime + elif mtime > last_wiki_mtime: + last_wiki_mtime = mtime + except asyncio.CancelledError: + raise + except Exception: # noqa: BLE001 + log.exception("CE viewer supervisor loop error") + await asyncio.sleep(_POLL_SEC) diff --git a/src/switchbay/daemon.py b/src/switchbay/daemon.py index 8c6da4f..2c7fdbb 100644 --- a/src/switchbay/daemon.py +++ b/src/switchbay/daemon.py @@ -37,7 +37,7 @@ ce_tools, command_palettes, commands, conversations, curation_history, dbintrospect, demo_workspace, - duckdb_starters, embed_proxy, file_state, fileops, llm_config, llmgateway, + duckdb_starters, embed_proxy, ce_viewer_supervisor, file_state, fileops, llm_config, llmgateway, localllm, orchestrator_fs, schedules, mcpstore, merging, model_cache, modestore, owid, packstore, pasteboard, permissions, plots, html_decks, library, local_models, media_settings, micro_edits, projects, proposals, protocol, rail, report_html, report_packages, reports, secrets, selection, service, share, sheets, @@ -959,7 +959,25 @@ async def handle_settings_post(request: web.Request) -> web.Response: # rebuilds the index if the vector space changed. conversations.reset_embedder() if "proxied_skill_embeds" in body: - app_settings.set_proxied_skill_embeds(bool(body["proxied_skill_embeds"])) + want = bool(body["proxied_skill_embeds"]) + had = app_settings.get_proxied_skill_embeds() + app_settings.set_proxied_skill_embeds(want) + if want and not had: + # Kick CE keep-alive immediately (supervisor also polls). + async def _kick() -> None: + out = await asyncio.to_thread( + ce_viewer_supervisor.start, workspace, + ) + if out.get("ok"): + log.info("CE viewer started via settings: %s", out) + else: + log.warning("CE viewer start via settings failed: %s", out) + asyncio.create_task(_kick()) + elif had and not want: + async def _halt() -> None: + out = await asyncio.to_thread(ce_viewer_supervisor.stop) + log.info("CE viewer stopped via settings: %s", out) + asyncio.create_task(_halt()) # Media: { media: { image?: {provider, model}|null, video?: …, voice?: … } } if "media" in body and isinstance(body["media"], dict): if not admin_policy.feature_enabled("media_generation"): @@ -17492,6 +17510,41 @@ async def _go() -> None: _app["_ce_skill_task"] = asyncio.create_task(_go()) app.on_startup.append(_ensure_ce_skill) + async def _start_ce_viewer_supervisor(_app: web.Application) -> None: + # Keep CE HTML viewer alive for /embed/ce when proxied embeds are on. + # Never block bind — spawn the poll loop as a background task. + async def _go() -> None: + # First-chance start if the flag is already on. + if app_settings.get_proxied_skill_embeds(): + ws = Path(_app["workspace"]) + out = await asyncio.to_thread(ce_viewer_supervisor.start, ws) + if out.get("ok"): + log.info("CE viewer supervisor initial start: %s", out) + else: + log.warning( + "CE viewer supervisor initial start failed: %s", + out.get("error"), + ) + await ce_viewer_supervisor.run_supervisor(_app) + + async def _stop(_a: web.Application) -> None: + t = _a.get("_ce_viewer_supervisor_task") + if t: + t.cancel() + try: + await t + except (asyncio.CancelledError, Exception): + pass + # Leave the CE process running across daemon restarts only if + # proxied embeds stay on; otherwise tear it down with the daemon. + if not app_settings.get_proxied_skill_embeds(): + await asyncio.to_thread(ce_viewer_supervisor.stop) + + _app["_ce_viewer_supervisor_task"] = asyncio.create_task(_go()) + _app.on_cleanup.append(_stop) + + app.on_startup.append(_start_ce_viewer_supervisor) + # The launch workspace is set directly (no _activate), so relocate # its rail-history DB to match the current setting on startup too. async def _relocate_rail_history(_app: web.Application) -> None: From 1bd8d6f03fbed8b1c32914cab8a53fe54a86d795 Mon Sep 17 00:00:00 2001 From: benjsmith Date: Sat, 19 Sep 2026 09:56:09 +0200 Subject: [PATCH 04/25] feat(embed): soft-reload proxied Graph on wiki files_changed ProxiedSkillPanel listens for sy:files-changed (dispatched from the WS files_changed handler), debounces, and soft-refetches /embed/* with cache-busting so the panel updates without a full loading flash. CE supervisor rebuilds on wiki mtime and broadcasts files_changed; embed proxy forces Cache-Control: no-store so stale HTML shells do not stick. --- frontend/src/App.tsx | 2 + .../src/widgets/embed/ProxiedSkillPanel.tsx | 44 ++++++++++++++++--- src/switchbay/ce_viewer_supervisor.py | 33 +++++++++++--- src/switchbay/daemon.py | 3 ++ src/switchbay/embed_proxy.py | 5 +++ 5 files changed, 77 insertions(+), 10 deletions(-) diff --git a/frontend/src/App.tsx b/frontend/src/App.tsx index c837a54..454fe1e 100644 --- a/frontend/src/App.tsx +++ b/frontend/src/App.tsx @@ -998,6 +998,8 @@ export default function App() { }); } else if (msg.type === "files_changed") { setFilesVersion((v) => v + 1); + // Proxied Graph/Agents panels listen for this to soft-refetch /embed/*. + window.dispatchEvent(new CustomEvent("sy:files-changed")); } else if (msg.type === "artifact") { // Pend a pulse badge (Zen) AND switch to the surface the // agent just wrote — sheet/plot requests used to finish diff --git a/frontend/src/widgets/embed/ProxiedSkillPanel.tsx b/frontend/src/widgets/embed/ProxiedSkillPanel.tsx index cf0d102..040b92f 100644 --- a/frontend/src/widgets/embed/ProxiedSkillPanel.tsx +++ b/frontend/src/widgets/embed/ProxiedSkillPanel.tsx @@ -1,4 +1,4 @@ -import { useCallback, useEffect, useState } from "react"; +import { useCallback, useEffect, useRef, useState } from "react"; /** * Phase 4a same-origin skill panel (NO iframe). @@ -7,6 +7,8 @@ import { useCallback, useEffect, useState } from "react"; * `/embed/ce/*` or `/embed/okstratr/*`. In-app path navigation uses * fetch + a same-document panel — never a nested frame. * + * Reloads automatically when the daemon broadcasts `files_changed` + * (wiki edits, curator, rescan) so the proxied Graph feels live. * Full atlas/observer chrome lands as those skills grow hosted-mode * fragment/API UIs; this panel proves the proxy path and stays usable * when upstream is down. @@ -55,13 +57,25 @@ export default function ProxiedSkillPanel({ kind }: Props) { const [path, setPath] = useState(DEFAULT_PATH[kind]); const [draft, setDraft] = useState(DEFAULT_PATH[kind]); const [state, setState] = useState({ status: "loading" }); + const [refreshing, setRefreshing] = useState(false); + const pathRef = useRef(path); + pathRef.current = path; const load = useCallback( - async (p: string) => { - setState({ status: "loading" }); - const url = prefix + (p.startsWith("/") ? p : `/${p}`); + async (p: string, opts?: { soft?: boolean }) => { + const soft = !!opts?.soft; + if (soft) { + setRefreshing(true); + } else { + setState({ status: "loading" }); + } + // Cache-bust so CE static bundle / proxy do not serve a stale shell. + const base = prefix + (p.startsWith("/") ? p : `/${p}`); + const sep = base.includes("?") ? "&" : "?"; + const url = `${base}${sep}_sb=${Date.now()}`; try { const r = await fetch(url, { + cache: "no-store", headers: { Accept: "text/html, application/json;q=0.9, */*;q=0.8" }, }); const ct = r.headers.get("content-type") || ""; @@ -85,6 +99,8 @@ export default function ProxiedSkillPanel({ kind }: Props) { status: "error", message: (e as Error).message || "fetch failed", }); + } finally { + if (soft) setRefreshing(false); } }, [prefix], @@ -94,6 +110,23 @@ export default function ProxiedSkillPanel({ kind }: Props) { void load(path); }, [load, path]); + // Live view: wiki / curator / rescan → daemon files_changed → soft refetch. + useEffect(() => { + let timer: ReturnType | null = null; + const onFiles = () => { + if (timer) clearTimeout(timer); + // Debounce bursts (curator multi-write) into one refetch. + timer = setTimeout(() => { + void load(pathRef.current, { soft: true }); + }, 400); + }; + window.addEventListener("sy:files-changed", onFiles); + return () => { + if (timer) clearTimeout(timer); + window.removeEventListener("sy:files-changed", onFiles); + }; + }, [load]); + const onNavigate = (ev: React.FormEvent) => { ev.preventDefault(); const next = draft.trim() || "/"; @@ -106,6 +139,7 @@ export default function ProxiedSkillPanel({ kind }: Props) { {LABEL[kind]} same-origin via {prefix} (no iframe) + {refreshing ? " · updating…" : ""}
{ setDraft(path); - void load(path); + void load(path, { soft: true }); }} > Reload diff --git a/src/switchbay/ce_viewer_supervisor.py b/src/switchbay/ce_viewer_supervisor.py index 59068c4..e9a4c8f 100644 --- a/src/switchbay/ce_viewer_supervisor.py +++ b/src/switchbay/ce_viewer_supervisor.py @@ -28,8 +28,8 @@ _PID_NAME = "ce-viewer.pid" _LOG_NAME = "ce-viewer.log" _HEALTH_PATH = "/" -_POLL_SEC = 8.0 -_WIKI_POLL_SEC = 15.0 +_POLL_SEC = 5.0 +_WIKI_POLL_SEC = 3.0 def _state_dir() -> Path: @@ -270,6 +270,22 @@ def _wiki_mtime(workspace: Path) -> float: return newest + +async def _broadcast_files_changed(app: Any) -> None: + """Notify connected clients (proxied panel listens via App → CustomEvent).""" + from . import protocol + + # Prefer the daemon's own _broadcast if the app stored a bound helper. + bc = app.get("_broadcast_fn") + if callable(bc): + await bc(app, protocol.files_changed()) + return + # Fallback: import daemon._broadcast (circular-import safe at call time). + from . import daemon as _daemon + + await _daemon._broadcast(app, protocol.files_changed()) + + async def run_supervisor(app: Any) -> None: """Background task: keep CE up while proxied embeds are enabled.""" last_wiki_mtime = 0.0 @@ -296,11 +312,18 @@ async def run_supervisor(app: Any) -> None: and (now - last_rebuild) > _WIKI_POLL_SEC ): last_rebuild = now + last_wiki_mtime = mtime reb = await asyncio.to_thread(rebuild_bundle, ws) log.info("CE viewer rebuild after wiki change: %s", reb) - if last_wiki_mtime == 0.0: - last_wiki_mtime = mtime - elif mtime > last_wiki_mtime: + if reb.get("ok"): + # Tell proxied Graph panels to soft-refetch /embed/ce. + try: + from . import protocol + + await _broadcast_files_changed(app) + except Exception: # noqa: BLE001 + log.exception("files_changed after CE rebuild failed") + elif last_wiki_mtime == 0.0: last_wiki_mtime = mtime except asyncio.CancelledError: raise diff --git a/src/switchbay/daemon.py b/src/switchbay/daemon.py index 2c7fdbb..05adbef 100644 --- a/src/switchbay/daemon.py +++ b/src/switchbay/daemon.py @@ -17513,6 +17513,9 @@ async def _go() -> None: async def _start_ce_viewer_supervisor(_app: web.Application) -> None: # Keep CE HTML viewer alive for /embed/ce when proxied embeds are on. # Never block bind — spawn the poll loop as a background task. + # Expose broadcast to CE supervisor without import cycles at module load. + _app["_broadcast_fn"] = _broadcast + async def _go() -> None: # First-chance start if the flag is already on. if app_settings.get_proxied_skill_embeds(): diff --git a/src/switchbay/embed_proxy.py b/src/switchbay/embed_proxy.py index e46a983..011ef10 100644 --- a/src/switchbay/embed_proxy.py +++ b/src/switchbay/embed_proxy.py @@ -173,7 +173,12 @@ def filter_response_headers(headers: Mapping[str, str]) -> dict[str, str]: for k, v in headers.items(): if k.lower() in _HOP_BY_HOP: continue + # Drop upstream cache directives — proxied Graph soft-reloads on + # wiki changes and must not keep a stale HTML shell. + if k.lower() in {"cache-control", "etag", "last-modified", "expires"}: + continue out[k] = v + out["Cache-Control"] = "no-store" return out From bf39c066e80004d2828852a9698506dcc753b5d2 Mon Sep 17 00:00:00 2001 From: benjsmith Date: Sat, 19 Sep 2026 08:16:36 +0000 Subject: [PATCH 05/25] feat(core-skills): always auto-start CE+okstratr; rail host_notify MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Ben lock / contract C1+C2 (Switchbay half): - okstratr_supervisor mirrors CE: spawn `okstratr serve --host/--port` (PATH binary or python -m; SWITCHBAY_OKSTRATR_BIN / UPSTREAM overrides), health /health, restart on death; pid/log under ~/.local/state/switchbay/ - Daemon on_startup always starts CE + okstratr supervisors (proxied_skill_embeds remains UI panel choice only; no longer gates keep-alive or tears down CE) - GET /api/core-skills/status → {ce, okstratr, wiki_build} contract shape - POST /api/okstratr/host-notify maps okstratr.host_notify → protocol.notice rail broadcast (path-native; no OS/email channel) - Unit tests for URL/port/argv helpers, notify mapping+HTTP, status shape; embed_proxy tests still green --- src/switchbay/ce_viewer_supervisor.py | 76 +++++- src/switchbay/core_skills.py | 23 ++ src/switchbay/daemon.py | 110 +++++--- src/switchbay/host_notify.py | 118 +++++++++ src/switchbay/okstratr_supervisor.py | 354 +++++++++++++++++++++++++ tests/unit/test_core_skills_status.py | 39 +++ tests/unit/test_host_notify.py | 104 ++++++++ tests/unit/test_okstratr_supervisor.py | 42 +++ 8 files changed, 818 insertions(+), 48 deletions(-) create mode 100644 src/switchbay/core_skills.py create mode 100644 src/switchbay/host_notify.py create mode 100644 src/switchbay/okstratr_supervisor.py create mode 100644 tests/unit/test_core_skills_status.py create mode 100644 tests/unit/test_host_notify.py create mode 100644 tests/unit/test_okstratr_supervisor.py diff --git a/src/switchbay/ce_viewer_supervisor.py b/src/switchbay/ce_viewer_supervisor.py index e9a4c8f..de6c1d9 100644 --- a/src/switchbay/ce_viewer_supervisor.py +++ b/src/switchbay/ce_viewer_supervisor.py @@ -31,6 +31,57 @@ _POLL_SEC = 5.0 _WIKI_POLL_SEC = 3.0 +# Contract C1 wiki_build + CE state (module-level; supervisor updates). +_CE_STATE: dict[str, str] = {"state": "stopped", "detail": ""} +_WIKI_BUILD: dict[str, object] = { + "state": "idle", # idle|building|failed + "pages": None, + "detail": "", +} + + +def _set_ce_state(state: str, detail: str = "") -> None: + _CE_STATE["state"] = state + _CE_STATE["detail"] = detail + + +def _set_wiki_build(state: str, *, pages=None, detail: str = "") -> None: + _WIKI_BUILD["state"] = state + _WIKI_BUILD["pages"] = pages + _WIKI_BUILD["detail"] = detail + + +def contract_slice(workspace: Path | None = None) -> dict[str, Any]: + """C1 status slice for CE.""" + healthy = is_healthy() + pid = _read_pid() + if healthy: + state = "healthy" + detail = _CE_STATE.get("detail") or "up" + elif _CE_STATE.get("state") == "starting": + state = "starting" + detail = _CE_STATE.get("detail") or "starting" + elif pid and _pid_alive(pid): + state = "unhealthy" + detail = _CE_STATE.get("detail") or "process up but health failing" + elif _CE_STATE.get("state") == "unhealthy": + state = "unhealthy" + detail = _CE_STATE.get("detail") or "unhealthy" + else: + state = "stopped" + detail = _CE_STATE.get("detail") or "stopped" + return {"state": state, "url": upstream_base(), "detail": detail} + + +def wiki_build_slice() -> dict[str, Any]: + """C1 wiki_build slice.""" + return { + "state": str(_WIKI_BUILD.get("state") or "idle"), + "pages": _WIKI_BUILD.get("pages"), + "detail": str(_WIKI_BUILD.get("detail") or ""), + } + + def _state_dir() -> Path: d = Path.home() / ".local" / "state" / "switchbay" @@ -99,6 +150,7 @@ def start(workspace: Path) -> dict[str, Any]: """Start CE viewer.sh serve if not healthy. Idempotent.""" workspace = Path(workspace).expanduser().resolve() if is_healthy(): + _set_ce_state("healthy", "already running") return { "ok": True, "already_running": True, @@ -132,6 +184,7 @@ def start(workspace: Path) -> dict[str, Any]: except OSError: pass + _set_ce_state("starting", "spawning viewer.sh serve") env = os.environ.copy() env["CE_PUBLIC_BASE"] = _public_base() # Prefer the same CE root Switchbay already resolved. @@ -229,12 +282,15 @@ def rebuild_bundle(workspace: Path) -> dict[str, Any]: """Run viewer.sh build so the served bundle picks up wiki edits.""" workspace = Path(workspace).expanduser().resolve() if not cebridge.has_wiki(workspace): + _set_wiki_build("failed", detail="no wiki/") return {"ok": False, "error": "no wiki/"} script = cebridge.ce_root() / "scripts" / "viewer.sh" if not script.is_file(): + _set_wiki_build("failed", detail=f"viewer.sh missing under {cebridge.ce_root()}") return {"ok": False, "error": f"viewer.sh missing under {cebridge.ce_root()}"} env = os.environ.copy() env["CE_PUBLIC_BASE"] = _public_base() + _set_wiki_build("building", detail="viewer.sh build") try: proc = subprocess.run( # noqa: S603 ["bash", str(script), "build"], @@ -246,10 +302,13 @@ def rebuild_bundle(workspace: Path) -> dict[str, Any]: check=False, ) except (OSError, subprocess.TimeoutExpired) as e: + _set_wiki_build("failed", detail=str(e)) return {"ok": False, "error": str(e)} if proc.returncode != 0: tail = (proc.stderr or proc.stdout or "")[-800:] + _set_wiki_build("failed", detail=f"build exit {proc.returncode}") return {"ok": False, "error": f"build exit {proc.returncode}", "detail": tail} + _set_wiki_build("idle", detail="rebuilt") return {"ok": True, "rebuilt": True} @@ -287,21 +346,27 @@ async def _broadcast_files_changed(app: Any) -> None: async def run_supervisor(app: Any) -> None: - """Background task: keep CE up while proxied embeds are enabled.""" + """Background task: always keep CE up (core experience; Ben lock). + + Proxied-embeds flag still chooses the Graph UI panel; supervisors bring + the skill up regardless so status is never an empty Graph. + """ last_wiki_mtime = 0.0 last_rebuild = 0.0 while True: try: - if not app_settings.get_proxied_skill_embeds(): - await asyncio.sleep(_POLL_SEC) - continue ws = Path(app["workspace"]) if not is_healthy(): log.info("CE viewer unhealthy — starting for %s", ws) + _set_ce_state("starting", "supervisor restart") out = await asyncio.to_thread(start, ws) if not out.get("ok"): + _set_ce_state("unhealthy", str(out.get("error") or "start failed")) log.warning("CE viewer start failed: %s", out.get("error")) + else: + _set_ce_state("healthy", "supervisor") else: + _set_ce_state("healthy", "supervisor") # Rebuild static bundle when wiki changes (Phase 4a panel # still needs a browser reload; keep-alive is the hard part). mtime = await asyncio.to_thread(_wiki_mtime, ws) @@ -318,8 +383,6 @@ async def run_supervisor(app: Any) -> None: if reb.get("ok"): # Tell proxied Graph panels to soft-refetch /embed/ce. try: - from . import protocol - await _broadcast_files_changed(app) except Exception: # noqa: BLE001 log.exception("files_changed after CE rebuild failed") @@ -329,4 +392,5 @@ async def run_supervisor(app: Any) -> None: raise except Exception: # noqa: BLE001 log.exception("CE viewer supervisor loop error") + _set_ce_state("unhealthy", "supervisor loop error") await asyncio.sleep(_POLL_SEC) diff --git a/src/switchbay/core_skills.py b/src/switchbay/core_skills.py new file mode 100644 index 0000000..68f17ec --- /dev/null +++ b/src/switchbay/core_skills.py @@ -0,0 +1,23 @@ +"""Combined core-skills status (CE + okstratr + wiki_build). + +Contract C1 shared status shape for UI / Herdr / CLI. +""" + +from __future__ import annotations + +from pathlib import Path +from typing import Any + +from . import ce_viewer_supervisor, okstratr_supervisor + + +def status(workspace: Path | None = None) -> dict[str, Any]: + """Return the C1 status object.""" + ce = ce_viewer_supervisor.contract_slice(workspace) + oks = okstratr_supervisor.contract_slice() + wiki = ce_viewer_supervisor.wiki_build_slice() + return { + "ce": ce, + "okstratr": oks, + "wiki_build": wiki, + } diff --git a/src/switchbay/daemon.py b/src/switchbay/daemon.py index 05adbef..fb25532 100644 --- a/src/switchbay/daemon.py +++ b/src/switchbay/daemon.py @@ -37,7 +37,7 @@ ce_tools, command_palettes, commands, conversations, curation_history, dbintrospect, demo_workspace, - duckdb_starters, embed_proxy, ce_viewer_supervisor, file_state, fileops, llm_config, llmgateway, + duckdb_starters, embed_proxy, ce_viewer_supervisor, okstratr_supervisor, core_skills, host_notify, file_state, fileops, llm_config, llmgateway, localllm, orchestrator_fs, schedules, mcpstore, merging, model_cache, modestore, owid, packstore, pasteboard, permissions, plots, html_decks, library, local_models, media_settings, micro_edits, projects, proposals, protocol, rail, report_html, report_packages, reports, secrets, selection, service, share, sheets, @@ -962,22 +962,18 @@ async def handle_settings_post(request: web.Request) -> web.Response: want = bool(body["proxied_skill_embeds"]) had = app_settings.get_proxied_skill_embeds() app_settings.set_proxied_skill_embeds(want) + # Proxied flag only switches Graph/Agents UI panels. Core skills + # (CE + okstratr) stay supervised either way (Ben lock / C1). if want and not had: - # Kick CE keep-alive immediately (supervisor also polls). async def _kick() -> None: - out = await asyncio.to_thread( + ce_out = await asyncio.to_thread( ce_viewer_supervisor.start, workspace, ) - if out.get("ok"): - log.info("CE viewer started via settings: %s", out) - else: - log.warning("CE viewer start via settings failed: %s", out) + oks_out = await asyncio.to_thread( + okstratr_supervisor.start, workspace, + ) + log.info("core skills kick via settings: ce=%s oks=%s", ce_out, oks_out) asyncio.create_task(_kick()) - elif had and not want: - async def _halt() -> None: - out = await asyncio.to_thread(ce_viewer_supervisor.stop) - log.info("CE viewer stopped via settings: %s", out) - asyncio.create_task(_halt()) # Media: { media: { image?: {provider, model}|null, video?: …, voice?: … } } if "media" in body and isinstance(body["media"], dict): if not admin_policy.feature_enabled("media_generation"): @@ -1013,6 +1009,24 @@ async def _halt() -> None: return await handle_settings_get(request) + +async def handle_core_skills_status(request: web.Request) -> web.Response: + """C1 shared status: CE + okstratr + wiki_build.""" + workspace: Path = request.app["workspace"] + return web.json_response(core_skills.status(workspace)) + + +async def handle_okstratr_host_notify(request: web.Request) -> web.Response: + """C2 path-native notify: okstratr → Switchbay rail (protocol.notice).""" + try: + body = await request.json() + except json.JSONDecodeError: + return web.json_response({"ok": False, "error": "invalid json"}, status=400) + result = await host_notify.apply_host_notify(request.app, body) + status = 200 if result.get("ok") else 400 + return web.json_response(result, status=status) + + async def handle_orchestration_policy_get(request: web.Request) -> web.Response: """Inspectable Auto-orchestration policy + hard bounds. No secrets. @@ -17395,6 +17409,8 @@ def build_app(workspace: Path) -> web.Application: app.router.add_post("/api/llm/refresh_models", handle_llm_refresh_models) app.router.add_get("/api/rail/events", handle_rail_events) app.router.add_get("/api/settings", handle_settings_get) + app.router.add_get("/api/core-skills/status", handle_core_skills_status) + app.router.add_post("/api/okstratr/host-notify", handle_okstratr_host_notify) app.router.add_get("/api/curator-profile", handle_curator_profile_get) app.router.add_post("/api/curator-profile", handle_curator_profile_post) app.router.add_post("/api/curator-profile/draft", handle_curator_profile_draft) @@ -17510,43 +17526,53 @@ async def _go() -> None: _app["_ce_skill_task"] = asyncio.create_task(_go()) app.on_startup.append(_ensure_ce_skill) - async def _start_ce_viewer_supervisor(_app: web.Application) -> None: - # Keep CE HTML viewer alive for /embed/ce when proxied embeds are on. - # Never block bind — spawn the poll loop as a background task. - # Expose broadcast to CE supervisor without import cycles at module load. + async def _start_core_skill_supervisors(_app: web.Application) -> None: + # C1: always auto-start CE + okstratr with the shell (Ben lock). + # Proxied embeds only choose which UI panels load /embed/* — supervisors + # still bring skills up. Never block bind — background tasks only. _app["_broadcast_fn"] = _broadcast - async def _go() -> None: - # First-chance start if the flag is already on. - if app_settings.get_proxied_skill_embeds(): - ws = Path(_app["workspace"]) - out = await asyncio.to_thread(ce_viewer_supervisor.start, ws) - if out.get("ok"): - log.info("CE viewer supervisor initial start: %s", out) - else: - log.warning( - "CE viewer supervisor initial start failed: %s", - out.get("error"), - ) + async def _go_ce() -> None: + ws = Path(_app["workspace"]) + out = await asyncio.to_thread(ce_viewer_supervisor.start, ws) + if out.get("ok"): + log.info("CE viewer supervisor initial start: %s", out) + else: + log.warning( + "CE viewer supervisor initial start failed: %s", + out.get("error"), + ) await ce_viewer_supervisor.run_supervisor(_app) - async def _stop(_a: web.Application) -> None: - t = _a.get("_ce_viewer_supervisor_task") - if t: - t.cancel() - try: - await t - except (asyncio.CancelledError, Exception): - pass - # Leave the CE process running across daemon restarts only if - # proxied embeds stay on; otherwise tear it down with the daemon. - if not app_settings.get_proxied_skill_embeds(): - await asyncio.to_thread(ce_viewer_supervisor.stop) + async def _go_oks() -> None: + ws = Path(_app["workspace"]) + out = await asyncio.to_thread(okstratr_supervisor.start, ws) + if out.get("ok"): + log.info("okstratr supervisor initial start: %s", out) + else: + log.warning( + "okstratr supervisor initial start failed: %s", + out.get("error"), + ) + await okstratr_supervisor.run_supervisor(_app) - _app["_ce_viewer_supervisor_task"] = asyncio.create_task(_go()) + async def _stop(_a: web.Application) -> None: + for key in ("_ce_viewer_supervisor_task", "_okstratr_supervisor_task"): + t = _a.get(key) + if t: + t.cancel() + try: + await t + except (asyncio.CancelledError, Exception): + pass + # Leave skill processes up across daemon restarts (supervisors + # re-attach via health check). Matches prior CE-when-proxied-on. + + _app["_ce_viewer_supervisor_task"] = asyncio.create_task(_go_ce()) + _app["_okstratr_supervisor_task"] = asyncio.create_task(_go_oks()) _app.on_cleanup.append(_stop) - app.on_startup.append(_start_ce_viewer_supervisor) + app.on_startup.append(_start_core_skill_supervisors) # The launch workspace is set directly (no _activate), so relocate # its rail-history DB to match the current setting on startup too. diff --git a/src/switchbay/host_notify.py b/src/switchbay/host_notify.py new file mode 100644 index 0000000..a3a5238 --- /dev/null +++ b/src/switchbay/host_notify.py @@ -0,0 +1,118 @@ +"""Map okstratr ``host_notify`` envelopes onto Switchbay rail notices. + +Contract C2: Switchbay sink is the **rail** only (path-native). okstratr +emits ``okstratr.host_notify``; this host maps to ``protocol.notice`` and +broadcasts on the agent chat stream. +""" + +from __future__ import annotations + +from typing import Any + +from . import protocol + +ENVELOPE_TYPE = "okstratr.host_notify" +ENVELOPE_V = 1 + +_KIND_LABELS = { + "schedule.start": "Schedule started", + "schedule.progress": "Schedule progress", + "schedule.done": "Schedule done", + "schedule.failed": "Schedule failed", + "desk.progress": "Desk progress", + "desk.done": "Desk done", +} + + +def validate_envelope(body: Any) -> tuple[dict[str, Any] | None, str | None]: + """Return (envelope, error). error is set when invalid.""" + if not isinstance(body, dict): + return None, "envelope must be a JSON object" + if body.get("type") != ENVELOPE_TYPE: + return None, f"type must be {ENVELOPE_TYPE!r}" + v = body.get("v", ENVELOPE_V) + try: + if int(v) != ENVELOPE_V: + return None, f"unsupported envelope v={v!r}" + except (TypeError, ValueError): + return None, f"unsupported envelope v={v!r}" + kind = body.get("kind") + if not isinstance(kind, str) or not kind.strip(): + return None, "kind is required" + return body, None + + +def format_rail_text(envelope: dict[str, Any]) -> str: + """Human rail line from a validated envelope.""" + kind = str(envelope.get("kind") or "") + label = _KIND_LABELS.get(kind, kind) + title = (envelope.get("title") or "").strip() + body = (envelope.get("body") or "").strip() + desk = envelope.get("desk") + schedule_id = envelope.get("schedule_id") + progress = envelope.get("progress") if isinstance(envelope.get("progress"), dict) else None + + parts: list[str] = [f"[okstratr] {label}"] + if title: + parts[0] = f"[okstratr] {label}: {title}" + meta: list[str] = [] + if desk not in (None, "", "null"): + meta.append(f"desk={desk}") + if schedule_id: + meta.append(f"id={schedule_id}") + if progress: + pct = progress.get("pct") + phase = progress.get("phase") + detail = progress.get("detail") + bits: list[str] = [] + if pct is not None: + bits.append(f"{pct}%") + if phase: + bits.append(str(phase)) + if detail: + bits.append(str(detail)) + if bits: + meta.append(" · ".join(bits)) + lines = [parts[0]] + if meta: + lines.append("(" + ", ".join(meta) + ")") + if body: + lines.append(body) + return "\n".join(lines) + + +def notice_from_envelope(envelope: dict[str, Any]) -> dict[str, Any]: + """Build a ``protocol.notice`` payload for rail broadcast.""" + return protocol.notice(format_rail_text(envelope), kind="okstratr") + + +def _notice_text(msg: dict[str, Any]) -> str | None: + """Extract human text from a protocol.notice CUSTOM envelope.""" + if "text" in msg and isinstance(msg.get("text"), str): + return msg["text"] + val = msg.get("value") + if isinstance(val, dict) and isinstance(val.get("text"), str): + return val["text"] + return None + + +async def apply_host_notify(app: Any, envelope: dict[str, Any]) -> dict[str, Any]: + """Validate + broadcast to rail. Returns result dict for HTTP handlers.""" + env, err = validate_envelope(envelope) + if err or env is None: + return {"ok": False, "error": err or "invalid envelope"} + msg = notice_from_envelope(env) + # Daemon stores the bound helper as ``_broadcast_fn`` (same as CE). + bc = app.get("_broadcast_fn") if hasattr(app, "get") else None + if callable(bc): + await bc(app, msg) + else: + from . import daemon as _daemon + + await _daemon._broadcast(app, msg) + return { + "ok": True, + "broadcast": True, + "text": _notice_text(msg) if isinstance(msg, dict) else None, + "notice": msg, + } diff --git a/src/switchbay/okstratr_supervisor.py b/src/switchbay/okstratr_supervisor.py new file mode 100644 index 0000000..713ad3c --- /dev/null +++ b/src/switchbay/okstratr_supervisor.py @@ -0,0 +1,354 @@ +"""Keep okstratr desk kernel alive for Switchbay (core skill). + +Mirrors ``ce_viewer_supervisor``: start on loopback (default +``http://127.0.0.1:8767``, override ``SWITCHBAY_OKSTRATR_UPSTREAM``), +health-check ``/health``, restart on death. Pid/log under +``~/.local/state/switchbay/``. + +Spawn: ``okstratr serve --host/--port`` on PATH (or +``SWITCHBAY_OKSTRATR_BIN`` / ``python -m okstratr serve``). Uses the +long-lived ``serve`` entry (same as okstratr ``lifecycle._spawn_serve``), +not ``start`` which exits after detaching. +""" + +from __future__ import annotations + +import asyncio +import logging +import os +import shutil +import signal +import subprocess +import sys +import time +from pathlib import Path +from typing import Any +from urllib.error import URLError +from urllib.parse import urlparse +from urllib.request import urlopen + +from . import embed_proxy + +log = logging.getLogger("switchbay.okstratr_supervisor") + +_PID_NAME = "okstratr.pid" +_LOG_NAME = "okstratr.log" +_HEALTH_PATH = "/health" +_POLL_SEC = 5.0 + +# Contract C1 state machine (shared with core_skills.status). +_STATE: dict[str, Any] = { + "state": "stopped", # starting|healthy|unhealthy|stopped + "detail": "", +} + + +def _state_dir() -> Path: + d = Path.home() / ".local" / "state" / "switchbay" + d.mkdir(parents=True, exist_ok=True) + return d + + +def pid_path() -> Path: + return _state_dir() / _PID_NAME + + +def log_path() -> Path: + return _state_dir() / _LOG_NAME + + +def upstream_base() -> str: + return embed_proxy.okstratr_upstream() + + +def upstream_port() -> int: + u = urlparse(upstream_base()) + if u.port: + return int(u.port) + return 443 if u.scheme == "https" else 80 + + +def upstream_host() -> str: + u = urlparse(upstream_base()) + return u.hostname or "127.0.0.1" + + +def health_url() -> str: + return upstream_base().rstrip("/") + _HEALTH_PATH + + +def is_healthy(*, timeout: float = 2.0) -> bool: + try: + with urlopen(health_url(), timeout=timeout) as resp: # noqa: S310 — loopback + if not (200 <= int(getattr(resp, "status", 200)) < 500): + return False + # okstratr /health returns {"ok": true}; tolerate plain 200 too. + raw = resp.read().decode("utf-8", errors="replace") + if not raw.strip(): + return True + import json + + try: + data = json.loads(raw) + except json.JSONDecodeError: + return True + if isinstance(data, dict) and "ok" in data: + return bool(data.get("ok")) + return True + except (URLError, OSError, TimeoutError, ValueError): + return False + + +def _read_pid() -> int | None: + p = pid_path() + if not p.is_file(): + return None + try: + raw = p.read_text(encoding="utf-8").strip() + return int(raw) if raw else None + except (OSError, ValueError): + return None + + +def _pid_alive(pid: int) -> bool: + if pid <= 0: + return False + try: + os.kill(pid, 0) + return True + except OSError: + return False + + +def _public_base() -> str: + return ( + os.environ.get("OKSTRATR_PUBLIC_BASE") or "/embed/okstratr" + ).rstrip("/") or "/embed/okstratr" + + +def resolve_okstratr_argv() -> list[str]: + """Return argv prefix to invoke okstratr CLI (no subcommand yet).""" + override = (os.environ.get("SWITCHBAY_OKSTRATR_BIN") or "").strip() + if override: + return [override] + exe = shutil.which("okstratr") + if exe: + return [exe] + return [sys.executable or "python3", "-m", "okstratr"] + + +def _set_state(state: str, detail: str = "") -> None: + _STATE["state"] = state + _STATE["detail"] = detail + + +def contract_slice() -> dict[str, Any]: + """C1 status slice for okstratr.""" + healthy = is_healthy() + pid = _read_pid() + if healthy: + state = "healthy" + detail = _STATE.get("detail") or "up" + elif _STATE.get("state") == "starting": + state = "starting" + detail = _STATE.get("detail") or "starting" + elif pid and _pid_alive(pid): + state = "unhealthy" + detail = _STATE.get("detail") or "process up but /health failing" + elif _STATE.get("state") == "unhealthy": + state = "unhealthy" + detail = _STATE.get("detail") or "unhealthy" + else: + state = "stopped" + detail = _STATE.get("detail") or "stopped" + return { + "state": state, + "url": upstream_base(), + "detail": detail, + } + + +def start(workspace: Path | None = None) -> dict[str, Any]: + """Start okstratr serve if not healthy. Idempotent.""" + if is_healthy(): + _set_state("healthy", "already running") + return { + "ok": True, + "already_running": True, + "healthy": True, + "url": upstream_base(), + "port": upstream_port(), + "pid": _read_pid(), + } + + _set_state("starting", "spawning okstratr serve") + old = _read_pid() + if old and not _pid_alive(old): + try: + pid_path().unlink(missing_ok=True) + except OSError: + pass + + host = upstream_host() + port = upstream_port() + # Use `serve` (long-lived HTTP), not `start` (launcher that exits after + # detaching). Same command lifecycle._spawn_serve uses. + argv = resolve_okstratr_argv() + [ + "serve", + "--host", + host, + "--port", + str(port), + ] + + env = os.environ.copy() + env["OKSTRATR_PUBLIC_BASE"] = _public_base() + # Prefer embed public base so hosted observer URLs match /embed/okstratr. + if workspace is not None: + env.setdefault("OKSTRATR_WORKSPACE", str(Path(workspace).expanduser())) + + logf = log_path().open("a", encoding="utf-8") + try: + proc = subprocess.Popen( # noqa: S603 + argv, + cwd=str(Path(workspace).expanduser()) if workspace else None, + env=env, + stdout=logf, + stderr=subprocess.STDOUT, + start_new_session=True, + ) + except OSError as e: + logf.close() + _set_state("unhealthy", f"failed to spawn: {e}") + return {"ok": False, "error": f"failed to spawn okstratr: {e}", "argv": argv} + + try: + pid_path().write_text(str(proc.pid) + "\n", encoding="utf-8") + except OSError as e: + log.warning("could not write okstratr pid: %s", e) + + deadline = time.time() + 15.0 + while time.time() < deadline: + if is_healthy(): + _set_state("healthy", "started") + return { + "ok": True, + "already_running": False, + "healthy": True, + "url": upstream_base(), + "port": port, + "pid": proc.pid, + "argv": argv, + "log": str(log_path()), + } + if proc.poll() is not None: + _set_state("unhealthy", f"serve exited code={proc.returncode}") + return { + "ok": False, + "error": f"okstratr serve exited early code={proc.returncode}", + "log": str(log_path()), + "argv": argv, + } + time.sleep(0.4) + + if is_healthy(): + _set_state("healthy", "started (late health)") + return { + "ok": True, + "already_running": False, + "healthy": True, + "url": upstream_base(), + "port": port, + "pid": proc.pid, + "argv": argv, + "log": str(log_path()), + } + + _set_state("unhealthy", "health check timed out") + return { + "ok": False, + "error": "okstratr started but health check timed out", + "pid": proc.pid, + "log": str(log_path()), + "url": upstream_base(), + "argv": argv, + } + + +def stop() -> dict[str, Any]: + """Stop the serve process we spawned (best-effort).""" + pid = _read_pid() + if not pid: + _set_state("stopped", "no pid file") + return {"ok": True, "stopped": False, "note": "no pid file"} + if not _pid_alive(pid): + try: + pid_path().unlink(missing_ok=True) + except OSError: + pass + _set_state("stopped", "pid gone") + return {"ok": True, "stopped": False, "note": "pid gone", "pid": pid} + try: + os.kill(pid, signal.SIGTERM) + except OSError as e: + try: + pid_path().unlink(missing_ok=True) + except OSError: + pass + _set_state("stopped", f"pid gone: {e}") + return {"ok": True, "stopped": False, "note": f"pid gone: {e}", "pid": pid} + deadline = time.time() + 5.0 + while time.time() < deadline and _pid_alive(pid): + time.sleep(0.2) + if _pid_alive(pid): + try: + os.kill(pid, signal.SIGKILL) + except OSError: + pass + try: + pid_path().unlink(missing_ok=True) + except OSError: + pass + _set_state("stopped", "stopped") + return {"ok": True, "stopped": True, "pid": pid} + + +def status(workspace: Path | None = None) -> dict[str, Any]: + pid = _read_pid() + return { + "ok": True, + "healthy": is_healthy(), + "url": upstream_base(), + "port": upstream_port(), + "pid": pid, + "pid_alive": _pid_alive(pid) if pid else False, + "public_base": _public_base(), + "argv_prefix": resolve_okstratr_argv(), + "workspace": str(workspace) if workspace else None, + "log": str(log_path()), + "contract": contract_slice(), + } + + +async def run_supervisor(app: Any) -> None: + """Background task: always keep okstratr up (core experience).""" + while True: + try: + if not is_healthy(): + ws = Path(app["workspace"]) + log.info("okstratr unhealthy — starting for %s", ws) + _set_state("starting", "supervisor restart") + out = await asyncio.to_thread(start, ws) + if not out.get("ok"): + _set_state("unhealthy", str(out.get("error") or "start failed")) + log.warning("okstratr start failed: %s", out.get("error")) + else: + _set_state("healthy", "supervisor") + else: + _set_state("healthy", "supervisor") + except asyncio.CancelledError: + raise + except Exception: # noqa: BLE001 + log.exception("okstratr supervisor loop error") + _set_state("unhealthy", "supervisor loop error") + await asyncio.sleep(_POLL_SEC) diff --git a/tests/unit/test_core_skills_status.py b/tests/unit/test_core_skills_status.py new file mode 100644 index 0000000..cc5ba84 --- /dev/null +++ b/tests/unit/test_core_skills_status.py @@ -0,0 +1,39 @@ +"""C1 combined core-skills status shape.""" + +from __future__ import annotations + +from switchbay import core_skills + + +def test_status_shape(monkeypatch): + from switchbay import ce_viewer_supervisor as ce + from switchbay import okstratr_supervisor as oks + + monkeypatch.setattr( + ce, + "contract_slice", + lambda _ws=None: { + "state": "stopped", + "url": "http://127.0.0.1:8766", + "detail": "x", + }, + ) + monkeypatch.setattr( + oks, + "contract_slice", + lambda: { + "state": "healthy", + "url": "http://127.0.0.1:8767", + "detail": "up", + }, + ) + monkeypatch.setattr( + ce, + "wiki_build_slice", + lambda: {"state": "idle", "pages": None, "detail": ""}, + ) + out = core_skills.status() + assert set(out) == {"ce", "okstratr", "wiki_build"} + assert out["ce"]["state"] == "stopped" + assert out["okstratr"]["state"] == "healthy" + assert out["wiki_build"]["state"] == "idle" diff --git a/tests/unit/test_host_notify.py b/tests/unit/test_host_notify.py new file mode 100644 index 0000000..46f2faf --- /dev/null +++ b/tests/unit/test_host_notify.py @@ -0,0 +1,104 @@ +"""okstratr.host_notify → rail protocol.notice mapping.""" + +from __future__ import annotations + +import pytest +from aiohttp import web +from aiohttp.test_utils import TestClient, TestServer + +from switchbay import host_notify + + +def _env(**overrides): + base = { + "type": "okstratr.host_notify", + "v": 1, + "kind": "schedule.start", + "schedule_id": "sched-1", + "desk": "auto", + "title": "Morning digest", + "body": "Running hedge-fund desk.", + "progress": {"pct": None, "phase": "start", "detail": ""}, + "ts": "2026-09-19T08:00:00Z", + } + base.update(overrides) + return base + + +def test_validate_ok(): + env, err = host_notify.validate_envelope(_env()) + assert err is None + assert env is not None + + +def test_validate_rejects_wrong_type(): + env, err = host_notify.validate_envelope(_env(type="other")) + assert env is None + assert "type" in (err or "") + + +def test_format_includes_title_and_meta(): + text = host_notify.format_rail_text(_env()) + assert "[okstratr]" in text + assert "Morning digest" in text + assert "desk=auto" in text + assert "id=sched-1" in text + assert "Running hedge-fund desk." in text + + +def test_format_progress(): + text = host_notify.format_rail_text( + _env( + kind="schedule.progress", + progress={"pct": 40, "phase": "fetch", "detail": "n=3"}, + ) + ) + assert "40%" in text + assert "fetch" in text + + +def test_notice_from_envelope_is_protocol_notice(): + msg = host_notify.notice_from_envelope(_env()) + text = host_notify._notice_text(msg) + assert text and "Morning digest" in text + + +@pytest.mark.asyncio +async def test_apply_broadcasts_via_app_helper(): + seen: list[dict] = [] + + async def _bc(app, msg): + seen.append(msg) + + app = {"_broadcast_fn": _bc} + result = await host_notify.apply_host_notify(app, _env()) + assert result["ok"] is True + assert result["broadcast"] is True + assert len(seen) == 1 + assert "Morning digest" in (result.get("text") or "") + + +@pytest.mark.asyncio +async def test_http_host_notify_endpoint_broadcasts(): + """POST /api/okstratr/host-notify posts envelope → rail notice.""" + from switchbay import daemon + + seen: list[dict] = [] + + async def _bc(app, msg): + seen.append(msg) + + app = web.Application() + app["workspace"] = "/tmp" + app["ws_clients"] = set() + app["_broadcast_fn"] = _bc + app.router.add_post("/api/okstratr/host-notify", daemon.handle_okstratr_host_notify) + + async with TestClient(TestServer(app)) as client: + resp = await client.post("/api/okstratr/host-notify", json=_env()) + assert resp.status == 200 + data = await resp.json() + assert data["ok"] is True + assert data["broadcast"] is True + assert "Morning digest" in (data.get("text") or "") + assert len(seen) == 1 diff --git a/tests/unit/test_okstratr_supervisor.py b/tests/unit/test_okstratr_supervisor.py new file mode 100644 index 0000000..0814ba5 --- /dev/null +++ b/tests/unit/test_okstratr_supervisor.py @@ -0,0 +1,42 @@ +"""URL/port helpers + argv resolution for okstratr supervisor.""" + +from __future__ import annotations + +from switchbay import okstratr_supervisor as oks + + +def test_upstream_defaults(monkeypatch): + monkeypatch.delenv("SWITCHBAY_OKSTRATR_UPSTREAM", raising=False) + assert oks.upstream_base() == "http://127.0.0.1:8767" + assert oks.upstream_port() == 8767 + assert oks.upstream_host() == "127.0.0.1" + assert oks.health_url().endswith("/health") + + +def test_upstream_env_override(monkeypatch): + monkeypatch.setenv("SWITCHBAY_OKSTRATR_UPSTREAM", "http://127.0.0.1:9876") + assert oks.upstream_base() == "http://127.0.0.1:9876" + assert oks.upstream_port() == 9876 + + +def test_resolve_argv_prefers_bin_override(monkeypatch, tmp_path): + fake = tmp_path / "okstratr" + fake.write_text("#!/bin/sh\n", encoding="utf-8") + fake.chmod(0o755) + monkeypatch.setenv("SWITCHBAY_OKSTRATR_BIN", str(fake)) + assert oks.resolve_okstratr_argv() == [str(fake)] + + +def test_resolve_argv_falls_back_to_module(monkeypatch): + monkeypatch.delenv("SWITCHBAY_OKSTRATR_BIN", raising=False) + monkeypatch.setattr(oks.shutil, "which", lambda _n: None) + argv = oks.resolve_okstratr_argv() + assert argv[-2:] == ["-m", "okstratr"] + + +def test_contract_slice_shape(monkeypatch): + monkeypatch.setattr(oks, "is_healthy", lambda **_: False) + monkeypatch.setattr(oks, "_read_pid", lambda: None) + slice_ = oks.contract_slice() + assert set(slice_) == {"state", "url", "detail"} + assert slice_["state"] in {"starting", "healthy", "unhealthy", "stopped"} From b2edd23f5d26ac2760068c9f11a9b0d1e9570b4a Mon Sep 17 00:00:00 2001 From: benjsmith Date: Sat, 19 Sep 2026 11:43:14 +0000 Subject: [PATCH 06/25] feat(embed): v2 same-document mount + core-skills wait chrome Execute proxied CE/okstratr scripts in-panel (no iframe): rewrite asset URLs under /embed/*, inject markup, createElement-append classic/module scripts, tear down on soft-reload. Poll /api/core-skills/status for starting/unhealthy/building-wiki/live banners and suppress 502 flash while supervisors are still starting. --- docs/ADR-004-same-origin-embed-proxy.md | 4 +- docs/ADR-004b-embed-v2-same-document-mount.md | 78 ++++ frontend/src/index.css | 33 +- .../src/widgets/embed/ProxiedSkillPanel.tsx | 290 +++++++++++---- frontend/src/widgets/embed/embedMount.test.ts | 145 ++++++++ frontend/src/widgets/embed/embedMount.ts | 345 ++++++++++++++++++ 6 files changed, 823 insertions(+), 72 deletions(-) create mode 100644 docs/ADR-004b-embed-v2-same-document-mount.md create mode 100644 frontend/src/widgets/embed/embedMount.test.ts create mode 100644 frontend/src/widgets/embed/embedMount.ts diff --git a/docs/ADR-004-same-origin-embed-proxy.md b/docs/ADR-004-same-origin-embed-proxy.md index 26a52f5..de15a54 100644 --- a/docs/ADR-004-same-origin-embed-proxy.md +++ b/docs/ADR-004-same-origin-embed-proxy.md @@ -1,6 +1,6 @@ # ADR-004: Same-origin embed reverse-proxy (no iframes) -- **Status:** Accepted (Phase 4a) +- **Status:** Accepted (Phase 4a); Embed v2 mount → [ADR-004b](./ADR-004b-embed-v2-same-document-mount.md) - **Date:** 2026-09-18 - **Deciders:** Ben / skill-shell rationalization charter @@ -32,7 +32,7 @@ frames**. 2. **No iframes** for Graph/Agents skill surfaces. Feature flag `proxied_skill_embeds` (default **false**) switches Graph → CE panel and Agents → okstratr panel that navigate `/embed/*` via same-origin - `fetch` + same-document rendering (script-stripped HTML / JSON). + `fetch` + same-document rendering. Phase 4a was script-stripped HTML / JSON; Embed v2 (ADR-004b) executes skill scripts in-panel. Built-in GraphTab / AgentDashboardTab / filebrowser remain the default and are **not deleted**. diff --git a/docs/ADR-004b-embed-v2-same-document-mount.md b/docs/ADR-004b-embed-v2-same-document-mount.md new file mode 100644 index 0000000..ec018ef --- /dev/null +++ b/docs/ADR-004b-embed-v2-same-document-mount.md @@ -0,0 +1,78 @@ +# ADR-004b: Embed v2 — same-document interactive mount + +- **Status:** Accepted (Embed v2) +- **Date:** 2026-09-19 +- **Parent:** [ADR-004](./ADR-004-same-origin-embed-proxy.md) +- **Contract:** umbrella `CONTRACT-EMBED-V2.md` + +## Context + +Phase 4a (`ProxiedSkillPanel`) fetched `/embed/*` HTML and injected a +**script-stripped** body. That proved the reverse-proxy path but left Graph +and Agents inert (no pan/zoom, no desk controls). Charter + ADR-004 forbid +`