1use crate::request::Request;
2use crate::websocket::{CloseCode, ErrorCode, Message, MessageMode, MessageType};
3use crate::{HttpError, Method, Uri};
4use hpack::Decoder;
5use log::{info, warn};
6use rustls::{ClientConnection, ServerConnection, StreamOwned};
7use std::io::{ErrorKind, Read, Write};
8use std::net::TcpStream;
9use std::sync::{Arc, Mutex};
10use std::thread;
11use std::time::Duration;
12
13#[derive(Debug, Clone)]
14pub enum Scheme {
15 Http(Arc<Mutex<TcpStream>>),
16 Https(Arc<Mutex<StreamOwned<ServerConnection, TcpStream>>>),
17}
18
19pub struct SchemeReader {
21 inner: SchemeReaderInner,
22 pending: Vec<u8>,
23}
24
25#[allow(dead_code)]
26enum SchemeReaderInner {
27 Http(TcpStream),
28 Https(Box<rustls::StreamOwned<ServerConnection, TcpStream>>),
29}
30
31pub struct SchemeWriter {
33 inner: SchemeWriterInner,
34}
35
36#[allow(dead_code)]
37enum SchemeWriterInner {
38 Http(TcpStream),
39 Https(Box<rustls::StreamOwned<ServerConnection, TcpStream>>),
40}
41
42impl Scheme {
43 pub fn split_for_websocket(
44 scheme: &Arc<Mutex<Scheme>>,
45 ) -> Result<(SchemeReader, SchemeWriter), HttpError> {
46 let guard = scheme
47 .lock()
48 .map_err(|e| HttpError::new(500, &format!("lock poisoned: {}", e)))?;
49 match &*guard {
50 Scheme::Http(stream) => {
51 let inner_guard = stream
52 .lock()
53 .map_err(|e| HttpError::new(500, &format!("lock poisoned: {}", e)))?;
54 let read_stream = inner_guard.try_clone().map_err(|e| {
55 HttpError::new(500, &format!("clone read stream failed: {}", e))
56 })?;
57 let write_stream = inner_guard.try_clone().map_err(|e| {
58 HttpError::new(500, &format!("clone write stream failed: {}", e))
59 })?;
60 Ok((
61 SchemeReader {
62 inner: SchemeReaderInner::Http(read_stream),
63 pending: vec![],
64 },
65 SchemeWriter {
66 inner: SchemeWriterInner::Http(write_stream),
67 },
68 ))
69 }
70 Scheme::Https(_) => Err(HttpError::new(
71 500,
72 "HTTPS split not supported, use shared mode",
73 )),
74 }
75 }
76}
77
78impl SchemeReader {
79 pub fn read_ws_data(
80 &mut self,
81 deflate: &crate::websocket::DeflateConfig,
82 ) -> Result<Message, HttpError> {
83 if !self.pending.is_empty() {
85 let message = Message::parse_message(&mut self.pending, deflate);
86 match message.message_type {
87 MessageType::TimeOut => {} _ => return Ok(message),
89 }
90 }
91
92 let mut buffer = vec![0u8; 1024 * 1024];
93
94 let res = match &mut self.inner {
95 SchemeReaderInner::Http(stream) => {
96 stream
97 .set_read_timeout(Some(Duration::from_millis(100)))
98 .ok();
99 let result = stream.read(&mut buffer);
100 stream.set_read_timeout(None).ok();
101 result
102 }
103 SchemeReaderInner::Https(stream) => {
104 stream
105 .get_mut()
106 .set_read_timeout(Some(Duration::from_millis(100)))
107 .ok();
108 let result = stream.read(&mut buffer);
109 stream.get_mut().set_read_timeout(None).ok();
110 result
111 }
112 };
113
114 match res {
115 Ok(0) => Ok(Message {
116 mode: MessageMode::Client,
117 message_type: MessageType::Close,
118 payload: vec![],
119 text: CloseCode::GoingAway.str(),
120 close: CloseCode::GoingAway,
121 error: ErrorCode::None,
122 }),
123 Ok(n) => {
124 self.pending.extend_from_slice(&buffer[..n]);
125 let start = std::time::Instant::now();
126 let total_timeout = Duration::from_secs(30);
127
128 loop {
129 if start.elapsed() > total_timeout {
130 log::warn!("等待 WebSocket 完整帧超时 (30s),关闭连接");
131 return Ok(Message {
132 mode: MessageMode::Client,
133 message_type: MessageType::Close,
134 payload: vec![],
135 text: "等待完整帧超时".to_string(),
136 close: CloseCode::ProtocolError,
137 error: ErrorCode::TimeOut,
138 });
139 }
140
141 let message = Message::parse_message(&mut self.pending, deflate);
142
143 match message.message_type {
144 MessageType::TimeOut => {
145 let mut more_buffer = vec![0u8; 1024 * 64];
146 let more_res = match &mut self.inner {
147 SchemeReaderInner::Http(stream) => {
148 stream
149 .set_read_timeout(Some(Duration::from_millis(100)))
150 .ok();
151 let result = stream.read(&mut more_buffer);
152 stream.set_read_timeout(None).ok();
153 result
154 }
155 SchemeReaderInner::Https(stream) => {
156 stream
157 .get_mut()
158 .set_read_timeout(Some(Duration::from_millis(100)))
159 .ok();
160 let result = stream.read(&mut more_buffer);
161 stream.get_mut().set_read_timeout(None).ok();
162 result
163 }
164 };
165 match more_res {
166 Ok(0) => {
167 return Ok(Message {
168 mode: MessageMode::Client,
169 message_type: MessageType::Close,
170 payload: vec![],
171 text: CloseCode::GoingAway.str(),
172 close: CloseCode::GoingAway,
173 error: ErrorCode::None,
174 });
175 }
176 Ok(m) => {
177 self.pending.extend_from_slice(&more_buffer[..m]);
178 continue;
179 }
180 Err(ref e)
181 if e.kind() == ErrorKind::WouldBlock
182 || e.kind() == ErrorKind::TimedOut =>
183 {
184 continue;
185 }
186 Err(_) => {
187 continue;
188 }
189 }
190 }
191 _ => return Ok(message),
192 }
193 }
194 }
195 Err(ref e) if e.kind() == ErrorKind::WouldBlock => Ok(Message {
196 mode: MessageMode::Client,
197 message_type: MessageType::TimeOut,
198 payload: vec![],
199 text: String::new(),
200 close: CloseCode::NormalClosure,
201 error: ErrorCode::TimeOut,
202 }),
203 Err(e) => Ok(Message {
204 mode: MessageMode::Client,
205 message_type: MessageType::Error,
206 payload: vec![],
207 text: e.to_string(),
208 close: CloseCode::Other(1011),
209 error: ErrorCode::Unknown,
210 }),
211 }
212 }
213}
214
215impl SchemeWriter {
216 pub fn write_all(&mut self, data: &[u8]) -> Result<(), HttpError> {
218 let result = match &mut self.inner {
219 SchemeWriterInner::Http(stream) => stream.write_all(data),
220 SchemeWriterInner::Https(stream) => stream.write_all(data),
221 };
222 match result {
223 Ok(()) => {
224 self.flush()?;
225 Ok(())
226 }
227 Err(e) => Err(HttpError::new(500, format!("write: {}", e).as_str())),
228 }
229 }
230
231 pub fn flush(&mut self) -> Result<(), HttpError> {
233 let result = match &mut self.inner {
234 SchemeWriterInner::Http(stream) => stream.flush(),
235 SchemeWriterInner::Https(stream) => stream.flush(),
236 };
237 match result {
238 Ok(()) => Ok(()),
239 Err(e) => Err(HttpError::new(500, format!("flush: {}", e).as_str())),
240 }
241 }
242}
243
244impl Scheme {
245 pub fn read(&mut self, data: &mut Vec<u8>) -> Result<(), HttpError> {
247 let mut buf = vec![0u8; 1024 * 1024];
248
249 let mut index = 2;
250 loop {
251 let result = match self {
252 Self::Http(stream) => stream.lock().unwrap().read(&mut buf),
253 Self::Https(stream) => stream.lock().unwrap().read(&mut buf),
254 };
255 return match result {
256 Ok(0) => Err(HttpError::new(500, "read: 客户端主动关闭")),
257 Ok(n) => {
258 data.extend(&buf[..n]);
259 return Ok(());
260 }
261 Err(ref e) if e.kind() == ErrorKind::Interrupted => {
262 if !data.is_empty() {
263 return Ok(());
264 }
265 if index > 0 {
266 index -= 1;
267 continue;
268 }
269 Err(HttpError::new(
270 500,
271 format!("read现在没数据可读: {}", e.to_string().as_str()).as_str(),
272 ))
273 }
274 Err(e) => Err(HttpError::new(
275 500,
276 format!("read: {}", e.to_string().as_str()).as_str(),
277 )),
278 };
279 }
280 }
281 fn read_data(&self, init_data: &mut Vec<u8>, length: usize) -> Result<(), HttpError> {
283 loop {
284 if init_data.len() >= length {
285 return Ok(());
286 }
287 let mut buf = vec![0u8; 1024 * 1024];
288 let result = match self {
289 Self::Http(stream) => stream.lock().unwrap().read(&mut buf),
290 Self::Https(stream) => stream.lock().unwrap().read(&mut buf),
291 };
292 return match result {
293 Ok(0) => Err(HttpError::new(500, "read_data: 客户端主动关闭")),
294 Ok(n) => {
295 init_data.extend(&buf[..n]);
296 Ok(())
297 }
298 Err(ref e) if e.kind() == ErrorKind::WouldBlock => {
299 thread::sleep(Duration::from_millis(100));
301 continue;
302 }
303 Err(e) => Err(HttpError::new(
304 500,
305 format!("read_data: {}", e.to_string().as_str()).as_str(),
306 )),
307 };
308 }
309 }
310
311 pub fn write(&mut self, data: &[u8]) -> Result<(), HttpError> {
313 let mut off = 0;
314 loop {
315 let result = match self {
316 Self::Http(stream) => stream.lock().unwrap().write(&data[off..]),
317 Self::Https(stream) => stream.lock().unwrap().get_mut().write(&data[off..]),
318 };
319 match result {
320 Ok(0) => return Err(HttpError::new(500, "write: 客户端主动关闭")),
321 Ok(e) => {
322 if e != data.len() {
323 off = e;
324 continue;
325 }
326 self.flush()?;
327 return Ok(());
328 }
329 Err(ref e)
330 if e.kind() == ErrorKind::WouldBlock || e.kind() == ErrorKind::Interrupted => {}
331 Err(e) => {
332 return Err(HttpError::new(
333 500,
334 format!("write: {}", e.to_string().as_str()).as_str(),
335 ))
336 }
337 };
338 }
339 }
340 pub fn write_all(&mut self, data: &[u8]) -> Result<(), HttpError> {
341 let result = match self {
342 Self::Http(stream) => stream.lock().unwrap().write_all(data),
343 Self::Https(stream) => stream.lock().unwrap().write_all(data),
344 };
345 match result {
346 Ok(()) => {
347 self.flush()?;
348 Ok(())
349 }
350 Err(e) => Err(HttpError::new(
351 500,
352 format!("write: {}", e.to_string().as_str()).as_str(),
353 )),
354 }
355 }
356
357 pub fn flush(&mut self) -> Result<(), HttpError> {
358 let result = match self {
359 Self::Http(stream) => stream.lock().unwrap().flush(),
360 Self::Https(stream) => stream.lock().unwrap().flush(),
361 };
362 match result {
363 Ok(()) => Ok(()),
364 Err(e) => Err(HttpError::new(
365 500,
366 format!("flush: {}", e.to_string().as_str()).as_str(),
367 )),
368 }
369 }
370
371 pub fn read_ws_data(
372 &mut self,
373 deflate: &crate::websocket::DeflateConfig,
374 ) -> Result<Message, HttpError> {
375 let mut response = vec![];
376 let mut buffer = vec![0u8; 1024 * 1024];
377 let res = match self {
378 Self::Http(stream) => {
379 let mut guard = stream.lock().unwrap();
380 guard
381 .set_read_timeout(Some(std::time::Duration::from_millis(100)))
382 .ok();
383 let result = guard.read(&mut buffer);
384 guard.set_read_timeout(None).ok();
385 result
386 }
387 Self::Https(ref mut stream) => {
388 let mut guard = stream.lock().unwrap();
389 guard
390 .get_mut()
391 .set_read_timeout(Some(std::time::Duration::from_millis(100)))
392 .ok();
393 let result = guard.read(&mut buffer);
394 guard.get_mut().set_read_timeout(None).ok();
395 result
396 }
397 };
398 match res {
399 Ok(0) => Ok(Message {
400 mode: MessageMode::Client,
401 message_type: MessageType::Close,
402 payload: vec![],
403 text: CloseCode::GoingAway.str(),
404 close: CloseCode::GoingAway,
405 error: ErrorCode::None,
406 }),
407 Ok(n) => {
408 response.extend(buffer[..n].to_vec());
409 let start = std::time::Instant::now();
410 let total_timeout = std::time::Duration::from_secs(30);
411
412 loop {
413 if start.elapsed() > total_timeout {
414 log::warn!("等待 WebSocket 完整帧超时 (30s),关闭连接");
415 return Ok(Message {
416 mode: MessageMode::Client,
417 message_type: MessageType::Close,
418 payload: vec![],
419 text: "等待完整帧超时".to_string(),
420 close: CloseCode::ProtocolError,
421 error: ErrorCode::TimeOut,
422 });
423 }
424
425 let message = Message::parse_message(&mut response, deflate);
426
427 match message.message_type {
428 MessageType::TimeOut => {
429 let mut more_buffer = vec![0u8; 1024 * 64];
430 let more_res = match self {
431 Self::Http(stream) => {
432 let mut guard = stream.lock().unwrap();
433 guard
434 .set_read_timeout(Some(std::time::Duration::from_millis(
435 100,
436 )))
437 .ok();
438 let result = guard.read(&mut more_buffer);
439 guard.set_read_timeout(None).ok();
440 result
441 }
442 Self::Https(ref mut stream) => {
443 let mut guard = stream.lock().unwrap();
444 guard
445 .get_mut()
446 .set_read_timeout(Some(std::time::Duration::from_millis(
447 100,
448 )))
449 .ok();
450 let result = guard.read(&mut more_buffer);
451 guard.get_mut().set_read_timeout(None).ok();
452 result
453 }
454 };
455 match more_res {
456 Ok(0) => {
457 return Ok(Message {
458 mode: MessageMode::Client,
459 message_type: MessageType::Close,
460 payload: vec![],
461 text: CloseCode::GoingAway.str(),
462 close: CloseCode::GoingAway,
463 error: ErrorCode::None,
464 });
465 }
466 Ok(m) => {
467 response.extend(more_buffer[..m].to_vec());
468 continue;
469 }
470 Err(ref e)
471 if e.kind() == ErrorKind::WouldBlock
472 || e.kind() == ErrorKind::TimedOut =>
473 {
474 continue;
475 }
476 Err(_) => {
477 continue;
478 }
479 }
480 }
481 _ => return Ok(message),
482 }
483 }
484 }
485 Err(ref e) if e.kind() == ErrorKind::WouldBlock => Ok(Message {
486 mode: MessageMode::Client,
487 message_type: MessageType::TimeOut,
488 payload: vec![],
489 text: String::new(),
490 close: CloseCode::NormalClosure,
491 error: ErrorCode::TimeOut,
492 }),
493 Err(e) => Ok(Message {
494 mode: MessageMode::Client,
495 message_type: MessageType::Error,
496 payload: vec![],
497 text: e.to_string(),
498 close: CloseCode::Other(1011),
499 error: ErrorCode::Unknown,
500 }),
501 }
502 }
503
504 pub fn client_ip(&mut self) -> String {
505 match self {
506 Self::Http(stream) => match stream.lock().unwrap().peer_addr() {
507 Ok(e) => e.ip().to_string(),
508 Err(_) => "未知".to_string(),
509 },
510 Self::Https(stream) => stream
511 .lock()
512 .unwrap()
513 .get_mut()
514 .peer_addr()
515 .unwrap()
516 .ip()
517 .to_string(),
518 }
519 }
520 pub fn server_ip(&mut self) -> String {
521 match self {
522 Self::Http(stream) => stream
523 .lock()
524 .unwrap()
525 .local_addr()
526 .unwrap()
527 .ip()
528 .to_string(),
529 Self::Https(stream) => stream
530 .lock()
531 .unwrap()
532 .get_mut()
533 .local_addr()
534 .unwrap()
535 .ip()
536 .to_string(),
537 }
538 }
539 pub fn http2_packet(
541 &mut self,
542 init_data: &mut Vec<u8>,
543 ) -> Result<(Vec<u8>, FrameType, u8, u32), HttpError> {
544 let bytes = init_data;
545 self.read_data(bytes, 9)?;
546 let headers = bytes.drain(..9).collect::<Vec<u8>>();
547 let length =
548 ((headers[0] as u32) << 16) | (u32::from(headers[1]) << 8) | u32::from(headers[2]);
549 let frame_type = headers[3];
550 let flags = headers[4];
551 let stream_id =
552 u32::from_be_bytes([headers[5], headers[6], headers[7], headers[8]]) & 0x7FFF_FFFF;
553 self.read_data(bytes, length as usize)?;
554 let payload = bytes.drain(..length as usize).collect::<Vec<u8>>();
555 Ok((payload, FrameType::from(frame_type), flags, stream_id))
556 }
557 pub fn http2_handle_header(
559 &mut self,
560 data: &mut Vec<u8>,
561 request: &mut Request,
562 ) -> Result<(), HttpError> {
563 loop {
564 let (payload, frame_type, flags, stream_id) = self.http2_packet(data)?;
565 if request.config.debug {
566 info!("http2_handle_header: frame_type: {frame_type:?} flags: {flags} stream_id: {stream_id} payload: {}", payload.len());
567 }
568 match frame_type {
569 FrameType::Settings => {
570 let is_ack = flags & 0x01 != 0;
571 if !is_ack {
572 self.http2_settings_ack()?;
573 }
574 }
575 FrameType::WindowUpdate => {
576 if payload.len() == 4 {
577 let raw = u32::from_be_bytes(payload.clone().try_into().unwrap());
578 let increment = raw & 0x7FFF_FFFF; if request.config.debug {
580 info!("WindowUpdate: increment = {} {:?}", increment, payload);
581 }
582 } else {
583 return Err(HttpError::new(
584 400,
585 format!("Invalid WindowUpdate frame length: {}", payload.len())
586 .as_str(),
587 ));
588 }
589 }
590 FrameType::Headers => {
591 let mut decoder = Decoder::new();
592 let headers = decoder.decode(&payload).unwrap();
593 if request.config.debug {
594 println!(
595 "=================请求头 {:?}=================",
596 thread::current().id()
597 );
598 }
599 for (name, value) in headers {
600 let header_name = String::from_utf8_lossy(name.as_slice());
601 let header_value = String::from_utf8_lossy(value.as_slice());
602 if request.config.debug {
603 println!("{header_name}: {header_value}");
604 }
605 match header_name.as_ref() {
606 ":method" => request.method = Method::from(header_value.as_ref()),
607 ":path" => request.uri = Uri::from(header_value.as_ref()),
608 ":scheme" => request.set_header("scheme", header_value.as_ref())?,
609 ":authority" => request.set_header("host", header_value.as_ref())?,
610 _ => request.set_header(&header_name, &header_value)?,
611 }
612 }
613 if request.config.debug {
614 println!("====================================================");
615 }
616 return Ok(());
617 }
618 _ => {
619 return Err(HttpError::new(
620 400,
621 format!("Invalid {frame_type:?}").as_str(),
622 ))
623 }
624 }
625 }
626 }
627 pub fn http2_handle_body(
629 &mut self,
630 data: &mut Vec<u8>,
631 request: Request,
632 ) -> Result<Vec<u8>, HttpError> {
633 let mut body = vec![];
634 loop {
635 let (payload, frame_type, flags, stream_id) = self.http2_packet(data)?;
636 if request.config.debug {
637 info!("http2_handle_body: frame_type: {frame_type:?} flags: {flags} stream_id: {stream_id} data: {}",payload.len());
638 }
639 match frame_type {
640 FrameType::Data => {
641 if flags == 1 {
642 body.extend(payload);
643 return Ok(body);
644 }
645 body.extend(payload);
646 }
647 FrameType::Headers => {}
648 FrameType::RstStream => {}
649 FrameType::Settings => {
650 if !payload.is_empty() {
651 self.http2_send_server_settings()?;
652 } else {
653 self.http2_settings_ack()?;
654 }
655 }
656 FrameType::Ping => {}
657 FrameType::Goaway => {
658 let text = String::from_utf8_lossy(&payload);
659 if request.config.debug {
660 warn!("Goaway: {text}");
661 }
662 return Ok(vec![]);
663 }
664 FrameType::WindowUpdate => {
665 if payload.len() == 4 {
666 let raw = u32::from_be_bytes(payload.clone().try_into().unwrap());
667 let increment = raw & 0x7FFF_FFFF; if request.config.debug {
669 info!("WindowUpdate: increment = {} {:?}", increment, payload);
670 }
671 } else {
672 return Err(HttpError::new(
673 400,
674 format!("Invalid WindowUpdate frame length: {}", payload.len())
675 .as_str(),
676 ));
677 }
678 }
679 FrameType::Continuation => {}
680 FrameType::None => {}
681 }
682 }
683 }
684 pub fn http2_send_server_settings(&mut self) -> Result<(), HttpError> {
686 let payload = {
687 let mut p = Vec::new();
688 p.extend_from_slice(&2u16.to_be_bytes());
690 p.extend_from_slice(&0u32.to_be_bytes());
691 p.extend_from_slice(&4u16.to_be_bytes());
693 p.extend_from_slice(&65_535u32.to_be_bytes());
694 p.extend_from_slice(&5u16.to_be_bytes());
696 p.extend_from_slice(&16_384u32.to_be_bytes());
697 p
699 };
700 let len = payload.len();
701 let mut f = Vec::with_capacity(9 + len);
702 f.extend_from_slice(&[(len >> 16) as u8, (len >> 8) as u8, len as u8]); f.push(0x04); f.push(0x00); f.extend_from_slice(&0u32.to_be_bytes()); f.extend_from_slice(&payload);
707 self.write_all(&f)?;
708 Ok(())
709 }
710 pub fn http2_settings_ack(&mut self) -> Result<(), HttpError> {
712 let f = [0x00, 0x00, 0x00, 0x04, 0x01, 0x00, 0x00, 0x00, 0x00];
713 self.write_all(&f)?;
714 Ok(())
715 }
716 pub fn http2_goaway(&mut self, last_stream_id: u32, error_code: u32) -> Result<(), HttpError> {
717 let mut frame = Vec::new();
719 frame.extend_from_slice(&[0x00, 0x00, 0x08]); frame.push(0x07); frame.push(0x00); frame.extend_from_slice(&[0x00, 0x00, 0x00, 0x00]); frame.extend_from_slice(&last_stream_id.to_be_bytes());
724 frame.extend_from_slice(&error_code.to_be_bytes());
725 self.write_all(frame.as_slice())?;
726 Ok(())
727 }
728}
729#[derive(Debug)]
730pub enum FrameType {
731 Data,
732 Headers,
733 RstStream,
734 Settings,
735 Ping,
736 Goaway,
737 WindowUpdate,
738 Continuation,
739 None,
740}
741impl FrameType {
742 pub fn from(code: u8) -> Self {
743 match code {
744 0x00 => Self::Data,
745 0x01 => Self::Headers,
746 0x03 => Self::RstStream,
747 0x04 => Self::Settings,
748 0x06 => Self::Ping,
749 0x07 => Self::Goaway,
750 0x08 => Self::WindowUpdate,
751 0x09 => Self::Continuation,
752 _ => Self::None,
753 }
754 }
755}
756
757pub enum ClientStream {
758 Http(TcpStream),
759 Https(Box<StreamOwned<ClientConnection, TcpStream>>),
760}
761impl ClientStream {
762 pub fn write_all(&mut self, data: &[u8]) -> std::io::Result<()> {
763 match self {
764 ClientStream::Http(e) => e.write_all(data),
765 ClientStream::Https(e) => e.write_all(data),
766 }
767 }
768 pub fn read(&mut self, buf: &mut [u8]) -> std::io::Result<usize> {
769 match self {
770 ClientStream::Http(e) => e.read(buf),
771 ClientStream::Https(e) => e.read(buf),
772 }
773 }
774 pub fn read_data(&mut self, buffer: &mut Vec<u8>) -> Result<(), String> {
775 let mut tmp = [0u8; 1024];
776 let n = self.read(&mut tmp).map_err(|e| e.to_string())?;
777 if n == 0 {
778 return Err("unexpected EOF while reading chunk data".to_string());
779 }
780 buffer.extend_from_slice(&tmp[..n]);
781 Ok(())
782 }
783 pub fn flush(&mut self) -> std::io::Result<()> {
785 match self {
786 ClientStream::Http(e) => e.flush(),
787 ClientStream::Https(e) => e.flush(),
788 }
789 }
790}