From 5ae1b23c6bbb0b7a68fa03d4fc9fb3f63b671e66 Mon Sep 17 00:00:00 2001 From: Graham Christensen Date: Tue, 14 Jul 2026 20:40:06 -0400 Subject: [PATCH 1/9] feat(tracing): instrument connections, queries, transactions, and migrations Add tracing spans and events across the core operational surface of all three backends (Postgres, MySQL, SQLite): - connection establish / close / ping - statement prepare, with cache hit/miss trace events - query execution (run / execute) - transaction begin / commit / rollback and implicit rollback on drop - pool acquire / connect - migration run / undo, with per-migration apply/skip/revert events Spans use OpenTelemetry-style field naming (db.system, server.address, db.operation.parameters, ...) and are gated at debug/trace level, with migrations logged at info. Co-Authored-By: Claude Opus 4.8 (1M context) --- sqlx-core/src/migrate/migrator.rs | 40 +++++++++++++++++ sqlx-core/src/pool/inner.rs | 18 ++++++++ sqlx-core/src/transaction.rs | 5 +++ sqlx-mysql/src/connection/establish.rs | 19 ++++++++ sqlx-mysql/src/connection/executor.rs | 21 +++++++++ sqlx-mysql/src/connection/mod.rs | 15 +++++++ sqlx-mysql/src/transaction.rs | 28 ++++++++++++ sqlx-postgres/src/connection/establish.rs | 22 +++++++++ sqlx-postgres/src/connection/executor.rs | 25 +++++++++++ sqlx-postgres/src/connection/mod.rs | 15 +++++++ sqlx-postgres/src/transaction.rs | 28 ++++++++++++ sqlx-sqlite/src/connection/worker.rs | 55 +++++++++++++++++++++++ 12 files changed, 291 insertions(+) diff --git a/sqlx-core/src/migrate/migrator.rs b/sqlx-core/src/migrate/migrator.rs index c64b0f211b..eeee10eb59 100644 --- a/sqlx-core/src/migrate/migrator.rs +++ b/sqlx-core/src/migrate/migrator.rs @@ -224,6 +224,17 @@ impl Migrator { // Getting around the annoying "implementation of `Acquire` is not general enough" error #[doc(hidden)] + #[tracing::instrument( + target = "sqlx::migrate", + name = "migrate.run", + skip_all, + fields( + migrate.target = target, + migrate.skip = skip, + migrate.table = %self.table_name, + migrate.known = self.migrations.len(), + ), + )] pub async fn run_direct( &self, target: Option, @@ -235,6 +246,7 @@ impl Migrator { { // lock the database for exclusive access by the migrator if self.locking { + tracing::debug!(target: "sqlx::migrate", "acquiring migration lock"); conn.lock().await?; } @@ -277,8 +289,21 @@ impl Migrator { } None => { if skip { + tracing::info!( + target: "sqlx::migrate", + version = migration.version, + description = %migration.description, + "skipping migration (marking as applied without running)", + ); conn.skip(&self.table_name, migration).await?; } else { + tracing::info!( + target: "sqlx::migrate", + version = migration.version, + description = %migration.description, + migration_type = ?migration.migration_type, + "applying migration", + ); conn.apply(&self.table_name, migration).await?; } } @@ -288,6 +313,7 @@ impl Migrator { // unlock the migrator to allow other migrators to run // but do nothing as we already migrated if self.locking { + tracing::debug!(target: "sqlx::migrate", "releasing migration lock"); conn.unlock().await?; } @@ -311,6 +337,12 @@ impl Migrator { /// # }) /// # } /// ``` + #[tracing::instrument( + target = "sqlx::migrate", + name = "migrate.undo", + skip_all, + fields(migrate.target = target, migrate.table = %self.table_name), + )] pub async fn undo<'a, A>(&self, migrator: A, target: i64) -> Result<(), MigrateError> where A: Acquire<'a>, @@ -320,6 +352,7 @@ impl Migrator { // lock the database for exclusive access by the migrator if self.locking { + tracing::debug!(target: "sqlx::migrate", "acquiring migration lock"); conn.lock().await?; } @@ -347,12 +380,19 @@ impl Migrator { .filter(|m| applied_migrations.contains_key(&m.version)) .filter(|m| m.version > target) { + tracing::info!( + target: "sqlx::migrate", + version = migration.version, + description = %migration.description, + "reverting migration", + ); conn.revert(&self.table_name, migration).await?; } // unlock the migrator to allow other migrators to run // but do nothing as we already migrated if self.locking { + tracing::debug!(target: "sqlx::migrate", "releasing migration lock"); conn.unlock().await?; } diff --git a/sqlx-core/src/pool/inner.rs b/sqlx-core/src/pool/inner.rs index b698dc9df0..54eb53c2c6 100644 --- a/sqlx-core/src/pool/inner.rs +++ b/sqlx-core/src/pool/inner.rs @@ -243,6 +243,17 @@ impl PoolInner { } } + #[tracing::instrument( + target = "sqlx::pool::acquire", + name = "pool.acquire", + skip_all, + fields( + pool.size = self.size(), + pool.idle = self.num_idle(), + pool.max = self.options.max_connections, + ), + level = "debug", + )] pub(super) async fn acquire(self: &Arc) -> Result>, Error> { if self.is_closed() { return Err(Error::PoolClosed); @@ -322,6 +333,13 @@ impl PoolInner { Ok(acquired) } + #[tracing::instrument( + target = "sqlx::pool::connect", + name = "pool.connect", + skip_all, + fields(pool.size = self.size()), + level = "debug", + )] pub(super) async fn connect( self: &Arc, deadline: Instant, diff --git a/sqlx-core/src/transaction.rs b/sqlx-core/src/transaction.rs index 46f89d25b2..1f463beec5 100644 --- a/sqlx-core/src/transaction.rs +++ b/sqlx-core/src/transaction.rs @@ -274,6 +274,11 @@ where // operation that will happen on the next asynchronous invocation of the underlying // connection (including if the connection is returned to a pool) + tracing::debug!( + target: "sqlx::transaction", + "transaction dropped without explicit commit/rollback; queueing implicit rollback", + ); + DB::TransactionManager::start_rollback(&mut self.connection); } } diff --git a/sqlx-mysql/src/connection/establish.rs b/sqlx-mysql/src/connection/establish.rs index f61654d876..96fd8836db 100644 --- a/sqlx-mysql/src/connection/establish.rs +++ b/sqlx-mysql/src/connection/establish.rs @@ -12,7 +12,21 @@ use crate::protocol::Capabilities; use crate::{MySqlConnectOptions, MySqlConnection, MySqlSslMode}; impl MySqlConnection { + #[tracing::instrument( + target = "sqlx::connect", + name = "mysql.establish", + skip_all, + fields( + db.system = "mysql", + server.address = %options.host, + server.port = options.port, + db.name = options.database.as_deref().unwrap_or_default(), + db.user = %options.username, + ), + level = "debug", + )] pub(crate) async fn establish(options: &MySqlConnectOptions) -> Result { + tracing::debug!("establishing MySQL connection"); let do_handshake = DoHandshake::new(options)?; let handshake = match &options.socket { @@ -22,6 +36,11 @@ impl MySqlConnection { let stream = handshake?; + tracing::debug!( + server_version = ?stream.server_version, + "MySQL connection established" + ); + Ok(Self { inner: Box::new(MySqlConnectionInner { stream, diff --git a/sqlx-mysql/src/connection/executor.rs b/sqlx-mysql/src/connection/executor.rs index ee59d03d0a..940132845a 100644 --- a/sqlx-mysql/src/connection/executor.rs +++ b/sqlx-mysql/src/connection/executor.rs @@ -26,6 +26,13 @@ use sqlx_core::sql_str::SqlStr; use std::{pin::pin, sync::Arc}; impl MySqlConnection { + #[tracing::instrument( + target = "sqlx::prepare", + name = "mysql.prepare", + skip_all, + fields(db.system = "mysql"), + level = "debug", + )] async fn prepare_statement( &mut self, sql: &str, @@ -78,10 +85,13 @@ impl MySqlConnection { sql: &str, ) -> Result<(u32, MySqlStatementMetadata), Error> { if let Some(statement) = self.inner.cache_statement.get_mut(sql) { + tracing::trace!(target: "sqlx::prepare", "prepared statement cache hit"); // is internally reference-counted return Ok((*statement).clone()); } + tracing::trace!(target: "sqlx::prepare", "prepared statement cache miss"); + let (id, metadata) = self.prepare_statement(sql).await?; // in case of the cache being full, close the least recently used statement @@ -100,6 +110,17 @@ impl MySqlConnection { } #[allow(clippy::needless_lifetimes)] + #[tracing::instrument( + target = "sqlx::query", + name = "mysql.run", + skip_all, + fields( + db.system = "mysql", + db.operation.parameters = arguments.as_ref().map_or(0, |a| a.types.len()), + db.mysql.prepared = arguments.is_some(), + ), + level = "debug", + )] pub(crate) async fn run<'e, 'c: 'e, 'q: 'e>( &'c mut self, sql: SqlStr, diff --git a/sqlx-mysql/src/connection/mod.rs b/sqlx-mysql/src/connection/mod.rs index 569ad32722..fd60237318 100644 --- a/sqlx-mysql/src/connection/mod.rs +++ b/sqlx-mysql/src/connection/mod.rs @@ -70,7 +70,15 @@ impl Connection for MySqlConnection { type Options = MySqlConnectOptions; + #[tracing::instrument( + target = "sqlx::connect", + name = "mysql.close", + skip_all, + fields(db.system = "mysql"), + level = "debug", + )] async fn close(mut self) -> Result<(), Error> { + tracing::debug!("closing MySQL connection gracefully"); self.inner.stream.send_packet(Quit).await?; self.inner.stream.shutdown().await?; @@ -82,6 +90,13 @@ impl Connection for MySqlConnection { Ok(()) } + #[tracing::instrument( + target = "sqlx::connect", + name = "mysql.ping", + skip_all, + fields(db.system = "mysql"), + level = "trace", + )] async fn ping(&mut self) -> Result<(), Error> { self.inner.stream.wait_until_ready().await?; self.inner.stream.send_packet(Ping).await?; diff --git a/sqlx-mysql/src/transaction.rs b/sqlx-mysql/src/transaction.rs index 18db30b183..992fe0d38e 100644 --- a/sqlx-mysql/src/transaction.rs +++ b/sqlx-mysql/src/transaction.rs @@ -14,6 +14,13 @@ pub struct MySqlTransactionManager; impl TransactionManager for MySqlTransactionManager { type Database = MySql; + #[tracing::instrument( + target = "sqlx::transaction", + name = "mysql.transaction.begin", + skip_all, + fields(db.system = "mysql", depth = conn.inner.transaction_depth), + level = "debug", + )] async fn begin(conn: &mut MySqlConnection, statement: Option) -> Result<(), Error> { let depth = conn.inner.transaction_depth; @@ -30,9 +37,18 @@ impl TransactionManager for MySqlTransactionManager { } conn.inner.transaction_depth += 1; + tracing::debug!("transaction/savepoint opened"); + Ok(()) } + #[tracing::instrument( + target = "sqlx::transaction", + name = "mysql.transaction.commit", + skip_all, + fields(db.system = "mysql", depth = conn.inner.transaction_depth), + level = "debug", + )] async fn commit(conn: &mut MySqlConnection) -> Result<(), Error> { let depth = conn.inner.transaction_depth; @@ -44,6 +60,13 @@ impl TransactionManager for MySqlTransactionManager { Ok(()) } + #[tracing::instrument( + target = "sqlx::transaction", + name = "mysql.transaction.rollback", + skip_all, + fields(db.system = "mysql", depth = conn.inner.transaction_depth), + level = "debug", + )] async fn rollback(conn: &mut MySqlConnection) -> Result<(), Error> { let depth = conn.inner.transaction_depth; @@ -59,6 +82,11 @@ impl TransactionManager for MySqlTransactionManager { let depth = conn.inner.transaction_depth; if depth > 0 { + tracing::debug!( + target: "sqlx::transaction", + depth, + "queueing implicit rollback for unfinished transaction/savepoint on drop", + ); conn.inner.stream.waiting.push_back(Waiting::Result); conn.inner.stream.sequence_id = 0; conn.inner diff --git a/sqlx-postgres/src/connection/establish.rs b/sqlx-postgres/src/connection/establish.rs index 3c2f516533..039aaf4df4 100644 --- a/sqlx-postgres/src/connection/establish.rs +++ b/sqlx-postgres/src/connection/establish.rs @@ -15,7 +15,23 @@ use super::PgConnectionInner; // https://www.postgresql.org/docs/current/protocol-flow.html#id-1.10.5.7.11 impl PgConnection { + #[tracing::instrument( + target = "sqlx::connect", + name = "postgres.establish", + skip_all, + fields( + db.system = "postgresql", + server.address = %options.host, + server.port = options.port, + db.name = options.database.as_deref().unwrap_or_default(), + db.user = %options.username, + // Recorded once the startup handshake completes. + db.postgresql.backend_pid = tracing::field::Empty, + ), + )] pub(crate) async fn establish(options: &PgConnectOptions) -> Result { + tracing::debug!("establishing PostgreSQL connection"); + // Upgrade to TLS if we were asked to and the server supports it let mut stream = PgStream::connect(options).await?; @@ -123,6 +139,12 @@ impl PgConnection { // start-up is completed. The frontend can now issue commands transaction_status = message.decode::()?.transaction_status; + tracing::Span::current().record("db.postgresql.backend_pid", process_id); + tracing::debug!( + backend_pid = process_id, + "PostgreSQL connection established" + ); + break; } diff --git a/sqlx-postgres/src/connection/executor.rs b/sqlx-postgres/src/connection/executor.rs index e0f4c3d44a..22f10ac227 100644 --- a/sqlx-postgres/src/connection/executor.rs +++ b/sqlx-postgres/src/connection/executor.rs @@ -20,6 +20,17 @@ use sqlx_core::sql_str::SqlStr; use sqlx_core::Either; use std::{pin::pin, sync::Arc}; +#[tracing::instrument( + target = "sqlx::prepare", + name = "postgres.prepare", + skip_all, + fields( + db.system = "postgresql", + db.operation.parameters = arg_types.len(), + db.postgresql.persistent = persistent, + ), + level = "debug", +)] async fn prepare( conn: &mut PgConnection, sql: &str, @@ -168,9 +179,12 @@ impl PgConnection { resolve_column_origin: bool, ) -> Result<(StatementId, Arc), Error> { if let Some(statement) = self.inner.cache_statement.get_mut(sql) { + tracing::trace!(target: "sqlx::prepare", "prepared statement cache hit"); return Ok((*statement).clone()); } + tracing::trace!(target: "sqlx::prepare", "prepared statement cache miss"); + let statement = prepare( self, sql, @@ -196,6 +210,17 @@ impl PgConnection { Ok(statement) } + #[tracing::instrument( + target = "sqlx::query", + name = "postgres.run", + skip_all, + fields( + db.system = "postgresql", + db.operation.parameters = arguments.as_ref().map_or(0, |a| a.len()), + db.postgresql.prepared = arguments.is_some(), + ), + level = "debug", + )] pub(crate) async fn run<'e, 'c: 'e, 'q: 'e>( &'c mut self, query: SqlStr, diff --git a/sqlx-postgres/src/connection/mod.rs b/sqlx-postgres/src/connection/mod.rs index d594585b6c..b4d1240d09 100644 --- a/sqlx-postgres/src/connection/mod.rs +++ b/sqlx-postgres/src/connection/mod.rs @@ -159,7 +159,15 @@ impl Connection for PgConnection { type Options = PgConnectOptions; + #[tracing::instrument( + target = "sqlx::connect", + name = "postgres.close", + skip_all, + fields(db.system = "postgresql", backend_pid = self.inner.process_id), + level = "debug", + )] async fn close(mut self) -> Result<(), Error> { + tracing::debug!("closing PostgreSQL connection gracefully"); // The normal, graceful termination procedure is that the frontend sends a Terminate // message and immediately closes the connection. @@ -177,6 +185,13 @@ impl Connection for PgConnection { Ok(()) } + #[tracing::instrument( + target = "sqlx::connect", + name = "postgres.ping", + skip_all, + fields(db.system = "postgresql"), + level = "trace", + )] async fn ping(&mut self) -> Result<(), Error> { // Users were complaining about this showing up in query statistics on the server. // By sending a comment we avoid an error if the connection was in the middle of a rowset diff --git a/sqlx-postgres/src/transaction.rs b/sqlx-postgres/src/transaction.rs index 3f4122ea82..aa36030e67 100644 --- a/sqlx-postgres/src/transaction.rs +++ b/sqlx-postgres/src/transaction.rs @@ -14,6 +14,13 @@ pub struct PgTransactionManager; impl TransactionManager for PgTransactionManager { type Database = Postgres; + #[tracing::instrument( + target = "sqlx::transaction", + name = "postgres.transaction.begin", + skip_all, + fields(db.system = "postgresql", depth = conn.inner.transaction_depth), + level = "debug", + )] async fn begin(conn: &mut PgConnection, statement: Option) -> Result<(), Error> { let depth = conn.inner.transaction_depth; @@ -34,9 +41,18 @@ impl TransactionManager for PgTransactionManager { rollback.conn.inner.transaction_depth += 1; rollback.defuse(); + tracing::debug!("transaction/savepoint opened"); + Ok(()) } + #[tracing::instrument( + target = "sqlx::transaction", + name = "postgres.transaction.commit", + skip_all, + fields(db.system = "postgresql", depth = conn.inner.transaction_depth), + level = "debug", + )] async fn commit(conn: &mut PgConnection) -> Result<(), Error> { if conn.inner.transaction_depth > 0 { conn.execute(commit_ansi_transaction_sql(conn.inner.transaction_depth)) @@ -48,6 +64,13 @@ impl TransactionManager for PgTransactionManager { Ok(()) } + #[tracing::instrument( + target = "sqlx::transaction", + name = "postgres.transaction.rollback", + skip_all, + fields(db.system = "postgresql", depth = conn.inner.transaction_depth), + level = "debug", + )] async fn rollback(conn: &mut PgConnection) -> Result<(), Error> { if conn.inner.transaction_depth > 0 { conn.execute(rollback_ansi_transaction_sql(conn.inner.transaction_depth)) @@ -61,6 +84,11 @@ impl TransactionManager for PgTransactionManager { fn start_rollback(conn: &mut PgConnection) { if conn.inner.transaction_depth > 0 { + tracing::debug!( + target: "sqlx::transaction", + depth = conn.inner.transaction_depth, + "queueing implicit rollback for unfinished transaction/savepoint on drop", + ); conn.queue_simple_query( rollback_ansi_transaction_sql(conn.inner.transaction_depth).as_str(), ) diff --git a/sqlx-sqlite/src/connection/worker.rs b/sqlx-sqlite/src/connection/worker.rs index 2585dc1312..9450f7092d 100644 --- a/sqlx-sqlite/src/connection/worker.rs +++ b/sqlx-sqlite/src/connection/worker.rs @@ -8,6 +8,7 @@ use futures_intrusive::sync::{Mutex, MutexGuard}; use sqlx_core::sql_str::SqlStr; use tracing::span::Span; +use sqlx_core::arguments::Arguments as _; use sqlx_core::error::Error; use sqlx_core::transaction::{ begin_ansi_transaction_sql, commit_ansi_transaction_sql, rollback_ansi_transaction_sql, @@ -102,7 +103,15 @@ enum Command { } impl ConnectionWorker { + #[tracing::instrument( + target = "sqlx::connect", + name = "sqlite.establish", + skip_all, + fields(db.system = "sqlite"), + level = "debug", + )] pub(crate) async fn establish(params: EstablishParams) -> Result { + tracing::debug!("establishing SQLite connection worker"); let (establish_tx, establish_rx) = oneshot::channel(); thread::Builder::new() @@ -347,6 +356,13 @@ impl ConnectionWorker { establish_rx.await.map_err(|_| Error::WorkerCrashed)? } + #[tracing::instrument( + target = "sqlx::prepare", + name = "sqlite.prepare", + skip_all, + fields(db.system = "sqlite"), + level = "debug", + )] pub(crate) async fn prepare(&mut self, query: SqlStr) -> Result { self.oneshot_cmd(|tx| Command::Prepare { query, tx }) .await? @@ -361,6 +377,17 @@ impl ConnectionWorker { .await? } + #[tracing::instrument( + target = "sqlx::query", + name = "sqlite.execute", + skip_all, + fields( + db.system = "sqlite", + db.operation.parameters = args.as_ref().map_or(0, |a| a.len()), + db.sqlite.persistent = persistent, + ), + level = "debug", + )] pub(crate) async fn execute( &mut self, query: SqlStr, @@ -388,16 +415,37 @@ impl ConnectionWorker { Ok(rx) } + #[tracing::instrument( + target = "sqlx::transaction", + name = "sqlite.transaction.begin", + skip_all, + fields(db.system = "sqlite", depth = self.shared.get_transaction_depth()), + level = "debug", + )] pub(crate) async fn begin(&mut self, statement: Option) -> Result<(), Error> { self.oneshot_cmd_with_ack(|tx| Command::Begin { tx, statement }) .await? } + #[tracing::instrument( + target = "sqlx::transaction", + name = "sqlite.transaction.commit", + skip_all, + fields(db.system = "sqlite", depth = self.shared.get_transaction_depth()), + level = "debug", + )] pub(crate) async fn commit(&mut self) -> Result<(), Error> { self.oneshot_cmd_with_ack(|tx| Command::Commit { tx }) .await? } + #[tracing::instrument( + target = "sqlx::transaction", + name = "sqlite.transaction.rollback", + skip_all, + fields(db.system = "sqlite", depth = self.shared.get_transaction_depth()), + level = "debug", + )] pub(crate) async fn rollback(&mut self) -> Result<(), Error> { self.oneshot_cmd_with_ack(|tx| Command::Rollback { tx: Some(tx) }) .await? @@ -409,6 +457,13 @@ impl ConnectionWorker { .map_err(|_| Error::WorkerCrashed) } + #[tracing::instrument( + target = "sqlx::connect", + name = "sqlite.ping", + skip_all, + fields(db.system = "sqlite"), + level = "trace", + )] pub(crate) async fn ping(&mut self) -> Result<(), Error> { self.oneshot_cmd(|tx| Command::Ping { tx }).await } From 22abbce36eddc898b1c5467643afc6a8db81d611 Mon Sep 17 00:00:00 2001 From: Graham Christensen Date: Tue, 14 Jul 2026 20:50:15 -0400 Subject: [PATCH 2/9] fix(tracing): gate postgres.establish span at debug level Review finding #2: postgres.establish omitted `level`, so it defaulted to INFO while mysql.establish and sqlite.establish both use debug. Bring Postgres in line so connection-establish spans are consistently debug-level across backends. Co-Authored-By: Claude Opus 4.8 (1M context) --- sqlx-postgres/src/connection/establish.rs | 1 + 1 file changed, 1 insertion(+) diff --git a/sqlx-postgres/src/connection/establish.rs b/sqlx-postgres/src/connection/establish.rs index 039aaf4df4..1a81df33cb 100644 --- a/sqlx-postgres/src/connection/establish.rs +++ b/sqlx-postgres/src/connection/establish.rs @@ -28,6 +28,7 @@ impl PgConnection { // Recorded once the startup handshake completes. db.postgresql.backend_pid = tracing::field::Empty, ), + level = "debug", )] pub(crate) async fn establish(options: &PgConnectOptions) -> Result { tracing::debug!("establishing PostgreSQL connection"); From 8fd6abd2d1b6f7156c25c87c324558203b8d812b Mon Sep 17 00:00:00 2001 From: Graham Christensen Date: Tue, 14 Jul 2026 20:51:24 -0400 Subject: [PATCH 3/9] refactor(tracing): drop redundant connect/close log events Review finding #3: each establish/close fn opened a same-level, same-target span whose name already announces the operation, then immediately logged a bare "establishing/closing ..." event saying the same thing. Remove those leading duplicates. The trailing "connection established" events are kept since they carry new data (server_version / backend_pid). Co-Authored-By: Claude Opus 4.8 (1M context) --- sqlx-mysql/src/connection/establish.rs | 1 - sqlx-mysql/src/connection/mod.rs | 1 - sqlx-postgres/src/connection/establish.rs | 2 -- sqlx-postgres/src/connection/mod.rs | 1 - sqlx-sqlite/src/connection/worker.rs | 1 - 5 files changed, 6 deletions(-) diff --git a/sqlx-mysql/src/connection/establish.rs b/sqlx-mysql/src/connection/establish.rs index 96fd8836db..8849e74be1 100644 --- a/sqlx-mysql/src/connection/establish.rs +++ b/sqlx-mysql/src/connection/establish.rs @@ -26,7 +26,6 @@ impl MySqlConnection { level = "debug", )] pub(crate) async fn establish(options: &MySqlConnectOptions) -> Result { - tracing::debug!("establishing MySQL connection"); let do_handshake = DoHandshake::new(options)?; let handshake = match &options.socket { diff --git a/sqlx-mysql/src/connection/mod.rs b/sqlx-mysql/src/connection/mod.rs index fd60237318..a1155317c6 100644 --- a/sqlx-mysql/src/connection/mod.rs +++ b/sqlx-mysql/src/connection/mod.rs @@ -78,7 +78,6 @@ impl Connection for MySqlConnection { level = "debug", )] async fn close(mut self) -> Result<(), Error> { - tracing::debug!("closing MySQL connection gracefully"); self.inner.stream.send_packet(Quit).await?; self.inner.stream.shutdown().await?; diff --git a/sqlx-postgres/src/connection/establish.rs b/sqlx-postgres/src/connection/establish.rs index 1a81df33cb..a5b23c26dc 100644 --- a/sqlx-postgres/src/connection/establish.rs +++ b/sqlx-postgres/src/connection/establish.rs @@ -31,8 +31,6 @@ impl PgConnection { level = "debug", )] pub(crate) async fn establish(options: &PgConnectOptions) -> Result { - tracing::debug!("establishing PostgreSQL connection"); - // Upgrade to TLS if we were asked to and the server supports it let mut stream = PgStream::connect(options).await?; diff --git a/sqlx-postgres/src/connection/mod.rs b/sqlx-postgres/src/connection/mod.rs index b4d1240d09..cff61bbb23 100644 --- a/sqlx-postgres/src/connection/mod.rs +++ b/sqlx-postgres/src/connection/mod.rs @@ -167,7 +167,6 @@ impl Connection for PgConnection { level = "debug", )] async fn close(mut self) -> Result<(), Error> { - tracing::debug!("closing PostgreSQL connection gracefully"); // The normal, graceful termination procedure is that the frontend sends a Terminate // message and immediately closes the connection. diff --git a/sqlx-sqlite/src/connection/worker.rs b/sqlx-sqlite/src/connection/worker.rs index 9450f7092d..7a1490ed4f 100644 --- a/sqlx-sqlite/src/connection/worker.rs +++ b/sqlx-sqlite/src/connection/worker.rs @@ -111,7 +111,6 @@ impl ConnectionWorker { level = "debug", )] pub(crate) async fn establish(params: EstablishParams) -> Result { - tracing::debug!("establishing SQLite connection worker"); let (establish_tx, establish_rx) = oneshot::channel(); thread::Builder::new() From a314465ea27f1ad085924e6888a0ad6e695fef2e Mon Sep 17 00:00:00 2001 From: Graham Christensen Date: Tue, 14 Jul 2026 20:52:35 -0400 Subject: [PATCH 4/9] refactor(tracing): use Arguments::len() in mysql query span Review finding #4: the mysql.run span computed parameter count via `a.types.len()`, reaching into a pub(crate) field, while postgres.run and sqlite.execute use the `Arguments::len()` trait method. Use the trait method for cross-backend consistency. Co-Authored-By: Claude Opus 4.8 (1M context) --- sqlx-mysql/src/connection/executor.rs | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/sqlx-mysql/src/connection/executor.rs b/sqlx-mysql/src/connection/executor.rs index 940132845a..633f8fc91c 100644 --- a/sqlx-mysql/src/connection/executor.rs +++ b/sqlx-mysql/src/connection/executor.rs @@ -21,6 +21,7 @@ use futures_core::future::BoxFuture; use futures_core::stream::BoxStream; use futures_core::Stream; use futures_util::TryStreamExt; +use sqlx_core::arguments::Arguments as _; use sqlx_core::column::{ColumnOrigin, TableColumn}; use sqlx_core::sql_str::SqlStr; use std::{pin::pin, sync::Arc}; @@ -116,7 +117,7 @@ impl MySqlConnection { skip_all, fields( db.system = "mysql", - db.operation.parameters = arguments.as_ref().map_or(0, |a| a.types.len()), + db.operation.parameters = arguments.as_ref().map_or(0, |a| a.len()), db.mysql.prepared = arguments.is_some(), ), level = "debug", From 253e7f7a8a113a9cc5542510ff9900e8818fdd58 Mon Sep 17 00:00:00 2001 From: Graham Christensen Date: Tue, 14 Jul 2026 20:53:45 -0400 Subject: [PATCH 5/9] fix(tracing): align backend_pid field name on postgres.close Review finding #5: postgres.establish records the backend PID as `db.postgresql.backend_pid` but postgres.close recorded the same value as plain `backend_pid`, so consumers filtering on the field name would miss one. Use `db.postgresql.backend_pid` in both spans. Co-Authored-By: Claude Opus 4.8 (1M context) --- sqlx-postgres/src/connection/mod.rs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/sqlx-postgres/src/connection/mod.rs b/sqlx-postgres/src/connection/mod.rs index cff61bbb23..debcf01cb3 100644 --- a/sqlx-postgres/src/connection/mod.rs +++ b/sqlx-postgres/src/connection/mod.rs @@ -163,7 +163,7 @@ impl Connection for PgConnection { target = "sqlx::connect", name = "postgres.close", skip_all, - fields(db.system = "postgresql", backend_pid = self.inner.process_id), + fields(db.system = "postgresql", db.postgresql.backend_pid = self.inner.process_id), level = "debug", )] async fn close(mut self) -> Result<(), Error> { From cbdefbd5de8d9c62f799e2952d9ee81f53c8e82e Mon Sep 17 00:00:00 2001 From: Graham Christensen Date: Tue, 14 Jul 2026 21:21:50 -0400 Subject: [PATCH 6/9] feat(tracing): add InstrumentedStream span combinator to sqlx-core Review finding #1 (shared mechanism): tracing 0.1's `Instrumented` only implements `Future`, not `Stream`, so a `#[tracing::instrument]` on an async fn that returns a stream closes its span before the stream is ever polled. Add a small pin-projected `InstrumentedStream` adapter and `InstrumentStream` extension trait that enters a `Span` on every `poll_next`, letting a span cover the full lifetime of a streamed query. pin-project-lite is already in the dependency tree. Co-Authored-By: Claude Opus 4.8 (1M context) --- sqlx-core/Cargo.toml | 1 + sqlx-core/src/instrument_stream.rs | 52 ++++++++++++++++++++++++++++++ sqlx-core/src/lib.rs | 1 + 3 files changed, 54 insertions(+) create mode 100644 sqlx-core/src/instrument_stream.rs diff --git a/sqlx-core/Cargo.toml b/sqlx-core/Cargo.toml index 90ed446b4b..c680e0ad01 100644 --- a/sqlx-core/Cargo.toml +++ b/sqlx-core/Cargo.toml @@ -88,6 +88,7 @@ futures-util = { version = "0.3.32", default-features = false, features = ["allo log = { version = "0.4.18", default-features = false } memchr = { version = "2.5.0", default-features = false } percent-encoding = "2.3.0" +pin-project-lite = "0.2.16" serde = { version = "1.0.219", features = ["derive", "rc"], optional = true } serde_json = { version = "1.0.142", features = ["raw_value"], optional = true } toml = { version = "0.8.16", optional = true } diff --git a/sqlx-core/src/instrument_stream.rs b/sqlx-core/src/instrument_stream.rs new file mode 100644 index 0000000000..2aa221d20d --- /dev/null +++ b/sqlx-core/src/instrument_stream.rs @@ -0,0 +1,52 @@ +//! Attach a [`tracing::Span`] to a [`Stream`] so the span stays open for the +//! whole stream rather than just its construction. + +use std::pin::Pin; +use std::task::{Context, Poll}; + +use futures_core::Stream; +use pin_project_lite::pin_project; +use tracing::Span; + +pin_project! { + /// A [`Stream`] adapter that enters `span` for the duration of every + /// [`poll_next`](Stream::poll_next). + /// + /// A plain `#[tracing::instrument]` on an `async fn` that *returns* a stream + /// only keeps its span entered while the stream is being built; the span is + /// then closed before the caller ever polls the stream. Wrapping the + /// returned stream with this adapter instead keeps the span open across row + /// fetching, so e.g. a `sqlx::query` span measures the whole query rather + /// than just its setup. + pub struct InstrumentedStream { + #[pin] + stream: S, + span: Span, + } +} + +impl Stream for InstrumentedStream { + type Item = S::Item; + + fn poll_next(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll> { + let this = self.project(); + let _entered = this.span.enter(); + this.stream.poll_next(cx) + } + + fn size_hint(&self) -> (usize, Option) { + self.stream.size_hint() + } +} + +/// Extension trait for attaching a [`Span`] to a [`Stream`]. +pub trait InstrumentStream: Stream + Sized { + /// Wrap this stream so that `span` is entered every time it is polled. + /// + /// See [`InstrumentedStream`]. + fn instrument_stream(self, span: Span) -> InstrumentedStream { + InstrumentedStream { stream: self, span } + } +} + +impl InstrumentStream for S {} diff --git a/sqlx-core/src/lib.rs b/sqlx-core/src/lib.rs index 494c41e9bf..b48ad5b333 100644 --- a/sqlx-core/src/lib.rs +++ b/sqlx-core/src/lib.rs @@ -66,6 +66,7 @@ pub mod describe; pub mod executor; pub mod from_row; pub mod fs; +pub mod instrument_stream; pub mod io; pub mod logger; pub mod net; From 10f89c9494e984548e8207d1c0b00f9c5337615f Mon Sep 17 00:00:00 2001 From: Graham Christensen Date: Tue, 14 Jul 2026 21:23:46 -0400 Subject: [PATCH 7/9] fix(tracing): keep postgres.run span open across row fetch Review finding #1: postgres.run used `#[tracing::instrument]`, whose span closed as soon as `run()` returned the result stream, so it covered only query setup (bind/execute/flush) and not the row fetching that dominates query latency. Create the span explicitly and attach it to the returned stream via `instrument_stream` so it stays entered while rows are received. Co-Authored-By: Claude Opus 4.8 (1M context) --- sqlx-postgres/src/connection/executor.rs | 26 +++++++++++++----------- 1 file changed, 14 insertions(+), 12 deletions(-) diff --git a/sqlx-postgres/src/connection/executor.rs b/sqlx-postgres/src/connection/executor.rs index 22f10ac227..e45417f1ce 100644 --- a/sqlx-postgres/src/connection/executor.rs +++ b/sqlx-postgres/src/connection/executor.rs @@ -16,6 +16,7 @@ use futures_core::stream::BoxStream; use futures_core::Stream; use futures_util::TryStreamExt; use sqlx_core::arguments::Arguments; +use sqlx_core::instrument_stream::InstrumentStream; use sqlx_core::sql_str::SqlStr; use sqlx_core::Either; use std::{pin::pin, sync::Arc}; @@ -210,17 +211,6 @@ impl PgConnection { Ok(statement) } - #[tracing::instrument( - target = "sqlx::query", - name = "postgres.run", - skip_all, - fields( - db.system = "postgresql", - db.operation.parameters = arguments.as_ref().map_or(0, |a| a.len()), - db.postgresql.prepared = arguments.is_some(), - ), - level = "debug", - )] pub(crate) async fn run<'e, 'c: 'e, 'q: 'e>( &'c mut self, query: SqlStr, @@ -228,6 +218,17 @@ impl PgConnection { persistent: bool, metadata_opt: Option>, ) -> Result, Error>> + 'e, Error> { + // The span is attached to the returned stream (see `instrument_stream`) + // rather than via `#[tracing::instrument]` so it stays open while rows + // are fetched, not just while the query is set up. + let span = tracing::debug_span!( + target: "sqlx::query", + "postgres.run", + db.system = "postgresql", + db.operation.parameters = arguments.as_ref().map_or(0, |a| a.len()), + db.postgresql.prepared = arguments.is_some(), + ); + let mut logger = QueryLogger::new(query, self.inner.log_settings.clone()); let sql = logger.sql().as_str(); @@ -397,7 +398,8 @@ impl PgConnection { } Ok(()) - }) + } + .instrument_stream(span)) } } From 37438bdb7a668ffda5492fc792806fae1c1c4f4d Mon Sep 17 00:00:00 2001 From: Graham Christensen Date: Tue, 14 Jul 2026 21:23:46 -0400 Subject: [PATCH 8/9] fix(tracing): keep mysql.run span open across row fetch MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Review finding #1: same defect as postgres.run — the mysql.run `#[tracing::instrument]` span closed when `run()` returned the result stream, so it did not cover row fetching. Create the span explicitly and attach it to the returned stream via `instrument_stream`. Co-Authored-By: Claude Opus 4.8 (1M context) --- sqlx-mysql/src/connection/executor.rs | 26 ++++++++++++++------------ 1 file changed, 14 insertions(+), 12 deletions(-) diff --git a/sqlx-mysql/src/connection/executor.rs b/sqlx-mysql/src/connection/executor.rs index 633f8fc91c..1ab00b22ab 100644 --- a/sqlx-mysql/src/connection/executor.rs +++ b/sqlx-mysql/src/connection/executor.rs @@ -23,6 +23,7 @@ use futures_core::Stream; use futures_util::TryStreamExt; use sqlx_core::arguments::Arguments as _; use sqlx_core::column::{ColumnOrigin, TableColumn}; +use sqlx_core::instrument_stream::InstrumentStream; use sqlx_core::sql_str::SqlStr; use std::{pin::pin, sync::Arc}; @@ -111,17 +112,6 @@ impl MySqlConnection { } #[allow(clippy::needless_lifetimes)] - #[tracing::instrument( - target = "sqlx::query", - name = "mysql.run", - skip_all, - fields( - db.system = "mysql", - db.operation.parameters = arguments.as_ref().map_or(0, |a| a.len()), - db.mysql.prepared = arguments.is_some(), - ), - level = "debug", - )] pub(crate) async fn run<'e, 'c: 'e, 'q: 'e>( &'c mut self, sql: SqlStr, @@ -129,6 +119,17 @@ impl MySqlConnection { persistent: bool, ) -> Result, Error>> + 'e, Error> { + // The span is attached to the returned stream (see `instrument_stream`) + // rather than via `#[tracing::instrument]` so it stays open while rows + // are fetched, not just while the query is set up. + let span = tracing::debug_span!( + target: "sqlx::query", + "mysql.run", + db.system = "mysql", + db.operation.parameters = arguments.as_ref().map_or(0, |a| a.len()), + db.mysql.prepared = arguments.is_some(), + ); + let mut logger = QueryLogger::new(sql, self.inner.log_settings.clone()); self.inner.stream.wait_until_ready().await?; @@ -288,7 +289,8 @@ impl MySqlConnection { r#yield!(v); } } - }) + } + .instrument_stream(span)) } } From 91b5f5e00bafa4d8df3625319277ddc3a04735d4 Mon Sep 17 00:00:00 2001 From: Graham Christensen Date: Thu, 13 Aug 2026 09:12:40 -0400 Subject: [PATCH 9/9] Allow type complexity on sqlite.execute, since tracing puts its complexity over-budget --- sqlx-sqlite/src/connection/worker.rs | 1 + 1 file changed, 1 insertion(+) diff --git a/sqlx-sqlite/src/connection/worker.rs b/sqlx-sqlite/src/connection/worker.rs index 7a1490ed4f..dc0de90afc 100644 --- a/sqlx-sqlite/src/connection/worker.rs +++ b/sqlx-sqlite/src/connection/worker.rs @@ -376,6 +376,7 @@ impl ConnectionWorker { .await? } + #[allow(clippy::type_complexity)] #[tracing::instrument( target = "sqlx::query", name = "sqlite.execute",