diff --git a/Cargo.lock b/Cargo.lock index 846dcfab3cd..34064d20eec 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -8118,6 +8118,7 @@ dependencies = [ "futures", "headers", "http", + "http-body", "http-body-util", "humantime", "hyper", diff --git a/Cargo.toml b/Cargo.toml index 842332340d8..b4b3f6286d3 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -223,6 +223,7 @@ hex = "0.4.3" home = "0.5" hostname = "^0.3" http = "1.0" +http-body = "1.0" http-body-util= "0.1.3" humantime = "2.3" hyper = "1.0" diff --git a/crates/client-api/Cargo.toml b/crates/client-api/Cargo.toml index afbb5a2a340..3cd1f3fb466 100644 --- a/crates/client-api/Cargo.toml +++ b/crates/client-api/Cargo.toml @@ -16,6 +16,7 @@ spacetimedb-paths.workspace = true spacetimedb-schema.workspace = true base64.workspace = true +http-body.workspace = true http-body-util.workspace = true tokio = { version = "1.2", features = ["full"] } lazy_static = "1.4.0" diff --git a/crates/client-api/src/lib.rs b/crates/client-api/src/lib.rs index 7fd6f238128..3529ed069d3 100644 --- a/crates/client-api/src/lib.rs +++ b/crates/client-api/src/lib.rs @@ -14,10 +14,12 @@ use spacetimedb::energy::{EnergyBalance, EnergyQuanta}; use spacetimedb::host::{HostController, MigratePlanResult, ModuleHost, NoSuchModule, UpdateDatabaseResult}; use spacetimedb::identity::{AuthCtx, Identity}; use spacetimedb::messages::control_db::{Database, HostType, Node, Replica}; +use spacetimedb::metrics::ENGINE_METRICS; use spacetimedb::sql; use spacetimedb_client_api_messages::http::{SqlStmtResult, SqlStmtStats}; use spacetimedb_client_api_messages::name::{DomainName, InsertDomainResult, RegisterTldResult, SetDomainsResult, Tld}; -use spacetimedb_lib::{ProductTypeElement, ProductValue}; +use spacetimedb_datastore::execution_context::WorkloadType; +use spacetimedb_lib::{bsatn, ProductTypeElement, ProductValue}; use spacetimedb_paths::server::ModuleLogsDir; use spacetimedb_schema::auto_migrate::{MigrationPolicy, PrettyPrintStyle}; use thiserror::Error; @@ -137,7 +139,7 @@ impl Host { pub async fn exec_sql( &self, auth: AuthCtx, - _database: Database, + database: Database, confirmed_read: bool, body: String, ) -> axum::response::Result>> { @@ -172,6 +174,12 @@ impl Host { drop(_guard); sql_span.record("total_duration", tracing::field::debug(total_duration)); + // Charge egress; each transport charges for the rows it sends. + ENGINE_METRICS + .bytes_sent_to_clients + .with_label_values(&WorkloadType::Sql, &database.database_identity) + .inc_by(sql_egress_bytes(&result.rows)); + let schema = header .into_iter() .map(|(col_name, col_type)| ProductTypeElement::new(col_type, Some(col_name))) @@ -205,6 +213,12 @@ impl Host { .await } } + +/// Uses the BSATN size for parity with WebSocket queries. +fn sql_egress_bytes(rows: &[ProductValue]) -> u64 { + rows.iter().map(|row| bsatn::to_len(row).unwrap_or(0) as u64).sum() +} + /// Parameters for publishing a database. /// /// See [`ControlStateDelegate::publish_database`]. diff --git a/crates/client-api/src/routes/mod.rs b/crates/client-api/src/routes/mod.rs index 0a649db75b4..bbb966b8c3b 100644 --- a/crates/client-api/src/routes/mod.rs +++ b/crates/client-api/src/routes/mod.rs @@ -1,4 +1,14 @@ +use std::pin::Pin; +use std::task::{Context, Poll}; + +use ::prometheus::IntCounter; +use axum::body::{Body, Bytes, HttpBody}; +use axum::extract::{MatchedPath, Request}; +use axum::middleware::Next; +use axum::response::Response; use http::header; +use http_body::Frame; +use spacetimedb::worker_metrics::WORKER_METRICS; use tower_http::cors; use crate::{Authorization, ControlStateDelegate, NodeDelegate}; @@ -19,6 +29,88 @@ use self::{database::DatabaseRoutes, identity::IdentityRoutes}; /// establish a connection to SpacetimeDB. This API call doesn't actually do anything. pub async fn ping(_auth: crate::auth::SpacetimeAuthHeader) {} +/// A body wrapper that counts bytes as they are actually transferred. +struct CountingBody { + inner: Body, + counter: IntCounter, +} + +impl HttpBody for CountingBody { + type Data = Bytes; + type Error = axum::Error; + + fn poll_frame(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll, Self::Error>>> { + let poll = Pin::new(&mut self.inner).poll_frame(cx); + if let Poll::Ready(Some(Ok(frame))) = &poll + && let Some(data) = frame.data_ref() + { + self.counter.inc_by(data.len() as u64); + } + poll + } + + fn size_hint(&self) -> http_body::SizeHint { + self.inner.size_hint() + } + + fn is_end_stream(&self) -> bool { + self.inner.is_end_stream() + } +} + +impl CountingBody { + fn wrap(counter: IntCounter) -> impl FnOnce(Body) -> Body { + move |inner| Body::new(CountingBody { inner, counter }) + } +} + +/// Returns the method as a static label value, +/// bucketing non-standard methods to keep label cardinality bounded. +fn method_label(method: &http::Method) -> &'static str { + match method.as_str() { + "GET" => "GET", + "POST" => "POST", + "PUT" => "PUT", + "DELETE" => "DELETE", + "HEAD" => "HEAD", + "OPTIONS" => "OPTIONS", + "PATCH" => "PATCH", + "CONNECT" => "CONNECT", + "TRACE" => "TRACE", + _ => "OTHER", + } +} + +/// Records request count, latency and body sizes per HTTP route. +async fn http_metrics_middleware(req: Request, next: Next) -> Response { + let Some(route) = req.extensions().get::().cloned() else { + return next.run(req).await; + }; + + let method = method_label(req.method()); + + let request_body_bytes = WORKER_METRICS.http_request_body_bytes.with_label_values(route.as_str()); + let response_body_bytes = WORKER_METRICS + .http_response_body_bytes + .with_label_values(route.as_str()); + + let req = req.map(CountingBody::wrap(request_body_bytes)); + + let start = std::time::Instant::now(); + let res = next.run(req).await; + + WORKER_METRICS + .http_requests + .with_label_values(route.as_str(), method, res.status().as_str()) + .inc(); + WORKER_METRICS + .http_request_duration + .with_label_values(route.as_str()) + .observe(start.elapsed().as_secs_f64()); + + res.map(CountingBody::wrap(response_body_bytes)) +} + #[allow(clippy::let_and_return)] pub fn router( ctx: &S, @@ -45,6 +137,11 @@ where .allow_origin(cors::Any); axum::Router::new() - .nest("/v1", router.layer(cors)) + .nest( + "/v1", + router + .layer(axum::middleware::from_fn(http_metrics_middleware)) + .layer(cors), + ) .nest("/internal", internal::router()) } diff --git a/crates/core/src/sql/execute.rs b/crates/core/src/sql/execute.rs index ddd57d780c5..65204d50e1c 100644 --- a/crates/core/src/sql/execute.rs +++ b/crates/core/src/sql/execute.rs @@ -49,6 +49,8 @@ pub struct SqlResult { /// If a `ModuleHost` is provided, the SQL query is executed via the module host, /// meaning the module’s core is used to run the statement. /// If no module host is provided, the SQL query is executed on the current thread. +/// +/// Callers that send the returned rows to a client must record `bytes_sent_to_clients` themselves. pub async fn run( db: Arc, sql_text: String, @@ -1686,4 +1688,37 @@ pub(crate) mod tests { Ok(()) } + + // Egress is charged by the transport that sends the rows; charging here would double count. + #[test] + fn test_select_does_not_charge_egress() -> ResultTest<()> { + let db = TestDB::durable()?; + + let table_id = db.create_table_for_test("T", &[("a", AlgebraicType::U8)], &[])?; + with_auto_commit(&db, |tx| -> Result<_, DBError> { + for i in 0..4u8 { + insert(&db, tx, table_id, &product!(i))?; + } + Ok(()) + })?; + + let rt = db.runtime().expect("runtime should be there"); + + let server = Identity::from_claims("issuer", "server"); + let auth = AuthCtx::new(server, server); + + let result = rt.block_on(run( + db.clone(), + "SELECT * FROM T".to_string(), + auth, + None, + None, + &mut vec![], + ))?; + + assert_eq!(result.rows.len(), 4); + assert_eq!(result.metrics.bytes_sent_to_clients, 0); + + Ok(()) + } } diff --git a/crates/core/src/worker_metrics/mod.rs b/crates/core/src/worker_metrics/mod.rs index c4cfb407187..582ed549aab 100644 --- a/crates/core/src/worker_metrics/mod.rs +++ b/crates/core/src/worker_metrics/mod.rs @@ -221,6 +221,26 @@ metrics_group!( #[labels(database_identity: Identity, protocol: str)] pub websocket_request_msg_size: HistogramVec, + #[name = spacetime_http_requests_total] + #[help = "The cumulative number of HTTP requests, by matched route, method and response status"] + #[labels(route: str, method: str, status: str)] + pub http_requests: IntCounterVec, + + #[name = spacetime_http_request_duration_sec] + #[help = "The time (in seconds) spent handling an HTTP request, by matched route"] + #[labels(route: str)] + pub http_request_duration: HistogramVec, + + #[name = spacetime_http_request_body_bytes_total] + #[help = "The cumulative number of HTTP request body bytes received, by matched route"] + #[labels(route: str)] + pub http_request_body_bytes: IntCounterVec, + + #[name = spacetime_http_response_body_bytes_total] + #[help = "The cumulative number of HTTP response body bytes sent, by matched route. Observability only, not billed"] + #[labels(route: str)] + pub http_response_body_bytes: IntCounterVec, + #[name = jemalloc_active_bytes] #[help = "Number of bytes in jemallocs heap"] #[labels(node_id: str)]