(
State(state): State<WsState>,
existing_user: Option<Extension<AuthedUser>>,
ws: WebSocketUpgrade,
ConnectInfo(addr): ConnectInfo<SocketAddr>,
tower_session: Option<Extension<Towe
| 289 | } |
| 290 | |
| 291 | pub(crate) async fn handle_sql_ws( |
| 292 | State(state): State<WsState>, |
| 293 | existing_user: Option<Extension<AuthedUser>>, |
| 294 | ws: WebSocketUpgrade, |
| 295 | ConnectInfo(addr): ConnectInfo<SocketAddr>, |
| 296 | tower_session: Option<Extension<TowerSession>>, |
| 297 | ) -> Result<impl IntoResponse, AuthError> { |
| 298 | let session = tower_session.map(|Extension(session)| session); |
| 299 | // The `x_materialize_user_header_auth` middleware may have already provided the user for us |
| 300 | let user = match existing_user { |
| 301 | Some(Extension(user)) => Some(ExistingUser::XMaterializeUserHeader(user)), |
| 302 | None => { |
| 303 | let session = maybe_get_authenticated_session(session.as_ref()).await; |
| 304 | if let Some((session, session_data)) = session { |
| 305 | let user = ensure_session_unexpired(session, session_data).await?; |
| 306 | Some(ExistingUser::Session(user)) |
| 307 | } else { |
| 308 | None |
| 309 | } |
| 310 | } |
| 311 | }; |
| 312 | |
| 313 | let addr = Box::new(addr.ip()); |
| 314 | Ok(ws |
| 315 | .max_message_size(MAX_REQUEST_SIZE) |
| 316 | .on_upgrade(|ws| async move { run_ws(state, user, *addr, ws).await })) |
| 317 | } |
| 318 | |
| 319 | #[derive(Serialize, Deserialize, Debug, PartialEq, Eq)] |
| 320 | #[serde(untagged)] |
nothing calls this directly
no test coverage detected