Perform the handshake with a custom `HelloFrame`. Returns `(stream, ack_frame)` on success, or the parsed `HelloErrorFrame` via `Err`.
(
addr: std::net::SocketAddr,
hello: &HelloFrame,
)
| 164 | /// Perform the handshake with a custom `HelloFrame`. |
| 165 | /// Returns `(stream, ack_frame)` on success, or the parsed `HelloErrorFrame` via `Err`. |
| 166 | async fn do_handshake( |
| 167 | addr: std::net::SocketAddr, |
| 168 | hello: &HelloFrame, |
| 169 | ) -> Result<(TcpStream, HelloAckFrame), HelloErrorFrame> { |
| 170 | let mut stream = TcpStream::connect(addr).await.expect("connect"); |
| 171 | stream |
| 172 | .write_all(&hello.encode()) |
| 173 | .await |
| 174 | .expect("write hello"); |
| 175 | stream.flush().await.expect("flush"); |
| 176 | |
| 177 | let mut magic_buf = [0u8; 4]; |
| 178 | stream.read_exact(&mut magic_buf).await.expect("read magic"); |
| 179 | let magic = u32::from_be_bytes(magic_buf); |
| 180 | |
| 181 | if magic == HELLO_ERROR_MAGIC_U32 { |
| 182 | // Read error code + msg_len + message. |
| 183 | let mut code_buf = [0u8; 1]; |
| 184 | stream.read_exact(&mut code_buf).await.expect("read code"); |
| 185 | let mut len_buf = [0u8; 1]; |
| 186 | stream.read_exact(&mut len_buf).await.expect("read msg_len"); |
| 187 | let msg_len = len_buf[0] as usize; |
| 188 | let mut msg = vec![0u8; msg_len]; |
| 189 | if msg_len > 0 { |
| 190 | stream.read_exact(&mut msg).await.expect("read msg"); |
| 191 | } |
| 192 | // Reassemble the full error frame bytes for HelloErrorFrame::decode. |
| 193 | let mut full = Vec::with_capacity(6 + msg_len); |
| 194 | full.extend_from_slice(b"NDBE"); |
| 195 | full.push(code_buf[0]); |
| 196 | full.push(len_buf[0]); |
| 197 | full.extend_from_slice(&msg); |
| 198 | let err_frame = HelloErrorFrame::decode(&full).expect("decode error frame"); |
| 199 | return Err(err_frame); |
| 200 | } |
| 201 | |
| 202 | assert_eq!(magic, HELLO_ACK_MAGIC, "expected HelloAck magic"); |
| 203 | |
| 204 | // Read fixed rest: proto_version(2) + capabilities(8) + sv_len(1). |
| 205 | let mut fixed_rest = [0u8; 11]; |
| 206 | stream |
| 207 | .read_exact(&mut fixed_rest) |
| 208 | .await |
| 209 | .expect("read fixed"); |
| 210 | let sv_len = fixed_rest[10] as usize; |
| 211 | let var_len = sv_len + 1 + 7 * 5; |
| 212 | let mut var_buf = vec![0u8; var_len]; |
| 213 | stream.read_exact(&mut var_buf).await.expect("read var"); |
| 214 | |
| 215 | let mut ack_buf = Vec::with_capacity(4 + 11 + var_len); |
| 216 | ack_buf.extend_from_slice(&magic_buf); |
| 217 | ack_buf.extend_from_slice(&fixed_rest); |
| 218 | ack_buf.extend_from_slice(&var_buf); |
| 219 | |
| 220 | let ack = HelloAckFrame::decode(&ack_buf).expect("decode ack"); |
| 221 | Ok((stream, ack)) |
| 222 | } |
| 223 |
no test coverage detected