Skip to content

Commit 7bc3133

Browse files
committed
feat: add an optional query logger (using log)
1 parent d06de14 commit 7bc3133

8 files changed

Lines changed: 111 additions & 65 deletions

File tree

sqlx-core/src/lib.rs

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -62,6 +62,7 @@ pub mod done;
6262
pub mod executor;
6363
pub mod from_row;
6464
mod io;
65+
mod logger;
6566
mod net;
6667
pub mod query_as;
6768
pub mod query_scalar;

sqlx-core/src/logger.rs

Lines changed: 79 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,79 @@
1+
use log::Level;
2+
use std::time::{Duration, Instant};
3+
4+
const SLOW_QUERY_THRESHOLD: Duration = Duration::from_secs(1);
5+
6+
pub(crate) struct QueryLogger<'q> {
7+
sql: &'q str,
8+
rows: usize,
9+
start: Instant,
10+
}
11+
12+
impl<'q> QueryLogger<'q> {
13+
pub(crate) fn new(sql: &'q str) -> Self {
14+
Self {
15+
sql,
16+
rows: 0,
17+
start: Instant::now(),
18+
}
19+
}
20+
21+
pub(crate) fn increment_rows(&mut self) {
22+
self.rows += 1;
23+
}
24+
25+
pub(crate) fn finish(&self) {
26+
let elapsed = self.start.elapsed();
27+
28+
let lvl = if elapsed >= SLOW_QUERY_THRESHOLD {
29+
Level::Warn
30+
} else {
31+
Level::Info
32+
};
33+
34+
if lvl <= log::STATIC_MAX_LEVEL && lvl <= log::max_level() {
35+
let mut summary = parse_query_summary(&self.sql);
36+
37+
let sql = if summary != self.sql {
38+
summary.push_str(" …");
39+
format!(
40+
"\n\n{}\n",
41+
sqlformat::format(
42+
&self.sql,
43+
&sqlformat::QueryParams::None,
44+
sqlformat::FormatOptions::default()
45+
)
46+
)
47+
} else {
48+
String::new()
49+
};
50+
51+
let rows = self.rows;
52+
53+
log::logger().log(
54+
&log::Record::builder()
55+
.args(format_args!(
56+
"{}; rows: {}, elapsed: {:.3?}{}",
57+
summary, rows, elapsed, sql
58+
))
59+
.level(lvl)
60+
.module_path_static(Some("sqlx::query"))
61+
.build(),
62+
);
63+
}
64+
}
65+
}
66+
67+
impl<'q> Drop for QueryLogger<'q> {
68+
fn drop(&mut self) {
69+
self.finish();
70+
}
71+
}
72+
73+
fn parse_query_summary(sql: &str) -> String {
74+
// For now, just take the first 4 words
75+
sql.split_whitespace()
76+
.take(4)
77+
.collect::<Vec<&str>>()
78+
.join(" ")
79+
}

sqlx-core/src/logging.rs

Lines changed: 0 additions & 49 deletions
This file was deleted.

sqlx-core/src/mssql/connection/executor.rs

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,7 @@
11
use crate::describe::Describe;
22
use crate::error::Error;
33
use crate::executor::{Execute, Executor};
4+
use crate::logger::QueryLogger;
45
use crate::mssql::connection::prepare::prepare;
56
use crate::mssql::protocol::col_meta_data::Flags;
67
use crate::mssql::protocol::done::Status;
@@ -77,6 +78,7 @@ impl<'c> Executor<'c> for &'c mut MssqlConnection {
7778
{
7879
let sql = query.sql();
7980
let arguments = query.take_arguments();
81+
let mut logger = QueryLogger::new(sql);
8082

8183
Box::pin(try_stream! {
8284
self.run(sql, arguments).await?;
@@ -89,6 +91,8 @@ impl<'c> Executor<'c> for &'c mut MssqlConnection {
8991
let columns = Arc::clone(&self.stream.columns);
9092
let column_names = Arc::clone(&self.stream.column_names);
9193

94+
logger.increment_rows();
95+
9296
r#yield!(Either::Right(MssqlRow { row, column_names, columns }));
9397
}
9498

sqlx-core/src/mysql/connection/executor.rs

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -3,6 +3,7 @@ use crate::describe::Describe;
33
use crate::error::Error;
44
use crate::executor::{Execute, Executor};
55
use crate::ext::ustr::UStr;
6+
use crate::logger::QueryLogger;
67
use crate::mysql::connection::stream::Busy;
78
use crate::mysql::io::MySqlBufExt;
89
use crate::mysql::protocol::response::Status;
@@ -88,6 +89,8 @@ impl MySqlConnection {
8889
arguments: Option<MySqlArguments>,
8990
persistent: bool,
9091
) -> Result<impl Stream<Item = Result<Either<MySqlDone, MySqlRow>, Error>> + 'e, Error> {
92+
let mut logger = QueryLogger::new(sql);
93+
9194
self.stream.wait_until_ready().await?;
9295
self.stream.busy = Busy::Result;
9396

@@ -195,6 +198,8 @@ impl MySqlConnection {
195198
column_names: Arc::clone(&column_names),
196199
});
197200

201+
logger.increment_rows();
202+
198203
r#yield!(v);
199204
}
200205
}

