From bbdbb8df43d22369a79aebf2ef65cc5643ed1091 Mon Sep 17 00:00:00 2001 From: Miles Ward Date: Thu, 21 May 2026 02:13:08 -0400 Subject: [PATCH] v0.3.2: OpenPXE naming cleanup and beta hardening - remove remaining PXEForge/Gate/anvil wording from code, docs, UI, and deployment examples - fix Queue tab view wiring and rename queue-facing terminal/API copy - harden raw ISO Range handling, iPXE fallback lines, and WinPE SMB reconnect behavior - bump workspace and deployment examples to 0.3.2 --- Cargo.toml | 2 +- README.md | 33 +-- crates/core/src/arch.rs | 5 +- crates/core/src/client.rs | 30 +- crates/core/src/config.rs | 18 +- crates/core/src/host_bindings.rs | 2 +- crates/core/src/lib.rs | 4 +- crates/core/src/metrics.rs | 73 ++++- crates/core/src/queue.rs | 52 ++-- crates/core/src/settings.rs | 22 +- crates/dhcp-proxy/src/reply.rs | 18 +- crates/dhcp-proxy/src/server.rs | 47 +++- crates/http-api/src/app.rs | 372 +++++++++++++++--------- crates/http-api/src/ipxe_script.rs | 190 ++++++++++--- crates/http-api/src/iso_fs.rs | 43 ++- crates/http-api/src/log_stream.rs | 16 +- crates/http-api/src/terminal.rs | 120 ++++---- crates/http-api/tests/full_flow.rs | 419 ++++++++++++++++++++-------- crates/ipxe-assets/src/lib.rs | 4 +- crates/iso-store/src/introspect.rs | 49 +++- crates/iso-store/src/lib.rs | 4 +- crates/iso-store/src/nfs.rs | 12 +- crates/iso-store/src/smb.rs | 97 +++++-- crates/iso-store/src/store.rs | 50 +++- crates/iso-store/src/windows.rs | 63 ++++- crates/openpxe/src/main.rs | 68 +++-- crates/tftp/src/server.rs | 64 ++++- crates/webui/src/app.css | 5 +- crates/webui/src/app.js | 14 +- crates/webui/src/index.html | 2 +- crates/webui/src/lib.rs | 26 +- deploy/openshift/30-deployment.yaml | 2 +- deploy/unraid/README.md | 14 +- deploy/unraid/openpxe.xml | 4 +- docs/NEXT_PHASE.md | 4 +- docs/architecture.md | 16 +- runbooks/linux-network-boot.md | 12 +- 37 files changed, 1387 insertions(+), 589 deletions(-) diff --git a/Cargo.toml b/Cargo.toml index 44f66b6..d15b52c 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -12,7 +12,7 @@ members = [ ] [workspace.package] -version = "0.3.1" +version = "0.3.2" edition = "2021" rust-version = "1.80" license = "MIT OR Apache-2.0" diff --git a/README.md b/README.md index f39f96e..d604fca 100644 --- a/README.md +++ b/README.md @@ -5,12 +5,12 @@ Container-native PXE boot server. A Rust reimplementation of for Docker/OCI and OpenShift. Upload `.iso` files via the web UI; network clients PXE-boot them. -> **Status:** v0.2.0 / pre-beta. Phases 1–5 complete: full PXE stack, +> **Status:** v0.3.2 / pre-beta. Phases 1–5 complete: full PXE stack, > Queued Deployment queue, NFS-share ISO sources, live tracing log + an > operator terminal, per-MAC host bindings (Tinkerbell-style), -> Prometheus `/metrics`, light/dark theme toggle, animated anvil -> imaging-progress widget. **66 tests passing**, clippy clean. Ready -> for real-hardware validation. +> Prometheus `/metrics`, light/dark theme toggle, animated OpenPXE +> imaging-progress widget, and per-ISO boot passwords. The test suite and +> clippy are part of the release checklist. Ready for real-hardware validation. ## Design non-negotiables @@ -41,13 +41,13 @@ clients PXE-boot them. Tools > Utilities / OpenPXE Shell / Network Card Info Queued Deployment ``` -6. **Queued Deployment queue** — the "horse race" launch flow. A client that +6. **Queued Deployment queue** — the coordinated launch flow. A client that selects *Queued Deployment* gets a numbered position and waits. The operator picks an ISO in the web UI and fires it to every waiting client simultaneously. -7. **Web UI** (Netbox-style): sidebar nav (Dashboard / Network / Forge - Gate / Storage / Hosts / Terminal / About), light + dark themes - (toggle top-right or press `T`), animated anvil "forge progress" +7. **Web UI** (Netbox-style): sidebar nav (Dashboard / Network / Queue / + Storage / Hosts / Terminal / About), light + dark themes + (toggle top-right or press `T`), animated OpenPXE progress widget when devices are imaging. All assets served from the binary — no external requests. 8. **Per-MAC host bindings.** Pin a MAC to a boot target and the client @@ -81,7 +81,7 @@ skip TFTP and respond with an HTTP URL. ./scripts/fetch-ipxe.sh # 2. Build the container image (~3 min first time). -docker buildx build -f deploy/docker/Dockerfile -t openpxe:0.1.0 --load . +docker buildx build -f deploy/docker/Dockerfile -t openpxe:0.3.2 --load . # 3. Run it on the box plugged into your PXE network. Set PUBLIC_IP to # this host's LAN address so advertised iPXE URLs are reachable. @@ -91,7 +91,7 @@ docker run -d --name openpxe \ -e OPENPXE_DHCP_MODE=proxy \ -v $PWD/data/isos:/var/lib/openpxe/isos \ -v $PWD/data/work:/var/lib/openpxe/work \ - openpxe:0.1.0 + openpxe:0.3.2 # 4. Open the UI and drop an ISO in. open http://10.0.0.5 @@ -122,7 +122,7 @@ docker buildx create --name openpxe-multi --driver docker-container --use # Build + push both linux/amd64 and linux/arm64 under one tag. docker buildx build --builder openpxe-multi \ --platform linux/amd64,linux/arm64 \ - -t ghcr.io/YOUR-ORG/openpxe:0.1.0 \ + -t ghcr.io/YOUR-ORG/openpxe:0.3.2 \ --push \ -f deploy/docker/Dockerfile . ``` @@ -155,10 +155,10 @@ docker run --rm \ -v /my/iso-library:/seed:ro \ -v openpxe-data:/var/lib/openpxe/isos \ -e OPENPXE_PUBLIC_IP=10.0.0.5 \ - openpxe:0.1.0 seed --from /seed + openpxe:0.3.2 seed --from /seed # Dry run first to see what would be imported: -docker run --rm -v /my/iso-library:/seed:ro openpxe:0.1.0 seed --from /seed --dry-run +docker run --rm -v /my/iso-library:/seed:ro openpxe:0.3.2 seed --from /seed --dry-run ``` ### Environment overrides @@ -273,7 +273,7 @@ operational constraints inherited from the design: ## Queued Deployment -The "horse race" launch flow, end to end: +The coordinated launch flow, end to end: 1. A client boots and picks **Queued Deployment** in the PXE menu (or falls through on timeout with the default `timeout_action`). @@ -285,8 +285,9 @@ The "horse race" launch flow, end to end: The server broadcasts the assignment to every queued client via a `tokio::sync::Notify`; each client's next poll returns the boot script for the chosen image. -5. Every client chains the same image at effectively the same moment — the - queue releases and the horses run together. +5. Every client chains the same image at effectively the same moment. The + queue stays visible until the operator releases entries, which keeps a + useful audit trail during hardware testing. No user-facing iPXE anywhere in this flow. The client only ever runs scripts we generate; the operator only interacts with the web UI. diff --git a/crates/core/src/arch.rs b/crates/core/src/arch.rs index 9170d87..7f8b65b 100644 --- a/crates/core/src/arch.rs +++ b/crates/core/src/arch.rs @@ -126,7 +126,10 @@ mod tests { fn bootfile_names_stable() { assert_eq!(ClientArch::LegacyX86.ipxe_bootfile(), Some("undionly.kpxe")); assert_eq!(ClientArch::X64Uefi.ipxe_bootfile(), Some("snponly.efi")); - assert_eq!(ClientArch::Arm64Uefi.ipxe_bootfile(), Some("snponly-arm64.efi")); + assert_eq!( + ClientArch::Arm64Uefi.ipxe_bootfile(), + Some("snponly-arm64.efi") + ); assert_eq!(ClientArch::Unknown(0xFFFF).ipxe_bootfile(), None); } diff --git a/crates/core/src/client.rs b/crates/core/src/client.rs index 5332e62..14c119f 100644 --- a/crates/core/src/client.rs +++ b/crates/core/src/client.rs @@ -56,19 +56,25 @@ impl ClientRegistry { ) { let mut guard = self.inner.write(); let now = OffsetDateTime::now_utc(); - let entry = guard.entry(mac.to_string()).or_insert_with(|| ClientSnapshot { - mac: mac.to_string(), - last_ip: ip, - arch, - hostname: None, - first_seen: now, - last_seen: now, - events: Vec::new(), - selected_target: None, - }); + let entry = guard + .entry(mac.to_string()) + .or_insert_with(|| ClientSnapshot { + mac: mac.to_string(), + last_ip: ip, + arch, + hostname: None, + first_seen: now, + last_seen: now, + events: Vec::new(), + selected_target: None, + }); entry.last_seen = now; - if ip.is_some() { entry.last_ip = ip; } - if arch.is_some() { entry.arch = arch; } + if ip.is_some() { + entry.last_ip = ip; + } + if arch.is_some() { + entry.arch = arch; + } entry.events.push((now, event)); // Cap event history per client to keep memory bounded. const MAX_EVENTS: usize = 64; diff --git a/crates/core/src/config.rs b/crates/core/src/config.rs index 25874eb..044843a 100644 --- a/crates/core/src/config.rs +++ b/crates/core/src/config.rs @@ -68,7 +68,7 @@ pub struct Paths { pub work_dir: PathBuf, /// Directory containing bundled iPXE binaries (undionly.kpxe, snponly.efi, ...). pub ipxe_dir: PathBuf, - /// Path to the wimboot binary for Windows ISOs (optional — feature-gated). + /// Path to the wimboot binary for Windows ISOs (optional — feature-controlled). pub wimboot_path: Option, /// Directory under which Windows ISOs are extracted and served via SMB. /// Only used when `settings.windows_enabled = true`. Defaults to @@ -128,16 +128,24 @@ impl Config { /// Call this after loading the TOML file so env takes precedence. pub fn apply_env(&mut self) { if let Ok(v) = std::env::var("OPENPXE_HTTP_PORT") { - if let Ok(p) = v.parse() { self.server.http_port = p; } + if let Ok(p) = v.parse() { + self.server.http_port = p; + } } if let Ok(v) = std::env::var("OPENPXE_TFTP_PORT") { - if let Ok(p) = v.parse() { self.server.tftp_port = p; } + if let Ok(p) = v.parse() { + self.server.tftp_port = p; + } } if let Ok(v) = std::env::var("OPENPXE_DHCP_PORT") { - if let Ok(p) = v.parse() { self.network.dhcp_port = p; } + if let Ok(p) = v.parse() { + self.network.dhcp_port = p; + } } if let Ok(v) = std::env::var("OPENPXE_PUBLIC_IP") { - if let Ok(ip) = v.parse() { self.server.public_ip = Some(ip); } + if let Ok(ip) = v.parse() { + self.server.public_ip = Some(ip); + } } if let Ok(v) = std::env::var("OPENPXE_DHCP_MODE") { self.network.dhcp_mode = match v.to_ascii_lowercase().as_str() { diff --git a/crates/core/src/host_bindings.rs b/crates/core/src/host_bindings.rs index 2d152bf..05a7194 100644 --- a/crates/core/src/host_bindings.rs +++ b/crates/core/src/host_bindings.rs @@ -29,7 +29,7 @@ pub struct HostBinding { /// don't have to worry about case. pub mac: String, /// Preferred boot entry id (matches a `BootEntry::id` in the iso - /// store) OR one of the reserved menu names: `_local`, `_gate`, + /// store) OR one of the reserved menu names: `_local`, `_queue`, /// `_tools_menu`. Empty string falls back to the menu. pub target: String, /// Optional human-readable label shown in the UI (`"Tom's laptop"`, diff --git a/crates/core/src/lib.rs b/crates/core/src/lib.rs index efc0283..3002f21 100644 --- a/crates/core/src/lib.rs +++ b/crates/core/src/lib.rs @@ -6,18 +6,18 @@ pub mod arch; pub mod client; pub mod config; pub mod error; -pub mod queue; pub mod host_bindings; pub mod log_bus; pub mod metrics; +pub mod queue; pub mod settings; pub use arch::{ClientArch, FirmwareClass}; pub use client::{ClientEvent, ClientRegistry, ClientSnapshot}; pub use config::{Config, DhcpMode, NetworkConfig, Paths, ServerConfig}; pub use error::{Error, Result}; -pub use queue::{Gate, DeploymentQueue}; pub use host_bindings::{normalize_mac, HostBinding, HostBindings}; pub use log_bus::{LogBus, LogBusLayer, LogLine}; pub use metrics::{HttpRoute, Metrics}; +pub use queue::{DeploymentQueue, QueueEntry}; pub use settings::{Settings, SettingsStore, TimeoutAction}; diff --git a/crates/core/src/metrics.rs b/crates/core/src/metrics.rs index 438dd1e..e7db480 100644 --- a/crates/core/src/metrics.rs +++ b/crates/core/src/metrics.rs @@ -87,7 +87,9 @@ impl Metrics { } pub fn record_tftp_err(&self) { - self.inner.tftp_transfers_err.fetch_add(1, Ordering::Relaxed); + self.inner + .tftp_transfers_err + .fetch_add(1, Ordering::Relaxed); } // ── HTTP ─────────────────────────────────────────────────────────── @@ -181,7 +183,10 @@ impl Metrics { "", ); - let _ = writeln!(out, "# HELP openpxe_tftp_transfers_total TFTP transfers, by status."); + let _ = writeln!( + out, + "# HELP openpxe_tftp_transfers_total TFTP transfers, by status." + ); let _ = writeln!(out, "# TYPE openpxe_tftp_transfers_total counter"); let _ = writeln!( out, @@ -201,7 +206,10 @@ impl Metrics { "", ); - let _ = writeln!(out, "# HELP openpxe_http_requests_total HTTP requests served, by route family."); + let _ = writeln!( + out, + "# HELP openpxe_http_requests_total HTTP requests served, by route family." + ); let _ = writeln!(out, "# TYPE openpxe_http_requests_total counter"); for (label, counter) in [ ("boot_script", &i.http_boot_script), @@ -218,14 +226,53 @@ impl Metrics { } // Gauges. - write_gauge(&mut out, "openpxe_iso_count", "ISOs currently registered (local + NFS).", i.iso_count.load(Ordering::Relaxed), ""); - write_gauge(&mut out, "openpxe_client_count", "PXE clients seen this process lifetime.", i.client_count.load(Ordering::Relaxed), ""); - write_gauge(&mut out, "openpxe_queue_count", "Clients currently waiting at the deployment queue.", i.queue_count.load(Ordering::Relaxed), ""); - write_gauge(&mut out, "openpxe_queue_imaging", "Clients currently imaging (queue + assigned target).", i.queue_imaging.load(Ordering::Relaxed), ""); - write_gauge(&mut out, "openpxe_nfs_mounts_active", "NFS shares currently mounted.", i.nfs_mounts_active.load(Ordering::Relaxed), ""); - write_gauge(&mut out, "openpxe_uptime_seconds", "Seconds since this OpenPXE instance started.", uptime_secs, ""); + write_gauge( + &mut out, + "openpxe_iso_count", + "ISOs currently registered (local + NFS).", + i.iso_count.load(Ordering::Relaxed), + "", + ); + write_gauge( + &mut out, + "openpxe_client_count", + "PXE clients seen this process lifetime.", + i.client_count.load(Ordering::Relaxed), + "", + ); + write_gauge( + &mut out, + "openpxe_queue_count", + "Clients currently waiting at the deployment queue.", + i.queue_count.load(Ordering::Relaxed), + "", + ); + write_gauge( + &mut out, + "openpxe_queue_imaging", + "Clients currently imaging (queue + assigned target).", + i.queue_imaging.load(Ordering::Relaxed), + "", + ); + write_gauge( + &mut out, + "openpxe_nfs_mounts_active", + "NFS shares currently mounted.", + i.nfs_mounts_active.load(Ordering::Relaxed), + "", + ); + write_gauge( + &mut out, + "openpxe_uptime_seconds", + "Seconds since this OpenPXE instance started.", + uptime_secs, + "", + ); - let _ = writeln!(out, "# HELP openpxe_build_info Build metadata. Always 1; the version is in the label."); + let _ = writeln!( + out, + "# HELP openpxe_build_info Build metadata. Always 1; the version is in the label." + ); let _ = writeln!(out, "# TYPE openpxe_build_info gauge"); let _ = writeln!(out, "openpxe_build_info{{version=\"{version}\"}} 1"); @@ -259,7 +306,11 @@ mod tests { m.record_http(HttpRoute::Api); m.set_iso_count(3); let out = m.render("0.2.0", 42); - assert_eq!(out.matches("# TYPE openpxe_dhcp_replies_total counter").count(), 1); + assert_eq!( + out.matches("# TYPE openpxe_dhcp_replies_total counter") + .count(), + 1 + ); assert_eq!(out.matches("# TYPE openpxe_iso_count gauge").count(), 1); assert!(out.contains("openpxe_dhcp_replies_total{arch=\"uefi\"} 1")); assert!(out.contains("openpxe_dhcp_replies_total{arch=\"bios\"} 1")); diff --git a/crates/core/src/queue.rs b/crates/core/src/queue.rs index 396a616..dcf9b29 100644 --- a/crates/core/src/queue.rs +++ b/crates/core/src/queue.rs @@ -1,16 +1,16 @@ //! Queued Deployment queue. //! //! When a client selects "Queued Deployment" at the PXE menu, iPXE POSTs to -//! `/api/queue/join` and receives a gate position. It then enters a poll +//! `/api/queue/join` and receives a queue position. It then enters a poll //! loop hitting `/api/queue/poll/`; the server holds the request open //! until either (a) the operator assigns an ISO from the WebUI, in which //! case the poll returns an iPXE `chain` URL, or (b) the poll times out //! (iPXE's HTTP client has its own timeout), in which case iPXE re-POSTs. //! -//! The WebUI shows the queue (`GET /api/gate`) and issues +//! The WebUI shows the queue (`GET /api/queue`) and issues //! `POST /api/queue/assign { iso_id, entry_ids: [...] }` to launch a single -//! ISO across many gated clients at once. This is the "horse-race gate" -//! UX the user asked for — every horse leaves the line simultaneously. +//! ISO across many queued clients at once. Every waiting machine receives +//! the assignment without operator visits at the rack. use parking_lot::RwLock; use serde::{Deserialize, Serialize}; @@ -23,11 +23,11 @@ use uuid::Uuid; use crate::ClientArch; -/// Per-gate state visible to the WebUI. +/// Per-client queue state visible to the WebUI. #[derive(Debug, Clone, Serialize, Deserialize)] -pub struct Gate { +pub struct QueueEntry { pub id: String, - /// 1-based race-gate position — position 1 is whoever got there first. + /// 1-based queue position — position 1 is whoever got there first. pub position: u32, pub mac: String, pub ip: Option, @@ -55,8 +55,8 @@ struct QueueEntryInner { } impl QueueEntryInner { - fn snapshot(&self) -> Gate { - Gate { + fn snapshot(&self) -> QueueEntry { + QueueEntry { id: self.id.clone(), position: self.position, mac: self.mac.clone(), @@ -80,21 +80,25 @@ impl DeploymentQueue { Arc::new(Self::default()) } - /// Add a client to the gate. Returns the new `Gate` snapshot. If the - /// MAC is already queued, the existing gate is returned unchanged — + /// Add a client to the queue. Returns the current queue snapshot. If the + /// MAC is already queued, the existing entry is returned unchanged — /// retrying iPXE clients don't duplicate their slot. - pub fn join(&self, mac: &str, ip: Option, arch: Option) -> Gate { + pub fn join(&self, mac: &str, ip: Option, arch: Option) -> QueueEntry { let now = OffsetDateTime::now_utc(); let mut guard = self.inner.write(); if let Some(existing) = guard.values_mut().find(|g| g.mac == mac) { existing.last_poll_at = now; - if ip.is_some() { existing.ip = ip; } - if arch.is_some() { existing.arch = arch; } + if ip.is_some() { + existing.ip = ip; + } + if arch.is_some() { + existing.arch = arch; + } return existing.snapshot(); } - // Race position = max(position) + 1, or 1 if empty. + // Queue position = max(position) + 1, or 1 if empty. let next_pos = guard.values().map(|g| g.position).max().unwrap_or(0) + 1; let id = Uuid::new_v4().to_string(); let inner = QueueEntryInner { @@ -113,24 +117,24 @@ impl DeploymentQueue { snap } - /// Look up the `Notify` primitive for a given gate id, for long-polling. + /// Look up the `Notify` primitive for a given queue entry id, for long-polling. #[must_use] pub fn notifier(&self, entry_id: &str) -> Option> { self.inner.read().get(entry_id).map(|g| g.notify.clone()) } - /// Update the last-poll timestamp (keeps the gate's "live" indicator + /// Update the last-poll timestamp (keeps the queue's "live" indicator /// fresh in the UI) and return the current snapshot. Returns None if - /// the gate was released/expired between requests. - pub fn touch(&self, entry_id: &str) -> Option { + /// the entry was released/expired between requests. + pub fn touch(&self, entry_id: &str) -> Option { let mut guard = self.inner.write(); let g = guard.get_mut(entry_id)?; g.last_poll_at = OffsetDateTime::now_utc(); Some(g.snapshot()) } - /// Operator assigns an ISO entry (boot_entry id) to one or more gates. - /// Returns the number of gates that were updated. Gates not in the + /// Operator assigns an ISO entry (boot_entry id) to one or more clients. + /// Returns the number of queue entries that were updated. Entries not in the /// queue are silently skipped. pub fn assign(&self, entry_ids: &[String], target: &str) -> usize { let mut guard = self.inner.write(); @@ -145,9 +149,9 @@ impl DeploymentQueue { updated } - /// Remove a gate and return its final snapshot. Called after the client + /// Remove a queue entry and return its final snapshot. Called after the client /// has successfully chained onto its assignment. - pub fn release(&self, entry_id: &str) -> Option { + pub fn release(&self, entry_id: &str) -> Option { let mut guard = self.inner.write(); let g = guard.remove(entry_id)?; g.notify.notify_waiters(); @@ -162,7 +166,7 @@ impl DeploymentQueue { } #[must_use] - pub fn list(&self) -> Vec { + pub fn list(&self) -> Vec { let guard = self.inner.read(); let mut v: Vec<_> = guard.values().map(QueueEntryInner::snapshot).collect(); v.sort_by_key(|g| g.position); diff --git a/crates/core/src/settings.rs b/crates/core/src/settings.rs index 9f744b0..75651c4 100644 --- a/crates/core/src/settings.rs +++ b/crates/core/src/settings.rs @@ -48,9 +48,9 @@ pub struct Settings { pub default_local_hdd: bool, /// When a client hits the Queued Deployment item, how long (seconds) to - /// hold it at the gate before giving up and falling back to the menu. + /// hold it in queue before giving up and falling back to the menu. /// 0 = forever. - pub gate_wait_max_secs: u32, + pub queue_wait_max_secs: u32, /// Optional DNS server advertised on the Network tab. Purely /// informational today — OpenPXE does not run a DNS server, but @@ -67,11 +67,8 @@ pub enum TimeoutAction { /// Chain the "Boot from Local HDD" entry. LocalHdd, /// Put the client into the deployment queue, waiting for operator - /// assignment. The serde alias keeps v0.2.0 settings.json files - /// readable after the v0.3.0 rename — old `"gated_deployment"` - /// values deserialize transparently. + /// assignment. #[default] - #[serde(alias = "gated_deployment")] QueuedDeployment, } @@ -84,7 +81,7 @@ impl Default for Settings { smb_host_override: String::new(), extra_kernel_args: String::new(), default_local_hdd: true, - gate_wait_max_secs: 0, + queue_wait_max_secs: 0, dns_server: String::new(), } } @@ -115,7 +112,10 @@ impl SettingsStore { }, Err(_) => Settings::default(), }; - Arc::new(Self { path, inner: RwLock::new(initial) }) + Arc::new(Self { + path, + inner: RwLock::new(initial), + }) } #[must_use] @@ -181,6 +181,12 @@ mod tests { assert!(s.windows_enabled); } + #[test] + fn settings_serialize_queue_naming() { + let text = serde_json::to_string(&Settings::default()).unwrap(); + assert!(text.contains("queue_wait_max_secs")); + } + #[test] fn corrupt_file_falls_back_to_default() { let dir = tempdir().unwrap(); diff --git a/crates/dhcp-proxy/src/reply.rs b/crates/dhcp-proxy/src/reply.rs index 1f9930c..a0386fc 100644 --- a/crates/dhcp-proxy/src/reply.rs +++ b/crates/dhcp-proxy/src/reply.rs @@ -56,11 +56,17 @@ pub fn decide(ctx: &ReplyContext<'_>) -> BootDirective { // it'll then do the same script-fetch the iPXE path does. let name = ctx.arch.ipxe_bootfile().unwrap_or("snponly.efi"); BootDirective::HttpScript { - url: format!("{}/ipxe/{}", ctx.public_base_url.trim_end_matches('/'), name), + url: format!( + "{}/ipxe/{}", + ctx.public_base_url.trim_end_matches('/'), + name + ), } } FirmwareClass::PxeClient => match ctx.arch.ipxe_bootfile() { - Some(name) => BootDirective::TftpIpxe { filename: name.to_string() }, + Some(name) => BootDirective::TftpIpxe { + filename: name.to_string(), + }, None => BootDirective::Ignore, }, FirmwareClass::Other => BootDirective::Ignore, @@ -107,12 +113,16 @@ pub fn build_reply(ctx: &ReplyContext<'_>, directive: &BootDirective) -> Option< match directive { BootDirective::TftpIpxe { filename } => { - opts.insert(DhcpOption::TFTPServerName(ctx.our_ip.to_string().into_bytes())); + opts.insert(DhcpOption::TFTPServerName( + ctx.our_ip.to_string().into_bytes(), + )); opts.insert(DhcpOption::BootfileName(filename.as_bytes().to_vec())); } BootDirective::HttpScript { url } => { opts.insert(DhcpOption::BootfileName(url.as_bytes().to_vec())); - opts.insert(DhcpOption::TFTPServerName(ctx.our_ip.to_string().into_bytes())); + opts.insert(DhcpOption::TFTPServerName( + ctx.our_ip.to_string().into_bytes(), + )); } BootDirective::Ignore => return None, } diff --git a/crates/dhcp-proxy/src/server.rs b/crates/dhcp-proxy/src/server.rs index 0bd98f4..9994cc5 100644 --- a/crates/dhcp-proxy/src/server.rs +++ b/crates/dhcp-proxy/src/server.rs @@ -4,9 +4,7 @@ use crate::reply::{build_reply, decide, BootDirective, ReplyContext}; use dhcproto::v4::{DhcpOption, Message, OptionCode}; use dhcproto::{Decodable, Decoder, Encodable, Encoder}; -use openpxe_core::{ - ClientArch, ClientEvent, ClientRegistry, FirmwareClass, -}; +use openpxe_core::{ClientArch, ClientEvent, ClientRegistry, FirmwareClass}; use socket2::{Domain, Protocol, Socket, Type}; use std::net::{IpAddr, Ipv4Addr, SocketAddr, SocketAddrV4}; use std::sync::Arc; @@ -86,11 +84,22 @@ impl DhcpProxyServer { ) -> anyhow::Result<()> { let request = Message::decode(&mut Decoder::new(data))?; - let vendor_class = request.opts().get(OptionCode::ClassIdentifier).and_then(|o| { - if let DhcpOption::ClassIdentifier(v) = o { Some(v.as_slice()) } else { None } - }); + let vendor_class = request + .opts() + .get(OptionCode::ClassIdentifier) + .and_then(|o| { + if let DhcpOption::ClassIdentifier(v) = o { + Some(v.as_slice()) + } else { + None + } + }); let user_class = request.opts().get(OptionCode::UserClass).and_then(|o| { - if let DhcpOption::UserClass(v) = o { Some(v.as_slice()) } else { None } + if let DhcpOption::UserClass(v) = o { + Some(v.as_slice()) + } else { + None + } }); let class = FirmwareClass::classify(vendor_class, user_class); if matches!(class, FirmwareClass::Other) { @@ -135,7 +144,9 @@ impl DhcpProxyServer { } self.metrics.record_dhcp_reply(arch.as_str()); - let Some(reply) = build_reply(&ctx, &directive) else { return Ok(()); }; + let Some(reply) = build_reply(&ctx, &directive) else { + return Ok(()); + }; let mut out = Vec::with_capacity(512); reply.encode(&mut Encoder::new(&mut out))?; @@ -203,7 +214,10 @@ fn bind_udp(bind: IpAddr, port: u16, broadcast: bool) -> anyhow::Result String { let take = chaddr.iter().take(6).copied().collect::>(); - take.iter().map(|b| format!("{b:02x}")).collect::>().join(":") + take.iter() + .map(|b| format!("{b:02x}")) + .collect::>() + .join(":") } /// Walk raw DHCP options looking for option 93 (Client System Architecture) @@ -217,10 +231,17 @@ fn extract_raw_arch(packet: &[u8]) -> Option { let mut i = 0; while i < opts.len() { let code = opts[i]; - if code == 0xff { return None; } // END - if code == 0x00 { i += 1; continue; } // PAD + if code == 0xff { + return None; + } // END + if code == 0x00 { + i += 1; + continue; + } // PAD i += 1; - if i >= opts.len() { return None; } + if i >= opts.len() { + return None; + } let len = opts[i] as usize; i += 1; if code == 93 && len >= 2 && i + 2 <= opts.len() { @@ -240,7 +261,7 @@ mod tests { // Minimal BOOTP header + magic cookie + option 93 (arch)=0x0007 + END. let mut pkt = vec![0u8; 240]; pkt[236..240].copy_from_slice(&[99, 130, 83, 99]); // magic cookie - pkt.extend_from_slice(&[53, 1, 1]); // option 53 DHCPDISCOVER + pkt.extend_from_slice(&[53, 1, 1]); // option 53 DHCPDISCOVER pkt.extend_from_slice(&[93, 2, 0x00, 0x07]); pkt.push(0xff); assert_eq!(extract_raw_arch(&pkt), Some(0x0007)); diff --git a/crates/http-api/src/app.rs b/crates/http-api/src/app.rs index e20688a..c4480ae 100644 --- a/crates/http-api/src/app.rs +++ b/crates/http-api/src/app.rs @@ -14,8 +14,8 @@ //! | `/api/*` | JSON/HTML API for the web UI | use crate::ipxe_script::{ - render_entry, render_family_menu, render_queue_entry, render_local_hdd, - render_menu, render_nic_info, render_shell, render_tools_menu, render_util, + render_entry, render_family_menu, render_local_hdd, render_menu, render_nic_info, + render_queue_entry, render_shell, render_tools_menu, render_util, }; use crate::iso_fs; use crate::log_stream; @@ -61,7 +61,7 @@ pub fn build_router(state: AppState) -> Router { // JSON API. .route("/api/isos", get(api_list_isos).post(api_upload_iso)) .route("/api/isos/:id", delete(api_delete_iso)) - // v0.3.1: per-ISO password gate. PUT body `{ "password": "..." }` + // Per-ISO password prompt. PUT body `{ "password": "..." }` // sets, `{ "password": null }` (or DELETE) clears. .route( "/api/isos/:id/password", @@ -105,42 +105,67 @@ pub fn build_router(state: AppState) -> Router { async fn index(State(state): State) -> Response { let html = openpxe_webui::index_html(&state.public_base_url); - ([(header::CONTENT_TYPE, HeaderValue::from_static("text/html; charset=utf-8"))], html) + ( + [( + header::CONTENT_TYPE, + HeaderValue::from_static("text/html; charset=utf-8"), + )], + html, + ) .into_response() } async fn ui_js() -> Response { ( - [(header::CONTENT_TYPE, HeaderValue::from_static("application/javascript"))], + [( + header::CONTENT_TYPE, + HeaderValue::from_static("application/javascript"), + )], openpxe_webui::app_js(), - ).into_response() + ) + .into_response() } async fn ui_css() -> Response { ( [(header::CONTENT_TYPE, HeaderValue::from_static("text/css"))], openpxe_webui::app_css(), - ).into_response() + ) + .into_response() } async fn ui_logo() -> Response { ( - [(header::CONTENT_TYPE, HeaderValue::from_static("image/svg+xml"))], + [( + header::CONTENT_TYPE, + HeaderValue::from_static("image/svg+xml"), + )], openpxe_webui::logo_svg(), - ).into_response() + ) + .into_response() } async fn ui_loader() -> Response { ( - [(header::CONTENT_TYPE, HeaderValue::from_static("image/svg+xml"))], + [( + header::CONTENT_TYPE, + HeaderValue::from_static("image/svg+xml"), + )], openpxe_webui::loader_svg(), - ).into_response() + ) + .into_response() } // ─── iPXE scripts ────────────────────────────────────────────────────────── fn text_plain(body: String) -> Response { - ([(header::CONTENT_TYPE, HeaderValue::from_static("text/plain; charset=utf-8"))], body) + ( + [( + header::CONTENT_TYPE, + HeaderValue::from_static("text/plain; charset=utf-8"), + )], + body, + ) .into_response() } @@ -148,11 +173,10 @@ fn text_plain(body: String) -> Response { /// requesting client carries a `?mac=...` query param (iPXE's `${mac}` /// substitution) and that MAC has a binding, we short-circuit straight /// to the bound target instead of rendering the menu. -async fn boot_top_menu( - State(state): State, - Query(p): Query, -) -> Response { - state.metrics.record_http(openpxe_core::HttpRoute::BootScript); +async fn boot_top_menu(State(state): State, Query(p): Query) -> Response { + state + .metrics + .record_http(openpxe_core::HttpRoute::BootScript); let isos = state.iso_store.list(); let settings = state.settings.snapshot(); let base = &state.public_base_url; @@ -212,19 +236,19 @@ async fn boot_sub( let settings = state.settings.snapshot(); let base = &state.public_base_url; let script = match name { - "_local" => render_local_hdd(base), - "_linux_menu" => render_family_menu(&isos, base, false), - "_windows_menu" => render_family_menu(&isos, base, true), - "_tools_menu" => render_tools_menu(base), - "_util" => render_util(base), - "_shell" => render_shell(base), - "_nic" => render_nic_info(base), - "_queue" => render_queue_entry(base), + "_local" => render_local_hdd(base), + "_linux_menu" => render_family_menu(&isos, base, false), + "_windows_menu" => render_family_menu(&isos, base, true), + "_tools_menu" => render_tools_menu(base), + "_util" => render_util(base), + "_shell" => render_shell(base), + "_nic" => render_nic_info(base), + "_queue" => render_queue_entry(base), other => { for iso in &isos { for entry in &iso.boot_entries { if entry.id == other { - // Password gate. If the ISO has a password set + // Password prompt. If the ISO has a password set // we block the actual boot script behind it: // - no token -> render a prompt // - wrong token -> render auth-fail @@ -235,38 +259,45 @@ async fn boot_sub( match p.token.as_deref() { None | Some("") => { return text_plain(crate::ipxe_script::render_password_prompt( - &entry.id, &iso.filename, base, + &entry.id, + &iso.filename, + base, )); } - Some(token) => match state.iso_store.verify_password(&iso.id, token) { - Ok(true) => { /* fall through to render the entry */ } - Ok(false) => { - // Don't log the candidate — just the - // mac (when iPXE supplies one) and - // the entry id, so an operator can - // see brute-force attempts in the - // live log. - tracing::warn!( - target: "openpxe::http::boot", - entry = %other, - "wrong password supplied for protected boot entry" - ); - return text_plain( - crate::ipxe_script::render_password_failed( - &entry.id, base, - ), - ); + Some(token) => { + match state.iso_store.verify_password(&iso.id, token) { + Ok(true) => { /* fall through to render the entry */ } + Ok(false) => { + // Don't log the candidate — just the + // mac (when iPXE supplies one) and + // the entry id, so an operator can + // see brute-force attempts in the + // live log. + tracing::warn!( + target: "openpxe::http::boot", + entry = %other, + "wrong password supplied for protected boot entry" + ); + return text_plain( + crate::ipxe_script::render_password_failed( + &entry.id, base, + ), + ); + } + Err(e) => { + tracing::error!( + target: "openpxe::http::boot", + entry = %other, error = %e, + "password verify failed unexpectedly" + ); + return ( + StatusCode::INTERNAL_SERVER_ERROR, + "password check failed", + ) + .into_response(); + } } - Err(e) => { - tracing::error!( - target: "openpxe::http::boot", - entry = %other, error = %e, - "password verify failed unexpectedly" - ); - return (StatusCode::INTERNAL_SERVER_ERROR, - "password check failed").into_response(); - } - }, + } } } return text_plain(render_entry(entry, &settings, base)); @@ -290,11 +321,15 @@ async fn ipxe_binary(AxumPath(name): AxumPath) -> Response { }; ( [ - (header::CONTENT_TYPE, HeaderValue::from_static("application/octet-stream")), + ( + header::CONTENT_TYPE, + HeaderValue::from_static("application/octet-stream"), + ), (header::CONTENT_LENGTH, HeaderValue::from(bytes.len())), ], bytes, - ).into_response() + ) + .into_response() } // ─── ISO streaming (raw + in-ISO) ───────────────────────────────────────── @@ -324,7 +359,9 @@ async fn iso_file( let p = iso_path.clone(); let in_path = format!("/{path}"); let loc = tokio::task::spawn_blocking(move || iso_fs::lookup(&p, &in_path)) - .await.ok().flatten(); + .await + .ok() + .flatten(); let Some(loc) = loc else { return (StatusCode::NOT_FOUND, "not found inside iso").into_response(); }; @@ -340,21 +377,43 @@ async fn stream_file_range( ) -> anyhow::Result { let meta = tokio::fs::metadata(path).await?; let total = meta.len(); - let (start, end, partial) = parse_range(range, total); + if total == 0 { + return Ok(Response::builder() + .status(StatusCode::OK) + .header(header::CONTENT_TYPE, "application/octet-stream") + .header(header::ACCEPT_RANGES, "bytes") + .header(header::CONTENT_LENGTH, 0) + .body(Body::empty()) + .unwrap()); + } + let Some((start, end, partial)) = parse_range(range, total) else { + return Ok(Response::builder() + .status(StatusCode::RANGE_NOT_SATISFIABLE) + .header(header::CONTENT_RANGE, format!("bytes */{total}")) + .body(Body::empty()) + .unwrap()); + }; let len = end - start + 1; let mut file = tokio::fs::File::open(path).await?; file.seek(std::io::SeekFrom::Start(start)).await?; let reader = file.take(len); let stream = tokio_util::io::ReaderStream::new(reader); let body = Body::from_stream(stream); - let status = if partial { StatusCode::PARTIAL_CONTENT } else { StatusCode::OK }; + let status = if partial { + StatusCode::PARTIAL_CONTENT + } else { + StatusCode::OK + }; let mut builder = Response::builder() .status(status) .header(header::CONTENT_TYPE, "application/octet-stream") .header(header::ACCEPT_RANGES, "bytes") .header(header::CONTENT_LENGTH, len); if partial { - builder = builder.header(header::CONTENT_RANGE, format!("bytes {start}-{end}/{total}")); + builder = builder.header( + header::CONTENT_RANGE, + format!("bytes {start}-{end}/{total}"), + ); } Ok(builder.body(body).unwrap()) } @@ -377,21 +436,40 @@ async fn stream_byte_range( .unwrap()) } -fn parse_range(h: Option<&HeaderValue>, total: u64) -> (u64, u64, bool) { - let Some(h) = h else { return (0, total.saturating_sub(1), false); }; - let Ok(s) = h.to_str() else { return (0, total.saturating_sub(1), false); }; - let Some(spec) = s.strip_prefix("bytes=") else { return (0, total.saturating_sub(1), false); }; +fn parse_range(h: Option<&HeaderValue>, total: u64) -> Option<(u64, u64, bool)> { + let Some(h) = h else { + return Some((0, total.saturating_sub(1), false)); + }; + let Ok(s) = h.to_str() else { + return Some((0, total.saturating_sub(1), false)); + }; + let Some(spec) = s.strip_prefix("bytes=") else { + return Some((0, total.saturating_sub(1), false)); + }; let spec = spec.split(',').next().unwrap_or("").trim(); if let Some(suffix) = spec.strip_prefix('-') { if let Ok(n) = suffix.parse::() { let n = n.min(total); - return (total.saturating_sub(n), total.saturating_sub(1), true); + return Some((total.saturating_sub(n), total.saturating_sub(1), true)); } } let mut parts = spec.splitn(2, '-'); - let start = parts.next().and_then(|s| s.parse::().ok()).unwrap_or(0); - let end = parts.next().and_then(|s| s.parse::().ok()).unwrap_or(total.saturating_sub(1)); - (start, end.min(total.saturating_sub(1)), true) + let start = parts + .next() + .and_then(|s| s.parse::().ok()) + .unwrap_or(0); + let end = parts + .next() + .and_then(|s| s.parse::().ok()) + .unwrap_or(total.saturating_sub(1)); + if start >= total { + return None; + } + let end = end.min(total.saturating_sub(1)); + if start > end { + return None; + } + Some((start, end, true)) } // ─── ISO upload / list / delete ─────────────────────────────────────────── @@ -442,8 +520,10 @@ async fn api_set_iso_password( ); StatusCode::NO_CONTENT.into_response() } - Err(pxeforge_error_invalid) if matches!(pxeforge_error_invalid, openpxe_core::Error::Invalid(_)) => { - (StatusCode::NOT_FOUND, format!("{pxeforge_error_invalid}")).into_response() + Err(openpxe_error_invalid) + if matches!(openpxe_error_invalid, openpxe_core::Error::Invalid(_)) => + { + (StatusCode::NOT_FOUND, format!("{openpxe_error_invalid}")).into_response() } Err(e) => (StatusCode::INTERNAL_SERVER_ERROR, format!("{e}")).into_response(), } @@ -465,12 +545,11 @@ async fn api_clear_iso_password( } } -async fn api_upload_iso( - State(state): State, - mut multipart: Multipart, -) -> Response { +async fn api_upload_iso(State(state): State, mut multipart: Multipart) -> Response { while let Ok(Some(mut field)) = multipart.next_field().await { - if field.name() != Some("file") { continue; } + if field.name() != Some("file") { + continue; + } let filename = field.file_name().unwrap_or("uploaded.iso").to_string(); if !filename.to_ascii_lowercase().ends_with(".iso") { return (StatusCode::BAD_REQUEST, "only .iso uploads accepted").into_response(); @@ -499,7 +578,11 @@ async fn api_upload_iso( async fn healthz() -> Response { // Simple liveness — HTTP task is responsive. Does not touch storage or // other subsystems so we never fail for downstream reasons. - ([(header::CONTENT_TYPE, HeaderValue::from_static("text/plain"))], "ok\n").into_response() + ( + [(header::CONTENT_TYPE, HeaderValue::from_static("text/plain"))], + "ok\n", + ) + .into_response() } async fn readyz(State(state): State) -> Response { @@ -517,7 +600,11 @@ async fn readyz(State(state): State) -> Response { problems.push("iso directory not readable"); } if problems.is_empty() { - ([(header::CONTENT_TYPE, HeaderValue::from_static("text/plain"))], "ready\n").into_response() + ( + [(header::CONTENT_TYPE, HeaderValue::from_static("text/plain"))], + "ready\n", + ) + .into_response() } else { let body = format!("not ready:\n- {}\n", problems.join("\n- ")); (StatusCode::SERVICE_UNAVAILABLE, body).into_response() @@ -536,17 +623,22 @@ async fn api_status(State(state): State) -> Json { let nfs_active = nfs.iter().filter(|m| m.mounted).count(); let isos = state.iso_store.list(); let clients = state.clients.list(); - let gates = state.queue.list(); - // Phase 4: dashboard tracks "imaging" as gates with an assignment + let queue_entries = state.queue.list(); + // Phase 4: dashboard tracks "imaging" as queue entries with an assignment // already issued — they're the ones actively chaining a boot script. - let imaging = gates.iter().filter(|g| g.assigned_target.is_some()).count(); - let waiting = gates.len() - imaging; + let imaging = queue_entries + .iter() + .filter(|entry| entry.assigned_target.is_some()) + .count(); + let waiting = queue_entries.len() - imaging; // Side-effect: push gauge values out to the Prometheus surface. // Doing it here (in the most-frequently-polled endpoint) keeps the // gauges fresh without a dedicated scrape-time hook. state.metrics.set_iso_count(isos.len() as u64); state.metrics.set_client_count(clients.len() as u64); - state.metrics.set_queue_counts(gates.len() as u64, imaging as u64); + state + .metrics + .set_queue_counts(queue_entries.len() as u64, imaging as u64); state.metrics.set_nfs_active(nfs_active as u64); state.metrics.record_http(openpxe_core::HttpRoute::Api); let now = time::OffsetDateTime::now_utc(); @@ -556,7 +648,7 @@ async fn api_status(State(state): State) -> Json { "public_base_url": state.public_base_url, "iso_count": isos.len(), "client_count": clients.len(), - "queue_count": gates.len(), + "queue_count": queue_entries.len(), "imaging_count": imaging, "waiting_count": waiting, "ipxe_assets": openpxe_ipxe_assets::list_assets(), @@ -591,7 +683,8 @@ async fn api_put_settings( "cannot enable Windows: 'wimboot' binary is not bundled. \ Place a signed wimboot build at assets/ipxe/wimboot and rebuild \ the container. See docs/architecture.md for details.", - ).into_response(); + ) + .into_response(); } } new.smb_host_override = new.smb_host_override.trim().to_string(); @@ -604,9 +697,15 @@ async fn api_put_settings( if let Some(smb) = &state.smb { match (was_enabled, want_enabled) { - (false, true) => { let _ = smb.start(); } - (true, false) => { smb.stop(); } - (true, true) => { let _ = smb.reconcile(); } + (false, true) => { + let _ = smb.start(); + } + (true, false) => { + smb.stop(); + } + (true, true) => { + let _ = smb.reconcile(); + } (false, false) => {} } } @@ -623,7 +722,7 @@ async fn api_list_queue(State(state): State) -> Json, @@ -634,7 +733,7 @@ struct GateJoinParams { /// until poll returns an actual boot script. async fn api_queue_join( State(state): State, - Query(p): Query, + Query(p): Query, headers: HeaderMap, ) -> Response { let mac = p.mac.unwrap_or_else(|| "unknown".to_string()); @@ -644,10 +743,14 @@ async fn api_queue_join( .and_then(|s| s.split(',').next()) .and_then(|s| s.trim().parse().ok()); - let gate = state.queue.join(&mac, ip, None); + let queue_entry = state.queue.join(&mac, ip, None); state.clients.record( - &mac, ip, None, - ClientEvent::HttpScriptFetch { target: "queue-join".into() }, + &mac, + ip, + None, + ClientEvent::HttpScriptFetch { + target: "queue-join".into(), + }, ); let base = &state.public_base_url; @@ -656,12 +759,12 @@ async fn api_queue_join( "#!ipxe\n\ echo\n\ echo ==========================================\n\ - echo Queued Deployment - Gate Position {}\n\ + echo Queued Deployment - Queue Position {}\n\ echo Waiting for operator to assign an image\n\ echo (Ctrl-B returns to the iPXE shell)\n\ echo ==========================================\n\ chain {base}/api/queue/poll/{}\n", - gate.position, gate.id + queue_entry.position, queue_entry.id ); text_plain(script) } @@ -674,7 +777,7 @@ async fn api_queue_poll( AxumPath(entry_id): AxumPath, ) -> Response { let Some(notify) = state.queue.notifier(&entry_id) else { - // Gate was released; send client back to the main menu. + // Queue entry was released; send client back to the main menu. let base = &state.public_base_url; return text_plain(format!("#!ipxe\nchain {base}/boot.ipxe\n")); }; @@ -687,7 +790,7 @@ async fn api_queue_poll( match snap { // Bind `target` directly so we can't observe an Option::None between // the guard and the unwrap (the old code had a race with concurrent - // `release`). We also do NOT release the gate here — the web UI + // `release`). We also do NOT release the queue entry here — the web UI // operator releases it explicitly, which keeps a record of "this // machine was assigned image X" visible until the client is known // to have started. Clients that retry on transient network errors @@ -698,20 +801,20 @@ async fn api_queue_poll( tracing::info!( target: "openpxe::queue", entry_id=%entry_id, mac=%g.mac, target=%target, - "gate assignment delivered" + "queue assignment delivered" ); text_plain(format!( "#!ipxe\n\ - echo Gate assignment received: {target}\n\ + echo Queue assignment received: {target}\n\ chain {base}/boot/{target}.ipxe || chain {base}/api/queue/poll/{entry_id}\n" )) } Some(g) => { // No assignment yet - loop and re-poll. Repaint position so the - // UI count stays accurate if other gates were released meanwhile. + // UI count stays accurate if other queue entries were released meanwhile. text_plain(format!( "#!ipxe\n\ - echo Gate Position {} - still waiting\n\ + echo Queue Position {} - still waiting\n\ chain {base}/api/queue/poll/{entry_id}\n", g.position )) @@ -721,27 +824,34 @@ async fn api_queue_poll( } #[derive(Debug, Deserialize)] -struct GateAssignBody { +struct QueueAssignBody { /// Boot entry id (from `BootEntry::id`). Same one used in /// `/boot/.ipxe`. target: String, - /// Gate ids to assign. Empty = assign to all currently queued gates. + /// Queue entry ids to assign. Empty = assign to all currently queued clients. entry_ids: Vec, } async fn api_queue_assign( State(state): State, - Json(body): Json, + Json(body): Json, ) -> Json { let ids = if body.entry_ids.is_empty() { - state.queue.list().into_iter().map(|g| g.id).collect::>() + state + .queue + .list() + .into_iter() + .map(|g| g.id) + .collect::>() } else { body.entry_ids }; // Guard: target must exist as a BootEntry id. - let found = state.iso_store.list().into_iter().any(|i| { - i.boot_entries.iter().any(|e| e.id == body.target) - }); + let found = state + .iso_store + .list() + .into_iter() + .any(|i| i.boot_entries.iter().any(|e| e.id == body.target)); if !found { return Json(json!({ "ok": false, "error": "unknown target" })); } @@ -765,10 +875,7 @@ async fn api_nfs_list(State(state): State) -> Json Json(json!({ "mounts": state.nfs.list() })) } -async fn api_nfs_add( - State(state): State, - Json(req): Json, -) -> Response { +async fn api_nfs_add(State(state): State, Json(req): Json) -> Response { match state.nfs.add(req).await { Ok(m) => (StatusCode::CREATED, Json(m)).into_response(), // Anything from the manager surfaces as a user-fixable validation @@ -779,20 +886,14 @@ async fn api_nfs_add( } } -async fn api_nfs_remove( - State(state): State, - AxumPath(id): AxumPath, -) -> Response { +async fn api_nfs_remove(State(state): State, AxumPath(id): AxumPath) -> Response { match state.nfs.remove(&id).await { Ok(()) => StatusCode::NO_CONTENT.into_response(), Err(e) => (StatusCode::INTERNAL_SERVER_ERROR, format!("{e}")).into_response(), } } -async fn api_nfs_scan( - State(state): State, - AxumPath(id): AxumPath, -) -> Response { +async fn api_nfs_scan(State(state): State, AxumPath(id): AxumPath) -> Response { match state.nfs.rescan(&id).await { Ok(n) => Json(json!({ "ok": true, "iso_count": n })).into_response(), Err(e) => (StatusCode::BAD_REQUEST, format!("{e}")).into_response(), @@ -899,11 +1000,14 @@ async fn api_metrics(State(state): State) -> Response { state .metrics .set_client_count(state.clients.list().len() as u64); - let gates = state.queue.list(); - let imaging = gates.iter().filter(|g| g.assigned_target.is_some()).count(); + let queue_entries = state.queue.list(); + let imaging = queue_entries + .iter() + .filter(|entry| entry.assigned_target.is_some()) + .count(); state .metrics - .set_queue_counts(gates.len() as u64, imaging as u64); + .set_queue_counts(queue_entries.len() as u64, imaging as u64); state .metrics .set_nfs_active(state.nfs.list().iter().filter(|m| m.mounted).count() as u64); @@ -927,25 +1031,39 @@ mod tests { #[test] fn range_full() { - let (s, e, p) = parse_range(None, 1000); + let (s, e, p) = parse_range(None, 1000).unwrap(); assert_eq!((s, e, p), (0, 999, false)); } #[test] fn range_open_ended() { let h = HeaderValue::from_static("bytes=500-"); - let (s, e, p) = parse_range(Some(&h), 1000); + let (s, e, p) = parse_range(Some(&h), 1000).unwrap(); assert_eq!((s, e, p), (500, 999, true)); } #[test] fn range_suffix() { let h = HeaderValue::from_static("bytes=-100"); - let (s, e, p) = parse_range(Some(&h), 1000); + let (s, e, p) = parse_range(Some(&h), 1000).unwrap(); assert_eq!((s, e, p), (900, 999, true)); } #[test] fn range_explicit() { let h = HeaderValue::from_static("bytes=10-99"); - let (s, e, p) = parse_range(Some(&h), 1000); + let (s, e, p) = parse_range(Some(&h), 1000).unwrap(); assert_eq!((s, e, p), (10, 99, true)); } + + #[test] + fn range_rejects_out_of_bounds_start() { + let h = HeaderValue::from_static("bytes=1000-"); + let got = parse_range(Some(&h), 1000); + assert_eq!(got, None); + } + + #[test] + fn range_rejects_start_after_end() { + let h = HeaderValue::from_static("bytes=99-10"); + let got = parse_range(Some(&h), 1000); + assert_eq!(got, None); + } } diff --git a/crates/http-api/src/ipxe_script.rs b/crates/http-api/src/ipxe_script.rs index 1860471..472ab14 100644 --- a/crates/http-api/src/ipxe_script.rs +++ b/crates/http-api/src/ipxe_script.rs @@ -23,8 +23,8 @@ //! There is intentionally no UI path to upload a custom `.ipxe` script. use openpxe_core::{Settings, TimeoutAction}; -use openpxe_iso_store::{BootEntry, BootKind, IsoMeta}; use openpxe_iso_store::introspect::DistroFamily; +use openpxe_iso_store::{BootEntry, BootKind, IsoMeta}; use std::fmt::Write as _; /// Top-level OpenPXE boot menu. Serialized identically for BIOS and UEFI @@ -49,9 +49,15 @@ pub fn render_menu(isos: &[IsoMeta], settings: &Settings, base_url: &str) -> Str let _ = writeln!(s, "set cls ${{esc:string}}[2J"); let _ = writeln!(s, ":menu"); let _ = writeln!(s, "menu OpenPXE - network boot menu"); - let _ = writeln!(s, "item --gap -- ------------------------- Default -------------------------"); + let _ = writeln!( + s, + "item --gap -- ------------------------- Default -------------------------" + ); let _ = writeln!(s, "item local Boot from Local HDD"); - let _ = writeln!(s, "item --gap -- ----------------------- Installers -----------------------"); + let _ = writeln!( + s, + "item --gap -- ----------------------- Installers -----------------------" + ); if has_family(isos, is_linux_family) { let _ = writeln!(s, "item linux Linux Installers >"); } else { @@ -64,9 +70,15 @@ pub fn render_menu(isos: &[IsoMeta], settings: &Settings, base_url: &str) -> Str } else { let _ = writeln!(s, "item --gap -- (Windows support disabled in Settings)"); } - let _ = writeln!(s, "item --gap -- -------------------------- Tools --------------------------"); + let _ = writeln!( + s, + "item --gap -- -------------------------- Tools --------------------------" + ); let _ = writeln!(s, "item tools Tools >"); - let _ = writeln!(s, "item --gap -- ---------------------- Queued Deployment ---------------------"); + let _ = writeln!( + s, + "item --gap -- ---------------------- Queued Deployment ---------------------" + ); let _ = writeln!(s, "item queue Queued Deployment (join queue)"); let _ = writeln!(s, "item --gap"); let _ = writeln!(s, "item --key x exit Exit iPXE"); @@ -74,18 +86,39 @@ pub fn render_menu(isos: &[IsoMeta], settings: &Settings, base_url: &str) -> Str if matches!(settings.timeout_action, TimeoutAction::Stay) { let _ = writeln!(s, "choose --default {default_item} target || goto menu"); } else { - let _ = writeln!(s, "choose --default {default_item} --timeout {timeout_ms} target || goto menu"); + let _ = writeln!( + s, + "choose --default {default_item} --timeout {timeout_ms} target || goto menu" + ); } // iPXE's `||` is strict about what follows. Each test uses `goto menu` // as the fallthrough target so the parser never sees a bare `||` with // trailing whitespace — some iPXE builds reject that. - let _ = writeln!(s, "iseq ${{target}} local && chain {base}/boot/_local.ipxe || goto menu"); - let _ = writeln!(s, "iseq ${{target}} linux && chain {base}/boot/_linux_menu.ipxe || goto menu"); - let _ = writeln!(s, "iseq ${{target}} windows && chain {base}/boot/_windows_menu.ipxe || goto menu"); - let _ = writeln!(s, "iseq ${{target}} tools && chain {base}/boot/_tools_menu.ipxe || goto menu"); - let _ = writeln!(s, "iseq ${{target}} queue && chain {base}/boot/_queue.ipxe || goto menu"); - let _ = writeln!(s, "iseq ${{target}} exit && exit || goto menu"); + let _ = writeln!( + s, + "iseq ${{target}} local && chain {base}/boot/_local.ipxe || goto menu" + ); + let _ = writeln!( + s, + "iseq ${{target}} linux && chain {base}/boot/_linux_menu.ipxe || goto menu" + ); + let _ = writeln!( + s, + "iseq ${{target}} windows && chain {base}/boot/_windows_menu.ipxe || goto menu" + ); + let _ = writeln!( + s, + "iseq ${{target}} tools && chain {base}/boot/_tools_menu.ipxe || goto menu" + ); + let _ = writeln!( + s, + "iseq ${{target}} queue && chain {base}/boot/_queue.ipxe || goto menu" + ); + let _ = writeln!( + s, + "iseq ${{target}} exit && exit || goto menu" + ); let _ = writeln!(s, "goto menu"); s } @@ -95,27 +128,41 @@ pub fn render_menu(isos: &[IsoMeta], settings: &Settings, base_url: &str) -> Str #[must_use] pub fn render_family_menu(isos: &[IsoMeta], base_url: &str, is_windows: bool) -> String { let base = base_url.trim_end_matches('/'); - let title = if is_windows { "Windows Installers" } else { "Linux Installers" }; + let title = if is_windows { + "Windows Installers" + } else { + "Linux Installers" + }; let label = if is_windows { "windows" } else { "linux" }; let mut s = String::new(); let _ = writeln!(s, "#!ipxe"); let _ = writeln!(s, "set base-url {base}"); let _ = writeln!(s, ":menu"); let _ = writeln!(s, "menu OpenPXE - {title}"); - let filter: fn(DistroFamily) -> bool = - if is_windows { is_windows_family } else { is_linux_family }; + let filter: fn(DistroFamily) -> bool = if is_windows { + is_windows_family + } else { + is_linux_family + }; let mut count = 0; for iso in isos { - if !filter(iso.introspection.family) { continue; } + if !filter(iso.introspection.family) { + continue; + } for entry in &iso.boot_entries { let size_label = fmt_size_mib(iso.size_bytes); let key = hotkey_for_index(count); // Visual hint: a leading `*` marks password-protected entries. // ASCII only — iPXE's menu console mangles non-ASCII on some // firmwares. - let lock = if iso.is_password_protected() { "*" } else { " " }; + let lock = if iso.is_password_protected() { + "*" + } else { + " " + }; let _ = writeln!( - s, "item {}{} {}[{:>6}] {}", + s, + "item {}{} {}[{:>6}] {}", key, entry.id, lock, @@ -131,7 +178,10 @@ pub fn render_family_menu(isos: &[IsoMeta], base_url: &str, is_windows: bool) -> let _ = writeln!(s, "item --gap"); let _ = writeln!(s, "item --key b back < Back to main menu"); let _ = writeln!(s, "choose target || goto menu"); - let _ = writeln!(s, "iseq ${{target}} back && chain {base}/boot.ipxe || goto menu"); + let _ = writeln!( + s, + "iseq ${{target}} back && chain {base}/boot.ipxe || goto menu" + ); let _ = writeln!(s, "chain {base}/boot/${{target}}.ipxe || goto menu"); s } @@ -171,12 +221,30 @@ pub fn render_tools_menu(base_url: &str) -> String { let _ = writeln!(s, "item --gap"); let _ = writeln!(s, "item --key b back < Back to main menu"); let _ = writeln!(s, "choose target || goto menu"); - let _ = writeln!(s, "iseq ${{target}} util && chain {base}/boot/_util.ipxe || goto menu"); - let _ = writeln!(s, "iseq ${{target}} shell && chain {base}/boot/_shell.ipxe || goto menu"); - let _ = writeln!(s, "iseq ${{target}} nic && chain {base}/boot/_nic.ipxe || goto menu"); - let _ = writeln!(s, "iseq ${{target}} reboot && reboot || goto menu"); - let _ = writeln!(s, "iseq ${{target}} firmware && exit 0 || goto menu"); - let _ = writeln!(s, "iseq ${{target}} back && chain {base}/boot.ipxe || goto menu"); + let _ = writeln!( + s, + "iseq ${{target}} util && chain {base}/boot/_util.ipxe || goto menu" + ); + let _ = writeln!( + s, + "iseq ${{target}} shell && chain {base}/boot/_shell.ipxe || goto menu" + ); + let _ = writeln!( + s, + "iseq ${{target}} nic && chain {base}/boot/_nic.ipxe || goto menu" + ); + let _ = writeln!( + s, + "iseq ${{target}} reboot && reboot || goto menu" + ); + let _ = writeln!( + s, + "iseq ${{target}} firmware && exit 0 || goto menu" + ); + let _ = writeln!( + s, + "iseq ${{target}} back && chain {base}/boot.ipxe || goto menu" + ); let _ = writeln!(s, "goto menu"); s } @@ -190,8 +258,15 @@ pub fn render_local_hdd(base_url: &str) -> String { let mut s = String::new(); let _ = writeln!(s, "#!ipxe"); let _ = writeln!(s, "# Boot from Local HDD - platform-sensitive"); - let _ = writeln!(s, "iseq ${{platform}} pcbios && sanboot --no-describe --drive 0x80 || "); - let _ = writeln!(s, "# UEFI path: fall through to the firmware's next boot entry"); + let _ = writeln!( + s, + "iseq ${{platform}} pcbios && sanboot --no-describe --drive 0x80 || goto uefi" + ); + let _ = writeln!(s, ":uefi"); + let _ = writeln!( + s, + "# UEFI path: fall through to the firmware's next boot entry" + ); let _ = writeln!(s, "exit 0"); let _ = writeln!(s, "# If the above exit returns, loop back to the main menu"); let _ = writeln!(s, "chain {base}/boot.ipxe"); @@ -212,8 +287,14 @@ pub fn render_util(base_url: &str) -> String { let _ = writeln!(s, "item --gap"); let _ = writeln!(s, "item back < Back"); let _ = writeln!(s, "choose target || goto menu"); - let _ = writeln!(s, "iseq ${{target}} memtest && chain {base}/ipxe/memtest.bin || "); - let _ = writeln!(s, "iseq ${{target}} back && chain {base}/boot/_tools_menu.ipxe || "); + let _ = writeln!( + s, + "iseq ${{target}} memtest && chain {base}/ipxe/memtest.bin || goto menu" + ); + let _ = writeln!( + s, + "iseq ${{target}} back && chain {base}/boot/_tools_menu.ipxe || goto menu" + ); let _ = writeln!(s, "goto menu"); s } @@ -260,7 +341,10 @@ pub fn render_queue_entry(base_url: &str) -> String { let base = base_url.trim_end_matches('/'); let mut s = String::new(); let _ = writeln!(s, "#!ipxe"); - let _ = writeln!(s, "# Queued Deployment - join the queue and wait for operator"); + let _ = writeln!( + s, + "# Queued Deployment - join the queue and wait for operator" + ); let _ = writeln!(s, "echo Joining deployment queue..."); // imgfetch writes the body to a file in iPXE's transient FS; we read // the queue entry id out of the Location-style header by asking the server @@ -277,7 +361,11 @@ pub fn render_entry(entry: &BootEntry, settings: &Settings, base_url: &str) -> S let _ = writeln!(s, "#!ipxe"); let _ = writeln!(s, "set base-url {base}"); match &entry.kind { - BootKind::LinuxKernel { kernel_url, initrd_urls, args } => { + BootKind::LinuxKernel { + kernel_url, + initrd_urls, + args, + } => { let mut cmdline = args.cmdline.replace("${base-url}", base); if !settings.extra_kernel_args.trim().is_empty() { cmdline.push(' '); @@ -328,7 +416,12 @@ fn has_family(isos: &[IsoMeta], pred: fn(DistroFamily) -> bool) -> bool { } fn escape_label(s: &str) -> String { - s.chars().map(|c| match c { '\n' | '\r' => ' ', c => c }).collect() + s.chars() + .map(|c| match c { + '\n' | '\r' => ' ', + c => c, + }) + .collect() } /// Render the password-prompt script for a protected boot entry. @@ -362,7 +455,10 @@ pub fn render_password_prompt(entry_id: &str, iso_filename: &str, base_url: &str let _ = writeln!(s, "echo =========================================="); let _ = writeln!(s, "set password "); let _ = writeln!(s, "read --secret password"); - let _ = writeln!(s, "iseq ${{password}} \"\" && chain {base}/boot.ipxe || goto submit"); + let _ = writeln!( + s, + "iseq ${{password}} \"\" && chain {base}/boot.ipxe || goto submit" + ); let _ = writeln!(s, ":submit"); let _ = writeln!(s, "echo Verifying..."); let _ = writeln!( @@ -384,7 +480,10 @@ pub fn render_password_failed(entry_id: &str, base_url: &str) -> String { let _ = writeln!(s, "echo"); let _ = writeln!(s, "echo Wrong password."); let _ = writeln!(s, "sleep 2"); - let _ = writeln!(s, "chain {base}/boot/{entry_id}.ipxe || chain {base}/boot.ipxe"); + let _ = writeln!( + s, + "chain {base}/boot/{entry_id}.ipxe || chain {base}/boot.ipxe" + ); s } @@ -414,4 +513,27 @@ mod password_tests { // Re-target the entry so the prompt flow runs again. assert!(s.contains("chain http://10.0.0.5/boot/alpha-linux.ipxe")); } + + #[test] + fn generated_scripts_do_not_emit_bare_or_trailing_fallbacks() { + let settings = Settings::default(); + let scripts = [ + render_menu(&[], &settings, "http://10.0.0.5"), + render_tools_menu("http://10.0.0.5"), + render_local_hdd("http://10.0.0.5"), + render_util("http://10.0.0.5"), + render_shell("http://10.0.0.5"), + render_nic_info("http://10.0.0.5"), + render_queue_entry("http://10.0.0.5"), + render_password_failed("alpha-linux", "http://10.0.0.5"), + ]; + for script in scripts { + for line in script.lines() { + assert!( + !line.trim_end().ends_with("||"), + "bare iPXE fallback operator in line: {line}\nscript:\n{script}" + ); + } + } + } } diff --git a/crates/http-api/src/iso_fs.rs b/crates/http-api/src/iso_fs.rs index 0ddaf76..12c2aa9 100644 --- a/crates/http-api/src/iso_fs.rs +++ b/crates/http-api/src/iso_fs.rs @@ -31,7 +31,9 @@ pub fn lookup(iso_path: &Path, in_iso_path: &str) -> Option { .split('/') .filter(|c| !c.is_empty()) .collect(); - if components.is_empty() { return None; } + if components.is_empty() { + return None; + } walk(&mut f, root.offset, root.length, &components) } @@ -40,11 +42,16 @@ fn read_root_directory(f: &mut std::fs::File) -> Option { let mut pvd = [0u8; 2048]; f.seek(SeekFrom::Start(16 * SECTOR)).ok()?; f.read_exact(&mut pvd).ok()?; - if pvd[0] != 0x01 || &pvd[1..6] != b"CD001" { return None; } + if pvd[0] != 0x01 || &pvd[1..6] != b"CD001" { + return None; + } // Root directory record is at offset 156, length 34. let rec = &pvd[156..156 + 34]; let (offset, length) = parse_dir_record_ext(rec)?; - Some(FileLocation { offset: offset * SECTOR, length }) + Some(FileLocation { + offset: offset * SECTOR, + length, + }) } /// Walk components down the directory tree starting at `dir_offset`. @@ -66,21 +73,29 @@ fn walk( if len == 0 { // Padding to sector boundary. let next = (i / SECTOR as usize + 1) * SECTOR as usize; - if next <= i { break; } + if next <= i { + break; + } i = next; continue; } - if i + len > dir.len() { break; } + if i + len > dir.len() { + break; + } let rec = &dir[i..i + len]; let name = dir_record_name(rec); let is_dir = (rec.get(25).copied().unwrap_or(0) & 0x02) != 0; // Skip "." (0x00) and ".." (0x01) pseudo-entries. - let is_pseudo = matches!(rec.get(32).copied(), Some(1)) && rec.get(33).copied() == Some(0x00) + let is_pseudo = matches!(rec.get(32).copied(), Some(1)) + && rec.get(33).copied() == Some(0x00) || matches!(rec.get(32).copied(), Some(1)) && rec.get(33).copied() == Some(0x01); if !is_pseudo && name.eq_ignore_ascii_case(target) { let (child_off, child_len) = parse_dir_record_ext(rec)?; if rest.is_empty() && !is_dir { - return Some(FileLocation { offset: child_off * SECTOR, length: child_len }); + return Some(FileLocation { + offset: child_off * SECTOR, + length: child_len, + }); } else if !rest.is_empty() && is_dir { return walk(f, child_off * SECTOR, child_len, rest); } @@ -94,7 +109,9 @@ fn walk( /// Layout per ISO9660: bytes 2..10 extent LBA (LE+BE duplicate), 10..18 /// data length (LE+BE duplicate). We trust the little-endian copy. fn parse_dir_record_ext(rec: &[u8]) -> Option<(u64, u64)> { - if rec.len() < 34 { return None; } + if rec.len() < 34 { + return None; + } let lba = u32::from_le_bytes(rec[2..6].try_into().ok()?) as u64; let len = u32::from_le_bytes(rec[10..14].try_into().ok()?) as u64; Some((lba, len)) @@ -104,9 +121,15 @@ fn parse_dir_record_ext(rec: &[u8]) -> Option<(u64, u64)> { /// `;1` version suffix. fn dir_record_name(rec: &[u8]) -> String { let name_len = *rec.get(32).unwrap_or(&0) as usize; - if name_len == 0 || rec.len() < 33 + name_len { return String::new(); } + if name_len == 0 || rec.len() < 33 + name_len { + return String::new(); + } let raw = &rec[33..33 + name_len]; let s = String::from_utf8_lossy(raw).to_string(); // Strip `;N` version suffix. - if let Some(i) = s.rfind(';') { s[..i].to_string() } else { s } + if let Some(i) = s.rfind(';') { + s[..i].to_string() + } else { + s + } } diff --git a/crates/http-api/src/log_stream.rs b/crates/http-api/src/log_stream.rs index 430bd58..cc1a4c2 100644 --- a/crates/http-api/src/log_stream.rs +++ b/crates/http-api/src/log_stream.rs @@ -36,9 +36,11 @@ pub async fn stream( let rx = state.log_bus.subscribe(); let live = BroadcastStream::new(rx).map(|res| match res { Ok(line) => Ok(Event::default().data(line_json(&line))), - Err(tokio_stream::wrappers::errors::BroadcastStreamRecvError::Lagged(n)) => Ok(Event::default() - .event("lagged") - .data(json!({ "skipped": n }).to_string())), + Err(tokio_stream::wrappers::errors::BroadcastStreamRecvError::Lagged(n)) => { + Ok(Event::default() + .event("lagged") + .data(json!({ "skipped": n }).to_string())) + } }); Sse::new(recent_stream.chain(live)) @@ -55,9 +57,11 @@ pub async fn recent(State(state): State) -> Json { /// keep streaming new lines as they arrive). pub async fn clear(State(state): State) -> Json { state.log_bus.clear(); - state - .log_bus - .push("info", "openpxe::terminal", "log buffer cleared by operator"); + state.log_bus.push( + "info", + "openpxe::terminal", + "log buffer cleared by operator", + ); Json(json!({ "ok": true })) } diff --git a/crates/http-api/src/terminal.rs b/crates/http-api/src/terminal.rs index 31e54f3..c0f48b8 100644 --- a/crates/http-api/src/terminal.rs +++ b/crates/http-api/src/terminal.rs @@ -31,12 +31,17 @@ pub async fn run_command( ) -> impl IntoResponse { let line = req.command.trim(); if line.is_empty() { - return (StatusCode::OK, Json(json!({ "output": HELP_TEXT, "ok": true }))); + return ( + StatusCode::OK, + Json(json!({ "output": HELP_TEXT, "ok": true })), + ); } // Echo the typed command into the live log so the Terminal tab shows // operator activity in-band with server-emitted log lines. - state.log_bus.push("info", "openpxe::terminal", format!("> {line}")); + state + .log_bus + .push("info", "openpxe::terminal", format!("> {line}")); let argv = shell_split(line); if argv.is_empty() { @@ -55,7 +60,11 @@ pub async fn run_command( // so reading the live tail tells the same story as scrolling the // terminal pane. let mirror = if output.len() > 1024 { - format!("{}\n... ({} bytes truncated)", &output[..1024], output.len() - 1024) + format!( + "{}\n... ({} bytes truncated)", + &output[..1024], + output.len() - 1024 + ) } else { output.clone() }; @@ -78,7 +87,7 @@ async fn dispatch(state: &AppState, argv: &[String]) -> Result { "status" => Ok(status_text(state)), "isos" | "images" => Ok(isos_text(state)), "clients" => Ok(clients_text(state)), - "queue" => gate_command(state, tail).await, + "queue" => queue_command(state, tail).await, "nfs" => nfs_command(state, tail).await, "smb" => smb_command(state, tail).await, "log" => log_command(state, tail), @@ -96,7 +105,7 @@ async fn dispatch(state: &AppState, argv: &[String]) -> Result { fn status_text(s: &AppState) -> String { let isos = s.iso_store.list(); let clients = s.clients.list(); - let gates = s.queue.list(); + let queue_entries = s.queue.list(); let smb = s.smb.as_ref().map(|m| m.snapshot()); let nfs = s.nfs.list(); let nfs_active = nfs.iter().filter(|m| m.mounted).count(); @@ -112,13 +121,23 @@ fn status_text(s: &AppState) -> String { nfs mounts: {n_total} configured ({n_active} active)\n", ver = env!("CARGO_PKG_VERSION"), base = s.public_base_url, - nic = if s.nic_name.is_empty() { "?" } else { s.nic_name.as_str() }, + nic = if s.nic_name.is_empty() { + "?" + } else { + s.nic_name.as_str() + }, up = uptime_string(s), n_isos = isos.len(), - n_local = isos.iter().filter(|i| matches!(i.source, openpxe_iso_store::IsoSource::Local)).count(), - n_nfs = isos.iter().filter(|i| !matches!(i.source, openpxe_iso_store::IsoSource::Local)).count(), + n_local = isos + .iter() + .filter(|i| matches!(i.source, openpxe_iso_store::IsoSource::Local)) + .count(), + n_nfs = isos + .iter() + .filter(|i| !matches!(i.source, openpxe_iso_store::IsoSource::Local)) + .count(), n_clients = clients.len(), - n_entries = gates.len(), + n_entries = queue_entries.len(), smb = smb.map_or_else(|| "(disabled)".into(), |s| format!("{s:?}")), n_total = nfs.len(), n_active = nfs_active, @@ -159,11 +178,7 @@ fn clients_text(s: &AppState) -> String { return "(no clients yet)".into(); } let mut out = String::new(); - let _ = writeln!( - out, - "{:<19} {:<16} {:<8} LAST SEEN", - "MAC", "IP", "EVENTS" - ); + let _ = writeln!(out, "{:<19} {:<16} {:<8} LAST SEEN", "MAC", "IP", "EVENTS"); for c in clients { let ip = c.last_ip.map_or_else(|| "-".into(), |i| i.to_string()); let _ = writeln!( @@ -180,35 +195,35 @@ fn clients_text(s: &AppState) -> String { out } -// ── gate ─────────────────────────────────────────────────────────────── +// ── queue ────────────────────────────────────────────────────────────── -// `async` for symmetry with the other dispatch helpers — gate operations +// `async` for symmetry with the other dispatch helpers — queue operations // are sync today but might grow to await on a database in a future phase. #[allow(clippy::unused_async)] -async fn gate_command(s: &AppState, args: &[String]) -> Result { +async fn queue_command(s: &AppState, args: &[String]) -> Result { match args.first().map(String::as_str) { None | Some("list") => { - let gs = s.queue.list(); - if gs.is_empty() { - return Ok("(no gates)".into()); + let entries = s.queue.list(); + if entries.is_empty() { + return Ok("(queue empty)".into()); } let mut out = String::new(); - for g in gs { + for entry in entries { let _ = writeln!( out, "#{:<3} {:<19} {:<16} target={}", - g.position, - g.mac, - g.id, - g.assigned_target.unwrap_or_else(|| "-".into()) + entry.position, + entry.mac, + entry.id, + entry.assigned_target.unwrap_or_else(|| "-".into()) ); } Ok(out) } Some("assign-all") => { - let target = args.get(1).ok_or_else(|| { - "usage: gate assign-all ".to_string() - })?; + let target = args + .get(1) + .ok_or_else(|| "usage: queue assign-all ".to_string())?; let found = s .iso_store .list() @@ -219,30 +234,32 @@ async fn gate_command(s: &AppState, args: &[String]) -> Result { } let ids: Vec<_> = s.queue.list().into_iter().map(|g| g.id).collect(); let n = s.queue.assign(&ids, target); - Ok(format!("assigned {n} gates -> {target}")) + Ok(format!("assigned {n} queue entries -> {target}")) } Some("assign") => { let entry_id = args .get(1) - .ok_or_else(|| "usage: gate assign ".to_string())?; + .ok_or_else(|| "usage: queue assign ".to_string())?; let target = args .get(2) - .ok_or_else(|| "usage: gate assign ".to_string())?; + .ok_or_else(|| "usage: queue assign ".to_string())?; let n = s.queue.assign(std::slice::from_ref(entry_id), target); if n == 0 { - return Err(format!("no such gate: {entry_id}")); + return Err(format!("no such queue entry: {entry_id}")); } - Ok(format!("assigned 1 gate -> {target}")) + Ok(format!("assigned 1 queue entry -> {target}")) } Some("release") => { - let entry_id = args.get(1).ok_or_else(|| "usage: gate release ".to_string())?; + let entry_id = args + .get(1) + .ok_or_else(|| "usage: queue release ".to_string())?; match s.queue.release(entry_id) { Some(_) => Ok(format!("released {entry_id}")), - None => Err(format!("no such gate: {entry_id}")), + None => Err(format!("no such queue entry: {entry_id}")), } } Some(other) => Err(format!( - "unknown gate subcommand: {other}\ntry: gate [list|assign-all|assign|release]" + "unknown queue subcommand: {other}\ntry: queue [list|assign-all|assign|release]" )), } } @@ -294,7 +311,9 @@ async fn nfs_command(s: &AppState, args: &[String]) -> Result { let version = match args.get(2).map(String::as_str) { Some("v3") => openpxe_iso_store::NfsVersion::V3, Some("v41") | None => openpxe_iso_store::NfsVersion::V41, - Some(other) => return Err(format!("unknown nfs version: {other} (expect v3 or v41)")), + Some(other) => { + return Err(format!("unknown nfs version: {other} (expect v3 or v41)")) + } }; let read_only = !matches!(args.get(3).map(String::as_str), Some("rw")); let req = openpxe_iso_store::NfsAddRequest { @@ -309,14 +328,18 @@ async fn nfs_command(s: &AppState, args: &[String]) -> Result { } } Some("unmount") => { - let id = args.get(1).ok_or_else(|| "usage: nfs unmount ".to_string())?; + let id = args + .get(1) + .ok_or_else(|| "usage: nfs unmount ".to_string())?; match s.nfs.remove(id).await { Ok(()) => Ok(format!("unmounted {id}")), Err(e) => Err(format!("unmount failed: {e}")), } } Some("scan") => { - let id = args.get(1).ok_or_else(|| "usage: nfs scan ".to_string())?; + let id = args + .get(1) + .ok_or_else(|| "usage: nfs scan ".to_string())?; match s.nfs.rescan(id).await { Ok(n) => Ok(format!("re-scanned {id}: {n} isos")), Err(e) => Err(format!("scan failed: {e}")), @@ -368,10 +391,7 @@ fn log_command(s: &AppState, args: &[String]) -> Result { Ok("log buffer cleared".into()) } Some("tail") => { - let n: usize = args - .get(1) - .and_then(|v| v.parse().ok()) - .unwrap_or(20); + let n: usize = args.get(1).and_then(|v| v.parse().ok()).unwrap_or(20); let lines = s.log_bus.recent(); let start = lines.len().saturating_sub(n); let mut out = String::new(); @@ -462,10 +482,10 @@ OpenPXE terminal — available commands: isos list registered ISOs clients list PXE clients seen this session - gate list list gated-deployment queue - gate assign assign one gate to a boot entry - gate assign-all assign every waiting gate - gate release release one gate + queue list list queued clients + queue assign assign one queued client to a boot entry + queue assign-all assign every waiting client + queue release release one queued client nfs list list NFS mounts nfs mount : [v3|v41] [ro|rw] add and mount an NFS share @@ -521,4 +541,10 @@ mod tests { assert_eq!(truncate("hi", 10), "hi"); assert_eq!(truncate("longerthanfive", 5), "long…"); } + + #[test] + fn help_uses_queue_language() { + assert!(HELP_TEXT.contains("queue list")); + assert!(HELP_TEXT.contains("queued clients")); + } } diff --git a/crates/http-api/tests/full_flow.rs b/crates/http-api/tests/full_flow.rs index 354f1d5..39d75f9 100644 --- a/crates/http-api/tests/full_flow.rs +++ b/crates/http-api/tests/full_flow.rs @@ -38,13 +38,12 @@ fn fake_alpine_iso() -> Vec { } fn multipart_iso_body(filename: &str, bytes: &[u8]) -> (String, Vec) { - let boundary = "----PxeForgeTestBoundary1234"; + let boundary = "----OpenPxeTestBoundary1234"; let mut body = Vec::new(); body.extend_from_slice(format!("--{boundary}\r\n").as_bytes()); body.extend_from_slice( - format!( - "Content-Disposition: form-data; name=\"file\"; filename=\"{filename}\"\r\n" - ).as_bytes(), + format!("Content-Disposition: form-data; name=\"file\"; filename=\"{filename}\"\r\n") + .as_bytes(), ); body.extend_from_slice(b"Content-Type: application/octet-stream\r\n\r\n"); body.extend_from_slice(bytes); @@ -60,7 +59,10 @@ async fn get(router: &axum::Router, path: &str) -> (StatusCode, Vec) { .await .unwrap(); let status = res.status(); - let body = axum::body::to_bytes(res.into_body(), usize::MAX).await.unwrap().to_vec(); + let body = axum::body::to_bytes(res.into_body(), usize::MAX) + .await + .unwrap() + .to_vec(); (status, body) } @@ -78,7 +80,10 @@ async fn post_json(router: &axum::Router, path: &str, body: &str) -> (StatusCode .await .unwrap(); let status = res.status(); - let body = axum::body::to_bytes(res.into_body(), usize::MAX).await.unwrap().to_vec(); + let body = axum::body::to_bytes(res.into_body(), usize::MAX) + .await + .unwrap() + .to_vec(); (status, body) } @@ -87,7 +92,7 @@ async fn build_state() -> (AppState, tempfile::TempDir) { let iso_store = IsoStore::new(dir.path().join("isos")); iso_store.ensure_dirs().await.unwrap(); let clients = ClientRegistry::new(); - let gates = DeploymentQueue::new(); + let queue = DeploymentQueue::new(); let settings = SettingsStore::load_or_default(dir.path()); let nfs = NfsManager::new(dir.path(), iso_store.clone()); iso_store.set_nfs_root(nfs.mount_root()); @@ -97,7 +102,7 @@ async fn build_state() -> (AppState, tempfile::TempDir) { let state = AppState { iso_store, clients, - queue: gates, + queue, settings, hosts, metrics, @@ -153,13 +158,21 @@ async fn upload_introspects_and_generates_boot_entry() { // Confirm the ISO shows up in the menu. let (_, menu) = get(&app, "/boot.ipxe").await; let menu = String::from_utf8(menu).unwrap(); - assert!(menu.contains("Linux Installers"), "menu missing Linux submenu:\n{menu}"); + assert!( + menu.contains("Linux Installers"), + "menu missing Linux submenu:\n{menu}" + ); let (_, linux) = get(&app, "/boot/_linux_menu.ipxe").await; let linux = String::from_utf8(linux).unwrap(); - assert!(linux.contains("fake-alpine-linux"), "linux submenu missing entry:\n{linux}"); - assert!(linux.contains("[ 0 MB]") || linux.contains("[ 0 MB]"), - "size label missing in {linux}"); + assert!( + linux.contains("fake-alpine-linux"), + "linux submenu missing entry:\n{linux}" + ); + assert!( + linux.contains("[ 0 MB]") || linux.contains("[ 0 MB]"), + "size label missing in {linux}" + ); // Per-entry boot script should include kernel + initrd URLs + boot. let (_, entry) = get(&app, "/boot/fake-alpine-linux.ipxe").await; @@ -178,9 +191,14 @@ async fn iso_range_request_slices_correctly() { app.clone() .oneshot( Request::builder() - .method("POST").uri("/api/isos") + .method("POST") + .uri("/api/isos") .header("content-type", ct) - .body(Body::from(body)).unwrap()).await.unwrap(); + .body(Body::from(body)) + .unwrap(), + ) + .await + .unwrap(); // Range bytes=0x8000-0x8005 should return the PVD signature byte. let res = app @@ -189,10 +207,15 @@ async fn iso_range_request_slices_correctly() { Request::builder() .uri("/iso/fake-alpine.iso") .header(header::RANGE, "bytes=32768-32773") - .body(Body::empty()).unwrap()) - .await.unwrap(); + .body(Body::empty()) + .unwrap(), + ) + .await + .unwrap(); assert_eq!(res.status(), StatusCode::PARTIAL_CONTENT); - let slice = axum::body::to_bytes(res.into_body(), usize::MAX).await.unwrap(); + let slice = axum::body::to_bytes(res.into_body(), usize::MAX) + .await + .unwrap(); assert_eq!(slice[0], 0x01); // PVD type assert_eq!(&slice[1..6], b"CD001"); } @@ -205,27 +228,40 @@ async fn queued_deployment_full_flow() { // Upload an ISO so the target exists. let (ct, body) = multipart_iso_body("fake-alpine.iso", &fake_alpine_iso()); app.clone() - .oneshot(Request::builder().method("POST").uri("/api/isos") - .header("content-type", ct).body(Body::from(body)).unwrap()) - .await.unwrap(); + .oneshot( + Request::builder() + .method("POST") + .uri("/api/isos") + .header("content-type", ct) + .body(Body::from(body)) + .unwrap(), + ) + .await + .unwrap(); // Two clients join. let (_, join1) = get(&app, "/api/queue/join?mac=aa:bb:cc:00:00:01").await; let (_, join2) = get(&app, "/api/queue/join?mac=aa:bb:cc:00:00:02").await; let s1 = String::from_utf8(join1).unwrap(); let s2 = String::from_utf8(join2).unwrap(); - assert!(s1.contains("Gate Position 1")); - assert!(s2.contains("Gate Position 2")); + assert!(s1.contains("Queue Position 1")); + assert!(s2.contains("Queue Position 2")); - let gate1_id = s1.lines().find_map(|l| l.strip_prefix("chain http://127.0.0.1/api/queue/poll/")) - .unwrap().to_string(); - let gate2_id = s2.lines().find_map(|l| l.strip_prefix("chain http://127.0.0.1/api/queue/poll/")) - .unwrap().to_string(); + let queue1_id = s1 + .lines() + .find_map(|l| l.strip_prefix("chain http://127.0.0.1/api/queue/poll/")) + .unwrap() + .to_string(); + let queue2_id = s2 + .lines() + .find_map(|l| l.strip_prefix("chain http://127.0.0.1/api/queue/poll/")) + .unwrap() + .to_string(); // Kick off a long-poll for client 1 in the background. Then assign. let app2 = app.clone(); let poll_future = tokio::spawn(async move { - let uri = format!("/api/queue/poll/{gate1_id}"); + let uri = format!("/api/queue/poll/{queue1_id}"); get(&app2, &uri).await }); @@ -233,13 +269,16 @@ async fn queued_deployment_full_flow() { tokio::time::sleep(std::time::Duration::from_millis(50)).await; // Operator assigns. - let body = format!(r#"{{"target":"fake-alpine-linux","entry_ids":["{gate2_id}"]}}"#); + let body = format!(r#"{{"target":"fake-alpine-linux","entry_ids":["{queue2_id}"]}}"#); let (s, b) = post_json(&app, "/api/queue/assign", &body).await; assert_eq!(s, StatusCode::OK); let assign_json = String::from_utf8(b).unwrap(); - assert!(assign_json.contains(r#""assigned":1"#), "assign response: {assign_json}"); + assert!( + assign_json.contains(r#""assigned":1"#), + "assign response: {assign_json}" + ); - // Now assign to gate 1 too so the background poll wakes. + // Now assign to queue entry 1 too so the background poll wakes. let body = r#"{"target":"fake-alpine-linux","entry_ids":[]}"#; post_json(&app, "/api/queue/assign", body).await; @@ -251,14 +290,23 @@ async fn queued_deployment_full_flow() { "poll response should chain the boot script:\n{poll_s}" ); // Retry-on-error fallback must be present. - assert!(poll_s.contains("|| chain http://127.0.0.1/api/queue/poll/"), - "retry fallback missing"); + assert!( + poll_s.contains("|| chain http://127.0.0.1/api/queue/poll/"), + "retry fallback missing" + ); // Bad target must be rejected. - let (_, bad) = post_json(&app, "/api/queue/assign", - r#"{"target":"does-not-exist","entry_ids":[]}"#).await; + let (_, bad) = post_json( + &app, + "/api/queue/assign", + r#"{"target":"does-not-exist","entry_ids":[]}"#, + ) + .await; let bad_s = String::from_utf8(bad).unwrap(); - assert!(bad_s.contains(r#""ok":false"#), "expected rejection: {bad_s}"); + assert!( + bad_s.contains(r#""ok":false"#), + "expected rejection: {bad_s}" + ); } #[tokio::test] @@ -272,8 +320,9 @@ async fn settings_put_persists_across_reads() { "smb_host_override": "", "extra_kernel_args": "console=ttyS0", "default_local_hdd": true, - "gate_wait_max_secs": 0 - }).to_string(); + "queue_wait_max_secs": 0 + }) + .to_string(); let res = app .clone() @@ -283,7 +332,8 @@ async fn settings_put_persists_across_reads() { .uri("/api/settings") .header("content-type", "application/json") .body(Body::from(body)) - .unwrap()) + .unwrap(), + ) .await .unwrap(); assert_eq!(res.status(), StatusCode::NO_CONTENT); @@ -297,10 +347,14 @@ async fn settings_put_persists_across_reads() { // And the menu should now use the new timeout. let (_, menu) = get(&app, "/boot.ipxe").await; let menu = String::from_utf8(menu).unwrap(); - assert!(menu.contains("--timeout 42000"), - "menu should reflect 42s timeout:\n{menu}"); - assert!(menu.contains("--default local"), - "menu should default to local:\n{menu}"); + assert!( + menu.contains("--timeout 42000"), + "menu should reflect 42s timeout:\n{menu}" + ); + assert!( + menu.contains("--default local"), + "menu should default to local:\n{menu}" + ); } #[tokio::test] @@ -309,14 +363,22 @@ async fn reboot_and_firmware_exit_in_tools_menu() { let app = build_router(state); let (_, tools) = get(&app, "/boot/_tools_menu.ipxe").await; let tools = String::from_utf8(tools).unwrap(); - assert!(tools.contains("Reboot Computer"), - "tools menu missing Reboot item:\n{tools}"); - assert!(tools.contains("Exit and continue BIOS boot"), - "tools menu missing firmware-exit item:\n{tools}"); - assert!(tools.contains("&& reboot"), - "reboot command not wired:\n{tools}"); - assert!(tools.contains("&& exit 0"), - "firmware exit command not wired:\n{tools}"); + assert!( + tools.contains("Reboot Computer"), + "tools menu missing Reboot item:\n{tools}" + ); + assert!( + tools.contains("Exit and continue BIOS boot"), + "tools menu missing firmware-exit item:\n{tools}" + ); + assert!( + tools.contains("&& reboot"), + "reboot command not wired:\n{tools}" + ); + assert!( + tools.contains("&& exit 0"), + "firmware exit command not wired:\n{tools}" + ); } #[tokio::test] @@ -334,11 +396,15 @@ async fn ui_assets_served_offline() { let res = app .clone() .oneshot(Request::builder().uri(path).body(Body::empty()).unwrap()) - .await.unwrap(); + .await + .unwrap(); assert_eq!(res.status(), StatusCode::OK, "{path} not 200"); - let got = res.headers() - .get(header::CONTENT_TYPE).unwrap() - .to_str().unwrap(); + let got = res + .headers() + .get(header::CONTENT_TYPE) + .unwrap() + .to_str() + .unwrap(); assert!(got.starts_with(ct), "{path} ct={got}, expected {ct}"); } } @@ -348,20 +414,38 @@ async fn no_external_urls_in_generated_ipxe() { // Sanity check that nothing we serve points off-server. let (state, _dir) = build_state().await; let app = build_router(state); - for path in ["/boot.ipxe", "/boot/_tools_menu.ipxe", "/boot/_linux_menu.ipxe", - "/boot/_shell.ipxe", "/boot/_nic.ipxe", "/boot/_local.ipxe"] { + for path in [ + "/boot.ipxe", + "/boot/_tools_menu.ipxe", + "/boot/_linux_menu.ipxe", + "/boot/_shell.ipxe", + "/boot/_nic.ipxe", + "/boot/_local.ipxe", + ] { let (_, body) = get(&app, path).await; let s = String::from_utf8(body).unwrap(); // The only URLs we should emit are relative to our own public_base_url. - for url in ["github.com", "googleapis", "cdn.", "cdnjs", "unpkg", "jsdelivr"] { - assert!(!s.contains(url), "{path} references external host {url}:\n{s}"); + for url in [ + "github.com", + "googleapis", + "cdn.", + "cdnjs", + "unpkg", + "jsdelivr", + ] { + assert!( + !s.contains(url), + "{path} references external host {url}:\n{s}" + ); } // Confirm URLs are all ours. for line in s.lines() { if let Some(idx) = line.find("http://") { let rest = &line[idx..]; - assert!(rest.starts_with("http://127.0.0.1"), - "{path} references non-public-base URL: {line}"); + assert!( + rest.starts_with("http://127.0.0.1"), + "{path} references non-public-base URL: {line}" + ); } } } @@ -384,7 +468,10 @@ async fn nfs_add_with_bad_export_is_rejected() { .await; assert_eq!(s, StatusCode::BAD_REQUEST); let msg = String::from_utf8_lossy(&b); - assert!(msg.contains("export"), "expected validation hint, got: {msg}"); + assert!( + msg.contains("export"), + "expected validation hint, got: {msg}" + ); } #[tokio::test] @@ -413,7 +500,10 @@ async fn terminal_help_and_status_round_trip() { assert_eq!(s, StatusCode::OK); let v: serde_json::Value = serde_json::from_slice(&b).unwrap(); let out = v["output"].as_str().unwrap(); - assert!(out.starts_with("OpenPXE"), "unexpected status output: {out}"); + assert!( + out.starts_with("OpenPXE"), + "unexpected status output: {out}" + ); assert!(out.contains("isos:"), "status missing iso line: {out}"); // Unknown command -> ok=false plus help hint. @@ -434,7 +524,10 @@ async fn log_recent_returns_buffered_lines() { assert_eq!(s, StatusCode::OK); let v: serde_json::Value = serde_json::from_slice(&b).unwrap(); let lines = v["lines"].as_array().expect("lines array"); - assert!(!lines.is_empty(), "log buffer should have at least one line"); + assert!( + !lines.is_empty(), + "log buffer should have at least one line" + ); // Every entry should have the canonical timestamp/level/target/message. for l in lines { for k in ["timestamp", "level", "target", "message"] { @@ -479,7 +572,7 @@ async fn windows_iso_renders_clean_wimboot_script_with_no_trust_store_writes() { .body(Body::from( r#"{"boot_menu_timeout_secs":600,"timeout_action":"queued_deployment", "windows_enabled":false,"smb_host_override":"","extra_kernel_args":"", - "default_local_hdd":true,"gate_wait_max_secs":0,"dns_server":""}"# + "default_local_hdd":true,"queue_wait_max_secs":0,"dns_server":""}"# .to_string(), )) .unwrap(), @@ -504,7 +597,9 @@ async fn windows_iso_renders_clean_wimboot_script_with_no_trust_store_writes() { .await .unwrap(); assert_eq!(res.status(), StatusCode::CREATED); - let body = axum::body::to_bytes(res.into_body(), usize::MAX).await.unwrap(); + let body = axum::body::to_bytes(res.into_body(), usize::MAX) + .await + .unwrap(); let meta: serde_json::Value = serde_json::from_slice(&body).unwrap(); assert_eq!(meta["introspection"]["family"], "windows_pe"); assert!( @@ -534,7 +629,10 @@ async fn windows_iso_renders_clean_wimboot_script_with_no_trust_store_writes() { assert_eq!(s, StatusCode::OK); let script = String::from_utf8(body).unwrap(); assert!(script.contains("kernel "), "missing kernel line:\n{script}"); - assert!(script.contains("ipxe/wimboot"), "missing wimboot loader:\n{script}"); + assert!( + script.contains("ipxe/wimboot"), + "missing wimboot loader:\n{script}" + ); for tag in ["bootmgr", "bootmgr.efi", "bcd", "boot.sdi", "boot.wim"] { assert!( script.contains(&format!("initrd --name {tag}")), @@ -544,8 +642,12 @@ async fn windows_iso_renders_clean_wimboot_script_with_no_trust_store_writes() { // Hard guarantees we never want to see in any client-facing script. let lower = script.to_lowercase(); for forbidden in [ - "bcdedit", "testsigning", "certutil", "test-signed", - "httpdisk", "/set testsigning", + "bcdedit", + "testsigning", + "certutil", + "test-signed", + "httpdisk", + "/set testsigning", ] { assert!( !lower.contains(forbidden), @@ -610,16 +712,25 @@ async fn metrics_endpoint_emits_prometheus_format() { let res = app .clone() - .oneshot(Request::builder().uri("/metrics").body(Body::empty()).unwrap()) + .oneshot( + Request::builder() + .uri("/metrics") + .body(Body::empty()) + .unwrap(), + ) .await .unwrap(); assert_eq!(res.status(), StatusCode::OK); - let ct = res.headers().get(header::CONTENT_TYPE).unwrap().to_str().unwrap(); - assert!( - ct.starts_with("text/plain"), - "wrong content-type: {ct}" - ); - let body = axum::body::to_bytes(res.into_body(), usize::MAX).await.unwrap(); + let ct = res + .headers() + .get(header::CONTENT_TYPE) + .unwrap() + .to_str() + .unwrap(); + assert!(ct.starts_with("text/plain"), "wrong content-type: {ct}"); + let body = axum::body::to_bytes(res.into_body(), usize::MAX) + .await + .unwrap(); let body = String::from_utf8(body.to_vec()).unwrap(); // Spot-check the must-have metric families. for name in [ @@ -633,10 +744,7 @@ async fn metrics_endpoint_emits_prometheus_format() { assert!(body.contains(name), "missing metric {name} in:\n{body}"); } // Each name appears exactly once as a `# TYPE` declaration. - for name in [ - "openpxe_dhcp_replies_total", - "openpxe_iso_count", - ] { + for name in ["openpxe_dhcp_replies_total", "openpxe_iso_count"] { let count = body.matches(&format!("# TYPE {name}")).count(); assert_eq!(count, 1, "{name} TYPE line appears {count} times"); } @@ -673,10 +781,10 @@ async fn network_endpoint_exposes_dns_round_trip() { assert_eq!(v["dns_server"], "10.0.0.1"); } -// ─── v0.3.1: per-ISO password gate ──────────────────────────────────────── +// ─── Per-ISO password prompt ────────────────────────────────────────────── #[tokio::test] -async fn iso_password_gate_blocks_until_correct_token() { +async fn iso_password_prompt_blocks_until_correct_token() { let (state, _dir) = build_state().await; let app = build_router(state); @@ -689,9 +797,14 @@ async fn iso_password_gate_blocks_until_correct_token() { .clone() .oneshot( Request::builder() - .method("POST").uri("/api/isos") + .method("POST") + .uri("/api/isos") .header("content-type", ct) - .body(Body::from(body)).unwrap()).await.unwrap(); + .body(Body::from(body)) + .unwrap(), + ) + .await + .unwrap(); assert_eq!(res.status(), StatusCode::CREATED); // 1. With NO password set, /boot/.ipxe returns the boot script @@ -702,62 +815,93 @@ async fn iso_password_gate_blocks_until_correct_token() { let (_, lm) = get(&app, "/boot/_linux_menu.ipxe").await; let lm = String::from_utf8(lm).unwrap(); assert!(lm.contains("fake-alpine-linux")); - assert!(!lm.contains("fake-alpine-linux *["), - "expected no lock marker in menu before password set:\n{lm}"); + assert!( + !lm.contains("fake-alpine-linux *["), + "expected no lock marker in menu before password set:\n{lm}" + ); // 2. Set a password. - let res = app.clone().oneshot( - Request::builder() - .method("PUT") - .uri("/api/isos/fake-alpine/password") - .header("content-type", "application/json") - .body(Body::from(r#"{"password":"hunter2"}"#)).unwrap() - ).await.unwrap(); + let res = app + .clone() + .oneshot( + Request::builder() + .method("PUT") + .uri("/api/isos/fake-alpine/password") + .header("content-type", "application/json") + .body(Body::from(r#"{"password":"hunter2"}"#)) + .unwrap(), + ) + .await + .unwrap(); assert_eq!(res.status(), StatusCode::NO_CONTENT); // The menu now shows the lock marker (`*` prefix on the size box). let (_, lm) = get(&app, "/boot/_linux_menu.ipxe").await; let lm = String::from_utf8(lm).unwrap(); - assert!(lm.contains("fake-alpine-linux *["), - "expected lock marker in menu after password set:\n{lm}"); + assert!( + lm.contains("fake-alpine-linux *["), + "expected lock marker in menu after password set:\n{lm}" + ); // 3. Without a token, /boot/.ipxe now returns the password // PROMPT script (read --secret), not the boot script. let (_, body) = get(&app, "/boot/fake-alpine-linux.ipxe").await; let s = String::from_utf8(body).unwrap(); - assert!(s.contains("read --secret password"), - "expected prompt script with no token, got:\n{s}"); + assert!( + s.contains("read --secret password"), + "expected prompt script with no token, got:\n{s}" + ); assert!(!s.contains("kernel "), "should not include kernel line yet"); // 4. Wrong token -> "Wrong password." script that chains back to the entry. let (_, body) = get(&app, "/boot/fake-alpine-linux.ipxe?token=wrongpw").await; let s = String::from_utf8(body).unwrap(); - assert!(s.contains("Wrong password."), "expected auth-fail script, got:\n{s}"); + assert!( + s.contains("Wrong password."), + "expected auth-fail script, got:\n{s}" + ); assert!(s.contains("/boot/fake-alpine-linux.ipxe")); assert!(!s.contains("kernel ")); // Critical: the WRONG token must NEVER be echoed back in the script. - assert!(!s.contains("wrongpw"), "wrong token must not appear in response"); + assert!( + !s.contains("wrongpw"), + "wrong token must not appear in response" + ); // 5. Correct token -> real boot script. let (_, body) = get(&app, "/boot/fake-alpine-linux.ipxe?token=hunter2").await; let s = String::from_utf8(body).unwrap(); - assert!(s.contains("kernel "), "expected boot script with correct token, got:\n{s}"); + assert!( + s.contains("kernel "), + "expected boot script with correct token, got:\n{s}" + ); // Don't echo the password into the boot script either. - assert!(!s.contains("hunter2"), "correct password must not leak into boot script"); + assert!( + !s.contains("hunter2"), + "correct password must not leak into boot script" + ); // 6. Clear the password (DELETE). - let res = app.clone().oneshot( - Request::builder() - .method("DELETE") - .uri("/api/isos/fake-alpine/password") - .body(Body::empty()).unwrap() - ).await.unwrap(); + let res = app + .clone() + .oneshot( + Request::builder() + .method("DELETE") + .uri("/api/isos/fake-alpine/password") + .body(Body::empty()) + .unwrap(), + ) + .await + .unwrap(); assert_eq!(res.status(), StatusCode::NO_CONTENT); // Boot is open again, no lock indicator. let (_, body) = get(&app, "/boot/fake-alpine-linux.ipxe").await; let s = String::from_utf8(body).unwrap(); - assert!(s.contains("kernel "), "expected boot script after clear, got:\n{s}"); + assert!( + s.contains("kernel "), + "expected boot script after clear, got:\n{s}" + ); let (_, lm) = get(&app, "/boot/_linux_menu.ipxe").await; let lm = String::from_utf8(lm).unwrap(); assert!(!lm.contains("fake-alpine-linux *[")); @@ -771,34 +915,63 @@ async fn iso_password_set_then_clear_via_null_body() { // Upload + set + clear via `{"password": null}` (alternative to DELETE). let iso = fake_alpine_iso(); let (ct, body) = multipart_iso_body("fake-alpine.iso", &iso); - let res = app.clone().oneshot( - Request::builder().method("POST").uri("/api/isos") - .header("content-type", ct) - .body(Body::from(body)).unwrap()).await.unwrap(); + let res = app + .clone() + .oneshot( + Request::builder() + .method("POST") + .uri("/api/isos") + .header("content-type", ct) + .body(Body::from(body)) + .unwrap(), + ) + .await + .unwrap(); assert_eq!(res.status(), StatusCode::CREATED); - for body in [r#"{"password":"x"}"#, r#"{"password":null}"#, r#"{"password":""}"#] { - let res = app.clone().oneshot( - Request::builder().method("PUT") - .uri("/api/isos/fake-alpine/password") - .header("content-type", "application/json") - .body(Body::from(body.to_string())).unwrap()).await.unwrap(); + for body in [ + r#"{"password":"x"}"#, + r#"{"password":null}"#, + r#"{"password":""}"#, + ] { + let res = app + .clone() + .oneshot( + Request::builder() + .method("PUT") + .uri("/api/isos/fake-alpine/password") + .header("content-type", "application/json") + .body(Body::from(body.to_string())) + .unwrap(), + ) + .await + .unwrap(); assert_eq!(res.status(), StatusCode::NO_CONTENT, "body={body}"); } // After the empty string, the entry should be unprotected. let (_, b) = get(&app, "/boot/fake-alpine-linux.ipxe").await; let s = String::from_utf8(b).unwrap(); - assert!(s.contains("kernel "), "should be unprotected after empty pw, got:\n{s}"); + assert!( + s.contains("kernel "), + "should be unprotected after empty pw, got:\n{s}" + ); } #[tokio::test] async fn set_password_for_unknown_iso_returns_404() { let (state, _dir) = build_state().await; let app = build_router(state); - let res = app.clone().oneshot( - Request::builder().method("PUT") - .uri("/api/isos/does-not-exist/password") - .header("content-type", "application/json") - .body(Body::from(r#"{"password":"x"}"#)).unwrap()).await.unwrap(); + let res = app + .clone() + .oneshot( + Request::builder() + .method("PUT") + .uri("/api/isos/does-not-exist/password") + .header("content-type", "application/json") + .body(Body::from(r#"{"password":"x"}"#)) + .unwrap(), + ) + .await + .unwrap(); assert_eq!(res.status(), StatusCode::NOT_FOUND); } diff --git a/crates/ipxe-assets/src/lib.rs b/crates/ipxe-assets/src/lib.rs index 37f7ee0..1227c6d 100644 --- a/crates/ipxe-assets/src/lib.rs +++ b/crates/ipxe-assets/src/lib.rs @@ -51,7 +51,9 @@ pub fn asset_slice(name: &str) -> Option> { /// Enumerate embedded asset filenames. Useful for startup logging so the /// operator can immediately tell which architectures will work. pub fn list_assets() -> Vec { - IpxeAssets::iter().map(std::borrow::Cow::into_owned).collect() + IpxeAssets::iter() + .map(std::borrow::Cow::into_owned) + .collect() } /// Log at startup which iPXE binaries are present and which are missing. diff --git a/crates/iso-store/src/introspect.rs b/crates/iso-store/src/introspect.rs index 8308521..c8ee826 100644 --- a/crates/iso-store/src/introspect.rs +++ b/crates/iso-store/src/introspect.rs @@ -78,7 +78,9 @@ pub fn introspect(path: &Path) -> IntrospectionReport { let mut haystack = Vec::with_capacity(scan_bytes.min(32 * 1024 * 1024)); while read_total < scan_bytes { let n = f.read(&mut buf).unwrap_or(0); - if n == 0 { break; } + if n == 0 { + break; + } haystack.extend_from_slice(&buf[..n]); read_total += n; } @@ -107,8 +109,11 @@ fn family_from_label(label: &str) -> DistroFamily { let l = label.to_ascii_lowercase(); if l.contains("ubuntu") || l.contains("debian") || l.contains("mint") { DistroFamily::DebianUbuntu - } else if l.contains("rhel") || l.contains("centos") || l.contains("fedora") - || l.contains("rocky") || l.contains("alma") + } else if l.contains("rhel") + || l.contains("centos") + || l.contains("fedora") + || l.contains("rocky") + || l.contains("alma") { DistroFamily::RhelFedora } else if l.contains("suse") || l.contains("opensuse") { @@ -127,17 +132,30 @@ fn family_from_label(label: &str) -> DistroFamily { fn guess_kernel_initrd(family: DistroFamily) -> (Option<&'static str>, Vec<&'static str>) { match family { DistroFamily::DebianUbuntu => (Some("/casper/vmlinuz"), vec!["/casper/initrd"]), - DistroFamily::RhelFedora => (Some("/images/pxeboot/vmlinuz"), vec!["/images/pxeboot/initrd.img"]), - DistroFamily::OpenSuse => (Some("/boot/x86_64/loader/linux"), vec!["/boot/x86_64/loader/initrd"]), - DistroFamily::Arch => (Some("/arch/boot/x86_64/vmlinuz-linux"), vec!["/arch/boot/x86_64/initramfs-linux.img"]), + DistroFamily::RhelFedora => ( + Some("/images/pxeboot/vmlinuz"), + vec!["/images/pxeboot/initrd.img"], + ), + DistroFamily::OpenSuse => ( + Some("/boot/x86_64/loader/linux"), + vec!["/boot/x86_64/loader/initrd"], + ), + DistroFamily::Arch => ( + Some("/arch/boot/x86_64/vmlinuz-linux"), + vec!["/arch/boot/x86_64/initramfs-linux.img"], + ), DistroFamily::Alpine => (Some("/boot/vmlinuz-lts"), vec!["/boot/initramfs-lts"]), DistroFamily::WindowsPe | DistroFamily::Unknown => (None, Vec::new()), } } fn contains_ascii(haystack: &[u8], needle: &[u8]) -> bool { - if needle.is_empty() || haystack.len() < needle.len() { return false; } - haystack.windows(needle.len()).any(|w| w.eq_ignore_ascii_case(needle)) + if needle.is_empty() || haystack.len() < needle.len() { + return false; + } + haystack + .windows(needle.len()) + .any(|w| w.eq_ignore_ascii_case(needle)) } #[cfg(test)] @@ -146,9 +164,18 @@ mod tests { #[test] fn label_matching() { - assert_eq!(family_from_label("Ubuntu 24.04"), DistroFamily::DebianUbuntu); - assert_eq!(family_from_label("Rocky-9-x86_64-dvd"), DistroFamily::RhelFedora); - assert_eq!(family_from_label("openSUSE-Leap-15.6"), DistroFamily::OpenSuse); + assert_eq!( + family_from_label("Ubuntu 24.04"), + DistroFamily::DebianUbuntu + ); + assert_eq!( + family_from_label("Rocky-9-x86_64-dvd"), + DistroFamily::RhelFedora + ); + assert_eq!( + family_from_label("openSUSE-Leap-15.6"), + DistroFamily::OpenSuse + ); assert_eq!(family_from_label("ARCH_202604"), DistroFamily::Arch); assert_eq!(family_from_label("weird-custom"), DistroFamily::Unknown); } diff --git a/crates/iso-store/src/lib.rs b/crates/iso-store/src/lib.rs index 19e6bc1..ad8965d 100644 --- a/crates/iso-store/src/lib.rs +++ b/crates/iso-store/src/lib.rs @@ -27,5 +27,7 @@ pub use entry::{BootEntry, BootKind, KernelArgs}; pub use introspect::{DistroFamily, IntrospectionReport}; pub use nfs::{NfsAddRequest, NfsManager, NfsMount, NfsVersion}; pub use smb::{extract_windows_iso, SmbManager, SmbState}; -pub use store::{generate_boot_entries_for, slugify_str, IsoMeta, IsoSource, IsoStore, UploadHandle}; +pub use store::{ + generate_boot_entries_for, slugify_str, IsoMeta, IsoSource, IsoStore, UploadHandle, +}; pub use windows::{WimPatcher, WinPatchState}; diff --git a/crates/iso-store/src/nfs.rs b/crates/iso-store/src/nfs.rs index bb37068..c1ac159 100644 --- a/crates/iso-store/src/nfs.rs +++ b/crates/iso-store/src/nfs.rs @@ -36,8 +36,8 @@ use crate::introspect::{introspect, IntrospectionReport}; use crate::store::{generate_boot_entries_for, slugify_str, IsoSource, IsoStore}; -use parking_lot::Mutex; use openpxe_core::{Error, Result}; +use parking_lot::Mutex; use serde::{Deserialize, Serialize}; use std::collections::HashMap; use std::path::{Path, PathBuf}; @@ -403,14 +403,8 @@ impl NfsManager { mount_id: m.id.clone(), relative_path: filename.clone(), }; - self.iso_store.register_external( - id, - filename, - size, - report, - boot_entries, - source, - ); + self.iso_store + .register_external(id, filename, size, report, boot_entries, source); count += 1; } Ok(count) diff --git a/crates/iso-store/src/smb.rs b/crates/iso-store/src/smb.rs index add3271..b88eb35 100644 --- a/crates/iso-store/src/smb.rs +++ b/crates/iso-store/src/smb.rs @@ -25,11 +25,11 @@ //! Samba), we return `SmbState::SmbdMissing` and the UI surfaces the //! gap. No panics, no retries, no silent failure. +use parking_lot::Mutex; use serde::{Deserialize, Serialize}; use std::path::{Path, PathBuf}; use std::process::{Child, Command, Stdio}; use std::sync::Arc; -use parking_lot::Mutex; #[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] #[serde(rename_all = "snake_case", tag = "state")] @@ -72,7 +72,9 @@ impl SmbManager { /// ISO under `smb_dir//` becomes a share named ``. Returns /// the sorted list. pub fn discover_shares(&self) -> Vec { - let Ok(rd) = std::fs::read_dir(&self.smb_dir) else { return vec![]; }; + let Ok(rd) = std::fs::read_dir(&self.smb_dir) else { + return vec![]; + }; let mut out: Vec = rd .flatten() .filter(|e| e.path().is_dir()) @@ -130,7 +132,9 @@ impl SmbManager { let shares = match self.write_conf() { Ok(v) => v, Err(e) => { - let s = SmbState::Failed { reason: format!("write smb.conf: {e}") }; + let s = SmbState::Failed { + reason: format!("write smb.conf: {e}"), + }; *self.state.lock() = s.clone(); return s; } @@ -139,7 +143,8 @@ impl SmbManager { .args([ "--foreground", "--no-process-group", - "--configfile", self.conf_path.to_str().unwrap_or(""), + "--configfile", + self.conf_path.to_str().unwrap_or(""), "--log-stdout", ]) .stdin(Stdio::null()) @@ -156,7 +161,9 @@ impl SmbManager { s } Err(e) => { - let s = SmbState::Failed { reason: format!("spawn smbd: {e}") }; + let s = SmbState::Failed { + reason: format!("spawn smbd: {e}"), + }; *self.state.lock() = s.clone(); s } @@ -168,11 +175,15 @@ impl SmbManager { #[allow(unsafe_code)] pub fn reconcile(&self) -> SmbState { let mut g = self.child.lock(); - if g.is_none() { return self.state.lock().clone(); } + if g.is_none() { + return self.state.lock().clone(); + } let shares = match self.write_conf() { Ok(v) => v, Err(e) => { - let s = SmbState::Failed { reason: format!("write smb.conf: {e}") }; + let s = SmbState::Failed { + reason: format!("write smb.conf: {e}"), + }; *self.state.lock() = s.clone(); return s; } @@ -191,8 +202,13 @@ impl SmbManager { // covers this is `nix`, which pulls ~40 transitive deps for a // single signal send. One documented unsafe call is the better // tradeoff for a container-first project. - unsafe { libc::kill(pid, libc::SIGHUP); } - let s = SmbState::Running { pid: pid as u32, shares }; + unsafe { + libc::kill(pid, libc::SIGHUP); + } + let s = SmbState::Running { + pid: pid as u32, + shares, + }; *self.state.lock() = s.clone(); s } else { @@ -212,15 +228,19 @@ impl SmbManager { } fn smbd_present() -> bool { - let Ok(paths) = std::env::var("PATH") else { return false; }; + let Ok(paths) = std::env::var("PATH") else { + return false; + }; for dir in std::env::split_paths(&paths) { - if dir.join("smbd").is_file() { return true; } + if dir.join("smbd").is_file() { + return true; + } } false } const SMB_CONF_GLOBAL: &str = r"[global] -workgroup = PXEFORGE +workgroup = OPENPXE server min protocol = SMB2 smb ports = 445 log level = 1 @@ -235,6 +255,15 @@ lock directory = /tmp state directory = /tmp cache directory = /tmp pid directory = /tmp +# WinPE reconnect hardening. Windows Setup can reboot mid-install and +# reconnect from the same IP; stale sessions/oplocks otherwise cause +# intermittent `net use` failures on the second stage. +reset on zero vc = yes +oplocks = no +kernel oplocks = no +level2 oplocks = no +strict locking = no +deadtime = 1 "; /// Extract a Windows ISO at `iso_path` into `smb_dir//`. Uses @@ -245,7 +274,11 @@ pid directory = /tmp /// Idempotent: if the target dir already contains `sources/boot.wim`, we /// skip extraction. Callers who want a forced re-extract should remove the /// dir first. -pub fn extract_windows_iso(iso_path: &Path, smb_dir: &Path, slug: &str) -> std::io::Result { +pub fn extract_windows_iso( + iso_path: &Path, + smb_dir: &Path, + slug: &str, +) -> std::io::Result { let target = smb_dir.join(slug); if target.join("sources").join("boot.wim").is_file() { tracing::debug!(target: "openpxe::smb", slug, "ISO already extracted, skipping"); @@ -263,7 +296,9 @@ pub fn extract_windows_iso(iso_path: &Path, smb_dir: &Path, slug: &str) -> std:: .stdout(Stdio::null()) .stderr(Stdio::piped()) .output()?; - if out.status.success() { return Ok(target); } + if out.status.success() { + return Ok(target); + } tracing::warn!( target: "openpxe::smb", stderr=%String::from_utf8_lossy(&out.stderr), @@ -278,7 +313,9 @@ pub fn extract_windows_iso(iso_path: &Path, smb_dir: &Path, slug: &str) -> std:: .args(["-C"]) .arg(&target) .output()?; - if out.status.success() { return Ok(target); } + if out.status.success() { + return Ok(target); + } return Err(std::io::Error::other(format!( "bsdtar failed: {}", String::from_utf8_lossy(&out.stderr) @@ -294,7 +331,9 @@ fn which(cmd: &str) -> Option { let paths = std::env::var_os("PATH")?; for dir in std::env::split_paths(&paths) { let p = dir.join(cmd); - if p.is_file() { return Some(p); } + if p.is_file() { + return Some(p); + } } None } @@ -320,7 +359,9 @@ mod tests { let m = SmbManager::new(dir.path().into()); let st = m.start(); // Restore PATH before asserting so any subsequent failure is legible. - if let Some(p) = saved { std::env::set_var("PATH", p); } + if let Some(p) = saved { + std::env::set_var("PATH", p); + } assert_eq!(st, SmbState::SmbdMissing); } @@ -347,5 +388,27 @@ mod tests { assert!(conf.contains("guest ok = yes")); assert!(conf.contains("read only = yes")); assert!(conf.contains("server min protocol = SMB2")); + assert!(conf.contains("workgroup = OPENPXE")); + } + + #[test] + fn write_conf_includes_winpe_reconnect_tuning() { + let dir = tempdir().unwrap(); + let m = SmbManager::new(dir.path().into()); + m.write_conf().unwrap(); + let conf = std::fs::read_to_string(dir.path().join("smb.conf")).unwrap(); + for expected in [ + "reset on zero vc = yes", + "oplocks = no", + "kernel oplocks = no", + "level2 oplocks = no", + "strict locking = no", + "deadtime = 1", + ] { + assert!( + conf.contains(expected), + "missing Windows reconnect Samba option {expected} in:\n{conf}" + ); + } } } diff --git a/crates/iso-store/src/store.rs b/crates/iso-store/src/store.rs index a87ce36..cc5d515 100644 --- a/crates/iso-store/src/store.rs +++ b/crates/iso-store/src/store.rs @@ -3,8 +3,8 @@ use crate::entry::{BootEntry, BootKind, KernelArgs}; use crate::introspect::{introspect, DistroFamily, IntrospectionReport}; use bytes::Bytes; -use parking_lot::RwLock; use openpxe_core::{Error, Result}; +use parking_lot::RwLock; use serde::{Deserialize, Serialize}; use sha2::{Digest, Sha256}; use std::collections::HashMap; @@ -170,7 +170,9 @@ impl IsoStore { let mut entries = tokio::fs::read_dir(self.iso_dir.as_path()).await?; while let Some(e) = entries.next_entry().await? { let p = e.path(); - if p.extension().and_then(|s| s.to_str()) != Some("json") { continue; } + if p.extension().and_then(|s| s.to_str()) != Some("json") { + continue; + } if !p .file_name() .and_then(|s| s.to_str()) @@ -276,7 +278,10 @@ impl IsoStore { /// share or unmount the NFS share entirely. pub async fn delete(&self, id: &str) -> Result<()> { let meta = self.get(id); - let is_local = matches!(meta.as_ref().map(|m| &m.source), Some(IsoSource::Local) | None); + let is_local = matches!( + meta.as_ref().map(|m| &m.source), + Some(IsoSource::Local) | None + ); if is_local { let iso = self.iso_path(id); let meta_path = self.meta_path(id); @@ -319,9 +324,9 @@ impl IsoStore { /// to clean out stale entries. pub fn drop_external_source(&self, mount_id: &str) { let mut g = self.inner.write(); - g.isos.retain(|_, m| { - !matches!(&m.source, IsoSource::Nfs { mount_id: mid, .. } if mid == mount_id) - }); + g.isos.retain( + |_, m| !matches!(&m.source, IsoSource::Nfs { mount_id: mid, .. } if mid == mount_id), + ); } /// Set or clear an ISO's boot password. @@ -371,7 +376,7 @@ impl IsoStore { /// Verify a candidate password against the stored bcrypt hash. /// Returns: - /// - `Ok(true)` — match (or the ISO has no password set; gate is open) + /// - `Ok(true)` — match (or the ISO has no password set; boot is open) /// - `Ok(false)` — mismatch /// - `Err(_)` — id not found, or bcrypt error pub fn verify_password(&self, id: &str, candidate: &str) -> Result { @@ -427,7 +432,10 @@ pub fn generate_boot_entries_for( /// Build `BootEntry`s from the introspection report. URLs are relative — /// the HTTP layer rewrites them with the public base URL per request. fn generate_boot_entries(id: &str, filename: &str, r: &IntrospectionReport) -> Vec { - let title = r.volume_label.clone().unwrap_or_else(|| filename.to_string()); + let title = r + .volume_label + .clone() + .unwrap_or_else(|| filename.to_string()); match r.family { DistroFamily::WindowsPe if r.has_boot_wim => { // Standard wimboot chain. Paths are in-ISO; the HTTP layer maps @@ -451,12 +459,22 @@ fn generate_boot_entries(id: &str, filename: &str, r: &IntrospectionReport) -> V fam if r.kernel_path.is_some() => { let base = format!("iso/{id}"); let kernel_url = format!("{base}{}", r.kernel_path.as_deref().unwrap_or("")); - let initrd_urls = r.initrd_paths.iter().map(|p| format!("{base}{p}")).collect(); - let args = KernelArgs { cmdline: linux_cmdline(fam, id) }; + let initrd_urls = r + .initrd_paths + .iter() + .map(|p| format!("{base}{p}")) + .collect(); + let args = KernelArgs { + cmdline: linux_cmdline(fam, id), + }; vec![BootEntry { id: format!("{id}-linux"), title, - kind: BootKind::LinuxKernel { kernel_url, initrd_urls, args }, + kind: BootKind::LinuxKernel { + kernel_url, + initrd_urls, + args, + }, }] } _ => { @@ -465,7 +483,9 @@ fn generate_boot_entries(id: &str, filename: &str, r: &IntrospectionReport) -> V vec![BootEntry { id: format!("{id}-sanboot"), title: format!("{title} (SAN boot — may fail for >1GiB ISOs)"), - kind: BootKind::SanBootIso { iso_url: format!("iso/{id}.iso") }, + kind: BootKind::SanBootIso { + iso_url: format!("iso/{id}.iso"), + }, }] } } @@ -536,7 +556,11 @@ mod tests { let dir = tempdir().unwrap(); let store = IsoStore::new(dir.path().to_path_buf()); store.ensure_dirs().await.unwrap(); - store.inner.write().isos.insert("alpha".into(), fake_meta("alpha")); + store + .inner + .write() + .isos + .insert("alpha".into(), fake_meta("alpha")); // No password set — verify_password returns Ok(true) for any input. assert!(store.verify_password("alpha", "anything").unwrap()); diff --git a/crates/iso-store/src/windows.rs b/crates/iso-store/src/windows.rs index 4fa23e5..5d534bc 100644 --- a/crates/iso-store/src/windows.rs +++ b/crates/iso-store/src/windows.rs @@ -54,7 +54,10 @@ pub struct WimPatcher { impl WimPatcher { #[must_use] pub fn new(smb_host: String, smb_share: String) -> Self { - Self { smb_host, smb_share } + Self { + smb_host, + smb_share, + } } /// Apply WinPE patches to `boot.wim` inside `extracted_iso_dir`. Returns @@ -72,31 +75,40 @@ impl WimPatcher { let work = match tempfile::tempdir() { Ok(d) => d, - Err(e) => return WinPatchState::Failed { reason: format!("tempdir: {e}") }, + Err(e) => { + return WinPatchState::Failed { + reason: format!("tempdir: {e}"), + } + } }; // Stage the two files we want present at /Windows/System32/. let staging = work.path().join("stage/Windows/System32"); if let Err(e) = std::fs::create_dir_all(&staging) { - return WinPatchState::Failed { reason: format!("staging mkdir: {e}") }; + return WinPatchState::Failed { + reason: format!("staging mkdir: {e}"), + }; } if let Err(e) = std::fs::write(staging.join("winpeshl.ini"), WINPESHL_INI) { - return WinPatchState::Failed { reason: format!("write winpeshl.ini: {e}") }; + return WinPatchState::Failed { + reason: format!("write winpeshl.ini: {e}"), + }; } let startnet = render_startnet(&self.smb_host, &self.smb_share); if let Err(e) = std::fs::write(staging.join("startnet.cmd"), startnet) { - return WinPatchState::Failed { reason: format!("write startnet.cmd: {e}") }; + return WinPatchState::Failed { + reason: format!("write startnet.cmd: {e}"), + }; } // Build a wimlib update command file: // add /Windows/System32 /Windows/System32 let update_file = work.path().join("update.cmd"); - let update_cmd = format!( - "add \"{}\" \"/Windows/System32\"\n", - staging.display() - ); + let update_cmd = format!("add \"{}\" \"/Windows/System32\"\n", staging.display()); if let Err(e) = std::fs::write(&update_file, update_cmd) { - return WinPatchState::Failed { reason: format!("write update.cmd: {e}") }; + return WinPatchState::Failed { + reason: format!("write update.cmd: {e}"), + }; } // Run wimlib-imagex update against image index 2 (WinPE). @@ -120,7 +132,9 @@ impl WimPatcher { String::from_utf8_lossy(&o.stderr) ), }, - Err(e) => WinPatchState::Failed { reason: format!("spawn wimlib-imagex: {e}") }, + Err(e) => WinPatchState::Failed { + reason: format!("spawn wimlib-imagex: {e}"), + }, } } } @@ -133,7 +147,9 @@ fn which(cmd: &str) -> Option { let paths = std::env::var_os("PATH")?; for dir in std::env::split_paths(&paths) { let p = dir.join(cmd); - if p.is_file() { return Some(p); } + if p.is_file() { + return Some(p); + } } None } @@ -174,7 +190,11 @@ fn render_startnet(host: &str, share: &str) -> String { ) .unwrap(); s.push_str(":havenet\r\n"); - writeln!(s, "echo Mapping install media from \\\\{host}\\{share}...\r").unwrap(); + writeln!( + s, + "echo Mapping install media from \\\\{host}\\{share}...\r" + ) + .unwrap(); writeln!( s, ":mapshare\r\nnet use Z: \\\\{host}\\{share} /user:guest \"\" /persistent:no && goto mapped\r\n\ @@ -205,6 +225,23 @@ mod tests { assert!(s.contains("setup.exe")); } + #[test] + fn startnet_primes_workstation_and_surfaces_mapping_errors() { + let s = render_startnet("10.0.0.5", "win11"); + assert!( + s.contains("net start Workstation"), + "WinPE should explicitly start the SMB client before net use:\n{s}" + ); + let net_use_line = s + .lines() + .find(|line| line.contains("net use Z:")) + .expect("net use line"); + assert!( + !net_use_line.contains(">nul"), + "net use errors must remain visible in WinPE console: {net_use_line}" + ); + } + #[test] fn patcher_reports_wimlib_missing_gracefully() { // We don't assume wimlib is present in CI; this checks the missing diff --git a/crates/openpxe/src/main.rs b/crates/openpxe/src/main.rs index 0302628..b7815fb 100644 --- a/crates/openpxe/src/main.rs +++ b/crates/openpxe/src/main.rs @@ -4,16 +4,16 @@ use clap::{Parser, Subcommand}; use openpxe_core::{ - ClientRegistry, Config, DhcpMode, DeploymentQueue, HostBindings, LogBus, LogBusLayer, Metrics, + ClientRegistry, Config, DeploymentQueue, DhcpMode, HostBindings, LogBus, LogBusLayer, Metrics, SettingsStore, }; use openpxe_dhcp_proxy::DhcpProxyServer; use openpxe_http_api::{build_router, AppState}; use openpxe_iso_store::{IsoStore, NfsManager, SmbManager}; -use std::sync::Arc; use openpxe_tftp::TftpServer; use std::net::{Ipv4Addr, SocketAddr}; use std::path::PathBuf; +use std::sync::Arc; use tokio::io::AsyncReadExt; #[derive(Debug, Parser)] @@ -39,7 +39,7 @@ enum Command { /// docker run --rm \ /// -v /my/isos:/seed:ro \ /// -v openpxe-data:/var/lib/openpxe/isos \ - /// openpxe:0.1.0 seed --from /seed + /// openpxe:0.3.2 seed --from /seed Seed { /// Source directory containing one or more `.iso` files. #[arg(long)] @@ -100,7 +100,7 @@ async fn main() -> anyhow::Result<()> { let iso_store = IsoStore::new(config.paths.iso_dir.clone()); iso_store.load_from_disk().await?; let clients = ClientRegistry::new(); - let gates = DeploymentQueue::new(); + let queue = DeploymentQueue::new(); let settings = SettingsStore::load_or_default(&config.paths.work_dir); let hosts = HostBindings::load_or_default(&config.paths.work_dir); let metrics = Metrics::new(); @@ -138,7 +138,7 @@ async fn main() -> anyhow::Result<()> { iso_store: iso_store.clone(), clients: clients.clone(), settings: settings.clone(), - queue: gates.clone(), + queue: queue.clone(), hosts: hosts.clone(), metrics: metrics.clone(), smb: Some(smb.clone()), @@ -210,7 +210,11 @@ async fn run_command(cmd: Command, config: Config) -> anyhow::Result<()> { /// Reuses `IsoStore::begin_upload` / `finish` so the resulting meta on disk /// is identical to a web upload — same slug rules, same introspection, same /// sha256. -async fn seed_from_dir(src: &std::path::Path, config: &Config, dry_run: bool) -> anyhow::Result<()> { +async fn seed_from_dir( + src: &std::path::Path, + config: &Config, + dry_run: bool, +) -> anyhow::Result<()> { let store = IsoStore::new(config.paths.iso_dir.clone()); store.load_from_disk().await?; let mut entries = tokio::fs::read_dir(src).await?; @@ -218,7 +222,12 @@ async fn seed_from_dir(src: &std::path::Path, config: &Config, dry_run: bool) -> let mut skipped = 0u32; while let Some(entry) = entries.next_entry().await? { let p = entry.path(); - if p.extension().and_then(|e| e.to_str()).map(str::to_ascii_lowercase).as_deref() != Some("iso") { + if p.extension() + .and_then(|e| e.to_str()) + .map(str::to_ascii_lowercase) + .as_deref() + != Some("iso") + { continue; } let filename = p @@ -226,8 +235,14 @@ async fn seed_from_dir(src: &std::path::Path, config: &Config, dry_run: bool) -> .and_then(|s| s.to_str()) .ok_or_else(|| anyhow::anyhow!("non-utf8 filename: {}", p.display()))? .to_string(); - println!(" {} ({} bytes)", filename, tokio::fs::metadata(&p).await?.len()); - if dry_run { continue; } + println!( + " {} ({} bytes)", + filename, + tokio::fs::metadata(&p).await?.len() + ); + if dry_run { + continue; + } let mut handle = match store.begin_upload(&filename).await { Ok(h) => h, @@ -242,15 +257,23 @@ async fn seed_from_dir(src: &std::path::Path, config: &Config, dry_run: bool) -> let mut buf = vec![0u8; 1024 * 1024]; loop { let n = file.read(&mut buf).await?; - if n == 0 { break; } + if n == 0 { + break; + } let chunk: bytes::Bytes = buf[..n].to_vec().into(); handle.write_chunk(&chunk).await?; } let meta = handle.finish(&store).await?; - println!(" -> id={} family={:?}", meta.id, meta.introspection.family); + println!( + " -> id={} family={:?}", + meta.id, meta.introspection.family + ); imported += 1; } - println!("\nimported={imported} skipped={skipped} {}", if dry_run { "(dry run)" } else { "" }); + println!( + "\nimported={imported} skipped={skipped} {}", + if dry_run { "(dry run)" } else { "" } + ); Ok(()) } @@ -294,9 +317,8 @@ fn hostname() -> std::io::Result { if let Ok(h) = std::fs::read_to_string("/proc/sys/kernel/hostname") { return Ok(h.trim().to_string()); } - std::env::var("HOSTNAME").map_err(|_| std::io::Error::new( - std::io::ErrorKind::NotFound, "no hostname", - )) + std::env::var("HOSTNAME") + .map_err(|_| std::io::Error::new(std::io::ErrorKind::NotFound, "no hostname")) } fn init_tracing(bus: Arc) { @@ -328,7 +350,10 @@ fn detect_network_info(our_ip: Ipv4Addr) -> NetworkInfo { // `ip -o -f inet addr show` lists every interface with its // `inet a.b.c.d/mask`. We match the line that mentions our IP. - if let Ok(out) = Command::new("ip").args(["-o", "-f", "inet", "addr", "show"]).output() { + if let Ok(out) = Command::new("ip") + .args(["-o", "-f", "inet", "addr", "show"]) + .output() + { if let Ok(text) = String::from_utf8(out.stdout) { for line in text.lines() { if !line.contains(&our_ip.to_string()) { @@ -353,7 +378,10 @@ fn detect_network_info(our_ip: Ipv4Addr) -> NetworkInfo { } // `ip route show default` -> "default via 10.0.0.1 dev enp1s0 ..." - if let Ok(out) = Command::new("ip").args(["route", "show", "default"]).output() { + if let Ok(out) = Command::new("ip") + .args(["route", "show", "default"]) + .output() + { if let Ok(text) = String::from_utf8(out.stdout) { if let Some(line) = text.lines().next() { let mut parts = line.split_whitespace(); @@ -374,7 +402,11 @@ fn detect_network_info(our_ip: Ipv4Addr) -> NetworkInfo { fn prefix_to_dotted(prefix: u8) -> String { let prefix = prefix.min(32); - let mask: u32 = if prefix == 0 { 0 } else { u32::MAX << (32 - prefix) }; + let mask: u32 = if prefix == 0 { + 0 + } else { + u32::MAX << (32 - prefix) + }; format!( "{}.{}.{}.{}", (mask >> 24) & 0xff, diff --git a/crates/tftp/src/server.rs b/crates/tftp/src/server.rs index 428fc9d..bacbcbd 100644 --- a/crates/tftp/src/server.rs +++ b/crates/tftp/src/server.rs @@ -45,7 +45,12 @@ impl TftpServer { clients: Arc, metrics: openpxe_core::Metrics, ) -> Self { - Self { bind, port, clients, metrics } + Self { + bind, + port, + clients, + metrics, + } } pub async fn run(self) -> anyhow::Result<()> { @@ -86,7 +91,9 @@ async fn handle_rrq( let Some(req) = parse_rrq(&packet) else { return Ok(()); }; - let Request { filename, options, .. } = req; + let Request { + filename, options, .. + } = req; // Per-transfer ephemeral socket. let sock = bind_udp(bind_ip, 0)?; @@ -98,7 +105,9 @@ async fn handle_rrq( &peer.ip().to_string(), Some(peer.ip()), None, - ClientEvent::TftpRead { file: filename.clone() }, + ClientEvent::TftpRead { + file: filename.clone(), + }, ); return Ok(()); }; @@ -112,7 +121,9 @@ async fn handle_rrq( &peer.ip().to_string(), Some(peer.ip()), None, - ClientEvent::TftpRead { file: filename.clone() }, + ClientEvent::TftpRead { + file: filename.clone(), + }, ); // Negotiate options. @@ -173,7 +184,9 @@ async fn handle_rrq( // Send one window worth of DATA. for _ in 0..window { - if offset >= total { break; } + if offset >= total { + break; + } let end = (offset + blksize).min(total); let chunk = &file_bytes[offset..end]; let pkt = encode_data(block_no, chunk); @@ -255,20 +268,30 @@ struct Request { } fn parse_rrq(pkt: &[u8]) -> Option { - if pkt.len() < 4 { return None; } + if pkt.len() < 4 { + return None; + } let op = u16::from_be_bytes([pkt[0], pkt[1]]); - if op != OP_RRQ { return None; } + if op != OP_RRQ { + return None; + } let mut rest = &pkt[2..]; let filename = read_cstr(&mut rest)?; let mode = read_cstr(&mut rest)?; let mut options = Vec::new(); while !rest.is_empty() { let Some(k) = read_cstr(&mut rest) else { break }; - if k.is_empty() { break; } + if k.is_empty() { + break; + } let v = read_cstr(&mut rest).unwrap_or_default(); options.push((k.to_ascii_lowercase(), v)); } - Some(Request { filename, mode, options }) + Some(Request { + filename, + mode, + options, + }) } fn read_cstr(buf: &mut &[u8]) -> Option { @@ -316,8 +339,12 @@ async fn recv_ack(sock: &UdpSocket, peer: SocketAddr) -> anyhow::Result { let mut buf = [0u8; 32]; loop { let (n, from) = sock.recv_from(&mut buf).await?; - if from.ip() != peer.ip() { continue; } - if n < 4 { continue; } + if from.ip() != peer.ip() { + continue; + } + if n < 4 { + continue; + } let op = u16::from_be_bytes([buf[0], buf[1]]); match op { OP_ACK => return Ok(u16::from_be_bytes([buf[2], buf[3]])), @@ -344,14 +371,19 @@ async fn wait_for_ack( Ok(Ok(_)) => {} Ok(Err(_)) | Err(_) => { tries += 1; - if tries > 5 { return Ok(false); } + if tries > 5 { + return Ok(false); + } } } } } fn bind_udp(bind: IpAddr, port: u16) -> anyhow::Result { - let domain = match bind { IpAddr::V4(_) => Domain::IPV4, IpAddr::V6(_) => Domain::IPV6 }; + let domain = match bind { + IpAddr::V4(_) => Domain::IPV4, + IpAddr::V6(_) => Domain::IPV6, + }; let sock = Socket::new(domain, Type::DGRAM, Some(Protocol::UDP))?; sock.set_reuse_address(true)?; sock.set_nonblocking(true)?; @@ -382,7 +414,9 @@ pub fn plan_window( let mut o = offset; let mut b = starting_block; for _ in 0..window { - if o >= total { break; } + if o >= total { + break; + } let end = (o + blksize).min(total); out.push((b, end - o)); o = end; @@ -453,7 +487,7 @@ mod tests { assert_eq!(p, vec![(65534, 1024), (65535, 1024)]); let p2 = plan_window(2048, 2048, 1024, 2, 0); assert!(p2.is_empty()); // nothing past EOF - // And a cross-boundary case: + // And a cross-boundary case: let p3 = plan_window(3072, 0, 1024, 3, 65535); assert_eq!(p3, vec![(65535, 1024), (0, 1024), (1, 1024)]); } diff --git a/crates/webui/src/app.css b/crates/webui/src/app.css index ee303af..85fd583 100644 --- a/crates/webui/src/app.css +++ b/crates/webui/src/app.css @@ -281,8 +281,7 @@ label.check input { accent-color: var(--accent); } /* ── Imaging progress widget ─────────────────────────────────────── Animated brand mark paired with a horizontal progress bar; surfaces - on Dashboard and the Queue tab. Renamed from `.forge-progress` in - v0.3.0 — the anvil-themed naming is gone with the rebrand. */ + on Dashboard and the Queue tab. */ .queue-progress { display: flex; align-items: center; gap: 16px; padding: 16px; @@ -380,7 +379,7 @@ tr.unbootable td:first-child { border-left: 3px solid var(--warn); } .form-row { display: grid; grid-template-columns: repeat(4, 1fr); gap: 10px 14px; } @media (max-width: 900px) { .form-row { grid-template-columns: 1fr; } } -/* ── Gate queue "horse race" visual ──────────────────────────────── */ +/* ── Queued deployment visual ────────────────────────────────────── */ .queue-track { display: grid; gap: 6px; padding: 10px 0; diff --git a/crates/webui/src/app.js b/crates/webui/src/app.js index bc4e2d4..f5eefb9 100644 --- a/crates/webui/src/app.js +++ b/crates/webui/src/app.js @@ -246,7 +246,7 @@ return el('div', {class:'grid'}, [networkCard]); }, - gate: async () => { + queue: async () => { const [{ entries = [] }, isos] = await Promise.all([ getJSON('/api/queue'), getJSON('/api/isos'), ]); @@ -267,7 +267,7 @@ if (!j.ok) { msg.textContent = 'Assign failed: ' + (j.error || 'unknown'); msg.className='msg err'; return; } msg.textContent = 'Launched ' + j.assigned + ' client' + (j.assigned===1?'':'s') + ' → ' + j.target; msg.className = 'msg ok'; - render('gate'); + render('queue'); }; const track = entries.length @@ -284,12 +284,12 @@ : el('span', {class:'tag accent'}, 'waiting')), el('button', {class:'ghost', onclick: async () => { await fetch('/api/queue/' + encodeURIComponent(g.id), {method:'DELETE'}); - render('gate'); + render('queue'); }}, 'Release'), ])) ) : el('div', {class:'empty'}, - 'No clients at the gate. Boot a client and choose "Queued Deployment" in the PXE menu.'); + 'No clients queued. Boot a client and choose "Queued Deployment" in the PXE menu.'); const imaging = entries.filter(g => g.assigned_target).length; @@ -299,7 +299,7 @@ queueProgressWidget(imaging, entries.length), ]), el('div', {class:'card'}, [ - el('header', {}, el('h2', {}, 'Launch an image across the gate')), + el('header', {}, el('h2', {}, 'Launch image for queued clients')), el('div', {class:'body'}, [ el('label', {class:'field'}, [ el('span', {class:'name'}, 'Target image'), @@ -312,7 +312,7 @@ ]), el('div', {class:'card'}, [ el('header', {}, [ - el('h2', {}, 'Gate positions'), + el('h2', {}, 'Queue positions'), el('span', {class:'sub'}, entries.length + ' waiting'), ]), el('div', {class:'body'}, track), @@ -854,7 +854,7 @@ const viewTitles = { dashboard: 'Dashboard', network: 'Network', - gate: 'Queue', + queue: 'Queue', storage: 'Storage', hosts: 'Hosts', terminal: 'Terminal', diff --git a/crates/webui/src/index.html b/crates/webui/src/index.html index 3cf7a3f..a236db8 100644 --- a/crates/webui/src/index.html +++ b/crates/webui/src/index.html @@ -29,7 +29,7 @@
OpenPXE -
v0.3.1
+
v0.3.2