mirror of
https://github.com/nestriness/nestri.git
synced 2026-09-19 17:25:19 +03:00
Four findings from review, all of them real. The relay's directory was mounted on the tree a session's shares live in. A fresh tmpfs there hides every directory the image prepared underneath it: the install, the user state, the work directory, and the mount point the log share is attached to from fstab. A box would have come up with a socket and without any of the places its workload looks for its files, and the exact-path check could not notice, because what fstab mounts is a directory inside that tree rather than the tree itself. It moves to /run, which is where a runtime socket belongs, is a tmpfs already, and has nothing else mounted inside it. It was also owned by this process and closed to everyone else, which stopped the workload traversing it to reach the relay at all. The directory is now readable and searchable, and still writable by nothing but this process, which is what makes the socket in it unreplaceable; the socket itself is what the workload is allowed to connect to. The permission belongs on the socket rather than on the path. The address served to a reader was built once at startup and served forever, so a reader that polls for a better one could only ever get the first. An endpoint does not know all of its own addresses when it binds: the first is the one that works on the same network and fails from anywhere else. It is now rebuilt per read, which is what makes polling for it worth doing. And the address was taken from whoever held a path in a directory the workload can write. Workload code could unlink the socket a service was listening on, bind its own, and every read afterwards would hand the client an address of its choosing -- a session given to somebody else rather than a session that fails. The peer's credentials are now checked before a byte is read, from the kernel rather than from anything the peer says about itself, and an address served by the workload's own user is refused and said loudly. That check is only worth something while the workload has a user of its own, so the image grows one. Two users, and they must stay two: one runs the services that ship in the image, the other is who a workload runs as. Sharing one does not weaken the check, it makes every session fail it. A workload running as root is every user at once and cannot be told apart from anything; the check stands down there and says so at boot instead, because refusing root would refuse whatever legitimately serves the address as well. Also bumps tinyvec by a patch release. It does not build on this toolchain -- `vec` resolves to the module and not the macro -- which made every crate that depends on an endpoint, including this one, unbuildable. Pre-existing and nothing to do with this change; the lockfile said the same version before it.
315 lines
11 KiB
Rust
315 lines
11 KiB
Rust
mod dgram;
|
|
mod ipc_listener;
|
|
mod keyframe;
|
|
mod screenshot;
|
|
mod session;
|
|
mod ticket;
|
|
|
|
use std::path::PathBuf;
|
|
use std::sync::Arc;
|
|
|
|
use anyhow::Result;
|
|
use clap::Parser;
|
|
use iroh::endpoint::presets;
|
|
|
|
use crate::session::SessionManager;
|
|
use crate::ticket::NestriTicket;
|
|
use nesprotocol::ALPN;
|
|
|
|
#[derive(Parser, Debug)]
|
|
#[command(name = "neshub")]
|
|
struct Args {
|
|
/// Relay mode: default, none, or a custom relay URL
|
|
#[arg(long, env = "NESTRI_RELAY", default_value = "default")]
|
|
relay: String,
|
|
|
|
/// Path for the video IPC socket (nescapture → neshub)
|
|
#[arg(
|
|
long,
|
|
env = "NESTRI_VIDEO_IPC",
|
|
default_value = "/tmp/nestri-video.sock"
|
|
)]
|
|
video_ipc: PathBuf,
|
|
|
|
/// Path for the audio IPC socket (neswire → neshub)
|
|
#[arg(
|
|
long,
|
|
env = "NESTRI_AUDIO_IPC",
|
|
default_value = "/tmp/nestri-audio.sock"
|
|
)]
|
|
audio_ipc: PathBuf,
|
|
|
|
/// Path for the input IPC socket (neshub → nescope).
|
|
#[arg(
|
|
long,
|
|
env = "NESTRI_INPUT_IPC",
|
|
default_value = "/tmp/nestri-input.sock"
|
|
)]
|
|
input_ipc: PathBuf,
|
|
|
|
/// Path for the stats IPC socket (nescapture → neshub stats).
|
|
#[arg(
|
|
long,
|
|
env = "NESTRI_STATS_IPC",
|
|
default_value = "/tmp/nestri-stats.sock"
|
|
)]
|
|
stats_ipc: PathBuf,
|
|
|
|
/// Socket the ticket is served on. neshub listens; nesinit dials and
|
|
/// carries the ticket to the host, because the person who needs it is
|
|
/// outside this VM and stdout here is a log file inside one.
|
|
#[arg(
|
|
long,
|
|
env = "NESTRI_TICKET_IPC",
|
|
default_value = "/tmp/nestri-ticket.sock"
|
|
)]
|
|
ticket_ipc: PathBuf,
|
|
|
|
/// Audio channels (from neswire config): 2 = stereo, 6 = 5.1, 8 = 7.1
|
|
#[arg(long, env = "NESTRI_AUDIO_CHANNELS", default_value_t = 2)]
|
|
audio_channels: u32,
|
|
|
|
/// Audio bitrate per channel in kbps
|
|
#[arg(long, env = "NESTRI_AUDIO_BITRATE", default_value_t = 64)]
|
|
audio_bitrate_per_channel: u32,
|
|
|
|
/// Socket nescope sends screenshots on. neshub listens; nescope dials out.
|
|
#[arg(
|
|
long,
|
|
env = "NESTRI_SCREENSHOT_IPC",
|
|
default_value = "/tmp/nestri-screenshot.sock"
|
|
)]
|
|
screenshot_ipc: PathBuf,
|
|
}
|
|
|
|
#[tokio::main]
|
|
async fn main() -> Result<()> {
|
|
tracing_subscriber::fmt()
|
|
.with_env_filter(
|
|
tracing_subscriber::EnvFilter::builder()
|
|
.with_default_directive(tracing_subscriber::filter::LevelFilter::INFO.into())
|
|
.from_env_lossy(),
|
|
)
|
|
.init();
|
|
|
|
let args = Args::parse();
|
|
|
|
let mut builder = iroh::Endpoint::builder(presets::N0)
|
|
.alpns(vec![ALPN.to_vec()])
|
|
.transport_config(crate::dgram::media_transport_config());
|
|
|
|
match args.relay.as_str() {
|
|
"default" | "" => {
|
|
builder = builder.relay_mode(iroh::endpoint::RelayMode::Default);
|
|
tracing::info!("using default n0-computer relays");
|
|
}
|
|
"none" | "off" | "disabled" => {
|
|
builder = builder.relay_mode(iroh::endpoint::RelayMode::Disabled);
|
|
tracing::info!("relays disabled (direct connections only)");
|
|
}
|
|
url => {
|
|
let relay_url: iroh::RelayUrl = url.parse()?;
|
|
let relay_map = iroh::RelayMap::empty();
|
|
relay_map.insert(
|
|
relay_url.clone(),
|
|
Arc::new(iroh::RelayConfig::new(relay_url, None)),
|
|
);
|
|
builder = builder.relay_mode(iroh::endpoint::RelayMode::Custom(relay_map));
|
|
tracing::info!("using custom relay: {url}");
|
|
}
|
|
}
|
|
|
|
let endpoint = builder.bind().await?;
|
|
let endpoint_addr = endpoint.addr();
|
|
let ep_id = endpoint_addr.id;
|
|
tracing::info!("endpoint online: {}", ep_id.fmt_short());
|
|
|
|
// Input broadcast channel: input reader -> input IPC listener -> nescope
|
|
let (input_broadcast_tx, _) = tokio::sync::broadcast::channel::<Vec<u8>>(256);
|
|
|
|
// Cursor channel: IPC listener (read side) -> client sessions -> desktop-app
|
|
let (cursor_tx, mut cursor_rx) = tokio::sync::mpsc::unbounded_channel::<Vec<u8>>();
|
|
|
|
// Nescope stats channel: IPC listener -> client sessions
|
|
let (nescope_stats_tx, mut nescope_stats_rx) =
|
|
tokio::sync::mpsc::unbounded_channel::<Vec<u8>>();
|
|
|
|
let session_manager = Arc::new(SessionManager::new());
|
|
|
|
// IDR / encode settings command channel: input reader → nescapture
|
|
let (cmd_tx, mut cmd_rx) = tokio::sync::mpsc::unbounded_channel::<Vec<u8>>();
|
|
{
|
|
tokio::spawn(async move {
|
|
let cmd_path = std::path::PathBuf::from("/tmp/nescapture-cmd.sock");
|
|
while let Some(bytes) = cmd_rx.recv().await {
|
|
if let Ok(sock) = std::os::unix::net::UnixDatagram::unbound() {
|
|
if sock.send_to(&bytes, &cmd_path).is_err() {
|
|
tracing::warn!("nescapture cmd send failed at {}", cmd_path.display());
|
|
}
|
|
}
|
|
}
|
|
});
|
|
}
|
|
|
|
// Spawn cursor relay
|
|
{
|
|
let mgr = session_manager.clone();
|
|
tokio::spawn(async move {
|
|
while let Some(data) = cursor_rx.recv().await {
|
|
mgr.broadcast_cursor(data).await;
|
|
}
|
|
});
|
|
}
|
|
|
|
// Spawn nescope stats relay
|
|
{
|
|
let mgr = session_manager.clone();
|
|
tokio::spawn(async move {
|
|
while let Some(data) = nescope_stats_rx.recv().await {
|
|
mgr.broadcast_stats(data).await;
|
|
}
|
|
});
|
|
}
|
|
|
|
// Spawn periodic hub stats
|
|
{
|
|
let mgr = session_manager.clone();
|
|
let audio_channels = args.audio_channels as u8;
|
|
// The configured target is worth saying once, here, where it is a fact
|
|
// about this hub's arguments. It is deliberately not what gets reported
|
|
// in the stats below -- see `SessionManager::audio_bitrate_kbps`.
|
|
tracing::info!(
|
|
"audio configured for {}ch at {}kbps/channel; stats report measured ingest",
|
|
args.audio_channels,
|
|
args.audio_bitrate_per_channel
|
|
);
|
|
tokio::spawn(async move {
|
|
let mut interval = tokio::time::interval(std::time::Duration::from_secs(1));
|
|
loop {
|
|
interval.tick().await;
|
|
let clients = mgr.client_count().await as u8;
|
|
let bitrate = mgr.video_bitrate_bps();
|
|
let audio_kbps = mgr.audio_bitrate_kbps();
|
|
let relay_ms = mgr.relay_ms();
|
|
let mut buf = Vec::with_capacity(15);
|
|
nesprotocol::stats::encode_hub_stats(
|
|
&mut buf,
|
|
clients,
|
|
bitrate,
|
|
relay_ms,
|
|
audio_kbps,
|
|
audio_channels,
|
|
);
|
|
mgr.broadcast_stats(buf).await;
|
|
}
|
|
});
|
|
}
|
|
|
|
// ── Accept mode: generate ticket, wait for desktop-app to connect ─────
|
|
let stream_name = ticket::generate_stream_name();
|
|
// For the log line below only. What a reader of the socket gets is built
|
|
// per read from the endpoint itself, because the addresses this can be
|
|
// reached at are not all known yet.
|
|
let ticket = NestriTicket::new(endpoint_addr, stream_name.clone());
|
|
|
|
tracing::info!("╔═══════════════╗");
|
|
tracing::info!("║ NESTRI TICKET ║");
|
|
tracing::info!("╚═══════════════╝");
|
|
tracing::info!("{ticket}\n");
|
|
|
|
// Spawn IPC listeners
|
|
|
|
let video_ipc = args.video_ipc.clone();
|
|
let audio_ipc = args.audio_ipc.clone();
|
|
let input_ipc = args.input_ipc.clone();
|
|
let stats_ipc = args.stats_ipc.clone();
|
|
let stats_tx_clone = nescope_stats_tx.clone();
|
|
tokio::spawn({
|
|
let mgr = session_manager.clone();
|
|
async move { ipc_listener::run_video_listener(video_ipc, mgr).await }
|
|
});
|
|
tokio::spawn({
|
|
let mgr = session_manager.clone();
|
|
async move { ipc_listener::run_audio_listener(audio_ipc, mgr).await }
|
|
});
|
|
let input_ipc_tx = input_broadcast_tx.clone();
|
|
let cursor_ipc_tx = cursor_tx.clone();
|
|
let ns_tx = nescope_stats_tx.clone();
|
|
tokio::spawn({
|
|
async move {
|
|
ipc_listener::run_input_ipc_listener(input_ipc, input_ipc_tx, cursor_ipc_tx, ns_tx)
|
|
.await
|
|
}
|
|
});
|
|
tokio::spawn({
|
|
let stx = stats_tx_clone.clone();
|
|
async move { ipc_listener::run_stats_ipc_listener(stats_ipc, stx).await }
|
|
});
|
|
|
|
let ticket_ipc = args.ticket_ipc.clone();
|
|
tokio::spawn({
|
|
// The endpoint rather than a ticket made from it: the addresses it can
|
|
// be reached at are not all known yet, and whoever reads this socket
|
|
// re-reads it so that a better one can replace the first.
|
|
let endpoint = endpoint.clone();
|
|
let stream_name = stream_name.clone();
|
|
async move { ipc_listener::run_ticket_ipc_listener(ticket_ipc, endpoint, stream_name).await }
|
|
});
|
|
|
|
// Accept loop
|
|
let ep = endpoint.clone();
|
|
let mgr = session_manager.clone();
|
|
let accept_handle = tokio::spawn(async move {
|
|
while let Some(incoming) = ep.accept().await {
|
|
match incoming.await {
|
|
Ok(conn) => {
|
|
let remote_id = conn.remote_id();
|
|
tracing::info!(remote = %remote_id.fmt_short(), "client connected");
|
|
let session = session::ClientSession::new(
|
|
conn.clone(),
|
|
input_broadcast_tx.clone(),
|
|
session_manager.relay_ms_atomic(),
|
|
cmd_tx.clone(),
|
|
);
|
|
mgr.add_session(remote_id, session).await;
|
|
let mgr_clone = mgr.clone();
|
|
let conn_clone = conn.clone();
|
|
tokio::spawn(async move {
|
|
conn_clone.closed().await;
|
|
mgr_clone.remove_session(&remote_id).await;
|
|
});
|
|
}
|
|
Err(e) => {
|
|
tracing::warn!("incoming connection failed: {e}");
|
|
}
|
|
}
|
|
}
|
|
tracing::info!("accept loop exited");
|
|
});
|
|
|
|
// Not wired to anything today. Kept because the capture works and "show me
|
|
// what the guest is displaying" is the first question when a payload
|
|
// renders black.
|
|
let _screenshots = match screenshot::listen(&args.screenshot_ipc) {
|
|
Ok(connection) => Some(connection),
|
|
Err(e) => {
|
|
tracing::warn!("screenshots unavailable: {e:#}");
|
|
None
|
|
}
|
|
};
|
|
|
|
tracing::info!("neshub running, ctrl+c to stop");
|
|
tokio::signal::ctrl_c().await?;
|
|
tracing::info!("shutting down..");
|
|
endpoint.close().await;
|
|
accept_handle.abort();
|
|
|
|
let _ = std::fs::remove_file(&args.video_ipc);
|
|
let _ = std::fs::remove_file(&args.audio_ipc);
|
|
let _ = std::fs::remove_file(&args.input_ipc);
|
|
let _ = std::fs::remove_file(&args.stats_ipc);
|
|
let _ = std::fs::remove_file("/tmp/nescapture-cmd.sock");
|
|
let _ = std::fs::remove_file(&args.ticket_ipc);
|
|
Ok(())
|
|
}
|