fix: synchronize RTMP publish lifecycle

pull/370/head
zijiren233 2 months ago
parent f7499889df
commit e7910ada54
No known key found for this signature in database
GPG Key ID: 534E082AAA9B39DC

@ -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<String>,
stream_name: impl Into<String>,
) -> Result<Self> {
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<String>,
stream_name: impl Into<String>,
) -> Result<Self> {
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<Self> {
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(

@ -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();

@ -327,6 +327,14 @@ impl FrameDataReceiver {
Self::Unbounded(receiver) => receiver.recv().await,
}
}
pub fn try_recv(&mut self) -> Result<FrameData, mpsc::error::TryRecvError> {
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 {

@ -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 {

@ -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!(

Loading…
Cancel
Save