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