Skip to main content

br_web_server/
stream.rs

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
19/// WebSocket 读取端 - 独立拥有读取半连接,无锁竞争
20pub 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
31/// WebSocket 写入端 - 独立拥有写入半连接,无锁竞争
32pub 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        // 先检查 pending 缓冲区是否有完整帧
84        if !self.pending.is_empty() {
85            let message = Message::parse_message(&mut self.pending, deflate);
86            match message.message_type {
87                MessageType::TimeOut => {} // 数据不完整,继续读取
88                _ => 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    /// 写入数据
217    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    /// 刷新缓冲区
232    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    /// 读取
246    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    /// 指定长度
282    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                    // ⚠️ 当前无数据可读,稍后再试
300                    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    /// 写入
312    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    /// 读取HTTP2
540    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    /// 读取HTTP2消息头
558    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; // 屏蔽最高位保留位
579                        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    /// 读取HTTP2消息体
628    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; // 屏蔽最高位保留位
668                        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    /// 发送HTTP2认证参数
685    pub fn http2_send_server_settings(&mut self) -> Result<(), HttpError> {
686        let payload = {
687            let mut p = Vec::new();
688            // SETTINGS_ENABLE_PUSH = 0  (对浏览器禁用推送)
689            p.extend_from_slice(&2u16.to_be_bytes());
690            p.extend_from_slice(&0u32.to_be_bytes());
691            // SETTINGS_INITIAL_WINDOW_SIZE = 65535
692            p.extend_from_slice(&4u16.to_be_bytes());
693            p.extend_from_slice(&65_535u32.to_be_bytes());
694            // SETTINGS_MAX_FRAME_SIZE = 16384
695            p.extend_from_slice(&5u16.to_be_bytes());
696            p.extend_from_slice(&16_384u32.to_be_bytes());
697            // 可按需再加 MAX_CONCURRENT_STREAMS 等
698            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]); // ✅ 正确长度
703        f.push(0x04); // type = SETTINGS
704        f.push(0x00); // flags = none
705        f.extend_from_slice(&0u32.to_be_bytes()); // sid=0
706        f.extend_from_slice(&payload);
707        self.write_all(&f)?;
708        Ok(())
709    }
710    /// 发送对“对端 SETTINGS”的 ACK(len=0, flags=ACK)
711    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        // 构造帧头
718        let mut frame = Vec::new();
719        frame.extend_from_slice(&[0x00, 0x00, 0x08]); // length
720        frame.push(0x07); // type = GOAWAY
721        frame.push(0x00); // flags = none
722        frame.extend_from_slice(&[0x00, 0x00, 0x00, 0x00]); // stream id = 0
723        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    /// 强制输出
784    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}