| 75 | |
| 76 | impl Message { |
| 77 | pub fn read(r: &mut impl io::BufRead) -> Result<Option<Message>, ParseError> { |
| 78 | let mut buf = String::new(); |
| 79 | // Consume all headers - either end of stream or "\r\n\r\n" |
| 80 | loop { |
| 81 | // No more content, therefore no message |
| 82 | if r.read_line(&mut buf)? == 0 { |
| 83 | return Ok(None); |
| 84 | } |
| 85 | if buf.ends_with("\r\n\r\n") { |
| 86 | break; |
| 87 | } |
| 88 | } |
| 89 | let mut headers = [EMPTY_HEADER; 2]; |
| 90 | if let httparse::Status::Complete((size, _)) = parse_headers(buf.as_bytes(), &mut headers)? |
| 91 | && size != buf.len() |
| 92 | { |
| 93 | Err(ParseError::HeaderDecodeMismatch(size, buf.len()))? |
| 94 | } |
| 95 | let mut content_length = 0; |
| 96 | for header in &headers { |
| 97 | if header.name.eq_ignore_ascii_case("content-length") { |
| 98 | content_length = std::str::from_utf8(header.value)?.parse::<usize>()?; |
| 99 | } else if header.name.eq_ignore_ascii_case("content-type") { |
| 100 | // ¯\_(ツ)_/¯ |
| 101 | } else if header != &EMPTY_HEADER { |
| 102 | Err(ParseError::InvalidHeader(header.name.to_owned()))? |
| 103 | } |
| 104 | } |
| 105 | if content_length == 0 { |
| 106 | Err(ParseError::NoLength)? |
| 107 | } |
| 108 | buf.clear(); |
| 109 | let mut buf = buf.into_bytes(); |
| 110 | buf.resize(content_length, 0); |
| 111 | r.read_exact(&mut buf)?; |
| 112 | let message: Message = serde_json::from_slice(buf.as_slice())?; |
| 113 | Ok(Some(message)) |
| 114 | } |
| 115 | |
| 116 | pub fn write(self, w: &mut impl io::Write) -> io::Result<()> { |
| 117 | #[derive(Serialize)] |