fix: harden node upgrades and prepare 1.9.0-alpha

This commit is contained in:
archipelago
2026-10-05 12:43:49 -04:00
parent 138a541d01
commit daac47cac4
129 changed files with 9910 additions and 794 deletions
+395 -69
View File
@@ -7,7 +7,8 @@
//! to a full chip erase before write.
//!
//! MeshCore and Meshtastic are flashed the same way: download a released
//! image, `esptool erase_flash`, then `esptool write_flash 0x0 <image>`.
//! image, verify it, then run `write_flash --erase-all 0x0 <image>` with
//! the packaged esptool executable.
//! Reticulum/RNode is different: `archy-rnodeconf --autoinstall` owns the
//! whole fetch+erase+flash+EEPROM-bootstrap sequence itself (confirmed live
//! via `archy-rnodeconf --help` — there is no raw esptool path exposed for
@@ -49,32 +50,18 @@ impl FlashBoard {
}
}
/// Map a detected USB vid:pid to a known flashable board, using the same
/// table as `image-recipe/configs/99-mesh-radio.rules`. CP2102 (10c4:ea60)
/// is confirmed there as Heltec V3's USB-UART bridge chip, and is safe to
/// auto-match since that vid:pid is bridge-chip-specific.
///
/// Heltec V4 is NOT auto-matchable and deliberately has no entry here: it
/// was confirmed live (real hardware, 2026-07-23) to use the ESP32-S3's
/// built-in native-USB JTAG/serial peripheral, reporting vid:pid 303a:1001
/// with product string "USB JTAG/serial debug unit" — that descriptor is
/// baked into the chip's ROM and is IDENTICAL across every ESP32-S3 board
/// with native USB enabled, not just Heltec V4. Adding `303a:1001 =>
/// HeltecV4` here would silently misidentify any other native-USB ESP32-S3
/// board (a T3-S3, a bare devkit, etc.) as a V4 and risk writing the wrong
/// board's image. Callers (the RPC layer / frontend) must let the user pick
/// the board manually whenever this returns `None`.
pub fn resolve_flash_board(info: &DetectedDeviceInfo) -> Option<FlashBoard> {
match (info.vid.as_deref(), info.pid.as_deref()) {
(Some("10c4"), Some("ea60")) => Some(FlashBoard::HeltecV3),
_ => None,
}
/// Generic CP2102 and native ESP32-S3 USB IDs identify adapters/chips, not
/// board wiring. Require an explicit board until a board-specific identity is
/// available; guessing a Heltec model can write incompatible firmware.
pub fn resolve_flash_board(_info: &DetectedDeviceInfo) -> Option<FlashBoard> {
None
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
#[serde(rename_all = "lowercase")]
pub enum FlashStage {
Downloading,
Preparing,
Erasing,
Writing,
Autoinstalling,
@@ -101,11 +88,10 @@ const LOG_TAIL_MAX: usize = 200;
/// start opening the port (which itself toggles DTR/RTS) again.
const POST_FLASH_SETTLE_DELAY: std::time::Duration = std::time::Duration::from_secs(5);
/// Absolute ceiling on a whole flash job (download + erase + write, or
/// autoinstall), regardless of what it's doing internally. Last-resort
/// safety net so a hang anywhere can't wedge the single-flash-job guard
/// forever — generous enough to never trigger on a legitimately slow
/// multi-hundred-MB transfer.
/// Deadline for preparation and warning threshold for the active flasher.
/// A download can be cancelled safely. An active write retains ownership until
/// its subprocess exits, even after this threshold: reporting an aborted job
/// while a detached writer continues would let a retry corrupt the device.
const MAX_JOB_DURATION: std::time::Duration = std::time::Duration::from_secs(15 * 60);
/// How long to wait for MeshService::stop() to release the serial port
@@ -247,6 +233,16 @@ fn firmware_cache_dir(data_dir: &Path) -> PathBuf {
data_dir.join("mesh").join("firmware-cache")
}
async fn invalidate_radio_settings_marker(data_dir: &Path) -> Result<()> {
// A full-chip flash erased the device's preferences. A previous host-side
// marker is no longer evidence that this radio has the requested settings.
match tokio::fs::remove_file(data_dir.join("meshcore-radio-params.json")).await {
Ok(()) => Ok(()),
Err(error) if error.kind() == std::io::ErrorKind::NotFound => Ok(()),
Err(error) => Err(error).context("Firmware written, but old radio-settings marker could not be cleared; reconnect is paused"),
}
}
/// No blanket `.timeout()` here on purpose: reqwest's request timeout covers
/// the *entire* request including streaming the response body, which would
/// kill a legitimate large download partway through (Meshtastic's esp32s3
@@ -314,6 +310,10 @@ pub async fn list_firmware(family: DeviceType) -> Result<Vec<String>> {
struct GithubAsset {
name: String,
browser_download_url: String,
#[serde(default)]
size: Option<u64>,
#[serde(default)]
digest: Option<String>,
}
#[derive(serde::Deserialize)]
@@ -322,8 +322,24 @@ struct GithubRelease {
assets: Vec<GithubAsset>,
}
async fn register_flash_job(handle: &FlashJobHandle, job: &Arc<FlashJob>) -> Result<()> {
// Keep checking and registering under one lock: concurrent RPCs must
// never start two writers on the same device.
let mut existing = handle.try_write().map_err(|_| anyhow::anyhow!(
"The radio is being inspected or reconfigured; wait for that operation to finish before flashing"
))?;
if let Some(current) = existing.as_ref() {
if !current.snapshot().await.done {
anyhow::bail!("A firmware flash is already in progress on this node");
}
}
*existing = Some(Arc::clone(job));
Ok(())
}
/// Start a flash job in the background. Returns as soon as the job has been
/// registered and the listener released — callers poll `FlashJobHandle` via
/// registered — preparation and listener shutdown run in the background.
/// Callers poll `FlashJobHandle` via
/// `mesh.flash-status` for progress. Only one job may be in flight at a time.
pub async fn start_flash_job(
handle: &FlashJobHandle,
@@ -333,21 +349,44 @@ pub async fn start_flash_job(
board: FlashBoard,
family: DeviceType,
) -> Result<()> {
{
let existing = handle.read().await;
if let Some(job) = existing.as_ref() {
if !job.snapshot().await.done {
anyhow::bail!("A firmware flash is already in progress on this node");
}
}
// Missing/corrupt tooling must fail before stopping a working radio or
// registering a job that can never reach the serial device.
if matches!(family, DeviceType::Meshtastic | DeviceType::Meshcore) {
preflight_esptool().await?;
}
let job = FlashJob::new(board, family, path.clone());
*handle.write().await = Some(Arc::clone(&job));
register_flash_job(handle, &job).await?;
let bg_job = Arc::clone(&job);
let bg_service = Arc::clone(mesh_service);
let task = tokio::spawn(async move {
// Downloads and checksum checks do not need exclusive serial access.
// Keep a working listener alive if upstream/download verification fails.
let prepared_image = if matches!(family, DeviceType::Meshcore | DeviceType::Meshtastic) {
match tokio::time::timeout(
MAX_JOB_DURATION,
fetch_esptool_image(board, family, &data_dir, &bg_job),
)
.await
{
Ok(Ok(image)) => Some(image),
result => {
let error = match result {
Ok(Err(error)) => error,
_ => anyhow::anyhow!(
"Firmware download exceeded the time limit; radio was not changed"
),
};
bg_job.fail(&error).await;
return;
}
}
} else {
None
};
// Cancellation is safe during download. Once listener shutdown starts,
// allow that bounded operation to finish instead of orphaning its task.
bg_job.set_stage(FlashStage::Preparing).await;
// esptool/archy-rnodeconf need exclusive serial access — release
// the listener's hold on the port before touching it. This USED
// TO run synchronously in start_flash_job before the job was even
@@ -404,17 +443,31 @@ pub async fn start_flash_job(
// subsequent mesh.flash-device call failed with "already in
// progress" until the service was restarted). Generous enough that
// a legitimately slow multi-hundred-MB transfer still completes.
let result = match tokio::time::timeout(
MAX_JOB_DURATION,
run_flash(board, family, &data_dir, &path, &bg_job),
)
.await
{
let flash = run_flash(
board,
family,
&data_dir,
&path,
prepared_image.as_deref(),
&bg_job,
);
tokio::pin!(flash);
let result = match tokio::time::timeout(MAX_JOB_DURATION, &mut flash).await {
Ok(inner) => inner,
Err(_) => Err(anyhow::anyhow!(
"Flash job exceeded the {}-minute ceiling — aborted",
MAX_JOB_DURATION.as_secs() / 60
)),
Err(_) if bg_job.snapshot().await.stage == FlashStage::Downloading => Err(
anyhow::anyhow!("Firmware download exceeded the time limit; radio was not written"),
),
Err(_) => {
// Dropping this future does NOT stop its flash subprocess.
// Keep the job busy until it exits, rather than allowing a
// second writer while the first one still owns the device.
bg_job.push_log("Flashing is taking longer than expected; waiting for the active tool to exit before allowing another operation").await;
flash.await
}
};
let result = match result {
Ok(()) => invalidate_radio_settings_marker(&data_dir).await,
error => error,
};
let succeeded = result.is_ok();
@@ -423,7 +476,6 @@ pub async fn start_flash_job(
bg_job
.push_log("Flash completed successfully".to_string())
.await;
bg_job.finish().await;
info!(path = %path, board = ?board, family = %family, "LoRa firmware flash succeeded");
}
Err(e) => {
@@ -434,7 +486,6 @@ pub async fn start_flash_job(
// erase_flash failed", no actual esptool stderr).
warn!(path = %path, error = %format!("{e:#}"), "LoRa firmware flash failed");
bg_job.push_log(format!("ERROR: {e:#}")).await;
bg_job.fail(e).await;
}
}
@@ -450,6 +501,9 @@ pub async fn start_flash_job(
}
if !succeeded {
if let Err(error) = &result {
bg_job.fail(error).await;
}
// Deliberately do NOT auto-restart the listener here. A failed
// flash means we can't vouch for the board's state — reopening
// the port immediately (esptool/rnodeconf's own reset sequence
@@ -499,6 +553,7 @@ pub async fn start_flash_job(
Err(e) => warn!(error = %e, "Failed to load mesh config after flash"),
}
}
bg_job.finish().await;
});
*job.abort_handle.write().await = Some(task.abort_handle());
@@ -510,12 +565,13 @@ async fn run_flash(
family: DeviceType,
data_dir: &Path,
path: &str,
prepared_image: Option<&Path>,
job: &Arc<FlashJob>,
) -> Result<()> {
match family {
DeviceType::Meshtastic | DeviceType::Meshcore => {
let image = fetch_esptool_image(board, family, data_dir, job).await?;
esptool_erase_and_write(path, &image, job).await
let image = prepared_image.context("Firmware was not prepared before serial access")?;
esptool_erase_and_write(path, image, job).await
}
DeviceType::Reticulum => {
let lora_region = super::load_config(data_dir)
@@ -664,15 +720,65 @@ async fn fetch_meshcore_image(
anyhow::anyhow!("No matching MeshCore image in release {}", release.tag_name)
})?;
anyhow::ensure!(
Path::new(&asset.name)
.file_name()
.and_then(|name| name.to_str())
== Some(asset.name.as_str()),
"Invalid firmware asset filename"
);
let out_path = cache.join(&asset.name);
if tokio::fs::metadata(&out_path).await.is_ok() {
job.push_log(format!("Using cached {}", asset.name)).await;
return Ok(out_path);
if verify_meshcore_asset(&out_path, asset).await.is_ok() {
job.push_log(format!("Using verified cached {}", asset.name))
.await;
return Ok(out_path);
}
tokio::fs::remove_file(&out_path)
.await
.context("Removing invalid cached firmware")?;
job.push_log("Cached firmware failed verification; downloading a fresh copy")
.await;
}
download_to_file(client, &asset.browser_download_url, &out_path, job).await?;
if let Err(error) = verify_meshcore_asset(&out_path, asset).await {
// A partial/unverified download must not become next attempt's cache.
let _ = tokio::fs::remove_file(&out_path).await;
return Err(error);
}
Ok(out_path)
}
async fn verify_meshcore_asset(path: &Path, asset: &GithubAsset) -> Result<()> {
use sha2::{Digest, Sha256};
let size = asset
.size
.context("MeshCore release did not provide firmware size")?;
anyhow::ensure!(
(1..=16 * 1024 * 1024).contains(&size),
"Invalid MeshCore firmware size"
);
anyhow::ensure!(
tokio::fs::metadata(path).await?.len() == size,
"Firmware size mismatch; radio was not written"
);
let expected = asset
.digest
.as_deref()
.and_then(|value| value.strip_prefix("sha256:"))
.context("MeshCore release did not provide a SHA-256 checksum")?;
anyhow::ensure!(
expected.len() == 64 && expected.bytes().all(|byte| byte.is_ascii_hexdigit()),
"Invalid MeshCore release checksum"
);
let actual = hex::encode(Sha256::digest(tokio::fs::read(path).await?));
anyhow::ensure!(
actual.eq_ignore_ascii_case(expected),
"Firmware checksum mismatch; radio was not written"
);
Ok(())
}
async fn download_to_file(
client: &reqwest::Client,
url: &str,
@@ -722,7 +828,8 @@ async fn download_to_file(
}
}
}
file.flush().await.ok();
file.flush().await.context("Flushing firmware download")?;
file.sync_all().await.context("Saving firmware download")?;
tokio::fs::rename(&tmp, dest)
.await
.context("Finalizing firmware download")?;
@@ -773,19 +880,10 @@ async fn esptool_erase_and_write(path: &str, image: &Path, job: &Arc<FlashJob>)
/// Building global args separately from subcommand args keeps this correct
/// by construction instead of relying on call-site ordering.
///
/// Normal stub-loader mode (no --no-stub) needs the esp32s3 stub flasher
/// blob at /usr/lib/python3/dist-packages/esptool/targets/stub_flasher/
/// stub_flasher_32s3.json — Debian's `esptool` package (4.7.0+dfsg-0.1)
/// ships without it (stripped for DFSG compliance: the prebuilt blob has no
/// buildable-from-source path Debian could verify), so scripts/self-update.sh
/// fetches the exact same file from the matching upstream esptool release
/// tag and installs it alongside the apt package (see the esptool install
/// step there). --no-stub (talk directly to the ROM bootloader, skip the
/// stub) was tried first and works for connecting, but the ROM bootloader
/// doesn't implement a full-chip-erase opcode at all — only the stub does —
/// so --no-stub broke our "always erase before write" default outright
/// rather than just being slower. Restoring the real stub file is the
/// correct fix, not routing around its absence.
/// Normal stub-loader mode is required for full-chip erase. The packaged
/// archy-esptool includes and self-tests Espressif's ESP32-S3 stub. Debian's
/// stripped esptool package alone did not provide that resource, so changing
/// flags to --no-stub cannot repair this operation.
fn esptool_global_args<'a>(path: &'a str, baud: Option<&'a str>) -> Vec<&'a str> {
let mut args = vec!["--chip", ESPTOOL_CHIP, "--port", path];
if let Some(b) = baud {
@@ -795,24 +893,96 @@ fn esptool_global_args<'a>(path: &'a str, baud: Option<&'a str>) -> Vec<&'a str>
args
}
fn esptool_executable() -> Result<PathBuf> {
let mut candidates = vec![PathBuf::from("/usr/local/bin/archy-esptool")];
if let Some(paths) = std::env::var_os("PATH") {
for directory in std::env::split_paths(&paths) {
for name in ["esptool", "esptool.py"] {
candidates.push(directory.join(name));
}
}
}
executable_from_candidates(candidates)
}
fn executable_from_candidates(candidates: impl IntoIterator<Item = PathBuf>) -> Result<PathBuf> {
use std::os::unix::fs::PermissionsExt;
candidates.into_iter().find(|path| {
path.is_file() && path.metadata().is_ok_and(|metadata| metadata.permissions().mode() & 0o111 != 0)
}).ok_or_else(|| anyhow::anyhow!(
"Radio flashing tools are missing. Install the complete Archipelago update and retry; the radio has not been changed."
))
}
async fn preflight_esptool() -> Result<PathBuf> {
let executable = esptool_executable()?;
check_esptool(&executable).await?;
Ok(executable)
}
async fn check_esptool(executable: &Path) -> Result<()> {
let mut command = Command::new(executable);
command
.arg(
if executable
.file_name()
.is_some_and(|name| name == "archy-esptool")
{
"--archy-self-test"
} else {
"version"
},
)
.kill_on_drop(true);
let output = tokio::time::timeout(std::time::Duration::from_secs(20), command.output())
.await
.context("Radio flashing tool did not respond; the radio has not been changed")?
.context("Radio flashing tool could not start; check its installation and permissions")?;
anyhow::ensure!(
output.status.success(),
"Radio flashing tool self-check failed; reinstall the complete update before retrying"
);
Ok(())
}
fn retryable_flash_error(error: &anyhow::Error) -> bool {
if error
.chain()
.any(|cause| cause.downcast_ref::<std::io::Error>().is_some())
{
return false;
}
let detail = format!("{error:#}").to_lowercase();
[
"failed to connect",
"timed out waiting",
"invalid head of packet",
"serial data stream stopped",
]
.iter()
.any(|message| detail.contains(message))
}
async fn esptool_with_retry(path: &str, subcommand: &[&str], job: &Arc<FlashJob>) -> Result<()> {
let mut cmd = Command::new("esptool");
let executable = preflight_esptool().await?;
let mut cmd = Command::new(&executable);
cmd.args(esptool_global_args(path, None));
cmd.args(subcommand);
match run_streamed(cmd, None, job).await {
Ok(()) => Ok(()),
Err(first_err) => {
Err(first_err) if retryable_flash_error(&first_err) => {
job.push_log(format!(
"First attempt failed ({first_err:#}); retrying once at {ESPTOOL_FALLBACK_BAUD} baud"
))
.await;
let mut retry = Command::new("esptool");
let mut retry = Command::new(&executable);
retry.args(esptool_global_args(path, Some(ESPTOOL_FALLBACK_BAUD)));
retry.args(subcommand);
run_streamed(retry, None, job)
.await
.context(format!("retry also failed (first attempt: {first_err:#})"))
}
Err(error) => Err(error),
}
}
@@ -990,3 +1160,159 @@ async fn run_streamed(mut cmd: Command, stdin: Option<Vec<u8>>, job: &Arc<FlashJ
}
Ok(())
}
#[cfg(test)]
mod flashing_regression_tests {
use super::*;
#[tokio::test]
async fn full_flash_invalidates_only_radio_settings_marker() {
let dir = tempfile::tempdir().unwrap();
let marker = dir.path().join("meshcore-radio-params.json");
tokio::fs::write(&marker, b"old settings").await.unwrap();
let other = dir.path().join("mesh-config.json");
tokio::fs::write(&other, b"preserved").await.unwrap();
invalidate_radio_settings_marker(dir.path()).await.unwrap();
assert!(!marker.exists());
assert_eq!(tokio::fs::read(&other).await.unwrap(), b"preserved");
invalidate_radio_settings_marker(dir.path()).await.unwrap();
tokio::fs::create_dir(&marker).await.unwrap();
assert!(invalidate_radio_settings_marker(dir.path()).await.is_err());
}
use std::os::unix::fs::PermissionsExt;
#[tokio::test]
async fn meshcore_cache_requires_release_size_and_checksum() {
use sha2::{Digest, Sha256};
let dir = tempfile::tempdir().unwrap();
let file = dir.path().join("fixture.bin");
tokio::fs::write(&file, b"firmware fixture").await.unwrap();
let mut asset = GithubAsset {
name: "fixture.bin".into(),
browser_download_url: "https://example.invalid/fixture.bin".into(),
size: Some(16),
digest: Some(format!(
"sha256:{}",
hex::encode(Sha256::digest(b"firmware fixture"))
)),
};
assert!(verify_meshcore_asset(&file, &asset).await.is_ok());
tokio::fs::write(&file, b"tampered fixture").await.unwrap();
assert!(verify_meshcore_asset(&file, &asset).await.is_err());
tokio::fs::write(&file, b"short").await.unwrap();
assert!(verify_meshcore_asset(&file, &asset).await.is_err());
tokio::fs::write(&file, b"firmware fixture").await.unwrap();
asset.digest = None;
assert!(verify_meshcore_asset(&file, &asset).await.is_err());
}
#[test]
fn generic_usb_ids_do_not_select_firmware() {
for (vid, pid) in [("10c4", "ea60"), ("303a", "1001")] {
let info = DetectedDeviceInfo {
path: "/dev/fixture".into(),
vid: Some(vid.into()),
pid: Some(pid.into()),
product: None,
manufacturer: None,
plugged_at: None,
};
assert_eq!(resolve_flash_board(&info), None);
}
}
#[test]
fn absent_nonexecutable_and_directory_candidates_are_rejected() {
let dir = tempfile::tempdir().unwrap();
let file = dir.path().join("flasher");
std::fs::write(&file, "#!/bin/sh\nexit 0\n").unwrap();
std::fs::set_permissions(&file, std::fs::Permissions::from_mode(0o600)).unwrap();
assert!(executable_from_candidates([
dir.path().join("absent"),
file.clone(),
dir.path().to_path_buf()
])
.is_err());
std::fs::set_permissions(&file, std::fs::Permissions::from_mode(0o700)).unwrap();
assert_eq!(executable_from_candidates([file.clone()]).unwrap(), file);
}
#[tokio::test]
async fn flasher_preflight_reports_missing_interpreter_and_failed_self_check() {
let dir = tempfile::tempdir().unwrap();
let file = dir.path().join("archy-esptool");
for script in ["#!/missing/python\n", "#!/bin/sh\nexit 7\n"] {
std::fs::write(&file, script).unwrap();
std::fs::set_permissions(&file, std::fs::Permissions::from_mode(0o700)).unwrap();
assert!(check_esptool(&file).await.is_err());
}
std::fs::write(&file, "#!/bin/sh\n[ \"$1\" = --archy-self-test ]\n").unwrap();
assert!(check_esptool(&file).await.is_ok());
}
#[tokio::test]
async fn flash_registration_excludes_other_writers_and_active_probes() {
let handle = new_job_handle();
let first = FlashJob::new(
FlashBoard::HeltecV3,
DeviceType::Meshcore,
"/dev/fixture".into(),
);
let second = FlashJob::new(
FlashBoard::HeltecV3,
DeviceType::Meshcore,
"/dev/fixture".into(),
);
let probe = handle.read().await;
assert!(register_flash_job(&handle, &first).await.is_err());
drop(probe);
let (a, b) = tokio::join!(
register_flash_job(&handle, &first),
register_flash_job(&handle, &second)
);
assert_eq!(usize::from(a.is_ok()) + usize::from(b.is_ok()), 1);
handle.read().await.as_ref().unwrap().finish().await;
let next = FlashJob::new(
FlashBoard::HeltecV3,
DeviceType::Meshcore,
"/dev/fixture".into(),
);
assert!(register_flash_job(&handle, &next).await.is_ok());
}
#[test]
fn retries_only_known_serial_transport_failures() {
for kind in [
std::io::ErrorKind::NotFound,
std::io::ErrorKind::PermissionDenied,
] {
assert!(!retryable_flash_error(
&anyhow::Error::from(std::io::Error::from(kind))
.context("Failed to start subprocess")
));
}
for message in [
"Wrong chip",
"Invalid image",
"module missing",
"permission denied",
"port is busy",
] {
assert!(!retryable_flash_error(&anyhow::anyhow!(message)));
}
assert!(retryable_flash_error(&anyhow::anyhow!(
"Failed to connect to ESP32-S3: timed out waiting for packet header"
)));
assert_eq!(
esptool_global_args("/dev/fixture", Some("115200")),
[
"--chip",
"esp32s3",
"--port",
"/dev/fixture",
"--baud",
"115200"
]
);
}
}
+54 -5
View File
@@ -1092,6 +1092,14 @@ impl MeshService {
let (new_tx, new_rx) = tokio::sync::mpsc::channel(32);
*self.state.cmd_tx.write().await = new_tx;
self.cmd_rx = Some(new_rx);
{
let mut status = self.state.status.write().await;
status.device_connected = false;
status.device_path = None;
status.firmware_version = None;
status.self_node_id = None;
status.peer_count = 0;
}
info!("Mesh service stopped");
}
@@ -2321,7 +2329,10 @@ impl MeshService {
}
}
let was_enabled = self.config.enabled;
let listener_running = self
.listener_handle
.as_ref()
.is_some_and(|handle| !handle.is_finished());
let needs_session_restart = session_config_changed(&self.config, &config);
self.config = config.clone();
@@ -2335,10 +2346,13 @@ impl MeshService {
.unwrap_or_else(|| "archipelago".to_string());
}
// If enabled state changed, start/stop the listener
if config.enabled && !was_enabled {
// Reconcile desired state with the actual task, not the old enabled
// flag. A failed flash can stop the task while keeping enabled=true;
// explicitly reconnecting with the same settings must restart it.
if config.enabled && !listener_running {
self.stop().await;
self.start()?;
} else if !config.enabled && was_enabled {
} else if !config.enabled {
self.stop().await;
// Clear connected state
let mut status = self.state.status.write().await;
@@ -2347,7 +2361,7 @@ impl MeshService {
status.firmware_version = None;
status.self_node_id = None;
status.peer_count = 0;
} else if config.enabled && was_enabled && needs_session_restart {
} else if needs_session_restart {
info!("Mesh session config changed — restarting listener to apply");
self.stop().await;
self.start()?;
@@ -2537,6 +2551,41 @@ mod tests {
}
use super::*;
#[tokio::test]
async fn reconnect_with_unchanged_enabled_config_restarts_stopped_listener() {
let dir = tempfile::tempdir().unwrap();
let key = SigningKey::from_bytes(&[7; 32]);
let public = hex::encode(key.verifying_key().to_bytes());
let did = crate::identity::did_key_from_pubkey_hex(&public).unwrap();
let config = MeshConfig {
enabled: true,
..Default::default()
};
save_config(dir.path(), &config).await.unwrap();
let mut service = MeshService::new(dir.path(), &key, &did, &public)
.await
.unwrap();
assert!(service.listener_handle.is_none());
service.configure(config.clone()).await.unwrap();
assert!(service.listener_handle.is_some());
service.stop().await;
assert!(service.config.enabled);
service.configure(config.clone()).await.unwrap();
assert!(service.listener_handle.is_some());
service.stop().await;
service.listener_handle = Some(tokio::spawn(async {}));
tokio::task::yield_now().await;
service.configure(config).await.unwrap();
assert!(!service.listener_handle.as_ref().unwrap().is_finished());
service.stop().await;
let disabled = MeshConfig {
enabled: false,
..Default::default()
};
service.configure(disabled).await.unwrap();
assert!(service.listener_handle.is_none());
}
#[test]
fn session_config_change_detection() {
let base = MeshConfig::default();
+42
View File
@@ -378,6 +378,19 @@ fn decode_mesh_name(bytes: &[u8], fallback: &str) -> String {
/// Parse RESP_DEVICE_INFO (0x0D) response.
/// Returns firmware version string and device capabilities.
pub fn parse_device_info(data: &[u8]) -> Result<(String, u16)> {
// Official companion v3+ binary layout: protocol, half-capacity,
// channels, PIN (4), build date (12), model (40), version (20).
// Verified against companion-v1.17.1 MyMesh.cpp and Heltec V3 readback.
if data
.first()
.is_some_and(|version| (3..=31).contains(version))
{
anyhow::ensure!(data.len() >= 79, "Truncated MeshCore device info");
return Ok((
decode_mesh_name(&data[59..79], "unknown"),
u16::from(data[1]) * 2,
));
}
// Device info format varies by firmware version.
// Minimum: firmware version string (null-terminated) + max_contacts (u16 LE)
if data.is_empty() {
@@ -404,6 +417,15 @@ pub fn parse_self_info(data: &[u8]) -> Result<(u32, String)> {
anyhow::bail!("Self info response too short: {} bytes", data.len());
}
// Current companions send type/power/max-power, public key (32),
// position (8), four preference bytes, RF parameters (10), then name.
if data.len() >= 57 {
let node_id = u32::from_le_bytes(data[3..7].try_into().unwrap());
return Ok((
node_id,
decode_mesh_name(&data[57..], &format!("node-{node_id:08x}")),
));
}
let node_id = u32::from_le_bytes([data[0], data[1], data[2], data[3]]);
// Name follows after fixed fields. A firmware whose fixed-field layout
@@ -966,4 +988,24 @@ mod tests {
fn test_parse_self_info_too_short() {
assert!(parse_self_info(&[0x01, 0x02]).is_err());
}
#[test]
fn current_companion_binary_metadata_is_not_a_name_or_version() {
let mut info = vec![0u8; 81];
info[0] = 13;
info[1] = 150;
info[19..28].copy_from_slice(b"Heltec V3");
info[59..74].copy_from_slice(b"v1.17.1-d929643");
assert_eq!(
parse_device_info(&info).unwrap(),
("v1.17.1-d929643".into(), 300)
);
assert!(parse_device_info(&info[..78]).is_err());
let mut own = vec![0u8; 57];
own[0..3].copy_from_slice(&[1, 22, 22]);
own[3..7].copy_from_slice(&42u32.to_le_bytes());
own[47..51].copy_from_slice(&869618u32.to_le_bytes());
own.extend_from_slice(b"My radio");
assert_eq!(parse_self_info(&own).unwrap(), (42, "My radio".into()));
}
}
+48 -11
View File
@@ -277,6 +277,19 @@ fn terminate_group(child: &Child) {
}
}
/// Own the daemon during its asynchronous handshake too. Cancellation drops
/// the future without executing an error branch; a bare Child would survive
/// and keep the serial port open even though the listener had stopped.
struct StartingDaemon(Option<Child>);
impl Drop for StartingDaemon {
fn drop(&mut self) {
if let Some(child) = self.0.as_ref() {
terminate_group(child);
}
}
}
/// One peer learned via an RNS announce (LXMF delivery destination).
#[derive(Clone)]
struct ReticulumPeer {
@@ -517,10 +530,10 @@ impl ReticulumLink {
let child = cmd
.spawn()
.context("Failed to spawn reticulum-daemon — is it installed/packaged?")?;
let mut starting = StartingDaemon(Some(child));
// Wait for the socket to appear, then for the daemon's "ready" event.
// Runs as a block so every failure path tears the just-spawned daemon
// group down via `terminate_group` (the child has no `kill_on_drop`).
// StartingDaemon tears the group down on errors AND cancellation.
let init = async {
let deadline = tokio::time::Instant::now() + Duration::from_secs(15);
let stream = loop {
@@ -560,20 +573,17 @@ impl ReticulumLink {
dest_hash = %dest_hash_hex,
"Reticulum daemon ready"
);
Ok((write_half, reader, dest_hash, display_name))
};
let (write_half, reader, dest_hash, display_name) = match init.await {
Ok(parts) => parts,
Err(e) => {
terminate_group(&child);
return Err(e);
}
Ok::<_, anyhow::Error>((write_half, reader, dest_hash, display_name))
};
let (write_half, reader, dest_hash, display_name) = init.await?;
let mut link = Self {
device_path: label,
socket_path,
child,
child: starting
.0
.take()
.expect("starting daemon is owned until ready"),
writer: write_half,
reader,
dest_hash,
@@ -1557,6 +1567,33 @@ impl Drop for ReticulumLink {
mod tests {
use super::*;
#[tokio::test]
async fn cancelled_daemon_handshake_terminates_child() {
let (started, ready) = tokio::sync::oneshot::channel();
let task = tokio::spawn(async move {
let child = Command::new("sleep")
.arg("60")
.process_group(0)
.spawn()
.unwrap();
let pid = child.id().unwrap();
let _starting = StartingDaemon(Some(child));
started.send(pid).unwrap();
std::future::pending::<()>().await;
});
let pid = ready.await.unwrap();
assert_eq!(unsafe { libc::kill(pid as i32, 0) }, 0);
task.abort();
assert!(task.await.unwrap_err().is_cancelled());
tokio::time::timeout(Duration::from_secs(3), async {
while unsafe { libc::kill(pid as i32, 0) } == 0 {
tokio::time::sleep(Duration::from_millis(20)).await;
}
})
.await
.expect("cancelled handshake left its child alive");
}
#[test]
fn announced_name_precedence() {
// Daemon-decoded LXMF name always wins.
+6 -2
View File
@@ -146,12 +146,16 @@ impl MeshcoreDevice {
}
}
let info = DeviceInfo {
firmware_version: name.clone(),
let mut info = DeviceInfo {
firmware_version: "unknown".to_string(),
node_id,
max_contacts: 100,
device_type: super::types::DeviceType::Meshcore,
};
if let Some((version, capacity)) = self.query_device_info().await {
info.firmware_version = version;
info.max_contacts = capacity;
}
self.device_info = Some(info.clone());
info!("Meshcore initialization complete on {}", self.device_path);