(options: &PgDumpOptions, serve: F)
| 154 | pub(crate) type PgDumpVirtualSocket = TcpSocketHalf; |
| 155 | |
| 156 | pub(crate) fn dump_direct_sql<F>(options: &PgDumpOptions, serve: F) -> Result<String> |
| 157 | where |
| 158 | F: FnOnce(PgDumpVirtualSocket) -> Result<()>, |
| 159 | { |
| 160 | options.validate()?; |
| 161 | let (socket_tx, socket_rx) = mpsc::sync_channel(1); |
| 162 | let networking = DirectPgDumpNetworking::new(socket_tx); |
| 163 | let runner_options = options.clone(); |
| 164 | let runner = thread::spawn(move || { |
| 165 | dump_sql_with_networking(DIRECT_PG_DUMP_ADDR, &runner_options, networking) |
| 166 | }); |
| 167 | |
| 168 | let accepted = receive_direct_pg_dump_socket(&socket_rx, &runner) |
| 169 | .context("accept direct pg_dump virtual protocol connection"); |
| 170 | let serve_result = match accepted { |
| 171 | Ok(socket) => serve(socket), |
| 172 | Err(err) => Err(err), |
| 173 | }; |
| 174 | let dump_result = runner |
| 175 | .join() |
| 176 | .map_err(|_| anyhow!("direct pg_dump runner thread panicked"))?; |
| 177 | |
| 178 | match (serve_result, dump_result) { |
| 179 | (Ok(()), Ok(sql)) => Ok(sql), |
| 180 | (Err(err), Ok(_)) => Err(err), |
| 181 | (Ok(()), Err(err)) => Err(err), |
| 182 | (Err(err), Err(dump_err)) => { |
| 183 | Err(err.context(format!("direct pg_dump runner also failed: {dump_err:#}"))) |
| 184 | } |
| 185 | } |
| 186 | } |
| 187 | |
| 188 | fn dump_sql_with_networking<N>( |
| 189 | addr: SocketAddr, |
no test coverage detected