sqlx-core/src/postgres/connection/executor.rs

Lines changed: 9 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,7 @@
11
use crate::describe::Describe;
22
use crate::error::Error;
33
use crate::executor::{Execute, Executor};
4+
use crate::logger::QueryLogger;
45
use crate::postgres::message::{
56
self, Bind, Close, CommandComplete, DataRow, MessageFormat, ParameterDescription, Parse, Query,
67
RowDescription,
@@ -190,14 +191,16 @@ impl PgConnection {
190191
Ok(statement)
191192
}
192193

193-
async fn run(
194-
&mut self,
195-
query: &str,
194+
async fn run<'e, 'c: 'e, 'q: 'e>(
195+
&'c mut self,
196+
query: &'q str,
196197
arguments: Option<PgArguments>,
197198
limit: u8,
198199
persistent: bool,
199200
metadata_opt: Option<Arc<PgStatementMetadata>>,
200-
) -> Result<impl Stream<Item = Result<Either<PgDone, PgRow>, Error>> + '_, Error> {
201+
) -> Result<impl Stream<Item = Result<Either<PgDone, PgRow>, Error>> + 'e, Error> {
202+
let mut logger = QueryLogger::new(query);
203+
201204
// before we continue, wait until we are "ready" to accept more queries
202205
self.wait_until_ready().await?;
203206

@@ -294,6 +297,8 @@ impl PgConnection {
294297
}
295298

296299
MessageFormat::DataRow => {
300+
logger.increment_rows();
301+
297302
// one of the set of rows returned by a SELECT, FETCH, etc query
298303
let data: DataRow = message.decode()?;
299304
let row = PgRow {

sqlx-core/src/postgres/connection/stream.rs

Lines changed: 9 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -128,18 +128,15 @@ impl PgStream {
128128
};
129129

130130
if lvl <= log::STATIC_MAX_LEVEL && lvl <= log::max_level() {
131-
// let message = notice.message();
132-
// let args = format_args!("{}", notice.message().to_string());
133-
let mut builder = log::Record::builder();
134-
let record = builder
135-
// .args(args)
136-
.level(lvl)
137-
.module_path_static(Some("sqlx::postgres::notice"))
138-
.file_static(Some(file!()))
139-
.line(Some(line!()))
140-
.build();
141-
142-
log::logger().log(&record);
131+
log::logger().log(
132+
&log::Record::builder()
133+
.args(format_args!("{}", notice.message()))
134+
.level(lvl)
135+
.module_path_static(Some("sqlx::postgres::notice"))
136+
.file_static(Some(file!()))
137+
.line(Some(line!()))
138+
.build(),
139+
);
143140
}
144141

145142
continue;

sqlx-core/src/sqlite/connection/executor.rs

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -2,6 +2,7 @@ use crate::common::StatementCache;
22
use crate::describe::Describe;
33
use crate::error::Error;
44
use crate::executor::{Execute, Executor};
5+
use crate::logger::QueryLogger;
56
use crate::sqlite::connection::describe::describe;
67
use crate::sqlite::statement::{StatementHandle, VirtualStatement};
78
use crate::sqlite::{
@@ -71,6 +72,7 @@ impl<'c> Executor<'c> for &'c mut SqliteConnection {
7172
E: Execute<'q, Self::Database>,
7273
{
7374
let sql = query.sql();
75+
let mut logger = QueryLogger::new(sql);
7476
let arguments = query.take_arguments();
7577
let persistent = query.persistent() && arguments.is_some();
7678

@@ -135,6 +137,8 @@ impl<'c> Executor<'c> for &'c mut SqliteConnection {
135137
let v = Either::Right(row);
136138
*last_row_values = Some(weak_values_ref);
137139

140+
logger.increment_rows();
141+
138142
r#yield!(v);
139143
}
140144
}

0 commit comments

Comments
 (0)