use crate::net::{self, Ev, Peer, Transport};
use anyhow::{anyhow, bail, Result};
use bytes::Bytes;
use serde_json::{json, Value};
use std::collections::{HashMap, VecDeque};
use std::sync::atomic::{AtomicU32, AtomicUsize, Ordering};
use std::sync::Arc;
use std::time::{Duration, Instant};
use tokio::io::{AsyncRead, AsyncReadExt, AsyncWrite, AsyncWriteExt};
use tokio::net::{TcpListener, TcpStream};
use tokio::sync::{mpsc, Mutex};
use tokio::task::AbortHandle;
pub const L2_SID_BASE: u32 = 0x8000_0000;
#[inline]
pub fn is_l2_sid(sid: u32) -> bool {
sid & L2_SID_BASE != 0
}
type PipeItem = Option<Bytes>;
type StreamTx = mpsc::Sender<PipeItem>;
struct StreamHandle {
tx: StreamTx,
read_pump: Option<AbortHandle>,
}
pub const MAX_STREAMS_PER_LINK: usize = 8;
pub const MAX_PTYS_GLOBAL: usize = 32;
pub static LIVE_PTYS: AtomicUsize = AtomicUsize::new(0);
pub struct PtyGuard;
impl PtyGuard {
pub fn try_acquire() -> Option<PtyGuard> {
let mut cur = LIVE_PTYS.load(Ordering::Relaxed);
loop {
if cur >= MAX_PTYS_GLOBAL {
return None;
}
match LIVE_PTYS.compare_exchange_weak(cur, cur + 1, Ordering::AcqRel, Ordering::Relaxed) {
Ok(_) => return Some(PtyGuard),
Err(actual) => cur = actual,
}
}
}
}
impl Drop for PtyGuard {
fn drop(&mut self) {
LIVE_PTYS.fetch_sub(1, Ordering::AcqRel);
}
}
pub struct Mux {
transport: Arc<dyn Transport>,
streams: Mutex<HashMap<u32, StreamHandle>>,
next_sid: AtomicU32,
accepted: Mutex<HashMap<u32, ()>>,
resizers: Mutex<HashMap<u32, mpsc::UnboundedSender<(u16, u16)>>>,
mount_ack_tx: Mutex<HashMap<u32, tokio::sync::oneshot::Sender<serde_json::Value>>>,
}
impl Mux {
pub fn new(t: Arc<dyn Transport>) -> Arc<Self> {
Arc::new(Mux {
transport: t,
streams: Mutex::new(HashMap::new()),
next_sid: AtomicU32::new(0),
accepted: Mutex::new(HashMap::new()),
resizers: Mutex::new(HashMap::new()),
mount_ack_tx: Mutex::new(HashMap::new()),
})
}
pub fn transport(&self) -> Arc<dyn Transport> {
self.transport.clone()
}
fn alloc_sid(&self) -> u32 {
let n = self.next_sid.fetch_add(1, Ordering::Relaxed) & 0x3FFF_FFFF;
let role = if self.transport.sid_answerer() { 0x4000_0000 } else { 0 };
n | L2_SID_BASE | role
}
async fn register(&self, sid: u32) -> mpsc::Receiver<PipeItem> {
let (tx, rx) = mpsc::channel::<PipeItem>(256);
self.streams
.lock()
.await
.insert(sid, StreamHandle { tx, read_pump: None });
rx
}
async fn set_read_pump(&self, sid: u32, h: AbortHandle) {
if let Some(s) = self.streams.lock().await.get_mut(&sid) {
s.read_pump = Some(h);
} else {
h.abort();
}
}
pub async fn register_stream(&self, sid: u32) -> mpsc::Receiver<PipeItem> {
self.register(sid).await
}
pub async fn live_streams(&self) -> usize {
self.streams.lock().await.len()
}
pub async fn at_stream_cap(&self) -> bool {
self.live_streams().await >= MAX_STREAMS_PER_LINK
}
async fn drop_stream(&self, sid: u32) {
self.resizers.lock().await.remove(&sid);
if let Some(s) = self.streams.lock().await.remove(&sid) {
if let Some(h) = s.read_pump {
h.abort();
}
}
}
pub async fn register_resizer(&self, sid: u32, tx: mpsc::UnboundedSender<(u16, u16)>) {
self.resizers.lock().await.insert(sid, tx);
}
pub async fn resize_pty(&self, sid: u32, cols: u16, rows: u16) {
if let Some(tx) = self.resizers.lock().await.get(&sid) {
let _ = tx.send((cols, rows));
}
}
pub async fn drop_pty(&self, sid: u32) {
self.resizers.lock().await.remove(&sid);
self.streams.lock().await.remove(&sid);
}
pub async fn on_frame(&self, sid: u32, payload: Bytes) {
let tx = self.streams.lock().await.get(&sid).map(|s| s.tx.clone());
if let Some(tx) = tx {
let msg = if payload.is_empty() { None } else { Some(payload) };
let _ = tx.send(msg).await; }
}
pub async fn on_open_ack(&self, sid: u32) {
let tx = self.streams.lock().await.get(&sid).map(|s| s.tx.clone());
if let Some(tx) = tx {
let _ = tx.send(Some(Bytes::new())).await;
}
}
async fn on_close(&self, sid: u32, _err: Option<&str>) {
self.mount_ack_tx.lock().await.remove(&sid);
self.drop_stream(sid).await;
}
pub async fn shutdown_all(&self) {
self.resizers.lock().await.clear(); self.mount_ack_tx.lock().await.clear();
let mut map = self.streams.lock().await;
for (_, s) in map.drain() {
if let Some(h) = s.read_pump {
h.abort();
}
}
}
}
async fn socket_to_dc<R: AsyncRead + Unpin>(
transport: Arc<dyn Transport>,
sid: u32,
mut rd: R,
) -> Result<()> {
let cap = transport.max_payload();
let mut buf = vec![0u8; cap];
loop {
let n = rd.read(&mut buf).await?;
if n == 0 {
transport.send_frame(sid, 0, &[]).await?; return Ok(());
}
transport.send_frame(sid, 0, &buf[..n]).await?;
}
}
async fn dc_to_socket<W: AsyncWrite + Unpin>(
mut rx: mpsc::Receiver<PipeItem>,
mut wr: W,
first: Option<PipeItem>,
) -> Result<()> {
if let Some(item) = first {
match item {
Some(bytes) => wr.write_all(&bytes).await?,
None => {
let _ = wr.shutdown().await;
return Ok(());
}
}
}
while let Some(item) = rx.recv().await {
match item {
Some(bytes) => wr.write_all(&bytes).await?,
None => {
let _ = wr.shutdown().await; return Ok(());
}
}
}
let _ = wr.shutdown().await; Ok(())
}
async fn serve_stream<S: AsyncRead + AsyncWrite + Unpin + Send + 'static>(
mux: Arc<Mux>,
sid: u32,
sock: S,
rx: mpsc::Receiver<PipeItem>,
send_close: bool,
first: Option<PipeItem>,
) {
let (rd, wr) = tokio::io::split(sock);
let mut writer = tokio::spawn(dc_to_socket(rx, wr, first));
let mut reader = tokio::spawn(socket_to_dc(mux.transport.clone(), sid, rd));
mux.set_read_pump(sid, reader.abort_handle()).await;
let mut ticker = tokio::time::interval(Duration::from_secs(2));
ticker.tick().await; let read_result;
loop {
tokio::select! {
r = &mut reader => { read_result = Some(r); writer.abort(); break; }
_ = &mut writer => { reader.abort(); read_result = None; break; }
_ = ticker.tick() => {
if !mux.transport.is_alive() {
reader.abort();
writer.abort();
read_result = None;
break;
}
}
}
}
mux.streams.lock().await.remove(&sid);
if send_close {
let close = match read_result {
Some(Ok(Ok(()))) => json!({ "type": "l2-close", "sid": sid }), Some(Ok(Err(e))) => json!({ "type": "l2-close", "sid": sid, "err": e.to_string() }),
Some(Err(_aborted)) => return, None => json!({ "type": "l2-close", "sid": sid }),
};
let _ = mux.transport.send_control(&close).await;
}
}
pub const SESSION_BUFFER_CAP: usize = 256 * 1024;
pub const PTY_REATTACH_RESET: &[u8] = b"\x1b[?1000l\x1b[?1002l\x1b[?1003l\x1b[?1006l\x1b[?1015l";
fn mouse_mode_changes(data: &[u8]) -> (Vec<u16>, Vec<u16>) {
let mut sets: Vec<u16> = Vec::new();
let mut resets: Vec<u16> = Vec::new();
let mut i = 0;
while i < data.len() {
if data[i] != b'\x1b' || i + 3 >= data.len() || data[i + 1] != b'[' || data[i + 2] != b'?' {
i += 1;
continue;
}
let start = i + 3;
let mut j = start;
while j < data.len() && (data[j].is_ascii_digit() || data[j] == b';') {
j += 1;
}
if j > start && j < data.len() && matches!(data[j], b'h' | b'l') {
let set = data[j] == b'h';
let nums: &str = std::str::from_utf8(&data[start..j]).unwrap_or("");
for part in nums.split(';') {
if let Ok(n) = part.parse::<u16>() {
if matches!(n, 1000 | 1002 | 1003 | 1006 | 1015) {
if set {
if !sets.contains(&n) { sets.push(n); }
} else {
if !resets.contains(&n) { resets.push(n); }
}
}
}
}
}
i = j + 1;
}
(sets, resets)
}
pub const SESSION_DETACHED_IDLE: Duration = Duration::from_secs(180);
pub const SESSION_MAX_LIFETIME: Duration = Duration::from_secs(8 * 3600);
struct OutBind {
transport: Arc<dyn Transport>,
sid: u32,
}
enum SessionCmd {
Attach { transport: Arc<dyn Transport>, sid: u32 },
Detach,
End,
}
#[derive(Clone)]
pub struct PtySessionHandle {
input_tx: std::sync::mpsc::Sender<Vec<u8>>,
resize_tx: mpsc::UnboundedSender<(u16, u16)>,
cmd_tx: mpsc::UnboundedSender<SessionCmd>,
dead: Arc<std::sync::atomic::AtomicBool>,
}
impl PtySessionHandle {
pub fn is_dead(&self) -> bool {
self.dead.load(Ordering::Acquire)
}
pub fn feed_input(&self, bytes: Vec<u8>) {
let _ = self.input_tx.send(bytes);
}
pub fn resize(&self, cols: u16, rows: u16) {
let _ = self.resize_tx.send((cols, rows));
}
pub fn attach(&self, transport: Arc<dyn Transport>, sid: u32) {
let _ = self.cmd_tx.send(SessionCmd::Attach { transport, sid });
}
pub fn detach(&self) {
let _ = self.cmd_tx.send(SessionCmd::Detach);
}
pub fn end(&self) {
let _ = self.cmd_tx.send(SessionCmd::End);
}
}
#[derive(Default)]
pub struct PtySessions {
map: Mutex<HashMap<String, PtySessionHandle>>,
}
impl PtySessions {
pub fn new() -> Arc<Self> {
Arc::new(Self::default())
}
pub async fn get_live(&self, id: &str) -> Option<PtySessionHandle> {
let mut map = self.map.lock().await;
match map.get(id) {
Some(h) if !h.is_dead() => Some(h.clone()),
Some(_) => {
map.remove(id);
None
}
None => None,
}
}
pub async fn insert(&self, id: String, h: PtySessionHandle) {
self.map.lock().await.insert(id, h);
}
pub async fn remove(&self, id: &str) {
self.map.lock().await.remove(id);
}
}
pub async fn spawn_pty_session(
sessions: Arc<PtySessions>,
session_id: String,
transport: Arc<dyn Transport>,
sid: u32,
cols: u16,
rows: u16,
term: &str,
argv: Vec<String>,
pty_guard: PtyGuard,
) -> Option<PtySessionHandle> {
use portable_pty::{native_pty_system, CommandBuilder, PtySize};
use std::io::{Read as _, Write as _};
let size = PtySize { rows: rows.max(1), cols: cols.max(1), pixel_width: 0, pixel_height: 0 };
let pair = match native_pty_system().openpty(size) {
Ok(p) => p,
Err(e) => {
let _ = transport.send_control(&json!({ "type": "l2-close", "sid": sid, "err": format!("pty: {e}") })).await;
return None;
}
};
let mut cmd = CommandBuilder::new(&argv[0]);
for a in &argv[1..] {
cmd.arg(a);
}
cmd.env("TERM", if term.is_empty() { "xterm-256color" } else { term });
cmd.env("COLORTERM", "truecolor");
cmd.cwd(crate::platform::Paths::home_dir());
let mut child = match pair.slave.spawn_command(cmd) {
Ok(c) => c,
Err(e) => {
let _ = transport.send_control(&json!({ "type": "l2-close", "sid": sid, "err": format!("spawn: {e}") })).await;
return None;
}
};
drop(pair.slave); let master = pair.master;
let mut reader = match master.try_clone_reader() {
Ok(r) => r,
Err(_) => return None,
};
let mut writer = match master.take_writer() {
Ok(w) => w,
Err(_) => return None,
};
let (otx, mut orx) = mpsc::channel::<Vec<u8>>(128);
std::thread::spawn(move || {
let mut buf = [0u8; 8192];
loop {
match reader.read(&mut buf) {
Ok(0) | Err(_) => break, Ok(n) => {
if otx.blocking_send(buf[..n].to_vec()).is_err() {
break;
}
}
}
}
});
let (input_tx, wrx) = std::sync::mpsc::channel::<Vec<u8>>();
std::thread::spawn(move || {
while let Ok(b) = wrx.recv() {
if writer.write_all(&b).is_err() {
break;
}
let _ = writer.flush();
}
});
let (resize_tx, mut resize_rx) = mpsc::unbounded_channel::<(u16, u16)>();
let (cmd_tx, mut cmd_rx) = mpsc::unbounded_channel::<SessionCmd>();
let dead = Arc::new(std::sync::atomic::AtomicBool::new(false));
let handle = PtySessionHandle {
input_tx,
resize_tx,
cmd_tx,
dead: dead.clone(),
};
let sessions_for_task = sessions.clone();
let session_id_for_task = session_id.clone();
tokio::spawn(async move {
let _guard = pty_guard; let mut bind: Option<OutBind> = Some(OutBind { transport, sid });
let mut ring: VecDeque<u8> = VecDeque::new();
let started = Instant::now();
let mut detached_since: Option<Instant> = None;
let mut active_mouse_modes: Vec<u16> = Vec::new();
let mut reaper = tokio::time::interval(Duration::from_secs(5));
reaper.tick().await;
loop {
tokio::select! {
out = orx.recv() => match out {
Some(bytes) => {
let (sets, resets) = mouse_mode_changes(&bytes);
for m in sets {
if !active_mouse_modes.contains(&m) { active_mouse_modes.push(m); }
}
active_mouse_modes.retain(|m| !resets.contains(m));
push_ring(&mut ring, &bytes, SESSION_BUFFER_CAP);
if let Some(b) = &bind {
for chunk in bytes.chunks(b.transport.max_payload().max(1)) {
if b.transport.send_frame(b.sid, 0, chunk).await.is_err() {
detached_since = Some(Instant::now());
bind = None;
break;
}
}
}
}
None => break, },
cmd = cmd_rx.recv() => match cmd {
Some(SessionCmd::Attach { transport, sid }) => {
let snapshot: Vec<u8> = ring.iter().copied().collect();
let mp = transport.max_payload().max(1);
let mut ok = true;
for chunk in snapshot.chunks(mp) {
if transport.send_frame(sid, 0, chunk).await.is_err() {
ok = false;
break;
}
}
if ok {
if transport.send_frame(sid, 0, PTY_REATTACH_RESET).await.is_err() {
ok = false;
}
for m in &active_mouse_modes {
let seq = format!("\x1b[?{m}h");
if ok && transport.send_frame(sid, 0, seq.as_bytes()).await.is_err() {
ok = false;
break;
}
}
}
if ok {
bind = Some(OutBind { transport, sid });
detached_since = None;
} else {
detached_since = Some(Instant::now());
bind = None;
}
}
Some(SessionCmd::Detach) => {
bind = None;
detached_since = Some(Instant::now());
}
Some(SessionCmd::End) => break, None => break, },
rs = resize_rx.recv() => {
if let Some((c, r)) = rs {
let _ = master.resize(PtySize { rows: r.max(1), cols: c.max(1), pixel_width: 0, pixel_height: 0 });
}
}
_ = reaper.tick() => {
let lifetime_up = started.elapsed() >= SESSION_MAX_LIFETIME;
let idle_up = detached_since
.map(|t| t.elapsed() >= SESSION_DETACHED_IDLE)
.unwrap_or(false);
if lifetime_up || idle_up {
break;
}
}
}
}
let _ = child.kill();
let _ = child.wait();
dead.store(true, Ordering::Release);
if let Some(b) = &bind {
let _ = b.transport.send_control(&json!({ "type": "l2-close", "sid": b.sid })).await;
}
sessions_for_task.remove(&session_id_for_task).await;
});
sessions.insert(session_id, handle.clone()).await;
Some(handle)
}
fn push_ring(ring: &mut VecDeque<u8>, bytes: &[u8], cap: usize) {
if bytes.len() >= cap {
ring.clear();
ring.extend(&bytes[bytes.len() - cap..]);
return;
}
ring.extend(bytes);
while ring.len() > cap {
ring.pop_front();
}
}
pub enum OpenVerdict {
Accept { sid: u32, host: String, port: u16, rx: mpsc::Receiver<PipeItem> },
Deny { sid: u32, err: &'static str },
Ignore,
}
impl Mux {
pub async fn accept_control(&self, v: &Value, trusted: bool, allow_nonloopback: bool) -> OpenVerdict {
match v["type"].as_str() {
Some("l2-open") => {
let Some(sid) = v["sid"].as_u64().map(|s| s as u32) else {
return OpenVerdict::Ignore;
};
if !is_l2_sid(sid) {
return OpenVerdict::Ignore; }
{
let mut acc = self.accepted.lock().await;
if acc.contains_key(&sid) {
return OpenVerdict::Ignore;
}
acc.insert(sid, ());
}
if !trusted {
return OpenVerdict::Deny { sid, err: "denied" };
}
let host = v["host"].as_str().unwrap_or("127.0.0.1").to_string();
let port = v["rport"].as_u64().or_else(|| v["port"].as_u64()).unwrap_or(0) as u16;
if port == 0 {
return OpenVerdict::Deny { sid, err: "bad port" };
}
if !host_is_loopback(&host) && !allow_nonloopback {
return OpenVerdict::Deny { sid, err: "non-loopback denied (not in l2-allow.json)" };
}
if self.at_stream_cap().await {
self.accepted.lock().await.remove(&sid);
return OpenVerdict::Deny { sid, err: "too many streams" };
}
let rx = self.register(sid).await; OpenVerdict::Accept { sid, host, port, rx }
}
Some("l2-close") => {
if let Some(sid) = v["sid"].as_u64() {
self.on_close(sid as u32, v["err"].as_str()).await;
}
OpenVerdict::Ignore
}
_ => OpenVerdict::Ignore,
}
}
pub async fn dial_and_serve(self: Arc<Self>, sid: u32, host: String, port: u16, rx: mpsc::Receiver<PipeItem>) {
match TcpStream::connect((host.as_str(), port)).await {
Ok(sock) => {
let _ = sock.set_nodelay(true);
let _ = self
.transport
.send_control(&json!({ "type": "l2-open-ack", "sid": sid, "credit": 0 }))
.await;
serve_stream(self.clone(), sid, sock, rx, true, None).await;
self.accepted.lock().await.remove(&sid);
}
Err(e) => {
self.drop_stream(sid).await;
self.accepted.lock().await.remove(&sid);
let _ = self
.transport
.send_control(&json!({ "type": "l2-close", "sid": sid, "err": e.to_string() }))
.await;
}
}
}
}
fn host_is_loopback(host: &str) -> bool {
if host.eq_ignore_ascii_case("localhost") {
return true;
}
host.parse::<std::net::IpAddr>().map(|ip| ip.is_loopback()).unwrap_or(false)
}
pub struct LinkGuard {
sio: Option<rust_socketio::asynchronous::Client>,
peer: Option<Arc<Peer>>,
}
impl LinkGuard {
fn forget(mut self) {
if let Some(sio) = self.sio.take() {
std::mem::forget(sio);
}
if let Some(p) = self.peer.take() {
std::mem::forget(p);
}
}
async fn close(mut self) {
if let Some(p) = self.peer.take() {
p.close().await;
}
if let Some(sio) = self.sio.take() {
let _ = sio.disconnect().await;
}
}
}
async fn bring_up_to_known(
server: &str,
peer_name: &str,
relay: bool,
role: &'static str,
) -> Result<(Arc<dyn Transport>, mpsc::UnboundedReceiver<Ev>, LinkGuard, crate::diag::Attempt)> {
let secret = crate::devices_load()
.into_iter()
.find(|(n, _)| n.eq_ignore_ascii_case(peer_name))
.map(|(_, s)| s)
.ok_or_else(|| anyhow!("no known device named '{peer_name}', run `filament pair` first (see `filament devices`)"))?;
let channel = crate::channel_of(&secret);
let mut diag = crate::diag::Attempt::new(server, &crate::diag::peer_hash_from_secret(&secret), role);
let mut entered_establishing = false;
let cfg = net::fetch_config(server).await?;
let (tx, mut rx) = mpsc::unbounded_channel::<Ev>();
let mut sio = net::connect_signaling(server, tx.clone()).await?;
let my_uid = crate::mk_uid("l2");
let solo = format!("l2-{}", crate::fresh_secret());
let join_payload =
json!({ "room": solo, "uid": my_uid, "name": crate::display_name() });
sio.emit("join", join_payload.clone()).await.ok();
let mut my_id: Option<String> = None;
let mut peer: Option<Arc<Peer>> = None;
let mut peer_uid: Option<String> = None;
let mut generation: u32 = 0;
let mut queue: VecDeque<(String, Option<String>)> = VecDeque::new();
let candidate_secs: u64 = std::env::var("FILAMENT_L2_CANDIDATE_SECS")
.ok()
.and_then(|v| v.parse::<u64>().ok())
.filter(|n| *n > 0)
.unwrap_or(7);
let mut endpoint: Option<quinn::Endpoint> = None;
let mut direct_cands: Option<Vec<String>> = None;
let mut direct_racing = false;
let spawn_timer = |pid: String, g: u32| {
let tx = tx.clone();
tokio::spawn(async move {
tokio::time::sleep(std::time::Duration::from_secs(candidate_secs)).await;
let _ = tx.send(Ev::Stuck(pid, g));
});
};
let silent = role.starts_with("reconnect");
if !silent {
crate::ui::say(&match role {
"bootstrap" => format!("filament: authenticating with '{peer_name}'..."),
_ => format!("\rfilament: waiting for known device '{peer_name}'..."),
});
}
let mut resubscribe = tokio::time::interval(Duration::from_millis(2000));
resubscribe.tick().await; let mut welcome_silent_ticks: u32 = 0;
let connect_started = tokio::time::Instant::now();
let mut heartbeat = tokio::time::interval(Duration::from_secs(7));
heartbeat.tick().await;
loop {
if peer.is_none() {
if let Some((pid, uid)) = queue.pop_front() {
if !entered_establishing {
diag.enter(crate::diag::Phase::Establishing);
entered_establishing = true;
}
let mine = my_id.clone().unwrap_or_default();
let polite = net::polite_role(&my_uid, uid.as_deref(), &mine, &pid);
generation += 1;
spawn_timer(pid.clone(), generation);
let p = Peer::connect(
pid.clone(), polite, cfg.ice_servers.clone(), relay,
sio.clone(), tx.clone(), generation,
)
.await?;
peer_uid = uid;
peer = Some(p);
if !direct_racing && crate::direct::direct_enabled() {
if endpoint.is_none() {
match crate::direct::bind_endpoint() {
Ok((ep, port)) => {
direct_cands =
Some(crate::direct::gather_candidates(server, port).await);
endpoint = Some(ep);
crate::ui::trace(&format!("filament: DIRECT-OFFER sent to '{peer_name}', port {port}"));
}
Err(e) => {
crate::ui::trace(&format!("filament: direct disabled (endpoint bind failed: {e}), WebRTC only"));
}
}
}
if endpoint.is_some() {
if let Some(c) = &direct_cands {
let offer =
json!({ "type": "transport-offer", "v": 1, "addrs": c });
sio.emit("signal", json!({ "to": pid, "data": offer })).await.ok();
}
}
}
}
}
let ev = tokio::select! {
ev = rx.recv() => match ev {
Some(ev) => ev,
None => break,
},
_ = heartbeat.tick() => {
if role != "doctor" {
crate::ui::say(&format!(
"filament: still reaching '{peer_name}'... ({}s)",
connect_started.elapsed().as_secs()
));
}
continue;
}
_ = resubscribe.tick() => {
if my_id.is_none() {
welcome_silent_ticks += 1;
if welcome_silent_ticks >= 2 {
welcome_silent_ticks = 0;
if let Ok(fresh) = net::reconnect_signaling(server, tx.clone()).await {
let _ = sio.disconnect().await;
sio = fresh;
}
}
sio.emit("join", join_payload.clone()).await.ok();
} else if peer.is_none() && queue.is_empty() {
net::subscribe_with_ack(&sio, vec![channel.clone()], tx.clone()).await;
}
continue;
}
};
match ev {
Ev::Welcome(v) => {
my_id = v["id"].as_str().map(|s| s.to_string());
diag.enter(crate::diag::Phase::Presence);
net::subscribe_with_ack(&sio, vec![channel.clone()], tx.clone()).await;
}
Ev::KnownPeer(v) => {
if v["channel"].as_str() != Some(channel.as_str()) {
continue;
}
let pid = match v["id"].as_str() {
Some(p) => p.to_string(),
None => continue,
};
if Some(pid.as_str()) == my_id.as_deref() {
continue;
}
if crate::is_self_uid(&my_uid, v["uid"].as_str()) {
continue;
}
if peer.as_ref().is_some_and(|p| p.id == pid)
|| queue.iter().any(|(q, _)| *q == pid)
{
continue;
}
queue.push_back((pid, v["uid"].as_str().map(|s| s.to_string())));
}
Ev::Signal(v) => {
let data = v["data"].clone();
if data["type"].as_str() == Some("transport-offer") {
if direct_racing {
continue; }
let ep = endpoint
.take()
.or_else(|| crate::direct::bind_endpoint().ok().map(|(ep, _)| ep));
if let Some(ep) = ep {
direct_racing = true;
let peer_cands: Vec<String> = data["addrs"]
.as_array()
.map(|a| a.iter().filter_map(|x| x.as_str().map(String::from)).collect())
.unwrap_or_default();
crate::ui::debug(&format!(
"filament: got transport-offer ({} cand), racing direct-quic",
peer_cands.len()
));
let secret = secret.clone();
let pid = v["from"].as_str().unwrap_or_default().to_string();
let tx = tx.clone();
tokio::spawn(async move {
if let Some(t) = crate::direct::race_connect_labeled(
ep, peer_cands, &secret, pid.clone(), tx.clone(), "direct-quic", false,
)
.await
{
let _ = tx.send(Ev::DirectReady(pid, t, "direct-quic"));
}
});
}
continue;
}
let from = v["from"].as_str().unwrap_or_default();
let Some(p) = &peer else { continue };
if p.id != from {
continue;
}
match p.handle_signal(data).await {
Ok(net::SignalOutcome::Handled) => {}
Ok(net::SignalOutcome::Glare(offer)) => {
let old = peer.take().unwrap();
let pid = old.id.clone();
old.mark_closed();
tokio::spawn(async move { old.close().await });
generation += 1;
spawn_timer(pid.clone(), generation);
let p = Peer::connect(
pid, true, cfg.ice_servers.clone(), relay,
sio.clone(), tx.clone(), generation,
)
.await?;
if let Err(e) = p.handle_signal(offer).await {
crate::ui::trace(&format!("filament: signal: {e}"));
}
peer = Some(p);
}
Err(e) => crate::ui::trace(&format!("filament: signal: {e}")),
}
}
Ev::DirectReady(_pid, t, route) => {
if !role.starts_with("reconnect") && role != "bootstrap" {
crate::ui::say(&format!("\rfilament: tunnel up to '{peer_name}' (route: {route})"));
}
diag.enter(crate::diag::Phase::Ready);
let guard = LinkGuard { sio: Some(sio), peer: peer.take() };
return Ok((t, rx, guard, diag));
}
Ev::Stuck(pid, g) => {
if g == generation && peer.as_ref().is_some_and(|p| p.id == pid) {
let p = peer.take().unwrap();
p.mark_closed();
tokio::spawn(async move { p.close().await });
diag.stall(crate::diag::Phase::Establishing, candidate_secs * 1000);
crate::ui::debug("filament: candidate unresponsive, rotating");
queue.push_back((pid, peer_uid.take()));
}
}
Ev::ChannelReady(pid, t) if peer.as_ref().is_some_and(|p| p.id == pid) => {
if let Some(p) = &peer {
if let Some((my_fp, their_fp)) = p.fingerprints().await {
let mac = crate::proof_for(
&secret, &my_uid, &my_uid,
peer_uid.as_deref().unwrap_or(""), &my_fp, &their_fp,
);
t.send_control(&json!({ "type": "pair-proof", "mac": mac })).await?;
}
}
if !role.starts_with("reconnect") && role != "bootstrap" {
crate::ui::say(&format!("filament: tunnel up to '{peer_name}'"));
}
diag.enter(crate::diag::Phase::Ready);
let guard = LinkGuard { sio: Some(sio), peer: peer.take() };
return Ok((t, rx, guard, diag));
}
Ev::PcState(pid, s) if s == "failed" || s == "closed" => {
if peer.as_ref().is_some_and(|p| p.id == pid) {
let p = peer.take().unwrap();
p.mark_closed();
tokio::spawn(async move { p.close().await });
crate::ui::debug(&format!("filament: connection {s}, rotating"));
queue.push_back((pid, peer_uid.take()));
}
}
_ => {}
}
}
diag.fail("signaling ended before a data channel came up");
Err(anyhow!("signaling ended before a data channel came up"))
}
pub struct ProbeOutcome {
pub timings: Vec<crate::diag::PhaseTiming>,
pub total_ms: u64,
pub established: bool,
pub failed_phase: Option<crate::diag::Phase>,
pub error: Option<String>,
pub path: Option<crate::net::PathInfo>,
}
pub async fn establish_probe(server: &str, peer: &str, relay: bool) -> Result<ProbeOutcome> {
let probe_secs: u64 = std::env::var("FILAMENT_DOCTOR_PROBE_SECS")
.ok()
.and_then(|v| v.parse::<u64>().ok())
.filter(|n| *n > 0)
.unwrap_or(30);
let deadline = std::time::Duration::from_secs(probe_secs);
match tokio::time::timeout(deadline, bring_up_to_known(server, peer, relay, "doctor")).await {
Ok(Ok((t, rx, guard, mut diag))) => {
let mux = Mux::new(t);
let pump = tokio::spawn(pump_initiator(rx, mux.clone()));
diag.enter(crate::diag::Phase::L2Open);
let probe_rport: u16 = 9; let (sid, _rx_pipe) = open_stream(&mux, probe_rport).await?;
diag.up("tunnel", "datachannel-or-direct");
let path = Some(
crate::net::describe_path(mux.transport().as_ref(), guard.peer.as_deref()).await,
);
mux.drop_stream(sid).await;
let _ = mux
.transport()
.send_control(&json!({ "type": "l2-close", "sid": sid }))
.await;
mux.shutdown_all().await;
pump.abort();
guard.close().await;
Ok(ProbeOutcome {
timings: diag.timings().to_vec(),
total_ms: diag.total_ms(),
established: true,
failed_phase: None,
error: None,
path,
})
}
Ok(Err(e)) => {
let (timings, failed_phase) = match crate::diag::latest_span_ladder() {
Some((t, p)) => (t, Some(p)),
None => (Vec::new(), None),
};
Ok(ProbeOutcome {
timings,
total_ms: 0,
established: false,
failed_phase,
error: Some(e.to_string()),
path: None,
})
}
Err(_elapsed) => {
let (timings, failed_phase) = match crate::diag::latest_span_ladder() {
Some((t, p)) => (t, Some(p)),
None => (Vec::new(), None),
};
Ok(ProbeOutcome {
timings,
total_ms: probe_secs * 1000,
established: false,
failed_phase,
error: Some(format!("establishment timed out after {probe_secs}s")),
path: None,
})
}
}
}
async fn pump_initiator(mut rx: mpsc::UnboundedReceiver<Ev>, mux: Arc<Mux>) {
while let Some(ev) = rx.recv().await {
match ev {
Ev::Control(_pid, v) => match v["type"].as_str() {
Some("l2-close") => {
if let Some(sid) = v["sid"].as_u64() {
mux.on_close(sid as u32, v["err"].as_str()).await;
}
}
Some("l2-open-ack") => {
if let Some(sid) = v["sid"].as_u64() {
mux.on_open_ack(sid as u32).await;
}
}
Some("mount-open-ack") => {
if let Some(sid) = v["sid"].as_u64() {
let mut ack_map = mux.mount_ack_tx.lock().await;
if let Some(tx) = ack_map.remove(&(sid as u32)) {
let caps = v.get("caps").cloned().unwrap_or(serde_json::Value::Null);
let _ = tx.send(caps);
}
}
}
_ => {}
},
Ev::Chunk(_pid, sid, _offset, data) if is_l2_sid(sid) => {
mux.on_frame(sid, data).await;
}
Ev::PcState(_, s) if s == "failed" || s == "closed" || s == "disconnected" => {
crate::ui::debug(&format!("filament: tunnel {s}, closing streams"));
mux.shutdown_all().await;
}
_ => {}
}
}
mux.shutdown_all().await;
}
pub(crate) async fn open_stream(mux: &Arc<Mux>, rport: u16) -> Result<(u32, mpsc::Receiver<PipeItem>)> {
let sid = mux.alloc_sid();
let rx = mux.register(sid).await;
let host = std::env::var("FILAMENT_L2_DIALHOST").unwrap_or_else(|_| "127.0.0.1".to_string());
mux.transport
.send_control(&json!({ "type": "l2-open", "sid": sid, "host": host, "rport": rport }))
.await?;
Ok((sid, rx))
}
#[cfg(unix)]
pub(crate) fn warm_verify_window() -> std::time::Duration {
let ms = std::env::var("FILAMENT_WARM_VERIFY_MS")
.ok()
.and_then(|s| s.parse::<u64>().ok())
.filter(|n| *n > 0)
.unwrap_or(2500);
std::time::Duration::from_millis(ms)
}
#[cfg(unix)]
async fn verify_first_frame(
mux: &Arc<Mux>,
sid: u32,
mut rx: mpsc::Receiver<PipeItem>,
verify: std::time::Duration,
) -> Result<(PipeItem, mpsc::Receiver<PipeItem>)> {
match tokio::time::timeout(verify, rx.recv()).await {
Ok(Some(first)) => Ok((first, rx)),
Ok(None) => Err(anyhow!("warm stream closed before any frame")),
Err(_) => {
mux.streams.lock().await.remove(&sid);
let _ = mux
.transport
.send_control(&json!({ "type": "l2-close", "sid": sid }))
.await;
Err(anyhow!("warm link unresponsive after {}ms (zombie)", verify.as_millis()))
}
}
}
#[cfg(unix)]
pub(crate) async fn open_stream_verified(
mux: &Arc<Mux>,
rport: u16,
verify: std::time::Duration,
) -> Result<(u32, PipeItem, mpsc::Receiver<PipeItem>)> {
let (sid, rx) = open_stream(mux, rport).await?;
let (first, rx) = verify_first_frame(mux, sid, rx, verify).await?;
Ok((sid, first, rx))
}
#[cfg(unix)]
pub(crate) async fn serve_verified_stream<S: AsyncRead + AsyncWrite + Unpin + Send + 'static>(
mux: Arc<Mux>,
sid: u32,
sock: S,
first: PipeItem,
rx: mpsc::Receiver<PipeItem>,
) {
serve_stream(mux, sid, sock, rx, true, Some(first)).await;
}
pub(crate) async fn open_mount_stream(
mux: &Arc<Mux>,
root: &str,
) -> Result<(u32, mpsc::Receiver<PipeItem>, crate::mount_proto::MountCaps)> {
let sid = mux.alloc_sid();
let rx = mux.register(sid).await;
let (tx, caps_rx) = tokio::sync::oneshot::channel();
mux.mount_ack_tx.lock().await.insert(sid, tx);
let encoded = crate::mount_proto::path_encode(std::path::Path::new(root));
mux.transport
.send_control(&json!({ "type": "mount-open", "sid": sid, "root": encoded }))
.await?;
let caps_value = tokio::time::timeout(
std::time::Duration::from_secs(10),
caps_rx,
)
.await
.map_err(|_| anyhow::anyhow!("mount-open-ack not received (timed out)"))?
.map_err(|_| anyhow::anyhow!("mount-open-ack channel closed"))?;
let caps = crate::mount_proto::parse_mount_caps(caps_value)?;
let _ = mux.transport.send_control(&json!({
"type": "mount-cap-ack", "sid": sid,
"binary_frames": caps.protocol_version >= 2
})).await;
Ok((sid, rx, caps))
}
pub(crate) async fn open_pty_stream(
mux: &Arc<Mux>,
session: &str,
cols: u16,
rows: u16,
term: &str,
cmd: &str,
) -> Result<(u32, mpsc::Receiver<PipeItem>)> {
let sid = mux.alloc_sid();
let rx = mux.register(sid).await;
let mut ctl = json!({
"type": "pty-open", "sid": sid, "session": session, "cols": cols, "rows": rows, "term": term
});
if !cmd.is_empty() {
ctl["cmd"] = json!(cmd);
}
mux.transport.send_control(&ctl).await?;
Ok((sid, rx))
}
#[cfg(unix)]
pub(crate) async fn open_pty_stream_verified(
mux: &Arc<Mux>,
session: &str,
cols: u16,
rows: u16,
term: &str,
cmd: &str,
verify: std::time::Duration,
) -> Result<(u32, PipeItem, mpsc::Receiver<PipeItem>)> {
let (sid, rx) = open_pty_stream(mux, session, cols, rows, term, cmd).await?;
let (first, rx) = verify_first_frame(mux, sid, rx, verify).await?;
Ok((sid, first, rx))
}
#[cfg(unix)]
pub(crate) async fn serve_opened_stream<S: AsyncRead + AsyncWrite + Unpin + Send + 'static>(
mux: Arc<Mux>,
sid: u32,
stream: S,
rx: mpsc::Receiver<PipeItem>,
) {
serve_stream(mux, sid, stream, rx, true, None).await;
}
#[cfg(unix)]
async fn pump_stdio_over(sock: tokio::net::UnixStream) -> Result<()> {
let (mut rd, mut wr) = tokio::io::split(sock);
let writer = tokio::spawn(async move {
let mut stdin = tokio::io::stdin();
let _ = tokio::io::copy(&mut stdin, &mut wr).await; let _ = wr.shutdown().await; });
let mut stdout = tokio::io::stdout();
tokio::io::copy(&mut rd, &mut stdout).await?;
let _ = stdout.flush().await;
writer.abort();
Ok(())
}
#[cfg(unix)]
async fn pump_warm_pty_stdio(
sock: tokio::net::UnixStream,
stdin_rx: &mut tokio::sync::mpsc::UnboundedReceiver<Vec<u8>>,
pending: &mut Option<Vec<u8>>,
) -> Result<()> {
let (mut rd, mut wr) = tokio::io::split(sock);
if let Some(buf) = pending.take() {
if wr.write_all(&buf).await.is_err() {
*pending = Some(buf);
return Ok(());
}
let _ = wr.flush().await;
}
let mut stdout = tokio::io::stdout();
let mut buf = [0u8; 16 * 1024];
loop {
tokio::select! {
r = rd.read(&mut buf) => match r {
Ok(0) | Err(_) => break, Ok(n) => { stdout.write_all(&buf[..n]).await?; stdout.flush().await?; }
},
chunk = stdin_rx.recv() => match chunk {
Some(c) if c.is_empty() => { let _ = wr.shutdown().await; } Some(c) => {
if wr.write_all(&c).await.is_err() { *pending = Some(c); break; }
let _ = wr.flush().await;
}
None => { let _ = wr.shutdown().await; }
},
}
}
let _ = stdout.flush().await;
Ok(())
}
#[cfg(unix)]
async fn pump_warm_pty_one_shot(
sock: tokio::net::UnixStream,
stdin_rx: &mut tokio::sync::mpsc::UnboundedReceiver<Vec<u8>>,
pending: &mut Option<Vec<u8>>,
) -> Result<()> {
let (mut rd, mut wr) = tokio::io::split(sock);
if let Some(buf) = pending.take() {
if wr.write_all(&buf).await.is_err() {
*pending = Some(buf);
return Ok(());
}
let _ = wr.flush().await;
}
let mut stdout = tokio::io::stdout();
let mut buf = [0u8; 16 * 1024];
loop {
tokio::select! {
r = rd.read(&mut buf) => match r {
Ok(0) | Err(_) => break, Ok(n) => {
stdout.write_all(&buf[..n]).await?;
stdout.flush().await?;
}
},
chunk = stdin_rx.recv() => match chunk {
Some(c) if c.is_empty() => { let _ = wr.shutdown().await; } Some(c) => {
if wr.write_all(&c).await.is_err() { *pending = Some(c); break; }
let _ = wr.flush().await;
}
None => { let _ = wr.shutdown().await; }
},
}
}
let _ = stdout.flush().await;
Ok(())
}
#[cfg(unix)]
pub async fn dial_cmd(peer: &str, port: u16) -> Result<()> {
match crate::ctl::try_dial(peer, port).await {
Some(sock) => pump_stdio_over(sock).await,
None => bail!(
"could not dial {peer}.mesh:{port} over the overlay (is the daemon up, the peer paired, and the port expose'd on it?)"
),
}
}
#[cfg(not(unix))]
pub async fn dial_cmd(_peer: &str, _port: u16) -> Result<()> {
bail!("filament dial needs the local daemon's control socket (unix only)")
}
pub async fn netcat_cmd(server: &str, peer: &str, rport: u16, relay: bool) -> Result<()> {
#[cfg(unix)]
if !relay {
if let Some(sock) = crate::ctl::try_open(peer, rport).await {
crate::ui::trace(&format!("filament: reusing warm link to '{peer}' (no establish)"));
return pump_stdio_over(sock).await;
}
}
let connect_secs: u64 = std::env::var("FILAMENT_CONNECT_SECS")
.ok()
.and_then(|s| s.parse().ok())
.filter(|n| *n > 0)
.unwrap_or(45);
let (t, rx, guard, mut diag) = match tokio::time::timeout(
std::time::Duration::from_secs(connect_secs),
bring_up_to_known(server, peer, relay, "init"),
)
.await
{
Ok(inner) => inner?,
Err(_) => {
crate::ui::problem(
&format!("filament netcat: can't reach '{peer}'"),
&format!(
"couldn't establish a link to '{peer}' in {connect_secs}s - it may be offline or unreachable from here."
),
&[
format!("check it's reachable: {}", crate::ui::paint(crate::ui::Tone::Brand, &format!("filament ping {peer}"))),
format!("diagnose the connect: {}", crate::ui::paint(crate::ui::Tone::Brand, &format!("filament doctor {peer}"))),
],
);
std::process::exit(1);
}
};
guard.forget(); let mux = Mux::new(t);
let pump = tokio::spawn(pump_initiator(rx, mux.clone()));
diag.enter(crate::diag::Phase::L2Open);
let (sid, mut rx_pipe) = open_stream(&mux, rport).await?;
diag.up("tunnel", "datachannel-or-direct");
let t_in = mux.transport();
let reader = tokio::spawn(async move {
let mut stdin = tokio::io::stdin();
let cap = t_in.max_payload();
let mut buf = vec![0u8; cap];
loop {
match stdin.read(&mut buf).await {
Ok(0) | Err(_) => {
let _ = t_in.send_frame(sid, 0, &[]).await; break;
}
Ok(n) => {
if t_in.send_frame(sid, 0, &buf[..n]).await.is_err() {
break;
}
}
}
}
});
let mut stdout = tokio::io::stdout();
while let Some(item) = rx_pipe.recv().await {
match item {
Some(bytes) => {
stdout.write_all(&bytes).await?;
stdout.flush().await?;
}
None => break, }
}
let _ = reader.await;
mux.drop_stream(sid).await;
let _ = mux
.transport()
.send_control(&json!({ "type": "l2-close", "sid": sid }))
.await;
pump.abort();
Ok(())
}
struct RawGuard {
active: bool,
}
impl RawGuard {
fn enable() -> Result<Self> {
crossterm::terminal::enable_raw_mode()?;
Ok(RawGuard { active: true })
}
}
impl Drop for RawGuard {
fn drop(&mut self) {
if self.active {
let _ = crossterm::terminal::disable_raw_mode();
crossterm::execute!(std::io::stderr(), crossterm::cursor::Show).ok();
eprint!("\r\n");
}
}
}
enum PtyOutcome {
Exited,
Dropped,
}
fn spawn_stdin_reader() -> tokio::sync::mpsc::UnboundedReceiver<Vec<u8>> {
let (tx, rx) = tokio::sync::mpsc::unbounded_channel::<Vec<u8>>();
std::thread::spawn(move || {
use std::io::Read;
let mut stdin = std::io::stdin().lock();
let mut buf = [0u8; 16 * 1024];
loop {
match stdin.read(&mut buf) {
Ok(0) | Err(_) => {
let _ = tx.send(Vec::new()); break;
}
Ok(n) => {
if tx.send(buf[..n].to_vec()).is_err() {
break; }
}
}
}
});
rx
}
async fn send_frames_chunked(t: &Arc<dyn Transport>, sid: u32, data: &[u8]) -> std::result::Result<(), ()> {
let cap = t.max_payload().max(1);
for chunk in data.chunks(cap) {
if t.send_frame(sid, 0, chunk).await.is_err() {
return Err(());
}
}
Ok(())
}
#[allow(clippy::too_many_arguments)]
async fn pty_attach_once(
server: &str,
peer: &str,
relay: bool,
role: &'static str,
session_id: &str,
term: &str,
cmd: &str,
interactive: bool,
resume: bool,
raw: &mut Option<RawGuard>,
stdin_rx: &mut tokio::sync::mpsc::UnboundedReceiver<Vec<u8>>,
pending: &mut Option<Vec<u8>>,
) -> Result<PtyOutcome> {
let connect_secs: u64 = std::env::var("FILAMENT_CONNECT_SECS")
.ok()
.and_then(|s| s.parse().ok())
.filter(|n| *n > 0)
.unwrap_or(45);
let (t, rx, guard, mut diag) = match tokio::time::timeout(
std::time::Duration::from_secs(connect_secs),
bring_up_to_known(server, peer, relay, role),
)
.await
{
Ok(inner) => inner?,
Err(_) => {
bail!("connect timeout: couldn't reach '{peer}' in {connect_secs}s");
}
};
guard.forget();
let mux = Mux::new(t);
let pump = tokio::spawn(pump_initiator(rx, mux.clone()));
diag.enter(crate::diag::Phase::L2Open);
let sid = mux.alloc_sid();
let mut rx_pipe = mux.register(sid).await;
let (cols, rows) = if interactive {
crossterm::terminal::size().unwrap_or((80, 24))
} else {
(
std::env::var("COLUMNS").ok().and_then(|s| s.parse().ok()).unwrap_or(80u16),
std::env::var("LINES").ok().and_then(|s| s.parse().ok()).unwrap_or(24u16),
)
};
mux.transport()
.send_control(&{
let mut ctl = json!({
"type": "pty-open", "sid": sid, "session": session_id,
"cols": cols, "rows": rows, "term": term, "resume": resume,
});
if !cmd.is_empty() {
ctl["cmd"] = json!(cmd);
}
ctl
})
.await?;
diag.up("tunnel", "datachannel-or-direct");
if interactive && raw.is_none() {
*raw = Some(RawGuard::enable()?);
}
#[cfg(unix)]
let winch = if interactive {
let t_resize = mux.transport();
Some(tokio::spawn(async move {
use tokio::signal::unix::{signal, SignalKind};
let mut sig = match signal(SignalKind::window_change()) {
Ok(s) => s,
Err(_) => return,
};
while sig.recv().await.is_some() {
if let Ok((c, r)) = crossterm::terminal::size() {
let _ = t_resize
.send_control(&json!({ "type": "pty-resize", "sid": sid, "cols": c, "rows": r }))
.await;
}
}
}))
} else {
None
};
let t_in = mux.transport();
if let Some(buf) = pending.take() {
if send_frames_chunked(&t_in, sid, &buf).await.is_err() {
*pending = Some(buf);
#[cfg(unix)]
if let Some(w) = winch {
w.abort();
}
mux.drop_stream(sid).await;
pump.abort();
return Ok(PtyOutcome::Dropped);
}
}
let mut stdout = tokio::io::stdout();
let mut ticker = tokio::time::interval(Duration::from_secs(2));
ticker.tick().await; let mut stdin_done = false; let dropped;
loop {
tokio::select! {
item = rx_pipe.recv() => match item {
Some(Some(bytes)) => {
stdout.write_all(&bytes).await?;
stdout.flush().await?;
}
_ => { dropped = !mux.transport().is_alive(); break; }
},
chunk = stdin_rx.recv(), if !stdin_done => match chunk {
Some(c) if c.is_empty() => { let _ = t_in.send_frame(sid, 0, &[]).await; stdin_done = true; }
Some(c) => {
if send_frames_chunked(&t_in, sid, &c).await.is_err() {
*pending = Some(c);
dropped = true;
break;
}
}
None => { stdin_done = true; } },
_ = ticker.tick() => {
if !mux.transport().is_alive() { dropped = true; break; }
}
}
}
#[cfg(unix)]
if let Some(w) = winch {
w.abort();
}
mux.drop_stream(sid).await;
if !dropped {
let _ = mux.transport().send_control(&json!({ "type": "l2-close", "sid": sid })).await;
}
pump.abort();
Ok(if dropped { PtyOutcome::Dropped } else { PtyOutcome::Exited })
}
#[cfg(unix)]
async fn try_warm_pty(
peer: &str,
session: &str,
term: &str,
cmd: &str,
interactive: bool,
raw: &mut Option<RawGuard>,
stdin_rx: &mut tokio::sync::mpsc::UnboundedReceiver<Vec<u8>>,
pending: &mut Option<Vec<u8>>,
) -> Option<Result<()>> {
let (cols, rows) = crossterm::terminal::size().unwrap_or((80, 24));
let sock = crate::ctl::try_pty(peer, session, cols, rows, term, cmd).await?; if interactive {
crate::ui::trace(&format!("filament: reusing warm link to '{peer}' for pty (no establish)"));
if raw.is_none() {
match RawGuard::enable() {
Ok(g) => *raw = Some(g),
Err(e) => return Some(Err(e)),
}
}
let session_owned = session.to_string();
let winch = tokio::spawn(async move {
use tokio::signal::unix::{signal, SignalKind};
let mut sig = match signal(SignalKind::window_change()) {
Ok(s) => s,
Err(_) => return,
};
while sig.recv().await.is_some() {
if let Ok((c, r)) = crossterm::terminal::size() {
crate::ctl::try_resize(&session_owned, c, r).await;
}
}
});
let r = pump_warm_pty_stdio(sock, stdin_rx, pending).await;
winch.abort();
Some(r)
} else {
crate::ui::trace(&format!("filament: reusing warm link to '{peer}' for one-shot pty"));
Some(pump_warm_pty_one_shot(sock, stdin_rx, pending).await)
}
}
pub async fn pty_cmd(server: &str, peer: &str, relay: bool, cmd: Vec<String>) -> Result<()> {
let one_shot = cmd.join(" ");
use std::io::IsTerminal;
let interactive = std::io::stdin().is_terminal() && std::io::stdout().is_terminal();
let session_id = crate::fresh_secret();
let term = std::env::var("TERM").ok().filter(|s| !s.is_empty()).unwrap_or_else(|| "xterm-256color".into());
let mut raw: Option<RawGuard> = None;
let mut stdin_rx = spawn_stdin_reader();
let mut pending: Option<Vec<u8>> = None;
let mut ever_connected = false;
let mut last_up = std::time::Instant::now();
let mut backoff = Duration::from_millis(300);
let mut role: &'static str = "init";
let mut warm_ended = false;
#[cfg(unix)]
if !relay && (interactive || !one_shot.is_empty()) {
match try_warm_pty(peer, &session_id, &term, &one_shot, interactive, &mut raw, &mut stdin_rx, &mut pending).await {
Some(Err(e)) => return Err(e),
Some(Ok(())) => {
warm_ended = true;
role = "reconnect";
}
None => {} }
}
loop {
let resume = ever_connected || warm_ended;
match pty_attach_once(server, peer, relay, role, &session_id, &term, &one_shot, interactive, resume, &mut raw, &mut stdin_rx, &mut pending).await {
Ok(PtyOutcome::Exited) => return Ok(()),
Ok(PtyOutcome::Dropped) => {
if !interactive {
return Ok(());
}
ever_connected = true;
last_up = std::time::Instant::now();
backoff = Duration::from_millis(300);
role = "reconnect";
eprint!("\r\n\x1b[2m[filament: link dropped, reconnecting...]\x1b[0m\r\n");
continue;
}
Err(e) => {
if !ever_connected && !warm_ended {
let connect_secs: u64 = std::env::var("FILAMENT_CONNECT_SECS")
.ok()
.and_then(|s| s.parse().ok())
.filter(|n| *n > 0)
.unwrap_or(45);
crate::ui::problem(
&format!("filament pty: can't reach '{peer}'"),
&format!(
"couldn't establish a link to '{peer}' in {connect_secs}s - it may be offline or unreachable from here."
),
&[
format!("check it's reachable: {}", crate::ui::paint(crate::ui::Tone::Brand, &format!("filament ping {peer}"))),
format!("diagnose the connect: {}", crate::ui::paint(crate::ui::Tone::Brand, &format!("filament doctor {peer}"))),
],
);
std::process::exit(1);
}
if last_up.elapsed() > Duration::from_secs(150) {
eprint!("\r\n\x1b[2m[filament: session expired, reconnect window passed]\x1b[0m\r\n");
return Ok(());
}
tokio::time::sleep(backoff).await;
backoff = (backoff * 2).min(Duration::from_secs(5));
role = "reconnect";
continue;
}
}
}
}
fn port_in_use_msg(lport: u16, peer: &str, rport: u16) -> String {
let suggested = lport.saturating_add(1);
format!(
"filament: local port {lport} is already in use, pick another (e.g. filament forward {suggested} {peer} {rport})"
)
}
#[cfg(unix)]
async fn bridge_streams(mut tcp: TcpStream, mut unix: tokio::net::UnixStream) -> std::io::Result<()> {
tokio::io::copy_bidirectional(&mut tcp, &mut unix).await.map(|_| ())
}
fn norm_device_name(s: &str) -> String {
s.chars().filter(|c| c.is_ascii_alphanumeric()).map(|c| c.to_ascii_lowercase()).collect()
}
fn forward_target_is_self(peer: &str) -> bool {
let p = norm_device_name(peer);
if p.is_empty() {
return false;
}
let dn = crate::display_name();
let host_part = dn.rsplit('@').next().unwrap_or("").to_string();
let hostname = std::fs::read_to_string("/etc/hostname").unwrap_or_default();
[dn.as_str(), host_part.as_str(), hostname.trim()]
.iter()
.any(|c| !c.is_empty() && norm_device_name(c) == p)
}
struct ForwardActivity {
active: std::sync::Arc<std::sync::atomic::AtomicU64>,
total: std::sync::Arc<std::sync::atomic::AtomicU64>,
first: std::sync::Arc<std::sync::atomic::AtomicBool>,
peer: String,
rport: u16,
}
impl ForwardActivity {
fn new(peer: &str, rport: u16) -> Self {
Self {
active: std::sync::Arc::new(std::sync::atomic::AtomicU64::new(0)),
total: std::sync::Arc::new(std::sync::atomic::AtomicU64::new(0)),
first: std::sync::Arc::new(std::sync::atomic::AtomicBool::new(false)),
peer: peer.to_string(),
rport,
}
}
fn line(&self) {
use std::sync::atomic::Ordering::Relaxed;
crate::ui::status(&format!(
"filament: forwarding to {}:{} - {} active, {} total",
self.peer,
self.rport,
self.active.load(Relaxed),
self.total.load(Relaxed)
));
}
fn begin(&self) -> ConnGuard {
use std::sync::atomic::Ordering::Relaxed;
self.total.fetch_add(1, Relaxed);
self.active.fetch_add(1, Relaxed);
if !self.first.swap(true, Relaxed) {
crate::ui::say(&format!(
"filament: first connection forwarded to {}:{} - the link is live",
self.peer, self.rport
));
}
self.line();
ConnGuard {
active: self.active.clone(),
total: self.total.clone(),
peer: self.peer.clone(),
rport: self.rport,
}
}
}
struct ConnGuard {
active: std::sync::Arc<std::sync::atomic::AtomicU64>,
total: std::sync::Arc<std::sync::atomic::AtomicU64>,
peer: String,
rport: u16,
}
impl Drop for ConnGuard {
fn drop(&mut self) {
use std::sync::atomic::Ordering::Relaxed;
self.active.fetch_sub(1, Relaxed);
crate::ui::status(&format!(
"filament: forwarding to {}:{} - {} active, {} total",
self.peer,
self.rport,
self.active.load(Relaxed),
self.total.load(Relaxed)
));
}
}
pub async fn mount_cmd(
server: &str,
peer: &str,
relay: bool,
root: &str,
) -> Result<crate::mount_proto::MountClient> {
let (t, rx, guard, _diag) = bring_up_to_known(server, peer, relay, "mount").await?;
guard.forget();
let mux = Mux::new(t.clone());
let _pump = tokio::spawn(pump_initiator(rx, mux.clone()));
let (sid, pipe_rx, caps) = open_mount_stream(&mux, root).await?;
let client = crate::mount_proto::MountClient::from_mux_v2(t.clone(), sid, pipe_rx, caps);
Ok(client)
}
pub async fn forward_cmd(server: &str, lport: u16, peer: &str, rport: u16, relay: bool) -> Result<()> {
if forward_target_is_self(peer) {
bail!(
"filament: '{peer}' is this device - a forward reaches a DIFFERENT machine. \
Whatever runs on this host's :{rport} is already here as 127.0.0.1:{rport}; \
point the forward at another peer (see `filament devices`)."
);
}
if !crate::devices_load().iter().any(|(n, _)| n.eq_ignore_ascii_case(peer)) {
bail!(
"filament: no known device named '{peer}'. Pair it first with `filament pair`, \
then `filament devices` shows who you can reach."
);
}
let listener = match TcpListener::bind(("127.0.0.1", lport)).await {
Ok(l) => l,
Err(e) if e.kind() == std::io::ErrorKind::AddrInUse => {
bail!("{}", port_in_use_msg(lport, peer, rport));
}
Err(e) if e.kind() == std::io::ErrorKind::PermissionDenied => {
bail!(
"filament: cannot bind 127.0.0.1:{lport}: permission denied. Local ports below \
1024 need root; pick a higher local port (e.g. `filament forward 8{lport:0>3} {peer} {rport}`) \
or run with sudo."
);
}
Err(e) => {
return Err(anyhow::Error::new(e).context(format!(
"filament: failed to bind 127.0.0.1:{lport} for forward to {peer}:127.0.0.1:{rport}"
)));
}
};
crate::ui::say(&format!("filament: forwarding 127.0.0.1:{lport} -> {peer}:127.0.0.1:{rport}"));
#[cfg(unix)]
let via_daemon = !relay && crate::ctl::daemon_present().await;
#[cfg(not(unix))]
let via_daemon = false;
let warm = via_daemon;
let mut cold_rx = if via_daemon {
match crate::ctl::try_ping(peer).await {
Some(facts) => {
let route = facts["route"].as_str().unwrap_or("link");
crate::ui::say(&format!(
"filament: ready - 127.0.0.1:{lport} -> {peer}:{rport} over the daemon's live {route} link (no extra presence on {peer})"
));
}
None => {
crate::ui::say(&format!(
"filament: listening on 127.0.0.1:{lport} -> {peer}:{rport} via the local daemon; no live link to {peer} yet - it opens on the first connection (check with `filament ping {peer}`)"
));
}
}
None
} else {
crate::ui::status(&format!("filament: bringing up the link to {peer} ..."));
let (tx, mut rx) = tokio::sync::watch::channel::<Option<Arc<Mux>>>(None);
{
let (server, peer_s) = (server.to_string(), peer.to_string());
tokio::spawn(async move { manage_cold_link(server, peer_s, relay, tx).await });
}
while rx.borrow().is_none() {
if rx.changed().await.is_err() {
bail!("filament: could not establish the link to {peer}; is it online? check with `filament ping {peer}` (or `filament devices`)");
}
}
crate::ui::say(&format!(
"filament: ready, listening on 127.0.0.1:{lport} -> {peer}:{rport} (connect to it to forward; run `filament up` here to avoid a separate presence on {peer})"
));
Some(rx)
};
let activity = ForwardActivity::new(peer, rport);
loop {
let sock = match listener.accept().await {
Ok((s, _)) => s,
Err(e) => {
crate::ui::status(&format!("filament: accept paused ({e}), retrying..."));
tokio::time::sleep(std::time::Duration::from_millis(200)).await;
continue;
}
};
let _ = sock.set_nodelay(true);
#[cfg(unix)]
if warm {
if let Some(usock) = crate::ctl::try_open(peer, rport).await {
let guard = activity.begin();
tokio::spawn(async move {
let _guard = guard; let _ = bridge_streams(sock, usock).await;
});
continue;
}
}
if cold_rx.is_none() {
crate::ui::debug(&format!(
"filament: no warm link to {peer}, using a direct link for forwarding"
));
let (tx, rx) = tokio::sync::watch::channel::<Option<Arc<Mux>>>(None);
let (server, peer_s) = (server.to_string(), peer.to_string());
tokio::spawn(async move { manage_cold_link(server, peer_s, relay, tx).await });
cold_rx = Some(rx);
}
let rx = cold_rx.clone().unwrap();
let guard = activity.begin();
let peer_c = peer.to_string();
tokio::spawn(async move {
let _guard = guard; serve_cold_connection(rx, sock, rport, peer_c).await;
});
}
}
pub async fn proxy_cmd(server: &str, bind: &str, port: u16, http_port: u16, relay: bool) -> Result<()> {
let listener = match TcpListener::bind((bind, port)).await {
Ok(l) => l,
Err(e) if e.kind() == std::io::ErrorKind::AddrInUse => {
bail!("filament: {bind}:{port} is already in use; pick another with --port");
}
Err(e) => {
return Err(anyhow::Error::new(e).context(format!("filament: failed to bind {bind}:{port}")));
}
};
crate::ui::say(&format!("filament: SOCKS5 proxy on {bind}:{port} (no TUN, no sudo)"));
crate::ui::say(&format!(
" point apps here; {}.mesh rides the mesh, everything else connects directly",
"<peer>"
));
crate::ui::say(&format!(" e.g. curl --socks5-hostname {bind}:{port} http://<peer>.mesh:8080/"));
#[cfg(unix)]
if !crate::ctl::daemon_present().await {
crate::ui::say(&crate::ui::paint(
crate::ui::Tone::Dim,
" note: no local daemon; each .mesh connection brings up its own link. `filament up` makes them instant.",
));
}
let cold: Arc<Mutex<HashMap<String, tokio::sync::watch::Receiver<Option<Arc<Mux>>>>>> =
Arc::new(Mutex::new(HashMap::new()));
if http_port > 0 {
let http_listener = match TcpListener::bind((bind, http_port)).await {
Ok(l) => l,
Err(e) if e.kind() == std::io::ErrorKind::AddrInUse => {
bail!("filament: {bind}:{http_port} is already in use; pick another with --http-port");
}
Err(e) => {
return Err(anyhow::Error::new(e).context(format!("filament: failed to bind {bind}:{http_port}")));
}
};
crate::ui::say(&format!("filament: HTTP CONNECT proxy on {bind}:{http_port}"));
crate::ui::say(&format!(
" PAC file: http://127.0.0.1:{http_port}/proxy.pac"
));
crate::ui::say(&format!(
" e.g. curl -x http://127.0.0.1:{http_port} https://<peer>.mesh"
));
let cold_http = cold.clone();
let server_http = server.to_string();
tokio::spawn(async move {
loop {
let sock = match http_listener.accept().await {
Ok((s, _)) => s,
Err(e) => {
crate::ui::status(&format!("filament: HTTP accept paused ({e}), retrying..."));
tokio::time::sleep(std::time::Duration::from_millis(200)).await;
continue;
}
};
let _ = sock.set_nodelay(true);
let (server, cold) = (server_http.clone(), cold_http.clone());
tokio::spawn(async move {
if let Err(e) = handle_http(sock, &server, port, relay, cold).await {
crate::ui::debug(&format!("filament: HTTP proxy connection ended: {e}"));
}
});
}
});
}
loop {
let sock = match listener.accept().await {
Ok((s, _)) => s,
Err(e) => {
crate::ui::status(&format!("filament: accept paused ({e}), retrying..."));
tokio::time::sleep(std::time::Duration::from_millis(200)).await;
continue;
}
};
let _ = sock.set_nodelay(true);
let (server, cold) = (server.to_string(), cold.clone());
tokio::spawn(async move {
if let Err(e) = handle_socks(sock, &server, relay, cold).await {
crate::ui::debug(&format!("filament: proxy connection ended: {e}"));
}
});
}
}
async fn socks_reply(sock: &mut TcpStream, code: u8) -> std::io::Result<()> {
sock.write_all(&[0x05, code, 0x00, 0x01, 0, 0, 0, 0, 0, 0]).await
}
async fn handle_socks(
mut sock: TcpStream,
server: &str,
relay: bool,
cold: Arc<Mutex<HashMap<String, tokio::sync::watch::Receiver<Option<Arc<Mux>>>>>>,
) -> Result<()> {
let mut greet = [0u8; 2];
sock.read_exact(&mut greet).await?;
if greet[0] != 0x05 {
bail!("not a SOCKS5 client");
}
let mut methods = vec![0u8; greet[1] as usize];
sock.read_exact(&mut methods).await?;
sock.write_all(&[0x05, 0x00]).await?;
let mut req = [0u8; 4];
sock.read_exact(&mut req).await?;
if req[0] != 0x05 {
bail!("bad SOCKS5 request");
}
let host = match req[3] {
0x01 => {
let mut a = [0u8; 4];
sock.read_exact(&mut a).await?;
std::net::Ipv4Addr::from(a).to_string()
}
0x04 => {
let mut a = [0u8; 16];
sock.read_exact(&mut a).await?;
std::net::Ipv6Addr::from(a).to_string()
}
0x03 => {
let mut len = [0u8; 1];
sock.read_exact(&mut len).await?;
let mut d = vec![0u8; len[0] as usize];
sock.read_exact(&mut d).await?;
String::from_utf8_lossy(&d).into_owned()
}
_ => {
socks_reply(&mut sock, 0x08).await?; return Ok(());
}
};
let mut pb = [0u8; 2];
sock.read_exact(&mut pb).await?;
let dport = u16::from_be_bytes(pb);
if req[1] != 0x01 {
socks_reply(&mut sock, 0x07).await?; return Ok(());
}
match host.strip_suffix(".mesh") {
Some(peer) => {
let peer = peer.to_string();
#[cfg(unix)]
if crate::ctl::daemon_present().await {
if let Some(usock) = crate::ctl::try_open(&peer, dport).await {
socks_reply(&mut sock, 0x00).await?;
return bridge_streams(sock, usock).await.map_err(Into::into);
}
if let Some(usock) = crate::ctl::try_dial(&peer, dport).await {
socks_reply(&mut sock, 0x00).await?;
return bridge_streams(sock, usock).await.map_err(Into::into);
}
}
let rx = {
let mut map = cold.lock().await;
if let Some(rx) = map.get(&peer) {
rx.clone()
} else {
let (tx, rx) = tokio::sync::watch::channel::<Option<Arc<Mux>>>(None);
let (s, pr) = (server.to_string(), peer.clone());
tokio::spawn(async move { manage_cold_link(s, pr, relay, tx).await });
map.insert(peer.clone(), rx.clone());
rx
}
};
socks_reply(&mut sock, 0x00).await?;
serve_cold_connection(rx, sock, dport, peer).await;
Ok(())
}
None => {
match TcpStream::connect((host.as_str(), dport)).await {
Ok(mut up) => {
let _ = up.set_nodelay(true);
socks_reply(&mut sock, 0x00).await?;
let _ = tokio::io::copy_bidirectional(&mut sock, &mut up).await;
Ok(())
}
Err(_) => {
socks_reply(&mut sock, 0x05).await?; Ok(())
}
}
}
}
}
async fn handle_http(
mut sock: TcpStream,
server: &str,
socks_port: u16,
relay: bool,
cold: Arc<Mutex<HashMap<String, tokio::sync::watch::Receiver<Option<Arc<Mux>>>>>>,
) -> Result<()> {
let mut buf = Vec::new();
let mut tmp = [0u8; 1];
loop {
sock.read_exact(&mut tmp).await?;
buf.push(tmp[0]);
if buf.len() >= 4 && &buf[buf.len()-4..] == b"\r\n\r\n" {
break;
}
if buf.len() > 8192 {
bail!("HTTP request too large");
}
}
let request = String::from_utf8_lossy(&buf);
let first_line = request.lines().next().unwrap_or("");
let mut parts = first_line.split_whitespace();
let method = parts.next().unwrap_or("");
let path = parts.next().unwrap_or("");
if method.eq_ignore_ascii_case("CONNECT") {
let host_port = path;
let (host, dport) = if let Some(colon) = host_port.rfind(':') {
let h = &host_port[..colon];
let p: u16 = host_port[colon+1..].parse().unwrap_or(0);
(h.to_string(), p)
} else {
(host_port.to_string(), 80)
};
match host.strip_suffix(".mesh") {
Some(peer) => {
let peer = peer.to_string();
#[cfg(unix)]
if crate::ctl::daemon_present().await {
if let Some(usock) = crate::ctl::try_open(&peer, dport).await {
let _ = sock.write_all(b"HTTP/1.1 200 Connection Established\r\n\r\n").await;
return bridge_streams(sock, usock).await.map_err(Into::into);
}
if let Some(usock) = crate::ctl::try_dial(&peer, dport).await {
let _ = sock.write_all(b"HTTP/1.1 200 Connection Established\r\n\r\n").await;
return bridge_streams(sock, usock).await.map_err(Into::into);
}
}
let rx = {
let mut map = cold.lock().await;
if let Some(rx) = map.get(&peer) {
rx.clone()
} else {
let (tx, rx) = tokio::sync::watch::channel::<Option<Arc<Mux>>>(None);
let (s, pr) = (server.to_string(), peer.clone());
tokio::spawn(async move { manage_cold_link(s, pr, relay, tx).await });
map.insert(peer.clone(), rx.clone());
rx
}
};
let _ = sock.write_all(b"HTTP/1.1 200 Connection Established\r\n\r\n").await;
serve_cold_connection(rx, sock, dport, peer).await;
Ok(())
}
None => {
match TcpStream::connect((host.as_str(), dport)).await {
Ok(mut up) => {
let _ = up.set_nodelay(true);
let _ = sock.write_all(b"HTTP/1.1 200 Connection Established\r\n\r\n").await;
let _ = tokio::io::copy_bidirectional(&mut sock, &mut up).await;
Ok(())
}
Err(_) => {
let _ = sock.write_all(b"HTTP/1.1 502 Bad Gateway\r\n\r\n").await;
Ok(())
}
}
}
}
} else if path == "/proxy.pac" || path == "/wpad.dat" {
let pac = format!(
r#"function FindProxyForURL(url, host) {{
if (dnsDomainIs(host, ".mesh") || shExpMatch(host, "*.mesh")) {{
return "SOCKS5 127.0.0.1:{socks_port}; DIRECT";
}}
return "DIRECT";
}}
"#
);
let response = format!(
"HTTP/1.1 200 OK\r\n\
Content-Type: application/x-ns-proxy-autoconfig\r\n\
Content-Length: {}\r\n\
Connection: close\r\n\
\r\n\
{}",
pac.len(),
pac
);
let _ = sock.write_all(response.as_bytes()).await;
Ok(())
} else {
let response = "HTTP/1.1 404 Not Found\r\n\
Content-Length: 0\r\n\
Connection: close\r\n\
\r\n";
let _ = sock.write_all(response.as_bytes()).await;
Ok(())
}
}
async fn serve_cold_connection(
mut rx: tokio::sync::watch::Receiver<Option<Arc<Mux>>>,
sock: TcpStream,
rport: u16,
peer: String,
) {
let mut tries = 0u32;
loop {
let mux = loop {
let cur = rx.borrow_and_update().clone();
match cur {
Some(m) if m.transport().is_alive() => break m,
_ => {
if rx.changed().await.is_err() {
return; }
}
}
};
match open_stream(&mux, rport).await {
Ok((sid, rx_pipe)) => {
serve_stream(mux, sid, sock, rx_pipe, true, None).await;
return;
}
Err(_) => {
tries += 1;
if tries >= 3 {
crate::ui::debug(&format!(
"filament: dropping a connection to {peer} (link recovering); the forward stays up"
));
return; }
}
}
}
}
async fn manage_cold_link(
server: String,
peer: String,
relay: bool,
tx: tokio::sync::watch::Sender<Option<Arc<Mux>>>,
) {
let mut backoff_ms = 500u64;
let mut had_link = false; loop {
let mux = match bring_up_to_known(&server, &peer, relay, "init").await {
Ok((t, rx, guard, mut diag)) => {
guard.forget();
diag.up("tunnel", "datachannel-or-direct");
let m = Mux::new(t);
tokio::spawn(pump_initiator(rx, m.clone()));
m
}
Err(e) => {
crate::ui::status(&format!("filament: reaching {peer} failed ({e}), retrying..."));
tokio::time::sleep(std::time::Duration::from_millis(backoff_ms)).await;
backoff_ms = (backoff_ms * 2).min(8000);
continue;
}
};
backoff_ms = 500;
if had_link {
crate::ui::say(&format!("filament: link to {peer} recovered"));
}
had_link = true;
if tx.send(Some(mux.clone())).is_err() {
return; }
loop {
tokio::time::sleep(std::time::Duration::from_secs(1)).await;
if tx.is_closed() {
return;
}
if !mux.transport().is_alive() {
crate::ui::status(&format!("filament: link to {peer} lost, reconnecting..."));
let _ = tx.send(None);
break;
}
}
}
}
struct BootstrapInfo {
hostkeys: Vec<String>,
user: Option<String>,
sshd: Option<bool>,
}
async fn shell_bootstrap(server: &str, peer: &str, relay: bool, ssh_port: u16) -> Result<BootstrapInfo> {
let pubkey = crate::sshkeys::ensure_managed_key()?;
let connect_secs: u64 = std::env::var("FILAMENT_SSH_CONNECT_SECS")
.ok()
.and_then(|s| s.parse().ok())
.filter(|n| *n > 0)
.unwrap_or(45);
let (t, mut rx, guard, mut diag) = match tokio::time::timeout(
std::time::Duration::from_secs(connect_secs),
bring_up_to_known(server, peer, relay, "bootstrap"),
)
.await
{
Ok(inner) => inner?,
Err(_) => {
crate::ui::problem(
&format!("filament ssh: can't reach '{peer}'"),
&format!(
"couldn't establish a link to '{peer}' in {connect_secs}s - it may be offline or unreachable from here."
),
&[
format!("check it's reachable: {}", crate::ui::paint(crate::ui::Tone::Brand, &format!("filament ping {peer}"))),
format!("diagnose the connect: {}", crate::ui::paint(crate::ui::Tone::Brand, &format!("filament doctor {peer}"))),
],
);
std::process::exit(1);
}
};
diag.up("tunnel", "datachannel-or-direct");
t.send_control(&json!({ "type": "shell-bootstrap", "v": 1, "pubkey": pubkey, "ssh_port": ssh_port })).await?;
let deadline = tokio::time::Instant::now() + std::time::Duration::from_secs(20);
let verdict: Result<BootstrapInfo> = loop {
let remaining = deadline.saturating_duration_since(tokio::time::Instant::now());
if remaining.is_zero() {
break Err(anyhow!(
"shell bootstrap timed out, is '{peer}' running `filament up` with shell access granted?"
));
}
match tokio::time::timeout(remaining, rx.recv()).await {
Ok(Some(Ev::Control(_pid, v))) => match v["type"].as_str() {
Some("shell-bootstrap-ack") => {
let hostkeys: Vec<String> = v["hostkeys"]
.as_array()
.map(|a| a.iter().filter_map(|k| k.as_str().map(String::from)).collect())
.unwrap_or_default();
let user = v["user"].as_str().map(String::from);
let sshd = v["sshd"].as_bool();
break Ok(BootstrapInfo { hostkeys, user, sshd });
}
Some("shell-bootstrap-deny") => {
let why = v["reason"].as_str().unwrap_or("shell capability not granted");
break Err(anyhow!(
"shell refused by '{peer}': {why}. Run `filament grant <this-device> shell` on '{peer}'."
));
}
_ => continue,
},
Ok(Some(_)) => continue, Ok(None) => break Err(anyhow!("channel closed before shell bootstrap completed")),
Err(_) => continue, }
};
drop(t);
guard.close().await;
verdict
}
async fn bootstrap_key(server: &str, peer: &str, relay: bool, ssh_port: u16) -> Result<BootstrapInfo> {
#[cfg(unix)]
if !relay {
let pubkey = crate::sshkeys::ensure_managed_key()?;
if let Some(v) = crate::ctl::try_bootstrap(peer, &pubkey, ssh_port).await {
let hostkeys: Vec<String> = v["hostkeys"]
.as_array()
.map(|a| a.iter().filter_map(|k| k.as_str().map(String::from)).collect())
.unwrap_or_default();
if !hostkeys.is_empty() {
crate::ui::trace(&format!(
"filament: reusing warm link to '{peer}' for ssh bootstrap (no establish)"
));
return Ok(BootstrapInfo {
hostkeys,
user: v["user"].as_str().map(String::from),
sshd: v["sshd"].as_bool(),
});
}
}
}
shell_bootstrap(server, peer, relay, ssh_port).await
}
fn resolve_login(remote_user: Option<String>) -> String {
std::env::var("FILAMENT_SSH_USER")
.ok()
.or(remote_user)
.or_else(|| std::env::var("USER").ok())
.unwrap_or_else(|| "root".into())
}
async fn run_ssh(
server: &str,
peer: &str,
relay: bool,
host: &str,
login: &str,
rport: u16,
extra: &[String],
revive: bool,
) -> Result<i32> {
#[cfg(not(target_os = "linux"))]
let _ = revive;
#[cfg(target_os = "linux")]
if std::env::var("FILAMENT_NO_L3_SSH").as_deref() != Ok("1") {
if let Some((mesh_host, addr)) = l3_mesh_addr(peer, rport) {
if probe_sshd(addr, std::time::Duration::from_millis(600)) {
crate::ui::debug(&format!(
"ssh over the L3 overlay ({mesh_host}) - survives link repairs"
));
let code = spawn_ssh_direct(login, &mesh_host, extra)?;
if code != 255 {
return Ok(code);
}
crate::ui::say("filament: L3 ssh failed, falling back to the tunnel");
} else if revive {
crate::ui::say(&format!(
"filament: L3 overlay to '{peer}' down, falling back to the tunnel; reviving in background"
));
let peer = peer.to_string();
tokio::spawn(async move {
revive_l3(&peer, rport).await;
});
}
}
}
spawn_ssh(server, peer, relay, host, login, rport, extra)
}
#[cfg(target_os = "linux")]
fn spawn_ssh_direct(login: &str, mesh_host: &str, extra: &[String]) -> Result<i32> {
let key = crate::sshkeys::managed_key_path();
let kh = crate::sshkeys::known_hosts_path();
let dest_token = format!("{login}@{mesh_host}");
let mut cmd = std::process::Command::new("ssh");
cmd.arg("-o").arg(format!("IdentityFile={}", key.display()))
.arg("-o").arg("IdentitiesOnly=yes")
.arg("-o").arg(format!("UserKnownHostsFile={}", kh.display()))
.arg("-o").arg("GlobalKnownHostsFile=/dev/null")
.arg("-o").arg("StrictHostKeyChecking=accept-new")
.arg("-o").arg("ConnectTimeout=10")
.arg("-o").arg("ServerAliveInterval=15")
.arg("-o").arg("ServerAliveCountMax=4");
let mut split = extra.len();
for (i, a) in extra.iter().enumerate() {
if !a.starts_with('-') {
split = i;
break;
}
}
for a in &extra[..split] {
cmd.arg(a);
}
cmd.arg(&dest_token);
for a in &extra[split..] {
cmd.arg(a);
}
Ok(cmd.status()?.code().unwrap_or(1))
}
fn spawn_ssh(
server: &str,
peer: &str,
relay: bool,
host: &str,
login: &str,
rport: u16,
extra: &[String],
) -> Result<i32> {
let exe = std::env::current_exe()?;
let exe = exe.to_string_lossy();
let mut proxy = format!("{exe} --server {server}");
if relay {
proxy.push_str(" --relay");
}
proxy.push_str(&format!(" netcat {peer} {rport}"));
let key = crate::sshkeys::managed_key_path();
let kh = crate::sshkeys::known_hosts_path();
let dest_token = format!("{login}@{host}");
let mut cmd = std::process::Command::new("ssh");
cmd.arg("-o").arg(format!("ProxyCommand={proxy}"))
.arg("-o").arg(format!("IdentityFile={}", key.display()))
.arg("-o").arg("IdentitiesOnly=yes")
.arg("-o").arg(format!("UserKnownHostsFile={}", kh.display()))
.arg("-o").arg("GlobalKnownHostsFile=/dev/null")
.arg("-o").arg("StrictHostKeyChecking=accept-new")
.arg("-o").arg("ConnectTimeout=25")
.arg("-o").arg("ServerAliveInterval=15")
.arg("-o").arg("ServerAliveCountMax=3");
let mut split = extra.len();
for (i, a) in extra.iter().enumerate() {
if !a.starts_with('-') {
split = i;
break;
}
}
for a in &extra[..split] {
cmd.arg(a);
}
cmd.arg(&dest_token);
for a in &extra[split..] {
cmd.arg(a);
}
Ok(cmd.status()?.code().unwrap_or(1))
}
#[cfg(target_os = "linux")]
fn l3_mesh_addr(peer: &str, port: u16) -> Option<(String, std::net::SocketAddr)> {
use std::net::ToSocketAddrs;
let name = format!("{}.mesh", crate::l3::sanitize_host(peer));
let addr = (name.as_str(), port).to_socket_addrs().ok()?.next()?;
Some((name, addr))
}
#[cfg(target_os = "linux")]
fn probe_sshd(addr: std::net::SocketAddr, to: std::time::Duration) -> bool {
std::net::TcpStream::connect_timeout(&addr, to).is_ok()
}
#[cfg(target_os = "linux")]
async fn revive_l3(peer: &str, rport: u16) {
let _ = tokio::time::timeout(std::time::Duration::from_secs(3), crate::ctl::try_open(peer, rport)).await;
}
pub(crate) struct PeerSshInfo {
pub login: String,
pub host: String,
pub rport: u16,
pub key_path: std::path::PathBuf,
pub known_hosts_path: std::path::PathBuf,
pub took_fast_path: bool,
}
pub(crate) async fn ensure_peer_bootstrap(server: &str, peer: &str, relay: bool) -> Result<PeerSshInfo> {
let peer = peer.strip_suffix(".mesh").unwrap_or(peer);
let _host = format!("filament-{peer}");
let rport: u16 =
std::env::var("FILAMENT_SSH_PORT").ok().and_then(|s| s.parse().ok()).unwrap_or(22);
ensure_peer_bootstrap_port(server, peer, relay, rport).await
}
pub(crate) async fn ensure_peer_bootstrap_port(server: &str, peer: &str, relay: bool, rport: u16) -> Result<PeerSshInfo> {
let peer = peer.strip_suffix(".mesh").unwrap_or(peer);
let host = format!("filament-{peer}");
let cached = if crate::sshkeys::host_pinned(&host) {
crate::sshkeys::bootstrap_cache_get(peer)
} else {
None
};
let (login, took_fast_path) = match cached {
Some(cached_user) => (resolve_login(cached_user), true),
None => {
let info = bootstrap_key(server, peer, relay, rport).await?;
ensure_sshd(peer, rport, info.sshd).await;
crate::sshkeys::pin_host_keys(&host, &info.hostkeys)?;
crate::sshkeys::bootstrap_cache_put(peer, info.user.as_deref());
(resolve_login(info.user), false)
}
};
Ok(PeerSshInfo {
login,
host,
rport,
key_path: crate::sshkeys::managed_key_path(),
known_hosts_path: crate::sshkeys::known_hosts_path(),
took_fast_path,
})
}
pub(crate) async fn rebootstrap_peer(server: &str, peer: &str, relay: bool) -> Result<PeerSshInfo> {
let peer = peer.strip_suffix(".mesh").unwrap_or(peer);
let host = format!("filament-{peer}");
let rport: u16 =
std::env::var("FILAMENT_SSH_PORT").ok().and_then(|s| s.parse().ok()).unwrap_or(22);
crate::sshkeys::bootstrap_cache_clear(peer);
let info = shell_bootstrap(server, peer, relay, rport).await?;
ensure_sshd(peer, rport, info.sshd).await;
crate::sshkeys::pin_host_keys(&host, &info.hostkeys)?;
crate::sshkeys::bootstrap_cache_put(peer, info.user.as_deref());
Ok(PeerSshInfo {
login: resolve_login(info.user),
host,
rport,
key_path: crate::sshkeys::managed_key_path(),
known_hosts_path: crate::sshkeys::known_hosts_path(),
took_fast_path: false,
})
}
pub(crate) fn ssh_transport_args(info: &PeerSshInfo, server: &str, peer: &str, relay: bool) -> Vec<String> {
let mut args = vec![
"-o".into(), format!("IdentityFile={}", info.key_path.display()),
"-o".into(), "IdentitiesOnly=yes".into(),
"-o".into(), format!("UserKnownHostsFile={}", info.known_hosts_path.display()),
"-o".into(), "GlobalKnownHostsFile=/dev/null".into(),
"-o".into(), "StrictHostKeyChecking=accept-new".into(),
"-o".into(), "ConnectTimeout=10".into(),
"-o".into(), "ServerAliveInterval=15".into(),
"-o".into(), "ServerAliveCountMax=4".into(),
];
let exe = std::env::current_exe().unwrap();
let exe = exe.to_string_lossy();
let mut proxy = format!("{exe} --server {server}");
if relay {
proxy.push_str(" --relay");
}
proxy.push_str(&format!(" netcat {peer} {}", info.rport));
args.push("-o".into());
args.push(format!("ProxyCommand={proxy}"));
args
}
#[cfg(not(target_os = "linux"))]
pub(crate) fn l3_dest(_info: &PeerSshInfo) -> Option<String> {
None
}
#[cfg(target_os = "linux")]
pub(crate) fn l3_dest(info: &PeerSshInfo) -> Option<String> {
let peer = info.host.strip_prefix("filament-").unwrap_or(&info.host);
let (mesh_host, addr) = l3_mesh_addr(peer, info.rport)?;
for attempt in 0..3 {
if probe_sshd(addr, std::time::Duration::from_millis(1000)) {
return Some(format!("{}@{mesh_host}", info.login));
}
if attempt < 2 {
std::thread::sleep(std::time::Duration::from_millis(500));
}
}
None
}
pub async fn ssh_cmd(server: &str, peer: &str, extra: &[String], relay: bool) -> Result<()> {
let peer = peer.strip_suffix(".mesh").unwrap_or(peer);
let info = ensure_peer_bootstrap(server, peer, relay).await?;
let code = run_ssh(server, peer, relay, &info.host, &info.login, info.rport, extra, true).await?;
if code == 255 && info.took_fast_path {
crate::ui::say(&format!("filament: re-authenticating with '{peer}'..."));
let retry = rebootstrap_peer(server, peer, relay).await?;
let code = run_ssh(server, peer, relay, &retry.host, &retry.login, retry.rport, extra, false).await?;
ssh_failed_hint(peer, code);
std::process::exit(code);
}
ssh_failed_hint(peer, code);
std::process::exit(code);
}
fn ssh_failed_hint(peer: &str, code: i32) {
if code == 255 {
crate::ui::say(&crate::ui::paint(
crate::ui::Tone::Dim,
&format!(" tip: `filament {peer}` opens a filament shell instead (no sshd needed, survives link repairs)"),
));
}
}
async fn ensure_sshd(peer: &str, rport: u16, reported: Option<bool>) {
let sshd = match reported {
Some(b) => Some(b),
#[cfg(unix)]
None => probe_sshd_warm(peer, rport).await,
#[cfg(not(unix))]
None => None,
};
if sshd != Some(false) {
return;
}
crate::ui::problem(
&format!("filament ssh: no sshd on '{peer}'"),
&format!(
"'{peer}' is reachable, but nothing is listening on localhost:{rport} for ssh. \
(sshd may be bound to a non-localhost address like the mesh ULA — \
`filament ssh` connects to localhost:{rport} on the peer.)",
),
&[
format!("start an sshd on '{peer}' listening on localhost (or all interfaces)"),
format!("set {} to a different port", crate::ui::paint(crate::ui::Tone::Brand, "FILAMENT_SSH_PORT")),
format!(
"use {} for a shell that needs no sshd",
crate::ui::paint(crate::ui::Tone::Brand, &format!("filament pty {peer}"))
),
],
);
std::process::exit(1);
}
#[cfg(unix)]
async fn probe_sshd_warm(peer: &str, rport: u16) -> Option<bool> {
use tokio::io::AsyncReadExt;
let mut s = crate::ctl::try_open(peer, rport).await?;
let mut buf = [0u8; 8];
match tokio::time::timeout(std::time::Duration::from_secs(3), s.read(&mut buf)).await {
Ok(Ok(0)) => Some(false), Ok(Ok(_)) => Some(true), Ok(Err(_)) => Some(false), Err(_) => None, }
}
#[cfg(test)]
mod h1_tests {
use super::*;
use async_trait::async_trait;
use std::sync::Mutex as StdMutex;
struct MockTransport {
controls: StdMutex<Vec<Value>>,
}
impl MockTransport {
fn new() -> Arc<Self> {
Arc::new(MockTransport { controls: StdMutex::new(Vec::new()) })
}
}
#[async_trait]
impl Transport for MockTransport {
async fn send_control(&self, msg: &Value) -> Result<()> {
self.controls.lock().unwrap().push(msg.clone());
Ok(())
}
async fn send_frame(&self, _sid: u32, _offset: u64, _payload: &[u8]) -> Result<()> {
Ok(())
}
async fn flush(&self) -> Result<()> {
Ok(())
}
fn max_payload(&self) -> usize {
1024
}
fn is_dead(&self) -> bool {
false
}
fn as_any(&self) -> &dyn std::any::Any {
self
}
}
fn open_msg(sid: u32) -> Value {
json!({ "type": "l2-open", "sid": sid, "host": "127.0.0.1", "rport": 9 })
}
#[tokio::test]
async fn pty_open_close_leaves_maps_empty() {
let start = LIVE_PTYS.load(Ordering::SeqCst);
let mux = Mux::new(MockTransport::new());
let n = 5u32;
for i in 0..n {
let sid = L2_SID_BASE | (1000 + i);
let guard = PtyGuard::try_acquire().expect("slot free");
let _rx = mux.register_stream(sid).await;
let (tx, _rrx) = mpsc::unbounded_channel::<(u16, u16)>();
mux.register_resizer(sid, tx).await;
assert_eq!(mux.live_streams().await, 1);
assert_eq!(mux.resizers.lock().await.len(), 1);
mux.on_close(sid, None).await;
drop(guard); assert_eq!(mux.live_streams().await, 0, "stream not freed on l2-close");
assert_eq!(mux.resizers.lock().await.len(), 0, "resizer leaked on l2-close");
}
for i in 0..n {
let sid = L2_SID_BASE | (2000 + i);
let guard = PtyGuard::try_acquire().expect("slot free");
let _rx = mux.register_stream(sid).await;
let (tx, _rrx) = mpsc::unbounded_channel::<(u16, u16)>();
mux.register_resizer(sid, tx).await;
mux.drop_pty(sid).await; drop(guard);
assert_eq!(mux.live_streams().await, 0, "stream not freed on drop_pty");
assert_eq!(mux.resizers.lock().await.len(), 0, "resizer leaked on drop_pty");
}
let mut guards = Vec::new();
for i in 0..n {
let sid = L2_SID_BASE | (3000 + i);
guards.push(PtyGuard::try_acquire().expect("slot free"));
let _rx = mux.register_stream(sid).await;
let (tx, _rrx) = mpsc::unbounded_channel::<(u16, u16)>();
mux.register_resizer(sid, tx).await;
}
assert_eq!(mux.live_streams().await, n as usize);
mux.shutdown_all().await;
drop(guards);
assert_eq!(mux.live_streams().await, 0, "streams leaked past shutdown_all");
assert_eq!(mux.resizers.lock().await.len(), 0, "resizers leaked past shutdown_all");
assert_eq!(LIVE_PTYS.load(Ordering::SeqCst), start, "global PTY count must return to baseline");
}
#[tokio::test]
async fn per_link_stream_cap_refuses_over_limit() {
let mux = Mux::new(MockTransport::new());
for i in 0..MAX_STREAMS_PER_LINK as u32 {
let sid = L2_SID_BASE | (i + 1);
match mux.accept_control(&open_msg(sid), true, false).await {
OpenVerdict::Accept { .. } => {}
other => panic!("expected Accept under cap, got {:?}", std::mem::discriminant(&other)),
}
}
assert_eq!(mux.live_streams().await, MAX_STREAMS_PER_LINK);
let over = L2_SID_BASE | 9999;
match mux.accept_control(&open_msg(over), true, false).await {
OpenVerdict::Deny { sid, err } => {
assert_eq!(sid, over);
assert_eq!(err, "too many streams");
}
other => panic!("expected Deny over cap, got {:?}", std::mem::discriminant(&other)),
}
assert_eq!(mux.live_streams().await, MAX_STREAMS_PER_LINK, "over-cap open must not register");
assert!(!mux.accepted.lock().await.contains_key(&over));
}
#[tokio::test]
async fn global_pty_cap_is_enforced() {
let mut held = Vec::new();
while LIVE_PTYS.load(Ordering::SeqCst) < MAX_PTYS_GLOBAL {
match PtyGuard::try_acquire() {
Some(g) => held.push(g),
None => break,
}
}
assert_eq!(LIVE_PTYS.load(Ordering::SeqCst), MAX_PTYS_GLOBAL);
assert!(PtyGuard::try_acquire().is_none(), "must refuse at global cap");
let before = held.len();
drop(held);
assert!(LIVE_PTYS.load(Ordering::SeqCst) <= MAX_PTYS_GLOBAL - before.min(1));
}
struct CapTransport {
frames: StdMutex<Vec<(u32, Vec<u8>)>>,
controls: StdMutex<Vec<Value>>,
}
impl CapTransport {
fn new() -> Arc<Self> {
Arc::new(CapTransport { frames: StdMutex::new(Vec::new()), controls: StdMutex::new(Vec::new()) })
}
fn bytes_for(&self, sid: u32) -> Vec<u8> {
self.frames
.lock()
.unwrap()
.iter()
.filter(|(s, _)| *s == sid)
.flat_map(|(_, b)| b.clone())
.collect()
}
}
#[async_trait]
impl Transport for CapTransport {
async fn send_control(&self, msg: &Value) -> Result<()> {
self.controls.lock().unwrap().push(msg.clone());
Ok(())
}
async fn send_frame(&self, sid: u32, _offset: u64, payload: &[u8]) -> Result<()> {
self.frames.lock().unwrap().push((sid, payload.to_vec()));
Ok(())
}
async fn flush(&self) -> Result<()> {
Ok(())
}
fn max_payload(&self) -> usize {
1024
}
fn is_dead(&self) -> bool {
false
}
fn as_any(&self) -> &dyn std::any::Any {
self
}
}
struct KillableTransport {
alive: std::sync::atomic::AtomicBool,
}
impl KillableTransport {
fn new() -> Arc<Self> {
Arc::new(KillableTransport { alive: std::sync::atomic::AtomicBool::new(true) })
}
fn kill(&self) {
self.alive.store(false, Ordering::SeqCst);
}
}
#[async_trait]
impl Transport for KillableTransport {
async fn send_control(&self, _msg: &Value) -> Result<()> {
Ok(())
}
async fn send_frame(&self, _sid: u32, _offset: u64, _payload: &[u8]) -> Result<()> {
Ok(())
}
async fn flush(&self) -> Result<()> {
Ok(())
}
fn max_payload(&self) -> usize {
1024
}
fn is_alive(&self) -> bool {
self.alive.load(Ordering::SeqCst)
}
fn is_dead(&self) -> bool {
!self.alive.load(Ordering::SeqCst)
}
fn as_any(&self) -> &dyn std::any::Any {
self
}
}
#[tokio::test]
async fn serve_stream_tears_down_on_transport_death_with_idle_client() {
let t = KillableTransport::new();
let mux = Mux::new(t.clone());
let sid = L2_SID_BASE | 7;
let rx = mux.register(sid).await;
let (client, server_side) = tokio::io::duplex(1024);
let bridge = tokio::spawn(serve_stream(mux.clone(), sid, server_side, rx, true, None));
tokio::time::sleep(Duration::from_millis(50)).await; t.kill(); let r = tokio::time::timeout(Duration::from_secs(4), bridge).await;
assert!(r.is_ok(), "serve_stream deadlocked on an idle client after transport death");
drop(client);
}
#[test]
fn ring_buffer_is_bounded_and_keeps_newest() {
let mut ring = VecDeque::new();
push_ring(&mut ring, b"hello", 8);
assert_eq!(ring.iter().copied().collect::<Vec<u8>>(), b"hello");
push_ring(&mut ring, b"world!!", 8);
assert_eq!(ring.len(), 8);
assert_eq!(ring.iter().copied().collect::<Vec<u8>>(), b"oworld!!");
}
#[test]
fn ring_buffer_oversized_write_keeps_tail() {
let mut ring = VecDeque::new();
push_ring(&mut ring, b"0123456789", 4);
assert_eq!(ring.iter().copied().collect::<Vec<u8>>(), b"6789");
}
#[tokio::test]
async fn session_survives_detach_and_reattaches_with_replay() {
if !std::path::Path::new("/bin/cat").exists() {
return; }
let sessions = PtySessions::new();
let ta = CapTransport::new();
let sid_a = L2_SID_BASE | 1;
let guard = PtyGuard::try_acquire().expect("slot");
let sess = spawn_pty_session(
sessions.clone(),
"sess-x".to_string(),
ta.clone(),
sid_a,
80,
24,
"xterm-256color",
vec!["/bin/cat".to_string()],
guard,
)
.await
.expect("spawn");
sess.feed_input(b"before-drop\n".to_vec());
tokio::time::sleep(Duration::from_millis(200)).await;
let a_out = String::from_utf8_lossy(&ta.bytes_for(sid_a)).to_string();
assert!(a_out.contains("before-drop"), "link A never saw the echo: {a_out:?}");
sess.detach();
tokio::time::sleep(Duration::from_millis(50)).await;
assert!(!sess.is_dead(), "detach must NOT kill the session");
let tb = CapTransport::new();
let sid_b = L2_SID_BASE | 2;
let live = sessions.get_live("sess-x").await.expect("session still live for reattach");
live.attach(tb.clone(), sid_b);
tokio::time::sleep(Duration::from_millis(150)).await;
let b_replay = String::from_utf8_lossy(&tb.bytes_for(sid_b)).to_string();
assert!(b_replay.contains("before-drop"), "reattach did not replay buffered output: {b_replay:?}");
live.feed_input(b"after-reconnect\n".to_vec());
tokio::time::sleep(Duration::from_millis(200)).await;
let b_out = String::from_utf8_lossy(&tb.bytes_for(sid_b)).to_string();
assert!(b_out.contains("after-reconnect"), "post-reattach input did not echo: {b_out:?}");
live.end();
tokio::time::sleep(Duration::from_millis(150)).await;
assert!(sessions.get_live("sess-x").await.is_none(), "ended session must leave the store");
}
#[test]
fn port_in_use_msg_names_port_and_suggests_forward() {
let msg = port_in_use_msg(8080, "laptop", 22);
assert!(msg.contains("8080"), "message must name the conflicting port: {msg}");
assert!(msg.contains("filament forward"), "message must suggest a filament forward retry: {msg}");
assert!(msg.contains("8081"), "suggested port should be lport+1: {msg}");
assert!(msg.contains("laptop"), "message should reference the peer: {msg}");
assert!(msg.contains("22"), "message should reference the rport: {msg}");
}
#[test]
fn port_in_use_msg_saturates_at_u16_max() {
let msg = port_in_use_msg(u16::MAX, "laptop", 22);
assert!(msg.contains(&format!("{}", u16::MAX)), "must name the conflicting port: {msg}");
assert!(!msg.contains("0 "), "saturating add must not wrap to 0: {msg}");
}
#[cfg(unix)]
#[tokio::test]
async fn warm_reuse_zombie_link_self_heals_instead_of_hanging() {
let t = CapTransport::new();
let mux = Mux::new(t.clone());
let verdict =
open_stream_verified(&mux, 22, std::time::Duration::from_millis(50)).await;
assert!(verdict.is_err(), "a link that delivers no inbound frame must be a zombie Err");
assert_eq!(mux.live_streams().await, 0, "zombie stream must be removed, not leaked");
let ctrls = t.controls.lock().unwrap();
assert!(
ctrls.iter().any(|c| c["type"] == "l2-open"),
"must have attempted the open"
);
assert!(
ctrls.iter().any(|c| c["type"] == "l2-close"),
"must send l2-close so the peer reaps its half of the zombie stream"
);
}
#[cfg(unix)]
#[tokio::test]
async fn warm_reuse_healthy_link_passes_and_replays_first_frame() {
let t = CapTransport::new();
let mux = Mux::new(t.clone());
let mux2 = mux.clone();
let h = tokio::spawn(async move {
open_stream_verified(&mux2, 22, std::time::Duration::from_secs(5)).await
});
let sid = loop {
if let Some(&sid) = mux.streams.lock().await.keys().next() {
break sid;
}
tokio::task::yield_now().await;
};
mux.on_frame(sid, Bytes::from_static(b"BANNER")).await;
let (got_sid, first, rx) = h.await.expect("task panicked").expect("healthy link must be Ok");
assert_eq!(got_sid, sid);
assert_eq!(first, Some(Bytes::from_static(b"BANNER")), "first frame must be preserved for replay");
let (mut client, srv) = tokio::io::duplex(1024);
let mux3 = mux.clone();
let s = tokio::spawn(async move { serve_verified_stream(mux3, sid, srv, first, rx).await });
let mut buf = [0u8; 6];
client.read_exact(&mut buf).await.expect("replayed frame must reach the client");
assert_eq!(&buf, b"BANNER", "the verified first frame must be replayed verbatim");
mux.on_frame(sid, Bytes::new()).await; drop(client); s.await.expect("serve task panicked");
}
#[cfg(unix)]
#[tokio::test]
async fn warm_reuse_client_first_link_passes_on_open_ack() {
let t = CapTransport::new();
let mux = Mux::new(t.clone());
let mux2 = mux.clone();
let h = tokio::spawn(async move {
open_stream_verified(&mux2, 80, std::time::Duration::from_secs(5)).await
});
let sid = loop {
if let Some(&sid) = mux.streams.lock().await.keys().next() {
break sid;
}
tokio::task::yield_now().await;
};
mux.on_open_ack(sid).await;
let (got_sid, first, rx) =
h.await.expect("task panicked").expect("open-ack must prove the link live");
assert_eq!(got_sid, sid);
assert_eq!(first, Some(Bytes::new()), "the ack is an empty liveness marker");
let (mut client, srv) = tokio::io::duplex(1024);
let mux3 = mux.clone();
let s = tokio::spawn(async move { serve_verified_stream(mux3, sid, srv, first, rx).await });
mux.on_frame(sid, Bytes::from_static(b"HTTP/1.1 200 OK")).await;
let mut buf = [0u8; 15];
client
.read_exact(&mut buf)
.await
.expect("real data must reach the client after the empty ack");
assert_eq!(&buf, b"HTTP/1.1 200 OK", "the empty ack must not corrupt the stream");
mux.on_frame(sid, Bytes::new()).await;
drop(client);
s.await.expect("serve task panicked");
}
}