Send a request and read the response.
(
&mut self,
op: OpCode,
fields: TextFields,
)
| 352 | |
| 353 | /// Send a request and read the response. |
| 354 | pub(crate) async fn send( |
| 355 | &mut self, |
| 356 | op: OpCode, |
| 357 | fields: TextFields, |
| 358 | ) -> NodeDbResult<NativeResponse> { |
| 359 | let req = NativeRequest { |
| 360 | op, |
| 361 | seq: self.next_seq(), |
| 362 | fields: RequestFields::Text(fields), |
| 363 | }; |
| 364 | |
| 365 | let payload = zerompk::to_msgpack_vec(&req) |
| 366 | .map_err(|e| NodeDbError::serialization("msgpack", format!("request encode: {e}")))?; |
| 367 | |
| 368 | let len = payload.len() as u32; |
| 369 | self.stream |
| 370 | .write_all(&len.to_be_bytes()) |
| 371 | .await |
| 372 | .map_err(io_err)?; |
| 373 | self.stream.write_all(&payload).await.map_err(io_err)?; |
| 374 | self.stream.flush().await.map_err(io_err)?; |
| 375 | |
| 376 | let mut combined_rows: Vec<Vec<nodedb_types::Value>> = Vec::new(); |
| 377 | let mut final_resp: Option<NativeResponse> = None; |
| 378 | |
| 379 | loop { |
| 380 | let mut len_buf = [0u8; FRAME_HEADER_LEN]; |
| 381 | self.stream.read_exact(&mut len_buf).await.map_err(io_err)?; |
| 382 | let resp_len = u32::from_be_bytes(len_buf); |
| 383 | if resp_len > MAX_FRAME_SIZE { |
| 384 | return Err(NodeDbError::internal(format!( |
| 385 | "response frame too large: {resp_len}" |
| 386 | ))); |
| 387 | } |
| 388 | |
| 389 | let mut resp_buf = vec![0u8; resp_len as usize]; |
| 390 | self.stream |
| 391 | .read_exact(&mut resp_buf) |
| 392 | .await |
| 393 | .map_err(io_err)?; |
| 394 | |
| 395 | let resp: NativeResponse = zerompk::from_msgpack(&resp_buf).map_err(|e| { |
| 396 | NodeDbError::serialization("msgpack", format!("response decode: {e}")) |
| 397 | })?; |
| 398 | |
| 399 | if resp.status == ResponseStatus::Partial { |
| 400 | if let Some(rows) = resp.rows { |
| 401 | combined_rows.extend(rows); |
| 402 | } |
| 403 | if final_resp.is_none() { |
| 404 | final_resp = Some(NativeResponse { rows: None, ..resp }); |
| 405 | } |
| 406 | } else { |
| 407 | if combined_rows.is_empty() { |
| 408 | final_resp = Some(resp); |
| 409 | } else { |
| 410 | if let Some(ref rows) = resp.rows { |
| 411 | combined_rows.extend(rows.iter().cloned()); |
no test coverage detected