Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
5da05a519d |
Generated
+387
-278
File diff suppressed because it is too large
Load Diff
+11
-7
@@ -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"] }
|
||||
|
||||
|
||||
@@ -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());
|
||||
}
|
||||
|
||||
@@ -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()
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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(),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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();
|
||||
|
||||
@@ -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;
|
||||
}
|
||||
|
||||
@@ -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
@@ -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;
|
||||
}
|
||||
|
||||
@@ -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();
|
||||
|
||||
@@ -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}"))?;
|
||||
|
||||
|
||||
@@ -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`.
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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]
|
||||
|
||||
@@ -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"));
|
||||
}
|
||||
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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
@@ -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,
|
||||
|
||||
Reference in New Issue
Block a user