diff --git a/synctv-core/testing/src/streaming.rs b/synctv-core/testing/src/streaming.rs index b13485f6..e51ca59c 100644 --- a/synctv-core/testing/src/streaming.rs +++ b/synctv-core/testing/src/streaming.rs @@ -9,6 +9,7 @@ use synctv_xiu::{ bytes_writer::AsyncBytesWriter, net_io::{TNetIO, TcpIO}, }, + flv::amf0::Amf0ValueType, rtmp::{ chunk::{ errors::UnpackErrorValue, @@ -16,7 +17,10 @@ use synctv_xiu::{ unpacketizer::{ChunkUnpacketizer, UnpackResult}, }, handshake::{define::ClientHandshakeState, handshake_client::SimpleHandshakeClient}, - messages::define::msg_type_id, + messages::{ + define::{msg_type_id, RtmpMessageData}, + parser::MessageParser, + }, netconnection::writer::{ConnectProperties, NetConnection}, netstream::writer::NetStreamWriter, protocol_control_messages::writer::ProtocolControlMessagesWriter, @@ -65,14 +69,37 @@ impl RtmpPublisher { app_name: impl Into, stream_name: impl Into, ) -> Result { - let app_name = app_name.into(); - let stream_name = stream_name.into(); + Self::connect_inner(address, app_name.into(), stream_name.into(), true).await + } + + /// Sends a publish request and returns before the server accepts it. + /// + /// Rejection-path tests use this to exercise sessions that never receive + /// `NetStream.Publish.Start`. + pub async fn connect_unconfirmed( + address: SocketAddr, + app_name: impl Into, + stream_name: impl Into, + ) -> Result { + Self::connect_inner(address, app_name.into(), stream_name.into(), false).await + } + + async fn connect_inner( + address: SocketAddr, + app_name: String, + stream_name: String, + wait_for_confirmation: bool, + ) -> Result { let (io, mut stream_writer) = connect_rtmp_session(address, &app_name, "SyncTV test publisher").await?; stream_writer .write_publish(&3.0, &stream_name, &"live".to_string()) .await?; + if wait_for_confirmation { + wait_for_publish_start(&io, Duration::from_secs(5)).await?; + } + let media = Common::new( Some(ChunkPacketizer::new(Arc::clone(&io))), // Common's send methods never use the StreamHub event sender. @@ -81,8 +108,6 @@ impl RtmpPublisher { None, ); - // Let the server install the publication before the first media packet. - tokio::time::sleep(Duration::from_millis(20)).await; Ok(Self { io, media }) } @@ -138,6 +163,74 @@ impl RtmpPublisher { } } +async fn wait_for_publish_start(io: &SharedIo, timeout: Duration) -> Result<()> { + let deadline = tokio::time::Instant::now() + timeout; + let mut unpacketizer = ChunkUnpacketizer::new(); + + loop { + let remaining = deadline + .checked_duration_since(tokio::time::Instant::now()) + .ok_or_else(|| anyhow::anyhow!("RTMP publish confirmation timed out"))?; + let data = io.lock().await.read_timeout(remaining).await?; + unpacketizer.extend_data(&data)?; + + loop { + let chunks = match unpacketizer.read_chunks() { + Ok(UnpackResult::Chunks(chunks)) => chunks, + Ok(_) => continue, + Err(error) + if matches!( + error.value, + UnpackErrorValue::CannotParse | UnpackErrorValue::MessageTooLarge(_, _) + ) => + { + return Err(error.into()); + } + Err(_) => break, + }; + + for chunk in chunks { + if chunk.message_header.msg_type_id == msg_type_id::SET_CHUNK_SIZE { + anyhow::ensure!(chunk.payload.len() >= 4, "truncated RTMP Set Chunk Size"); + let chunk_size = u32::from_be_bytes([ + chunk.payload[0], + chunk.payload[1], + chunk.payload[2], + chunk.payload[3], + ]) & 0x7fff_ffff; + unpacketizer.update_max_chunk_size(usize::try_from(chunk_size)?); + continue; + } + + if chunk.message_header.msg_type_id != msg_type_id::COMMAND_AMF0 + && chunk.message_header.msg_type_id != msg_type_id::COMMAND_AMF3 + { + continue; + } + let Some(RtmpMessageData::Amf0Command(command)) = + MessageParser::new(chunk).parse()? + else { + continue; + }; + if command.command_name != Amf0ValueType::UTF8String("onStatus".to_string()) { + continue; + } + let Some(Amf0ValueType::Object(status)) = command.others.first() else { + continue; + }; + let Some(Amf0ValueType::UTF8String(code)) = status.get("code") else { + continue; + }; + anyhow::ensure!( + code == "NetStream.Publish.Start", + "RTMP publish rejected with status {code}" + ); + return Ok(()); + } + } + } +} + impl RtmpPlayer { /// Connects to an RTMP server and starts a live play session. pub async fn connect( diff --git a/synctv-xiu/src/rtmp/session/server_session.rs b/synctv-xiu/src/rtmp/session/server_session.rs index 8353abbf..ca307cba 100644 --- a/synctv-xiu/src/rtmp/session/server_session.rs +++ b/synctv-xiu/src/rtmp/session/server_session.rs @@ -997,24 +997,6 @@ impl ServerSession { return Err(e.into()); } - let mut netstream = NetStreamWriter::new(Arc::clone(&self.io)); - if let Err(e) = netstream - .write_on_status(transaction_id, "status", "NetStream.Publish.Start", "") - .await - { - tracing::error!( - "Failed to send NetStream.Publish.Start after successful auth, cleaning up: {}", - e - ); - cleanup_auth().await; - return Err(e.into()); - } - tracing::info!( - "[ S->C ] [NetStream.Publish.Start] app_name: {}, stream_name: {}", - self.app_name, - self.stream_name - ); - if let Err(e) = self .common .publish_to_stream_hub( @@ -1035,6 +1017,25 @@ impl ServerSession { self.is_publishing = true; + let mut netstream = NetStreamWriter::new(Arc::clone(&self.io)); + if let Err(error) = netstream + .write_on_status(transaction_id, "status", "NetStream.Publish.Start", "") + .await + { + tracing::error!( + "Failed to confirm StreamHub publication to the RTMP publisher: {error}" + ); + if let Err(cleanup_error) = self.teardown_active_stream().await { + tracing::warn!(%cleanup_error, "RTMP publication cleanup failed after confirmation error"); + } + return Err(error.into()); + } + tracing::info!( + "[ S->C ] [NetStream.Publish.Start] app_name: {}, stream_name: {}", + self.app_name, + self.stream_name + ); + // Notify publisher start via callback if let Some(cb) = &self.callbacks.on_publisher_start { cb(); diff --git a/synctv-xiu/src/streamhub/define.rs b/synctv-xiu/src/streamhub/define.rs index 060a9046..811ac7b8 100644 --- a/synctv-xiu/src/streamhub/define.rs +++ b/synctv-xiu/src/streamhub/define.rs @@ -327,6 +327,14 @@ impl FrameDataReceiver { Self::Unbounded(receiver) => receiver.recv().await, } } + + pub fn try_recv(&mut self) -> Result { + match self { + Self::Bounded(receiver) => receiver.try_recv(), + Self::Budgeted { receiver } => receiver.try_recv().map(|item| item.data), + Self::Unbounded(receiver) => receiver.try_recv(), + } + } } fn frame_data_bytes(data: &FrameData) -> usize { diff --git a/synctv-xiu/src/streamhub/transceiver.rs b/synctv-xiu/src/streamhub/transceiver.rs index 657c675f..df0e2e86 100644 --- a/synctv-xiu/src/streamhub/transceiver.rs +++ b/synctv-xiu/src/streamhub/transceiver.rs @@ -327,6 +327,17 @@ impl StreamDataTransceiver { } } _ = exit.recv() => { + while let Ok(data) = receiver.try_recv() { + if let Err(error) = Self::receive_frame_data( + Some(data), + &context, + &mut cached_snapshot, + &mut cached_gen, + ).await { + tracing::warn!(%error, "buffered publisher frame rejected during shutdown"); + break; + } + } break; } } @@ -413,6 +424,17 @@ impl StreamDataTransceiver { ).await; } _ = exit.recv() => { + while let Ok(data) = receiver.try_recv() { + Self::receive_packet_data( + Some(data), + &packet_senders, + &generation, + &mut cached_snapshot, + &mut cached_gen, + &statistics_data, + publisher_activity.as_deref(), + ).await; + } break; } } @@ -744,7 +766,6 @@ impl StreamDataTransceiver { let event_result = event_handle .await .map_err(|error| map_task_join_error("event loop", &error)); - tasks.abort_all(); let mut first_error = event_result.err(); while let Some(join_result) = tasks.join_next().await { diff --git a/synctv-xiu/tests/streaming_e2e_tests.rs b/synctv-xiu/tests/streaming_e2e_tests.rs index c514d412..7166c920 100644 --- a/synctv-xiu/tests/streaming_e2e_tests.rs +++ b/synctv-xiu/tests/streaming_e2e_tests.rs @@ -2653,7 +2653,7 @@ async fn duplicate_rtmp_publisher_cannot_replace_active_stream() -> Result<()> { first.send_video(0, true).await?; first.send_video(1, true).await?; - let mut duplicate = RtmpPublisher::connect(server.address, APP, STREAM).await?; + let mut duplicate = RtmpPublisher::connect_unconfirmed(server.address, APP, STREAM).await?; let duplicate_marker = [0x17, 0x01, 0, 0, 0, 0, 0, 0, 4, 0x65, 0xaa, 0xbb, 0xcc]; let _ = duplicate.send_raw_video(2, &duplicate_marker).await; tokio::time::sleep(Duration::from_millis(30)).await; @@ -2701,7 +2701,7 @@ async fn rejected_publish_releases_session_and_same_key_can_publish_next() -> Re }); let server = start_rtmp_hub_server_with_auth(Some(auth)).await?; - let rejected = RtmpPublisher::connect(server.address, APP, STREAM).await?; + let rejected = RtmpPublisher::connect_unconfirmed(server.address, APP, STREAM).await?; let deadline = Instant::now() + Duration::from_secs(2); while publish_attempts.load(Ordering::SeqCst) < 1 { anyhow::ensure!(