use anyhow::Result;
use axum::extract::{Path, Request, State};
use axum::middleware::{self, Next};
use axum::response::sse::{Event, KeepAlive, Sse};
use axum::response::{Html, IntoResponse, Response};
use axum::routing::{get, post};
use axum::{http::StatusCode, Json, Router};
use serde::{Deserialize, Serialize};
use std::convert::Infallible;
use std::path::PathBuf;
use std::sync::Arc;
use tokio::sync::{broadcast, mpsc, Mutex, RwLock};
use tokio_stream::wrappers::UnboundedReceiverStream;
use tokio_stream::StreamExt;
use crate::config::Config;
use crate::providers::CatalogModel;
use crate::runtime::{AgentEvent, Runtime};
use crate::session::{Session, SessionMeta};
#[derive(Clone, Default)]
pub(crate) struct Snapshot {
pub(crate) status: serde_json::Value,
pub(crate) models: Vec<CatalogModel>,
pub(crate) busy: bool,
}
#[derive(Clone)]
pub struct AppState {
pub(crate) inner: Arc<Mutex<Inner>>,
pub(crate) snap: Arc<RwLock<Snapshot>>,
pub(crate) dashboard_token: Option<String>,
pub(crate) allowed_hosts: Arc<Vec<String>>,
pub events: broadcast::Sender<String>,
}
pub(crate) struct Inner {
pub(crate) runtime: Runtime,
pub(crate) session: Session,
}
pub fn event_frame(session: &str, ev: &AgentEvent) -> String {
serde_json::json!({
"type": "event",
"session": session,
"event": { "kind": ev.kind, "text": ev.text },
"ts": chrono::Utc::now().to_rfc3339(),
})
.to_string()
}
pub(crate) fn broadcast_event(tx: &broadcast::Sender<String>, session: &str, ev: &AgentEvent) {
let _ = tx.send(event_frame(session, ev));
}
impl AppState {
pub async fn turn(&self, message: &str) -> anyhow::Result<String> {
{
let mut s = self.snap.write().await;
s.busy = true;
}
let mut g = self.inner.lock().await;
let Inner { runtime, session } = &mut *g;
let sid = session.id().to_string();
let tx = self.events.clone();
let result = runtime
.turn(session, message, |ev| broadcast_event(&tx, &sid, &ev))
.await;
let status = runtime.status_json(Some(session));
drop(g);
let mut s = self.snap.write().await;
s.status = status;
s.busy = false;
result
}
pub(crate) async fn sessions_payload(&self) -> serde_json::Value {
sessions_payload()
}
pub(crate) async fn session_payload(&self, id: &str) -> serde_json::Value {
let bus = crate::mailbox::Bus::default();
match Session::load(id) {
Ok(s) => session_value(s, id, &bus),
Err(e) => serde_json::json!({"ok": false, "error": e.to_string()}),
}
}
pub(crate) async fn status_payload(&self) -> serde_json::Value {
let snap = self.snap.read().await;
let v = snap.status.clone();
if v.is_null() {
drop(snap);
let g = self.inner.lock().await;
return g.runtime.status_json(Some(&g.session));
}
overlay_busy(v, snap.busy)
}
#[allow(dead_code)]
pub(crate) async fn message_session(&self, id: &str, text: &str) -> Response {
message_session_send(id, text).await
}
}
#[derive(Deserialize)]
pub struct ChatReq {
pub message: String,
pub session_id: Option<String>,
pub model: Option<String>,
}
#[derive(Deserialize)]
pub struct CompletionMessage {
pub role: String,
#[serde(default)]
pub content: String,
}
#[derive(Deserialize)]
pub struct CompletionReq {
pub messages: Vec<CompletionMessage>,
pub model: Option<String>,
}
pub fn render_completion_prompt(messages: &[CompletionMessage]) -> Option<String> {
let mut system = Vec::new();
let mut transcript = Vec::new();
for m in messages {
let text = m.content.trim();
if text.is_empty() {
continue;
}
match m.role.as_str() {
"system" | "developer" => system.push(text.to_string()),
"assistant" => transcript.push(format!("[you] {text}")),
_ => transcript.push(text.to_string()),
}
}
if transcript.is_empty() {
return None;
}
let mut prompt = String::from("[external chat request]\n");
if !system.is_empty() {
prompt.push_str("Caller instructions:\n");
prompt.push_str(&system.join("\n"));
prompt.push_str("\n\n");
}
prompt.push_str("Conversation so far (oldest first):\n");
prompt.push_str(&transcript.join("\n"));
prompt.push_str("\n\nReply to the latest message. Your reply is posted verbatim.");
Some(prompt)
}
#[derive(Serialize)]
pub struct ChatRes {
pub session_id: String,
pub reply: String,
pub events: Vec<String>,
}
pub async fn serve(cfg: Config, cwd: PathBuf) -> Result<()> {
Config::ensure_home()?;
let is_loopback = matches!(
cfg.dashboard_host.as_str(),
"127.0.0.1" | "localhost" | "::1"
);
if !is_loopback && cfg.dashboard_token.is_none() {
anyhow::bail!("remote dashboard bind requires VARYNTH_DASHBOARD_TOKEN");
}
tokio::spawn(crate::automation::run_loop(cfg.clone(), cwd.clone()));
let mut runtime = Runtime::new(cfg.clone(), cwd.clone())?;
let (relay_tx, relay_rx) = tokio::sync::broadcast::channel::<crate::approval_relay::RelayRequest>(
crate::gateway::RELAY_CAPACITY,
);
runtime.remote_approval_events = Some(relay_tx);
let session = Session::new(&cwd.display().to_string(), &runtime.cfg.model)?;
let host = cfg.dashboard_host.clone();
let port = cfg.dashboard_port;
let models = runtime.list_models().await.unwrap_or_default();
let status = runtime.status_json(Some(&session));
let (events_tx, _) = broadcast::channel(crate::gateway::EVENT_CAPACITY);
let state = AppState {
inner: Arc::new(Mutex::new(Inner { runtime, session })),
snap: Arc::new(RwLock::new(Snapshot {
status,
models,
busy: false,
})),
dashboard_token: cfg.dashboard_token.clone(),
allowed_hosts: Arc::new(host_allowlist(&cfg)),
events: events_tx,
};
let shared = Arc::new(state.clone());
let protected = Router::new()
.route("/api/status", get(api_status))
.route("/api/models", get(api_models))
.route("/api/sessions", get(api_sessions))
.route("/api/session/new", post(api_session_new))
.route("/api/session/{id}", get(api_session))
.route("/api/session/{id}/messages", post(api_session_message))
.route("/api/chat", post(api_chat))
.route("/api/chat/stream", post(api_chat_stream))
.route("/api/channels", get(api_channels))
.route("/api/doctor", get(api_doctor))
.route("/v1/chat/completions", post(v1_chat_completions))
.with_state(state.clone())
.layer(middleware::from_fn_with_state(
state.clone(),
require_dashboard_auth,
));
let app = host_guard(
Router::new()
.route("/", get(index))
.route("/health", get(|| async { "ok" }))
.route("/diff", get(diff_page))
.route("/diff.html", get(diff_page))
.route("/GATEWAY.md", get(gateway_doc))
.with_state(state.clone())
.merge(protected)
.merge(crate::gateway::router(shared.clone())),
state.allowed_hosts.clone(),
);
crate::telegram::spawn(cfg.clone(), state.clone());
tokio::spawn(crate::gateway::control_event_forwarder(shared.clone()));
tokio::spawn(crate::gateway::approval_relay_forwarder(shared, relay_rx));
let addr = format!("{host}:{port}");
let listener = tokio::net::TcpListener::bind(&addr).await?;
eprintln!("varynth dashboard → http://{addr}");
if host == "0.0.0.0" || host == "::" {
if let Ok(ip) = local_ipv4() {
eprintln!("varynth dashboard LAN → http://{ip}:{port}");
}
}
axum::serve(listener, app).await?;
Ok(())
}
fn host_guard<S>(routes: Router<S>, allowed_hosts: Arc<Vec<String>>) -> Router<S>
where
S: Clone + Send + Sync + 'static,
{
routes.layer(middleware::from_fn_with_state(
allowed_hosts,
require_dashboard_host,
))
}
async fn require_dashboard_host(
State(allowed_hosts): State<Arc<Vec<String>>>,
request: Request,
next: Next,
) -> Response {
let host = request
.headers()
.get(axum::http::header::HOST)
.and_then(|v| v.to_str().ok())
.unwrap_or_default();
if host_header_allowed(host, &allowed_hosts) {
next.run(request).await
} else {
(StatusCode::BAD_REQUEST, "invalid host header").into_response()
}
}
fn host_allowlist(cfg: &Config) -> Vec<String> {
let port = cfg.dashboard_port;
[cfg.dashboard_host.as_str(), "127.0.0.1", "localhost"]
.into_iter()
.map(|h| format!("{h}:{port}"))
.collect()
}
fn host_header_allowed(header: &str, allowed: &[String]) -> bool {
allowed.iter().any(|h| h.eq_ignore_ascii_case(header))
}
pub(crate) fn same_origin(headers: &axum::http::HeaderMap) -> bool {
let Some(origin) = headers.get(axum::http::header::ORIGIN) else {
return true;
};
let Some(host) = headers
.get(axum::http::header::HOST)
.and_then(|v| v.to_str().ok())
else {
return false;
};
let Ok(origin) = origin.to_str() else {
return false;
};
origin.eq_ignore_ascii_case(&format!("http://{host}"))
|| origin.eq_ignore_ascii_case(&format!("https://{host}"))
}
pub(crate) fn tokens_equal(provided: &str, expected: &str) -> bool {
let a = provided.as_bytes();
let b = expected.as_bytes();
let mut diff = (a.len() ^ b.len()) as u64;
for i in 0..a.len().max(b.len()) {
diff |= (a.get(i).copied().unwrap_or(0) ^ b.get(i).copied().unwrap_or(0)) as u64;
}
diff == 0
}
async fn require_dashboard_auth(
State(state): State<AppState>,
request: Request,
next: Next,
) -> Response {
if !same_origin(request.headers()) {
return (StatusCode::FORBIDDEN, "foreign origin rejected").into_response();
}
let Some(expected) = state.dashboard_token.as_deref() else {
return next.run(request).await;
};
let authorized = request
.headers()
.get(axum::http::header::AUTHORIZATION)
.and_then(|value| value.to_str().ok())
.and_then(|value| value.strip_prefix("Bearer "))
.is_some_and(|provided| tokens_equal(provided, expected));
if authorized {
next.run(request).await
} else {
(
StatusCode::UNAUTHORIZED,
"dashboard authentication required",
)
.into_response()
}
}
fn local_ipv4() -> Result<std::net::Ipv4Addr> {
let sock = std::net::UdpSocket::bind("0.0.0.0:0")?;
sock.connect("8.8.8.8:80")?;
match sock.local_addr()?.ip() {
std::net::IpAddr::V4(v) => Ok(v),
_ => anyhow::bail!("no ipv4"),
}
}
async fn index() -> Html<&'static str> {
Html(INDEX)
}
async fn diff_page() -> Html<&'static str> {
Html(include_str!("../web/diff.html"))
}
async fn gateway_doc() -> &'static str {
include_str!("../web/GATEWAY.md")
}
async fn api_status(State(st): State<AppState>) -> impl IntoResponse {
Json(st.status_payload().await)
}
async fn api_models(State(st): State<AppState>) -> impl IntoResponse {
let snap = st.snap.read().await;
if !snap.models.is_empty() {
return Json(serde_json::json!({"ok": true, "models": snap.models}));
}
drop(snap);
let g = st.inner.lock().await;
match g.runtime.list_models().await {
Ok(m) => Json(serde_json::json!({"ok": true, "models": m})),
Err(e) => Json(serde_json::json!({"ok": false, "error": e.to_string()})),
}
}
async fn api_session_new(State(st): State<AppState>) -> impl IntoResponse {
let mut g = st.inner.lock().await;
let cwd = g.runtime.jail.cwd.display().to_string();
let model = g.runtime.cfg.model.clone();
match Session::new(&cwd, &model) {
Ok(s) => {
let id = s.id().to_string();
g.session = s;
let status = g.runtime.status_json(Some(&g.session));
drop(g);
let mut snap = st.snap.write().await;
snap.status = status;
snap.busy = false;
Json(serde_json::json!({"ok": true, "session_id": id}))
}
Err(e) => Json(serde_json::json!({"ok": false, "error": e.to_string()})),
}
}
async fn api_sessions(State(st): State<AppState>) -> impl IntoResponse {
Json(st.sessions_payload().await)
}
pub(crate) fn sessions_payload() -> serde_json::Value {
let bus = crate::mailbox::Bus::default();
match Session::list() {
Ok(sessions) => sessions_value(sessions, &bus),
Err(e) => serde_json::json!({"ok": false, "error": e.to_string()}),
}
}
fn sessions_value(sessions: Vec<SessionMeta>, bus: &crate::mailbox::Bus) -> serde_json::Value {
let sessions: Vec<serde_json::Value> = sessions
.into_iter()
.map(|meta| {
let mut v = serde_json::to_value(&meta).unwrap_or(serde_json::Value::Null);
if let Some(obj) = v.as_object_mut() {
obj.insert(
"unread".into(),
serde_json::json!(bus.unread(&meta.id).len()),
);
}
v
})
.collect();
serde_json::json!({"ok": true, "sessions": sessions})
}
fn session_value(s: Session, id: &str, bus: &crate::mailbox::Bus) -> serde_json::Value {
serde_json::json!({
"ok": true,
"meta": s.meta,
"messages": s.messages,
"unread": bus.unread(id).len()
})
}
fn overlay_busy(mut status: serde_json::Value, busy: bool) -> serde_json::Value {
if let Some(obj) = status.as_object_mut() {
obj.insert("busy".into(), serde_json::json!(busy));
}
status
}
async fn message_session_send(id: &str, text: &str) -> Response {
if text.trim().is_empty() {
return (
StatusCode::BAD_REQUEST,
Json(serde_json::json!({"ok": false, "error": "text is required"})),
)
.into_response();
}
match crate::mailbox::Bus::default().send("dashboard", id, text) {
Ok(_) => Json(serde_json::json!({"ok": true, "to": id})).into_response(),
Err(e) => (
StatusCode::INTERNAL_SERVER_ERROR,
Json(serde_json::json!({"ok": false, "error": e.to_string()})),
)
.into_response(),
}
}
async fn api_session(State(st): State<AppState>, Path(id): Path<String>) -> impl IntoResponse {
Json(st.session_payload(&id).await)
}
async fn api_session_message(
Path(id): Path<String>,
Json(body): Json<serde_json::Value>,
) -> Response {
let text = body
.get("text")
.and_then(|t| t.as_str())
.unwrap_or_default();
message_session_send(&id, text).await
}
async fn api_channels() -> impl IntoResponse {
Json(crate::channels::v2_status())
}
async fn api_doctor(State(st): State<AppState>) -> impl IntoResponse {
let g = st.inner.lock().await;
match crate::doctor::run(&g.runtime.cfg).await {
Ok(v) => Json(v),
Err(e) => Json(serde_json::json!({"ok": false, "error": e.to_string()})),
}
}
async fn api_chat_stream(
State(st): State<AppState>,
Json(req): Json<ChatReq>,
) -> Sse<impl tokio_stream::Stream<Item = Result<Event, Infallible>>> {
let (tx, rx) = mpsc::unbounded_channel::<AgentEvent>();
tokio::spawn(async move {
{
let mut s = st.snap.write().await;
s.busy = true;
}
let mut g = st.inner.lock().await;
if let Some(id) = &req.session_id {
if id != g.session.id() {
if let Ok(s) = Session::load(id) {
g.session = s;
}
}
}
if let Some(model) = req.model {
if !model.is_empty() {
g.runtime.cfg.model = model;
}
}
let sid = g.session.id().to_string();
let tx_ev = tx.clone();
let ev_tx = st.events.clone();
let mut next = req.message.clone();
let mut last_ok: Option<String> = None;
let mut last_err: Option<String> = None;
loop {
let Inner { runtime, session } = &mut *g;
let result = runtime
.turn(session, &next, |ev| {
let _ = tx_ev.send(ev.clone());
broadcast_event(&ev_tx, &sid, &ev);
})
.await;
let cont = runtime.wants_goal_continue();
match result {
Ok(reply) => {
last_ok = Some(reply);
if !cont {
break;
}
let _ = tx.send(AgentEvent {
kind: "system".into(),
text: "goal still active — continuing".into(),
});
next = "Goal still active. Continue uninterrupted. Do not ask what to do. End with a line that is exactly GOAL_COMPLETE only when the condition is fully met.".into();
}
Err(e) => {
last_err = Some(e.to_string());
break;
}
}
}
let status = {
let Inner { runtime, session } = &mut *g;
runtime.status_json(Some(session))
};
drop(g);
{
let mut s = st.snap.write().await;
s.status = status;
s.busy = false;
}
if let Some(e) = last_err {
let _ = tx.send(AgentEvent {
kind: "error".into(),
text: e,
});
} else if let Some(reply) = last_ok {
let _ = tx.send(AgentEvent {
kind: "done".into(),
text: reply,
});
let _ = tx.send(AgentEvent {
kind: "session".into(),
text: sid,
});
}
});
let stream = UnboundedReceiverStream::new(rx).map(|ev| {
Ok(Event::default()
.event(ev.kind.clone())
.json_data(&ev)
.unwrap_or_else(|_| Event::default().data(ev.text)))
});
Sse::new(stream).keep_alive(KeepAlive::default())
}
async fn api_chat(State(st): State<AppState>, Json(req): Json<ChatReq>) -> impl IntoResponse {
let mut g = st.inner.lock().await;
if let Some(id) = &req.session_id {
if id != g.session.id() {
if let Ok(s) = Session::load(id) {
g.session = s;
}
}
}
if let Some(model) = req.model {
if !model.is_empty() {
g.runtime.cfg.model = model;
}
}
let mut events = Vec::new();
let Inner { runtime, session } = &mut *g;
let sid = session.id().to_string();
let tx = st.events.clone();
let result = runtime
.turn(session, &req.message, |ev| {
events.push(format!("{}: {}", ev.kind, ev.text));
broadcast_event(&tx, &sid, &ev);
})
.await;
match result {
Ok(reply) => Json(serde_json::json!({
"ok": true,
"session_id": session.id(),
"reply": reply,
"events": events,
})),
Err(e) => Json(serde_json::json!({
"ok": false,
"error": e.to_string(),
"events": events,
})),
}
}
fn completion_error(status: StatusCode, message: String) -> Response {
(
status,
Json(serde_json::json!({
"error": {"message": message, "type": "varynth_error"}
})),
)
.into_response()
}
async fn v1_chat_completions(
State(st): State<AppState>,
Json(req): Json<CompletionReq>,
) -> Response {
let Some(prompt) = render_completion_prompt(&req.messages) else {
return completion_error(
StatusCode::BAD_REQUEST,
"messages must contain at least one non-system message".into(),
);
};
{
let mut s = st.snap.write().await;
s.busy = true;
}
let mut g = st.inner.lock().await;
let cwd = g.runtime.jail.cwd.display().to_string();
let model = g.runtime.cfg.model.clone();
let result = match Session::new(&cwd, &model) {
Ok(mut session) => {
let sid = session.id().to_string();
let tx = st.events.clone();
g.runtime
.turn(&mut session, &prompt, |ev| broadcast_event(&tx, &sid, &ev))
.await
}
Err(e) => Err(e),
};
drop(g);
{
let mut s = st.snap.write().await;
s.busy = false;
}
match result {
Ok(reply) => Json(serde_json::json!({
"id": format!("chatcmpl-{}", uuid::Uuid::new_v4()),
"object": "chat.completion",
"created": chrono::Utc::now().timestamp(),
"model": req.model.unwrap_or(model),
"choices": [{
"index": 0,
"message": {"role": "assistant", "content": reply},
"finish_reason": "stop"
}]
}))
.into_response(),
Err(e) => completion_error(StatusCode::BAD_GATEWAY, e.to_string()),
}
}
const INDEX: &str = r##"<!doctype html>
<html lang="en">
<head>
<meta charset="utf-8">
<meta name="viewport" content="width=device-width,initial-scale=1">
<title>Varynth Relay Deck</title>
<style>
:root {
--void: #07090b;
--panel: #10161c;
--line: #1d3a3a;
--teal: #2ec4b6;
--coral: #ff6b4a;
--lime: #c8f542;
--mist: #c9d6d4;
--dim: #6f8582;
--ink: #e8f2ef;
}
* { box-sizing: border-box; }
html, body { height: 100%; margin: 0; background: var(--void); color: var(--ink);
font-family: "IBM Plex Sans", "Segoe UI", sans-serif; }
body {
display: grid;
grid-template-columns: 280px 1fr 300px;
grid-template-rows: 64px 1fr 88px;
grid-template-areas:
"brand top top"
"rail main side"
"rail compose compose";
min-height: 100%;
}
@media (max-width: 900px) {
body { grid-template-columns: 1fr; grid-template-rows: auto auto 1fr auto auto;
grid-template-areas: "brand" "top" "main" "compose" "side"; }
.rail { display: none; }
}
.brand { grid-area: brand; display: flex; align-items: center; gap: 12px; padding: 0 18px;
border-right: 1px solid var(--line); border-bottom: 1px solid var(--line); }
.mark { width: 28px; height: 28px; }
.brand b { font-weight: 600; letter-spacing: .08em; font-size: 13px; }
.brand span { color: var(--dim); font-size: 11px; display: block; }
.top { grid-area: top; display: flex; align-items: center; justify-content: space-between;
padding: 0 20px; border-bottom: 1px solid var(--line); }
.route { font-family: "IBM Plex Mono", ui-monospace, monospace; font-size: 12px; color: var(--teal); }
.pills { display: flex; gap: 8px; flex-wrap: wrap; }
.pill { border: 1px solid var(--line); padding: 4px 9px; font-size: 11px; color: var(--mist); }
.pill.on { border-color: var(--teal); color: var(--teal); }
.rail { grid-area: rail; border-right: 1px solid var(--line); overflow: auto; padding: 12px; }
.rail h2 { font-size: 11px; color: var(--dim); font-weight: 500; margin: 8px 4px 12px; }
.sess { padding: 10px 10px 12px; border-left: 2px solid transparent; cursor: pointer; margin-bottom: 4px; }
.sess:hover { background: #0c1418; }
.sess.active { border-left-color: var(--coral); background: #0c1418; }
.sess .t { font-size: 13px; white-space: nowrap; overflow: hidden; text-overflow: ellipsis; }
.sess .m { font-size: 11px; color: var(--dim); margin-top: 4px; font-family: ui-monospace, monospace; }
.main { grid-area: main; overflow: auto; padding: 24px 28px 12px;
background-image: linear-gradient(180deg, transparent 0, transparent 31px, #122 32px);
background-size: 100% 32px; }
.msg { max-width: 72ch; margin: 0 0 22px; }
.msg .who { font-size: 11px; color: var(--dim); margin-bottom: 6px; font-family: ui-monospace, monospace; }
.msg.user .who { color: var(--coral); }
.msg.assistant .who { color: var(--teal); }
.msg.tool .who { color: var(--lime); }
.msg pre { white-space: pre-wrap; word-break: break-word; margin: 0; font: 14px/1.55 "IBM Plex Sans", sans-serif; }
.side { grid-area: side; border-left: 1px solid var(--line); padding: 16px; overflow: auto; }
.side h2 { font-size: 11px; color: var(--dim); font-weight: 500; }
.side ul { list-style: none; padding: 0; margin: 8px 0 20px; font-size: 12px; color: var(--mist); }
.side li { padding: 4px 0; border-bottom: 1px dashed #152; }
.compose { grid-area: compose; display: flex; gap: 10px; padding: 16px 20px;
border-top: 1px solid var(--line); background: var(--panel); }
textarea { flex: 1; resize: none; height: 56px; background: #0a1014; color: var(--ink);
border: 1px solid var(--line); padding: 10px 12px; font: 14px/1.4 inherit; }
textarea:focus { outline: 1px solid var(--teal); }
button.go { background: var(--coral); color: #140804; border: 0; padding: 0 22px;
font-weight: 650; letter-spacing: .04em; cursor: pointer; }
button.go:disabled { opacity: .4; }
button.new { width: 100%; margin: 0 0 12px; background: transparent; color: var(--teal);
border: 1px solid var(--line); padding: 8px; cursor: pointer; font: 12px inherit; }
button.new:hover { border-color: var(--teal); }
select { background: #0a1014; color: var(--mist); border: 1px solid var(--line); padding: 6px 8px; max-width: 280px; }
.err { color: var(--coral); font-size: 13px; }
.dot { width: 8px; height: 8px; border-radius: 50%; background: #3a4442; display: inline-block; margin-right: 6px; }
.dot.on { background: var(--teal); box-shadow: 0 0 6px var(--teal); }
.dot.err { background: var(--coral); }
.gwrow { display: flex; gap: 6px; align-items: center; margin-bottom: 8px; flex-wrap: wrap; }
.gwtok { flex: 1; min-width: 0; background: #0a1014; color: var(--mist); border: 1px solid var(--line);
padding: 6px 8px; font: 12px ui-monospace, monospace; }
.gwtok:focus { outline: 1px solid var(--teal); }
.gwb { background: transparent; color: var(--mist); border: 1px solid var(--line); padding: 4px 8px;
font-size: 11px; cursor: pointer; text-decoration: none; }
.gwb:hover { border-color: var(--teal); color: var(--teal); }
.gwstate { font-size: 11px; color: var(--dim); margin-right: auto; }
.gwfeed { list-style: none; margin: 8px 0 0; padding: 0; font-size: 11px; max-height: 260px; overflow: auto; }
.gwfeed li { padding: 4px 0; border-bottom: 1px dashed #152; color: var(--mist);
word-break: break-word; display: flex; gap: 6px; align-items: baseline; }
.badge { flex: 0 0 auto; font: 10px ui-monospace, monospace; padding: 1px 5px;
border: 1px solid var(--line); color: var(--dim); }
.badge.k-tool { color: var(--lime); border-color: var(--lime); }
.badge.k-assistant, .badge.k-done { color: var(--teal); border-color: var(--teal); }
.badge.k-error { color: var(--coral); border-color: var(--coral); }
.badge.k-approval_request { color: #7E57C2; border-color: #7E57C2; }
.approval { border: 1px solid #7E57C2; padding: 8px; margin: 8px 0; font-size: 12px; }
.approval .apid { font: 11px ui-monospace, monospace; color: #b39ddb; cursor: pointer;
word-break: break-all; margin-top: 4px; }
.approval .abtns { display: flex; gap: 6px; margin-top: 6px; }
.approval button { flex: 1; padding: 4px 0; cursor: pointer; font-size: 11px;
border: 1px solid var(--line); background: transparent; color: var(--mist); }
.approval button.once:hover, .approval button.always:hover { border-color: var(--teal); color: var(--teal); }
.approval button.deny:hover { border-color: var(--coral); color: var(--coral); }
</style>
</head>
<body>
<div class="brand">
<svg class="mark" viewBox="0 0 32 32" aria-hidden="true">
<path d="M4 28 L16 4 L28 28" fill="none" stroke="#2ec4b6" stroke-width="2.2"/>
<circle cx="16" cy="18" r="3" fill="#ff6b4a"/>
</svg>
<div><b>VARYNTH</b><span>local agent deck</span></div>
</div>
<header class="top">
<div class="route" id="route">route · connecting</div>
<div class="pills">
<span class="pill on" id="pill-dash">dashboard</span>
<span class="pill" id="pill-tg">telegram</span>
<span class="pill" id="pill-dc">discord v2</span>
<span class="pill" id="pill-wa">whatsapp v2</span>
</div>
</header>
<aside class="rail">
<h2>sessions</h2>
<button type="button" class="new" id="new-session">new session</button>
<div id="sessions"></div>
</aside>
<main class="main" id="log"></main>
<aside class="side">
<h2>models</h2>
<select id="model"></select>
<h2 style="margin-top:22px">skills</h2>
<ul id="skills"></ul>
<h2>doctor</h2>
<ul id="doctor"></ul>
<h2 style="margin-top:22px">gateway</h2>
<div class="gwrow">
<input id="gw-token" class="gwtok" placeholder="gateway token" autocomplete="off">
</div>
<div class="gwrow">
<span class="dot" id="gw-dot"></span><span id="gw-state" class="gwstate">offline</span>
<button type="button" class="gwb" id="gw-conn">connect</button>
</div>
<div class="gwrow">
<button type="button" class="gwb" id="gw-ping">ping</button>
<button type="button" class="gwb" id="gw-list">sessions</button>
<button type="button" class="gwb" id="gw-status">status</button>
<a class="gwb" href="/diff" target="_blank">diff</a>
</div>
<div id="gw-approvals"></div>
<ul id="gw-feed" class="gwfeed"></ul>
</aside>
<form class="compose" id="form">
<textarea id="input" placeholder="Task for this machine. /skill-name works."></textarea>
<button class="go" type="submit">SEND</button>
</form>
<script>
const log = document.getElementById('log');
const sessionsEl = document.getElementById('sessions');
const modelEl = document.getElementById('model');
let sessionId = null;
function addMsg(role, text) {
const d = document.createElement('div');
d.className = 'msg ' + role;
d.innerHTML = '<div class="who">' + role + '</div><pre></pre>';
d.querySelector('pre').textContent = text;
log.appendChild(d);
log.scrollTop = log.scrollHeight;
}
async function refresh() {
const st = await (await fetch('/api/status')).json();
document.getElementById('route').textContent =
(st.busy ? 'busy · ' : '') + (st.provider || 'proxy') + ' · ' + (st.model || '') + ' · ' + (st.sandbox || '');
sessionId = st.session;
const skills = document.getElementById('skills');
skills.innerHTML = (st.skills || []).map(s => '<li>/' + s + '</li>').join('') || '<li>none</li>';
const models = await (await fetch('/api/models')).json();
if (models.ok) {
const groups = {};
(models.models || []).forEach(m => {
const g = m.owned_by || (m.id.split(':')[0]) || 'other';
(groups[g] = groups[g] || []).push(m);
});
modelEl.innerHTML = Object.keys(groups).sort().map(g => {
const opts = groups[g].map(m =>
'<option value="' + m.id + '"' + (m.id === st.model ? ' selected' : '') + '>' +
(m.display_name ? m.display_name + ' — ' : '') + m.id + '</option>').join('');
return '<optgroup label="' + g + '">' + opts + '</optgroup>';
}).join('');
}
const sess = await (await fetch('/api/sessions')).json();
if (sess.ok) {
sessionsEl.innerHTML = sess.sessions.map(s =>
'<div class="sess' + (s.id === sessionId ? ' active' : '') + '" data-id="' + s.id + '">' +
'<div class="t">' + (s.title || s.id) + '</div>' +
'<div class="m">' + (s.model || '') + '</div></div>').join('');
sessionsEl.querySelectorAll('.sess').forEach(el => el.onclick = async () => {
sessionId = el.dataset.id;
const data = await (await fetch('/api/session/' + sessionId)).json();
log.innerHTML = '';
(data.messages || []).forEach(m => addMsg(m.role, m.content));
refresh();
});
}
const doc = await (await fetch('/api/doctor')).json();
document.getElementById('doctor').innerHTML =
(doc.checks || []).map(c => '<li>' + (c.ok ? 'ok' : 'fail') + ' · ' + c.name + ' · ' + c.detail + '</li>').join('');
const ch = await (await fetch('/api/channels')).json();
const tg = document.getElementById('pill-tg');
if (tg) {
if (ch.telegram && String(ch.telegram).indexOf('live') >= 0) tg.classList.add('on');
else tg.classList.remove('on');
}
}
document.getElementById('form').onsubmit = async (e) => {
e.preventDefault();
const input = document.getElementById('input');
const text = input.value.trim();
if (!text) return;
addMsg('user', text);
input.value = '';
const btn = document.querySelector('.go');
btn.disabled = true;
try {
const res = await fetch('/api/chat/stream', {
method: 'POST',
headers: {'content-type': 'application/json', 'accept': 'text/event-stream'},
body: JSON.stringify({ message: text, session_id: sessionId, model: modelEl.value })
});
if (!res.body) throw new Error('no stream');
const reader = res.body.getReader();
const dec = new TextDecoder();
let buf = '';
while (true) {
const { value, done } = await reader.read();
if (done) break;
buf += dec.decode(value, { stream: true });
const parts = buf.split('\n\n');
buf = parts.pop();
for (const block of parts) {
let ev = 'message';
let data = '';
for (const line of block.split('\n')) {
if (line.startsWith('event:')) ev = line.slice(6).trim();
if (line.startsWith('data:')) data += line.slice(5).trim();
}
if (!data) continue;
let payload = data;
try { const j = JSON.parse(data); payload = j.text || data; ev = j.kind || ev; } catch (_) {}
if (ev === 'session') sessionId = payload;
else if (ev === 'done') addMsg('assistant', payload);
else if (ev === 'assistant') {}
else if (ev === 'error') addMsg('assistant', 'error: ' + payload);
else addMsg('tool', ev + ': ' + payload);
}
}
} catch (err) {
addMsg('assistant', String(err));
}
btn.disabled = false;
refresh();
};
document.getElementById('new-session').onclick = async () => {
const res = await fetch('/api/session/new', { method: 'POST' });
const data = await res.json();
if (data.ok) {
sessionId = data.session_id;
log.innerHTML = '';
addMsg('system', 'new session ' + sessionId);
refresh();
}
};
// ---- gateway control plane ----
const gwFeed = document.getElementById('gw-feed');
const gwDot = document.getElementById('gw-dot');
const gwState = document.getElementById('gw-state');
const gwApprovals = document.getElementById('gw-approvals');
const gwTokenBox = document.getElementById('gw-token');
const gwConnBtn = document.getElementById('gw-conn');
gwTokenBox.value = localStorage.getItem('varynth_gateway_token') || '';
let gws = null, gwSeq = 1;
const gwPending = {};
function gwSetState(state, label) {
gwDot.className = 'dot' + (state === true ? ' on' : state === false ? ' err' : '');
gwState.textContent = label;
}
function gwFeedItem(kind, text) {
const li = document.createElement('li');
const badge = document.createElement('span');
badge.className = 'badge k-' + String(kind).replace(/[^a-z_]/g, '');
badge.textContent = kind;
const body = document.createElement('span');
body.textContent = text;
li.appendChild(badge);
li.appendChild(body);
gwFeed.appendChild(li);
while (gwFeed.children.length > 80) gwFeed.removeChild(gwFeed.firstChild);
gwFeed.scrollTop = gwFeed.scrollHeight;
}
function gwSend(method, params, cb) {
if (!gws || gws.readyState !== 1) {
gwFeedItem('system', 'gateway offline — connect first');
return;
}
const id = 'r' + (gwSeq++);
gwPending[id] = cb || function () {};
gws.send(JSON.stringify({ jsonrpc: '2.0', id: id, method: method, params: params || {} }));
}
// Relay ids travel as event.id or as an [id:<uuid>] token in the text.
const GW_ID_RE = /\[id:([0-9a-fA-F]{8}-[0-9a-fA-F]{4}-[0-9a-fA-F]{4}-[0-9a-fA-F]{4}-[0-9a-fA-F]{12})\]/;
function gwAskDecision(frame) {
const ev = frame.event || {};
const raw = String(ev.text || '');
const match = raw.match(GW_ID_RE);
const rid = ev.id || (match && match[1]) || null;
const shown = raw.replace(/\s*\[id:[0-9a-fA-F-]+\]\s*$/, '');
const card = document.createElement('div');
card.className = 'approval';
const title = document.createElement('div');
title.textContent = 'approval needed · ' + shown;
card.appendChild(title);
if (rid) {
const idEl = document.createElement('div');
idEl.className = 'apid';
idEl.title = 'click to copy · telegram: approve ' + rid + ' once|always|deny';
idEl.textContent = rid;
idEl.onclick = function () {
const mark = function () {
idEl.textContent = rid + ' — copied';
setTimeout(function () { idEl.textContent = rid; }, 1200);
};
if (navigator.clipboard && navigator.clipboard.writeText) {
navigator.clipboard.writeText(rid).then(mark, mark);
} else { mark(); }
};
card.appendChild(idEl);
} else {
const note = document.createElement('div');
note.textContent = 'no relay id in event — approve in the TUI';
note.style.color = 'var(--dim)';
card.appendChild(note);
}
const btns = document.createElement('div');
btns.className = 'abtns';
['once', 'always', 'deny'].forEach(function (decision) {
const b = document.createElement('button');
b.type = 'button';
b.className = decision;
b.textContent = decision;
b.disabled = !rid;
b.onclick = function () {
gwSend('approval.respond', { id: rid, decision: decision }, function (res) {
if (res && res.result === true) {
card.remove();
gwFeedItem('system', 'approval ' + decision + ' sent · ' + rid);
} else {
gwFeedItem('error', 'approval.respond failed (slot consumed or gone)');
}
});
};
btns.appendChild(b);
});
card.appendChild(btns);
gwApprovals.appendChild(card);
while (gwApprovals.children.length > 6) gwApprovals.removeChild(gwApprovals.firstChild);
}
function gwHandleFrame(p) {
if (!p || p.type !== 'event') {
if (p && p.type === 'lagged') gwFeedItem('system', p.text || ('lagged, missed ' + p.missed));
else gwFeedItem('system', JSON.stringify(p).slice(0, 160));
return;
}
const ev = p.event || {};
if (ev.kind === 'approval_request') {
gwFeedItem('approval_request', String(ev.text || '').slice(0, 120));
gwAskDecision(p);
return;
}
gwFeedItem(ev.kind || 'event', String(ev.text || '').slice(0, 200));
}
function gwConnect() {
const token = gwTokenBox.value.trim();
localStorage.setItem('varynth_gateway_token', token);
const proto = location.protocol === 'https:' ? 'wss://' : 'ws://';
const url = proto + location.host + '/ws/gateway' + (token ? '?token=' + encodeURIComponent(token) : '');
gwSetState(null, 'connecting');
try { gws = new WebSocket(url); } catch (e) { gwSetState(false, String(e)); return; }
gws.onopen = function () { gwSetState(true, 'live'); };
gws.onclose = function () {
gws = null;
gwSetState(false, 'closed');
gwConnBtn.textContent = 'connect';
};
gws.onerror = function () { gwSetState(false, 'error'); };
gws.onmessage = function (m) {
let msg;
try { msg = JSON.parse(m.data); } catch (_) { return; }
if (msg.method === 'event') { gwHandleFrame(msg.params); return; }
const cb = gwPending[msg.id];
if (!cb) return;
delete gwPending[msg.id];
cb(msg);
};
}
gwConnBtn.onclick = function () {
if (gws) { const s = gws; gws = null; s.close(); gwSetState(false, 'closed'); this.textContent = 'connect'; }
else { gwConnect(); this.textContent = 'disconnect'; }
};
document.getElementById('gw-ping').onclick = function () {
gwSend('ping', {}, function (r) { gwFeedItem('system', 'ping → ' + JSON.stringify(r.result)); });
};
document.getElementById('gw-list').onclick = function () {
gwSend('sessions.list', {}, function (r) {
const n = ((r.result && r.result.sessions) || []).length;
gwFeedItem('system', 'sessions → ' + n);
});
};
document.getElementById('gw-status').onclick = function () {
gwSend('status.get', {}, function (r) {
gwFeedItem('system', 'status → ' + JSON.stringify(r.result || r.error).slice(0, 140));
});
};
refresh();
setInterval(refresh, 15000);
</script>
</body>
</html>
"##;
#[cfg(test)]
mod tests {
use super::*;
use axum::body::Body;
use tower::ServiceExt;
#[test]
fn token_compare_is_exact_and_length_sensitive() {
assert!(tokens_equal("secret", "secret"));
assert!(tokens_equal("", ""));
assert!(!tokens_equal("secret", "secreT"));
assert!(!tokens_equal("secret", "secret "));
assert!(!tokens_equal("secret", "secrets"));
assert!(!tokens_equal("", "secret"));
}
#[test]
fn event_frame_shape_and_timestamp() {
let ev = AgentEvent {
kind: "tool".into(),
text: "ls -la".into(),
};
let frame: serde_json::Value = serde_json::from_str(&event_frame("sess-1", &ev)).unwrap();
assert_eq!(frame["type"], "event");
assert_eq!(frame["session"], "sess-1");
assert_eq!(frame["event"]["kind"], "tool");
assert_eq!(frame["event"]["text"], "ls -la");
chrono::DateTime::parse_from_rfc3339(frame["ts"].as_str().unwrap()).unwrap();
}
#[tokio::test]
async fn broadcast_event_fans_out_to_every_subscriber() {
let (tx, _) = broadcast::channel(256);
let mut first = tx.subscribe();
let mut second = tx.subscribe();
let ev = AgentEvent {
kind: "assistant".into(),
text: "hi".into(),
};
broadcast_event(&tx, "s-9", &ev);
for rx in [&mut first, &mut second] {
let frame: serde_json::Value = serde_json::from_str(&rx.try_recv().unwrap()).unwrap();
assert_eq!(frame["type"], "event");
assert_eq!(frame["session"], "s-9");
assert_eq!(frame["event"]["text"], "hi");
}
let (lonely, _dropped) = broadcast::channel::<String>(4);
broadcast_event(&lonely, "s-9", &ev);
}
#[test]
fn chat_fanout_hook_keeps_inner_sink_and_lands_frames() {
let (tx, mut rx) = broadcast::channel::<String>(16);
let sid = "sess-42".to_string();
let mut events: Vec<String> = Vec::new();
let mut hook = |ev: AgentEvent| {
events.push(format!("{}: {}", ev.kind, ev.text));
broadcast_event(&tx, &sid, &ev);
};
hook(AgentEvent {
kind: "tool".into(),
text: "ls".into(),
});
hook(AgentEvent {
kind: "assistant".into(),
text: "done".into(),
});
assert_eq!(
events,
vec!["tool: ls".to_string(), "assistant: done".to_string()]
);
let f1: serde_json::Value = serde_json::from_str(&rx.try_recv().unwrap()).unwrap();
assert_eq!(f1["type"], "event");
assert_eq!(f1["session"], "sess-42");
assert_eq!(f1["event"]["kind"], "tool");
assert_eq!(f1["event"]["text"], "ls");
let f2: serde_json::Value = serde_json::from_str(&rx.try_recv().unwrap()).unwrap();
assert_eq!(f2["event"]["kind"], "assistant");
assert_eq!(f2["event"]["text"], "done");
}
#[test]
fn payload_helpers_return_sane_shapes() {
let dir = tempfile::tempdir().unwrap();
let bus = crate::mailbox::Bus::at(dir.path().join("bus.jsonl"));
let meta = SessionMeta {
id: "sess-a".into(),
created_at: chrono::Utc::now(),
updated_at: chrono::Utc::now(),
cwd: dir.path().display().to_string(),
model: "mock".into(),
title: "temp".into(),
title_explicit: false,
};
let v = sessions_value(vec![meta], &bus);
assert_eq!(v["ok"], true);
let arr = v["sessions"].as_array().unwrap();
assert_eq!(arr.len(), 1);
assert_eq!(arr[0]["id"], "sess-a");
assert_eq!(arr[0]["unread"], 0);
let s = Session::create_at(
"sess-b".into(),
dir.path().join("sess-b.jsonl"),
&dir.path().display().to_string(),
"mock",
)
.unwrap();
let v = session_value(s, "sess-b", &bus);
assert_eq!(v["ok"], true);
assert_eq!(v["meta"]["id"], "sess-b");
assert_eq!(v["messages"], serde_json::json!([]));
assert_eq!(v["unread"], 0);
let v = overlay_busy(serde_json::json!({"model": "m", "session": "s"}), true);
assert_eq!(v["busy"], true);
assert_eq!(v["model"], "m");
assert_eq!(
overlay_busy(serde_json::json!({"model": "m"}), false)["busy"],
false
);
}
#[test]
fn host_header_check_rejects_foreign_hosts() {
let cfg = Config {
dashboard_host: "127.0.0.1".into(),
dashboard_port: 7420,
..Config::default()
};
let allowed = host_allowlist(&cfg);
assert!(host_header_allowed("127.0.0.1:7420", &allowed));
assert!(host_header_allowed("LOCALHOST:7420", &allowed));
assert!(!host_header_allowed("evil.example:7420", &allowed));
assert!(!host_header_allowed("127.0.0.1:9999", &allowed));
assert!(!host_header_allowed("localhost", &allowed));
assert!(!host_header_allowed("", &allowed));
}
#[tokio::test]
async fn message_handler_rejects_missing_or_blank_text() {
for body in [
serde_json::json!({}),
serde_json::json!({"text": " "}),
serde_json::json!({"text": 42}),
] {
let res = api_session_message(Path("s-x".into()), Json(body)).await;
assert_eq!(res.status(), StatusCode::BAD_REQUEST);
let bytes = axum::body::to_bytes(res.into_body(), usize::MAX)
.await
.unwrap();
let v: serde_json::Value = serde_json::from_slice(&bytes).unwrap();
assert_eq!(v["ok"], false);
assert!(v["error"].as_str().unwrap().contains("text"));
}
}
#[tokio::test]
async fn messages_route_requires_dashboard_host() {
let cfg = Config {
dashboard_host: "127.0.0.1".into(),
dashboard_port: 7420,
..Config::default()
};
let app = host_guard(
Router::new().route("/api/session/{id}/messages", post(api_session_message)),
Arc::new(host_allowlist(&cfg)),
);
let evil = Request::builder()
.method("POST")
.uri("/api/session/s-x/messages")
.header("host", "evil.example")
.header("content-type", "application/json")
.body(Body::from(r#"{"text":"hi"}"#))
.unwrap();
let res = app.clone().oneshot(evil).await.unwrap();
assert_eq!(res.status(), StatusCode::BAD_REQUEST);
let bytes = axum::body::to_bytes(res.into_body(), usize::MAX)
.await
.unwrap();
assert_eq!(
String::from_utf8(bytes.to_vec()).unwrap(),
"invalid host header"
);
let local = Request::builder()
.method("POST")
.uri("/api/session/s-x/messages")
.header("host", "127.0.0.1:7420")
.header("content-type", "application/json")
.body(Body::from(r#"{"text":" "}"#))
.unwrap();
let res = app.oneshot(local).await.unwrap();
assert_eq!(res.status(), StatusCode::BAD_REQUEST);
let bytes = axum::body::to_bytes(res.into_body(), usize::MAX)
.await
.unwrap();
let v: serde_json::Value = serde_json::from_slice(&bytes).unwrap();
assert_eq!(v["ok"], false);
assert!(v["error"].as_str().unwrap().contains("text"));
}
}