1use 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#[non_exhaustive]
28#[derive(Debug, Clone)]
29pub struct ClientConfig {
30 pub chunk_size: u32,
32 pub window_ack_size: u32,
34 pub app: String,
36 pub stream_key: String,
38 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#[non_exhaustive]
56#[derive(Debug, Clone, PartialEq)]
57pub enum ClientEvent {
58 Connected {
60 properties: Vec<(String, Amf0Value)>,
62 },
63 StreamCreated {
65 stream_id: u32,
67 },
68 Publishing,
70 Error {
72 code: String,
74 description: String,
76 },
77 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#[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#[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 #[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 pub fn start(&mut self) -> Vec<u8> {
207 self.handshake.start()
208 }
209
210 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 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 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 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 #[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 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 #[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 let (s0_s1_s2, server_events) = server.handle_data(&c0_c1).unwrap();
665 assert!(server_events.is_empty());
666
667 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 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 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 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 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 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 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}