From 48d120aa4aef9b7dab576af13fa2290580b84865 Mon Sep 17 00:00:00 2001 From: Yaroslav Bolyukin Date: Thu, 18 Jun 2026 23:15:29 +0000 Subject: [PATCH] ci: use in-tree remowt-agent --- --- a/Cargo.lock +++ b/Cargo.lock @@ -722,9 +722,9 @@ [[package]] name = "camino" -version = "1.2.2" +version = "1.2.3" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "e629a66d692cb9ff1a1c664e41771b3dcaf961985a9774c0eb0bd1b51cf60a48" +checksum = "b4ce8d3bd5823c7504d3f579f13e7b2f3da252fcb938c594d5680ee508bf846f" dependencies = [ "serde_core", ] @@ -1899,9 +1899,9 @@ "fleet-shared", "hex", "hmac 0.13.0", - "pbkdf2 0.12.2", + "pbkdf2 0.13.0", "rand 0.10.1", - "sha2 0.10.9", + "sha2 0.11.0", "x25519-dalek", ] @@ -5011,6 +5011,7 @@ "nix-eval", "remowt-client", "remowt-endpoints", + "remowt-link-shared", "serde", "thiserror 2.0.18", "tokio", @@ -5031,13 +5032,13 @@ "iroh-base", "n0-watcher", "noq-udp", - "remowt-ui-prompt", "serde", "serde_json", "thiserror 2.0.18", "tokio", "tokio-util", "tracing", + "uuid", ] [[package]] @@ -5093,11 +5094,11 @@ "anyhow", "bifrostlink", "bifrostlink-macros", + "remowt-link-shared", "serde", "thiserror 2.0.18", "tokio", "tracing", - "zbus", ] [[package]] --- a/Cargo.toml +++ b/Cargo.toml @@ -15,11 +15,11 @@ remowt-fleet = { path = "./crates/remowt-fleet" } remowt-client = { version = "0.1.9", path = "remowt/crates/remowt-client" } -remowt-polkit-shared = { version = "0.1.9", path = "remowt/crates/polkit-shared" } +remowt-endpoints = { version = "0.1.9", path = "remowt/crates/remowt-endpoints" } remowt-link-shared = { version = "0.1.9", path = "remowt/crates/remowt-link-shared" } remowt-plugin = { version = "0.1.9", path = "remowt/crates/remowt-plugin" } +remowt-polkit-shared = { version = "0.1.9", path = "remowt/crates/polkit-shared" } remowt-ui-prompt = { version = "0.1.9", path = "remowt/crates/remowt-ui-prompt" } -remowt-endpoints = { version = "0.1.9", path = "remowt/crates/remowt-endpoints" } bifrostlink = "0.2.0" bifrostlink-macros = "0.2.0" @@ -27,16 +27,9 @@ iroh = { version = "1.0.0", features = ["unstable-custom-transports"] } iroh-base = "1.0.0" -noq-udp = { version = "1.0.0", default-features = false } n0-watcher = "1.0.0" +noq-udp = { version = "1.0.0", default-features = false } -uuid = { version = "1", features = ["v4"] } -russh = { version = "0.61.2", default-features = false, features = [ - "ring", - "flate2", - "rsa", -] } -russh-config = "0.58.0" age = { version = "0.11", features = ["plugin", "ssh"] } anyhow = "1.0" base64 = "0.22.1" @@ -64,13 +57,15 @@ opentelemetry-appender-tracing = "0.32.0" opentelemetry-otlp = { version = "0.32.0", features = ["grpc-tonic", "gzip-tonic", "http-json", "reqwest-rustls"] } opentelemetry_sdk = "0.32.0" -pbkdf2 = "0.12" +pbkdf2 = "0.13" peg = "0.8.5" pkg-config = "0.3.30" rand = "0.10.0" +russh = { version = "0.61.2", default-features = false, features = ["flate2", "ring", "rsa"] } +russh-config = "0.58.0" serde = { version = "1.0", features = ["derive"] } serde_json = "1.0" -sha2 = "0.10" +sha2 = "0.11" shlex = "2.0.1" tabled = "0.21.0" tempfile = "3.20" @@ -81,14 +76,15 @@ tracing = "0.1" tracing-indicatif = "0.3.13" tracing-opentelemetry = "0.33.0" +uuid = { version = "1", features = ["v4"] } # For fixed coloring needs to be updated to https://github.com/tokio-rs/tracing/pull/3484 +tokio-util = "0.7.11" tracing-subscriber = { version = "0.3.19", features = ["env-filter", "fmt"] } unicode_categories = "0.1.1" vte = { version = "0.15.0", features = ["ansi"] } x25519-dalek = { version = "2.0.1", features = ["getrandom"] } zbus = "5.16.0" zbus_polkit = "5.0.0" -tokio-util = "0.7.11" [profile.dev] panic = "abort" --- a/cmds/fleet/src/cmds/build_systems.rs +++ b/cmds/fleet/src/cmds/build_systems.rs @@ -5,8 +5,9 @@ use clap::Parser; use fleet_base::{ deploy::{DeployAction, deploy_task, upload_task}, - host::{Config, DeployKind, GenerationStorage}, + host::{Config, DeployKind}, opts::FleetOpts, + pins::GenerationStorage, }; use futures::{StreamExt as _, stream::FuturesUnordered}; use nix_eval::nix_go; @@ -145,8 +146,7 @@ } let remote_path = - match upload_task(&config, &host, GenerationStorage::Deployer, built).await - { + match upload_task(&host, GenerationStorage::Deployer, built).await { Ok(v) => v, Err(e) => { error!("upload failed: {e}"); --- a/cmds/fleet/src/cmds/rollback.rs +++ b/cmds/fleet/src/cmds/rollback.rs @@ -4,8 +4,9 @@ use clap::Parser; use fleet_base::{ deploy::{DeployAction, deploy_task, upload_task}, - host::{Config, ConfigHost, Generation, GenerationStorage}, + host::{Config, ConfigHost}, opts::FleetOpts, + pins::{Generation, GenerationStorage}, }; use tabled::Table; use tracing::{info, warn}; @@ -103,13 +104,8 @@ Table::new(&generations) ); }; - let remote_path = upload_task( - config, - &host, - generation.location, - generation.store_path.clone(), - ) - .await?; + let remote_path = + upload_task(&host, generation.location, generation.store_path.clone()).await?; deploy_task( action, --- a/cmds/fleet/src/extra_args.rs +++ /dev/null @@ -1,18 +0,0 @@ -use std::ffi::{OsStr, OsString}; - -use anyhow::{Result, anyhow}; - -pub fn parse_os(os: &OsStr) -> Result> { - Ok(shlex::bytes::split(os.as_encoded_bytes()) - .ok_or_else(|| anyhow!("invalid arguments"))? - .into_iter() - .map(|a| { - // Unpaired surrogates are not touched - unsafe { OsString::from_encoded_bytes_unchecked(a) } - }) - .collect()) -} -// pub fn parse(s: &str) -> Result> { -// let osstr = OsString::try_from(s)?; -// parse_os(&osstr) -// } --- a/cmds/fleet/src/main.rs +++ b/cmds/fleet/src/main.rs @@ -1,8 +1,6 @@ #![recursion_limit = "512"] pub(crate) mod cmds; -// pub(crate) mod command; -pub(crate) mod extra_args; use std::{process::ExitCode, sync::Arc}; --- a/cmds/generator-helper/src/main.rs +++ b/cmds/generator-helper/src/main.rs @@ -14,7 +14,7 @@ use clap::{Parser, ValueEnum}; use ed25519_dalek::SecretKey; use fleet_shared::SecretData; -use hmac::Mac as _; +use hmac::{KeyInit as _, Mac as _}; use rand::{ Rng as _, distr::{Alphanumeric, Distribution, SampleString, Uniform}, @@ -338,8 +338,8 @@ ); type HmacSha256 = hmac::Hmac; - let mut mac = ::new_from_slice(&salted) - .expect("HMAC accepts any key length"); + let mut mac = + HmacSha256::new_from_slice(&salted).expect("HMAC accepts any key length"); mac.update(b"Client Key"); let client_key = mac.finalize().into_bytes(); @@ -347,8 +347,8 @@ hasher.update(client_key); let stored_key = hasher.finalize(); - let mut mac = ::new_from_slice(&salted) - .expect("HMAC accepts any key length"); + let mut mac = + HmacSha256::new_from_slice(&salted).expect("HMAC accepts any key length"); mac.update(b"Server Key"); let server_key = mac.finalize().into_bytes(); --- a/cmds/remowt-plugin-fleet/Cargo.toml +++ b/cmds/remowt-plugin-fleet/Cargo.toml @@ -2,7 +2,7 @@ name = "remowt-plugin-fleet" description = "Remowt plugin exposing fleet's Nix endpoint to a running agent" version.workspace = true -edition = "2021" +edition.workspace = true [dependencies] remowt-fleet.workspace = true --- a/crates/fleet-base/src/deploy.rs +++ b/crates/fleet-base/src/deploy.rs @@ -11,7 +11,8 @@ use tokio::time::sleep; use tracing::{Instrument as _, error, info, info_span, warn}; -use crate::host::{Config, ConfigHost, DeployKind, Generation, GenerationStorage}; +use crate::host::{ConfigHost, DeployKind}; +use crate::pins::{Generation, GenerationStorage}; #[derive(ValueEnum, Clone, Copy)] pub enum DeployAction { @@ -260,13 +261,12 @@ { // It is ok, if there was no reboot - then timer might not be running. } - if action.should_schedule_rollback_run() { - if let Err(e) = elevated_systemd + if action.should_schedule_rollback_run() + && let Err(e) = elevated_systemd .stop("rollback-watchdog-run.timer".to_owned()) .await - { - error!("failed to disarm rollback run: {e}"); - } + { + error!("failed to disarm rollback run: {e}"); } } else if let Err(_e) = elevated_fs .rm_file(Utf8PathBuf::from("/etc/fleet_rollback_marker")) @@ -280,7 +280,6 @@ } pub async fn upload_task( - config: &Config, host: &ConfigHost, location: GenerationStorage, generation: Utf8PathBuf, --- a/crates/fleet-base/src/fleetdata.rs +++ b/crates/fleet-base/src/fleetdata.rs @@ -6,6 +6,7 @@ }, fmt, io::{self, Cursor}, + str::FromStr, sync::RwLock, }; @@ -92,8 +93,9 @@ #[serde(skip_serializing)] host_secrets: BTreeMap>, } -impl FleetData { - pub fn from_str(s: &str) -> anyhow::Result { +impl FromStr for FleetData { + type Err = anyhow::Error; + fn from_str(s: &str) -> anyhow::Result { let mut data: Self = nixlike::parse_str(s)?; if !data.host_secrets.is_empty() { info!("migrating host secrets into shared secrets structure"); @@ -294,15 +296,16 @@ /// Drop expired distributions fn prune_expired(&mut self, now: DateTime) { for ele in self.distributions_mut() { - if let Some(expires_at) = ele.secret.expires_at { - if expires_at < now { - ele.prune(format!("expired during check at {now}")); - } + if let Some(expires_at) = ele.secret.expires_at + && expires_at < now + { + ele.prune(format!("expired during check at {now}")); } } } /// Perform all pruning relevant to shared secrets /// Also see expected_owner_removed + #[allow(clippy::too_many_arguments)] pub fn prune_shared( &mut self, expected_owners: &BTreeSet, @@ -323,14 +326,15 @@ let mut to_add = expected_owners.difference(¤t_owners); if to_add.next().is_some() && unique && regenerate_on_owner_added { for dist in self.distributions_mut() { - dist.prune(format!( + dist.prune( "owners missing, can't add new distribution, regeneration preferred" - )); + .to_string(), + ); } return; } - for to_remove in current_owners.difference(&expected_owners) { + for to_remove in current_owners.difference(expected_owners) { self.entry(to_remove.clone()).remove( regenerate_on_owner_removed, "owner was removed from expected owners list, regenerate_on_owner_removed is set" @@ -362,9 +366,7 @@ ) -> Option { self.distributions() .enumerate() - .max_by(|(_, a), (_, b)| { - compare_dists(&a, &b, prefer_identities, include_pruned_owners) - }) + .max_by(|(_, a), (_, b)| compare_dists(a, b, prefer_identities, include_pruned_owners)) .map(|(p, _)| p) } /// Secret wants to be the same on all hosts, leave only one unpruned version of it @@ -410,10 +412,10 @@ filter_owner: Option<&SecretOwner>, ) { 'dist: for ele in self.distributions_mut() { - if let Some(filter_owner) = filter_owner { - if !ele.owners.contains(filter_owner) { - continue; - } + if let Some(filter_owner) = filter_owner + && !ele.owners.contains(filter_owner) + { + continue; // Note: secret still can have multiple owners even if it is host-owned // in this case we expect that all owners using the same generator, so we can prune distribution for all of them } @@ -442,10 +444,10 @@ filter_owner: Option<&SecretOwner>, ) { for ele in self.distributions_mut() { - if let Some(filter_owner) = filter_owner { - if !ele.owners.contains(filter_owner) { - continue; - } + if let Some(filter_owner) = filter_owner + && !ele.owners.contains(filter_owner) + { + continue; // Note: secret still can have multiple owners even if it is host-owned // in this case we expect that all owners using the same generator, so we can prune distribution for all of them } @@ -519,7 +521,7 @@ } } -struct OccupiedDistEntry<'d> { +pub struct OccupiedDistEntry<'d> { distributions: &'d mut FleetSecretDistributions, idx: usize, owners: BTreeSet, @@ -537,16 +539,16 @@ owners: self.owners, } } - fn set(self, secret: FleetSecretData, reason: String) -> Self { + pub fn set(self, secret: FleetSecretData, reason: String) -> Self { self.remove(false, reason).set(secret) } } -struct VacantDistEntry<'d> { +pub struct VacantDistEntry<'d> { distributions: &'d mut FleetSecretDistributions, owners: BTreeSet, } impl<'d> VacantDistEntry<'d> { - fn set(self, secret: FleetSecretData) -> OccupiedDistEntry<'d> { + pub fn set(self, secret: FleetSecretData) -> OccupiedDistEntry<'d> { let Self { distributions, owners, @@ -568,18 +570,18 @@ } } -enum DistEntry<'d> { +pub enum DistEntry<'d> { Vacant(VacantDistEntry<'d>), Occupied(OccupiedDistEntry<'d>), } impl DistEntry<'_> { - fn remove(self, whole_dist: bool, reason: String) -> Self { + pub fn remove(self, whole_dist: bool, reason: String) -> Self { match self { DistEntry::Vacant(_) => self, DistEntry::Occupied(o) => Self::Vacant(o.remove(whole_dist, reason)), } } - fn set(self, secret: FleetSecretData, reason: String) -> Self { + pub fn set(self, secret: FleetSecretData, reason: String) -> Self { Self::Occupied(match self { DistEntry::Vacant(e) => e.set(secret), DistEntry::Occupied(e) => e.set(secret, reason), @@ -683,7 +685,7 @@ } impl FleetSecrets { - pub fn keys(&self) -> btree_map::Keys { + pub fn keys(&self) -> btree_map::Keys<'_, String, FleetSecretDistributions> { self.0.keys() } @@ -713,9 +715,7 @@ } pub fn get_or_create(&mut self, secret: &str) -> &mut FleetSecretDistributions { - self.0 - .entry(secret.to_owned()) - .or_insert(FleetSecretDistributions::default()) + self.0.entry(secret.to_owned()).or_default() } pub fn contains(&self, secret: &str) -> bool { --- a/crates/fleet-base/src/host.rs +++ b/crates/fleet-base/src/host.rs @@ -13,18 +13,17 @@ use camino::{Utf8Path, Utf8PathBuf}; use chrono::{DateTime, Utc}; use fleet_shared::SecretData; -use nix_eval::{Store, Value, eval_store, nix_go, nix_go_json, util::assert_warn}; +use nix_eval::{Store, Value, drv::DrvGraph, eval_store, nix_go, nix_go_json, util::assert_warn}; use remowt_client::{AgentBundle, Remowt}; use remowt_endpoints::fs::FsClient; -use remowt_link_shared::Address; +use remowt_fleet::NixClient; +use remowt_link_shared::BifConfig; use remowt_ui_prompt::auto::AutoPrompter; use remowt_ui_prompt::bifrost::PromptEndpoints; use remowt_ui_prompt::{PrependSourcePrompter, Source}; -use tabled::Tabled; use tempfile::NamedTempFile; -use time::UtcDateTime; -use tokio::task::spawn_blocking; -use tracing::warn; +use tokio::{sync::OnceCell, task::spawn_blocking}; +use tracing::{info, warn}; use crate::fleetdata::{ FleetData, FleetSecretData, FleetSecretDistribution, FleetSecretPart, SecretOwner, @@ -101,7 +100,7 @@ groups: OnceLock>, // TODO: Both of those values are taken from host opts, there should be a cleaner way to specify it - deploy_kind: OnceLock, + deploy_kind: OnceCell, session_destination: OnceLock, legacy_ssh_store: OnceLock, @@ -114,7 +113,7 @@ pub local: bool, pub remowt: OnceLock, nix_store: OnceLock>, - nix_plugin: tokio::sync::OnceCell<()>, + nix_plugin: OnceCell<()>, } const NIX_PLUGIN_ID: u16 = 2; @@ -130,75 +129,16 @@ fn agent_bundle() -> Result { AgentBundle::from_dir(agents_dir()?) -} - -#[derive(Debug, Clone, Copy)] -pub enum GenerationStorage { - Deployer, - Machine, - Pusher, -} -impl GenerationStorage { - fn prefix(&self) -> &'static str { - match self { - GenerationStorage::Deployer => "deployer.", - GenerationStorage::Machine => "", - GenerationStorage::Pusher => "pusher.", - } - } } -#[derive(Tabled, Debug)] -pub struct Generation { - #[tabled(rename = "ID", format("{}", self.rollback_id()))] - pub id: u32, - #[tabled(rename = "Current")] - pub current: bool, - #[tabled(rename = "Created at")] - pub datetime: UtcDateTime, - #[tabled(format = "{:?}")] - pub store_path: Utf8PathBuf, - #[tabled(skip)] - pub location: GenerationStorage, -} -impl Generation { - pub fn rollback_id(&self) -> String { - format!("{}{}", self.location.prefix(), self.id) - } +enum RemoteDerivationMode { + CopySigned, + Copy, + Attic, + Cachix, } impl ConfigHost { - pub async fn list_generations(&self, profile: &str) -> Result> { - let plugin_id = self.ensure_nix_plugin().await?; - let nix = self - .remowt() - .await? - .plugin_endpoints::>(plugin_id); - let raw = nix - .list_generations(profile.to_owned()) - .await - .map_err(|e| anyhow!("{e:?}"))? - .map_err(|e| anyhow!("{e}"))?; - raw.into_iter() - .map(|g| { - let id: u32 = - g.id.try_into() - .with_context(|| format!("generation id {} doesn't fit in u32", g.id))?; - let datetime = UtcDateTime::from_unix_timestamp(g.creation_time_unix) - .with_context(|| { - format!("invalid generation timestamp {}", g.creation_time_unix) - })?; - Ok(Generation { - id, - current: g.current, - datetime, - store_path: g.store_path, - location: GenerationStorage::Machine, - }) - }) - .collect() - } - pub fn set_session_destination(&self, dest: String) { self.session_destination .set(dest) @@ -215,34 +155,31 @@ .expect("legacy ssh store is already set") } pub async fn deploy_kind(&self) -> Result { - if let Some(kind) = self.deploy_kind.get() { - return Ok(*kind); - } - let remowt = self.remowt().await?; - let fs = remowt.endpoints::>(); - let is_fleet_managed = match fs.file_exists(Utf8PathBuf::from("/etc/FLEET_HOST")).await { - Ok(v) => v, - Err(e) => { - bail!("failed to query remote system kind: {e}"); + self.deploy_kind.get_or_try_init(|| async { + let remowt = self.remowt().await?; + let fs = remowt.endpoints::>(); + let is_fleet_managed = match fs.file_exists(Utf8PathBuf::from("/etc/FLEET_HOST")).await { + Ok(v) => v, + Err(e) => { + bail!("failed to query remote system kind: {e}"); + } + }; + if !is_fleet_managed { + bail!( + "{}", + indoc::indoc! {" + host is not marked as managed by fleet + if you're not trying to lustrate/install system from scratch, + you should either + 1. manually create /etc/FLEET_HOST file on the target host, + 2. use ?deploy_kind=fleet host argument if you're upgrading from older version of fleet + 3. use ?deploy_kind=upgrade_to_fleet if you're upgrading from plain nixos to fleet-managed nixos + for installation use ?deploy_kind=nixos_install / ?deploy_kind=nixos_lustrate + "} + ); } - }; - if !is_fleet_managed { - bail!( - "{}", - indoc::indoc! {" - host is not marked as managed by fleet - if you're not trying to lustrate/install system from scratch, - you should either - 1. manually create /etc/FLEET_HOST file on the target host, - 2. use ?deploy_kind=fleet host argument if you're upgrading from older version of fleet - 3. use ?deploy_kind=upgrade_to_fleet if you're upgrading from plain nixos to fleet-managed nixos - for installation use ?deploy_kind=nixos_install / ?deploy_kind=nixos_lustrate - "} - ); - } - // TOCTOU is possible - let _ = self.deploy_kind.set(DeployKind::Fleet); - Ok(*self.deploy_kind.get().expect("deploy kind is just set")) + Ok(DeployKind::Fleet) + }).await.copied() } async fn connection(&self) -> Result { if let Some(conn) = self.remowt.get() { @@ -285,6 +222,11 @@ self.connection().await } + pub async fn nix_client(&self) -> Result> { + let remowt = self.remowt().await?; + let plugin_id = self.ensure_nix_plugin().await?; + Ok(remowt.plugin_endpoints(plugin_id)) + } pub fn ensure_nix_plugin(&self) -> Pin> + Send + '_>> { Box::pin(async { self.nix_plugin @@ -305,12 +247,6 @@ .run0_load_plugin_path(NIX_PLUGIN_ID, bin.as_str()) .await .context("failed to load the fleet nix plugin")?; - self.remowt() - .await? - .rpc() - .wait_for_connection_to(Address::Plugin(NIX_PLUGIN_ID)) - .await - .map_err(|_| anyhow!("failed to wait for plugin"))?; anyhow::Ok(()) }) .await?; @@ -407,6 +343,10 @@ } let sign: Pin> + Send>> = { let path = path.clone(); + let graph = DrvGraph::resolve(&eval_store(), &path) + .context("failed to resolve graph to be uploaded")?; + info!("signing {} paths", graph.nodes.len()); + Box::pin(async move { let local = self.config.local_host(); let plugin_id = local.ensure_nix_plugin().await?; @@ -564,8 +504,8 @@ local: true, remowt: OnceLock::new(), nix_store: OnceLock::new(), - nix_plugin: tokio::sync::OnceCell::new(), - deploy_kind: OnceLock::new(), + nix_plugin: OnceCell::new(), + deploy_kind: OnceCell::new(), session_destination: OnceLock::new(), legacy_ssh_store: OnceLock::new(), }) @@ -607,8 +547,8 @@ local: self.localhost == name, remowt: OnceLock::new(), nix_store: OnceLock::new(), - nix_plugin: tokio::sync::OnceCell::new(), - deploy_kind: OnceLock::new(), + nix_plugin: OnceCell::new(), + deploy_kind: OnceCell::new(), session_destination: OnceLock::new(), legacy_ssh_store: OnceLock::new(), }) --- a/crates/fleet-base/src/lib.rs +++ b/crates/fleet-base/src/lib.rs @@ -3,4 +3,5 @@ pub mod host; mod keys; pub mod opts; +pub mod pins; pub mod primops; --- /dev/null +++ b/crates/fleet-base/src/pins.rs @@ -0,0 +1,69 @@ +use anyhow::{Context as _, Result, anyhow}; +use camino::Utf8PathBuf; +use tabled::Tabled; +use time::UtcDateTime; + +use crate::host::ConfigHost; + +#[derive(Debug, Clone, Copy)] +pub enum GenerationStorage { + Deployer, + Machine, + Pusher, +} +impl GenerationStorage { + fn prefix(&self) -> &'static str { + match self { + GenerationStorage::Deployer => "deployer.", + GenerationStorage::Machine => "", + GenerationStorage::Pusher => "pusher.", + } + } +} + +#[derive(Tabled, Debug)] +pub struct Generation { + #[tabled(rename = "ID", format("{}", self.rollback_id()))] + pub id: u32, + #[tabled(rename = "Current")] + pub current: bool, + #[tabled(rename = "Created at")] + pub datetime: UtcDateTime, + #[tabled(format = "{:?}")] + pub store_path: Utf8PathBuf, + #[tabled(skip)] + pub location: GenerationStorage, +} +impl Generation { + pub fn rollback_id(&self) -> String { + format!("{}{}", self.location.prefix(), self.id) + } +} +impl ConfigHost { + pub async fn list_generations(&self, profile: &str) -> Result> { + let nix = self.nix_client().await?; + let raw = nix + .list_generations(profile.to_owned()) + .await + .map_err(|e| anyhow!("{e:?}"))? + .map_err(|e| anyhow!("{e}"))?; + raw.into_iter() + .map(|g| { + let id: u32 = + g.id.try_into() + .with_context(|| format!("generation id {} doesn't fit in u32", g.id))?; + let datetime = UtcDateTime::from_unix_timestamp(g.creation_time_unix) + .with_context(|| { + format!("invalid generation timestamp {}", g.creation_time_unix) + })?; + Ok(Generation { + id, + current: g.current, + datetime, + store_path: g.store_path, + location: GenerationStorage::Machine, + }) + }) + .collect() + } +} --- a/crates/nixlike/fuzz/Cargo.toml +++ b/crates/nixlike/fuzz/Cargo.toml @@ -4,7 +4,7 @@ version = "0.0.0" authors = ["Automatically generated"] publish = false -edition = "2024" +edition.workspace = true [package.metadata] cargo-fuzz = true --- a/crates/remowt-fleet/Cargo.toml +++ b/crates/remowt-fleet/Cargo.toml @@ -2,7 +2,7 @@ name = "remowt-fleet" description = "Remowt endpoints for fleet" version.workspace = true -edition = "2021" +edition.workspace = true [dependencies] bifrostlink.workspace = true @@ -18,3 +18,4 @@ tokio = { workspace = true, features = ["io-util", "net", "process", "rt"] } tracing.workspace = true uuid.workspace = true +remowt-link-shared.workspace = true --- a/crates/remowt-fleet/src/lib.rs +++ b/crates/remowt-fleet/src/lib.rs @@ -7,6 +7,7 @@ use nix_eval::eval_store; use remowt_client::Remowt; use remowt_endpoints::nix_daemon::NixDaemonClient; +use remowt_link_shared::iroh_tunnel::TunnelAddr; use serde::{Deserialize, Serialize}; use tokio::net::UnixListener; use tokio::task::spawn_blocking; @@ -98,8 +99,10 @@ continue; } }; - let sock_str = remote_sock.as_str().to_owned(); - match nix.serve_store(store.clone(), sock_str).await { + match nix + .serve_store(store.clone(), TunnelAddr::Unix(remote_sock)) + .await + { Ok(Ok(())) => {} Ok(Err(e)) => { error!("nix bridge: {e}"); --- a/flake.lock +++ b/flake.lock @@ -171,35 +171,6 @@ "type": "github" } }, - "remowt-agents": { - "inputs": { - "crane": [ - "crane" - ], - "flake-parts": [ - "flake-parts" - ], - "nixpkgs": [ - "nixpkgs" - ], - "rust-overlay": "rust-overlay", - "shelly": [ - "shelly" - ] - }, - "locked": { - "lastModified": 1781491763, - "ref": "refs/heads/trunk", - "rev": "8b18e8b236d428ed8fbf5818272d17f5f538c13c", - "revCount": 40, - "type": "git", - "url": "https://gerrit.delta.rocks/remowt" - }, - "original": { - "type": "git", - "url": "https://gerrit.delta.rocks/remowt" - } - }, "root": { "inputs": { "crane": "crane", @@ -207,33 +178,12 @@ "fleet-tf": "fleet-tf", "nix": "nix", "nixpkgs": "nixpkgs", - "remowt-agents": "remowt-agents", - "rust-overlay": "rust-overlay_2", + "rust-overlay": "rust-overlay", "shelly": "shelly", "treefmt-nix": "treefmt-nix" } }, "rust-overlay": { - "inputs": { - "nixpkgs": [ - "remowt-agents", - "nixpkgs" - ] - }, - "locked": { - "lastModified": 1769309768, - "owner": "oxalica", - "repo": "rust-overlay", - "rev": "140c9dc582cb73ada2d63a2180524fcaa744fad5", - "type": "github" - }, - "original": { - "owner": "oxalica", - "repo": "rust-overlay", - "type": "github" - } - }, - "rust-overlay_2": { "inputs": { "nixpkgs": [ "nixpkgs" --- a/flake.nix +++ b/flake.nix @@ -23,13 +23,6 @@ url = "github:numtide/treefmt-nix"; inputs.nixpkgs.follows = "nixpkgs"; }; - remowt-agents = { - url = "git+https://gerrit.delta.rocks/remowt"; - inputs.nixpkgs.follows = "nixpkgs"; - inputs.flake-parts.follows = "flake-parts"; - inputs.crane.follows = "crane"; - inputs.shelly.follows = "shelly"; - }; # DeterminateSystem's nix fork is controversial, but I don't mind it, # and it has lazy-trees support which is useful for fleet. nix = { --- a/pkgs/default.nix +++ b/pkgs/default.nix @@ -3,9 +3,17 @@ craneLib, inputs, }: +let + remowt-agents-bundle = callPackage ./remowt-agents-bundle.nix { inherit craneLib; }; +in { fleet = callPackage ./fleet.nix { inherit craneLib inputs; }; - remowt-plugin-fleet = callPackage ./remowt-plugin-fleet.nix { inherit craneLib inputs; }; fleet-install-secrets = callPackage ./fleet-install-secrets.nix { inherit craneLib; }; fleet-generator-helper = callPackage ./fleet-generator-helper.nix { inherit craneLib; }; + + inherit remowt-agents-bundle; + remowt-plugin-fleet = callPackage ./remowt-plugin-fleet.nix { + inherit craneLib inputs remowt-agents-bundle; + }; + remowt-ssh = callPackage ./remowt-ssh.nix { inherit craneLib remowt-agents-bundle; }; } --- a/pkgs/fleet.nix +++ b/pkgs/fleet.nix @@ -3,6 +3,7 @@ craneLib, installShellFiles, inputs, + remowt-agents-bundle, stdenv, pkg-config, @@ -26,7 +27,7 @@ cargoExtraArgs = "--locked -p ${pname}"; - REMOWT_AGENTS_DIR = "${inputs.remowt-agents.packages.${system}.remowt-agents}"; + REMOWT_AGENTS_DIR = "${remowt-agents-bundle}"; # TODO: built-in fleet prompter should be a prodash widget, or it should require # tty remowt prompter running on host machine idk. ROFI = "${rofi}/bin/rofi"; --- /dev/null +++ b/pkgs/remowt-agents-bundle.nix @@ -0,0 +1,73 @@ +{ + craneLib, + lib, + pkgs, + runCommandLocal, +}: +let + crateName = "remowt-agent"; + + buildFor = + { + target, + crossPkgs, + }: + let + cc = crossPkgs.stdenv.cc; + ccBin = "${cc}/bin/${cc.targetPrefix}"; + ut = builtins.replaceStrings [ "-" ] [ "_" ] target; + linkerEnv = "CARGO_TARGET_${lib.toUpper ut}_LINKER"; + in + craneLib.buildPackage ( + { + src = craneLib.cleanCargoSource ../.; + pname = "${crateName}-${target}"; + + cargoExtraArgs = "--locked -p ${crateName}"; + + CARGO_BUILD_TARGET = target; + CARGO_BUILD_RUSTFLAGS = "-C target-feature=+crt-static"; + + depsBuildBuild = [ cc ]; + doCheck = false; + } + // { + ${linkerEnv} = "${ccBin}cc"; + "CC_${ut}" = "${ccBin}cc"; + "CXX_${ut}" = "${ccBin}c++"; + "AR_${ut}" = "${ccBin}ar"; + } + ); + x86_64 = buildFor { + target = "x86_64-unknown-linux-musl"; + crossPkgs = pkgs; + }; + aarch64 = buildFor { + target = "aarch64-unknown-linux-musl"; + crossPkgs = pkgs.pkgsCross.aarch64-multiplatform-musl; + }; + armv7l = buildFor { + target = "armv7-unknown-linux-musleabihf"; + crossPkgs = pkgs.pkgsCross.armv7l-hf-multiplatform.pkgsMusl; + }; +in +runCommandLocal "remowt-agents-bundle" + { + passthru = { + perArch = { + inherit x86_64 aarch64 armv7l; + }; + }; + } + '' + mkdir -p $out + cp ${x86_64}/bin/remowt-agent $out/remowt-agent-x86_64 + cp ${aarch64}/bin/remowt-agent $out/remowt-agent-aarch64 + cp ${armv7l}/bin/remowt-agent $out/remowt-agent-armv7l + chmod +w $out/remowt-agent-* + + for arch in x86_64 aarch64 armv7l; do + hash=$(sha256sum "$out/remowt-agent-$arch" | cut -d' ' -f1) + printf '%s %s\n' "$arch" "$hash" >> $out/hashes + done + '' --- /dev/null +++ b/pkgs/remowt-ssh.nix @@ -0,0 +1,17 @@ +{ + craneLib, + remowt-agents-bundle, + rofi, +}: +let + pname = "remowt-ssh"; +in +craneLib.buildPackage { + inherit pname; + src = craneLib.cleanCargoSource ../.; + + cargoExtraArgs = "--locked -p ${pname}"; + + REMOWT_AGENTS_DIR = "${remowt-agents-bundle}"; + ROFI = "${rofi}/bin/rofi"; +} --- a/remowt/cmds/remowt-agent/Cargo.toml +++ b/remowt/cmds/remowt-agent/Cargo.toml @@ -2,7 +2,7 @@ name = "remowt-agent" description = "remowt on-host agent serving fs/pty/systemd endpoints over bifrostlink" version.workspace = true -edition = "2021" +edition.workspace = true license.workspace = true [dependencies] --- a/remowt/cmds/remowt-agent/src/askpass.rs +++ b/remowt/cmds/remowt-agent/src/askpass.rs @@ -1,47 +1,22 @@ use std::borrow::Cow; use std::io::Write as _; -use anyhow::Context as _; -use remowt_ui_prompt::dbus::{DbusPrompterInterface, DbusPrompterProxy, BUS_NAME, PROMPTER_PATH}; +use bifrostlink::declarative::RemoteEndpoints as _; +use remowt_link_shared::{gateway, Address, BifConfig}; +use remowt_ui_prompt::bifrost::PromptEndpointsClient; use remowt_ui_prompt::{Prompter, Source}; -use tracing::debug; -use zbus::Connection; - -pub async fn serve

