Skip to main content

moteus_protocol/
multiplex.rs

1// Copyright 2026 mjbots Robotic Systems, LLC.  info@mjbots.com
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15//! Multiplex protocol encoding and decoding.
16//!
17//! The moteus multiplex protocol is a register-based protocol that runs over
18//! CAN-FD. It supports reading and writing multiple registers in a single
19//! frame using efficient variable-length encoding.
20
21use crate::frame::CanFdFrame;
22use crate::resolution::Resolution;
23use crate::scaling::{self, saturate_i16, saturate_i32, saturate_i8, Scaling};
24
25/// Write Int8 values
26pub const WRITE_INT8: u8 = 0x00;
27/// Write Int16 values
28pub const WRITE_INT16: u8 = 0x04;
29/// Write Int32 values
30pub const WRITE_INT32: u8 = 0x08;
31/// Write Float values
32pub const WRITE_FLOAT: u8 = 0x0c;
33
34/// Read Int8 values
35pub const READ_INT8: u8 = 0x10;
36/// Read Int16 values
37pub const READ_INT16: u8 = 0x14;
38/// Read Int32 values
39pub const READ_INT32: u8 = 0x18;
40/// Read Float values
41pub const READ_FLOAT: u8 = 0x1c;
42
43/// Reply Int8 values
44pub const REPLY_INT8: u8 = 0x20;
45/// Reply Int16 values
46pub const REPLY_INT16: u8 = 0x24;
47/// Reply Int32 values
48pub const REPLY_INT32: u8 = 0x28;
49/// Reply Float values
50pub const REPLY_FLOAT: u8 = 0x2c;
51
52/// Write error
53pub const WRITE_ERROR: u8 = 0x30;
54/// Read error
55pub const READ_ERROR: u8 = 0x31;
56
57/// Tunneled stream: client to server
58pub const CLIENT_TO_SERVER: u8 = 0x40;
59/// Tunneled stream: server to client
60pub const SERVER_TO_CLIENT: u8 = 0x41;
61/// Tunneled stream: client poll server
62pub const CLIENT_POLL_SERVER: u8 = 0x42;
63/// Tunneled stream: server to client, with flow control packet number
64pub const SERVER_TO_CLIENT_FLOW: u8 = 0x43;
65/// Tunneled stream: client poll server, acknowledging a flow control
66/// packet number
67pub const CLIENT_POLL_SERVER_FLOW: u8 = 0x44;
68
69/// No operation
70pub const NOP: u8 = 0x50;
71
72/// Reads a varuint from a slice, returning the value and the remaining
73/// bytes.
74pub(crate) fn read_varuint_slice(mut data: &[u8]) -> Option<(u32, &[u8])> {
75    let mut value = 0u32;
76    let mut shift = 0;
77    loop {
78        let (byte, rest) = data.split_first()?;
79        data = rest;
80        value |= ((byte & 0x7f) as u32) << shift;
81        if byte & 0x80 == 0 {
82            return Some((value, data));
83        }
84        shift += 7;
85        if shift > 28 {
86            return None;
87        }
88    }
89}
90
91/// Writer for appending data to a CAN-FD frame.
92pub struct WriteCanData<'a> {
93    data: &'a mut [u8; 64],
94    size: &'a mut u8,
95}
96
97impl<'a> core::fmt::Debug for WriteCanData<'a> {
98    fn fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
99        f.debug_struct("WriteCanData")
100            .field("offset", &(*self.size as usize))
101            .field("remaining", &(64 - *self.size as usize))
102            .finish()
103    }
104}
105
106impl<'a> WriteCanData<'a> {
107    /// Creates a new writer from a frame.
108    pub fn new(frame: &'a mut CanFdFrame) -> Self {
109        WriteCanData {
110            data: &mut frame.data,
111            size: &mut frame.size,
112        }
113    }
114
115    /// Creates a writer from raw pointers (for compatibility).
116    pub fn from_parts(data: &'a mut [u8; 64], size: &'a mut u8) -> Self {
117        WriteCanData { data, size }
118    }
119
120    /// Returns the current size of written data.
121    pub fn size(&self) -> u8 {
122        *self.size
123    }
124
125    /// Returns the remaining capacity.
126    pub fn remaining(&self) -> usize {
127        64 - *self.size as usize
128    }
129
130    /// Writes a single byte.
131    pub fn write_u8(&mut self, value: u8) {
132        if (*self.size as usize) < 64 {
133            self.data[*self.size as usize] = value;
134            *self.size += 1;
135        }
136    }
137
138    /// Writes a signed byte.
139    pub fn write_i8(&mut self, value: i8) {
140        self.write_u8(value as u8);
141    }
142
143    /// Writes a little-endian i16.
144    pub fn write_i16(&mut self, value: i16) {
145        let bytes = value.to_le_bytes();
146        self.write_bytes(&bytes);
147    }
148
149    /// Writes a little-endian i32.
150    pub fn write_i32(&mut self, value: i32) {
151        let bytes = value.to_le_bytes();
152        self.write_bytes(&bytes);
153    }
154
155    /// Writes a little-endian f32.
156    pub fn write_f32(&mut self, value: f32) {
157        let bytes = value.to_le_bytes();
158        self.write_bytes(&bytes);
159    }
160
161    /// Writes raw bytes.
162    pub fn write_bytes(&mut self, bytes: &[u8]) {
163        let start = *self.size as usize;
164        let end = start + bytes.len();
165        if end <= 64 {
166            self.data[start..end].copy_from_slice(bytes);
167            *self.size = end as u8;
168        }
169    }
170
171    /// Writes a variable-length unsigned integer (varuint).
172    pub fn write_varuint(&mut self, mut value: u32) {
173        loop {
174            let mut this_byte = (value & 0x7f) as u8;
175            value >>= 7;
176            if value != 0 {
177                this_byte |= 0x80;
178            }
179            self.write_u8(this_byte);
180            if value == 0 {
181                break;
182            }
183        }
184    }
185
186    /// Writes an integer value with the specified resolution.
187    pub fn write_int(&mut self, value: i32, res: Resolution) {
188        match res {
189            Resolution::Int8 => {
190                let clamped = value.clamp(-127, 127) as i8;
191                self.write_i8(clamped);
192            }
193            Resolution::Int16 => {
194                let clamped = value.clamp(-32767, 32767) as i16;
195                self.write_i16(clamped);
196            }
197            Resolution::Int32 => {
198                self.write_i32(value);
199            }
200            Resolution::Float => {
201                self.write_f32(value as f32);
202            }
203            Resolution::Ignore => {}
204        }
205    }
206
207    /// Writes a scaled value with the specified resolution and scaling.
208    pub fn write_mapped(&mut self, value: f32, scaling: &Scaling, res: Resolution) {
209        match res {
210            Resolution::Int8 => {
211                self.write_i8(saturate_i8(value, scaling.int8));
212            }
213            Resolution::Int16 => {
214                self.write_i16(saturate_i16(value, scaling.int16));
215            }
216            Resolution::Int32 => {
217                self.write_i32(saturate_i32(value, scaling.int32));
218            }
219            Resolution::Float => {
220                self.write_f32(value);
221            }
222            Resolution::Ignore => {}
223        }
224    }
225
226    // === Convenience methods for common register types ===
227
228    /// Writes a position value (revolutions).
229    pub fn write_position(&mut self, value: f32, res: Resolution) {
230        self.write_mapped(value, &scaling::POSITION, res);
231    }
232
233    /// Writes a velocity value (revolutions/second).
234    pub fn write_velocity(&mut self, value: f32, res: Resolution) {
235        self.write_mapped(value, &scaling::VELOCITY, res);
236    }
237
238    /// Writes an acceleration value (revolutions/second^2).
239    pub fn write_accel(&mut self, value: f32, res: Resolution) {
240        self.write_mapped(value, &scaling::ACCELERATION, res);
241    }
242
243    /// Writes a torque value (Nm).
244    pub fn write_torque(&mut self, value: f32, res: Resolution) {
245        self.write_mapped(value, &scaling::TORQUE, res);
246    }
247
248    /// Writes a PWM/normalized value (0-1).
249    pub fn write_pwm(&mut self, value: f32, res: Resolution) {
250        self.write_mapped(value, &scaling::PWM, res);
251    }
252
253    /// Writes a voltage value (V).
254    pub fn write_voltage(&mut self, value: f32, res: Resolution) {
255        self.write_mapped(value, &scaling::VOLTAGE, res);
256    }
257
258    /// Writes a temperature value (C).
259    pub fn write_temperature(&mut self, value: f32, res: Resolution) {
260        self.write_mapped(value, &scaling::TEMPERATURE, res);
261    }
262
263    /// Writes a time value (seconds).
264    pub fn write_time(&mut self, value: f32, res: Resolution) {
265        self.write_mapped(value, &scaling::TIME, res);
266    }
267
268    /// Writes a current value (A).
269    pub fn write_current(&mut self, value: f32, res: Resolution) {
270        self.write_mapped(value, &scaling::CURRENT, res);
271    }
272}
273
274/// Combines consecutive register writes of the same resolution for efficiency.
275///
276/// This helper determines how to group registers when encoding them to minimize
277/// the required bytes in the frame. It tracks state and writes framing bytes
278/// when the resolution changes.
279pub struct WriteCombiner<'a> {
280    base_command: u8,
281    start_register: u16,
282    resolutions: &'a [Resolution],
283    current_resolution: Resolution,
284    offset: usize,
285    reply_size: u8,
286}
287
288impl<'a> core::fmt::Debug for WriteCombiner<'a> {
289    fn fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
290        f.debug_struct("WriteCombiner")
291            .field("base_command", &self.base_command)
292            .field("start_register", &self.start_register)
293            .field("offset", &self.offset)
294            .field("current_resolution", &self.current_resolution)
295            .field("reply_size", &self.reply_size)
296            .finish()
297    }
298}
299
300impl<'a> WriteCombiner<'a> {
301    /// Creates a new WriteCombiner.
302    ///
303    /// # Arguments
304    /// * `base_command` - Base command (0x00 for write, 0x10 for read)
305    /// * `start_register` - First register address in the sequence
306    /// * `resolutions` - Resolution for each register in sequence
307    pub fn new(base_command: u8, start_register: u16, resolutions: &'a [Resolution]) -> Self {
308        WriteCombiner {
309            base_command,
310            start_register,
311            resolutions,
312            current_resolution: Resolution::Ignore,
313            offset: 0,
314            reply_size: 0,
315        }
316    }
317
318    /// Returns the expected reply size so far.
319    pub fn reply_size(&self) -> u8 {
320        self.reply_size
321    }
322
323    /// Advances to the next register and writes framing if needed.
324    ///
325    /// Returns `true` if the caller should write the register value.
326    /// Returns `false` if the register should be skipped (Ignore resolution).
327    pub fn maybe_write(&mut self, frame: &mut WriteCanData) -> bool {
328        let this_offset = self.offset;
329        self.offset += 1;
330
331        if this_offset >= self.resolutions.len() {
332            return false;
333        }
334
335        let new_resolution = self.resolutions[this_offset];
336
337        // Same resolution as before - no framing needed
338        if self.current_resolution == new_resolution {
339            return new_resolution != Resolution::Ignore;
340        }
341
342        // Update current resolution
343        self.current_resolution = new_resolution;
344
345        // Ignore means skip this register
346        if new_resolution == Resolution::Ignore {
347            return false;
348        }
349
350        // Count how many consecutive registers have this resolution
351        let mut count = 1i16;
352        for i in (this_offset + 1)..self.resolutions.len() {
353            if self.resolutions[i] == new_resolution {
354                count += 1;
355            } else {
356                break;
357            }
358        }
359
360        // Build the command byte
361        let write_command = self.base_command + new_resolution.type_code();
362
363        let start_size = frame.size();
364
365        if count <= 3 {
366            // Short form: command encodes count in lower 2 bits
367            frame.write_u8(write_command + count as u8);
368        } else {
369            // Long form: count follows command byte
370            frame.write_u8(write_command);
371            frame.write_u8(count as u8);
372        }
373
374        // Write register address
375        frame.write_varuint((self.start_register + this_offset as u16) as u32);
376
377        let end_size = frame.size();
378
379        // Update expected reply size
380        self.reply_size += end_size - start_size;
381        self.reply_size += (count as u8) * (new_resolution.size() as u8);
382
383        true
384    }
385}
386
387/// A decoded register value from the multiplex protocol.
388#[non_exhaustive]
389#[derive(Debug, Clone, Copy)]
390pub enum Value {
391    /// 8-bit signed integer
392    Int8(i8),
393    /// 16-bit signed integer
394    Int16(i16),
395    /// 32-bit signed integer
396    Int32(i32),
397    /// 32-bit IEEE 754 float
398    Float(f32),
399}
400
401impl Value {
402    /// Returns the raw integer value.
403    ///
404    /// Float values are cast to i32.  This should only be used for
405    /// registers which are intended to hold an integer type, and thus
406    /// can not store a non-finite value.
407    pub fn to_i32(&self) -> i32 {
408        match *self {
409            Value::Int8(v) => v as i32,
410            Value::Int16(v) => v as i32,
411            Value::Int32(v) => v,
412            Value::Float(v) => v as i32,
413        }
414    }
415
416    /// Returns the value as f32, applying NaN mapping and scaling.
417    ///
418    /// Integer minimum values (e.g., -128 for Int8) are mapped to NaN.
419    /// Integer values are multiplied by the appropriate scaling factor.
420    /// Float values are returned directly.
421    ///
422    /// # Examples
423    ///
424    /// ```
425    /// use moteus_protocol::Value;
426    /// use moteus_protocol::scaling;
427    ///
428    /// let raw = Value::Int16(5000);
429    /// let position = raw.to_f32(&scaling::POSITION);
430    /// assert!((position - 0.5).abs() < 0.001); // 5000 * 0.0001 = 0.5 rev
431    /// ```
432    pub fn to_f32(&self, scaling: &Scaling) -> f32 {
433        use crate::scaling::{nanify_i16, nanify_i32, nanify_i8};
434
435        match *self {
436            Value::Int8(v) => nanify_i8(v) * scaling.int8,
437            Value::Int16(v) => nanify_i16(v) * scaling.int16,
438            Value::Int32(v) => nanify_i32(v) * scaling.int32,
439            Value::Float(v) => v,
440        }
441    }
442}
443
444/// The type of a parsed subframe.
445#[non_exhaustive]
446#[derive(Debug, Clone, Copy, PartialEq, Eq)]
447pub enum SubframeType {
448    /// Write register values
449    Write,
450    /// Read register request
451    Read,
452    /// Response to a read request
453    Response,
454    /// Write error
455    WriteError,
456    /// Read error
457    ReadError,
458    /// Tunneled stream: client to server
459    StreamClientToServer,
460    /// Tunneled stream: server to client
461    StreamServerToClient,
462    /// Tunneled stream: client poll server
463    StreamClientPollServer,
464    /// Tunneled stream: server to client, with flow control
465    StreamServerToClientFlow,
466    /// Tunneled stream: client poll server, with flow control
467    StreamClientPollServerFlow,
468}
469
470/// A single parsed subframe from a multiplex protocol frame.
471#[non_exhaustive]
472#[derive(Debug, Clone)]
473pub enum Subframe<'a> {
474    /// A register read, write, or response subframe.
475    Register {
476        /// The type of operation
477        subframe_type: SubframeType,
478        /// The register address
479        register: u16,
480        /// The resolution of the value
481        resolution: Resolution,
482        /// The decoded value (None for Read requests)
483        value: Option<Value>,
484    },
485    /// An error subframe.
486    Error {
487        /// The type of error (WriteError or ReadError)
488        subframe_type: SubframeType,
489        /// The register that caused the error
490        register: u16,
491        /// The error code
492        error_code: u16,
493    },
494    /// A tunneled stream subframe.
495    Stream {
496        /// The stream direction
497        subframe_type: SubframeType,
498        /// The stream channel number
499        channel: u16,
500        /// The stream data (borrowed from the input)
501        data: &'a [u8],
502    },
503    /// A tunneled stream subframe with a flow control packet number.
504    StreamFlow {
505        /// The stream direction
506        subframe_type: SubframeType,
507        /// The stream channel number
508        channel: u16,
509        /// The flow control packet number
510        packet_number: u8,
511        /// For `StreamServerToClientFlow`, the stream data; empty for
512        /// `StreamClientPollServerFlow` polls (whose final byte is the
513        /// requested maximum length, not data)
514        data: &'a [u8],
515    },
516}
517
518/// Iterator over subframes in a multiplex protocol frame.
519///
520/// Handles all opcode families: WRITE (0x00-0x0f), READ (0x10-0x1f),
521/// REPLY (0x20-0x2f), ERROR (0x30-0x31), STREAM (0x40-0x44), NOP (0x50).
522///
523/// See [`parse_frame`] for a convenience wrapper.
524pub struct FrameParser<'a> {
525    data: &'a [u8],
526    offset: usize,
527    remaining: u16,
528    current_register: u16,
529    current_resolution: Resolution,
530    current_type: SubframeType,
531}
532
533impl<'a> core::fmt::Debug for FrameParser<'a> {
534    fn fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
535        f.debug_struct("FrameParser")
536            .field("data_len", &self.data.len())
537            .field("offset", &self.offset)
538            .field("remaining", &self.remaining)
539            .field("current_register", &self.current_register)
540            .field("current_resolution", &self.current_resolution)
541            .field("current_type", &self.current_type)
542            .finish()
543    }
544}
545
546impl<'a> FrameParser<'a> {
547    /// Creates a parser from a byte slice.
548    pub fn new(data: &'a [u8]) -> Self {
549        FrameParser {
550            data,
551            offset: 0,
552            remaining: 0,
553            current_register: 0,
554            current_resolution: Resolution::Ignore,
555            current_type: SubframeType::Response,
556        }
557    }
558
559    /// Creates a parser from a frame.
560    pub fn from_frame(frame: &'a CanFdFrame) -> Self {
561        Self::new(&frame.data[..frame.size as usize])
562    }
563
564    /// Reads a varuint from the current position.
565    fn read_varuint(&mut self) -> u16 {
566        let mut result: u16 = 0;
567        let mut shift = 0;
568
569        for _ in 0..5 {
570            if self.offset >= self.data.len() {
571                return result;
572            }
573
574            let byte = self.data[self.offset];
575            self.offset += 1;
576
577            result |= ((byte & 0x7f) as u16) << shift;
578            shift += 7;
579
580            if byte & 0x80 == 0 {
581                return result;
582            }
583        }
584
585        result
586    }
587
588    /// Reads a value at the current offset with the given resolution.
589    fn read_value(&mut self, res: Resolution) -> Option<Value> {
590        match res {
591            Resolution::Int8 => {
592                if self.offset < self.data.len() {
593                    let v = self.data[self.offset] as i8;
594                    self.offset += 1;
595                    Some(Value::Int8(v))
596                } else {
597                    None
598                }
599            }
600            Resolution::Int16 => {
601                if self.offset + 2 <= self.data.len() {
602                    let bytes = [self.data[self.offset], self.data[self.offset + 1]];
603                    self.offset += 2;
604                    Some(Value::Int16(i16::from_le_bytes(bytes)))
605                } else {
606                    None
607                }
608            }
609            Resolution::Int32 => {
610                if self.offset + 4 <= self.data.len() {
611                    let bytes = [
612                        self.data[self.offset],
613                        self.data[self.offset + 1],
614                        self.data[self.offset + 2],
615                        self.data[self.offset + 3],
616                    ];
617                    self.offset += 4;
618                    Some(Value::Int32(i32::from_le_bytes(bytes)))
619                } else {
620                    None
621                }
622            }
623            Resolution::Float => {
624                if self.offset + 4 <= self.data.len() {
625                    let bytes = [
626                        self.data[self.offset],
627                        self.data[self.offset + 1],
628                        self.data[self.offset + 2],
629                        self.data[self.offset + 3],
630                    ];
631                    self.offset += 4;
632                    Some(Value::Float(f32::from_le_bytes(bytes)))
633                } else {
634                    None
635                }
636            }
637            Resolution::Ignore => None,
638        }
639    }
640
641    /// Yields the next register subframe from the current block.
642    fn next_from_block(&mut self) -> Option<Subframe<'a>> {
643        if self.remaining == 0 {
644            return None;
645        }
646
647        self.remaining -= 1;
648        let register = self.current_register;
649        self.current_register += 1;
650
651        // For Read requests, no value data
652        if self.current_type == SubframeType::Read {
653            return Some(Subframe::Register {
654                subframe_type: self.current_type,
655                register,
656                resolution: self.current_resolution,
657                value: None,
658            });
659        }
660
661        // For Write/Response, read the value
662        let value = self.read_value(self.current_resolution);
663        if value.is_none() {
664            // Not enough data, stop iteration
665            self.offset = self.data.len();
666            self.remaining = 0;
667            return None;
668        }
669
670        Some(Subframe::Register {
671            subframe_type: self.current_type,
672            register,
673            resolution: self.current_resolution,
674            value,
675        })
676    }
677
678    /// Parses a register block header (WRITE, READ, or REPLY).
679    fn parse_register_block(&mut self, cmd: u8) -> Option<Subframe<'a>> {
680        let family = cmd & 0xf0;
681        let subframe_type = match family {
682            0x00 => SubframeType::Write,
683            0x10 => SubframeType::Read,
684            0x20 => SubframeType::Response,
685            _ => return None,
686        };
687
688        let resolution = match (cmd >> 2) & 0x03 {
689            0 => Resolution::Int8,
690            1 => Resolution::Int16,
691            2 => Resolution::Int32,
692            3 => Resolution::Float,
693            _ => Resolution::Int8,
694        };
695
696        let mut count = (cmd & 0x03) as u16;
697        if count == 0 {
698            if self.offset >= self.data.len() {
699                return None;
700            }
701            count = self.data[self.offset] as u16;
702            self.offset += 1;
703        }
704
705        if count == 0 {
706            return None;
707        }
708
709        let register = self.read_varuint();
710
711        self.current_type = subframe_type;
712        self.current_resolution = resolution;
713        self.current_register = register + 1;
714        self.remaining = count - 1;
715
716        // For Read requests, no value data
717        if subframe_type == SubframeType::Read {
718            return Some(Subframe::Register {
719                subframe_type,
720                register,
721                resolution,
722                value: None,
723            });
724        }
725
726        // For Write/Response, read the value
727        let value = self.read_value(resolution);
728        if value.is_none() {
729            self.offset = self.data.len();
730            self.remaining = 0;
731            return None;
732        }
733
734        Some(Subframe::Register {
735            subframe_type,
736            register,
737            resolution,
738            value,
739        })
740    }
741
742    /// Parses an error subframe (0x30 or 0x31).
743    fn parse_error(&mut self, cmd: u8) -> Option<Subframe<'a>> {
744        let subframe_type = match cmd {
745            WRITE_ERROR => SubframeType::WriteError,
746            READ_ERROR => SubframeType::ReadError,
747            _ => return None,
748        };
749
750        let register = self.read_varuint();
751        let error_code = self.read_varuint();
752
753        Some(Subframe::Error {
754            subframe_type,
755            register,
756            error_code,
757        })
758    }
759
760    /// Parses a stream subframe (0x40-0x44).
761    fn parse_stream(&mut self, cmd: u8) -> Option<Subframe<'a>> {
762        let subframe_type = match cmd {
763            CLIENT_TO_SERVER => SubframeType::StreamClientToServer,
764            SERVER_TO_CLIENT => SubframeType::StreamServerToClient,
765            CLIENT_POLL_SERVER => SubframeType::StreamClientPollServer,
766            SERVER_TO_CLIENT_FLOW => SubframeType::StreamServerToClientFlow,
767            CLIENT_POLL_SERVER_FLOW => SubframeType::StreamClientPollServerFlow,
768            _ => {
769                // Unknown stream subframe: the operand layout is not
770                // known, so the rest of the frame cannot be parsed.
771                self.offset = self.data.len();
772                return None;
773            }
774        };
775
776        let channel = self.read_varuint();
777
778        // The flow control variants carry a packet number between the
779        // channel and the size/max_length byte.
780        let packet_number = match subframe_type {
781            SubframeType::StreamServerToClientFlow | SubframeType::StreamClientPollServerFlow => {
782                if self.offset >= self.data.len() {
783                    return None;
784                }
785                let value = self.data[self.offset];
786                self.offset += 1;
787                Some(value)
788            }
789            _ => None,
790        };
791
792        let size = self.read_varuint() as usize;
793
794        // Poll subframes carry no data: their final byte is the
795        // requested maximum length.
796        let data = if subframe_type == SubframeType::StreamClientPollServerFlow {
797            &[]
798        } else {
799            if self.offset + size > self.data.len() {
800                self.offset = self.data.len();
801                return None;
802            }
803            let data = &self.data[self.offset..self.offset + size];
804            self.offset += size;
805            data
806        };
807
808        match packet_number {
809            Some(packet_number) => Some(Subframe::StreamFlow {
810                subframe_type,
811                channel,
812                packet_number,
813                data,
814            }),
815            None => Some(Subframe::Stream {
816                subframe_type,
817                channel,
818                data,
819            }),
820        }
821    }
822}
823
824impl<'a> Iterator for FrameParser<'a> {
825    type Item = Subframe<'a>;
826
827    fn next(&mut self) -> Option<Subframe<'a>> {
828        // Continue yielding from current block
829        if self.remaining > 0 {
830            return self.next_from_block();
831        }
832
833        // Look for next command
834        while self.offset < self.data.len() {
835            let cmd = self.data[self.offset];
836            self.offset += 1;
837
838            // NOP: skip silently
839            if cmd == NOP {
840                continue;
841            }
842
843            match cmd & 0xf0 {
844                // WRITE (0x00-0x0f), READ (0x10-0x1f), REPLY (0x20-0x2f)
845                0x00 | 0x10 | 0x20 => {
846                    if let Some(subframe) = self.parse_register_block(cmd) {
847                        return Some(subframe);
848                    }
849                    // parse_register_block returned None (e.g., count=0), try next cmd
850                    continue;
851                }
852                // ERROR (0x30-0x31)
853                0x30 => {
854                    if let Some(subframe) = self.parse_error(cmd) {
855                        return Some(subframe);
856                    }
857                    continue;
858                }
859                // STREAM (0x40-0x44)
860                0x40 => {
861                    if let Some(subframe) = self.parse_stream(cmd) {
862                        return Some(subframe);
863                    }
864                    continue;
865                }
866                // Unknown: stop iteration
867                _ => {
868                    self.offset = self.data.len();
869                    return None;
870                }
871            }
872        }
873
874        None
875    }
876}
877
878/// Creates a parser that iterates over all subframes in a multiplex protocol frame.
879///
880/// # Examples
881///
882/// ```
883/// use moteus_protocol::{parse_frame, Subframe};
884///
885/// // A response frame with mode=10 (position) at register 0
886/// let data = [0x21, 0x00, 0x0a];
887/// for subframe in parse_frame(&data) {
888///     match subframe {
889///         Subframe::Register { register, value, .. } => {
890///             println!("Register 0x{:03x} = {:?}", register, value);
891///         }
892///         _ => {}
893///     }
894/// }
895/// ```
896pub fn parse_frame(data: &[u8]) -> FrameParser<'_> {
897    FrameParser::new(data)
898}
899
900#[cfg(test)]
901mod tests {
902    use super::*;
903
904    #[test]
905    fn test_parse_stream_flow_server_to_client() {
906        // 0x43: channel, packet_number, size, data
907        let frame = [0x43, 0x01, 0x07, 0x02, 0xAA, 0xBB];
908        let mut parser = parse_frame(&frame);
909        match parser.next().unwrap() {
910            Subframe::StreamFlow {
911                subframe_type,
912                channel,
913                packet_number,
914                data,
915            } => {
916                assert_eq!(subframe_type, SubframeType::StreamServerToClientFlow);
917                assert_eq!(channel, 1);
918                assert_eq!(packet_number, 0x07);
919                assert_eq!(data, &[0xAA, 0xBB]);
920            }
921            other => panic!("unexpected subframe: {:?}", other),
922        }
923        assert!(parser.next().is_none());
924    }
925
926    #[test]
927    fn test_parse_stream_flow_poll() {
928        // 0x44: channel, packet_number, max_length — no data follows,
929        // so a register subframe after it must still parse.
930        let frame = [0x44, 0x01, 0x07, 0x30, 0x50, 0x50];
931        let mut parser = parse_frame(&frame);
932        match parser.next().unwrap() {
933            Subframe::StreamFlow {
934                subframe_type,
935                channel,
936                packet_number,
937                data,
938            } => {
939                assert_eq!(subframe_type, SubframeType::StreamClientPollServerFlow);
940                assert_eq!(channel, 1);
941                assert_eq!(packet_number, 0x07);
942                assert!(data.is_empty());
943            }
944            other => panic!("unexpected subframe: {:?}", other),
945        }
946        // The trailing NOPs are skipped cleanly.
947        assert!(parser.next().is_none());
948    }
949
950    #[test]
951    fn test_parse_stream_flow_does_not_misparse_payload() {
952        // Before flow support, the parser treated the bytes after an
953        // unknown 0x43 command as new subframes; ensure the payload is
954        // consumed as data instead.
955        let frame = [
956            0x43, 0x01, 0x07, 0x04, 0x21, 0x00, 0x01, 0x02, // flow data
957            0x41, 0x01, 0x01, 0x99, // plain stream subframe after
958        ];
959        let mut parser = parse_frame(&frame);
960        assert!(matches!(
961            parser.next().unwrap(),
962            Subframe::StreamFlow { data, .. } if data == [0x21, 0x00, 0x01, 0x02]
963        ));
964        assert!(matches!(
965            parser.next().unwrap(),
966            Subframe::Stream { data, .. } if data == [0x99]
967        ));
968        assert!(parser.next().is_none());
969    }
970
971    #[test]
972    fn test_write_varuint() {
973        // Test varuint 0
974        let mut frame = CanFdFrame::new();
975        {
976            let mut writer = WriteCanData::new(&mut frame);
977            writer.write_varuint(0);
978        }
979        assert_eq!(frame.data[0], 0x00);
980        assert_eq!(frame.size, 1);
981
982        // Test varuint 127
983        frame.size = 0;
984        {
985            let mut writer = WriteCanData::new(&mut frame);
986            writer.write_varuint(127);
987        }
988        assert_eq!(frame.data[0], 0x7f);
989
990        // Test varuint 128
991        frame.size = 0;
992        {
993            let mut writer = WriteCanData::new(&mut frame);
994            writer.write_varuint(128);
995        }
996        assert_eq!(frame.data[0], 0x80);
997        assert_eq!(frame.data[1], 0x01);
998        assert_eq!(frame.size, 2);
999    }
1000
1001    #[test]
1002    fn test_write_combiner() {
1003        let mut frame = CanFdFrame::new();
1004        let mut writer = WriteCanData::new(&mut frame);
1005
1006        let resolutions = [Resolution::Float, Resolution::Float, Resolution::Ignore];
1007        let mut combiner = WriteCombiner::new(0x00, 0x020, &resolutions);
1008
1009        // First register should trigger framing
1010        assert!(combiner.maybe_write(&mut writer));
1011        writer.write_f32(1.0);
1012        // Second register same resolution - no framing
1013        assert!(combiner.maybe_write(&mut writer));
1014        writer.write_f32(2.0);
1015        // Third register is Ignore
1016        assert!(!combiner.maybe_write(&mut writer));
1017
1018        // Expected: [0x0e] [0x20] [1.0 as f32 LE] [2.0 as f32 LE]
1019        //   0x0e = Write Float (0x0c) + count 2
1020        //   0x20 = register 0x020 as varuint
1021        assert_eq!(frame.size, 10); // 1 cmd + 1 reg + 4 float + 4 float
1022        assert_eq!(frame.data[0], 0x0e);
1023        assert_eq!(frame.data[1], 0x20);
1024        assert_eq!(
1025            f32::from_le_bytes(frame.data[2..6].try_into().unwrap()),
1026            1.0
1027        );
1028        assert_eq!(
1029            f32::from_le_bytes(frame.data[6..10].try_into().unwrap()),
1030            2.0
1031        );
1032
1033        // Verify reply_size: 2 framing bytes + 2 * 4 value bytes = 10
1034        assert_eq!(combiner.reply_size(), 10);
1035    }
1036
1037    #[test]
1038    fn test_parser_basic() {
1039        // Build a simple reply frame: mode=10 (position mode)
1040        let data = [
1041            0x21, // Reply Int8, count=1
1042            0x00, // Register 0 (MODE)
1043            0x0a, // Value 10 (Position mode)
1044        ];
1045
1046        let mut iter = parse_frame(&data);
1047
1048        let subframe = iter.next().unwrap();
1049        match subframe {
1050            Subframe::Register {
1051                subframe_type,
1052                register,
1053                resolution,
1054                value,
1055            } => {
1056                assert_eq!(subframe_type, SubframeType::Response);
1057                assert_eq!(register, 0);
1058                assert_eq!(resolution, Resolution::Int8);
1059                assert_eq!(value.unwrap().to_i32(), 10);
1060            }
1061            _ => panic!("Expected Register subframe"),
1062        }
1063
1064        assert!(iter.next().is_none());
1065    }
1066
1067    #[test]
1068    fn test_parser_write_subframes() {
1069        // Write Int16, count=2, register 0x20, values 100 and 200
1070        let data = [
1071            0x06, // Write Int16 (0x04), count=2
1072            0x20, // Register 0x20
1073            0x64, 0x00, // 100 as i16 LE
1074            0xc8, 0x00, // 200 as i16 LE
1075        ];
1076
1077        let mut iter = parse_frame(&data);
1078
1079        match iter.next().unwrap() {
1080            Subframe::Register {
1081                subframe_type,
1082                register,
1083                resolution,
1084                value,
1085            } => {
1086                assert_eq!(subframe_type, SubframeType::Write);
1087                assert_eq!(register, 0x20);
1088                assert_eq!(resolution, Resolution::Int16);
1089                assert_eq!(value.unwrap().to_i32(), 100);
1090            }
1091            _ => panic!("Expected Register subframe"),
1092        }
1093
1094        match iter.next().unwrap() {
1095            Subframe::Register {
1096                subframe_type,
1097                register,
1098                value,
1099                ..
1100            } => {
1101                assert_eq!(subframe_type, SubframeType::Write);
1102                assert_eq!(register, 0x21);
1103                assert_eq!(value.unwrap().to_i32(), 200);
1104            }
1105            _ => panic!("Expected Register subframe"),
1106        }
1107
1108        assert!(iter.next().is_none());
1109    }
1110
1111    #[test]
1112    fn test_parser_read_subframes() {
1113        // Read Float, count=1, register 5
1114        let data = [
1115            0x1d, // Read Float (0x1c), count=1
1116            0x05, // Register 5
1117        ];
1118
1119        let mut iter = parse_frame(&data);
1120
1121        match iter.next().unwrap() {
1122            Subframe::Register {
1123                subframe_type,
1124                register,
1125                resolution,
1126                value,
1127            } => {
1128                assert_eq!(subframe_type, SubframeType::Read);
1129                assert_eq!(register, 5);
1130                assert_eq!(resolution, Resolution::Float);
1131                assert!(value.is_none());
1132            }
1133            _ => panic!("Expected Register subframe"),
1134        }
1135
1136        assert!(iter.next().is_none());
1137    }
1138
1139    #[test]
1140    fn test_parser_nop_skipped() {
1141        let data = [
1142            0x50, // NOP
1143            0x21, // Reply Int8, count=1
1144            0x00, // Register 0
1145            0x05, // Value 5
1146        ];
1147
1148        let mut iter = parse_frame(&data);
1149
1150        match iter.next().unwrap() {
1151            Subframe::Register {
1152                register, value, ..
1153            } => {
1154                assert_eq!(register, 0);
1155                assert_eq!(value.unwrap().to_i32(), 5);
1156            }
1157            _ => panic!("Expected Register subframe"),
1158        }
1159
1160        assert!(iter.next().is_none());
1161    }
1162
1163    #[test]
1164    fn test_parser_multiple_blocks() {
1165        // Reply Int8 count=1 reg 0 value 10, then Reply Float count=1 reg 1 value 0.5
1166        let data = [0x21, 0x00, 0x0a, 0x2d, 0x01, 0x00, 0x00, 0x00, 0x3f];
1167
1168        let mut iter = parse_frame(&data);
1169
1170        match iter.next().unwrap() {
1171            Subframe::Register {
1172                register, value, ..
1173            } => {
1174                assert_eq!(register, 0);
1175                assert_eq!(value.unwrap().to_i32(), 10);
1176            }
1177            _ => panic!("Expected Register subframe"),
1178        }
1179
1180        match iter.next().unwrap() {
1181            Subframe::Register {
1182                register, value, ..
1183            } => {
1184                assert_eq!(register, 1);
1185                let val = match value.unwrap() {
1186                    Value::Float(f) => f,
1187                    _ => panic!("Expected Float"),
1188                };
1189                assert!((val - 0.5).abs() < 0.001);
1190            }
1191            _ => panic!("Expected Register subframe"),
1192        }
1193
1194        assert!(iter.next().is_none());
1195    }
1196
1197    #[test]
1198    fn test_value_to_f32() {
1199        // Int8 with position scaling: 50 * 0.01 = 0.5
1200        let v = Value::Int8(50);
1201        assert!((v.to_f32(&scaling::POSITION) - 0.5).abs() < 1e-5);
1202
1203        // Int8 min = NaN
1204        let v = Value::Int8(i8::MIN);
1205        assert!(v.to_f32(&scaling::POSITION).is_nan());
1206
1207        // Float passes through
1208        let v = Value::Float(1.5);
1209        assert!((v.to_f32(&scaling::POSITION) - 1.5).abs() < 1e-5);
1210    }
1211}