(
mut stream: S,
store: Arc<RedbChangeStore>,
shutdown: Arc<Notify>,
)
| 1305 | } |
| 1306 | |
| 1307 | async fn handle_connection<S>( |
| 1308 | mut stream: S, |
| 1309 | store: Arc<RedbChangeStore>, |
| 1310 | shutdown: Arc<Notify>, |
| 1311 | ) -> anyhow::Result<()> |
| 1312 | where |
| 1313 | S: AsyncRead + AsyncWrite + Unpin, |
| 1314 | { |
| 1315 | let request: RequestFrame = read_frame(&mut stream).await?; |
| 1316 | let (response, should_shutdown) = handle_request(&store, request); |
| 1317 | write_frame(&mut stream, &response).await?; |
| 1318 | stream.shutdown().await?; |
| 1319 | if should_shutdown { |
| 1320 | // `notify_one` retains a permit if the accept loop is between polls, |
| 1321 | // preventing a shutdown request from being acknowledged but lost. |
| 1322 | shutdown.notify_one(); |
| 1323 | } |
| 1324 | Ok(()) |
| 1325 | } |
| 1326 | |
| 1327 | async fn read_frame<T, S>(stream: &mut S) -> anyhow::Result<T> |
| 1328 | where |
nothing calls this directly
no test coverage detected