Support stream EOF half closes
ci / test (push) Canceled after 0s
ci / fuzz-smoke (push) Canceled after 0s
ci / windows-client (push) Canceled after 0s
ci / package-release (linux-x86_64, ubuntu-latest) (push) Canceled after 0s
ci / package-release (macos-aarch64, macos-14) (push) Canceled after 0s
ci / package-release (macos-x86_64, macos-13) (push) Canceled after 0s
ci / package-release (windows-x86_64, windows-latest) (push) Canceled after 0s
ci / remote-bench (push) Canceled after 0s
ci / publish-gitea-release (push) Canceled after 0s

This commit is contained in:
DuProcess
2026-07-13 00:20:52 -04:00
parent dccd648ac9
commit bf50b38ce2
3 changed files with 499 additions and 26 deletions
+52 -2
View File
@@ -17,7 +17,7 @@
use crate::protocol::{
self, CLIENT_TO_SERVER, PacketKind, ReplayWindow, SERVER_TO_CLIENT, StreamClose, StreamData,
StreamOpen, StreamOpenOk, StreamOpenReject, StreamWindowAdjust,
StreamEof, StreamOpen, StreamOpenOk, StreamOpenReject, StreamWindowAdjust,
};
use crate::udp::{is_transient_udp_error, is_transient_udp_send_error};
use anyhow::{Result, anyhow, bail};
@@ -144,6 +144,9 @@ pub enum TransportEvent {
stream_id: u64,
flushed: Vec<OutgoingStreamPacket>,
},
Eof {
stream_id: u64,
},
Close {
stream_id: u64,
},
@@ -437,6 +440,17 @@ impl StreamMux {
encode_packet(PacketKind::StreamClose, &StreamClose { stream_id })
}
pub fn eof_stream(&self, stream_id: u64) -> Result<OutgoingStreamPacket> {
if !self.has_stream_state(stream_id) {
bail!("stream {stream_id} is not open");
}
encode_packet(PacketKind::StreamEof, &StreamEof { stream_id })
}
pub fn handle_eof(&self, eof: StreamEof) -> bool {
self.has_stream_state(eof.stream_id)
}
pub fn handle_close(&mut self, close: StreamClose) -> bool {
if !self.has_stream_state(close.stream_id) {
return false;
@@ -507,6 +521,15 @@ impl StreamMux {
Ok(TransportEvent::Ignored { stream_id })
}
}
PacketKind::StreamEof => {
let eof: StreamEof = protocol::from_body(body)?;
let stream_id = eof.stream_id;
if self.handle_eof(eof) {
Ok(TransportEvent::Eof { stream_id })
} else {
Ok(TransportEvent::Ignored { stream_id })
}
}
_ => bail!("packet kind {kind:?} is not a stream transport packet"),
}
}
@@ -1004,6 +1027,11 @@ impl DoshTransport {
Ok(())
}
pub async fn eof(&mut self, stream_id: u64) -> Result<()> {
let packet = self.mux.eof_stream(stream_id)?;
self.send_outgoing(packet).await
}
pub async fn close(&mut self, stream_id: u64) -> Result<()> {
let packet = self.mux.close_stream(stream_id)?;
self.send_outgoing(packet).await
@@ -1098,7 +1126,8 @@ impl DoshTransport {
| PacketKind::StreamOpenReject
| PacketKind::StreamData
| PacketKind::StreamWindowAdjust
| PacketKind::StreamClose => {
| PacketKind::StreamClose
| PacketKind::StreamEof => {
let event = self.mux.handle_packet(kind, plain)?;
self.send_followups(&event).await?;
Ok(SessionEvent::Stream(event))
@@ -1123,6 +1152,7 @@ impl DoshTransport {
}
TransportEvent::Open(_)
| TransportEvent::OpenReject { .. }
| TransportEvent::Eof { .. }
| TransportEvent::Close { .. }
| TransportEvent::Ignored { .. } => {}
}
@@ -1724,6 +1754,26 @@ mod tests {
assert!(mux.send_data(99, b"lost".to_vec()).is_err());
}
#[test]
fn stream_eof_is_delivered_without_closing_stream_state() {
let mut mux = StreamMux::new(TransportConfig::default());
mux.open_stream(1, "@dosh-test", 0).unwrap();
mux.handle_open_ok(StreamOpenOk { stream_id: 1 })
.unwrap()
.unwrap();
let eof = mux.eof_stream(1).unwrap();
assert_eq!(eof.kind, PacketKind::StreamEof);
assert!(mux.is_open(1));
let body = protocol::to_body(&StreamEof { stream_id: 1 }).unwrap();
assert_eq!(
mux.handle_packet(PacketKind::StreamEof, &body).unwrap(),
TransportEvent::Eof { stream_id: 1 }
);
assert!(mux.is_open(1));
}
#[test]
fn rejecting_incoming_stream_retires_local_state() {
let mut mux = StreamMux::new(TransportConfig::default());