umsh_cli/
session.rs

1//! `CliSession` — the driver object that owns a clone of `LocalNode<M>`,
2//! a `CliIo`, and a `CliLogger`, plus the in-session state tables.
3//!
4//! The session holds a *clone* of `LocalNode<M>`. Other firmware components
5//! hold their own clones of the same underlying node state and continue to
6//! function independently.
7//!
8//! ## Driver pattern
9//!
10//! ```text
11//! tokio::select! {
12//!     r = host.run()  => { r?; }   // existing loop over pump_once()
13//!     r = cli.run()   => { r?; }   // loops read_line + service_events
14//! }
15//! ```
16//!
17//! Inside `cli.run()`:
18//! ```text
19//! loop {
20//!     select(read_line(&mut buf), wake.wait()):
21//!         Left(line)  → parse + execute (may call MAC-I/O directly)
22//!         Right(wake) → service_events() only
23//!     service_events() always runs after each select
24//! }
25//! ```
26
27use alloc::rc::Rc;
28use alloc::string::String;
29use alloc::vec::Vec;
30use core::cell::RefCell;
31use core::fmt::Write as _;
32
33use heapless::{Deque, FnvIndexMap, String as HString, Vec as HVec};
34use umsh_core::PublicKey;
35use umsh_mac::SendOptions;
36use umsh_node::{LocalNode, MacBackend, NodeError, OwnedMacCommand, Subscription};
37use umsh_text::UnicastTextChatWrapper;
38
39use crate::commands::{Command, ParseError, parse};
40use crate::events::{CliEvent, EVENT_LINE_MAX, EVENT_RAW_MAX};
41use crate::io::{CliInput, CliOutput};
42use crate::logger::{CliLogger, LogLevel};
43use crate::settings::SessionSettings;
44use crate::stats::Stats;
45use umsh_hal::{ChannelStore, PeerStore, PowerControl};
46use umsh_sync::AsyncCondition;
47
48/// Errors surfaced from `CliSession::run`.
49///
50/// `Node` flattens any `NodeError<…>` to a `String` via `format!("{:?}", e)`
51/// because the underlying error is generic over the `MacBackend`'s send and
52/// capacity error types and exposing it would force the CLI's error type to
53/// inherit those generics. Callers can read the message but cannot react
54/// programmatically to specific node-layer failure modes.
55///
56/// TODO: revisit this if a CLI consumer needs to branch on node errors.
57/// Options include adding a second generic parameter that carries the
58/// underlying `NodeError<E, C>`, or replacing the loss with a richer
59/// enum that mirrors the variants we actually care about at this layer.
60#[derive(Debug)]
61pub enum CliError<IoErr: core::fmt::Debug> {
62    Io(IoErr),
63    Node(String),
64}
65
66impl<IoErr: core::fmt::Debug> From<IoErr> for CliError<IoErr> {
67    fn from(e: IoErr) -> Self {
68        CliError::Io(e)
69    }
70}
71
72type SharedQueue<T> = Rc<RefCell<T>>;
73
74/// Entry in the CLI peer table.
75#[derive(Debug, Clone)]
76pub struct PeerEntry {
77    pub key: PublicKey,
78    pub alias: Option<HString<16>>,
79}
80
81/// Entry in the CLI channel table. Stores the raw key bytes so the `Channel`
82/// descriptor can be reconstructed for leave/bind operations.
83#[derive(Debug, Clone)]
84pub struct ChannelEntry {
85    pub name: HString<16>,
86    pub key_bytes: [u8; 32],
87}
88
89/// Outcome returned by `execute`. Not exposed outside this module.
90#[derive(Debug, PartialEq, Eq)]
91enum ExecOutcome {
92    Continue,
93    Quit,
94}
95
96/// The CLI driver. Owns the output half of the CLI transport; the input half
97/// is passed to [`CliSession::run`] so the driver can hold a long-lived read
98/// future across wake events without blocking writes to `out`.
99pub struct CliSession<
100    M,
101    OUT,
102    LOG,
103    PS,
104    CS,
105    PC,
106    const N_PEERS: usize,
107    const N_ALIASES: usize,
108    const N_CHANNELS: usize,
109    const N_EVENTS: usize,
110    const LINE_MAX: usize,
111> where
112    M: MacBackend,
113    M::SendError: core::fmt::Debug,
114    M::CapacityError: core::fmt::Debug,
115    OUT: CliOutput,
116    LOG: CliLogger,
117    PS: PeerStore,
118    CS: ChannelStore,
119    PC: PowerControl,
120{
121    pub(crate) node: LocalNode<M>,
122    pub(crate) local_key: PublicKey,
123    pub(crate) out: OUT,
124    pub(crate) logger: LOG,
125    pub(crate) peer_store: PS,
126    pub(crate) channel_store: CS,
127    pub(crate) power: PC,
128    pub(crate) peers: FnvIndexMap<PublicKey, PeerEntry, N_PEERS>,
129    pub(crate) aliases: FnvIndexMap<HString<16>, PublicKey, N_ALIASES>,
130    pub(crate) channels: FnvIndexMap<HString<16>, ChannelEntry, N_CHANNELS>,
131    pub(crate) events: SharedQueue<Deque<CliEvent, N_EVENTS>>,
132    pub(crate) events_dropped: SharedQueue<u64>,
133    pub(crate) stats: SharedQueue<Stats>,
134    pub(crate) wake: Rc<AsyncCondition>,
135    pub(crate) current_peer: Option<PublicKey>,
136    pub(crate) settings: SessionSettings,
137    /// Kept alive to prevent subscription teardown.
138    _subs: Vec<Subscription>,
139}
140
141impl<
142    M,
143    OUT,
144    LOG,
145    PS,
146    CS,
147    PC,
148    const N_PEERS: usize,
149    const N_ALIASES: usize,
150    const N_CHANNELS: usize,
151    const N_EVENTS: usize,
152    const LINE_MAX: usize,
153> CliSession<M, OUT, LOG, PS, CS, PC, N_PEERS, N_ALIASES, N_CHANNELS, N_EVENTS, LINE_MAX>
154where
155    M: MacBackend,
156    M::SendError: core::fmt::Debug,
157    M::CapacityError: core::fmt::Debug,
158    OUT: CliOutput,
159    LOG: CliLogger,
160    PS: PeerStore,
161    CS: ChannelStore,
162    PC: PowerControl,
163{
164    /// Construct a new session around a cloned `LocalNode<M>`. `local_key`
165    /// is passed in because the caller already has it; no node-crate accessor
166    /// for the local key is added. The session owns the `out` half of the
167    /// CLI transport; the input half is supplied to [`Self::run`].
168    ///
169    /// `peer_store` and `channel_store` are called on join/leave and add/rm
170    /// to persist changes. Pass [`umsh_hal::NoPeerStore`] /
171    /// [`umsh_hal::NoChannelStore`] when persistence is not needed.
172    ///
173    /// `power` is invoked by `/poweroff`; pass [`umsh_hal::NoPowerControl`]
174    /// on targets without a power-off path.
175    pub fn new(
176        node: LocalNode<M>,
177        local_key: PublicKey,
178        out: OUT,
179        logger: LOG,
180        peer_store: PS,
181        channel_store: CS,
182        power: PC,
183    ) -> Self {
184        let events: SharedQueue<Deque<CliEvent, N_EVENTS>> = Rc::new(RefCell::new(Deque::new()));
185        let events_dropped: SharedQueue<u64> = Rc::new(RefCell::new(0));
186        let stats: SharedQueue<Stats> = Rc::new(RefCell::new(Stats::default()));
187        let wake = Rc::new(AsyncCondition::new());
188
189        let subs = register_subscriptions(
190            &node,
191            events.clone(),
192            events_dropped.clone(),
193            stats.clone(),
194            wake.clone(),
195        );
196
197        Self {
198            node,
199            local_key,
200            out,
201            logger,
202            peer_store,
203            channel_store,
204            power,
205            peers: FnvIndexMap::new(),
206            aliases: FnvIndexMap::new(),
207            channels: FnvIndexMap::new(),
208            events,
209            events_dropped,
210            stats,
211            wake,
212            current_peer: None,
213            settings: SessionSettings::default(),
214            _subs: subs,
215        }
216    }
217
218    /// Resolve a `<peer-ref>` token to a `PublicKey`. Alias table first,
219    /// then encoded forms (hex / base58 / base64).
220    pub fn resolve_peer(&self, token: &str) -> Option<PublicKey> {
221        if let Ok(s) = HString::<16>::try_from(token) {
222            if let Some(k) = self.aliases.get(&s) {
223                return Some(*k);
224            }
225        }
226        crate::peer_ref::try_parse_pubkey(token)
227    }
228
229    /// Pre-register a peer at startup (e.g. from `--peer` CLI args).
230    /// Returns `false` if either the peer or alias table is full, or if the
231    /// MAC rejects the key (full peer table at the MAC layer).
232    ///
233    /// Also registers the peer in the MAC-layer peer table via `node.peer`
234    /// so inbound frames from this peer can be validated before the user
235    /// initiates any outbound traffic.
236    pub async fn register_peer(&mut self, key: PublicKey, alias: Option<&str>) -> bool {
237        if self.peers.contains_key(&key) {
238            return true;
239        }
240        let alias_heap = alias.and_then(|a| HString::<16>::try_from(a).ok());
241        if self
242            .peers
243            .insert(
244                key,
245                PeerEntry {
246                    key,
247                    alias: alias_heap.clone(),
248                },
249            )
250            .is_err()
251        {
252            return false;
253        }
254        if let Some(a) = alias_heap {
255            if self.aliases.insert(a, key).is_err() {
256                let _ = self.peers.remove(&key);
257                return false;
258            }
259        }
260        // Register at the MAC layer too — otherwise inbound unicast/auth
261        // packets from this peer would be dropped for missing keys.
262        if self.node.peer(key).await.is_err() {
263            let _ = self.peers.remove(&key);
264            if let Some(a) = alias.and_then(|s| HString::<16>::try_from(s).ok()) {
265                let _ = self.aliases.remove(&a);
266            }
267            return false;
268        }
269        true
270    }
271
272    /// Pre-register a channel at startup (e.g. from persisted storage).
273    ///
274    /// Joins the channel at the MAC layer and inserts it into the session
275    /// channel table. Returns `false` if the channel or MAC table is full, or
276    /// if the key is invalid. Does NOT call `channel_store` — use this only
277    /// when restoring data that is already durable.
278    #[cfg(feature = "software-crypto")]
279    pub async fn register_channel(&mut self, name: &str, key_bytes: [u8; 32]) -> bool {
280        let hname = match HString::<16>::try_from(name) {
281            Ok(s) => s,
282            Err(_) => return false,
283        };
284        if self.channels.contains_key(&hname) {
285            return true;
286        }
287        let channel_key = umsh_core::ChannelKey(key_bytes);
288        let channel = umsh_node::Channel::private(channel_key, name);
289        if self.node.join(&channel).await.is_err() {
290            return false;
291        }
292        let entry = ChannelEntry {
293            name: hname.clone(),
294            key_bytes,
295        };
296        if self.channels.insert(hname, entry).is_err() {
297            let _ = self.node.leave(&channel).await;
298            return false;
299        }
300        true
301    }
302
303    /// Populate the in-session peer and channel tables from persistent storage.
304    ///
305    /// Loads every peer and channel record from the backing stores and registers
306    /// them with both the MAC layer (idempotent) and the CLI display tables.
307    /// Call this once before [`run`](Self::run), or rely on `run` calling it
308    /// automatically at startup.
309    pub async fn load_from_stores(&mut self) {
310        // Collect into a local Vec first — the callback is sync, but
311        // `register_peer`/`register_channel` are async.
312        let mut peers: Vec<([u8; 32], Option<HString<16>>)> = Vec::new();
313        let _ = self
314            .peer_store
315            .for_each_peer(&mut |pk, alias| {
316                let alias_str = alias
317                    .and_then(|a| core::str::from_utf8(a).ok())
318                    .and_then(|s| HString::<16>::try_from(s).ok());
319                let _ = peers.push((*pk, alias_str));
320            })
321            .await;
322        for (pk, alias) in peers {
323            let _ = self
324                .register_peer(PublicKey(pk), alias.as_ref().map(|s| s.as_str()))
325                .await;
326        }
327
328        #[cfg(feature = "software-crypto")]
329        {
330            let mut channels: Vec<(HString<16>, [u8; 32])> = Vec::new();
331            let _ = self
332                .channel_store
333                .for_each_channel(&mut |name, key| {
334                    if let Ok(s) = core::str::from_utf8(name) {
335                        if let Ok(h) = HString::<16>::try_from(s) {
336                            let _ = channels.push((h, *key));
337                        }
338                    }
339                })
340                .await;
341            for (name, key_bytes) in channels {
342                let _ = self.register_channel(name.as_str(), key_bytes).await;
343            }
344        }
345    }
346
347    /// Build `SendOptions` from current CLI preferences.
348    fn send_opts(&self) -> SendOptions {
349        SendOptions::default()
350            .with_flood_hops(self.settings.flood_hops)
351            .with_ack_requested(self.settings.ack_requested)
352    }
353
354    /// Drive the CLI until `/quit` or EOF.
355    ///
356    /// Each outer iteration arms one long-lived `read_line` future and races
357    /// it against `wake` in an inner loop. The read future is only dropped
358    /// when it completes — wake-driven iterations preserve it across
359    /// `select!` rearms, so implementations of [`CliInput`] do not need to
360    /// be cancel-safe.
361    ///
362    /// `input` is borrowed for the duration of each read; it's a separate
363    /// parameter from `self` so the driver can continue writing to `self.out`
364    /// while a read is outstanding.
365    pub async fn run<IN>(&mut self, input: &mut IN) -> Result<(), CliError<OUT::Error>>
366    where
367        IN: CliInput<Error = OUT::Error>,
368    {
369        use futures::future::{Either, select};
370
371        self.load_from_stores().await;
372
373        // Service any events that arrived before first user input
374        // (e.g. from a previously started host).
375        self.service_events().await?;
376
377        let mut buf = [0u8; LINE_MAX];
378        let wake = self.wake.clone();
379        loop {
380            // Arm a single read future and keep it alive across wake events.
381            let read_fut = input.read_line(&mut buf);
382            futures::pin_mut!(read_fut);
383            let owned: Option<String> = loop {
384                let wait_fut = wake.wait();
385                futures::pin_mut!(wait_fut);
386                match select(read_fut.as_mut(), wait_fut).await {
387                    Either::Left((result, _)) => match result {
388                        Ok(Some(s)) => break Some(String::from(s)),
389                        Ok(None) => break None,
390                        Err(e) => return Err(CliError::Io(e)),
391                    },
392                    Either::Right(((), _)) => {
393                        // Wake fired; drain events without dropping read_fut.
394                        self.service_events().await?;
395                    }
396                }
397            };
398
399            match owned {
400                None => return Ok(()), // EOF
401                Some(line) => match parse(&line) {
402                    Err(ParseError::Empty) => {}
403                    Err(e) => {
404                        let msg = format_parse_error(&e);
405                        self.write_err(&msg).await?;
406                    }
407                    Ok(cmd) => match self.execute(cmd).await? {
408                        ExecOutcome::Continue => {}
409                        ExecOutcome::Quit => return Ok(()),
410                    },
411                },
412            }
413
414            // Drain events queued by execute() (e.g. send receipts).
415            self.service_events().await?;
416        }
417    }
418
419    // ─── Event-queue drain ──────────────────────────────────────────────────
420
421    async fn service_events(&mut self) -> Result<(), CliError<OUT::Error>> {
422        loop {
423            // Take one event per iteration (avoids holding RefCell across await).
424            let event = self.events.borrow_mut().pop_front();
425            let Some(event) = event else { break };
426            self.handle_event(event).await?;
427        }
428        Ok(())
429    }
430
431    async fn handle_event(&mut self, event: CliEvent) -> Result<(), CliError<OUT::Error>> {
432        match event {
433            CliEvent::Received {
434                from,
435                hops: _,
436                rssi,
437                snr,
438                prefix,
439            } => {
440                // Try to decode as text.
441                if let Some([pt, rest @ ..]) = prefix.as_slice().get(0..) {
442                    if *pt == umsh_core::PayloadType::TextMessage as u8 {
443                        if let Ok(msg) = umsh_text::parse_text_message(rest) {
444                            let alias = self.peer_alias_display(&from);
445                            let mut line: HString<EVENT_LINE_MAX> = HString::new();
446                            let _ = write!(
447                                &mut line,
448                                "<{}> {}",
449                                alias,
450                                msg.body_str().unwrap_or("<invalid utf-8>")
451                            );
452                            self.out.write_line(&line).await?;
453                            if self.settings.show_hex {
454                                let mut hex: HString<EVENT_LINE_MAX> = HString::new();
455                                let _ = write!(&mut hex, "  hex:");
456                                for b in prefix.iter() {
457                                    let _ = write!(&mut hex, " {:02x}", b);
458                                }
459                                self.out.write_line(&hex).await?;
460                            }
461                            return Ok(());
462                        }
463                    }
464                }
465                // Unknown payload — show origin + hint.
466                let alias = self.peer_alias_display(&from);
467                let mut line: HString<EVENT_LINE_MAX> = HString::new();
468                let _ = write!(
469                    &mut line,
470                    "[pkt from {} rssi={:?} snr={:?}]",
471                    alias, rssi, snr
472                );
473                self.out.write_line(&line).await?;
474            }
475
476            CliEvent::AckReceived { peer } => {
477                let alias = self.peer_alias_display(&peer);
478                let mut line: HString<EVENT_LINE_MAX> = HString::new();
479                let _ = write!(&mut line, "[ack from {}]", alias);
480                self.out.write_line(&line).await?;
481            }
482
483            CliEvent::AckTimeout { peer } => {
484                let alias = self.peer_alias_display(&peer);
485                let mut line: HString<EVENT_LINE_MAX> = HString::new();
486                let _ = write!(&mut line, "[ack timeout → {}]", alias);
487                self.out.write_line(&line).await?;
488            }
489
490            CliEvent::NodeDiscovered { from, name } => {
491                let alias = self.peer_alias_display(&from);
492                let mut line: HString<EVENT_LINE_MAX> = HString::new();
493                let _ = write!(
494                    &mut line,
495                    "[node discovered: {}{}]",
496                    alias,
497                    name.as_deref().unwrap_or(""),
498                );
499                self.out.write_line(&line).await?;
500            }
501
502            CliEvent::Beacon { hint, from } => {
503                let mut line: HString<EVENT_LINE_MAX> = HString::new();
504                match from {
505                    Some(key) => {
506                        let alias = self.peer_alias_display(&key);
507                        let _ = write!(&mut line, "[beacon from {}]", alias);
508                    }
509                    None => {
510                        // Hint-only beacon (sender didn't include the full pubkey).
511                        let _ = write!(&mut line, "[beacon hint:{}]", hint);
512                    }
513                }
514                self.out.write_line(&line).await?;
515            }
516
517            CliEvent::PfsEstablished { peer } => {
518                let alias = self.peer_alias_display(&peer);
519                let mut line: HString<EVENT_LINE_MAX> = HString::new();
520                let _ = write!(&mut line, "[pfs established with {}]", alias);
521                self.out.write_line(&line).await?;
522            }
523
524            CliEvent::PfsEnded { peer } => {
525                let alias = self.peer_alias_display(&peer);
526                let mut line: HString<EVENT_LINE_MAX> = HString::new();
527                let _ = write!(&mut line, "[pfs ended with {}]", alias);
528                self.out.write_line(&line).await?;
529            }
530
531            CliEvent::PfsFailed { peer, reason } => {
532                let alias = self.peer_alias_display(&peer);
533                let mut line: HString<EVENT_LINE_MAX> = HString::new();
534                let _ = write!(
535                    &mut line,
536                    "[pfs failed with {}: {}]",
537                    alias,
538                    pfs_failure_str(reason)
539                );
540                self.out.write_line(&line).await?;
541            }
542
543            CliEvent::Pong { peer, rtt_ms } => {
544                let alias = self.peer_alias_display(&peer);
545                let mut line: HString<EVENT_LINE_MAX> = HString::new();
546                let _ = write!(&mut line, "pong {} rtt={} ms", alias, rtt_ms);
547                self.out.write_line(&line).await?;
548            }
549
550            CliEvent::PingTimeout { peer } => {
551                let alias = self.peer_alias_display(&peer);
552                let mut line: HString<EVENT_LINE_MAX> = HString::new();
553                let _ = write!(&mut line, "[ping timeout: {}]", alias);
554                self.out.write_line(&line).await?;
555            }
556
557            CliEvent::UnknownMacCmdIn { peer, cmd_id } => {
558                let alias = self.peer_alias_display(&peer);
559                let mut line: HString<EVENT_LINE_MAX> = HString::new();
560                let _ = write!(&mut line, "[mac cmd 0x{:02x} from {}]", cmd_id, alias);
561                self.out.write_line(&line).await?;
562            }
563
564            CliEvent::OutputLine { line } => {
565                self.out.write_line(&line).await?;
566            }
567
568            CliEvent::RawTx { bytes } => {
569                if self.settings.show_raw {
570                    let mut line: HString<EVENT_LINE_MAX> = HString::new();
571                    let _ = write!(&mut line, "tx");
572                    for b in bytes.iter() {
573                        let _ = write!(&mut line, " {:02x}", b);
574                    }
575                    self.out.write_line(&line).await?;
576                }
577            }
578
579            CliEvent::RawRx { bytes } => {
580                if self.settings.show_raw {
581                    let mut line: HString<EVENT_LINE_MAX> = HString::new();
582                    let _ = write!(&mut line, "rx");
583                    for b in bytes.iter() {
584                        let _ = write!(&mut line, " {:02x}", b);
585                    }
586                    self.out.write_line(&line).await?;
587                }
588            }
589
590            // Outbound variants are dead in this implementation — `execute()`
591            // calls MAC I/O directly and never pushes these. See the TODO in
592            // `events.rs` for context on the intended future model.
593            CliEvent::SendText { .. }
594            | CliEvent::StartPfs { .. }
595            | CliEvent::EndPfs { .. }
596            | CliEvent::ChannelSend { .. }
597            | CliEvent::SendBeacon
598            | CliEvent::SendRaw { .. } => {}
599        }
600        Ok(())
601    }
602
603    // ─── Command dispatcher ──────────────────────────────────────────────────
604
605    async fn execute(&mut self, cmd: Command<'_>) -> Result<ExecOutcome, CliError<OUT::Error>> {
606        match cmd {
607            Command::Quit => Ok(ExecOutcome::Quit),
608            Command::Help(topic) => self.cmd_help(topic).await.map(|_| ExecOutcome::Continue),
609            Command::WhoAmI => self.cmd_whoami().await.map(|_| ExecOutcome::Continue),
610            Command::PeerAdd { pubkey, alias } => self
611                .cmd_peer_add(pubkey, alias)
612                .await
613                .map(|_| ExecOutcome::Continue),
614            Command::PeerAlias { peer, alias } => self
615                .cmd_peer_alias(peer, alias)
616                .await
617                .map(|_| ExecOutcome::Continue),
618            Command::PeerRm { peer } => self.cmd_peer_rm(peer).await.map(|_| ExecOutcome::Continue),
619            Command::Peers => self.cmd_peers().await.map(|_| ExecOutcome::Continue),
620            Command::Query { peer } => self.cmd_query(peer).await.map(|_| ExecOutcome::Continue),
621            Command::Set { var, val } => {
622                self.cmd_set(var, val).await.map(|_| ExecOutcome::Continue)
623            }
624            Command::SetShow => self.cmd_set_show().await.map(|_| ExecOutcome::Continue),
625            Command::Log { level } => self.cmd_log(level).await.map(|_| ExecOutcome::Continue),
626            Command::Stats => self.cmd_stats().await.map(|_| ExecOutcome::Continue),
627            Command::Counters => self.cmd_counters().await.map(|_| ExecOutcome::Continue),
628            Command::Channels => self.cmd_channels().await.map(|_| ExecOutcome::Continue),
629            Command::PfsStatus { peer } => self
630                .cmd_pfs_status(peer)
631                .await
632                .map(|_| ExecOutcome::Continue),
633            Command::Msg { peer, text } => self
634                .cmd_msg(peer, text)
635                .await
636                .map(|_| ExecOutcome::Continue),
637            Command::Text { body } => self.cmd_text(body).await.map(|_| ExecOutcome::Continue),
638            Command::Me { action } => self.cmd_me(action).await.map(|_| ExecOutcome::Continue),
639            Command::Ping { peer, bytes } => self
640                .cmd_ping(peer, bytes)
641                .await
642                .map(|_| ExecOutcome::Continue),
643            Command::PfsStart { peer, minutes } => self
644                .cmd_pfs_start(peer, minutes)
645                .await
646                .map(|_| ExecOutcome::Continue),
647            Command::PfsEnd { peer } => self.cmd_pfs_end(peer).await.map(|_| ExecOutcome::Continue),
648            Command::Beacon => self.cmd_beacon().await.map(|_| ExecOutcome::Continue),
649            Command::ChannelJoin { name, key } => self
650                .cmd_channel_join(name, key)
651                .await
652                .map(|_| ExecOutcome::Continue),
653            Command::ChannelLeave { name } => self
654                .cmd_channel_leave(name)
655                .await
656                .map(|_| ExecOutcome::Continue),
657            Command::ChannelSend { name, text } => self
658                .cmd_channel_send(name, text)
659                .await
660                .map(|_| ExecOutcome::Continue),
661            Command::Raw { peer, hex } => {
662                self.cmd_raw(peer, hex).await.map(|_| ExecOutcome::Continue)
663            }
664            Command::PowerOff => self.cmd_power_off().await.map(|_| ExecOutcome::Continue),
665            Command::Reboot => self.cmd_reboot().await.map(|_| ExecOutcome::Continue),
666        }
667    }
668
669    async fn cmd_power_off(&mut self) -> Result<(), CliError<OUT::Error>> {
670        self.out.write_line("powering off").await?;
671        self.power.request_power_off();
672        Ok(())
673    }
674
675    async fn cmd_reboot(&mut self) -> Result<(), CliError<OUT::Error>> {
676        self.out.write_line("rebooting").await?;
677        self.power.request_reboot();
678        Ok(())
679    }
680
681    // ─── Local-state commands ────────────────────────────────────────────────
682
683    async fn cmd_help(&mut self, topic: Option<&str>) -> Result<(), CliError<OUT::Error>> {
684        if let Some(t) = topic {
685            return self.cmd_help_topic(t).await;
686        }
687        let lines: &[&str] = &[
688            "session:",
689            "  /help [command]             show help, or detailed help for <command>",
690            "  /quit                       exit the CLI",
691            "  /whoami                     print the local public key",
692            "  /log <level>                set verbosity: error|warn|info|debug|trace",
693            "  /poweroff                   request a hardware power-off (alias: /off)",
694            "  /reboot                     request a soft reboot",
695            "",
696            "peers:",
697            "  /peer add <pubkey> [alias]  register a peer (base58/base64/hex, 32 bytes)",
698            "  /peer alias <peer> <alias>  rename a registered peer",
699            "  /peer rm <peer-ref>         remove a peer",
700            "  /peers                      list registered peers with full public key",
701            "  /query <peer-ref>           set the current peer for bare text",
702            "",
703            "messaging:",
704            "  /msg <peer-ref> <text>      send text to peer",
705            "  <text>                      bare text goes to the current peer",
706            "  /me <action>                emote (e.g. /me waves)",
707            "  /raw <peer-ref> <hex>       send raw hex payload bytes",
708            "",
709            "diagnostics:",
710            "  /ping <peer-ref> [bytes]    send an EchoRequest (default 8 bytes)",
711            "  /beacon                     broadcast a beacon",
712            "  /stats                      show TX/RX counters, RSSI, queue depth",
713            "  /counters                   show frame-counter state (local TX + per-peer RX)",
714            "",
715            "pfs (perfect forward secrecy):",
716            "  /pfs start <peer-ref> [min] request a PFS session (default 60 min)",
717            "  /pfs end <peer-ref>         end a PFS session",
718            "  /pfs status [peer-ref]      show PFS session state",
719            "",
720            "channels:",
721            "  /channel join <name> <key>  join a channel (key is base58)",
722            "  /channel leave <name>       leave a channel",
723            "  /channel send <name> <txt>  send multicast text",
724            "  /channels                   list joined channels",
725            "",
726            "settings:",
727            "  /set                        show current settings",
728            "  /set <var> <val>            flood_hops|ack_requested|show_hex|show_raw",
729        ];
730        for l in lines {
731            self.out.write_line(l).await?;
732        }
733        Ok(())
734    }
735
736    async fn cmd_help_topic(&mut self, topic: &str) -> Result<(), CliError<OUT::Error>> {
737        let t = topic.trim().trim_start_matches('/');
738        let detail: &[&str] = match t {
739            "quit" => &["/quit — exit the CLI (EOF does the same)."],
740            "help" => &[
741                "/help [command] — list all commands, or show detailed help for one.",
742                "  example: /help ping",
743            ],
744            "whoami" => &["/whoami — print the local public key as hex."],
745            "log" => &[
746                "/log <level> — set log verbosity.",
747                "  levels: error, warn, info, debug, trace",
748            ],
749            "poweroff" | "off" => &[
750                "/poweroff — request a controlled power-off.",
751                "  Persists any pending counters, sleeps the display, drops the",
752                "  peripheral rail, and enters System OFF. On supported boards a",
753                "  button press resumes the device (via reset + reboot).",
754                "  Alias: /off.",
755            ],
756            "reboot" => &[
757                "/reboot — request a soft reboot.",
758                "  Persists any pending counters, then triggers a system reset.",
759                "  The device comes back up running the same firmware image.",
760            ],
761            "peer" => &[
762                "/peer add <pubkey> [alias] — register a peer.",
763                "  <pubkey> accepts base58, base64, or hex (32-byte Ed25519 key).",
764                "  Also registers the peer at the MAC layer so inbound frames validate.",
765                "/peer alias <peer-ref> <alias> — rename a registered peer.",
766                "/peer rm <peer-ref> — remove a peer. <peer-ref> is an alias or full key.",
767            ],
768            "peers" => &["/peers — list registered peers: alias and full public key (hex)."],
769            "query" => &[
770                "/query <peer-ref> — set the current peer for bare-text sends.",
771                "  After /query bob, a bare line is sent to bob as a text message.",
772            ],
773            "msg" => &["/msg <peer-ref> <text> — send a text message to a peer."],
774            "me" => &[
775                "/me <action> — send an emote to the current peer.",
776                "  example: /me waves  →  sent as \"* waves\"",
777            ],
778            "raw" => &[
779                "/raw <peer-ref> <hex> — send raw payload bytes as a unicast packet.",
780                "  <hex> is an even-length hex string (no 0x prefix, no spaces).",
781            ],
782            "ping" => &[
783                "/ping <peer-ref> [bytes] — send a MAC-level EchoRequest.",
784                "  [bytes] is the total payload size (2..=60, default 8).",
785                "  The first 2 bytes are a nonce used to match the response.",
786                "  Prints \"pong <peer> rtt=<ms>\" when the reply arrives.",
787            ],
788            "beacon" => &["/beacon — broadcast a beacon frame announcing this node."],
789            "stats" => &[
790                "/stats — print counters maintained by the CLI.",
791                "  TX/RX packets, ACK outcomes, last RSSI/SNR, event-queue depth,",
792                "  and events_dropped (non-zero if the inbound queue overflowed).",
793            ],
794            "counters" => &[
795                "/counters — show the local TX frame counter and each known peer's RX counter (live and persisted boundaries).",
796            ],
797            "pfs" => &[
798                "/pfs start <peer-ref> [minutes] — request a PFS session.",
799                "  [minutes] is the requested lifetime (default 60).",
800                "/pfs end <peer-ref> — tear down an active PFS session.",
801                "/pfs status [peer-ref] — show PFS state for one peer or all.",
802            ],
803            "channel" | "channels" => &[
804                "/channel join <name> <key-b58> — bind a channel by name + shared key.",
805                "/channel leave <name> — leave a channel.",
806                "/channel send <name> <text> — send a multicast text message.",
807                "/channels — list currently joined channels.",
808            ],
809            "set" => &[
810                "/set — show current CLI-local settings.",
811                "/set <var> <val> — change one setting (resets on exit).",
812                "  flood_hops     u8 in 0..=15 (default 5) — max FHOPS_REM; a known route narrows it",
813                "  ack_requested  bool (default true)     — request MAC acks on unicast",
814                "  show_hex       bool (default false)    — also print inbound bytes as hex",
815                "  show_raw       bool (default false)    — log every TX/RX packet as hex",
816            ],
817            other => {
818                let mut msg: HString<EVENT_LINE_MAX> = HString::new();
819                let _ = write!(
820                    &mut msg,
821                    "no help for '{}' — try /help for the full list",
822                    other
823                );
824                return self.write_err(&msg).await;
825            }
826        };
827        for l in detail {
828            self.out.write_line(l).await?;
829        }
830        Ok(())
831    }
832
833    async fn cmd_whoami(&mut self) -> Result<(), CliError<OUT::Error>> {
834        let mut line: HString<EVENT_LINE_MAX> = HString::new();
835        let _ = write!(&mut line, "local: {}", self.local_key);
836        self.out.write_line(&line).await?;
837        Ok(())
838    }
839
840    async fn cmd_peer_add(
841        &mut self,
842        pubkey: &str,
843        alias: Option<&str>,
844    ) -> Result<(), CliError<OUT::Error>> {
845        let key = match crate::peer_ref::try_parse_pubkey(pubkey) {
846            Some(k) => k,
847            None => return self.write_err("invalid pubkey").await,
848        };
849        let alias_heap = match alias {
850            Some(a) => match HString::<16>::try_from(a) {
851                Ok(s) => Some(s),
852                Err(_) => return self.write_err("alias too long (max 16 chars)").await,
853            },
854            None => None,
855        };
856        if self.peers.contains_key(&key) {
857            return self.write_err("peer already registered").await;
858        }
859        if self
860            .peers
861            .insert(
862                key,
863                PeerEntry {
864                    key,
865                    alias: alias_heap.clone(),
866                },
867            )
868            .is_err()
869        {
870            return self.write_err("peer table full").await;
871        }
872        if let Some(a) = alias_heap.clone() {
873            if self.aliases.insert(a, key).is_err() {
874                let _ = self.peers.remove(&key);
875                return self.write_err("alias table full").await;
876            }
877        }
878        // Register at MAC layer so inbound frames from this peer validate.
879        if let Err(e) = self.node.peer(key).await {
880            let _ = self.peers.remove(&key);
881            if let Some(a) = alias_heap {
882                let _ = self.aliases.remove(&a);
883            }
884            let msg = node_err_str(&e);
885            return self.write_err(&msg).await;
886        }
887        // Persist to flash (best-effort; the peer is usable even if storage fails).
888        let alias_bytes = alias_heap.as_ref().map(|s: &HString<16>| s.as_bytes());
889        let _ = self.peer_store.store_peer(&key.0, alias_bytes).await;
890        self.out.write_line("ok").await?;
891        Ok(())
892    }
893
894    async fn cmd_peer_rm(&mut self, peer: &str) -> Result<(), CliError<OUT::Error>> {
895        let Some(key) = self.resolve_peer(peer) else {
896            return self.write_err("unknown peer").await;
897        };
898        let entry = match self.peers.remove(&key) {
899            Some(e) => e,
900            None => return self.write_err("peer not in table").await,
901        };
902        if let Some(a) = entry.alias {
903            let _ = self.aliases.remove(&a);
904        }
905        if self.current_peer == Some(key) {
906            self.current_peer = None;
907        }
908        // Remove from persistent storage (best-effort).
909        let _ = self.peer_store.delete_peer(&key.0).await;
910        self.out.write_line("ok").await?;
911        Ok(())
912    }
913
914    async fn cmd_peer_alias(
915        &mut self,
916        peer: &str,
917        new_alias: &str,
918    ) -> Result<(), CliError<OUT::Error>> {
919        let Some(key) = self.resolve_peer(peer) else {
920            return self.write_err("unknown peer").await;
921        };
922        let new_alias_heap = match HString::<16>::try_from(new_alias) {
923            Ok(s) => s,
924            Err(_) => return self.write_err("alias too long (max 16 chars)").await,
925        };
926        // Remove old alias from lookup table.
927        if let Some(entry) = self.peers.get(&key) {
928            if let Some(old_alias) = entry.alias.clone() {
929                let _ = self.aliases.remove(&old_alias);
930            }
931        }
932        // Insert new alias into lookup table.
933        if self.aliases.insert(new_alias_heap.clone(), key).is_err() {
934            return self.write_err("alias table full").await;
935        }
936        // Update the peer entry.
937        if let Some(entry) = self.peers.get_mut(&key) {
938            entry.alias = Some(new_alias_heap.clone());
939        }
940        // Persist (best-effort).
941        let _ = self
942            .peer_store
943            .store_peer(&key.0, Some(new_alias_heap.as_bytes()))
944            .await;
945        self.out.write_line("ok").await?;
946        Ok(())
947    }
948
949    async fn cmd_peers(&mut self) -> Result<(), CliError<OUT::Error>> {
950        if self.peers.is_empty() {
951            self.out.write_line("(no peers)").await?;
952            return Ok(());
953        }
954        let mut lines: Vec<String> = Vec::new();
955        for (_k, entry) in self.peers.iter() {
956            let mut line: HString<EVENT_LINE_MAX> = HString::new();
957            let alias = entry.alias.as_deref().unwrap_or("-");
958            let _ = write!(&mut line, "{:16} {}", alias, entry.key);
959            lines.push(String::from(line.as_str()));
960        }
961        for l in lines {
962            self.out.write_line(&l).await?;
963        }
964        Ok(())
965    }
966
967    async fn cmd_query(&mut self, peer: &str) -> Result<(), CliError<OUT::Error>> {
968        let Some(key) = self.resolve_peer(peer) else {
969            return self.write_err("unknown peer").await;
970        };
971        self.current_peer = Some(key);
972        self.out.write_line("ok").await?;
973        Ok(())
974    }
975
976    async fn cmd_set_show(&mut self) -> Result<(), CliError<OUT::Error>> {
977        let mut line: HString<EVENT_LINE_MAX> = HString::new();
978        let _ = write!(
979            &mut line,
980            "flood_hops={} ack_requested={} show_hex={} show_raw={}",
981            self.settings.flood_hops,
982            self.settings.ack_requested,
983            self.settings.show_hex,
984            self.settings.show_raw,
985        );
986        self.out.write_line(&line).await?;
987        Ok(())
988    }
989
990    async fn cmd_set(&mut self, var: &str, val: &str) -> Result<(), CliError<OUT::Error>> {
991        match var {
992            "flood_hops" => match val.parse::<u8>() {
993                Ok(v) if v <= 15 => {
994                    self.settings.flood_hops = v;
995                    self.out.write_line("ok").await?;
996                }
997                _ => self.write_err("flood_hops must be 0..=15").await?,
998            },
999            "ack_requested" => match parse_bool(val) {
1000                Some(b) => {
1001                    self.settings.ack_requested = b;
1002                    self.out.write_line("ok").await?;
1003                }
1004                None => self.write_err("ack_requested: expected true|false").await?,
1005            },
1006            "show_hex" => match parse_bool(val) {
1007                Some(b) => {
1008                    self.settings.show_hex = b;
1009                    self.out.write_line("ok").await?;
1010                }
1011                None => self.write_err("show_hex: expected true|false").await?,
1012            },
1013            "show_raw" => match parse_bool(val) {
1014                Some(b) => {
1015                    self.settings.show_raw = b;
1016                    self.out.write_line("ok").await?;
1017                }
1018                None => self.write_err("show_raw: expected true|false").await?,
1019            },
1020            _ => {
1021                self.write_err("unknown setting (flood_hops / ack_requested / show_hex / show_raw)")
1022                    .await?
1023            }
1024        }
1025        Ok(())
1026    }
1027
1028    async fn cmd_log(&mut self, level: &str) -> Result<(), CliError<OUT::Error>> {
1029        let lvl = match level {
1030            "error" => LogLevel::Error,
1031            "warn" => LogLevel::Warn,
1032            "info" => LogLevel::Info,
1033            "debug" => LogLevel::Debug,
1034            "trace" => LogLevel::Trace,
1035            _ => return self.write_err("level: error|warn|info|debug|trace").await,
1036        };
1037        self.logger.set_level(lvl);
1038        self.out.write_line("ok").await?;
1039        Ok(())
1040    }
1041
1042    async fn cmd_stats(&mut self) -> Result<(), CliError<OUT::Error>> {
1043        let s = self.stats.borrow().clone();
1044        let dropped = *self.events_dropped.borrow();
1045        let depth = self.events.borrow().len();
1046        let mut line: HString<EVENT_LINE_MAX> = HString::new();
1047        let _ = write!(
1048            &mut line,
1049            "rx={} tx={} ack_ok={} ack_timeout={} beacons={} discovered={} \
1050             event_depth={} dropped_events={}",
1051            s.packets_rx,
1052            s.packets_tx,
1053            s.acks_ok,
1054            s.acks_timeout,
1055            s.beacons_rx,
1056            s.nodes_discovered,
1057            depth,
1058            dropped,
1059        );
1060        self.out.write_line(&line).await?;
1061        if let Some(rssi) = s.last_rssi {
1062            let mut line: HString<EVENT_LINE_MAX> = HString::new();
1063            let _ = write!(&mut line, "last_rssi={} last_snr={:?}", rssi, s.last_snr);
1064            self.out.write_line(&line).await?;
1065        }
1066        Ok(())
1067    }
1068
1069    async fn cmd_counters(&mut self) -> Result<(), CliError<OUT::Error>> {
1070        let tx = self.node.frame_counter().await.unwrap_or(0);
1071        let persisted = self.node.persisted_frame_counter().await.unwrap_or(0);
1072        let mut line: HString<EVENT_LINE_MAX> = HString::new();
1073        let _ = write!(&mut line, "local: tx={} (persisted {})", tx, persisted);
1074        self.out.write_line(&line).await?;
1075
1076        // Collect first (callback can't await), then format with alias lookup.
1077        let mut entries: alloc::vec::Vec<(PublicKey, u32, u32)> = alloc::vec::Vec::new();
1078        self.node
1079            .for_each_peer_counter(&mut |pk, last_accepted, persisted_rx| {
1080                entries.push((pk, last_accepted, persisted_rx));
1081            })
1082            .await;
1083
1084        for (pk, last_accepted, persisted_rx) in entries {
1085            let alias = self.peer_alias_display(&pk);
1086            let mut s: HString<EVENT_LINE_MAX> = HString::new();
1087            let _ = write!(
1088                &mut s,
1089                "{:16} rx={} (persisted {})",
1090                alias, last_accepted, persisted_rx
1091            );
1092            self.out.write_line(&s).await?;
1093        }
1094
1095        Ok(())
1096    }
1097
1098    async fn cmd_channels(&mut self) -> Result<(), CliError<OUT::Error>> {
1099        if self.channels.is_empty() {
1100            self.out.write_line("(no channels)").await?;
1101            return Ok(());
1102        }
1103        let mut lines: Vec<String> = Vec::new();
1104        for (_k, entry) in self.channels.iter() {
1105            lines.push(String::from(entry.name.as_str()));
1106        }
1107        for l in lines {
1108            self.out.write_line(&l).await?;
1109        }
1110        Ok(())
1111    }
1112
1113    // ─── Async MAC-I/O commands ──────────────────────────────────────────────
1114
1115    async fn cmd_msg(&mut self, peer: &str, text: &str) -> Result<(), CliError<OUT::Error>> {
1116        let Some(key) = self.resolve_peer(peer) else {
1117            return self.write_err("unknown peer").await;
1118        };
1119        let pc = match self.node.peer(key).await {
1120            Ok(p) => p,
1121            Err(e) => return self.write_err(&node_err_str(&e)).await,
1122        };
1123        let chat = UnicastTextChatWrapper::from_peer(&pc);
1124        let opts = self.send_opts();
1125        match chat.send_text(text, &opts).await {
1126            Ok(_) => {
1127                self.stats.borrow_mut().packets_tx += 1;
1128                self.out.write_line("ok").await?;
1129            }
1130            Err(e) => {
1131                self.write_err(&alloc::format!("{:?}", e)).await?;
1132            }
1133        }
1134        Ok(())
1135    }
1136
1137    async fn cmd_text(&mut self, body: &str) -> Result<(), CliError<OUT::Error>> {
1138        let key = match self.current_peer {
1139            Some(k) => k,
1140            None => {
1141                return self
1142                    .write_err("no current peer — use /query <peer-ref> first")
1143                    .await;
1144            }
1145        };
1146        let pc = match self.node.peer(key).await {
1147            Ok(p) => p,
1148            Err(e) => return self.write_err(&node_err_str(&e)).await,
1149        };
1150        let chat = UnicastTextChatWrapper::from_peer(&pc);
1151        let opts = self.send_opts();
1152        match chat.send_text(body, &opts).await {
1153            Ok(_) => {
1154                self.stats.borrow_mut().packets_tx += 1;
1155            }
1156            Err(e) => {
1157                self.write_err(&alloc::format!("{:?}", e)).await?;
1158            }
1159        }
1160        Ok(())
1161    }
1162
1163    async fn cmd_me(&mut self, action: &str) -> Result<(), CliError<OUT::Error>> {
1164        let mut body: HString<EVENT_LINE_MAX> = HString::new();
1165        let _ = write!(&mut body, "* {}", action);
1166        let owned = String::from(body.as_str());
1167        // Send as a bare text line directed at current peer.
1168        self.cmd_text(&owned).await
1169    }
1170
1171    async fn cmd_ping(
1172        &mut self,
1173        peer: &str,
1174        bytes: Option<u16>,
1175    ) -> Result<(), CliError<OUT::Error>> {
1176        let Some(key) = self.resolve_peer(peer) else {
1177            return self.write_err("unknown peer").await;
1178        };
1179        let pc = match self.node.peer(key).await {
1180            Ok(p) => p,
1181            Err(e) => return self.write_err(&node_err_str(&e)).await,
1182        };
1183        let total = bytes.unwrap_or(8).min(60) as usize;
1184        let extra_bytes = total.saturating_sub(2);
1185        let opts = self.send_opts().with_mic_size(umsh_node::PING_MIC_SIZE);
1186        match pc.ping(extra_bytes, &opts, 30_000).await {
1187            Ok(_) => {
1188                self.stats.borrow_mut().packets_tx += 1;
1189                let alias_str = self.peer_alias_display(&key);
1190                let mut line: HString<EVENT_LINE_MAX> = HString::new();
1191                let _ = write!(&mut line, "ping {} ({} bytes)", alias_str, total.max(2));
1192                self.out.write_line(&line).await?;
1193            }
1194            Err(e) => {
1195                self.write_err(&node_err_str(&e)).await?;
1196            }
1197        }
1198        Ok(())
1199    }
1200
1201    async fn cmd_pfs_start(
1202        &mut self,
1203        peer: &str,
1204        minutes: Option<u16>,
1205    ) -> Result<(), CliError<OUT::Error>> {
1206        #[cfg(feature = "software-crypto")]
1207        {
1208            let Some(key) = self.resolve_peer(peer) else {
1209                return self.write_err("unknown peer").await;
1210            };
1211            let minutes = minutes.unwrap_or(60);
1212            let opts = self.send_opts();
1213            match self.node.request_pfs(&key, minutes, &opts).await {
1214                Ok(_) => self.out.write_line("pfs request sent").await?,
1215                Err(e) => self.write_err(&node_err_str(&e)).await?,
1216            }
1217        }
1218        #[cfg(not(feature = "software-crypto"))]
1219        {
1220            let _ = (peer, minutes);
1221            self.write_err("pfs requires software-crypto feature")
1222                .await?;
1223        }
1224        Ok(())
1225    }
1226
1227    async fn cmd_pfs_end(&mut self, peer: &str) -> Result<(), CliError<OUT::Error>> {
1228        #[cfg(feature = "software-crypto")]
1229        {
1230            let Some(key) = self.resolve_peer(peer) else {
1231                return self.write_err("unknown peer").await;
1232            };
1233            let opts = self.send_opts();
1234            match self.node.end_pfs(&key, &opts).await {
1235                Ok(_) => self.out.write_line("pfs ended").await?,
1236                Err(e) => self.write_err(&node_err_str(&e)).await?,
1237            }
1238        }
1239        #[cfg(not(feature = "software-crypto"))]
1240        {
1241            let _ = peer;
1242            self.write_err("pfs requires software-crypto feature")
1243                .await?;
1244        }
1245        Ok(())
1246    }
1247
1248    async fn cmd_pfs_status(&mut self, peer: Option<&str>) -> Result<(), CliError<OUT::Error>> {
1249        #[cfg(feature = "software-crypto")]
1250        {
1251            let targets: Vec<PublicKey> = match peer {
1252                Some(r) => match self.resolve_peer(r) {
1253                    Some(k) => alloc::vec![k],
1254                    None => return self.write_err("unknown peer").await,
1255                },
1256                None => self.peers.keys().copied().collect(),
1257            };
1258            for key in targets {
1259                // Compute owned values before write_line to avoid borrow conflicts.
1260                let result = self.node.pfs_status(&key).await;
1261                let alias = self.peer_alias_display(&key);
1262                match result {
1263                    Ok(s) => {
1264                        let mut line: HString<EVENT_LINE_MAX> = HString::new();
1265                        let _ = write!(&mut line, "{}: {:?}", alias, s);
1266                        self.out.write_line(&line).await?;
1267                    }
1268                    Err(e) => {
1269                        let msg = node_err_str(&e);
1270                        self.write_err(&msg).await?;
1271                    }
1272                }
1273            }
1274        }
1275        #[cfg(not(feature = "software-crypto"))]
1276        {
1277            let _ = peer;
1278            self.write_err("pfs requires software-crypto feature")
1279                .await?;
1280        }
1281        Ok(())
1282    }
1283
1284    async fn cmd_beacon(&mut self) -> Result<(), CliError<OUT::Error>> {
1285        // A beacon is an empty-payload broadcast. The MAC silently ignores
1286        // the `encrypted` / `ack_requested` flags for broadcasts, so reusing
1287        // `send_opts()` is fine.
1288        use umsh_node::Transport as _;
1289        let opts = self.send_opts();
1290        match self.node.send_all(&[], &opts).await {
1291            Ok(_) => {
1292                self.stats.borrow_mut().packets_tx += 1;
1293                self.out.write_line("beacon sent").await?;
1294            }
1295            Err(e) => self.write_err(&node_err_str(&e)).await?,
1296        }
1297        Ok(())
1298    }
1299
1300    async fn cmd_channel_join(
1301        &mut self,
1302        name: &str,
1303        _key_b58: &str,
1304    ) -> Result<(), CliError<OUT::Error>> {
1305        #[cfg(feature = "software-crypto")]
1306        {
1307            let hname = match HString::<16>::try_from(name) {
1308                Ok(s) => s,
1309                Err(_) => return self.write_err("channel name too long (max 16 chars)").await,
1310            };
1311            if self.channels.contains_key(&hname) {
1312                return self.write_err("already joined").await;
1313            }
1314            // Decode the b58 channel key.
1315            let key_bytes = match umsh_core::base58::decode(_key_b58.as_bytes()) {
1316                Ok(bytes) => bytes,
1317                Err(_) => {
1318                    return self
1319                        .write_err("key must be a 44-character base58 channel key")
1320                        .await;
1321                }
1322            };
1323            let channel_key = umsh_core::ChannelKey(key_bytes);
1324            let channel = umsh_node::Channel::private(channel_key, name);
1325            match self.node.join(&channel).await {
1326                Ok(_) => {
1327                    let entry = ChannelEntry {
1328                        name: HString::try_from(name).unwrap(),
1329                        key_bytes,
1330                    };
1331                    if self.channels.insert(hname, entry).is_err() {
1332                        let _ = self.node.leave(&channel).await;
1333                        return self.write_err("channel table full").await;
1334                    }
1335                    // Persist (best-effort).
1336                    let _ = self
1337                        .channel_store
1338                        .store_channel(name.as_bytes(), &key_bytes)
1339                        .await;
1340                    self.out.write_line("joined").await?;
1341                }
1342                Err(e) => self.write_err(&node_err_str(&e)).await?,
1343            }
1344        }
1345        #[cfg(not(feature = "software-crypto"))]
1346        {
1347            let _ = (name, _key_b58);
1348            self.write_err("channels require software-crypto feature")
1349                .await?;
1350        }
1351        Ok(())
1352    }
1353
1354    async fn cmd_channel_leave(&mut self, name: &str) -> Result<(), CliError<OUT::Error>> {
1355        #[cfg(feature = "software-crypto")]
1356        {
1357            let hname = match HString::<16>::try_from(name) {
1358                Ok(s) => s,
1359                Err(_) => return self.write_err("unknown channel").await,
1360            };
1361            let key_bytes = match self.channels.remove(&hname) {
1362                Some(e) => e.key_bytes,
1363                None => return self.write_err("not joined to that channel").await,
1364            };
1365            let channel_key = umsh_core::ChannelKey(key_bytes);
1366            let channel = umsh_node::Channel::private(channel_key, name);
1367            let _ = self.node.leave(&channel).await;
1368            // Persist deletion (best-effort).
1369            let _ = self.channel_store.delete_channel(name.as_bytes()).await;
1370            self.out.write_line("left").await?;
1371        }
1372        #[cfg(not(feature = "software-crypto"))]
1373        {
1374            let _ = name;
1375            self.write_err("channels require software-crypto feature")
1376                .await?;
1377        }
1378        Ok(())
1379    }
1380
1381    async fn cmd_channel_send(
1382        &mut self,
1383        name: &str,
1384        text: &str,
1385    ) -> Result<(), CliError<OUT::Error>> {
1386        #[cfg(feature = "software-crypto")]
1387        {
1388            let hname = match HString::<16>::try_from(name) {
1389                Ok(s) => s,
1390                Err(_) => return self.write_err("unknown channel").await,
1391            };
1392            let key_bytes = match self.channels.get(&hname) {
1393                Some(e) => e.key_bytes,
1394                None => return self.write_err("not joined to that channel").await,
1395            };
1396            let channel_key = umsh_core::ChannelKey(key_bytes);
1397            let channel = umsh_node::Channel::private(channel_key, name);
1398            let bound = match self.node.bound_channel(&channel) {
1399                Some(b) => b,
1400                None => return self.write_err("channel no longer active").await,
1401            };
1402            let wrapper = umsh_text::MulticastTextChatWrapper::new(bound);
1403            let opts = self.send_opts();
1404            match wrapper.send_text(text, &opts).await {
1405                Ok(_) => {
1406                    self.stats.borrow_mut().packets_tx += 1;
1407                    self.out.write_line("ok").await?;
1408                }
1409                Err(e) => self.write_err(&alloc::format!("{:?}", e)).await?,
1410            }
1411        }
1412        #[cfg(not(feature = "software-crypto"))]
1413        {
1414            let _ = (name, text);
1415            self.write_err("channels require software-crypto feature")
1416                .await?;
1417        }
1418        Ok(())
1419    }
1420
1421    async fn cmd_raw(&mut self, peer: &str, hex: &str) -> Result<(), CliError<OUT::Error>> {
1422        let Some(key) = self.resolve_peer(peer) else {
1423            return self.write_err("unknown peer").await;
1424        };
1425        // Decode hex into a bounded buffer.
1426        let mut bytes: HVec<u8, 128> = HVec::new();
1427        let hex = hex.trim_start_matches("0x");
1428        if hex.len() % 2 != 0 {
1429            return self.write_err("hex must have even length").await;
1430        }
1431        for chunk in hex.as_bytes().chunks(2) {
1432            let hi = hex_nib(chunk[0]);
1433            let lo = hex_nib(chunk[1]);
1434            match (hi, lo) {
1435                (Some(h), Some(l)) => {
1436                    if bytes.push((h << 4) | l).is_err() {
1437                        return self.write_err("hex too long (max 128 bytes)").await;
1438                    }
1439                }
1440                _ => return self.write_err("invalid hex digit").await,
1441            }
1442        }
1443        let pc = match self.node.peer(key).await {
1444            Ok(p) => p,
1445            Err(e) => return self.write_err(&node_err_str(&e)).await,
1446        };
1447        let opts = self.send_opts();
1448        match pc.send(&bytes, &opts).await {
1449            Ok(_) => {
1450                self.stats.borrow_mut().packets_tx += 1;
1451                if self.settings.show_hex {
1452                    let mut line: HString<EVENT_LINE_MAX> = HString::new();
1453                    let _ = write!(&mut line, "sent {} bytes:", bytes.len());
1454                    for b in bytes.iter() {
1455                        let _ = write!(&mut line, " {:02x}", b);
1456                    }
1457                    self.out.write_line(&line).await?;
1458                } else {
1459                    self.out.write_line("ok").await?;
1460                }
1461            }
1462            Err(e) => self.write_err(&node_err_str(&e)).await?,
1463        }
1464        Ok(())
1465    }
1466
1467    // ─── Helpers ─────────────────────────────────────────────────────────────
1468
1469    async fn write_err(&mut self, msg: &str) -> Result<(), CliError<OUT::Error>> {
1470        let mut line: HString<EVENT_LINE_MAX> = HString::new();
1471        let _ = write!(&mut line, "error: {}", msg);
1472        self.out.write_line(&line).await?;
1473        Ok(())
1474    }
1475
1476    fn peer_alias_display(&self, key: &PublicKey) -> String {
1477        if let Some(entry) = self.peers.get(key) {
1478            if let Some(a) = &entry.alias {
1479                return String::from(a.as_str());
1480            }
1481        }
1482        // Fall back to the spec's star-truncated base58 hint rendering.
1483        let mut s = String::new();
1484        let _ = write!(&mut s, "{}", key.hint());
1485        s
1486    }
1487}
1488
1489// ─── Subscription wiring (free function to avoid borrow conflicts) ────────────
1490
1491fn push_event<const N: usize>(
1492    events: &SharedQueue<Deque<CliEvent, N>>,
1493    dropped: &SharedQueue<u64>,
1494    wake: &Rc<AsyncCondition>,
1495    event: CliEvent,
1496) {
1497    let mut q = events.borrow_mut();
1498    if q.push_back(event).is_err() {
1499        *dropped.borrow_mut() += 1;
1500    }
1501    drop(q);
1502    wake.trigger();
1503}
1504
1505fn register_subscriptions<M, const N: usize>(
1506    node: &LocalNode<M>,
1507    events: SharedQueue<Deque<CliEvent, N>>,
1508    dropped: SharedQueue<u64>,
1509    stats: SharedQueue<Stats>,
1510    wake: Rc<AsyncCondition>,
1511) -> Vec<Subscription>
1512where
1513    M: MacBackend,
1514{
1515    let mut subs: Vec<Subscription> = Vec::new();
1516
1517    // on_receive — raw inbound packets
1518    {
1519        let ev = events.clone();
1520        let dr = dropped.clone();
1521        let st = stats.clone();
1522        let wk = wake.clone();
1523        subs.push(node.on_receive(move |pkt| {
1524            let from = match pkt.from_key() {
1525                Some(k) => k,
1526                None => return false,
1527            };
1528            let rssi = pkt.rssi().unwrap_or(0);
1529            let snr = pkt.snr().map(|s| s.as_centibels()).unwrap_or(0);
1530            let hops = pkt.flood_hops().map(|fh| fh.accumulated()).unwrap_or(0);
1531            {
1532                let mut s = st.borrow_mut();
1533                s.packets_rx += 1;
1534                s.last_rssi = pkt.rssi();
1535                s.last_snr = pkt.snr().map(|snr| snr.as_centibels());
1536            }
1537            let mut prefix: HVec<u8, 64> = HVec::new();
1538            let payload = pkt.payload_bytes();
1539            let n = payload.len().min(64);
1540            let _ = prefix.extend_from_slice(&payload[..n]);
1541            let wire = pkt.wire_bytes();
1542            let raw_n = wire.len().min(EVENT_RAW_MAX);
1543            let mut raw_bytes: HVec<u8, EVENT_RAW_MAX> = HVec::new();
1544            let _ = raw_bytes.extend_from_slice(&wire[..raw_n]);
1545            push_event(&ev, &dr, &wk, CliEvent::RawRx { bytes: raw_bytes });
1546            push_event(
1547                &ev,
1548                &dr,
1549                &wk,
1550                CliEvent::Received {
1551                    from,
1552                    hops,
1553                    rssi,
1554                    snr,
1555                    prefix,
1556                },
1557            );
1558            false // don't consume — let other handlers see it too
1559        }));
1560    }
1561
1562    // on_transmitted — raw outbound MAC frames
1563    {
1564        let ev = events.clone();
1565        let dr = dropped.clone();
1566        let wk = wake.clone();
1567        subs.push(node.on_transmitted(move |wire: &[u8]| {
1568            let raw_n = wire.len().min(EVENT_RAW_MAX);
1569            let mut raw_bytes: HVec<u8, EVENT_RAW_MAX> = HVec::new();
1570            let _ = raw_bytes.extend_from_slice(&wire[..raw_n]);
1571            push_event(&ev, &dr, &wk, CliEvent::RawTx { bytes: raw_bytes });
1572        }));
1573    }
1574
1575    // on_ack_received
1576    {
1577        let ev = events.clone();
1578        let dr = dropped.clone();
1579        let st = stats.clone();
1580        let wk = wake.clone();
1581        subs.push(node.on_ack_received(move |peer, _token| {
1582            st.borrow_mut().acks_ok += 1;
1583            push_event(&ev, &dr, &wk, CliEvent::AckReceived { peer });
1584        }));
1585    }
1586
1587    // on_ack_timeout
1588    {
1589        let ev = events.clone();
1590        let dr = dropped.clone();
1591        let st = stats.clone();
1592        let wk = wake.clone();
1593        subs.push(node.on_ack_timeout(move |peer, _token| {
1594            st.borrow_mut().acks_timeout += 1;
1595            push_event(&ev, &dr, &wk, CliEvent::AckTimeout { peer });
1596        }));
1597    }
1598
1599    // on_node_discovered
1600    {
1601        let ev = events.clone();
1602        let dr = dropped.clone();
1603        let st = stats.clone();
1604        let wk = wake.clone();
1605        subs.push(node.on_node_discovered(move |peer, name| {
1606            st.borrow_mut().nodes_discovered += 1;
1607            let name_h = name.and_then(|n| HString::<32>::try_from(n).ok());
1608            push_event(
1609                &ev,
1610                &dr,
1611                &wk,
1612                CliEvent::NodeDiscovered {
1613                    from: peer,
1614                    name: name_h,
1615                },
1616            );
1617        }));
1618    }
1619
1620    // on_beacon
1621    {
1622        let ev = events.clone();
1623        let dr = dropped.clone();
1624        let st = stats.clone();
1625        let wk = wake.clone();
1626        subs.push(node.on_beacon(move |hint, key| {
1627            st.borrow_mut().beacons_rx += 1;
1628            push_event(&ev, &dr, &wk, CliEvent::Beacon { hint, from: key });
1629        }));
1630    }
1631
1632    // on_pfs_established
1633    {
1634        let ev = events.clone();
1635        let dr = dropped.clone();
1636        let wk = wake.clone();
1637        subs.push(node.on_pfs_established(move |peer| {
1638            push_event(&ev, &dr, &wk, CliEvent::PfsEstablished { peer });
1639        }));
1640    }
1641
1642    // on_pfs_ended
1643    {
1644        let ev = events.clone();
1645        let dr = dropped.clone();
1646        let wk = wake.clone();
1647        subs.push(node.on_pfs_ended(move |peer| {
1648            push_event(&ev, &dr, &wk, CliEvent::PfsEnded { peer });
1649        }));
1650    }
1651
1652    // on_pfs_failed
1653    {
1654        let ev = events.clone();
1655        let dr = dropped.clone();
1656        let wk = wake.clone();
1657        subs.push(node.on_pfs_failed(move |peer, reason| {
1658            push_event(&ev, &dr, &wk, CliEvent::PfsFailed { peer, reason });
1659        }));
1660    }
1661
1662    // on_mac_command — EchoRequest and EchoResponse are both silently ignored
1663    // here. EchoRequest is auto-replied by the MAC coordinator. EchoResponse
1664    // is handled by on_pong below (match_pong is called in host.rs before
1665    // dispatch_mac_command fires).
1666    {
1667        let ev = events.clone();
1668        let dr = dropped.clone();
1669        let wk = wake.clone();
1670        subs.push(node.on_mac_command(move |peer, cmd| {
1671            let event = match cmd {
1672                OwnedMacCommand::EchoRequest { .. } | OwnedMacCommand::EchoResponse { .. } => {
1673                    return;
1674                }
1675                other => {
1676                    let cmd_id = mac_cmd_id(other);
1677                    CliEvent::UnknownMacCmdIn { peer, cmd_id }
1678                }
1679            };
1680            push_event(&ev, &dr, &wk, event);
1681        }));
1682    }
1683
1684    // on_pong — fired by node-layer ping tracking when an EchoResponse matches.
1685    {
1686        let ev = events.clone();
1687        let dr = dropped.clone();
1688        let wk = wake.clone();
1689        subs.push(node.on_pong(move |peer, rtt_ms| {
1690            push_event(&ev, &dr, &wk, CliEvent::Pong { peer, rtt_ms });
1691        }));
1692    }
1693
1694    // on_ping_timeout — fired when a pending ping exceeds its deadline.
1695    {
1696        let ev = events.clone();
1697        let dr = dropped.clone();
1698        let wk = wake.clone();
1699        subs.push(node.on_ping_timeout(move |peer| {
1700            push_event(&ev, &dr, &wk, CliEvent::PingTimeout { peer });
1701        }));
1702    }
1703
1704    subs
1705}
1706
1707fn pfs_failure_str(reason: umsh_node::PfsFailure) -> &'static str {
1708    use umsh_node::PfsFailure;
1709    match reason {
1710        PfsFailure::Capacity => "no ephemeral identity slot",
1711        PfsFailure::SessionMissing => "no matching session",
1712        PfsFailure::Crypto => "crypto error",
1713        PfsFailure::Send => "send failed",
1714        PfsFailure::Timeout => "no response (timed out)",
1715        PfsFailure::Other => "error",
1716    }
1717}
1718
1719fn mac_cmd_id(cmd: &OwnedMacCommand) -> u8 {
1720    use umsh_node::CommandId;
1721    match cmd {
1722        OwnedMacCommand::IdentityRequest { .. } => CommandId::IdentityRequest as u8,
1723        OwnedMacCommand::SignalReportRequest => CommandId::SignalReportRequest as u8,
1724        OwnedMacCommand::SignalReportResponse { .. } => CommandId::SignalReportResponse as u8,
1725        OwnedMacCommand::EchoRequest { .. } => CommandId::EchoRequest as u8,
1726        OwnedMacCommand::EchoResponse { .. } => CommandId::EchoResponse as u8,
1727        OwnedMacCommand::PfsSessionRequest { .. } => CommandId::PfsSessionRequest as u8,
1728        OwnedMacCommand::PfsSessionResponse { .. } => CommandId::PfsSessionResponse as u8,
1729        OwnedMacCommand::EndPfsSession => CommandId::EndPfsSession as u8,
1730    }
1731}
1732
1733fn node_err_str<M>(e: &NodeError<M>) -> String
1734where
1735    M: MacBackend,
1736    M::SendError: core::fmt::Debug,
1737    M::CapacityError: core::fmt::Debug,
1738{
1739    alloc::format!("{:?}", e)
1740}
1741
1742fn parse_bool(s: &str) -> Option<bool> {
1743    match s {
1744        "true" | "1" | "yes" | "on" => Some(true),
1745        "false" | "0" | "no" | "off" => Some(false),
1746        _ => None,
1747    }
1748}
1749
1750fn format_parse_error(e: &ParseError) -> String {
1751    match e {
1752        ParseError::Empty => String::from("empty"),
1753        ParseError::UnknownCommand(name) => {
1754            let mut s = String::from("unknown command: ");
1755            s.push_str(name);
1756            s
1757        }
1758        ParseError::MissingArg(n) => {
1759            let mut s = String::from("missing argument: ");
1760            s.push_str(n);
1761            s
1762        }
1763        ParseError::BadNumber => String::from("bad number"),
1764    }
1765}
1766
1767fn hex_nib(b: u8) -> Option<u8> {
1768    match b {
1769        b'0'..=b'9' => Some(b - b'0'),
1770        b'a'..=b'f' => Some(b - b'a' + 10),
1771        b'A'..=b'F' => Some(b - b'A' + 10),
1772        _ => None,
1773    }
1774}