diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 8157737..bec2d0a 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -27,9 +27,6 @@ jobs: - uses: Swatinem/rust-cache@v2 - - name: Format - run: cargo fmt --all -- --check - - name: Fetch wp-proto sibling repo run: | set -euo pipefail @@ -38,6 +35,9 @@ jobs: git clone --depth 1 https://github.com/dcc-bigfred/wireless-programmer.git "$dest" fi + - name: Format + run: cargo fmt --all -- --check + - name: Clippy run: cargo clippy --all-targets -- -D warnings diff --git a/Cargo.lock b/Cargo.lock index e6e2a9b..03f5263 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -162,6 +162,7 @@ dependencies = [ "url", "uuid", "wp-proto", + "z21-lan", ] [[package]] @@ -2025,7 +2026,7 @@ checksum = "d6bbff5f0aada427a1e5a6da5f1f98158182f26556f345ac9e04d36d0ebed650" [[package]] name = "wp-proto" -version = "0.1.0" +version = "0.2.0" dependencies = [ "serde", "serde_json", @@ -2061,6 +2062,14 @@ dependencies = [ "synstructure", ] +[[package]] +name = "z21-lan" +version = "0.1.0" +dependencies = [ + "thiserror 1.0.69", + "tokio", +] + [[package]] name = "zerocopy" version = "0.8.56" diff --git a/Cargo.toml b/Cargo.toml index 3009a82..b9f7a68 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -6,6 +6,9 @@ description = "Party/event kiosk for BigFred (SSO, accounts, handset pairing, CV license = "MIT" publish = false +[workspace] +members = [".", "crates/z21-lan"] + [[bin]] name = "bigfred-wizard" path = "src/main.rs" @@ -34,6 +37,7 @@ url = "2" qrcode = { version = "0.14.1", default-features = false, features = ["svg"] } # Wire types for the wireless-programmer Unix socket (sibling repo). wp-proto = { path = "../wireless-programmer/crates/wp-proto" } +z21-lan = { path = "crates/z21-lan" } [profile.release] lto = true diff --git a/Makefile b/Makefile index cce4020..e27cb43 100644 --- a/Makefile +++ b/Makefile @@ -25,7 +25,7 @@ else endif .PHONY: all build web-build release-musl host test test-release-assertions \ - fmt clippy clean dist dev-backend dev-web + fmt clippy clean dist dev-backend dev-web deploy-hub all: build @@ -95,3 +95,34 @@ clean: $(CARGO) clean rm -rf "$(WEB_DIR)/node_modules" dist find "$(WEB_DIR)/dist" -mindepth 1 ! -name .gitkeep -exec rm -rf {} + 2>/dev/null || true + +# --- Hub deploy (RO rootfs: binary lives on /data) ------------------------ +# Hub runs Dropbear. Older images lack /usr/libexec/sftp-server; -O uses +# legacy scp. Harmless on images that ship openssh sftp-server. +# +# make deploy-hub +# make deploy-hub HUB=192.168.0.10 +# +# /etc/init.d/bigfred-wizard prefers /data/opt/bigfred/bin/bigfred-wizard +# over the image copy in /usr/sbin. +HUB ?= 192.168.0.1 +HUB_USER ?= root +HUB_SSH ?= $(HUB_USER)@$(HUB) +SCP ?= scp +SCP_OPTS ?= -O +SSH ?= ssh +DIST_ARM64 ?= dist/bigfred-wizard-linux-arm64 +HUB_BIN_DIR ?= /data/opt/bigfred/bin + +# Upload next to the target and rename: writing in place fails with ETXTBSY +# once the hub is running the /data copy, and rename(2) swaps the inode +# atomically. +deploy-hub: release-musl + @test -f $(DIST_ARM64) || { echo "error: $(DIST_ARM64) missing — run make release-musl" >&2; exit 1; } + $(SSH) $(HUB_SSH) 'mkdir -p $(HUB_BIN_DIR)' + $(SCP) $(SCP_OPTS) $(DIST_ARM64) $(HUB_SSH):$(HUB_BIN_DIR)/.bigfred-wizard.new + $(SSH) $(HUB_SSH) 'set -e; \ + cd $(HUB_BIN_DIR); \ + chmod 755 .bigfred-wizard.new; \ + mv -f .bigfred-wizard.new bigfred-wizard; \ + microinit restart bigfred-wizard' diff --git a/crates/z21-lan/Cargo.toml b/crates/z21-lan/Cargo.toml new file mode 100644 index 0000000..5b47414 --- /dev/null +++ b/crates/z21-lan/Cargo.toml @@ -0,0 +1,18 @@ +[package] +name = "z21-lan" +version = "0.1.0" +edition = "2021" +description = "Z21 LAN (UDP) framing and locomotive CV programming client" +license = "MIT" +publish = false + +[lib] +name = "z21_lan" +path = "src/lib.rs" + +[dependencies] +thiserror = "1" +tokio = { version = "1", features = ["net", "time", "sync", "rt", "macros"] } + +[dev-dependencies] +tokio = { version = "1", features = ["net", "time", "sync", "rt", "macros", "rt-multi-thread"] } diff --git a/crates/z21-lan/src/client.rs b/crates/z21-lan/src/client.rs new file mode 100644 index 0000000..a212f1e --- /dev/null +++ b/crates/z21-lan/src/client.rs @@ -0,0 +1,95 @@ +//! Connected-UDP Z21 client with a single-flight send-and-await window. + +use std::net::SocketAddr; +use std::time::Duration; + +use tokio::net::UdpSocket; +use tokio::sync::Mutex; +use tokio::time::{timeout, Instant}; + +use crate::packets::{cv_read, cv_write, parse_cv_reply, pom_read, pom_write, CvReply}; + +const DEFAULT_TIMEOUT: Duration = Duration::from_secs(10); + +#[derive(Debug, thiserror::Error)] +pub enum CvError { + #[error("udp: {0}")] + Io(#[from] std::io::Error), + #[error("timeout waiting for CV reply")] + Timeout, + #[error("decoder NACK")] + Nack, + #[error("decoder short circuit")] + ShortCircuit, +} + +/// One connected UDP socket talking to a Z21 / RailBOX LAN command station. +pub struct Z21Client { + sock: UdpSocket, + io: Mutex<()>, + timeout: Duration, +} + +impl Z21Client { + pub async fn connect(addr: SocketAddr) -> Result { + let sock = UdpSocket::bind("0.0.0.0:0").await?; + sock.connect(addr).await?; + Ok(Self { + sock, + io: Mutex::new(()), + timeout: DEFAULT_TIMEOUT, + }) + } + + pub fn set_timeout(&mut self, timeout: Duration) { + if !timeout.is_zero() { + self.timeout = timeout; + } + } + + pub async fn read_cv(&self, cv: u16) -> Result { + let _g = self.io.lock().await; + self.await_result(&cv_read(cv), cv).await + } + + pub async fn write_cv(&self, cv: u16, value: u8) -> Result { + let _g = self.io.lock().await; + self.await_result(&cv_write(cv, value), cv).await + } + + pub async fn read_cv_pom(&self, addr: u16, cv: u16) -> Result { + let _g = self.io.lock().await; + self.await_result(&pom_read(addr, cv), cv).await + } + + /// POM write has no Z21 reply (spec §6.6). + pub async fn write_cv_pom(&self, addr: u16, cv: u16, value: u8) -> Result<(), CvError> { + let _g = self.io.lock().await; + self.sock.send(&pom_write(addr, cv, value)).await?; + Ok(()) + } + + async fn await_result(&self, req: &[u8], cv: u16) -> Result { + self.sock.send(req).await?; + let deadline = Instant::now() + self.timeout; + let mut buf = [0u8; 1500]; + loop { + let remaining = deadline.saturating_duration_since(Instant::now()); + if remaining.is_zero() { + return Err(CvError::Timeout); + } + let n = match timeout(remaining, self.sock.recv(&mut buf)).await { + Ok(Ok(n)) => n, + Ok(Err(err)) => return Err(err.into()), + Err(_) => return Err(CvError::Timeout), + }; + match parse_cv_reply(&buf[..n]) { + Some(CvReply::Result { cv: got, value }) if got == cv => return Ok(value), + Some(CvReply::Result { .. }) => continue, + Some(CvReply::Nack) => return Err(CvError::Nack), + Some(CvReply::NackShortCircuit) => return Err(CvError::ShortCircuit), + None => continue, + } + } + } +} diff --git a/crates/z21-lan/src/lib.rs b/crates/z21-lan/src/lib.rs new file mode 100644 index 0000000..967d7da --- /dev/null +++ b/crates/z21-lan/src/lib.rs @@ -0,0 +1,14 @@ +//! Z21 LAN (UDP) framing and locomotive CV / POM helpers. +//! +//! Packet layouts follow Roco Z21 LAN protocol §6. CV numbers on the wire +//! are zero-based (`0` = CV1). This crate is intentionally small: no radio, +//! no mDNS, no LocoNet dispatch. + +mod client; +mod packets; + +pub use client::{CvError, Z21Client}; +pub use packets::{ + cv_read, cv_write, encode_xbus, parse_cv_reply, parse_records, pom_read, pom_write, xor_sum, + CvReply, HEADER_XBUS, Z21_UDP_PORT, +}; diff --git a/crates/z21-lan/src/packets.rs b/crates/z21-lan/src/packets.rs new file mode 100644 index 0000000..29dc36e --- /dev/null +++ b/crates/z21-lan/src/packets.rs @@ -0,0 +1,170 @@ +//! Encode / parse Z21 LAN records used for decoder CV programming. + +/// Default Z21 LAN port (spec §1.1). +pub const Z21_UDP_PORT: u16 = 21105; +/// X-BUS tunnel header (`LAN_X_*`). +pub const HEADER_XBUS: u16 = 0x0040; + +/// XOR of every X-BUS byte except the checksum itself. +pub fn xor_sum(x: &[u8]) -> u8 { + x.iter().fold(0, |a, b| a ^ b) +} + +/// `DataLen LE | Header 0x0040 | xbus | xor(xbus)`. +pub fn encode_xbus(xbus: &[u8]) -> Vec { + let data_len = 4u16 + u16::try_from(xbus.len()).unwrap_or(u16::MAX) + 1; + let mut out = Vec::with_capacity(usize::from(data_len)); + out.extend_from_slice(&data_len.to_le_bytes()); + out.extend_from_slice(&HEADER_XBUS.to_le_bytes()); + out.extend_from_slice(xbus); + out.push(xor_sum(xbus)); + out +} + +/// NMRA CV number (1-based) to the Z21 wire value (`0` = CV1). +fn cv_wire(cv: u16) -> u16 { + cv.saturating_sub(1) +} + +fn loco_addr_bytes(addr: u16) -> (u8, u8) { + let mut msb = ((addr >> 8) & 0x3F) as u8; + if addr >= 128 { + msb |= 0xC0; + } + (msb, (addr & 0xFF) as u8) +} + +/// `LAN_X_CV_READ` (§6.1) — programming track, direct mode. +pub fn cv_read(cv: u16) -> Vec { + let w = cv_wire(cv); + encode_xbus(&[0x23, 0x11, (w >> 8) as u8, (w & 0xFF) as u8]) +} + +/// `LAN_X_CV_WRITE` (§6.2) — programming track, direct mode. +pub fn cv_write(cv: u16, value: u8) -> Vec { + let w = cv_wire(cv); + encode_xbus(&[0x24, 0x12, (w >> 8) as u8, (w & 0xFF) as u8, value]) +} + +/// `LAN_X_CV_POM_READ_BYTE` (§6.8). +pub fn pom_read(addr: u16, cv: u16) -> Vec { + let w = cv_wire(cv); + let (msb, lsb) = loco_addr_bytes(addr); + let db3 = 0xE4 | ((w >> 8) & 0x03) as u8; + encode_xbus(&[0xE6, 0x30, msb, lsb, db3, (w & 0xFF) as u8, 0x00]) +} + +/// `LAN_X_CV_POM_WRITE_BYTE` (§6.6). No Z21 reply. +pub fn pom_write(addr: u16, cv: u16, value: u8) -> Vec { + let w = cv_wire(cv); + let (msb, lsb) = loco_addr_bytes(addr); + let db3 = 0xEC | ((w >> 8) & 0x03) as u8; + encode_xbus(&[0xE6, 0x30, msb, lsb, db3, (w & 0xFF) as u8, value]) +} + +/// Parsed CV programming reply. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum CvReply { + Result { cv: u16, value: u8 }, + Nack, + NackShortCircuit, +} + +/// Walk concatenated Z21 records and return the first CV reply, if any. +pub fn parse_cv_reply(buf: &[u8]) -> Option { + for rec in parse_records(buf) { + if rec.header != HEADER_XBUS { + continue; + } + let d = rec.data; + if d.len() >= 6 && d[0] == 0x64 && d[1] == 0x14 { + let wire = (u16::from(d[2]) << 8) | u16::from(d[3]); + return Some(CvReply::Result { + cv: wire.saturating_add(1), + value: d[4], + }); + } + if d.len() >= 2 && d[0] == 0x61 && d[1] == 0x13 { + return Some(CvReply::Nack); + } + if d.len() >= 2 && d[0] == 0x61 && d[1] == 0x12 { + return Some(CvReply::NackShortCircuit); + } + } + None +} + +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct Record { + pub header: u16, + pub data: Vec, +} + +/// Walk concatenated Z21 records in one UDP datagram. +pub fn parse_records(buf: &[u8]) -> Vec { + let mut out = Vec::new(); + let mut off = 0usize; + while off + 4 <= buf.len() { + let data_len = u16::from_le_bytes([buf[off], buf[off + 1]]) as usize; + let header = u16::from_le_bytes([buf[off + 2], buf[off + 3]]); + if data_len < 4 || off + data_len > buf.len() { + break; + } + out.push(Record { + header, + data: buf[off + 4..off + data_len].to_vec(), + }); + off += data_len; + } + out +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn cv_read_matches_spec() { + let pkt = cv_read(1); + assert_eq!( + pkt, + vec![0x09, 0x00, 0x40, 0x00, 0x23, 0x11, 0x00, 0x00, 0x32] + ); + } + + #[test] + fn cv_write_matches_spec() { + let pkt = cv_write(8, 0x20); + assert_eq!( + pkt, + vec![0x0A, 0x00, 0x40, 0x00, 0x24, 0x12, 0x00, 0x07, 0x20, 0x11] + ); + } + + #[test] + fn parse_cv_result() { + let pkt = encode_xbus(&[0x64, 0x14, 0x00, 0x07, 0x20]); + assert_eq!( + parse_cv_reply(&pkt), + Some(CvReply::Result { cv: 8, value: 0x20 }) + ); + } + + #[test] + fn parse_nack() { + let pkt = encode_xbus(&[0x61, 0x13]); + assert_eq!(parse_cv_reply(&pkt), Some(CvReply::Nack)); + let sc = encode_xbus(&[0x61, 0x12]); + assert_eq!(parse_cv_reply(&sc), Some(CvReply::NackShortCircuit)); + } + + #[test] + fn pom_read_long_address() { + let pkt = pom_read(128, 1); + assert_eq!(pkt[0], 0x0C); + assert_eq!(&pkt[4..6], &[0xE6, 0x30]); + assert_eq!(pkt[6], 0xC0); + assert_eq!(pkt[7], 0x80); + assert_eq!(pkt[8], 0xE4); + } +} diff --git a/dev-config.json.example b/dev-config.json.example index a3a6488..1261169 100644 --- a/dev-config.json.example +++ b/dev-config.json.example @@ -24,5 +24,9 @@ "wirelessProgrammerSocket": "/data/run/wireless-programmer/wireless-programmer.sock", "fredProgramming": { "z21": { "address": "", "port": 0 } + }, + "locoProgramming": { + "mode": "bigfred", + "z21": { "address": "", "port": 0 } } } diff --git a/src/config.rs b/src/config.rs index 9b73841..0cb372c 100644 --- a/src/config.rs +++ b/src/config.rs @@ -89,6 +89,10 @@ pub struct Config { /// Digitrax FRED programming via a physical Z21 LAN command station. #[serde(default)] pub fred_programming: FredProgrammingConfig, + /// Loco CV / address programming: via BigFred dcc-bus (default) or + /// UDP straight to a Z21 / RailBOX. + #[serde(default)] + pub loco_programming: LocoProgrammingConfig, } /// FRED programming options (public; no secrets). @@ -116,6 +120,43 @@ impl FredZ21Config { } } +/// How the wizard programs locomotive decoders. +#[derive(Debug, Clone, Deserialize, Serialize, Default, PartialEq, Eq)] +#[serde(rename_all = "camelCase")] +pub enum LocoProgrammingMode { + #[default] + Bigfred, + Direct, +} + +/// Loco CV programming options (public; no secrets). +#[derive(Debug, Clone, Deserialize, Serialize, Default)] +#[serde(rename_all = "camelCase", default)] +pub struct LocoProgrammingConfig { + /// `bigfred` (default) uses dcc-bus; `direct` talks UDP to a Z21. + pub mode: LocoProgrammingMode, + /// Used when [`LocoProgrammingMode::Direct`]. + pub z21: FredZ21Config, + /// When both this and [`Self::layout_id`] are set, skip catalogue autodetection. + pub dcc_bus_id: Option, + /// Layout id paired with [`Self::dcc_bus_id`]. Ignored unless both are set. + pub layout_id: Option, +} + +impl LocoProgrammingConfig { + /// Hardcoded (command-station id, layout id) when both are present and non-zero. + pub fn fixed_dcc_bus(&self) -> Option<(u64, u64)> { + match (self.dcc_bus_id, self.layout_id) { + (Some(cs), Some(layout)) if cs > 0 && layout > 0 => Some((cs, layout)), + _ => None, + } + } + + pub fn is_direct(&self) -> bool { + self.mode == LocoProgrammingMode::Direct + } +} + fn default_throttle_host() -> String { "bigfred.local".to_string() } @@ -161,6 +202,7 @@ impl Default for Config { throttle_server_automatic: default_true(), wireless_programmer_socket: default_wireless_socket(), fred_programming: FredProgrammingConfig::default(), + loco_programming: LocoProgrammingConfig::default(), } } } @@ -186,6 +228,7 @@ pub struct PublicConfig { pub throttle_server_port: u16, pub throttle_server_automatic: bool, pub fred_programming: FredProgrammingConfig, + pub loco_programming: LocoProgrammingConfig, } #[derive(Debug, thiserror::Error)] @@ -261,6 +304,15 @@ impl Config { port: self.fred_programming.z21.port, }, }, + loco_programming: LocoProgrammingConfig { + mode: self.loco_programming.mode.clone(), + z21: FredZ21Config { + address: self.loco_programming.z21.address.trim().to_string(), + port: self.loco_programming.z21.port, + }, + dcc_bus_id: self.loco_programming.dcc_bus_id, + layout_id: self.loco_programming.layout_id, + }, } } @@ -501,6 +553,8 @@ mod tests { assert_eq!(json["throttleServerPort"], 12090); assert_eq!(json["fredProgramming"]["z21"]["address"], ""); assert_eq!(json["fredProgramming"]["z21"]["port"], 0); + assert_eq!(json["locoProgramming"]["mode"], "bigfred"); + assert_eq!(json["locoProgramming"]["z21"]["port"], 0); } #[test] @@ -518,6 +572,18 @@ mod tests { .skip_scan()); } + #[test] + fn loco_programming_fixed_ids_require_both() { + let mut cfg = LocoProgrammingConfig::default(); + assert!(cfg.fixed_dcc_bus().is_none()); + cfg.dcc_bus_id = Some(1); + assert!(cfg.fixed_dcc_bus().is_none()); + cfg.layout_id = Some(1); + assert_eq!(cfg.fixed_dcc_bus(), Some((1, 1))); + cfg.mode = LocoProgrammingMode::Direct; + assert!(cfg.is_direct()); + } + #[test] fn public_marks_psk_absent_when_empty() { let cfg = Config { diff --git a/src/dccbus_client.rs b/src/dccbus_client.rs index c6e0fe9..5934cdd 100644 --- a/src/dccbus_client.rs +++ b/src/dccbus_client.rs @@ -92,7 +92,7 @@ pub struct Ack { } /// One row of `GET /api/v1/command-stations/catalogue`. -#[derive(Debug, Clone, Deserialize)] +#[derive(Debug, Clone, Default, Deserialize)] #[serde(rename_all = "camelCase")] pub struct CommandStation { pub id: u64, @@ -367,7 +367,18 @@ impl DccBusClient { } /// Picks the programming-capable command station with the lowest id. + /// When `locoProgramming.dccBusId` and `layoutId` are both set, those + /// values are used instead of catalogue autodetection. pub async fn pick_station(&self, token: &str) -> Result { + let fixed = self.cfg.read().await.loco_programming.fixed_dcc_bus(); + if let Some((cs_id, _)) = fixed { + return Ok(CommandStation { + id: cs_id, + name: format!("dcc-bus #{cs_id}"), + programming: true, + ..CommandStation::default() + }); + } let api_base = self.cfg.read().await.bigfred_api_base(); let url = format!("{api_base}/api/v1/command-stations/catalogue"); let res = self diff --git a/src/main.rs b/src/main.rs index b07c00b..e170aa2 100644 --- a/src/main.rs +++ b/src/main.rs @@ -19,6 +19,7 @@ mod pin_api; mod programming_api; mod qr; mod wireless_api; +mod z21_direct; use std::net::SocketAddr; use std::path::PathBuf; @@ -38,6 +39,7 @@ use tower_http::trace::TraceLayer; use crate::config::{Config, PublicConfig}; use crate::dccbus_client::DccBusClient; +use crate::z21_direct::Z21DirectClient; /// Production SPA bundle. `make web-build` fills this directory before /// cargo runs; the placeholder keeps a fresh checkout compiling. @@ -64,6 +66,7 @@ pub struct AppState { pub cfg: Arc>, pub http: reqwest::Client, pub dcc: Arc, + pub z21: Arc, pub pulse_locks: Arc, pub wireless: Arc, } @@ -114,6 +117,7 @@ async fn main() -> Result<(), Box> { cfg: Arc::clone(&cfg), http: http.clone(), dcc: Arc::new(DccBusClient::new(Arc::clone(&cfg), http)), + z21: Arc::new(Z21DirectClient::new(Arc::clone(&cfg))), pulse_locks: Arc::new(programming_api::PulseLocks::default()), wireless: Arc::new(wireless_api::WirelessClient::new(Arc::clone(&cfg))), }; diff --git a/src/programming_api.rs b/src/programming_api.rs index 2fa5852..f91f049 100644 --- a/src/programming_api.rs +++ b/src/programming_api.rs @@ -131,6 +131,16 @@ pub async fn cvs_read( return Err(ApiError::bad_request("invalid_cvs")); } let mode = validate_mode(body.mode)?; + if state.config().await.loco_programming.is_direct() { + let ack = state + .z21 + .read_cvs(body.address, &body.cvs, mode.as_deref()) + .await?; + return Ok(Json(ProgrammingResponse { + ack, + command_station_id: None, + })); + } let payload = json!({ "address": body.address, "cvs": body.cvs, "mode": mode }); run(&state, &token, FRAME_CV_READ, payload).await } @@ -146,6 +156,16 @@ pub async fn cvs_write( return Err(ApiError::bad_request("invalid_cvs")); } let mode = validate_mode(body.mode)?; + if state.config().await.loco_programming.is_direct() { + let ack = state + .z21 + .write_cvs(body.address, &body.cvs, mode.as_deref()) + .await?; + return Ok(Json(ProgrammingResponse { + ack, + command_station_id: None, + })); + } let payload = json!({ "address": body.address, "cvs": body.cvs, "mode": mode }); run(&state, &token, FRAME_CV_WRITE, payload).await } @@ -166,6 +186,16 @@ pub async fn address_get( } } let payload = json!({ "address": body.address.unwrap_or(0), "mode": mode }); + if state.config().await.loco_programming.is_direct() { + let ack = state + .z21 + .addr_get(body.address.unwrap_or(0), mode.as_deref()) + .await?; + return Ok(Json(ProgrammingResponse { + ack, + command_station_id: None, + })); + } run(&state, &token, FRAME_ADDR_GET, payload).await } @@ -182,6 +212,16 @@ pub async fn address_set( "mode": mode, "verify": body.verify.unwrap_or(false), }); + if state.config().await.loco_programming.is_direct() { + let ack = state + .z21 + .addr_set(body.address, mode.as_deref(), body.verify.unwrap_or(false)) + .await?; + return Ok(Json(ProgrammingResponse { + ack, + command_station_id: None, + })); + } run(&state, &token, FRAME_ADDR_SET, payload).await } @@ -233,6 +273,9 @@ pub async fn function_pulse( /// Reports the socket state plus the station the wizard would program on. pub async fn status(State(state): State, headers: HeaderMap) -> ApiResult> { + if state.config().await.loco_programming.is_direct() { + return Ok(Json(state.z21.status())); + } let mut status = state.dcc.status(); if let Ok(token) = bearer(&headers) { if let Ok(station) = state.dcc.pick_station(&token).await { @@ -256,6 +299,9 @@ pub async fn connect(State(state): State, headers: HeaderMap) -> ApiRe "wizard_disabled", )); } + if state.config().await.loco_programming.is_direct() { + return Ok(Json(state.z21.ensure_connected().await?)); + } Ok(Json(state.dcc.ensure_connected(&token).await?)) } diff --git a/src/z21_direct.rs b/src/z21_direct.rs new file mode 100644 index 0000000..9c0812e --- /dev/null +++ b/src/z21_direct.rs @@ -0,0 +1,276 @@ +//! Direct UDP programming against a Z21 / RailBOX, bypassing dcc-bus. + +use std::net::SocketAddr; +use std::sync::Arc; +use std::time::Duration; + +use tokio::sync::{Mutex, RwLock}; + +use crate::config::Config; +use crate::dccbus_client::{Ack, CvEntry, Status}; +use crate::error::ApiError; + +const SETTLE: Duration = Duration::from_millis(300); +const ADDR_CVS: [u16; 4] = [1, 17, 18, 29]; + +pub struct Z21DirectClient { + cfg: Arc>, + inner: Mutex>, + status: std::sync::Mutex, +} + +impl Z21DirectClient { + pub fn new(cfg: Arc>) -> Self { + Self { + cfg, + inner: Mutex::new(None), + status: std::sync::Mutex::new(Status::default()), + } + } + + pub fn status(&self) -> Status { + self.status.lock().map(|g| g.clone()).unwrap_or_default() + } + + pub async fn ensure_connected(&self) -> Result { + let _ = self.ensure_inner().await?; + Ok(self.status()) + } + + pub async fn read_cvs( + &self, + address: u16, + cvs: &[u16], + mode: Option<&str>, + ) -> Result { + let guard = self.ensure_inner().await?; + let client = guard + .as_ref() + .ok_or_else(|| ApiError::unavailable("z21_unreachable"))?; + let pom = is_pom(mode); + let mut out = Vec::with_capacity(cvs.len()); + for (i, cv) in cvs.iter().copied().enumerate() { + if i > 0 { + tokio::time::sleep(SETTLE).await; + } + let value = if pom { + client.read_cv_pom(address, cv).await.map_err(map_cv_err)? + } else { + client.read_cv(cv).await.map_err(map_cv_err)? + }; + out.push(CvEntry { cv, value }); + } + Ok(Ack { + ok: true, + cvs: Some(out), + ..Ack::default() + }) + } + + pub async fn write_cvs( + &self, + address: u16, + cvs: &[CvEntry], + mode: Option<&str>, + ) -> Result { + let guard = self.ensure_inner().await?; + let client = guard + .as_ref() + .ok_or_else(|| ApiError::unavailable("z21_unreachable"))?; + let pom = is_pom(mode); + for (i, entry) in cvs.iter().enumerate() { + if i > 0 { + tokio::time::sleep(SETTLE).await; + } + if pom { + client + .write_cv_pom(address, entry.cv, entry.value) + .await + .map_err(map_cv_err)?; + } else { + client + .write_cv(entry.cv, entry.value) + .await + .map_err(map_cv_err)?; + } + } + Ok(Ack { + ok: true, + cvs: Some(cvs.to_vec()), + ..Ack::default() + }) + } + + pub async fn addr_get(&self, address: u16, mode: Option<&str>) -> Result { + let read = self.read_cvs(address, &ADDR_CVS, mode).await?; + let cvs = read.cvs.unwrap_or_default(); + let mut values = [0u8; 30]; + for e in &cvs { + if (e.cv as usize) < values.len() { + values[e.cv as usize] = e.value; + } + } + let (loco_address, long_address) = + address_from_cvs(values[1], values[17], values[18], values[29])?; + Ok(Ack { + ok: true, + cvs: Some(cvs), + loco_address: Some(loco_address), + long_address: Some(long_address), + ..Ack::default() + }) + } + + pub async fn addr_set( + &self, + address: u16, + mode: Option<&str>, + verify: bool, + ) -> Result { + let cv29 = self.read_cvs(address, &[29], mode).await?; + let current = cv29 + .cvs + .as_ref() + .and_then(|c| c.first()) + .map(|e| e.value) + .ok_or_else(|| ApiError::unavailable("programming_failed"))?; + let (writes, long) = address_cv_writes(address, current)?; + self.write_cvs(address, &writes, mode).await?; + if verify { + let got = self.addr_get(address, mode).await?; + if got.loco_address != Some(address) { + return Err( + ApiError::unavailable("programming_failed").with_detail("verify mismatch") + ); + } + } + Ok(Ack { + ok: true, + cvs: Some(writes), + loco_address: Some(address), + long_address: Some(long), + ..Ack::default() + }) + } + + async fn ensure_inner( + &self, + ) -> Result>, ApiError> { + let mut guard = self.inner.lock().await; + if guard.is_none() { + let cfg = self.cfg.read().await; + let z21 = &cfg.loco_programming.z21; + if !z21.skip_scan() { + return Err(ApiError::unavailable("z21_not_configured").with_detail( + "locoProgramming.z21.address and port are required in direct mode", + )); + } + let addr: SocketAddr = format!("{}:{}", z21.address.trim(), z21.port) + .parse() + .map_err(|err: std::net::AddrParseError| { + ApiError::bad_request("invalid_z21_address").with_detail(err.to_string()) + })?; + drop(cfg); + let client = z21_lan::Z21Client::connect(addr).await.map_err(|err| { + ApiError::unavailable("z21_unreachable").with_detail(err.to_string()) + })?; + *guard = Some(client); + if let Ok(mut status) = self.status.lock() { + status.connected = true; + status.last_error = None; + status.command_station_name = Some("Z21 (direct)".into()); + } + } + Ok(guard) + } +} + +fn is_pom(mode: Option<&str>) -> bool { + mode.map(str::trim) == Some("pom") +} + +fn map_cv_err(err: z21_lan::CvError) -> ApiError { + match err { + z21_lan::CvError::Timeout => ApiError::unavailable("programming_timeout"), + z21_lan::CvError::Nack | z21_lan::CvError::ShortCircuit => { + ApiError::unavailable("programming_failed").with_detail(err.to_string()) + } + z21_lan::CvError::Io(err) => { + ApiError::unavailable("z21_unreachable").with_detail(err.to_string()) + } + } +} + +fn address_from_cvs(cv1: u8, cv17: u8, cv18: u8, cv29: u8) -> Result<(u16, bool), ApiError> { + if cv29 & 0x20 != 0 { + let addr = (u16::from(cv17 & 0x3F) << 8) | u16::from(cv18); + Ok((addr, true)) + } else { + Ok((u16::from(cv1), false)) + } +} + +fn address_cv_writes(addr: u16, cv29: u8) -> Result<(Vec, bool), ApiError> { + if addr == 0 { + return Err(ApiError::bad_request("invalid_address")); + } + if addr <= 127 { + Ok(( + vec![ + CvEntry { + cv: 1, + value: addr as u8, + }, + CvEntry { + cv: 29, + value: cv29 & !0x20, + }, + ], + false, + )) + } else { + Ok(( + vec![ + CvEntry { + cv: 17, + value: ((addr >> 8) as u8) | 0xC0, + }, + CvEntry { + cv: 18, + value: addr as u8, + }, + CvEntry { + cv: 29, + value: cv29 | 0x20, + }, + ], + true, + )) + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn short_address_clears_long_bit() { + let (writes, long) = address_cv_writes(7, 0x26).unwrap(); + assert!(!long); + assert_eq!(writes[0].cv, 1); + assert_eq!(writes[0].value, 7); + assert_eq!(writes[1].cv, 29); + assert_eq!(writes[1].value, 0x06); + } + + #[test] + fn long_address_sets_cv17_18() { + let (writes, long) = address_cv_writes(1234, 0x06).unwrap(); + assert!(long); + assert_eq!(writes[0].cv, 17); + assert_eq!(writes[0].value, 0xC4); + assert_eq!(writes[1].cv, 18); + assert_eq!(writes[1].value, 0xD2); + assert_eq!(writes[2].value, 0x26); + } +} diff --git a/web/src/api/types.ts b/web/src/api/types.ts index be1dc4c..c34497b 100644 --- a/web/src/api/types.ts +++ b/web/src/api/types.ts @@ -21,6 +21,12 @@ export interface WizardConfig { fredProgramming: { z21: { address: string; port: number }; }; + locoProgramming: { + mode: "bigfred" | "direct"; + z21: { address: string; port: number }; + dccBusId?: number; + layoutId?: number; + }; } export interface HandsetSetup { diff --git a/web/src/i18n/de.json b/web/src/i18n/de.json index da3fe99..76f8d94 100644 --- a/web/src/i18n/de.json +++ b/web/src/i18n/de.json @@ -492,6 +492,8 @@ "dcc_pool_exhausted": "Keine freien DCC-Adressen mehr im Pool.", "invalid_address": "Ungültige DCC-Adresse.", "bigfred_unreachable": "Der BigFred-Server ist nicht erreichbar.", + "z21_not_configured": "Direktes Z21-Programmieren braucht Adresse und Port in der Wizard-Konfiguration.", + "z21_unreachable": "Die Z21-Zentrale hat im LAN nicht geantwortet.", "qr_url_unset": "Keine URL für diesen QR-Code konfiguriert.", "busy": "Der Regler ist beschäftigt — bitte gleich nochmal versuchen.", "noCandidates": "Kein Regler gefunden.", diff --git a/web/src/i18n/en.json b/web/src/i18n/en.json index 514db7c..a3ec80b 100644 --- a/web/src/i18n/en.json +++ b/web/src/i18n/en.json @@ -492,6 +492,8 @@ "dcc_pool_exhausted": "No free DCC addresses left in the pool.", "invalid_address": "Invalid DCC address.", "bigfred_unreachable": "The BigFred server is unreachable.", + "z21_not_configured": "Direct Z21 programming needs address and port in the wizard config.", + "z21_unreachable": "The Z21 command station did not answer on the LAN.", "qr_url_unset": "No URL configured for this QR code.", "busy": "The handset is busy — try again in a moment.", "noCandidates": "No handset found.", diff --git a/web/src/i18n/pl.json b/web/src/i18n/pl.json index a4171b9..7bbfb8f 100644 --- a/web/src/i18n/pl.json +++ b/web/src/i18n/pl.json @@ -492,6 +492,8 @@ "dcc_pool_exhausted": "Brak wolnych adresów DCC w puli.", "invalid_address": "Nieprawidłowy adres DCC.", "bigfred_unreachable": "Serwer BigFred jest niedostępny.", + "z21_not_configured": "Tryb direct wymaga adresu i portu Z21 w konfiguracji wizarda.", + "z21_unreachable": "Centralka Z21 nie odpowiedziała w sieci LAN.", "qr_url_unset": "Brak adresu do kodu QR w konfiguracji.", "busy": "Pilot jest zajęty, spróbuj za chwilę.", "noCandidates": "Nie znaleziono pilota.",