use std::sync::Arc;
use axum::extract::ws::Message;
use opentelemetry::Context as TelemetryContext;
use surrealdb_core::rpc::DbResponse;
use surrealdb_core::rpc::format::Format;
use tokio::sync::mpsc::Sender;
use tracing::Span;
use crate::rpc::format::WsFormat;
use crate::telemetry::metrics::ws::record_rpc;
pub async fn send(
response: DbResponse,
cx: Arc<TelemetryContext>,
fmt: Format,
chn: Sender<Message>,
) {
let id = response.id.clone();
let session_id = response.session_id;
let span = Span::current();
debug!("Process RPC response");
let is_error = response.result.is_err();
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());
}
let (len, msg) = match fmt.res_ws(response) {
Ok((l, m)) => (l, m),
Err(err) => fmt
.res_ws(DbResponse::failure(id, session_id, err))
.expect("Serialising internal error should always succeed"),
};
if chn.send(msg).await.is_ok() {
record_rpc(cx.as_ref(), len, is_error);
};
}