42 files changed
--- 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<Vec<OsString>> {
- 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<Vec<OsString>> {
-// 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<sha2::Sha256>;
- let mut mac = <HmacSha256 as hmac::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 = <HmacSha256 as hmac::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<SecretOwner, BTreeMap<String, FleetSecretDistribution>>,
}
-impl FleetData {
- pub fn from_str(s: &str) -> anyhow::Result<Self> {
+impl FromStr for FleetData {
+ type Err = anyhow::Error;
+ fn from_str(s: &str) -> anyhow::Result<Self> {
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<Utc>) {
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<SecretOwner>,
@@ -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<usize> {
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<SecretOwner>,
@@ -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<SecretOwner>,
}
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<String, FleetSecretDistributions> {
+ 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 {
1use std::{2 collections::{BTreeMap, BTreeSet, HashSet},3 future::Future,4 io::Write,5 ops::Deref,6 path::PathBuf,7 pin::Pin,8 str::FromStr,9 sync::{Arc, OnceLock},10};1112use anyhow::{Context, Result, anyhow, bail, ensure};13use camino::{Utf8Path, Utf8PathBuf};14use chrono::{DateTime, Utc};15use fleet_shared::SecretData;16use nix_eval::{Store, Value, eval_store, nix_go, nix_go_json, util::assert_warn};17use remowt_client::{AgentBundle, Remowt};18use remowt_endpoints::fs::FsClient;19use remowt_link_shared::Address;20use remowt_ui_prompt::auto::AutoPrompter;21use remowt_ui_prompt::bifrost::PromptEndpoints;22use remowt_ui_prompt::{PrependSourcePrompter, Source};23use tabled::Tabled;24use tempfile::NamedTempFile;25use time::UtcDateTime;26use tokio::task::spawn_blocking;27use tracing::warn;2829use crate::fleetdata::{30 FleetData, FleetSecretData, FleetSecretDistribution, FleetSecretPart, SecretOwner,31};3233pub struct FleetConfigInternals {34 pub prefer_identities: BTreeSet<SecretOwner>,35 pub now: DateTime<Utc>,3637 38 pub directory: PathBuf,39 40 pub local_system: String,41 pub data: Arc<FleetData>,42 43 pub config_field: Value,44 45 pub flake_outputs: Value,46 47 pub localhost: String,4849 50 pub default_pkgs: Value,51 52 pub nixpkgs: Value,5354 pub local_host: OnceLock<Arc<ConfigHost>>,55}565758#[derive(Clone)]59pub struct Config(pub Arc<FleetConfigInternals>);6061impl Deref for Config {62 type Target = FleetConfigInternals;6364 fn deref(&self) -> &Self::Target {65 &self.066 }67}6869#[derive(Clone, PartialEq, Copy, Debug)]70pub enum DeployKind {71 72 UpgradeToFleet,73 74 Fleet,75 76 77 NixosInstall,78 79 80 81 NixosLustrate,82}8384impl FromStr for DeployKind {85 type Err = anyhow::Error;86 fn from_str(s: &str) -> std::result::Result<Self, Self::Err> {87 match s {88 "upgrade-to-fleet" => Ok(Self::UpgradeToFleet),89 "fleet" => Ok(Self::Fleet),90 "nixos-install" => Ok(Self::NixosInstall),91 "nixos-lustrate" => Ok(Self::NixosLustrate),92 v => bail!(93 "unknown deploy_kind: {v}; expected on of \"upgrade-to-fleet\", \"fleet\", \"nixos-install\", \"nixos-lustrate\""94 ),95 }96 }97}98pub struct ConfigHost {99 config: Config,100 pub name: String,101 groups: OnceLock<Vec<String>>,102103 104 deploy_kind: OnceLock<DeployKind>,105 session_destination: OnceLock<String>,106 legacy_ssh_store: OnceLock<bool>,107108 pub host_config: Option<Value>,109 pub nixos_config: OnceLock<Value>,110 pub nixos_unchecked_config: OnceLock<Value>,111 pub pkgs_override: Option<Value>,112113 114 pub local: bool,115 pub remowt: OnceLock<Remowt>,116 nix_store: OnceLock<Arc<Store>>,117 nix_plugin: tokio::sync::OnceCell<()>,118}119120const NIX_PLUGIN_ID: u16 = 2;121122fn agents_dir() -> Result<PathBuf> {123 std::env::var_os("REMOWT_AGENTS_DIR")124 .map(PathBuf::from)125 .or_else(|| option_env!("REMOWT_AGENTS_DIR").map(PathBuf::from))126 .ok_or_else(|| {127 anyhow!("no remowt-agents bundle; set REMOWT_AGENTS_DIR to a remowt-agents output")128 })129}130131fn agent_bundle() -> Result<AgentBundle> {132 AgentBundle::from_dir(agents_dir()?)133}134135#[derive(Debug, Clone, Copy)]136pub enum GenerationStorage {137 Deployer,138 Machine,139 Pusher,140}141impl GenerationStorage {142 fn prefix(&self) -> &'static str {143 match self {144 GenerationStorage::Deployer => "deployer.",145 GenerationStorage::Machine => "",146 GenerationStorage::Pusher => "pusher.",147 }148 }149}150151#[derive(Tabled, Debug)]152pub struct Generation {153 #[tabled(rename = "ID", format("{}", self.rollback_id()))]154 pub id: u32,155 #[tabled(rename = "Current")]156 pub current: bool,157 #[tabled(rename = "Created at")]158 pub datetime: UtcDateTime,159 #[tabled(format = "{:?}")]160 pub store_path: Utf8PathBuf,161 #[tabled(skip)]162 pub location: GenerationStorage,163}164impl Generation {165 pub fn rollback_id(&self) -> String {166 format!("{}{}", self.location.prefix(), self.id)167 }168}169170impl ConfigHost {171 pub async fn list_generations(&self, profile: &str) -> Result<Vec<Generation>> {172 let plugin_id = self.ensure_nix_plugin().await?;173 let nix = self174 .remowt()175 .await?176 .plugin_endpoints::<remowt_fleet::NixClient<_>>(plugin_id);177 let raw = nix178 .list_generations(profile.to_owned())179 .await180 .map_err(|e| anyhow!("{e:?}"))?181 .map_err(|e| anyhow!("{e}"))?;182 raw.into_iter()183 .map(|g| {184 let id: u32 =185 g.id.try_into()186 .with_context(|| format!("generation id {} doesn't fit in u32", g.id))?;187 let datetime = UtcDateTime::from_unix_timestamp(g.creation_time_unix)188 .with_context(|| {189 format!("invalid generation timestamp {}", g.creation_time_unix)190 })?;191 Ok(Generation {192 id,193 current: g.current,194 datetime,195 store_path: g.store_path,196 location: GenerationStorage::Machine,197 })198 })199 .collect()200 }201202 pub fn set_session_destination(&self, dest: String) {203 self.session_destination204 .set(dest)205 .expect("session destination is already set")206 }207 pub fn set_deploy_kind(&self, kind: DeployKind) {208 self.deploy_kind209 .set(kind)210 .expect("deploy kind is already set");211 }212 pub fn set_legacy_ssh_store(&self, legacy: bool) {213 self.legacy_ssh_store214 .set(legacy)215 .expect("legacy ssh store is already set")216 }217 pub async fn deploy_kind(&self) -> Result<DeployKind> {218 if let Some(kind) = self.deploy_kind.get() {219 return Ok(*kind);220 }221 let remowt = self.remowt().await?;222 let fs = remowt.endpoints::<FsClient<_>>();223 let is_fleet_managed = match fs.file_exists(Utf8PathBuf::from("/etc/FLEET_HOST")).await {224 Ok(v) => v,225 Err(e) => {226 bail!("failed to query remote system kind: {e}");227 }228 };229 if !is_fleet_managed {230 bail!(231 "{}",232 indoc::indoc! {"233 host is not marked as managed by fleet234 if you're not trying to lustrate/install system from scratch,235 you should either236 1. manually create /etc/FLEET_HOST file on the target host,237 2. use ?deploy_kind=fleet host argument if you're upgrading from older version of fleet238 3. use ?deploy_kind=upgrade_to_fleet if you're upgrading from plain nixos to fleet-managed nixos239 for installation use ?deploy_kind=nixos_install / ?deploy_kind=nixos_lustrate 240 "}241 );242 }243 244 let _ = self.deploy_kind.set(DeployKind::Fleet);245 Ok(*self.deploy_kind.get().expect("deploy kind is just set"))246 }247 async fn connection(&self) -> Result<Remowt> {248 if let Some(conn) = self.remowt.get() {249 return Ok(conn.clone());250 }251 let bundle = agent_bundle()?;252 let conn = if self.local {253 Remowt::connect_local(&bundle, "remowt-fleet".to_owned())254 .await255 .context("starting local remowt agent")?256 } else {257 let dest = self258 .session_destination259 .get()260 .cloned()261 .unwrap_or_else(|| self.name.clone());262 Remowt::connect(&dest, &bundle, "remowt-fleet".to_owned())263 .await264 .map_err(|e| anyhow!("remowt error while connecting to {}: {e:#?}", self.name))?265 };266 PromptEndpoints(PrependSourcePrompter {267 prompter: AutoPrompter::new().await,268 source: if self.local {269 vec![]270 } else {271 vec![Source(std::borrow::Cow::Owned(format!(272 "ssh host: {}",273 self.name274 )))]275 },276 description: "".to_owned(),277 })278 .register_endpoints(&mut conn.rpc());279 let _ = self.remowt.set(conn);280 Ok(self.remowt.get().expect("just set").clone())281 }282283 284 pub async fn remowt(&self) -> Result<Remowt> {285 self.connection().await286 }287288 pub fn ensure_nix_plugin(&self) -> Pin<Box<dyn Future<Output = Result<u16>> + Send + '_>> {289 Box::pin(async {290 self.nix_plugin291 .get_or_try_init(|| async {292 let pkgs = self.pkgs()?;293 let name = "remowt-plugin-fleet";294 let plugin = nix_go!(pkgs[{ name }]);295 let built = plugin296 .build("out")297 .context("failed to build the fleet nix plugin")?;298 let copied = self299 .remote_derivation(&built)300 .await301 .context("failed to copy the fleet nix plugin to the host store")?;302 let bin = copied.join("bin/remowt-plugin-fleet");303 self.remowt()304 .await?305 .run0_load_plugin_path(NIX_PLUGIN_ID, bin.as_str())306 .await307 .context("failed to load the fleet nix plugin")?;308 self.remowt()309 .await?310 .rpc()311 .wait_for_connection_to(Address::Plugin(NIX_PLUGIN_ID))312 .await313 .map_err(|_| anyhow!("failed to wait for plugin"))?;314 anyhow::Ok(())315 })316 .await?;317 Ok(NIX_PLUGIN_ID)318 })319 }320321 async fn nix_store(&self) -> Result<Arc<Store>> {322 if let Some(store) = self.nix_store.get() {323 return Ok(store.clone());324 }325 let conn = self.connection().await?;326 let socket = match self.deploy_kind().await? {327 DeployKind::NixosInstall => {328 remowt_fleet::nix_store_socket(conn, "/mnt?require-sigs=false").await?329 }330 _ => remowt_fleet::nix_store_socket(conn, "auto").await?,331 };332 let uri = format!("unix://{}", socket.display());333 let store = Arc::new(Store::open(&uri)?);334 let _ = self.nix_store.set(store);335 Ok(self.nix_store.get().expect("just set").clone())336 }337338 pub async fn decrypt(&self, data: SecretData) -> Result<Vec<u8>> {339 ensure!(data.encrypted, "secret is not encrypted");340 let remowt = self.remowt().await?;341 let mut cmd = remowt.cmd("fleet-install-secrets");342 cmd.arg("decrypt").eqarg("--secret", data.to_string());343 let encoded = cmd344 .sudo()345 .run_string()346 .await347 .context("failed to call remote host for decrypt")?;348 let data: SecretData = encoded.parse().map_err(|e| anyhow!("{e}"))?;349 ensure!(!data.encrypted, "secret came out encrypted");350 Ok(data.data)351 }352 pub async fn reencrypt_distribution(353 &self,354 data: &FleetSecretDistribution,355 targets: BTreeSet<SecretOwner>,356 now: DateTime<Utc>,357 ) -> Result<FleetSecretDistribution> {358 let mut parts = BTreeMap::new();359 for (part_name, part) in &data.secret.parts {360 parts.insert(361 part_name.clone(),362 if part.raw.encrypted {363 FleetSecretPart {364 raw: self.reencrypt(part.raw.clone(), targets.clone()).await?,365 }366 } else {367 part.clone()368 },369 );370 }371 let secret = FleetSecretData {372 created_at: data.secret.created_at,373 expires_at: data.secret.expires_at,374 generation_data: data.secret.generation_data.clone(),375 parts,376 };377 Ok(FleetSecretDistribution::new(targets, secret, now))378 }379 pub async fn reencrypt(380 &self,381 data: SecretData,382 targets: BTreeSet<SecretOwner>,383 ) -> Result<SecretData> {384 let remowt = self.remowt().await?;385 ensure!(data.encrypted, "secret is not encrypted");386 let mut cmd = remowt.cmd("fleet-install-secrets");387 cmd.arg("reencrypt").eqarg("--secret", data.to_string());388 for target in targets {389 let key = self.config.key(&target).await?;390 cmd.eqarg("--targets", key);391 }392 let encoded = cmd393 .sudo()394 .run_string()395 .await396 .context("failed to call remote host for decrypt")?;397 let data: SecretData = encoded.parse().map_err(|e| anyhow!("{e}"))?;398 ensure!(data.encrypted, "secret came out not encrypted");399 Ok(data)400 }401 402 pub async fn remote_derivation(&self, path: impl AsRef<Utf8Path>) -> Result<Utf8PathBuf> {403 let path = path.as_ref().to_owned();404 if self.local {405 406 return Ok(path);407 }408 let sign: Pin<Box<dyn Future<Output = Result<()>> + Send>> = {409 let path = path.clone();410 Box::pin(async move {411 let local = self.config.local_host();412 let plugin_id = local.ensure_nix_plugin().await?;413 let nix = local414 .remowt()415 .await?416 .plugin_endpoints::<remowt_fleet::NixClient<_>>(plugin_id);417 nix.sign_closure(path, Utf8PathBuf::from("/etc/nix/private-key"))418 .await419 .map_err(|e| anyhow!("{e:?}"))?420 .map_err(|e| anyhow!("{e}"))?;421 Ok(())422 })423 };424 if let Err(e) = sign.await {425 warn!("failed to sign store paths: {e}");426 }427 let store = self.nix_store().await?;428 let eval_store = eval_store();429 {430 let path = path.clone();431 spawn_blocking(move || eval_store.copy_to(&store, path.as_ref()))432 .await433 .expect("copy_to panicked")434 .context("copying closure to remote store")?;435 }436 Ok(path)437 }438}439440struct HostSecretDefinition(Value);441442impl ConfigHost {443 444 445 pub fn tags(&self) -> Result<Vec<String>> {446 if let Some(v) = self.groups.get() {447 return Ok(v.clone());448 }449 let Some(host_config) = &self.host_config else {450 return Ok(vec![]);451 };452 let tags: Vec<String> = nix_go_json!(host_config.tags);453454 let _ = self.groups.set(tags.clone());455456 Ok(tags)457 }458 pub fn nixos_config(&self) -> Result<Value> {459 if let Some(v) = self.nixos_config.get() {460 return Ok(v.clone());461 }462 let Some(host_config) = &self.host_config else {463 bail!("local host has no nixos_config");464 };465 let nixos_config = nix_go!(host_config.nixos.config);466 assert_warn("nixos config evaluation", &nixos_config)?;467468 let _ = self.nixos_config.set(nixos_config.clone());469470 Ok(nixos_config)471 }472 pub fn nixos_unchecked_config(&self) -> Result<Value> {473 if let Some(v) = self.nixos_unchecked_config.get() {474 return Ok(v.clone());475 }476 let Some(host_config) = &self.host_config else {477 bail!("local host has no nixos_config");478 };479 let nixos_config = nix_go!(host_config.nixos_unchecked.config);480481 let _ = self.nixos_unchecked_config.set(nixos_config.clone());482483 Ok(nixos_config)484 }485486 pub fn list_defined_secrets(&self) -> Result<Vec<String>> {487 let nixos = self.nixos_unchecked_config()?;488 let secrets = nix_go!(nixos.secrets);489 secrets.list_fields()490 }491492 493 pub fn pkgs(&self) -> Result<Value> {494 if let Some(value) = &self.pkgs_override {495 return Ok(value.clone());496 }497 let Some(host_config) = &self.host_config else {498 bail!("local host has no host_config");499 };500 501 Ok(nix_go!(host_config.nixos.options._module.args.value.pkgs))502 }503}504505#[derive(Clone)]506pub struct SharedSecretDefinition(Value);507impl SharedSecretDefinition {508 pub fn expected_owners(&self) -> Result<BTreeSet<SecretOwner>> {509 let secret = &self.0;510 Ok(nix_go_json!(secret.expectedOwners))511 }512 pub fn allow_different(&self) -> Result<bool> {513 let secret = &self.0;514 Ok(nix_go_json!(secret.allowDifferent))515 }516 pub fn regenerate_on_owner_added(&self) -> Result<bool> {517 let secret = &self.0;518 Ok(nix_go_json!(secret.regenerateOnOwnerAdded))519 }520 pub fn regenerate_on_owner_removed(&self) -> Result<bool> {521 let secret = &self.0;522 Ok(nix_go_json!(secret.regenerateOnOwnerRemoved))523 }524 pub fn generator(&self) -> Result<Value> {525 let secret = &self.0;526 Ok(nix_go!(secret.generator))527 }528}529530impl Config {531 pub fn tagged_hostnames(&self, tag: &str) -> Result<Vec<String>> {532 let config = &self.config_field;533 let tagged: Vec<String> = nix_go_json!(config.taggedWith[{ tag }]);534 Ok(tagged)535 }536 pub fn expand_owner_set(&self, owners: Vec<String>) -> Result<BTreeSet<String>> {537 let mut out = BTreeSet::new();538 for owner in owners {539 if let Some(tag) = owner.strip_prefix('@') {540 let hosts = self.tagged_hostnames(tag)?;541 out.extend(hosts);542 } else {543 out.insert(owner);544 }545 }546 Ok(out)547 }548 pub fn local_host(&self) -> Arc<ConfigHost> {549 self.local_host550 .get_or_init(|| {551 Arc::new(ConfigHost {552 config: self.clone(),553 name: "<virtual localhost>".to_owned(),554 host_config: None,555 nixos_config: OnceLock::new(),556 nixos_unchecked_config: OnceLock::new(),557 groups: {558 let cell = OnceLock::new();559 let _ = cell.set(vec![]);560 cell561 },562 pkgs_override: Some(self.default_pkgs.clone()),563564 local: true,565 remowt: OnceLock::new(),566 nix_store: OnceLock::new(),567 nix_plugin: tokio::sync::OnceCell::new(),568 deploy_kind: OnceLock::new(),569 session_destination: OnceLock::new(),570 legacy_ssh_store: OnceLock::new(),571 })572 })573 .clone()574 }575576 pub fn preferred_hosts(577 &self,578 filter: impl Fn(&str) -> bool,579 ) -> Result<impl Iterator<Item = Result<ConfigHost>>> {580 let prefer = self581 .prefer_identities582 .iter()583 .filter_map(|v| v.as_host())584 .collect::<HashSet<_>>();585 let config = &self.config_field;586 let mut names = nix_go!(config.hosts).list_fields()?;587 names.retain(|s| filter(s));588 names.sort_by_key(|h| prefer.contains(h.as_str()));589590 Ok(names.into_iter().map(|h| self.host(&h)))591 }592593 pub fn host(&self, name: &str) -> Result<ConfigHost> {594 let config = &self.config_field;595 let host_config = nix_go!(config.hosts[{ name }]);596597 Ok(ConfigHost {598 config: self.clone(),599 name: name.to_owned(),600 host_config: Some(host_config),601 nixos_config: OnceLock::new(),602 nixos_unchecked_config: OnceLock::new(),603 groups: OnceLock::new(),604 pkgs_override: None,605606 607 local: self.localhost == name,608 remowt: OnceLock::new(),609 nix_store: OnceLock::new(),610 nix_plugin: tokio::sync::OnceCell::new(),611 deploy_kind: OnceLock::new(),612 session_destination: OnceLock::new(),613 legacy_ssh_store: OnceLock::new(),614 })615 }616 pub fn list_hosts(&self) -> Result<Vec<ConfigHost>> {617 let config = &self.config_field;618 let names = nix_go!(config.hosts).list_fields()?;619 let mut out = vec![];620 for name in names {621 out.push(self.host(&name)?);622 }623 Ok(out)624 }625 626 pub fn system_config(&self, host: &str) -> Result<Value> {627 let fleet_field = &self.config_field;628 Ok(nix_go!(fleet_field.hosts[{ host }].nixos.config))629 }630631 pub fn secret_definition(&self, secret: &str) -> Result<Option<SharedSecretDefinition>> {632 let config = &self.config_field;633 let shared_secrets = nix_go!(config.secrets);634 if !shared_secrets.has_field(secret)? {635 return Ok(None);636 }637 Ok(Some(SharedSecretDefinition(nix_go!(638 shared_secrets[secret]639 ))))640 }641642 pub fn save(&self) -> Result<()> {643 let mut tempfile = NamedTempFile::new_in(self.directory.clone()).context("failed to create updated version of fleet.nix in the same directory as original.\nDo you have write access to it? Access only to the fleet.nix won't be enough, the directory is used for atomic overwrite operation.\nIt is not recommended to use fleet by root anyway, move fleet project to your home directory.")?;644 let data = nixlike::serialize(&*self.data)?;645 tempfile.write_all(646 format!(647 "# This file contains fleet state and shouldn't be edited by hand\n\n{data}\n\n# vim: ts=2 et nowrap\n"648 )649 .as_bytes(),650 )?;651 let mut fleet_data_path = self.directory.clone();652 fleet_data_path.push("fleet.nix");653 tempfile.persist(fleet_data_path)?;654 Ok(())655 }656}
1use std::{2 collections::{BTreeMap, BTreeSet, HashSet},3 future::Future,4 io::Write,5 ops::Deref,6 path::PathBuf,7 pin::Pin,8 str::FromStr,9 sync::{Arc, OnceLock},10};1112use anyhow::{Context, Result, anyhow, bail, ensure};13use camino::{Utf8Path, Utf8PathBuf};14use chrono::{DateTime, Utc};15use fleet_shared::SecretData;16use nix_eval::{Store, Value, drv::DrvGraph, eval_store, nix_go, nix_go_json, util::assert_warn};17use remowt_client::{AgentBundle, Remowt};18use remowt_endpoints::fs::FsClient;19use remowt_fleet::NixClient;20use remowt_link_shared::BifConfig;21use remowt_ui_prompt::auto::AutoPrompter;22use remowt_ui_prompt::bifrost::PromptEndpoints;23use remowt_ui_prompt::{PrependSourcePrompter, Source};24use tempfile::NamedTempFile;25use tokio::{sync::OnceCell, task::spawn_blocking};26use tracing::{info, warn};2728use crate::fleetdata::{29 FleetData, FleetSecretData, FleetSecretDistribution, FleetSecretPart, SecretOwner,30};3132pub struct FleetConfigInternals {33 pub prefer_identities: BTreeSet<SecretOwner>,34 pub now: DateTime<Utc>,3536 37 pub directory: PathBuf,38 39 pub local_system: String,40 pub data: Arc<FleetData>,41 42 pub config_field: Value,43 44 pub flake_outputs: Value,45 46 pub localhost: String,4748 49 pub default_pkgs: Value,50 51 pub nixpkgs: Value,5253 pub local_host: OnceLock<Arc<ConfigHost>>,54}555657#[derive(Clone)]58pub struct Config(pub Arc<FleetConfigInternals>);5960impl Deref for Config {61 type Target = FleetConfigInternals;6263 fn deref(&self) -> &Self::Target {64 &self.065 }66}6768#[derive(Clone, PartialEq, Copy, Debug)]69pub enum DeployKind {70 71 UpgradeToFleet,72 73 Fleet,74 75 76 NixosInstall,77 78 79 80 NixosLustrate,81}8283impl FromStr for DeployKind {84 type Err = anyhow::Error;85 fn from_str(s: &str) -> std::result::Result<Self, Self::Err> {86 match s {87 "upgrade-to-fleet" => Ok(Self::UpgradeToFleet),88 "fleet" => Ok(Self::Fleet),89 "nixos-install" => Ok(Self::NixosInstall),90 "nixos-lustrate" => Ok(Self::NixosLustrate),91 v => bail!(92 "unknown deploy_kind: {v}; expected on of \"upgrade-to-fleet\", \"fleet\", \"nixos-install\", \"nixos-lustrate\""93 ),94 }95 }96}97pub struct ConfigHost {98 config: Config,99 pub name: String,100 groups: OnceLock<Vec<String>>,101102 103 deploy_kind: OnceCell<DeployKind>,104 session_destination: OnceLock<String>,105 legacy_ssh_store: OnceLock<bool>,106107 pub host_config: Option<Value>,108 pub nixos_config: OnceLock<Value>,109 pub nixos_unchecked_config: OnceLock<Value>,110 pub pkgs_override: Option<Value>,111112 113 pub local: bool,114 pub remowt: OnceLock<Remowt>,115 nix_store: OnceLock<Arc<Store>>,116 nix_plugin: OnceCell<()>,117}118119const NIX_PLUGIN_ID: u16 = 2;120121fn agents_dir() -> Result<PathBuf> {122 std::env::var_os("REMOWT_AGENTS_DIR")123 .map(PathBuf::from)124 .or_else(|| option_env!("REMOWT_AGENTS_DIR").map(PathBuf::from))125 .ok_or_else(|| {126 anyhow!("no remowt-agents bundle; set REMOWT_AGENTS_DIR to a remowt-agents output")127 })128}129130fn agent_bundle() -> Result<AgentBundle> {131 AgentBundle::from_dir(agents_dir()?)132}133134enum RemoteDerivationMode {135 CopySigned,136 Copy,137 Attic,138 Cachix,139}140141impl ConfigHost {142 pub fn set_session_destination(&self, dest: String) {143 self.session_destination144 .set(dest)145 .expect("session destination is already set")146 }147 pub fn set_deploy_kind(&self, kind: DeployKind) {148 self.deploy_kind149 .set(kind)150 .expect("deploy kind is already set");151 }152 pub fn set_legacy_ssh_store(&self, legacy: bool) {153 self.legacy_ssh_store154 .set(legacy)155 .expect("legacy ssh store is already set")156 }157 pub async fn deploy_kind(&self) -> Result<DeployKind> {158 self.deploy_kind.get_or_try_init(|| async {159 let remowt = self.remowt().await?;160 let fs = remowt.endpoints::<FsClient<_>>();161 let is_fleet_managed = match fs.file_exists(Utf8PathBuf::from("/etc/FLEET_HOST")).await {162 Ok(v) => v,163 Err(e) => {164 bail!("failed to query remote system kind: {e}");165 }166 };167 if !is_fleet_managed {168 bail!(169 "{}",170 indoc::indoc! {"171 host is not marked as managed by fleet172 if you're not trying to lustrate/install system from scratch,173 you should either174 1. manually create /etc/FLEET_HOST file on the target host,175 2. use ?deploy_kind=fleet host argument if you're upgrading from older version of fleet176 3. use ?deploy_kind=upgrade_to_fleet if you're upgrading from plain nixos to fleet-managed nixos177 for installation use ?deploy_kind=nixos_install / ?deploy_kind=nixos_lustrate 178 "}179 );180 }181 Ok(DeployKind::Fleet)182 }).await.copied()183 }184 async fn connection(&self) -> Result<Remowt> {185 if let Some(conn) = self.remowt.get() {186 return Ok(conn.clone());187 }188 let bundle = agent_bundle()?;189 let conn = if self.local {190 Remowt::connect_local(&bundle, "remowt-fleet".to_owned())191 .await192 .context("starting local remowt agent")?193 } else {194 let dest = self195 .session_destination196 .get()197 .cloned()198 .unwrap_or_else(|| self.name.clone());199 Remowt::connect(&dest, &bundle, "remowt-fleet".to_owned())200 .await201 .map_err(|e| anyhow!("remowt error while connecting to {}: {e:#?}", self.name))?202 };203 PromptEndpoints(PrependSourcePrompter {204 prompter: AutoPrompter::new().await,205 source: if self.local {206 vec![]207 } else {208 vec![Source(std::borrow::Cow::Owned(format!(209 "ssh host: {}",210 self.name211 )))]212 },213 description: "".to_owned(),214 })215 .register_endpoints(&mut conn.rpc());216 let _ = self.remowt.set(conn);217 Ok(self.remowt.get().expect("just set").clone())218 }219220 221 pub async fn remowt(&self) -> Result<Remowt> {222 self.connection().await223 }224225 pub async fn nix_client(&self) -> Result<NixClient<BifConfig>> {226 let remowt = self.remowt().await?;227 let plugin_id = self.ensure_nix_plugin().await?;228 Ok(remowt.plugin_endpoints(plugin_id))229 }230 pub fn ensure_nix_plugin(&self) -> Pin<Box<dyn Future<Output = Result<u16>> + Send + '_>> {231 Box::pin(async {232 self.nix_plugin233 .get_or_try_init(|| async {234 let pkgs = self.pkgs()?;235 let name = "remowt-plugin-fleet";236 let plugin = nix_go!(pkgs[{ name }]);237 let built = plugin238 .build("out")239 .context("failed to build the fleet nix plugin")?;240 let copied = self241 .remote_derivation(&built)242 .await243 .context("failed to copy the fleet nix plugin to the host store")?;244 let bin = copied.join("bin/remowt-plugin-fleet");245 self.remowt()246 .await?247 .run0_load_plugin_path(NIX_PLUGIN_ID, bin.as_str())248 .await249 .context("failed to load the fleet nix plugin")?;250 anyhow::Ok(())251 })252 .await?;253 Ok(NIX_PLUGIN_ID)254 })255 }256257 async fn nix_store(&self) -> Result<Arc<Store>> {258 if let Some(store) = self.nix_store.get() {259 return Ok(store.clone());260 }261 let conn = self.connection().await?;262 let socket = match self.deploy_kind().await? {263 DeployKind::NixosInstall => {264 remowt_fleet::nix_store_socket(conn, "/mnt?require-sigs=false").await?265 }266 _ => remowt_fleet::nix_store_socket(conn, "auto").await?,267 };268 let uri = format!("unix://{}", socket.display());269 let store = Arc::new(Store::open(&uri)?);270 let _ = self.nix_store.set(store);271 Ok(self.nix_store.get().expect("just set").clone())272 }273274 pub async fn decrypt(&self, data: SecretData) -> Result<Vec<u8>> {275 ensure!(data.encrypted, "secret is not encrypted");276 let remowt = self.remowt().await?;277 let mut cmd = remowt.cmd("fleet-install-secrets");278 cmd.arg("decrypt").eqarg("--secret", data.to_string());279 let encoded = cmd280 .sudo()281 .run_string()282 .await283 .context("failed to call remote host for decrypt")?;284 let data: SecretData = encoded.parse().map_err(|e| anyhow!("{e}"))?;285 ensure!(!data.encrypted, "secret came out encrypted");286 Ok(data.data)287 }288 pub async fn reencrypt_distribution(289 &self,290 data: &FleetSecretDistribution,291 targets: BTreeSet<SecretOwner>,292 now: DateTime<Utc>,293 ) -> Result<FleetSecretDistribution> {294 let mut parts = BTreeMap::new();295 for (part_name, part) in &data.secret.parts {296 parts.insert(297 part_name.clone(),298 if part.raw.encrypted {299 FleetSecretPart {300 raw: self.reencrypt(part.raw.clone(), targets.clone()).await?,301 }302 } else {303 part.clone()304 },305 );306 }307 let secret = FleetSecretData {308 created_at: data.secret.created_at,309 expires_at: data.secret.expires_at,310 generation_data: data.secret.generation_data.clone(),311 parts,312 };313 Ok(FleetSecretDistribution::new(targets, secret, now))314 }315 pub async fn reencrypt(316 &self,317 data: SecretData,318 targets: BTreeSet<SecretOwner>,319 ) -> Result<SecretData> {320 let remowt = self.remowt().await?;321 ensure!(data.encrypted, "secret is not encrypted");322 let mut cmd = remowt.cmd("fleet-install-secrets");323 cmd.arg("reencrypt").eqarg("--secret", data.to_string());324 for target in targets {325 let key = self.config.key(&target).await?;326 cmd.eqarg("--targets", key);327 }328 let encoded = cmd329 .sudo()330 .run_string()331 .await332 .context("failed to call remote host for decrypt")?;333 let data: SecretData = encoded.parse().map_err(|e| anyhow!("{e}"))?;334 ensure!(data.encrypted, "secret came out not encrypted");335 Ok(data)336 }337 338 pub async fn remote_derivation(&self, path: impl AsRef<Utf8Path>) -> Result<Utf8PathBuf> {339 let path = path.as_ref().to_owned();340 if self.local {341 342 return Ok(path);343 }344 let sign: Pin<Box<dyn Future<Output = Result<()>> + Send>> = {345 let path = path.clone();346 let graph = DrvGraph::resolve(&eval_store(), &path)347 .context("failed to resolve graph to be uploaded")?;348 info!("signing {} paths", graph.nodes.len());349350 Box::pin(async move {351 let local = self.config.local_host();352 let plugin_id = local.ensure_nix_plugin().await?;353 let nix = local354 .remowt()355 .await?356 .plugin_endpoints::<remowt_fleet::NixClient<_>>(plugin_id);357 nix.sign_closure(path, Utf8PathBuf::from("/etc/nix/private-key"))358 .await359 .map_err(|e| anyhow!("{e:?}"))?360 .map_err(|e| anyhow!("{e}"))?;361 Ok(())362 })363 };364 if let Err(e) = sign.await {365 warn!("failed to sign store paths: {e}");366 }367 let store = self.nix_store().await?;368 let eval_store = eval_store();369 {370 let path = path.clone();371 spawn_blocking(move || eval_store.copy_to(&store, path.as_ref()))372 .await373 .expect("copy_to panicked")374 .context("copying closure to remote store")?;375 }376 Ok(path)377 }378}379380struct HostSecretDefinition(Value);381382impl ConfigHost {383 384 385 pub fn tags(&self) -> Result<Vec<String>> {386 if let Some(v) = self.groups.get() {387 return Ok(v.clone());388 }389 let Some(host_config) = &self.host_config else {390 return Ok(vec![]);391 };392 let tags: Vec<String> = nix_go_json!(host_config.tags);393394 let _ = self.groups.set(tags.clone());395396 Ok(tags)397 }398 pub fn nixos_config(&self) -> Result<Value> {399 if let Some(v) = self.nixos_config.get() {400 return Ok(v.clone());401 }402 let Some(host_config) = &self.host_config else {403 bail!("local host has no nixos_config");404 };405 let nixos_config = nix_go!(host_config.nixos.config);406 assert_warn("nixos config evaluation", &nixos_config)?;407408 let _ = self.nixos_config.set(nixos_config.clone());409410 Ok(nixos_config)411 }412 pub fn nixos_unchecked_config(&self) -> Result<Value> {413 if let Some(v) = self.nixos_unchecked_config.get() {414 return Ok(v.clone());415 }416 let Some(host_config) = &self.host_config else {417 bail!("local host has no nixos_config");418 };419 let nixos_config = nix_go!(host_config.nixos_unchecked.config);420421 let _ = self.nixos_unchecked_config.set(nixos_config.clone());422423 Ok(nixos_config)424 }425426 pub fn list_defined_secrets(&self) -> Result<Vec<String>> {427 let nixos = self.nixos_unchecked_config()?;428 let secrets = nix_go!(nixos.secrets);429 secrets.list_fields()430 }431432 433 pub fn pkgs(&self) -> Result<Value> {434 if let Some(value) = &self.pkgs_override {435 return Ok(value.clone());436 }437 let Some(host_config) = &self.host_config else {438 bail!("local host has no host_config");439 };440 441 Ok(nix_go!(host_config.nixos.options._module.args.value.pkgs))442 }443}444445#[derive(Clone)]446pub struct SharedSecretDefinition(Value);447impl SharedSecretDefinition {448 pub fn expected_owners(&self) -> Result<BTreeSet<SecretOwner>> {449 let secret = &self.0;450 Ok(nix_go_json!(secret.expectedOwners))451 }452 pub fn allow_different(&self) -> Result<bool> {453 let secret = &self.0;454 Ok(nix_go_json!(secret.allowDifferent))455 }456 pub fn regenerate_on_owner_added(&self) -> Result<bool> {457 let secret = &self.0;458 Ok(nix_go_json!(secret.regenerateOnOwnerAdded))459 }460 pub fn regenerate_on_owner_removed(&self) -> Result<bool> {461 let secret = &self.0;462 Ok(nix_go_json!(secret.regenerateOnOwnerRemoved))463 }464 pub fn generator(&self) -> Result<Value> {465 let secret = &self.0;466 Ok(nix_go!(secret.generator))467 }468}469470impl Config {471 pub fn tagged_hostnames(&self, tag: &str) -> Result<Vec<String>> {472 let config = &self.config_field;473 let tagged: Vec<String> = nix_go_json!(config.taggedWith[{ tag }]);474 Ok(tagged)475 }476 pub fn expand_owner_set(&self, owners: Vec<String>) -> Result<BTreeSet<String>> {477 let mut out = BTreeSet::new();478 for owner in owners {479 if let Some(tag) = owner.strip_prefix('@') {480 let hosts = self.tagged_hostnames(tag)?;481 out.extend(hosts);482 } else {483 out.insert(owner);484 }485 }486 Ok(out)487 }488 pub fn local_host(&self) -> Arc<ConfigHost> {489 self.local_host490 .get_or_init(|| {491 Arc::new(ConfigHost {492 config: self.clone(),493 name: "<virtual localhost>".to_owned(),494 host_config: None,495 nixos_config: OnceLock::new(),496 nixos_unchecked_config: OnceLock::new(),497 groups: {498 let cell = OnceLock::new();499 let _ = cell.set(vec![]);500 cell501 },502 pkgs_override: Some(self.default_pkgs.clone()),503504 local: true,505 remowt: OnceLock::new(),506 nix_store: OnceLock::new(),507 nix_plugin: OnceCell::new(),508 deploy_kind: OnceCell::new(),509 session_destination: OnceLock::new(),510 legacy_ssh_store: OnceLock::new(),511 })512 })513 .clone()514 }515516 pub fn preferred_hosts(517 &self,518 filter: impl Fn(&str) -> bool,519 ) -> Result<impl Iterator<Item = Result<ConfigHost>>> {520 let prefer = self521 .prefer_identities522 .iter()523 .filter_map(|v| v.as_host())524 .collect::<HashSet<_>>();525 let config = &self.config_field;526 let mut names = nix_go!(config.hosts).list_fields()?;527 names.retain(|s| filter(s));528 names.sort_by_key(|h| prefer.contains(h.as_str()));529530 Ok(names.into_iter().map(|h| self.host(&h)))531 }532533 pub fn host(&self, name: &str) -> Result<ConfigHost> {534 let config = &self.config_field;535 let host_config = nix_go!(config.hosts[{ name }]);536537 Ok(ConfigHost {538 config: self.clone(),539 name: name.to_owned(),540 host_config: Some(host_config),541 nixos_config: OnceLock::new(),542 nixos_unchecked_config: OnceLock::new(),543 groups: OnceLock::new(),544 pkgs_override: None,545546 547 local: self.localhost == name,548 remowt: OnceLock::new(),549 nix_store: OnceLock::new(),550 nix_plugin: OnceCell::new(),551 deploy_kind: OnceCell::new(),552 session_destination: OnceLock::new(),553 legacy_ssh_store: OnceLock::new(),554 })555 }556 pub fn list_hosts(&self) -> Result<Vec<ConfigHost>> {557 let config = &self.config_field;558 let names = nix_go!(config.hosts).list_fields()?;559 let mut out = vec![];560 for name in names {561 out.push(self.host(&name)?);562 }563 Ok(out)564 }565 566 pub fn system_config(&self, host: &str) -> Result<Value> {567 let fleet_field = &self.config_field;568 Ok(nix_go!(fleet_field.hosts[{ host }].nixos.config))569 }570571 pub fn secret_definition(&self, secret: &str) -> Result<Option<SharedSecretDefinition>> {572 let config = &self.config_field;573 let shared_secrets = nix_go!(config.secrets);574 if !shared_secrets.has_field(secret)? {575 return Ok(None);576 }577 Ok(Some(SharedSecretDefinition(nix_go!(578 shared_secrets[secret]579 ))))580 }581582 pub fn save(&self) -> Result<()> {583 let mut tempfile = NamedTempFile::new_in(self.directory.clone()).context("failed to create updated version of fleet.nix in the same directory as original.\nDo you have write access to it? Access only to the fleet.nix won't be enough, the directory is used for atomic overwrite operation.\nIt is not recommended to use fleet by root anyway, move fleet project to your home directory.")?;584 let data = nixlike::serialize(&*self.data)?;585 tempfile.write_all(586 format!(587 "# This file contains fleet state and shouldn't be edited by hand\n\n{data}\n\n# vim: ts=2 et nowrap\n"588 )589 .as_bytes(),590 )?;591 let mut fleet_data_path = self.directory.clone();592 fleet_data_path.push("fleet.nix");593 tempfile.persist(fleet_data_path)?;594 Ok(())595 }596}
--- 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<Vec<Generation>> {
+ 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<P>(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::<BifConfig>::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<PrivateBus> {
- 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<BifConfig>,
-}
-#[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<u16> {
- 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<u16> {
- 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<BifConfig>,
-) -> 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<u16>;
- async fn forward_udp(&self, addr: &str) -> fdo::Result<u16>;
-}
-
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::<BifConfig>::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::<BifConfig>::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::<BifConfig>::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<Handle<SshHandler>>,
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<BifConfig> = 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<TempDir>)> {
- 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<ListenerForYourRequestHasBeenDeadError> for Error {
@@ -38,8 +37,8 @@
}
}
impl From<serde_json::Error> for Error {
- fn from(_value: serde_json::Error) -> Self {
- Self::ListenerDead
+ fn from(e: serde_json::Error) -> Self {
+ Self::Response(format!("{e}"))
}
}
impl From<Error> 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<BifConfig>) {
let host = Host {
@@ -25,7 +26,7 @@
}
impl Host {
- fn spawn(&self, id: u16, path: impl AsRef<OsStr>) -> Result<(), Error> {
+ async fn spawn(&self, id: u16, path: impl AsRef<OsStr>) -> 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<Error>,
{
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",
+]