use std::collections::BTreeMap;
use anyhow::Context;
use axum::extract::ws::{Message, WebSocket};
use axum::extract::{DefaultBodyLimit, Query, WebSocketUpgrade};
use axum::response::IntoResponse;
use axum::routing::options;
use axum::{Extension, Router};
use axum_extra::TypedHeader;
use bytes::Bytes;
use futures::{SinkExt, StreamExt};
use surrealdb_core::dbs::Session;
use surrealdb_core::dbs::capabilities::RouteTarget;
use surrealdb_types::{Array, SurrealValue, Value, Variables};
use tower_http::limit::RequestBodyLimitLayer;
use super::AppState;
use super::error::ResponseError;
use super::headers::Accept;
use super::output::Output;
use crate::cnf::{
HTTP_MAX_SQL_BODY_SIZE, WEBSOCKET_MAX_MESSAGE_SIZE, WEBSOCKET_MAX_WRITE_BUFFER_SIZE,
WEBSOCKET_READ_BUFFER_SIZE, WEBSOCKET_WRITE_BUFFER_SIZE,
};
use crate::ntw::error::Error as NetError;
use crate::ntw::input::bytes_to_utf8;
pub fn router<S>() -> Router<S>
where
S: Clone + Send + Sync + 'static,
{
Router::new()
.route("/sql", options(|| async {}).get(get_handler).post(post_handler))
.route_layer(DefaultBodyLimit::disable())
.layer(RequestBodyLimitLayer::new(*HTTP_MAX_SQL_BODY_SIZE))
}
async fn post_handler(
Extension(state): Extension<AppState>,
Extension(session): Extension<Session>,
output: Option<TypedHeader<Accept>>,
Query(params): Query<BTreeMap<String, String>>,
sql: Bytes,
) -> Result<Output, ResponseError> {
let vars = Variables::from(params);
let db = &state.datastore;
if !db.allows_http_route(&RouteTarget::Sql) {
warn!("Capabilities denied HTTP route request attempt, target: '{}'", &RouteTarget::Sql);
return Err(NetError::ForbiddenRoute(RouteTarget::Sql.to_string()).into());
}
if !db.allows_query_by_subject(session.au.as_ref()) {
return Err(NetError::ForbiddenRoute(RouteTarget::Sql.to_string()).into());
}
let sql = bytes_to_utf8(&sql).context("Non UTF-8 request body").map_err(ResponseError)?;
match db.execute(sql, &session, Some(vars)).await {
Ok(res) => match output.as_deref() {
None | Some(Accept::ApplicationJson) => {
let v = Value::Array(Array::from(
res.into_iter().map(|x| x.into_value()).collect::<Vec<Value>>(),
));
Ok(Output::json_value(&v))
}
Some(Accept::ApplicationCbor) => {
let v = Value::Array(Array::from(
res.into_iter().map(|x| x.into_value()).collect::<Vec<Value>>(),
));
Ok(Output::cbor(v))
}
Some(Accept::ApplicationFlatbuffers) => {
let v = res.into_value();
Ok(Output::flatbuffers(&v))
}
Some(_) => Err(NetError::InvalidType.into()),
},
Err(err) => Err(ResponseError(err.into())),
}
}
async fn get_handler(
ws: WebSocketUpgrade,
Extension(state): Extension<AppState>,
Extension(sess): Extension<Session>,
) -> impl IntoResponse {
let db = &state.datastore;
if !db.allows_http_route(&RouteTarget::Sql) {
warn!("Capabilities denied HTTP route request attempt, target: '{}'", &RouteTarget::Sql);
return NetError::ForbiddenRoute(RouteTarget::Sql.to_string()).into_response();
}
if !db.allows_query_by_subject(sess.au.as_ref()) {
return NetError::ForbiddenRoute(RouteTarget::Sql.to_string()).into_response();
}
ws
.max_frame_size(*WEBSOCKET_MAX_MESSAGE_SIZE)
.max_message_size(*WEBSOCKET_MAX_MESSAGE_SIZE)
.read_buffer_size(*WEBSOCKET_READ_BUFFER_SIZE)
.write_buffer_size(*WEBSOCKET_WRITE_BUFFER_SIZE)
.max_write_buffer_size(*WEBSOCKET_MAX_WRITE_BUFFER_SIZE)
.on_failed_upgrade(|err| {
warn!("Failed to upgrade WebSocket connection: {err}");
})
.on_upgrade(move |socket| handle_socket(state, socket, sess))
}
async fn handle_socket(state: AppState, ws: WebSocket, session: Session) {
let (mut tx, mut rx) = ws.split();
while let Some(res) = rx.next().await {
if let Ok(msg) = res
&& let Ok(sql) = msg.to_text()
{
let db = &state.datastore;
let _ = match db.execute(sql, &session, None).await {
Ok(v) => match surrealdb_core::rpc::format::json::encode_str(Value::Array(
Array::from(v.into_iter().map(|x| x.into_value()).collect::<Vec<_>>()),
)) {
Ok(v) => tx.send(Message::Text(v.into())).await,
Err(e) => {
tx.send(Message::Text(format!("Failed to parse JSON: {e}",).into())).await
}
},
Err(e) => tx.send(Message::Text(e.to_string().into())).await,
};
}
}
}