(conn: &Connection, prompter: P) -> anyhow::Result<()> -where - P: Prompter + 'static, -{ - conn.object_server() - .at(PROMPTER_PATH, DbusPrompterInterface(prompter)) - .await?; - match conn.request_name(BUS_NAME).await { - Ok(()) => {} - Err(zbus::Error::NameTaken) => { - debug!("{BUS_NAME} already owned, chaining to upstream"); - } - Err(e) => return Err(e.into()), - } - Ok(()) -} pub async fn ask(prompt: &str, description: String) -> anyhow::Result<()> { - let conn = Connection::session() - .await - .context("connecting to the session bus (DBUS_SESSION_BUS_ADDRESS)")?; - let proxy = DbusPrompterProxy::builder(&conn) - .destination(BUS_NAME)? - .path(PROMPTER_PATH)? - .build() - .await?; - - let password = proxy - .prompt_text( - false, - prompt, - &description, - &[Source(Cow::Borrowed("remowt-askpass"))], - ) - .await?; + let rpc = gateway::connect(&gateway::socket_path()?).await?; + let prompter = PromptEndpointsClient::::wrap(rpc.remote(Address::User)); + let password = Prompter::prompt_text( + &prompter, + false, + prompt, + &description, + &[Source(Cow::Borrowed("remowt-askpass"))], + ) + .await?; let mut out = std::io::stdout().lock(); out.write_all(password.as_bytes())?; --- a/remowt/cmds/remowt-agent/src/bus.rs +++ /dev/null @@ -1,40 +0,0 @@ -use std::process::Stdio; - -use anyhow::Context as _; -use futures::StreamExt as _; -use tokio::process::{Child, Command}; -use tokio_util::codec::{FramedRead, LinesCodec}; -use zbus::Connection; - -pub struct PrivateBus { - pub address: String, - pub conn: Connection, - _child: Child, -} - -pub async fn spawn() -> anyhow::Result { - let mut child = Command::new("dbus-daemon") - .args(["--session", "--nofork", "--print-address"]) - .stdout(Stdio::piped()) - .kill_on_drop(true) - .spawn() - .context("spawning dbus-daemon for the private bus")?; - - let stdout = child.stdout.take().expect("piped"); - let address = FramedRead::new(stdout, LinesCodec::new()) - .next() - .await - .context("dbus-daemon exited before printing its address")? - .context("reading dbus-daemon address")?; - - let conn = zbus::connection::Builder::address(address.as_str())? - .build() - .await - .context("connecting to the private bus")?; - - Ok(PrivateBus { - address, - conn, - _child: child, - }) -} --- a/remowt/cmds/remowt-agent/src/editor.rs +++ b/remowt/cmds/remowt-agent/src/editor.rs @@ -3,87 +3,23 @@ use std::time::Duration; use std::{fs, io}; -use anyhow::{bail, Context as _}; +use anyhow::{anyhow, bail, Context as _}; +use bifrostlink::declarative::RemoteEndpoints as _; use nix::libc; use remowt_link_shared::editor::EditorEndpointsClient; +use remowt_link_shared::{gateway, Address, BifConfig}; use tokio::process::Command; -use zbus::{fdo, interface, proxy, Connection}; - -use remowt_link_shared::BifConfig; - -const BUS_NAME: &str = "lach.RemowtEditor"; -const SERVICE_PATH: &str = "/lach/Editor"; - -pub struct EditorService { - editor: EditorEndpointsClient, -} -#[interface(name = "lach.RemowtEditor")] -impl EditorService { - /// Attach the User's GUI to the nvim server at `socket_path` (on the remote), - /// blocking until the user is done. - async fn edit(&self, socket_path: String) -> fdo::Result<()> { - self.editor - .open_editor(socket_path) - .await - .map_err(|e| fdo::Error::Failed(format!("requesting editor on the User: {e}")))? - .map_err(|e| fdo::Error::Failed(format!("editor failed: {e}")))?; - Ok(()) - } - - async fn forward_tcp(&self, addr: String) -> fdo::Result { - let local = self - .editor - .expose_tcp(addr) - .await - .map_err(|e| fdo::Error::Failed(format!("requesting tcp forward on the User: {e}")))? - .map_err(|e| fdo::Error::Failed(format!("tcp forward failed: {e}")))?; - Ok(local) - } - - async fn forward_udp(&self, addr: String) -> fdo::Result { - let local = self - .editor - .expose_udp(addr) - .await - .map_err(|e| fdo::Error::Failed(format!("requesting udp forward on the User: {e}")))? - .map_err(|e| fdo::Error::Failed(format!("udp forward failed: {e}")))?; - Ok(local) - } -} - -pub async fn serve( - conn: &Connection, - editor: EditorEndpointsClient, -) -> anyhow::Result<()> { - conn.object_server() - .at(SERVICE_PATH, EditorService { editor }) - .await?; - conn.request_name(BUS_NAME).await?; - Ok(()) -} - -#[proxy(interface = "lach.RemowtEditor")] -trait RemowtEditor { - async fn edit(&self, socket_path: &str) -> fdo::Result<()>; - async fn forward_tcp(&self, addr: &str) -> fdo::Result; - async fn forward_udp(&self, addr: &str) -> fdo::Result; -} - pub async fn forward(udp: bool, addr: String) -> anyhow::Result<()> { - let conn = Connection::session() - .await - .context("connecting to the session bus (DBUS_SESSION_BUS_ADDRESS)")?; - let proxy = RemowtEditorProxy::builder(&conn) - .destination(BUS_NAME)? - .path(SERVICE_PATH)? - .build() - .await?; + let rpc = gateway::connect(&gateway::socket_path()?).await?; + let editor = EditorEndpointsClient::::wrap(rpc.remote(Address::User)); let local = if udp { - proxy.forward_udp(&addr).await? + editor.expose_udp(addr).await } else { - proxy.forward_tcp(&addr).await? - }; + editor.expose_tcp(addr).await + } + .map_err(|e| anyhow!("requesting forward on the User: {e}"))? + .map_err(|e| anyhow!("forward failed: {e}"))?; println!("{local}"); Ok(()) } @@ -125,15 +61,13 @@ .await .context("nvim did not start its server")?; - let conn = Connection::session() + let rpc = gateway::connect(&gateway::socket_path()?).await?; + let editor = EditorEndpointsClient::::wrap(rpc.remote(Address::User)); + let result = editor + .open_editor(sock_str) .await - .context("connecting to the session bus (DBUS_SESSION_BUS_ADDRESS)")?; - let proxy = RemowtEditorProxy::builder(&conn) - .destination(BUS_NAME)? - .path(SERVICE_PATH)? - .build() - .await?; - let result = proxy.edit(&sock_str).await; + .map_err(|e| anyhow!("requesting editor on the User: {e}")) + .and_then(|r| r.map_err(|e| anyhow!("editor failed: {e}"))); if tokio::time::timeout(Duration::from_secs(2), child.wait()) .await @@ -143,8 +77,7 @@ } let _ = fs::remove_file(&sock); - result?; - Ok(()) + result } async fn wait_for_socket(path: &Path) -> anyhow::Result<()> { --- a/remowt/cmds/remowt-agent/src/main.rs +++ b/remowt/cmds/remowt-agent/src/main.rs @@ -7,8 +7,8 @@ use std::path::PathBuf; use std::sync::{Arc, Mutex, OnceLock}; +use bifrostlink::Rpc; use bifrostlink::declarative::RemoteEndpoints; -use bifrostlink::Rpc; use bifrostlink_ports::stdio::from_stdio; use bifrostlink_ports::unix_socket::from_socket; use clap::Parser; @@ -17,9 +17,9 @@ subprocess::Subprocess, systemd::Systemd, }; use remowt_link_shared::iroh_tunnel::TunnelDialer; -use remowt_link_shared::{editor::EditorEndpointsClient, Address, BifConfig}; -use remowt_polkit_shared::{emphasize, Identity, PidDisplay}; -use remowt_ui_prompt::bifrost::PromptEndpointsClient; +use remowt_link_shared::{Address, BifConfig, gateway}; +use remowt_polkit_shared::{Identity, PidDisplay, emphasize}; +use remowt_ui_prompt::bifrost::{PromptEndpointsClient, serve_prompts}; use remowt_ui_prompt::rofi::RofiPrompter; use remowt_ui_prompt::{PrependSourcePrompter, Prompter, Source}; use tokio::fs; @@ -29,13 +29,12 @@ use tracing::{debug, trace}; use zbus::fdo; use zbus::zvariant::{OwnedValue, Str}; -use zbus::{interface, Connection}; +use zbus::{Connection, interface}; use zbus_polkit::policykit1::Subject; use self::helper::{Helper, SocketHelper, SuidHelper}; pub mod askpass; -pub mod bus; pub mod editor; pub mod helper; @@ -132,19 +131,18 @@ 0 => { return Err(fdo::Error::AuthFailed( "no identity to authenticate as".to_owned(), - )) + )); } 1 => 0, - _ => { - prompter - .prompt_enum( - "Identity", - "Select identity to use for polkit authorization", - &identity_displays, - &[], - ) - .await? - } + _ => prompter + .prompt_enum( + "Identity", + "Select identity to use for polkit authorization", + &identity_displays, + &[], + ) + .await + .map_err(prompt_err)?, }; debug!("identity chosen"); @@ -198,6 +196,15 @@ } } +fn prompt_err(value: remowt_ui_prompt::Error) -> fdo::Error { + use remowt_ui_prompt::Error; + match value { + Error::Cancel => fdo::Error::NoReply("input was cancelled".to_owned()), + Error::Remote(e) => fdo::Error::NoReply(format!("remote error occured: {e}")), + Error::InputError(e) => fdo::Error::Failed(e), + } +} + const OBJ_PATH: &str = "/org/freedesktop/PolicyKit1/AuthenticationAgent"; #[derive(Parser)] @@ -268,10 +275,12 @@ }; register_auth_agent(&system_conn, Agent::new(helper, RofiPrompter)).await?; - let session_conn = Connection::session().await?; - askpass::serve(&session_conn, RofiPrompter).await?; + let mut rpc = Rpc::::new(Address::User); + serve_prompts(&mut rpc, RofiPrompter); + + gateway::serve(rpc.clone(), &gateway::local_socket()?).await?; - let _keep_alive = (system_conn, session_conn); + let _keep_alive = (system_conn, rpc); pending().await } async fn main_real_agent( @@ -298,13 +307,9 @@ remowt_plugin::host::serve(&mut rpc); let user_prompter = PromptEndpointsClient::wrap(rpc.remote(Address::User)); - let editor_client = EditorEndpointsClient::wrap(rpc.remote(Address::User)); - - let bus = bus::spawn().await?; - askpass::serve(&bus.conn, user_prompter.clone()).await?; - editor::serve(&bus.conn, editor_client).await?; let helpers = tempfile::Builder::new().prefix("remowt-path.").tempdir()?; + let gateway_socket = helpers.path().join(gateway::SOCKET_NAME); let exe = std::env::current_exe()?; let askpass_helper = helpers.path().join("remowt-askpass"); let editor_helper = helpers.path().join("remowt-editor"); @@ -342,7 +347,7 @@ std::env::set_var("SSH_ASKPASS_REQUIRE", "force"); std::env::set_var("EDITOR", &editor_helper); std::env::set_var("VISUAL", &editor_helper); - std::env::set_var("DBUS_SESSION_BUS_ADDRESS", &bus.address); + std::env::set_var(gateway::SOCKET_ENV, &gateway_socket); } let port = match path { @@ -351,10 +356,9 @@ }; rpc.add_direct(Address::User, port, bifrostlink::Rtt(0)); + gateway::serve(rpc.clone(), &gateway_socket).await?; + let polkit_conn = if !privileged && !local { - // The unprivileged agent doubles as a polkit authentication agent so - // `run0` (e.g. our own elevation) routes its prompt to the User over - // bifrost instead of failing on a tty-less session. let conn = Connection::system().await?; let helper = SocketHelper { fallback: SuidHelper, @@ -365,7 +369,7 @@ None }; - let _keep_alive = (bus, helpers, polkit_conn); + let _keep_alive = (helpers, polkit_conn); pending().await } --- a/remowt/cmds/remowt-ssh/Cargo.toml +++ b/remowt/cmds/remowt-ssh/Cargo.toml @@ -2,7 +2,7 @@ name = "remowt-ssh" description = "SSH transport client for connecting to a remowt agent" version.workspace = true -edition = "2021" +edition.workspace = true license.workspace = true [dependencies] --- a/remowt/cmds/remowt-ssh/src/main.rs +++ b/remowt/cmds/remowt-ssh/src/main.rs @@ -75,14 +75,8 @@ description: "".to_owned(), }, ); - if let Some(sess) = conn.ssh() { - serve_editor( - &mut rpc, - SshEditor { - sess, - conn: conn.clone(), - }, - ); + if conn.ssh().is_some() { + serve_editor(&mut rpc, SshEditor { conn: conn.clone() }); } debug!("entering shell"); --- a/remowt/crates/polkit-shared/Cargo.toml +++ b/remowt/crates/polkit-shared/Cargo.toml @@ -2,7 +2,7 @@ name = "remowt-polkit-shared" description = "Shared polkit/PAM types for remowt" version.workspace = true -edition = "2021" +edition.workspace = true license.workspace = true [dependencies] --- a/remowt/crates/remowt-client/Cargo.toml +++ b/remowt/crates/remowt-client/Cargo.toml @@ -2,7 +2,7 @@ name = "remowt-client" description = "russh-based client connection to a remowt agent" version.workspace = true -edition = "2021" +edition.workspace = true license.workspace = true [dependencies] --- a/remowt/crates/remowt-client/src/editor.rs +++ b/remowt/crates/remowt-client/src/editor.rs @@ -6,14 +6,12 @@ use remowt_endpoints::forward::ForwardClient; use remowt_link_shared::editor::{EditorBackend, Error}; use remowt_link_shared::BifConfig; -use russh::client::Handle; use tokio::net::{TcpListener, UdpSocket, UnixListener}; use tracing::error; -use crate::{Remowt, SshHandler}; +use crate::Remowt; pub struct SshEditor { - pub sess: Arc>, pub conn: Remowt, } impl EditorBackend for SshEditor { @@ -22,21 +20,41 @@ let _ = std::fs::remove_file(&local); let listener = UnixListener::bind(&local).map_err(|e| Error::Failed(e.to_string()))?; - let sess = self.sess.clone(); + let conn = self.conn.clone(); let forward = tokio::spawn(async move { loop { let Ok((mut stream, _)) = listener.accept().await else { break; }; - let sess = sess.clone(); + let conn = conn.clone(); let remote = socket_path.clone(); tokio::spawn(async move { - match sess.channel_open_direct_streamlocal(remote).await { - Ok(ch) => { - let mut remote = ch.into_stream(); + // Rides the iroh fast tunnel when established, else an + // ssh-forwarded unix socket. + let (forwarded, tunnel) = match conn.bind_fast_tunnel("editor", false).await { + Ok(v) => v, + Err(e) => { + error!("editor: bind tunnel failed: {e}"); + return; + } + }; + let fclient: ForwardClient = conn.endpoints(); + match fclient.connect_unix(tunnel, remote).await { + Ok(Ok(())) => {} + Ok(Err(e)) => { + error!("editor: agent connect_unix failed: {e}"); + return; + } + Err(e) => { + error!("editor: connect_unix rpc failed: {e}"); + return; + } + } + match forwarded.accept().await { + Ok(mut remote) => { let _ = tokio::io::copy_bidirectional(&mut stream, &mut remote).await; } - Err(e) => error!("opening direct-streamlocal to nvim failed: {e}"), + Err(e) => error!("editor: accept tunnel failed: {e}"), } }); } --- a/remowt/crates/remowt-client/src/lib.rs +++ b/remowt/crates/remowt-client/src/lib.rs @@ -4,7 +4,7 @@ use std::sync::atomic::AtomicU64; use std::sync::{Arc, Mutex}; -use anyhow::{anyhow, bail, ensure, Context as _, Result}; +use anyhow::{Context as _, Result, anyhow, bail, ensure}; use bifrostlink::declarative::RemoteEndpoints; use bifrostlink::{Remote, Rpc, Rtt}; use camino::{Utf8Path, Utf8PathBuf}; @@ -12,12 +12,12 @@ use remowt_link_shared::plugin::PluginEndpointsClient; use remowt_link_shared::port::child_port; use remowt_link_shared::{Address, BifConfig}; -use russh::client::{connect, Config, Handle, Handler, Msg, Session}; +use russh::Channel; +use russh::client::{Config, Handle, Handler, Msg, Session, connect}; +use russh::keys::agent::AgentIdentity; use russh::keys::agent::client::AgentClient; -use russh::keys::agent::AgentIdentity; use russh::keys::check_known_hosts; use russh::keys::ssh_key::PublicKey; -use russh::Channel; use tempfile::TempDir; use tokio::io::AsyncRead; use tokio::net::UnixListener; @@ -287,8 +287,12 @@ .unwrap_or_else(|| env::var("USER").unwrap_or_else(|_| "root".to_owned())); let subs: Subs = Arc::new(Mutex::new(HashMap::new())); + let config = Config { + nodelay: true, + ..Config::default() + }; let mut sess = connect( - Arc::new(Config::default()), + Arc::new(config), (hostname.clone(), port), SshHandler { host: hostname, @@ -549,9 +553,7 @@ } } - /// Bind a data tunnel, preferring the iroh fast path when it is up. Escalated tunnels - /// (the privileged agent, a separate process with no iroh connection of its own) and the - /// local transport always use the ssh/unix path. + /// Bind a data tunnel, preferring the iroh fast path when it is up. pub async fn bind_fast_tunnel( &self, hint: &str, @@ -586,7 +588,7 @@ async fn try_setup_iroh(&self) -> Result<()> { use remowt_endpoints::iroh_tunnel::IrohTunnelClient; - use remowt_link_shared::iroh_tunnel::{build_endpoint, ssh_custom_addr, REMOWT_ALPN}; + use remowt_link_shared::iroh_tunnel::{REMOWT_ALPN, build_endpoint, ssh_custom_addr}; let (listener, sock) = self.bind_runtime_unix("iroh-xport").await?; let secret = iroh::SecretKey::generate(); @@ -605,7 +607,7 @@ ); let conn = ep.connect(addr, REMOWT_ALPN).await?; ensure!(conn.remote_id() == agent_id, "iroh peer identity mismatch"); - info!("iroh fast tunnel established"); + debug!("iroh fast tunnel established"); let subs = self.0.iroh.subs.clone(); let accept_conn = conn.clone(); @@ -633,16 +635,6 @@ debug!("iroh accept loop ended: {e}"); break; } - } - } - }); - - let log_conn = conn.clone(); - tokio::spawn(async move { - tokio::time::sleep(std::time::Duration::from_secs(2)).await; - for p in log_conn.paths().iter() { - if p.is_selected() { - info!(rtt = ?p.rtt(), "iroh selected path: {}", p.remote_addr()); } } }); @@ -688,10 +680,10 @@ } fn local_runtime_dir() -> Result<(Utf8PathBuf, Option)> { - if let Ok(dir) = env::var("XDG_RUNTIME_DIR") { - if !dir.is_empty() { - return Ok((Utf8PathBuf::from(dir), None)); - } + if let Ok(dir) = env::var("XDG_RUNTIME_DIR") + && !dir.is_empty() + { + return Ok((Utf8PathBuf::from(dir), None)); } let tmp = tempfile::Builder::new() .prefix("remowt.") --- a/remowt/crates/remowt-endpoints/Cargo.toml +++ b/remowt/crates/remowt-endpoints/Cargo.toml @@ -2,7 +2,7 @@ name = "remowt-endpoints" description = "Nix daemon proxy" version.workspace = true -edition = "2021" +edition.workspace = true license.workspace = true [dependencies] --- a/remowt/crates/remowt-endpoints/src/forward.rs +++ b/remowt/crates/remowt-endpoints/src/forward.rs @@ -3,12 +3,12 @@ use std::sync::atomic::{AtomicU64, Ordering}; use std::sync::{Arc, Mutex}; +use bifrostlink::Config; use bifrostlink::declarative::endpoints; -use bifrostlink::Config; use remowt_link_shared::iroh_tunnel::{TunnelAddr, TunnelDialer}; use serde::{Deserialize, Serialize}; use std::result::Result; -use tokio::net::{TcpStream, UdpSocket}; +use tokio::net::{TcpStream, UdpSocket, UnixStream}; use tracing::warn; #[derive(Serialize, Deserialize, Debug, thiserror::Error)] @@ -107,6 +107,24 @@ Ok(session) } + + #[endpoints(id = 3)] + async fn connect_unix(&self, tunnel: TunnelAddr, path: String) -> Result<(), Error> { + let stream = self + .dialer + .connect_tunnel(&tunnel) + .await + .map_err(|e| Error::Tunnel(e.to_string()))?; + let unix = UnixStream::connect(&path) + .await + .map_err(|e| Error::Connect(path, e.to_string()))?; + tokio::spawn(async move { + let mut stream = stream; + let mut unix = unix; + let _ = tokio::io::copy_bidirectional(&mut stream, &mut unix).await; + }); + Ok(()) + } } fn unspecified_for(target: &SocketAddr) -> SocketAddr { --- a/remowt/crates/remowt-link-shared/Cargo.toml +++ b/remowt/crates/remowt-link-shared/Cargo.toml @@ -2,7 +2,7 @@ name = "remowt-link-shared" description = "Shared bifrostlink endpoint wiring for remowt" version.workspace = true -edition = "2021" +edition.workspace = true license.workspace = true [dependencies] @@ -13,13 +13,13 @@ serde = { workspace = true, features = ["derive"] } serde_json.workspace = true thiserror.workspace = true -tokio = { workspace = true, features = ["fs", "io-util", "macros", "net", "sync", "rt"] } +tokio = { workspace = true, features = ["fs", "io-util", "macros", "net", "rt", "sync"] } tokio-util = { workspace = true, features = ["codec"] } tracing.workspace = true -remowt-ui-prompt.workspace = true +uuid = { workspace = true, features = ["v4"] } +futures.workspace = true iroh.workspace = true iroh-base.workspace = true -noq-udp.workspace = true n0-watcher.workspace = true -futures.workspace = true +noq-udp.workspace = true --- a/remowt/crates/remowt-link-shared/src/lib.rs +++ b/remowt/crates/remowt-link-shared/src/lib.rs @@ -7,6 +7,7 @@ use serde::{Deserialize, Serialize}; pub mod editor; +pub mod gateway; pub mod iroh_tunnel; pub mod port; @@ -16,6 +17,7 @@ Agent, AgentPrivileged, Plugin(u16), + Ephemeral(uuid::Uuid), } impl AddressT for Address {} @@ -27,9 +29,6 @@ ListenerDead, #[error("response: {0}")] Response(String), - - #[error(transparent)] - Ui(#[from] remowt_ui_prompt::Error), } impl From for Error { @@ -38,8 +37,8 @@ } } impl From for Error { - fn from(_value: serde_json::Error) -> Self { - Self::ListenerDead + fn from(e: serde_json::Error) -> Self { + Self::Response(format!("{e}")) } } impl From for ResponseError { --- a/remowt/crates/remowt-link-shared/src/plugin.rs +++ b/remowt/crates/remowt-link-shared/src/plugin.rs @@ -10,6 +10,8 @@ BadName, #[error("spawning plugin failed: {0}")] Spawn(String), + #[error("failed to wait for connection, plugin process has died?")] + ConnectionWaitError, #[error("agent is shutting down")] Gone, } --- a/remowt/crates/remowt-plugin/Cargo.toml +++ b/remowt/crates/remowt-plugin/Cargo.toml @@ -2,7 +2,7 @@ name = "remowt-plugin" description = "Plugin host and protocol for remowt agents" version.workspace = true -edition = "2021" +edition.workspace = true license.workspace = true [dependencies] --- a/remowt/crates/remowt-plugin/src/host.rs +++ b/remowt/crates/remowt-plugin/src/host.rs @@ -8,6 +8,7 @@ use remowt_link_shared::plugin::{Error, PluginEndpoints, PluginHost}; use remowt_link_shared::port::child_port; use remowt_link_shared::{Address, BifConfig}; +use tokio::select; pub fn serve(rpc: &mut Rpc) { let host = Host { @@ -25,7 +26,7 @@ } impl Host { - fn spawn(&self, id: u16, path: impl AsRef) -> Result<(), Error> { + async fn spawn(&self, id: u16, path: impl AsRef) -> Result<(), Error> { let rpc = self.rpc.clone().upgrade().ok_or(Error::Gone)?; let mut child = Command::new(path) @@ -40,8 +41,20 @@ let stdout = child.stdout.take().expect("stdout piped"); let addr = Address::Plugin(id); - rpc.add_direct(addr, child_port(stdout, stdin), Rtt(0)); + rpc.add_direct(addr.clone(), child_port(stdout, stdin), Rtt(0)); + + select! { + e = rpc.wait_for_connection_to(addr) => { + if e.is_err() { + return Err(Error::ConnectionWaitError) + } + }, + _ = child.wait() => { + return Err(Error::ConnectionWaitError) + } + }; self.children.lock().expect("not poisoned").push(child); + Ok(()) } } @@ -58,13 +71,13 @@ let dir = exe .parent() .ok_or_else(|| Error::Spawn("primary agent has no parent directory".to_owned()))?; - self.spawn(id, dir.join(&name)) + self.spawn(id, dir.join(&name)).await } async fn load_plugin_path(&self, id: u16, path: String) -> Result<(), Error> { if path.is_empty() || path.contains('\0') { return Err(Error::BadName); } - self.spawn(id, path) + self.spawn(id, path).await } } --- a/remowt/crates/remowt-ui-prompt/src/bifrost.rs +++ b/remowt/crates/remowt-ui-prompt/src/bifrost.rs @@ -103,7 +103,6 @@ where P: Prompter + Send + Sync + 'static, C: Config, - C::Error: From, { PromptEndpoints(prompt).register_endpoints(rpc); } --- a/rust-toolchain.toml +++ b/rust-toolchain.toml @@ -1,3 +1,8 @@ [toolchain] channel = "stable" components = ["rustfmt", "clippy", "rust-analyzer", "rust-src"] +targets = [ + "x86_64-unknown-linux-musl", + "aarch64-unknown-linux-musl", + "armv7-unknown-linux-musleabihf", +] -- gitstuff