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
This commit is contained in:
Miles Ward
2026-05-21 02:13:08 -04:00
parent 9c6903351f
commit bbdbb8df43
37 changed files with 1387 additions and 589 deletions
+4 -1
View File
@@ -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);
}
+18 -12
View File
@@ -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;
+13 -5
View File
@@ -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<PathBuf>,
/// 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() {
+1 -1
View File
@@ -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"`,
+2 -2
View File
@@ -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};
+62 -11
View File
@@ -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"));
+28 -24
View File
@@ -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/<id>`; 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<IpAddr>,
@@ -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<IpAddr>, arch: Option<ClientArch>) -> Gate {
pub fn join(&self, mac: &str, ip: Option<IpAddr>, arch: Option<ClientArch>) -> 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<Arc<Notify>> {
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<Gate> {
/// the entry was released/expired between requests.
pub fn touch(&self, entry_id: &str) -> Option<QueueEntry> {
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<Gate> {
pub fn release(&self, entry_id: &str) -> Option<QueueEntry> {
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<Gate> {
pub fn list(&self) -> Vec<QueueEntry> {
let guard = self.inner.read();
let mut v: Vec<_> = guard.values().map(QueueEntryInner::snapshot).collect();
v.sort_by_key(|g| g.position);
+14 -8
View File
@@ -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();
+14 -4
View File
@@ -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,
}
+34 -13
View File
@@ -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<UdpSocke
fn format_mac(chaddr: &[u8]) -> String {
let take = chaddr.iter().take(6).copied().collect::<Vec<_>>();
take.iter().map(|b| format!("{b:02x}")).collect::<Vec<_>>().join(":")
take.iter()
.map(|b| format!("{b:02x}"))
.collect::<Vec<_>>()
.join(":")
}
/// Walk raw DHCP options looking for option 93 (Client System Architecture)
@@ -217,10 +231,17 @@ fn extract_raw_arch(packet: &[u8]) -> Option<u16> {
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));
+245 -127
View File
@@ -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<AppState>) -> 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<AppState>,
Query(p): Query<BootMenuParams>,
) -> Response {
state.metrics.record_http(openpxe_core::HttpRoute::BootScript);
async fn boot_top_menu(State(state): State<AppState>, Query(p): Query<BootMenuParams>) -> 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<String>) -> 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<Response> {
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::<u64>() {
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::<u64>().ok()).unwrap_or(0);
let end = parts.next().and_then(|s| s.parse::<u64>().ok()).unwrap_or(total.saturating_sub(1));
(start, end.min(total.saturating_sub(1)), true)
let start = parts
.next()
.and_then(|s| s.parse::<u64>().ok())
.unwrap_or(0);
let end = parts
.next()
.and_then(|s| s.parse::<u64>().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<AppState>,
mut multipart: Multipart,
) -> Response {
async fn api_upload_iso(State(state): State<AppState>, 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<AppState>) -> Response {
@@ -517,7 +600,11 @@ async fn readyz(State(state): State<AppState>) -> 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<AppState>) -> Json<serde_json::Value> {
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<AppState>) -> Json<serde_json::Value> {
"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<AppState>) -> Json<serde_json::Value
}
#[derive(Debug, Deserialize)]
struct GateJoinParams {
struct QueueJoinParams {
/// Client MAC from iPXE's `${mac}` variable. iPXE substitutes before
/// the HTTP request so we receive a plain colon-separated MAC.
mac: Option<String>,
@@ -634,7 +733,7 @@ struct GateJoinParams {
/// until poll returns an actual boot script.
async fn api_queue_join(
State(state): State<AppState>,
Query(p): Query<GateJoinParams>,
Query(p): Query<QueueJoinParams>,
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<String>,
) -> 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/<id>.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<String>,
}
async fn api_queue_assign(
State(state): State<AppState>,
Json(body): Json<GateAssignBody>,
Json(body): Json<QueueAssignBody>,
) -> Json<serde_json::Value> {
let ids = if body.entry_ids.is_empty() {
state.queue.list().into_iter().map(|g| g.id).collect::<Vec<_>>()
state
.queue
.list()
.into_iter()
.map(|g| g.id)
.collect::<Vec<_>>()
} 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<AppState>) -> Json<serde_json::Value>
Json(json!({ "mounts": state.nfs.list() }))
}
async fn api_nfs_add(
State(state): State<AppState>,
Json(req): Json<NfsAddRequest>,
) -> Response {
async fn api_nfs_add(State(state): State<AppState>, Json(req): Json<NfsAddRequest>) -> 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<AppState>,
AxumPath(id): AxumPath<String>,
) -> Response {
async fn api_nfs_remove(State(state): State<AppState>, AxumPath(id): AxumPath<String>) -> 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<AppState>,
AxumPath(id): AxumPath<String>,
) -> Response {
async fn api_nfs_scan(State(state): State<AppState>, AxumPath(id): AxumPath<String>) -> 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<AppState>) -> 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);
}
}
+156 -34
View File
@@ -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}"
);
}
}
}
}
+33 -10
View File
@@ -31,7 +31,9 @@ pub fn lookup(iso_path: &Path, in_iso_path: &str) -> Option<FileLocation> {
.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<FileLocation> {
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
}
}
+10 -6
View File
@@ -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<AppState>) -> Json<serde_json::Value> {
/// keep streaming new lines as they arrive).
pub async fn clear(State(state): State<AppState>) -> Json<serde_json::Value> {
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 }))
}
+73 -47
View File
@@ -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<String, String> {
"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<String, String> {
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<String, String> {
async fn queue_command(s: &AppState, args: &[String]) -> Result<String, String> {
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 <iso_boot_entry_id>".to_string()
})?;
let target = args
.get(1)
.ok_or_else(|| "usage: queue assign-all <iso_boot_entry_id>".to_string())?;
let found = s
.iso_store
.list()
@@ -219,30 +234,32 @@ async fn gate_command(s: &AppState, args: &[String]) -> Result<String, String> {
}
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 <entry_id> <iso_boot_entry_id>".to_string())?;
.ok_or_else(|| "usage: queue assign <entry_id> <iso_boot_entry_id>".to_string())?;
let target = args
.get(2)
.ok_or_else(|| "usage: gate assign <entry_id> <iso_boot_entry_id>".to_string())?;
.ok_or_else(|| "usage: queue assign <entry_id> <iso_boot_entry_id>".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 <entry_id>".to_string())?;
let entry_id = args
.get(1)
.ok_or_else(|| "usage: queue release <entry_id>".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<String, String> {
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<String, String> {
}
}
Some("unmount") => {
let id = args.get(1).ok_or_else(|| "usage: nfs unmount <id>".to_string())?;
let id = args
.get(1)
.ok_or_else(|| "usage: nfs unmount <id>".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 <id>".to_string())?;
let id = args
.get(1)
.ok_or_else(|| "usage: nfs scan <id>".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<String, String> {
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 <entry_id> <target> assign one gate to a boot entry
gate assign-all <target> assign every waiting gate
gate release <entry_id> release one gate
queue list list queued clients
queue assign <entry_id> <target> assign one queued client to a boot entry
queue assign-all <target> assign every waiting client
queue release <entry_id> release one queued client
nfs list list NFS mounts
nfs mount <s>:<e> [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"));
}
}
+296 -123
View File
@@ -38,13 +38,12 @@ fn fake_alpine_iso() -> Vec<u8> {
}
fn multipart_iso_body(filename: &str, bytes: &[u8]) -> (String, Vec<u8>) {
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<u8>) {
.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/<id>.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/<id>.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);
}
+3 -1
View File
@@ -51,7 +51,9 @@ pub fn asset_slice(name: &str) -> Option<std::borrow::Cow<'static, [u8]>> {
/// Enumerate embedded asset filenames. Useful for startup logging so the
/// operator can immediately tell which architectures will work.
pub fn list_assets() -> Vec<String> {
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.
+38 -11
View File
@@ -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);
}
+3 -1
View File
@@ -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};
+3 -9
View File
@@ -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)
+80 -17
View File
@@ -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/<slug>/` becomes a share named `<slug>`. Returns
/// the sorted list.
pub fn discover_shares(&self) -> Vec<String> {
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<String> = 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/<slug>/`. 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<PathBuf> {
pub fn extract_windows_iso(
iso_path: &Path,
smb_dir: &Path,
slug: &str,
) -> std::io::Result<PathBuf> {
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<PathBuf> {
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}"
);
}
}
}
+37 -13
View File
@@ -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<bool> {
@@ -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<BootEntry> {
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());
+50 -13
View File
@@ -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 <stage>/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<PathBuf> {
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
+50 -18
View File
@@ -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<String> {
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<LogBus>) {
@@ -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,
+49 -15
View File
@@ -45,7 +45,12 @@ impl TftpServer {
clients: Arc<ClientRegistry>,
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<Request> {
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<String> {
@@ -316,8 +339,12 @@ async fn recv_ack(sock: &UdpSocket, peer: SocketAddr) -> anyhow::Result<u16> {
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<UdpSocket> {
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)]);
}
+2 -3
View File
@@ -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;
+7 -7
View File
@@ -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',
+1 -1
View File
@@ -29,7 +29,7 @@
<img src="/assets/logo.svg" alt="" />
<div>
<strong>OpenPXE</strong>
<div class="sub">v<span data-bind="version">0.3.1</span></div>
<div class="sub">v<span data-bind="version">0.3.2</span></div>
</div>
</div>
<nav>
+17 -9
View File
@@ -15,22 +15,30 @@ pub fn index_html(base_url: &str) -> String {
}
#[must_use]
pub fn app_js() -> &'static str { APP_JS }
pub fn app_js() -> &'static str {
APP_JS
}
#[must_use]
pub fn app_css() -> &'static str { APP_CSS }
pub fn app_css() -> &'static str {
APP_CSS
}
#[must_use]
pub fn logo_svg() -> &'static str { LOGO_SVG }
pub fn logo_svg() -> &'static str {
LOGO_SVG
}
/// Larger, faster-cycling rainbow disc — used for the page-load
/// transition and the imaging-progress widget on Dashboard / Queue.
/// Pure SVG + SMIL, no JS, no GIF.
#[must_use]
pub fn loader_svg() -> &'static str { LOADER_SVG }
pub fn loader_svg() -> &'static str {
LOADER_SVG
}
const INDEX_HTML: &str = include_str!("index.html");
const APP_CSS: &str = include_str!("app.css");
const APP_JS: &str = include_str!("app.js");
const LOGO_SVG: &str = include_str!("logo.svg");
const LOADER_SVG: &str = include_str!("loader.svg");
const INDEX_HTML: &str = include_str!("index.html");
const APP_CSS: &str = include_str!("app.css");
const APP_JS: &str = include_str!("app.js");
const LOGO_SVG: &str = include_str!("logo.svg");
const LOADER_SVG: &str = include_str!("loader.svg");