Compare commits

...
1 Commits
Author SHA1 Message Date
Miles WardandClaude Opus 4.8 5da05a519d v0.6.2: Mythos Validation — full-codebase polish, hot-path optimizations, dhcproto 0.15
Codebase-wide review pass: finish or remove every loose end, take the
safe performance wins on the serving hot paths, and refresh the
dependency tree for reliability. No behavior changes for working
clients; legacy clients get clearer protocol errors.

Finalize / cleanup:
- Remove mac_allowlist/subnet_allowlist config fields — parsed but never
  enforced since introduction; the operator wants line-of-sight serving,
  so the honest fix is deletion, not wiring.
- Remove dead ClientRegistry API (get, set_selected_target,
  always-None selected_target field, never-emitted DhcpRequest/
  HttpIsoAsset events).
- TFTP: reject WRQ with ERR_ILLEGAL_OP and non-octet modes with a clear
  error instead of silent timeouts (legacy-client friendliness); fold
  plan_window into cfg(test); drop the unused-constant keep-alive hack.
- rustfmt sweep over the six files with accumulated drift.

Hot-path optimizations (all behavior-preserving):
- Serve embedded iPXE binaries zero-copy (Cow over rodata) on both TFTP
  and HTTP — was a ~1 MiB heap copy per boot file request.
- Cache the composited PXE boot-menu background PNG keyed on the
  branding logo revision — was ~50-200 ms of image work per booting
  client; now one compose per logo change.
- Run bcrypt verify/hash on the blocking pool (boot password gate,
  login, setup, credential rotation) so CPU-heavy auth can't stall the
  workers streaming ISO ranges to imaging machines.
- iso_raw: reuse the already-cloned IsoMeta for path resolution instead
  of a second registry lock + deep clone per range request.
- DriverEscalation: amortize the TTL sweep (1-min interval + inline
  staleness check) instead of an O(map) retain per DHCP packet.
- format_mac: one allocation instead of four per datagram.
- Introspection haystack sized to min(scan cap, file size) — was
  guaranteed a 32 MiB realloc on every large-ISO probe.

Robustness:
- parse_range: malformed Range headers are now ignored per RFC 7233
  (200 + full body) instead of answered with a bogus 206.

Dependencies:
- dhcproto 0.12 -> 0.15: drops the deprecated/unmaintained
  trust-dns-proto from the tree (hickory-proto), three releases of DHCP
  option coverage. Compiles + passes the full suite unchanged.
- socket2 0.6 (dedupes tree), bcrypt 0.19, tower-http 0.6.11 (sheds
  iri-string), tokio 1.52.3 / hyper 1.10 lockfile refresh; dead nom
  workspace entry removed; requested versions synced to shipped reality.

