Merge #c1156d20: Terminal developer environment, resumable TUI, and Arc…

Terminal developer environment, resumable TUI, and Archipelago skills

nostr:nevent1qqsvz9tdypdv9cmtjmsz6azv8yujk0xezwv0j4wpa57qnk60nfyj28gpz3mhxue69uhhyetvv9ujumn8d96zuer9wcffu24e

PR-Author: Personal
nostr:npub1w3sqdkrhn0gyuvsex32effzgnfpyde6qrrc4u467flg5e9txh4wsfn5vjg

PR description:

Adds the authenticated resumable terminal session backend, xterm-based terminal UI with bounded scrolling and app-style controls, tmux key translation, Codex/node bootstrap, preloaded Archipelago skills, developer app starter, and UAT runbook.

Validation: npm run build:production; cargo check -p archipelago; scripts/tests/archy-developer-setup-test.sh. Framework UAT health and unauthenticated terminal endpoint checks pass.
This commit is contained in:
archipelago
2026-10-09 15:08:54 -04:00
35 changed files with 2089 additions and 266 deletions
+20
View File
@@ -14,6 +14,7 @@ mod remote_input;
mod remote_relay;
mod rental_playback;
mod routstr_proxy;
mod terminal;
mod websocket;
use crate::api::rpc::RpcHandler;
@@ -426,6 +427,16 @@ impl ApiHandler {
.await;
}
// Owner terminal attachment — the browser socket is disposable; the
// authenticated tmux session survives reconnects and browser closes.
if method == Method::GET && path == "/ws/terminal" {
if !self.is_authenticated(req.headers()).await {
tracing::warn!("401 WebSocket /ws/terminal — session invalid or missing");
return Ok(Self::unauthorized());
}
return Self::handle_terminal_websocket(req).await;
}
// Remote input WebSocket — companion app sends keyboard/mouse events
if method == Method::GET && path == "/ws/remote-input" {
if !self.is_authenticated(req.headers()).await {
@@ -544,6 +555,15 @@ impl ApiHandler {
.unwrap())
}
(Method::GET, "/api/terminal/sessions") => {
if !self.is_authenticated(&headers).await { return Ok(Self::unauthorized()); }
terminal::list_response().await
}
(Method::POST, "/api/terminal/sessions") => {
if !self.is_authenticated(&headers).await { return Ok(Self::unauthorized()); }
terminal::create(&body_bytes).await
}
// Node message — P2P endpoint (authenticated by source validation, not cookie)
(Method::POST, "/archipelago/node-message") => {
Self::handle_node_message(body_bytes).await
@@ -0,0 +1,279 @@
//! Owner-authenticated terminal sessions backed by private tmux processes.
//!
//! Browser connections are disposable attachments. The tmux process and its
//! metadata remain on the node so a reconnect resumes the same workspace.
use anyhow::{anyhow, Result};
use futures_util::{SinkExt, StreamExt};
use hyper::{Request, Response, StatusCode};
use serde::{Deserialize, Serialize};
use std::path::{Path, PathBuf};
use tokio::process::Command;
use tokio_tungstenite::tungstenite::Message;
fn tmux_key_for_input(data: &str) -> Option<&'static str> {
match data {
"\u{3}" => Some("C-c"),
"\u{4}" => Some("C-d"),
"\r" | "\n" => Some("Enter"),
"\u{7f}" => Some("BSpace"),
"\t" => Some("Tab"),
"\u{1b}[A" => Some("Up"),
"\u{1b}[B" => Some("Down"),
"\u{1b}[C" => Some("Right"),
"\u{1b}[D" => Some("Left"),
"\u{1b}[H" => Some("Home"),
"\u{1b}[F" => Some("End"),
"\u{1b}[3~" => Some("DC"),
_ => None,
}
}
use uuid::Uuid;
use super::{build_response, ApiHandler};
#[derive(Debug, Clone, Serialize, Deserialize)]
pub(crate) struct SessionRecord {
pub schema: u8,
pub id: String,
pub name: String,
pub workspace: String,
pub state: String,
pub tmux: String,
pub updated_at: i64,
}
#[derive(Debug, Deserialize)]
struct CreateRequest {
name: Option<String>,
workspace: Option<String>,
}
#[derive(Debug, Deserialize)]
struct ClientMessage {
#[serde(rename = "type")]
kind: String,
data: Option<String>,
cols: Option<u16>,
rows: Option<u16>,
}
pub(crate) fn state_dir() -> PathBuf {
std::env::var_os("ARCHY_SESSION_STATE_DIR")
.map(PathBuf::from)
.unwrap_or_else(|| {
if let Some(xdg) = std::env::var_os("XDG_STATE_HOME") {
PathBuf::from(xdg).join("archipelago/sessions")
} else if let Some(home) = std::env::var_os("HOME") {
let user_dir = PathBuf::from(home).join(".local/state/archipelago/sessions");
if user_dir.exists() {
user_dir
} else {
PathBuf::from("/var/lib/archipelago/sessions")
}
} else {
PathBuf::from("/var/lib/archipelago/sessions")
}
})
}
fn valid_id(id: &str) -> bool {
!id.is_empty()
&& id.len() <= 64
&& id
.bytes()
.all(|b| b.is_ascii_alphanumeric() || b == b'.' || b == b'_' || b == b'-')
}
async fn read_record(dir: &Path, id: &str) -> Result<SessionRecord> {
if !valid_id(id) {
return Err(anyhow!("invalid session id"));
}
let bytes = tokio::fs::read(dir.join(format!("{id}.json"))).await?;
Ok(serde_json::from_slice(&bytes)?)
}
async fn tmux_alive(name: &str) -> bool {
Command::new("tmux")
.args(["has-session", "-t", name])
.output()
.await
.map(|out| out.status.success())
.unwrap_or(false)
}
async fn write_record(dir: &Path, record: &SessionRecord) -> Result<()> {
tokio::fs::create_dir_all(dir).await?;
let tmp = dir.join(format!(".{}.tmp-{}", record.id, Uuid::new_v4()));
let final_path = dir.join(format!("{}.json", record.id));
tokio::fs::write(&tmp, serde_json::to_vec_pretty(record)?).await?;
tokio::fs::rename(tmp, final_path).await?;
Ok(())
}
pub(crate) async fn list(dir: &Path) -> Result<Vec<SessionRecord>> {
let mut out = Vec::new();
let mut entries = match tokio::fs::read_dir(dir).await {
Ok(entries) => entries,
Err(error) if error.kind() == std::io::ErrorKind::NotFound => return Ok(out),
Err(error) => return Err(error.into()),
};
while let Some(entry) = entries.next_entry().await? {
if entry.path().extension().and_then(|s| s.to_str()) != Some("json") {
continue;
}
let Ok(bytes) = tokio::fs::read(entry.path()).await else {
continue;
};
let Ok(mut record) = serde_json::from_slice::<SessionRecord>(&bytes) else {
continue;
};
record.state = if tmux_alive(&record.tmux).await {
"detached".into()
} else if record.state == "running" {
"interrupted".into()
} else {
record.state.clone()
};
out.push(record);
}
out.sort_by(|a, b| b.updated_at.cmp(&a.updated_at));
Ok(out)
}
pub(crate) async fn list_response() -> Result<Response<hyper::Body>> {
Ok(Response::builder()
.status(StatusCode::OK)
.header("Content-Type", "application/json")
.body(hyper::Body::from(serde_json::to_vec(
&list(&state_dir()).await?,
)?))?)
}
pub(crate) async fn create(body: &[u8]) -> Result<Response<hyper::Body>> {
let request: CreateRequest = serde_json::from_slice(body).unwrap_or(CreateRequest {
name: None,
workspace: None,
});
let workspace_was_requested = request.workspace.is_some();
let workspace = request.workspace.unwrap_or_else(|| {
std::env::var_os("HOME")
.map(PathBuf::from)
.unwrap_or_else(|| PathBuf::from("/tmp"))
.join("Work")
.to_string_lossy()
.into_owned()
});
let workspace_path = PathBuf::from(&workspace);
if !workspace_was_requested {
tokio::fs::create_dir_all(&workspace_path).await?;
}
if !workspace_path.is_absolute() || !workspace_path.is_dir() || workspace_path == Path::new("/")
{
return Ok(build_response(
StatusCode::BAD_REQUEST,
"application/json",
hyper::Body::from(r#"{"error":"workspace must be an existing non-root directory"}"#),
));
}
let id = format!("s-{}", Uuid::new_v4().simple());
let tmux_name = format!("archy-{id}");
let output = Command::new("tmux")
.args(["new-session", "-d", "-s", &tmux_name, "-c", &workspace])
.output()
.await?;
if !output.status.success() {
return Ok(build_response(
StatusCode::SERVICE_UNAVAILABLE,
"application/json",
hyper::Body::from(r#"{"error":"tmux could not start the session"}"#),
));
}
let name = request
.name
.filter(|n| !n.trim().is_empty())
.unwrap_or_else(|| "Work".into());
let record = SessionRecord {
schema: 1,
id,
name,
workspace,
state: "running".into(),
tmux: tmux_name,
updated_at: chrono::Utc::now().timestamp(),
};
write_record(&state_dir(), &record).await?;
Ok(Response::builder()
.status(StatusCode::CREATED)
.header("Content-Type", "application/json")
.body(hyper::Body::from(serde_json::to_vec(&record)?))?)
}
pub(crate) async fn websocket(req: Request<hyper::Body>) -> Result<Response<hyper::Body>> {
let id = req
.uri()
.query()
.and_then(|query| {
query
.split('&')
.find_map(|part| part.strip_prefix("session="))
})
.unwrap_or("")
.to_string();
let record = read_record(&state_dir(), &id).await?;
if !tmux_alive(&record.tmux).await {
return Ok(build_response(
StatusCode::CONFLICT,
"application/json",
hyper::Body::from(r#"{"error":"session is not running"}"#),
));
}
let (response, ws_fut) =
hyper_ws_listener::create_ws(req).map_err(|e| anyhow!("WebSocket upgrade failed: {e}"))?;
if let Some(ws_fut) = ws_fut {
tokio::spawn(async move {
let Ok(Ok(stream)) = ws_fut.await else { return };
let (mut tx, mut rx) = stream.split();
let mut interval = tokio::time::interval(std::time::Duration::from_millis(150));
let mut last = String::new();
loop {
tokio::select! {
_ = interval.tick() => {
// Capture the visible pane only. Asking tmux for a large
// historical range injects hundreds of blank rows into
// xterm on every reconnect, producing a misleading
// giant initial scrollbar. xterm owns live scrollback.
let output = Command::new("tmux").args(["capture-pane", "-p", "-e", "-t", &record.tmux]).output().await;
if let Ok(output) = output {
let text = String::from_utf8_lossy(&output.stdout).into_owned();
if text != last { last = text.clone(); if tx.send(Message::Text(serde_json::json!({"type":"output", "data":text}).to_string())).await.is_err() { break; } }
} else { break; }
}
message = rx.next() => match message {
Some(Ok(Message::Text(text))) => {
let Ok(message) = serde_json::from_str::<ClientMessage>(&text) else { continue };
match message.kind.as_str() {
"input" => if let Some(data) = message.data { let mut command = Command::new("tmux"); command.args(["send-keys", "-t", &record.tmux]); if let Some(key) = tmux_key_for_input(&data) { command.arg(key); } else { command.args(["-l", "--", &data]); } let _ = command.output().await; },
"resize" => if let (Some(cols), Some(rows)) = (message.cols, message.rows) { let _ = Command::new("tmux").args(["resize-window", "-t", &record.tmux, "-x", &cols.to_string(), "-y", &rows.to_string()]).output().await; },
"ping" => { let _ = tx.send(Message::Text(r#"{"type":"pong"}"#.into())).await; },
_ => {}
}
}
Some(Ok(Message::Close(_))) | None => break,
Some(Err(_)) => break,
_ => {}
}
}
}
});
}
Ok(response)
}
impl ApiHandler {
pub(super) async fn handle_terminal_websocket(
req: Request<hyper::Body>,
) -> Result<Response<hyper::Body>> {
websocket(req).await
}
}