(
ws: &mut WebSocket<MaybeTlsStream<TcpStream>>,
auth_message: Message,
)
| 1727 | } |
| 1728 | |
| 1729 | pub fn auth_with_ws_impl( |
| 1730 | ws: &mut WebSocket<MaybeTlsStream<TcpStream>>, |
| 1731 | auth_message: Message, |
| 1732 | ) -> Result<Vec<WebSocketResponse>, anyhow::Error> { |
| 1733 | ws.send(auth_message)?; |
| 1734 | |
| 1735 | // Wait for initial ready response. |
| 1736 | let mut msgs = Vec::new(); |
| 1737 | loop { |
| 1738 | let resp = ws.read()?; |
| 1739 | match resp { |
| 1740 | Message::Text(msg) => { |
| 1741 | let msg: WebSocketResponse = serde_json::from_str(&msg).unwrap(); |
| 1742 | match msg { |
| 1743 | WebSocketResponse::ReadyForQuery(_) => break, |
| 1744 | msg => { |
| 1745 | msgs.push(msg); |
| 1746 | } |
| 1747 | } |
| 1748 | } |
| 1749 | Message::Ping(_) => continue, |
| 1750 | Message::Close(None) => return Err(anyhow!("ws closed after auth")), |
| 1751 | Message::Close(Some(close_frame)) => { |
| 1752 | return Err(anyhow!("ws closed after auth").context(close_frame)); |
| 1753 | } |
| 1754 | _ => panic!("unexpected response: {:?}", resp), |
| 1755 | } |
| 1756 | } |
| 1757 | Ok(msgs) |
| 1758 | } |
| 1759 | |
| 1760 | pub fn make_header<H: Header>(h: H) -> HeaderMap { |
| 1761 | let mut map = HeaderMap::new(); |
no test coverage detected