Co-Authored-By: Claude Opus 4.8 (1M context) <[email protected]>
2026-06-09 16:44:55 -04:00
21 changed files with 659 additions and 430 deletions
Generated
+387 -278
View File
File diff suppressed because it is too large Load Diff
+11 -7
View File
@@ -12,7 +12,7 @@ members = [
]
[workspace.package]
version = "0.6.1"
version = "0.6.2"
edition = "2021"
rust-version = "1.95"
license = "MIT OR Apache-2.0"
@@ -20,21 +20,23 @@ repository = "https://gitea.milesward.dev/mward4/OpenPXE"
authors = ["OpenPXE contributors"]
[workspace.dependencies]
tokio = { version = "1.40", features = ["full"] }
tokio = { version = "1.52", features = ["full"] }
tokio-util = { version = "0.7", features = ["io"] }
tokio-stream = { version = "0.1", features = ["sync"] }
futures = "0.3"
async-trait = "0.1"
dhcproto = "0.12"
socket2 = { version = "0.5", features = ["all"] }
# v0.6.2: dhcproto 0.15 drops the deprecated trust-dns-proto dependency
# (replaced by hickory-proto) and carries three releases of DHCP option
# coverage accumulated upstream — both directly relevant to the proxy core.
dhcproto = "0.15"
socket2 = { version = "0.6", features = ["all"] }
bytes = "1.7"
nom = "7.1"
axum = { version = "0.7", features = ["macros", "multipart", "http2"] }
tower = "0.5"
tower-http = { version = "0.6", features = ["fs", "trace", "cors", "limit"] }
hyper = "1.4"
hyper = "1.9"
reqwest = { version = "0.12", default-features = false, features = ["rustls-tls", "stream", "json"] }
serde = { version = "1.0", features = ["derive"] }
@@ -52,9 +54,11 @@ thiserror = "2.0"
clap = { version = "4.5", features = ["derive", "env"] }
uuid = { version = "1.10", features = ["v4", "serde"] }
time = { version = "0.3", features = ["serde", "serde-human-readable", "formatting", "macros"] }
# sha2 stays 0.10 deliberately: bergshamra-crypto requires ^0.10, and
# bumping to 0.11 would split the RustCrypto digest stack in the tree.
sha2 = "0.10"
hex = "0.4"
bcrypt = "0.15"
bcrypt = "0.19"
parking_lot = "0.12"
rust-embed = { version = "8.5", features = ["include-exclude"] }
+2 -4
View File
@@ -125,9 +125,7 @@ impl AdminStore {
{
let mut g = self.inner.write();
if g.admin.is_some() {
return Err(Error::Invalid(
"admin account already configured".into(),
));
return Err(Error::Invalid("admin account already configured".into()));
}
g.admin = Some(admin.clone());
}
@@ -346,7 +344,7 @@ mod tests {
assert!(s.bootstrap("", "hunter2hunter2").is_err());
assert!(s.bootstrap("ad:min", "hunter2hunter2").is_err()); // ':' reserved
assert!(s.bootstrap("admin", "short").is_err()); // <8 chars
// 65-char username is too long.
// 65-char username is too long.
let long = "a".repeat(65);
assert!(s.bootstrap(&long, "hunter2hunter2").is_err());
}
-17
View File
@@ -12,11 +12,9 @@ use time::OffsetDateTime;
#[derive(Debug, Clone, Serialize, Deserialize)]
pub enum ClientEvent {
DhcpDiscover,
DhcpRequest,
PxeBootServerRequest,
TftpRead { file: String },
HttpScriptFetch { target: String },
HttpIsoAsset { file: String },
}
#[derive(Debug, Clone, Serialize, Deserialize)]
@@ -32,8 +30,6 @@ pub struct ClientSnapshot {
// Events are left with default serialization (9-tuple) — they're
// diagnostic only and not consumed by the UI today.
pub events: Vec<(OffsetDateTime, ClientEvent)>,
/// The boot target (ISO id) last selected via the iPXE menu, if any.
pub selected_target: Option<String>,
}
#[derive(Debug, Default)]
@@ -66,7 +62,6 @@ impl ClientRegistry {
first_seen: now,
last_seen: now,
events: Vec::new(),
selected_target: None,
});
entry.last_seen = now;
if ip.is_some() {
@@ -84,13 +79,6 @@ impl ClientRegistry {
}
}
pub fn set_selected_target(&self, mac: &str, target: Option<String>) {
let mut guard = self.inner.write();
if let Some(c) = guard.get_mut(mac) {
c.selected_target = target;
}
}
#[must_use]
pub fn list(&self) -> Vec<ClientSnapshot> {
let guard = self.inner.read();
@@ -99,9 +87,4 @@ impl ClientRegistry {
v.sort_by_key(|c| std::cmp::Reverse(c.last_seen));
v
}
#[must_use]
pub fn get(&self, mac: &str) -> Option<ClientSnapshot> {
self.inner.read().get(mac).cloned()
}
}
-6
View File
@@ -40,10 +40,6 @@ pub struct NetworkConfig {
pub dhcp_port: u16,
/// UDP port for PXE Boot Server discovery. Standard is 4011.
pub pxe_port: u16,
/// Optional allowlist of client MAC prefixes (OUI). Empty = serve everyone.
pub mac_allowlist: Vec<String>,
/// Optional allowlist of subnets (CIDR). Empty = serve everyone.
pub subnet_allowlist: Vec<String>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize, Default)]
@@ -105,8 +101,6 @@ impl Default for NetworkConfig {
dhcp_bind: IpAddr::V4(Ipv4Addr::UNSPECIFIED),
dhcp_port: 67,
pxe_port: 4011,
mac_allowlist: Vec::new(),
subnet_allowlist: Vec::new(),
}
}
}
+4 -1
View File
@@ -303,7 +303,10 @@ mod tests {
smtp_host: "smtp.example.com".into(),
..Default::default()
});
assert!(matches!(r, Err(Error::Invalid(_))), "missing recipient should reject");
assert!(
matches!(r, Err(Error::Invalid(_))),
"missing recipient should reject"
);
s.replace(NotifyConfig {
enabled: true,
kind: NotifyKind::Smtp,
+3 -5
View File
@@ -181,10 +181,7 @@ mod tests {
Ipv4Addr::new(192, 168, 1, 255)
);
assert_eq!(
subnet_broadcast(
Ipv4Addr::new(10, 5, 3, 7),
Ipv4Addr::new(255, 255, 0, 0)
),
subnet_broadcast(Ipv4Addr::new(10, 5, 3, 7), Ipv4Addr::new(255, 255, 0, 0)),
Ipv4Addr::new(10, 5, 255, 255)
);
}
@@ -196,7 +193,8 @@ mod tests {
// and confirm send_magic transmits the exact 102-byte packet.
let rx = UdpSocket::bind(SocketAddrV4::new(Ipv4Addr::LOCALHOST, 0)).unwrap();
let port = rx.local_addr().unwrap().port();
rx.set_read_timeout(Some(std::time::Duration::from_secs(2))).unwrap();
rx.set_read_timeout(Some(std::time::Duration::from_secs(2)))
.unwrap();
let packet = magic_packet([0x0a, 0x1b, 0x2c, 0x3d, 0x4e, 0x5f]);
let sent = send_magic(&packet, &[Ipv4Addr::LOCALHOST], port).unwrap();
+42 -7
View File
@@ -35,6 +35,13 @@ const ENTRY_TTL: Duration = Duration::from_mins(30);
/// entry — escalation is best-effort, never a memory-growth vector.
const MAX_ENTRIES: usize = 4096;
/// How often (at most) the whole map is swept for expired entries.
/// Correctness doesn't depend on the sweep — a stale entry is also
/// detected inline when its MAC next appears — so the sweep only bounds
/// memory for MACs that never return, and amortizing it keeps the
/// per-packet path O(1) instead of O(map).
const PRUNE_INTERVAL: Duration = Duration::from_mins(1);
#[derive(Debug, Clone, Copy)]
struct Entry {
mode: DriverMode,
@@ -45,10 +52,26 @@ struct Entry {
last_seen: Instant,
}
#[derive(Debug)]
struct Inner {
map: HashMap<String, Entry>,
/// When the last full TTL sweep ran — see [`PRUNE_INTERVAL`].
last_prune: Instant,
}
impl Default for Inner {
fn default() -> Self {
Self {
map: HashMap::new(),
last_prune: Instant::now(),
}
}
}
/// Tracks per-MAC driver-mode escalation. Cheap to share via `Arc`.
#[derive(Debug, Default)]
pub struct DriverEscalation {
inner: Mutex<HashMap<String, Entry>>,
inner: Mutex<Inner>,
}
impl DriverEscalation {
@@ -75,11 +98,23 @@ impl DriverEscalation {
fn decide_at(&self, mac: &str, primary: bool, now: Instant) -> DriverMode {
let mut g = self.inner.lock();
g.retain(|_, e| now.duration_since(e.last_seen) < ENTRY_TTL);
if now.duration_since(g.last_prune) >= PRUNE_INTERVAL {
g.map
.retain(|_, e| now.duration_since(e.last_seen) < ENTRY_TTL);
g.last_prune = now;
}
// Inline staleness check: a MAC whose entry outlived the TTL starts
// fresh even when the amortized sweep above hasn't caught it yet.
if g.map
.get(mac)
.is_some_and(|e| now.duration_since(e.last_seen) >= ENTRY_TTL)
{
g.map.remove(mac);
}
match g.get_mut(mac) {
match g.map.get_mut(mac) {
None => {
g.insert(
g.map.insert(
mac.to_owned(),
Entry {
mode: DriverMode::Firmware,
@@ -88,8 +123,8 @@ impl DriverEscalation {
last_seen: now,
},
);
if g.len() > MAX_ENTRIES {
evict_oldest(&mut g);
if g.map.len() > MAX_ENTRIES {
evict_oldest(&mut g.map);
}
DriverMode::Firmware
}
@@ -114,7 +149,7 @@ impl DriverEscalation {
fn confirm_at(&self, mac: &str, now: Instant) {
let mut g = self.inner.lock();
if let Some(e) = g.get_mut(mac) {
if let Some(e) = g.map.get_mut(mac) {
e.awaiting_confirm = false;
e.last_seen = now;
}
+10 -5
View File
@@ -237,11 +237,16 @@ 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(":")
use std::fmt::Write;
// One allocation — this runs for every PXE datagram we answer.
let mut s = String::with_capacity(17);
for (i, b) in chaddr.iter().take(6).enumerate() {
if i > 0 {
s.push(':');
}
let _ = write!(s, "{b:02x}");
}
s
}
/// Walk raw DHCP options looking for option 93 (Client System Architecture)
+71 -14
View File
@@ -35,7 +35,7 @@ use openpxe_core::{
encoding::pct_encode, ext_for_mime, wol, BootEvent, ClientEvent, DeployProfile, Error,
LogoSlot, NotifyConfig, Settings, SsoConfig, ALLOWED_LOGO_MIMES, MAX_LOGO_BYTES,
};
use openpxe_ipxe_assets::asset_bytes;
use openpxe_ipxe_assets::asset_slice;
use openpxe_iso_store::{
render_template, IsoCategory, IsoMeta, IsoSource, NfsAddRequest, SftpAddRequest, SmbAddRequest,
SmbState, UnattendedKind, UnattendedMeta,
@@ -406,6 +406,22 @@ fn bundled_logo_response() -> Response {
/// brand mark falls back to the *default* background for the PXE screen
/// (the WebUI still renders the SVG natively in the top-left).
async fn ui_pxe_logo(State(state): State<AppState>) -> Response {
// The composite is a pure function of the uploaded logo, so the
// encoded PNG is cached keyed on the branding revision — an upload
// or clear bumps the rev and invalidates it. The response headers
// stay `no-cache` (clients must refetch); only the server-side
// ~50-200 ms decode/compose/encode is skipped per boot.
let rev = state.branding.logo_rev();
let cached = state
.pxe_bg_cache
.lock()
.as_ref()
.filter(|(r, _)| *r == rev)
.map(|(_, png)| png.clone());
if let Some(png) = cached {
return pxe_png_response(png);
}
// Resolve the operator's raster upload, if any and if it's a format
// iPXE/our compositor can consume. SVG (or a missing/unreadable
// file) yields `None`, which composes the default background.
@@ -450,6 +466,12 @@ async fn ui_pxe_logo(State(state): State<AppState>) -> Response {
.into_response();
}
};
let png = bytes::Bytes::from(composed);
*state.pxe_bg_cache.lock() = Some((rev, png.clone()));
pxe_png_response(png)
}
fn pxe_png_response(png: bytes::Bytes) -> Response {
(
[
(header::CONTENT_TYPE, HeaderValue::from_static("image/png")),
@@ -460,7 +482,7 @@ async fn ui_pxe_logo(State(state): State<AppState>) -> Response {
HeaderValue::from_static("no-cache, max-age=0"),
),
],
composed,
png,
)
.into_response()
}
@@ -638,7 +660,23 @@ async fn boot_sub(
));
}
Some(token) => {
match state.iso_store.verify_password(&iso.id, token) {
// bcrypt verify costs ~100-200 ms of pure
// CPU and this path is unauthenticated —
// run it on the blocking pool so password
// probes can't stall the workers that are
// streaming ISO bytes to imaging machines.
let store = state.iso_store.clone();
let iso_id = iso.id.clone();
let tok = token.to_string();
let verdict = match tokio::task::spawn_blocking(move || {
store.verify_password(&iso_id, &tok)
})
.await
{
Ok(v) => v,
Err(e) => Err(openpxe_core::Error::Other(e.into())),
};
match verdict {
Ok(true) => { /* fall through to render the entry */ }
Ok(false) => {
// Don't log the candidate — just the
@@ -735,9 +773,15 @@ async fn ipxe_binary(AxumPath(name): AxumPath<String>) -> Response {
if name.contains('/') || name.contains('\\') {
return (StatusCode::BAD_REQUEST, "invalid name").into_response();
}
let Some(bytes) = asset_bytes(&name) else {
let Some(data) = asset_slice(&name) else {
return (StatusCode::NOT_FOUND, "no such ipxe asset").into_response();
};
// Release builds embed the asset in rodata — serve it without the
// ~1 MiB per-request heap copy `into_owned` would cost.
let bytes = match data {
std::borrow::Cow::Borrowed(b) => bytes::Bytes::from_static(b),
std::borrow::Cow::Owned(v) => bytes::Bytes::from(v),
};
(
[
(
@@ -768,7 +812,10 @@ async fn iso_raw(
};
match &meta.source {
IsoSource::Local => {
let Some(path) = state.iso_store.iso_path_for(id) else {
// `local_path(&meta)` reuses the meta we already cloned —
// `iso_path_for(id)` would re-lock and deep-clone it again,
// hundreds of times per sanboot install.
let Some(path) = state.iso_store.local_path(&meta) else {
return (StatusCode::NOT_FOUND, "no such iso").into_response();
};
match stream_file_range(&path, headers.get(header::RANGE)).await {
@@ -1027,15 +1074,25 @@ fn parse_range(h: Option<&HeaderValue>, total: u64) -> Option<(u64, u64, bool)>
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));
// RFC 7233 §3.1: a Range header we can't parse is *ignored* (200 +
// full body), never coerced into a bogus 206 claiming the whole
// file. Only `first-pos[-last-pos]` with numeric positions reaches
// the partial path; `None` is reserved for syntactically valid but
// unsatisfiable ranges (→ 416).
let full = Some((0, total.saturating_sub(1), false));
let Some((start_s, end_s)) = spec.split_once('-') else {
return full;
};
let Ok(start) = start_s.trim().parse::<u64>() else {
return full;
};
let end = if end_s.trim().is_empty() {
total.saturating_sub(1)
} else if let Ok(e) = end_s.trim().parse::<u64>() {
e
} else {
return full;
};
if start >= total {
return None;
}
+33 -9
View File
@@ -154,10 +154,15 @@ pub fn session_cookie(session: &str) -> String {
fn parse_cookie(headers: &axum::http::HeaderMap) -> Option<String> {
// `Cookie: a=b; c=d` parsing — small enough not to drag in a crate.
// Two-step strip (name, then '=') keeps this allocation-free per
// candidate and can't match a longer cookie name sharing the prefix.
let raw = headers.get(header::COOKIE)?.to_str().ok()?;
for part in raw.split(';') {
let part = part.trim();
if let Some(v) = part.strip_prefix(&format!("{SESSION_COOKIE}=")) {
if let Some(v) = part
.strip_prefix(SESSION_COOKIE)
.and_then(|rest| rest.strip_prefix('='))
{
return Some(v.to_string());
}
}
@@ -246,7 +251,14 @@ pub async fn api_setup(State(state): State<AppState>, Json(body): Json<SetupBody
)
.into_response();
}
match state.admin.bootstrap(&body.username, &body.password) {
// bcrypt hashing is ~100-200 ms of pure CPU (and `bootstrap` also
// persists to disk synchronously) — keep it off the async workers.
let admin = state.admin.clone();
let result =
tokio::task::spawn_blocking(move || admin.bootstrap(&body.username, &body.password))
.await
.unwrap_or_else(|e| Err(openpxe_core::Error::Other(e.into())));
match result {
Ok(pub_) => {
let session = state.sessions.create(&pub_.username);
login_response(StatusCode::CREATED, &pub_, &session)
@@ -271,8 +283,13 @@ pub struct LoginBody {
pub async fn api_login(State(state): State<AppState>, Json(body): Json<LoginBody>) -> Response {
// Brief, deliberately vague — "invalid credentials" rather than
// "no such user" / "wrong password". Same anti-enumeration posture
// as Sonarr/Radarr.
let pub_ = match state.admin.verify(&body.username, &body.password) {
// as Sonarr/Radarr. The bcrypt verify is ~100-200 ms of pure CPU on
// an unauthenticated endpoint, so it runs on the blocking pool.
let admin = state.admin.clone();
let verdict = tokio::task::spawn_blocking(move || admin.verify(&body.username, &body.password))
.await
.unwrap_or_else(|e| Err(openpxe_core::Error::Other(e.into())));
let pub_ = match verdict {
Ok(Some(u)) => u,
Ok(None) => {
return (
@@ -394,11 +411,18 @@ pub async fn api_update_credentials(
)
.into_response();
}
let result = state.admin.update_credentials(
&body.current_password,
body.new_username.as_deref(),
body.new_password.as_deref(),
);
// Two bcrypt operations (verify current + hash new) plus a sync disk
// persist — run the lot on the blocking pool.
let admin = state.admin.clone();
let result = tokio::task::spawn_blocking(move || {
admin.update_credentials(
&body.current_password,
body.new_username.as_deref(),
body.new_password.as_deref(),
)
})
.await
.unwrap_or_else(|e| Err(openpxe_core::Error::Other(e.into())));
match result {
Ok(pub_) => {
state.sessions.revoke_all();
+5 -1
View File
@@ -105,7 +105,11 @@ async fn send_email(cfg: &NotifyConfig, subject: &str, body: &str) -> Result<(),
.trim()
.parse()
.map_err(|e| format!("invalid To address '{}': {e}", cfg.smtp_to))?)
.subject(if subject.is_empty() { "OpenPXE" } else { subject })
.subject(if subject.is_empty() {
"OpenPXE"
} else {
subject
})
.body(body.to_string())
.map_err(|e| format!("could not build email: {e}"))?;
+10
View File
@@ -11,6 +11,10 @@ use openpxe_iso_store::{
use std::sync::Arc;
use time::OffsetDateTime;
/// Cached composited PXE boot-menu background: `(logo_rev, encoded PNG)`.
/// See `AppState::pxe_bg_cache`.
pub type PxeBgCache = Arc<parking_lot::Mutex<Option<(u64, bytes::Bytes)>>>;
#[derive(Clone)]
pub struct AppState {
pub iso_store: IsoStore,
@@ -29,6 +33,12 @@ pub struct AppState {
/// operator hasn't uploaded anything, the WebUI serves the bundled
/// rainbow-horizon mark.
pub branding: BrandingStore,
/// v0.6.2: cache of the composited PXE boot-menu background PNG,
/// keyed on the branding logo revision. Composing costs ~50-200 ms
/// of image decode/encode and **every** booting client fetches it
/// for `console --picture` — caching makes that one compose per
/// logo change instead of one per boot.
pub pxe_bg_cache: PxeBgCache,
/// Forms-auth admin record + first-run bootstrap state. When
/// `admin.is_configured() == false`, the auth middleware passes
/// every request through and `/api/me` reports `setup_required`.
+1
View File
@@ -116,6 +116,7 @@ async fn build_state() -> (AppState, tempfile::TempDir) {
hosts,
boot_log,
branding,
pxe_bg_cache: openpxe_http_api::state::PxeBgCache::default(),
admin,
sessions,
sso,
+5 -17
View File
@@ -35,23 +35,11 @@ use rust_embed::Embed;
#[include = "wimboot"]
pub struct IpxeAssets;
/// Return the embedded iPXE binary for `arch`, or `None` if we didn't bundle
/// one for that architecture.
#[must_use]
pub fn bootfile_bytes(arch: ClientArch) -> Option<Vec<u8>> {
let name = arch.ipxe_bootfile()?;
IpxeAssets::get(name).map(|f| f.data.into_owned())
}
/// Return a named asset directly (e.g. `wimboot`, or a fallback `ipxe.efi`).
#[must_use]
pub fn asset_bytes(name: &str) -> Option<Vec<u8>> {
IpxeAssets::get(name).map(|f| f.data.into_owned())
}
/// Same as [`asset_bytes`] but returns the embedded slice directly,
/// avoiding the heap copy when the caller only needs to read the
/// payload. Falls back to None for unknown names.
/// Return a named embedded asset (e.g. `snponly.efi`, `wimboot`) as a
/// `Cow` over the embedded bytes. In release builds the data is borrowed
/// straight from the binary's rodata — **zero copy** — which matters
/// because the TFTP and HTTP serving paths hit this for every boot
/// (`ipxe.efi` is ~1 MiB). Debug builds read from disk and return Owned.
#[must_use]
pub fn asset_slice(name: &str) -> Option<std::borrow::Cow<'static, [u8]>> {
IpxeAssets::get(name).map(|f| f.data)
+7 -1
View File
@@ -103,7 +103,13 @@ pub fn introspect(path: &Path) -> IntrospectionReport {
let scan_bytes = 64 * 1024 * 1024;
let mut buf = vec![0u8; 1024 * 1024];
let mut read_total = 0usize;
let mut haystack = Vec::with_capacity(scan_bytes.min(32 * 1024 * 1024));
// Size the haystack to what will actually be read — the scan cap or
// the file itself, whichever is smaller — so the fill never reallocs
// and a small ISO doesn't reserve the full 64 MiB.
let file_len = f.metadata().map_or(usize::MAX, |m| {
usize::try_from(m.len()).unwrap_or(usize::MAX)
});
let mut haystack = Vec::with_capacity(scan_bytes.min(file_len));
while read_total < scan_bytes {
let n = f.read(&mut buf).unwrap_or(0);
if n == 0 {
+17 -18
View File
@@ -63,8 +63,8 @@ use bytes::Bytes;
use nfs3_client::tokio::TokioConnector;
use nfs3_client::Nfs3ConnectionBuilder;
use nfs3_types::nfs3::{
self as nfs3, diropargs3, entry3, filename3, nfs_fh3, GETATTR3args, LOOKUP3args,
Nfs3Result, READ3args, READDIR3args,
self as nfs3, diropargs3, entry3, filename3, nfs_fh3, GETATTR3args, LOOKUP3args, Nfs3Result,
READ3args, READDIR3args,
};
use nfs3_types::rpc::{auth_unix, opaque_auth};
use nfs3_types::xdr_codec::Opaque;
@@ -220,10 +220,7 @@ impl NfsShareManager {
/// Register an NFS share. Validates, probes connectivity by
/// performing a real MOUNT3 + READDIR3, and registers the
/// resulting ISOs with the store.
pub async fn add(
&self,
req: NfsAddRequest,
) -> std::result::Result<NfsShare, NfsShareError> {
pub async fn add(&self, req: NfsAddRequest) -> std::result::Result<NfsShare, NfsShareError> {
let server = normalize_server(&req.server);
let export = req.export.trim().to_string();
if server.is_empty() {
@@ -258,7 +255,9 @@ impl NfsShareManager {
if let Err(e) = self.rescan_inner(&id).await {
let m = self.get(&id);
return Err(NfsShareError {
error: m.as_ref().and_then(|m| m.last_error.clone())
error: m
.as_ref()
.and_then(|m| m.last_error.clone())
.unwrap_or_else(|| e.to_string()),
stderr: String::new(),
hint: m.and_then(|m| m.last_hint),
@@ -333,8 +332,7 @@ impl NfsShareManager {
return Err(Error::Invalid(format!("invalid filename '{filename}'")));
}
let (tx, rx) =
tokio::sync::mpsc::channel::<std::io::Result<Bytes>>(STREAM_BUFFER_DEPTH);
let (tx, rx) = tokio::sync::mpsc::channel::<std::io::Result<Bytes>>(STREAM_BUFFER_DEPTH);
let server = share.server.clone();
let export = share.export.clone();
let port = share.port;
@@ -358,16 +356,11 @@ impl NfsShareManager {
if let Err(e) = result {
// Best-effort signal of the error to the consumer.
// If the receiver has already dropped we just exit.
let _ = tx
.send(Err(std::io::Error::other(e.to_string())))
.await;
let _ = tx.send(Err(std::io::Error::other(e.to_string()))).await;
}
});
Ok(NfsStream {
rx,
_task: task,
})
Ok(NfsStream { rx, _task: task })
}
// ── internals ─────────────────────────────────────────────────────
@@ -982,8 +975,14 @@ mod tests {
// the allow-list and the secure/insecure angle.
let h = hint_for("connect failed: MNT3ERR_ACCES").unwrap();
let lc = h.to_lowercase();
assert!(lc.contains("insecure") || lc.contains("privileged"), "got: {h}");
assert!(lc.contains("allow") || lc.contains("permission"), "got: {h}");
assert!(
lc.contains("insecure") || lc.contains("privileged"),
"got: {h}"
);
assert!(
lc.contains("allow") || lc.contains("permission"),
"got: {h}"
);
}
#[test]
+10 -23
View File
@@ -231,10 +231,7 @@ impl SmbShareManager {
/// Add or refresh a share. Validates the input, writes a creds
/// file, probes connectivity, and scans for ISOs.
pub async fn add(
&self,
req: SmbAddRequest,
) -> std::result::Result<SmbShare, SmbShareError> {
pub async fn add(&self, req: SmbAddRequest) -> std::result::Result<SmbShare, SmbShareError> {
let server = normalize_server(&req.server);
let share = req.share.trim().trim_start_matches('/').to_string();
if server.is_empty() {
@@ -309,7 +306,9 @@ impl SmbShareManager {
if let Err(e) = self.rescan_inner(&id).await {
let m = self.get(&id);
return Err(SmbShareError {
error: m.as_ref().and_then(|m| m.last_error.clone())
error: m
.as_ref()
.and_then(|m| m.last_error.clone())
.unwrap_or_else(|| e.to_string()),
stderr: String::new(),
hint: m.and_then(|m| m.last_hint),
@@ -365,11 +364,7 @@ impl SmbShareManager {
/// throttling concurrent smbclients) would need to await without
/// changing the call sites.
#[allow(clippy::unused_async)]
pub async fn stream_iso(
&self,
share_id: &str,
filename: &str,
) -> Result<SmbStream> {
pub async fn stream_iso(&self, share_id: &str, filename: &str) -> Result<SmbStream> {
let share = self
.get(share_id)
.ok_or_else(|| Error::Invalid(format!("no such SMB share '{share_id}'")))?;
@@ -377,9 +372,7 @@ impl SmbShareManager {
// share root. smbclient itself accepts only filenames at the
// share root in our `get` form, but belt-and-suspenders.
if filename.contains('/') || filename.contains('\\') || filename.contains("..") {
return Err(Error::Invalid(format!(
"invalid filename '{filename}'"
)));
return Err(Error::Invalid(format!("invalid filename '{filename}'")));
}
let creds = share
.creds_path
@@ -537,10 +530,7 @@ impl SmbShareManager {
} else {
String::new()
};
return Err((
format!("could not exec smbclient: {e}"),
stderr,
));
return Err((format!("could not exec smbclient: {e}"), stderr));
}
};
if !output.status.success() {
@@ -751,10 +741,7 @@ fn parse_ls_iso(out: &str) -> Vec<SmbListEntry> {
/// Pre-flight TCP probe to `server:port`. Format matches v0.4.64 NFS
/// probe so the UI banner reads consistently.
async fn tcp_probe(
server: &str,
port: u16,
) -> std::result::Result<(), (String, String)> {
async fn tcp_probe(server: &str, port: u16) -> std::result::Result<(), (String, String)> {
use tokio::net::TcpStream;
let addr = format!("{server}:{port}");
match tokio::time::timeout(PROBE_TIMEOUT, TcpStream::connect(&addr)).await {
@@ -931,8 +918,8 @@ mod tests {
// exec error in `error` plus an empty `stderr`. The
// SmbShareError constructor's hint_for fallback checks error
// too, so this pattern needs to translate as well.
let h2 = hint_for("could not exec smbclient: No such file or directory (os error 2)")
.unwrap();
let h2 =
hint_for("could not exec smbclient: No such file or directory (os error 2)").unwrap();
assert!(h2.contains("smbclient"));
}
+10 -1
View File
@@ -342,9 +342,18 @@ impl IsoStore {
/// SMB sources or when the file is missing.
pub fn iso_path_for(&self, id: &str) -> Option<PathBuf> {
let meta = self.get(id)?;
self.local_path(&meta)
}
/// Same resolution as [`Self::iso_path_for`], but for a meta the
/// caller already holds — skips the second registry lock + deep
/// clone, which matters on the per-range-request ISO serving path
/// (a sanboot install issues hundreds of those).
#[must_use]
pub fn local_path(&self, meta: &IsoMeta) -> Option<PathBuf> {
match &meta.source {
IsoSource::Local => {
let path = self.iso_path(id);
let path = self.iso_path(&meta.id);
if path.exists() {
Some(path)
} else {
+1
View File
@@ -192,6 +192,7 @@ async fn main() -> anyhow::Result<()> {
hosts: hosts.clone(),
boot_log: boot_log.clone(),
branding: branding.clone(),
pxe_bg_cache: openpxe_http_api::state::PxeBgCache::default(),
admin: admin.clone(),
sessions: sessions.clone(),
sso: sso.clone(),
+30 -16
View File
@@ -7,12 +7,12 @@
//! `tftpd`/`in.tftpd` works and is why TFTP is awkward behind stateful NAT:
//! the ephemeral ports must be reachable from the client.
//!
//! We only serve files from `openpxe_ipxe_assets::asset_bytes` — that is,
//! We only serve files from `openpxe_ipxe_assets::asset_slice` — that is,
//! the bundled iPXE binaries and wimboot. No filesystem is ever opened, so
//! `../` path traversal attempts simply return ENOENT.
use openpxe_core::{ClientEvent, ClientRegistry};
use openpxe_ipxe_assets::asset_bytes;
use openpxe_ipxe_assets::asset_slice;
use socket2::{Domain, Protocol, Socket, Type};
use std::net::{IpAddr, SocketAddr};
use std::sync::Arc;
@@ -21,6 +21,7 @@ use tokio::net::UdpSocket;
// TFTP opcodes.
const OP_RRQ: u16 = 1;
const OP_WRQ: u16 = 2;
const OP_DATA: u16 = 3;
const OP_ACK: u16 = 4;
const OP_ERROR: u16 = 5;
@@ -89,16 +90,34 @@ async fn handle_rrq(
metrics: openpxe_core::Metrics,
) -> anyhow::Result<()> {
let Some(req) = parse_rrq(&packet) else {
// Not a well-formed RRQ. A WRQ deserves an explicit refusal —
// legacy clients retry a silently-dropped write until they time
// out; an ERROR packet fails them fast with a readable reason.
if packet.len() >= 2 && u16::from_be_bytes([packet[0], packet[1]]) == OP_WRQ {
let sock = bind_udp(bind_ip, 0)?;
let _ = send_error(&sock, peer, ERR_ILLEGAL_OP, "writes not supported").await;
}
return Ok(());
};
let Request {
filename, options, ..
filename,
mode,
options,
} = req;
// Per-transfer ephemeral socket.
let sock = bind_udp(bind_ip, 0)?;
let Some(file_bytes) = asset_bytes(&filename) else {
// We serve binary boot artifacts; netascii line-ending translation
// would corrupt them. Refuse loudly instead of timing out silently —
// matters for legacy clients that default to netascii.
if !mode.eq_ignore_ascii_case("octet") {
let _ = send_error(&sock, peer, ERR_NOT_DEFINED, "only octet mode is supported").await;
tracing::info!(target: "openpxe::tftp", peer=%peer, %mode, "rejected non-octet transfer");
return Ok(());
}
let Some(file_bytes) = asset_slice(&filename) else {
let _ = send_error(&sock, peer, ERR_FILE_NOT_FOUND, "no such file").await;
tracing::info!(target: "openpxe::tftp", peer=%peer, file=%filename, "404");
clients.record(
@@ -262,7 +281,6 @@ async fn handle_rrq(
#[derive(Debug)]
struct Request {
filename: String,
#[allow(dead_code)]
mode: String,
options: Vec<(String, String)>,
}
@@ -393,17 +411,13 @@ fn bind_udp(bind: IpAddr, port: u16) -> anyhow::Result<UdpSocket> {
Ok(UdpSocket::from_std(std_sock)?)
}
#[allow(dead_code)]
const _UNUSED: (u16, u16) = (ERR_NOT_DEFINED, ERR_ILLEGAL_OP);
/// Pure-logic helper used by the unit tests below and (in a refactor) by
/// `handle_rrq`. Given a position in the file and the window, return the
/// (block_no, chunk_len) list this window will emit. Useful as a sanity
/// check that our windowing math matches the wire behavior the spec
/// requires — tested against edge cases (exact-blksize tail, short tail,
/// single-block window).
#[must_use]
pub fn plan_window(
/// Pure-logic mirror of `handle_rrq`'s windowing math, exercised by the
/// unit tests below. Given a position in the file and the window, return
/// the (block_no, chunk_len) list this window will emit — tested against
/// edge cases (exact-blksize tail, short tail, single-block window,
/// block-number wraparound).
#[cfg(test)]
fn plan_window(
total: usize,
offset: usize,
blksize: usize,