use axum::extract::ws::Message;
use surrealdb_core::rpc::format::Format;
use surrealdb_rpc::DbResponse;
use tokio::sync::mpsc::Sender;
use tokio::sync::mpsc::error::TrySendError;
use tracing::Span;
use crate::rpc::format::WsFormat;
fn encode(response: DbResponse, fmt: Format) -> Message {
let id = response.id.clone();
let session_id = response.session_id;
let span = Span::current();
debug!("Process RPC response");
if let Err(err) = &response.result {
span.record("otel.status_code", "ERROR");
span.record("rpc.error_kind", format!("{:?}", err.kind_str()));
span.record("rpc.error_message", err.message());
}
match fmt.res_ws(response) {
Ok((_len, msg)) => msg,
Err(err) => {
let (_len, msg) = fmt
.res_ws(DbResponse::failure(id, session_id, err))
.expect("Serialising internal error should always succeed");
msg
}
}
}
pub async fn send(response: DbResponse, fmt: Format, chn: Sender<Message>) {
let _ = chn.send(encode(response, fmt)).await;
}
pub fn try_send(
response: DbResponse,
fmt: Format,
chn: &Sender<Message>,
) -> Result<(), TrySendError<Message>> {
chn.try_send(encode(response, fmt))
}