Skip to main content

rtmp_runtime/
client.rs

1//! RTMP client session engine — `connect` → `createStream` → `publish`
2//! (Adobe RTMP 1.0 §7.2, `NetConnection`/`NetStream` commands).
3//!
4//! [`ClientSession`] mirrors [`crate::server::ServerSession`]: a sans-IO
5//! state machine for the **client publish** role. Feed inbound bytes via
6//! [`ClientSession::handle_data`], get back outbound bytes + typed
7//! [`ClientEvent`]s. Auto-advances: a successful `connect` `_result`
8//! automatically emits `createStream`, and a successful `createStream`
9//! `_result` emits `publish`.
10
11use broadcast_common::{Parse, Serialize};
12
13use crate::RtmpError;
14use crate::amf0::{Amf0Value, Command};
15use crate::chunk::{ChunkAssembler, ChunkWriter, Message};
16use crate::handshake::{
17    EchoPacket, HANDSHAKE_PACKET_LEN, HandshakePacket, RTMP_VERSION, Version, default_random_fill,
18};
19use crate::message::{ProtocolControl, msg_type};
20
21type Result<T> = core::result::Result<T, RtmpError>;
22
23const COMMAND_CHUNK_STREAM_ID: u32 = 3;
24const VERSION_LEN: usize = 1;
25
26/// Configuration for a [`ClientSession`].
27#[non_exhaustive]
28#[derive(Debug, Clone)]
29pub struct ClientConfig {
30    /// Outbound chunk size (§5.4.1).
31    pub chunk_size: u32,
32    /// Window Acknowledgement Size (§5.4.4).
33    pub window_ack_size: u32,
34    /// RTMP `app` name for the `connect` command.
35    pub app: String,
36    /// Publishing name (stream key) for the `publish` command.
37    pub stream_key: String,
38    /// `tcUrl` for the `connect` command (e.g. `rtmp://host/app`).
39    pub tc_url: Option<String>,
40}
41
42impl Default for ClientConfig {
43    fn default() -> Self {
44        Self {
45            chunk_size: 4096,
46            window_ack_size: 2_500_000,
47            app: String::new(),
48            stream_key: String::new(),
49            tc_url: None,
50        }
51    }
52}
53
54/// Typed events [`ClientSession::handle_data`] surfaces.
55#[non_exhaustive]
56#[derive(Debug, Clone, PartialEq)]
57pub enum ClientEvent {
58    /// Server accepted `connect`.
59    Connected {
60        /// Properties from `_result`.
61        properties: Vec<(String, Amf0Value)>,
62    },
63    /// Server allocated a stream id.
64    StreamCreated {
65        /// The allocated message stream id.
66        stream_id: u32,
67    },
68    /// Server accepted `publish` — session is now publishing.
69    Publishing,
70    /// Server returned an error.
71    Error {
72        /// Status code.
73        code: String,
74        /// Description.
75        description: String,
76    },
77    /// Session closed.
78    Closed,
79}
80
81#[derive(Debug, Clone, Copy, PartialEq)]
82enum ClientState {
83    Init,
84    HandshakeDone,
85    ConnectSent { txn_id: f64 },
86    Connected,
87    CreateStreamSent { txn_id: f64 },
88    StreamCreated,
89    PublishSent,
90    Publishing,
91}
92
93/// Client-side handshake: generate C0+C1, consume S0+S1+S2, produce C2.
94#[derive(Debug)]
95struct ClientHandshake {
96    state: ClientHandshakeState,
97    local_time: u32,
98    local_random: [u8; 1528],
99}
100
101#[derive(Debug)]
102enum ClientHandshakeState {
103    SendC0C1,
104    WaitS0S1S2,
105    Done,
106}
107
108impl ClientHandshake {
109    fn new() -> Self {
110        Self {
111            state: ClientHandshakeState::SendC0C1,
112            local_time: 0,
113            local_random: default_random_fill(),
114        }
115    }
116
117    fn start(&mut self) -> Vec<u8> {
118        let c0 = Version(RTMP_VERSION);
119        let c1 = HandshakePacket {
120            time: self.local_time,
121            zero: 0,
122            random: self.local_random,
123        };
124        let mut out = vec![0u8; VERSION_LEN + HANDSHAKE_PACKET_LEN];
125        c0.serialize_into(&mut out[..VERSION_LEN]).unwrap();
126        c1.serialize_into(&mut out[VERSION_LEN..]).unwrap();
127        self.state = ClientHandshakeState::WaitS0S1S2;
128        out
129    }
130
131    fn read(&mut self, input: &[u8]) -> Result<(Vec<u8>, usize, bool)> {
132        match &self.state {
133            ClientHandshakeState::SendC0C1 => Ok((Vec::new(), 0, false)),
134            ClientHandshakeState::WaitS0S1S2 => {
135                let need = VERSION_LEN + HANDSHAKE_PACKET_LEN + HANDSHAKE_PACKET_LEN;
136                if input.len() < need {
137                    return Err(RtmpError::BufferTooShort {
138                        need,
139                        have: input.len(),
140                        what: "S0+S1+S2",
141                    });
142                }
143                let _s0 = Version::parse(&input[..VERSION_LEN])?;
144                let s1 = HandshakePacket::parse(
145                    &input[VERSION_LEN..VERSION_LEN + HANDSHAKE_PACKET_LEN],
146                )?;
147                let _s2 = EchoPacket::parse(&input[VERSION_LEN + HANDSHAKE_PACKET_LEN..need])?;
148
149                let c2 = EchoPacket {
150                    time: s1.time,
151                    time2: 0,
152                    random_echo: s1.random,
153                };
154                let mut reply = vec![0u8; HANDSHAKE_PACKET_LEN];
155                c2.serialize_into(&mut reply)?;
156
157                self.state = ClientHandshakeState::Done;
158                Ok((reply, need, true))
159            }
160            ClientHandshakeState::Done => Ok((Vec::new(), 0, true)),
161        }
162    }
163
164    fn is_done(&self) -> bool {
165        matches!(self.state, ClientHandshakeState::Done)
166    }
167}
168
169/// Sans-IO RTMP client publish session.
170#[derive(Debug)]
171pub struct ClientSession {
172    config: ClientConfig,
173    handshake: ClientHandshake,
174    handshake_buf: Vec<u8>,
175    assembler: ChunkAssembler,
176    writer: ChunkWriter,
177    state: ClientState,
178    stream_id: Option<u32>,
179    next_txn_id: f64,
180    ack_threshold: u32,
181    bytes_received: u64,
182    bytes_acked: u64,
183}
184
185impl ClientSession {
186    /// Create a new client session.
187    #[must_use]
188    pub fn new(config: ClientConfig) -> Self {
189        let ack_threshold = config.window_ack_size;
190        Self {
191            config,
192            handshake: ClientHandshake::new(),
193            handshake_buf: Vec::new(),
194            assembler: ChunkAssembler::new(),
195            writer: ChunkWriter::new(),
196            state: ClientState::Init,
197            stream_id: None,
198            next_txn_id: 1.0,
199            ack_threshold,
200            bytes_received: 0,
201            bytes_acked: 0,
202        }
203    }
204
205    /// Produce C0+C1 handshake bytes. Call once at session start.
206    pub fn start(&mut self) -> Vec<u8> {
207        self.handshake.start()
208    }
209
210    /// Feed inbound bytes. Returns `(outbound bytes, events)`.
211    pub fn handle_data(&mut self, input: &[u8]) -> Result<(Vec<u8>, Vec<ClientEvent>)> {
212        let mut out = Vec::new();
213        let mut events = Vec::new();
214
215        let chunk_input = match self.drive_handshake(input, &mut out, &mut events)? {
216            Some(bytes) => bytes,
217            None => return Ok((out, events)),
218        };
219
220        self.bytes_received = self.bytes_received.saturating_add(chunk_input.len() as u64);
221
222        self.assembler.feed(chunk_input);
223        while let Some(msg) = self.assembler.next_message()? {
224            self.dispatch_message(&msg, &mut out, &mut events)?;
225        }
226
227        self.maybe_ack(&mut out);
228        Ok((out, events))
229    }
230
231    /// Encode an audio message (type 8). Only valid in `Publishing` state.
232    pub fn send_audio(&mut self, timestamp: u32, data: &[u8]) -> Result<Vec<u8>> {
233        if self.state != ClientState::Publishing {
234            return Err(RtmpError::Malformed {
235                what: "send_audio: not in Publishing state",
236            });
237        }
238        let msg = Message {
239            chunk_stream_id: 4,
240            timestamp,
241            message_type_id: msg_type::AUDIO,
242            message_stream_id: self.stream_id.unwrap_or(1),
243            payload: data.to_vec(),
244        };
245        Ok(self.writer.write(&msg))
246    }
247
248    /// Encode a video message (type 9). Only valid in `Publishing` state.
249    pub fn send_video(&mut self, timestamp: u32, data: &[u8]) -> Result<Vec<u8>> {
250        if self.state != ClientState::Publishing {
251            return Err(RtmpError::Malformed {
252                what: "send_video: not in Publishing state",
253            });
254        }
255        let msg = Message {
256            chunk_stream_id: 5,
257            timestamp,
258            message_type_id: msg_type::VIDEO,
259            message_stream_id: self.stream_id.unwrap_or(1),
260            payload: data.to_vec(),
261        };
262        Ok(self.writer.write(&msg))
263    }
264
265    /// Encode a metadata message (`@setDataFrame`/`onMetaData`, type 18).
266    pub fn send_metadata(&mut self, metadata: &[(String, Amf0Value)]) -> Result<Vec<u8>> {
267        if self.state != ClientState::Publishing {
268            return Err(RtmpError::Malformed {
269                what: "send_metadata: not in Publishing state",
270            });
271        }
272        let mut payload = Amf0Value::String("@setDataFrame".to_string()).to_bytes();
273        payload.extend(Amf0Value::String("onMetaData".to_string()).to_bytes());
274        let obj = Amf0Value::Object(metadata.to_vec());
275        payload.extend(obj.to_bytes());
276        let msg = Message {
277            chunk_stream_id: 4,
278            timestamp: 0,
279            message_type_id: msg_type::DATA_AMF0,
280            message_stream_id: self.stream_id.unwrap_or(1),
281            payload,
282        };
283        Ok(self.writer.write(&msg))
284    }
285
286    /// Whether the session has reached `Publishing` state.
287    #[must_use]
288    pub fn is_publishing(&self) -> bool {
289        self.state == ClientState::Publishing
290    }
291
292    fn drive_handshake<'a>(
293        &mut self,
294        input: &'a [u8],
295        out: &mut Vec<u8>,
296        events: &mut Vec<ClientEvent>,
297    ) -> Result<Option<&'a [u8]>> {
298        if self.handshake.is_done() && self.state != ClientState::Init {
299            return Ok(Some(input));
300        }
301
302        self.handshake_buf.extend_from_slice(input);
303
304        if !self.handshake.is_done() {
305            match self.handshake.read(&self.handshake_buf) {
306                Ok((reply, consumed, done)) => {
307                    out.extend_from_slice(&reply);
308                    self.handshake_buf.drain(..consumed);
309                    if done {
310                        self.state = ClientState::HandshakeDone;
311                        out.extend_from_slice(&self.build_connect());
312                    }
313                }
314                Err(RtmpError::BufferTooShort { .. }) => {
315                    return Ok(None);
316                }
317                Err(e) => return Err(e),
318            }
319        }
320
321        if self.handshake.is_done() && !self.handshake_buf.is_empty() {
322            let remaining = std::mem::take(&mut self.handshake_buf);
323            self.bytes_received = self.bytes_received.saturating_add(remaining.len() as u64);
324            self.assembler.feed(&remaining);
325            while let Some(msg) = self.assembler.next_message()? {
326                self.dispatch_message(&msg, out, events)?;
327            }
328        }
329
330        Ok(None)
331    }
332
333    fn build_connect(&mut self) -> Vec<u8> {
334        let txn_id = self.next_txn_id;
335        self.next_txn_id += 1.0;
336        self.state = ClientState::ConnectSent { txn_id };
337
338        let tc_url = self
339            .config
340            .tc_url
341            .clone()
342            .unwrap_or_else(|| format!("rtmp://localhost/{}", self.config.app));
343
344        let cmd = Command {
345            name: "connect".to_string(),
346            transaction_id: txn_id,
347            arguments: vec![Amf0Value::Object(vec![
348                (
349                    "app".to_string(),
350                    Amf0Value::String(self.config.app.clone()),
351                ),
352                ("tcUrl".to_string(), Amf0Value::String(tc_url)),
353                (
354                    "type".to_string(),
355                    Amf0Value::String("nonprivate".to_string()),
356                ),
357            ])],
358        };
359
360        let mut out = Vec::new();
361
362        let set_chunk_size = ProtocolControl::SetChunkSize(self.config.chunk_size);
363        out.extend_from_slice(&self.writer.write(&set_chunk_size.to_message()));
364        self.writer.set_chunk_size(self.config.chunk_size);
365
366        let window_ack = ProtocolControl::WindowAckSize(self.config.window_ack_size);
367        out.extend_from_slice(&self.writer.write(&window_ack.to_message()));
368
369        let msg = Message {
370            chunk_stream_id: COMMAND_CHUNK_STREAM_ID,
371            timestamp: 0,
372            message_type_id: msg_type::COMMAND_AMF0,
373            message_stream_id: 0,
374            payload: cmd.to_body(),
375        };
376        out.extend_from_slice(&self.writer.write(&msg));
377        out
378    }
379
380    fn build_create_stream(&mut self) -> Vec<u8> {
381        let txn_id = self.next_txn_id;
382        self.next_txn_id += 1.0;
383        self.state = ClientState::CreateStreamSent { txn_id };
384
385        let cmd = Command {
386            name: "createStream".to_string(),
387            transaction_id: txn_id,
388            arguments: vec![Amf0Value::Null],
389        };
390        let msg = Message {
391            chunk_stream_id: COMMAND_CHUNK_STREAM_ID,
392            timestamp: 0,
393            message_type_id: msg_type::COMMAND_AMF0,
394            message_stream_id: 0,
395            payload: cmd.to_body(),
396        };
397        self.writer.write(&msg)
398    }
399
400    fn build_publish(&mut self) -> Vec<u8> {
401        self.state = ClientState::PublishSent;
402
403        let cmd = Command {
404            name: "publish".to_string(),
405            transaction_id: 0.0,
406            arguments: vec![
407                Amf0Value::Null,
408                Amf0Value::String(self.config.stream_key.clone()),
409                Amf0Value::String("live".to_string()),
410            ],
411        };
412        let msg = Message {
413            chunk_stream_id: COMMAND_CHUNK_STREAM_ID,
414            timestamp: 0,
415            message_type_id: msg_type::COMMAND_AMF0,
416            message_stream_id: self.stream_id.unwrap_or(1),
417            payload: cmd.to_body(),
418        };
419        self.writer.write(&msg)
420    }
421
422    fn dispatch_message(
423        &mut self,
424        msg: &Message,
425        out: &mut Vec<u8>,
426        events: &mut Vec<ClientEvent>,
427    ) -> Result<()> {
428        if let Some(pc) = ProtocolControl::from_message(msg)? {
429            match pc {
430                ProtocolControl::SetChunkSize(size) => {
431                    self.assembler.set_chunk_size(size);
432                }
433                ProtocolControl::WindowAckSize(size) => {
434                    self.ack_threshold = size;
435                }
436                ProtocolControl::SetPeerBandwidth {
437                    ack_window_size, ..
438                } => {
439                    self.ack_threshold = ack_window_size;
440                    let ack = ProtocolControl::WindowAckSize(ack_window_size);
441                    out.extend_from_slice(&self.writer.write(&ack.to_message()));
442                }
443                _ => {}
444            }
445            return Ok(());
446        }
447        if msg.message_type_id == msg_type::COMMAND_AMF0 {
448            let cmd = Command::parse(&msg.payload)?;
449            self.handle_command(&cmd, out, events)?;
450        }
451        Ok(())
452    }
453
454    fn handle_command(
455        &mut self,
456        cmd: &Command,
457        out: &mut Vec<u8>,
458        events: &mut Vec<ClientEvent>,
459    ) -> Result<()> {
460        match cmd.name.as_str() {
461            "_result" => self.handle_result(cmd, out, events),
462            "_error" => {
463                let (code, desc) = extract_status(&cmd.arguments);
464                events.push(ClientEvent::Error {
465                    code,
466                    description: desc,
467                });
468                Ok(())
469            }
470            "onStatus" => self.handle_on_status(cmd, events),
471            _ => Ok(()),
472        }
473    }
474
475    fn handle_result(
476        &mut self,
477        cmd: &Command,
478        out: &mut Vec<u8>,
479        events: &mut Vec<ClientEvent>,
480    ) -> Result<()> {
481        match self.state {
482            ClientState::ConnectSent { txn_id } if (cmd.transaction_id - txn_id).abs() < 0.5 => {
483                let props = match cmd.arguments.first() {
484                    Some(Amf0Value::Object(pairs)) => pairs.clone(),
485                    _ => Vec::new(),
486                };
487                events.push(ClientEvent::Connected { properties: props });
488                self.state = ClientState::Connected;
489                out.extend_from_slice(&self.build_create_stream());
490            }
491            ClientState::CreateStreamSent { txn_id }
492                if (cmd.transaction_id - txn_id).abs() < 0.5 =>
493            {
494                let sid = match cmd.arguments.last() {
495                    Some(Amf0Value::Number(n)) => *n as u32,
496                    _ => 1,
497                };
498                self.stream_id = Some(sid);
499                events.push(ClientEvent::StreamCreated { stream_id: sid });
500                self.state = ClientState::StreamCreated;
501                out.extend_from_slice(&self.build_publish());
502            }
503            _ => {}
504        }
505        Ok(())
506    }
507
508    fn handle_on_status(&mut self, cmd: &Command, events: &mut Vec<ClientEvent>) -> Result<()> {
509        let (code, desc) = extract_status(&cmd.arguments);
510        if code == "NetStream.Publish.Start" {
511            self.state = ClientState::Publishing;
512            events.push(ClientEvent::Publishing);
513        } else {
514            events.push(ClientEvent::Error {
515                code,
516                description: desc,
517            });
518        }
519        Ok(())
520    }
521
522    fn maybe_ack(&mut self, out: &mut Vec<u8>) {
523        if self.ack_threshold == 0 {
524            return;
525        }
526        if self.bytes_received.wrapping_sub(self.bytes_acked) >= u64::from(self.ack_threshold) {
527            let seq = (self.bytes_received & 0xFFFF_FFFF) as u32;
528            let ack = ProtocolControl::Acknowledgement(seq);
529            out.extend_from_slice(&self.writer.write(&ack.to_message()));
530            self.bytes_acked = self.bytes_received;
531        }
532    }
533}
534
535fn extract_status(args: &[Amf0Value]) -> (String, String) {
536    for arg in args {
537        if let Amf0Value::Object(pairs) = arg {
538            let code = pairs
539                .iter()
540                .find(|(k, _)| k == "code")
541                .and_then(|(_, v)| match v {
542                    Amf0Value::String(s) => Some(s.clone()),
543                    _ => None,
544                })
545                .unwrap_or_default();
546            let desc = pairs
547                .iter()
548                .find(|(k, _)| k == "description")
549                .and_then(|(_, v)| match v {
550                    Amf0Value::String(s) => Some(s.clone()),
551                    _ => None,
552                })
553                .unwrap_or_default();
554            if !code.is_empty() {
555                return (code, desc);
556            }
557        }
558    }
559    (String::new(), String::new())
560}
561
562#[cfg(test)]
563mod tests {
564    use super::*;
565    use crate::handshake;
566    use crate::server::{ServerConfig, ServerSession};
567
568    #[test]
569    fn client_handshake_round_trip() {
570        let mut client_hs = ClientHandshake::new();
571        let mut server_hs = handshake::Handshake::new();
572
573        let c0_c1 = client_hs.start();
574        assert!(!client_hs.is_done());
575
576        let (s0_s1_s2, consumed, _done) = server_hs.read(&c0_c1).unwrap();
577        assert_eq!(consumed, c0_c1.len());
578
579        let (c2, consumed2, done2) = client_hs.read(&s0_s1_s2).unwrap();
580        assert_eq!(consumed2, s0_s1_s2.len());
581        assert!(done2);
582        assert!(client_hs.is_done());
583        assert_eq!(c2.len(), HANDSHAKE_PACKET_LEN);
584
585        let (_reply, consumed3, done3) = server_hs.read(&c2).unwrap();
586        assert_eq!(consumed3, HANDSHAKE_PACKET_LEN);
587        assert!(done3);
588    }
589
590    #[test]
591    fn client_connect_flow() {
592        let config = ClientConfig {
593            app: "live".to_string(),
594            stream_key: "test_key".to_string(),
595            ..ClientConfig::default()
596        };
597        let mut client = ClientSession::new(config);
598        let c0_c1 = client.start();
599
600        let mut server = ServerSession::new(
601            ServerConfig::default().with_expected_stream_key(Some("test_key".to_string())),
602        );
603
604        let (s0_s1_s2, server_events) = server.handle_data(&c0_c1).unwrap();
605        assert!(server_events.is_empty());
606
607        let (client_out, _client_events) = client.handle_data(&s0_s1_s2).unwrap();
608        assert!(
609            !client_out.is_empty(),
610            "client should emit C2 + connect + protocol control"
611        );
612
613        let (server_reply, server_events) = server.handle_data(&client_out).unwrap();
614        assert!(
615            server_events
616                .iter()
617                .any(|e| matches!(e, crate::server::ServerEvent::Connected { .. })),
618            "server should see connect: {server_events:?}"
619        );
620
621        let (client_out2, client_events2) = client.handle_data(&server_reply).unwrap();
622        assert!(
623            client_events2
624                .iter()
625                .any(|e| matches!(e, ClientEvent::Connected { .. })),
626            "client should see Connected: {client_events2:?}"
627        );
628        // `client_out2` alone is not a reliable auto-advance signal: the
629        // server's connect reply also carries a `SetPeerBandwidth` protocol
630        // control message, which the client acks unconditionally — so bytes
631        // come back even if `createStream` itself never fires. Check the
632        // actual state transition instead (issue #933 finding M2).
633        let _ = client_out2;
634        assert!(
635            matches!(client.state, ClientState::CreateStreamSent { .. }),
636            "client should auto-advance from Connected to CreateStreamSent: {:?}",
637            client.state
638        );
639    }
640
641    /// Drives the full client auto-advance —
642    /// `connect` -> `createStream` -> `publish` — end to end against a real
643    /// [`ServerSession`], and asserts every state transition actually
644    /// happens (issue #933 finding M2: the prior version of this test
645    /// buried its post-`Connected` assertions inside
646    /// `if !client_out2.is_empty() { .. }` guards, so a broken auto-advance
647    /// that emitted no bytes would make the test vacuously pass instead of
648    /// fail).
649    #[test]
650    fn client_full_publish_flow() {
651        let config = ClientConfig {
652            app: "live".to_string(),
653            stream_key: "test_key".to_string(),
654            ..ClientConfig::default()
655        };
656        let mut client = ClientSession::new(config);
657        let c0_c1 = client.start();
658
659        let mut server = ServerSession::new(
660            ServerConfig::default().with_expected_stream_key(Some("test_key".to_string())),
661        );
662
663        // Handshake.
664        let (s0_s1_s2, server_events) = server.handle_data(&c0_c1).unwrap();
665        assert!(server_events.is_empty());
666
667        // Client completes the handshake and sends C2 + connect.
668        let (client_out, client_events) = client.handle_data(&s0_s1_s2).unwrap();
669        assert!(
670            !client_out.is_empty(),
671            "client should emit C2 + connect + protocol control"
672        );
673        assert!(
674            client_events.is_empty(),
675            "no client events expected before the server replies: {client_events:?}"
676        );
677        assert_eq!(client.state, ClientState::ConnectSent { txn_id: 1.0 });
678
679        // Server completes the handshake and accepts connect.
680        let (server_reply, server_events) = server.handle_data(&client_out).unwrap();
681        assert!(
682            server_events
683                .iter()
684                .any(|e| matches!(e, crate::server::ServerEvent::Connected { .. })),
685            "server should see connect: {server_events:?}"
686        );
687
688        // Client sees Connected and auto-advances: emits createStream.
689        let (client_out2, client_events2) = client.handle_data(&server_reply).unwrap();
690        assert!(
691            client_events2
692                .iter()
693                .any(|e| matches!(e, ClientEvent::Connected { .. })),
694            "client should see Connected: {client_events2:?}"
695        );
696        assert!(
697            !client_out2.is_empty(),
698            "client should auto-advance from Connected and emit createStream"
699        );
700        assert!(matches!(client.state, ClientState::CreateStreamSent { .. }));
701
702        // Server replies to createStream with an allocated stream id.
703        let (server_reply2, _server_events2) = server.handle_data(&client_out2).unwrap();
704        assert!(
705            !server_reply2.is_empty(),
706            "server should reply to createStream"
707        );
708
709        // Client sees StreamCreated and auto-advances: emits publish.
710        let (client_out3, client_events3) = client.handle_data(&server_reply2).unwrap();
711        assert!(
712            client_events3
713                .iter()
714                .any(|e| matches!(e, ClientEvent::StreamCreated { .. })),
715            "client should see StreamCreated: {client_events3:?}"
716        );
717        assert!(
718            !client_out3.is_empty(),
719            "client should auto-advance from StreamCreated and emit publish"
720        );
721        assert_eq!(client.state, ClientState::PublishSent);
722
723        // Server accepts publish.
724        let (server_reply3, server_events3) = server.handle_data(&client_out3).unwrap();
725        assert!(
726            server_events3
727                .iter()
728                .any(|e| matches!(e, crate::server::ServerEvent::Publish { .. })),
729            "server should see publish: {server_events3:?}"
730        );
731        assert!(
732            !server_reply3.is_empty(),
733            "server should reply onStatus for publish"
734        );
735
736        // Client sees Publishing — the auto-advance has run to completion.
737        let (_client_out4, client_events4) = client.handle_data(&server_reply3).unwrap();
738        assert!(
739            client_events4
740                .iter()
741                .any(|e| matches!(e, ClientEvent::Publishing)),
742            "client should see Publishing: {client_events4:?}"
743        );
744        assert_eq!(client.state, ClientState::Publishing);
745        assert!(client.is_publishing());
746    }
747
748    #[test]
749    fn send_audio_video_while_publishing() {
750        let config = ClientConfig {
751            app: "live".to_string(),
752            stream_key: "test".to_string(),
753            ..ClientConfig::default()
754        };
755        let mut client = ClientSession::new(config);
756        client.state = ClientState::Publishing;
757        client.stream_id = Some(1);
758
759        let audio = client.send_audio(100, &[0xAA, 0xBB]).unwrap();
760        assert!(!audio.is_empty());
761
762        let video = client.send_video(200, &[0xCC, 0xDD, 0xEE]).unwrap();
763        assert!(!video.is_empty());
764
765        let meta = client
766            .send_metadata(&[
767                ("width".to_string(), Amf0Value::Number(1920.0)),
768                ("height".to_string(), Amf0Value::Number(1080.0)),
769            ])
770            .unwrap();
771        assert!(!meta.is_empty());
772    }
773
774    #[test]
775    fn send_audio_before_publishing_fails() {
776        let config = ClientConfig {
777            app: "live".to_string(),
778            stream_key: "test".to_string(),
779            ..ClientConfig::default()
780        };
781        let mut client = ClientSession::new(config);
782        assert!(client.send_audio(0, &[0x00]).is_err());
783    }
784
785    #[test]
786    fn connect_error_produces_event() {
787        let config = ClientConfig {
788            app: "live".to_string(),
789            stream_key: "test".to_string(),
790            ..ClientConfig::default()
791        };
792        let mut client = ClientSession::new(config);
793        client.handshake.state = ClientHandshakeState::Done;
794        client.state = ClientState::ConnectSent { txn_id: 1.0 };
795
796        let error_cmd = Command {
797            name: "_error".to_string(),
798            transaction_id: 1.0,
799            arguments: vec![
800                Amf0Value::Null,
801                Amf0Value::Object(vec![
802                    (
803                        "code".to_string(),
804                        Amf0Value::String("NetConnection.Connect.Rejected".to_string()),
805                    ),
806                    (
807                        "description".to_string(),
808                        Amf0Value::String("Connection refused".to_string()),
809                    ),
810                ]),
811            ],
812        };
813        let msg = Message {
814            chunk_stream_id: COMMAND_CHUNK_STREAM_ID,
815            timestamp: 0,
816            message_type_id: msg_type::COMMAND_AMF0,
817            message_stream_id: 0,
818            payload: error_cmd.to_body(),
819        };
820        let bytes = ChunkWriter::new().write(&msg);
821
822        let (_, events) = client.handle_data(&bytes).unwrap();
823        assert!(
824            events.iter().any(|e| matches!(
825                e,
826                ClientEvent::Error { code, .. } if code == "NetConnection.Connect.Rejected"
827            )),
828            "expected Error event: {events:?}"
829        );
830    }
831}