Task that reads [`OutboundWsMessage`]s from `messages`, encodes them, and sends the resulting [`Frame`]s to `outgoing_frames`. Meant to be [`tokio::spawn`]ed. The function also takes care of reusing serialization buffers and reporting metrics via [`SendMetrics`].
(
metrics: SendMetrics,
config: ClientConfig,
mut messages: mpsc::UnboundedReceiver<OutboundWsMessage>,
outgoing_frames: mpsc::UnboundedSender<Frame>,
bsatn_rlb_pool: BsatnRowListB
| 1428 | /// The function also takes care of reusing serialization buffers and reporting |
| 1429 | /// metrics via [`SendMetrics`]. |
| 1430 | async fn ws_encode_task( |
| 1431 | metrics: SendMetrics, |
| 1432 | config: ClientConfig, |
| 1433 | mut messages: mpsc::UnboundedReceiver<OutboundWsMessage>, |
| 1434 | outgoing_frames: mpsc::UnboundedSender<Frame>, |
| 1435 | bsatn_rlb_pool: BsatnRowListBuilderPool, |
| 1436 | ) { |
| 1437 | let mut encoder = WsEncoder { |
| 1438 | config, |
| 1439 | buffers: SerializeBufferPool::new(config), |
| 1440 | metrics: &metrics, |
| 1441 | outgoing_frames: &outgoing_frames, |
| 1442 | bsatn_rlb_pool: &bsatn_rlb_pool, |
| 1443 | binary_server_messages: Vec::new(), |
| 1444 | }; |
| 1445 | let mut message_batch = Vec::new(); |
| 1446 | while messages.recv_many(&mut message_batch, ENCODE_BATCH_SIZE).await != 0 { |
| 1447 | log::trace!("encoding batch of {} websocket messages", message_batch.len()); |
| 1448 | // `encode_batch` drains `message_batch` on success. If forwarding to |
| 1449 | // the websocket send loop fails, the receiver is gone, so the encode |
| 1450 | // task can terminate. |
| 1451 | if encoder.encode_batch(&mut message_batch).await.is_err() { |
| 1452 | break; |
| 1453 | } |
| 1454 | } |
| 1455 | } |
| 1456 | |
| 1457 | /// Stateful websocket encoder for one client connection. |
| 1458 | /// |
no test coverage detected
searching dependent graphs…