diff --git a/src/common/time/src/timestamp.rs b/src/common/time/src/timestamp.rs index 7daa793baf..88cac7f4fc 100644 --- a/src/common/time/src/timestamp.rs +++ b/src/common/time/src/timestamp.rs @@ -483,6 +483,31 @@ impl Timestamp { ParseTimestampSnafu { raw: s }.fail() } + /// Interprets a timezone-less [`NaiveDateTime`] in the given timezone and + /// returns the corresponding timestamp. + /// + /// Datetimes that fall into a DST gap (a local time that does not exist) + /// are rejected, while ambiguous datetimes (from a repeated local time) + /// are resolved to the earlier instant. This policy is shared with + /// [`Timestamp::from_str`] so that the text and binary protocols interpret + /// datetimes consistently. + pub fn from_naive_datetime( + datetime: NaiveDateTime, + timezone: &Timezone, + ) -> crate::error::Result { + match datetime_to_utc(&datetime, timezone) { + LocalResult::Single(utc) | LocalResult::Ambiguous(utc, _) => { + Timestamp::from_chrono_datetime(utc).context(ParseTimestampSnafu { + raw: format!("{datetime} (timezone {timezone})"), + }) + } + LocalResult::None => ParseTimestampSnafu { + raw: format!("{datetime} (timezone {timezone})"), + } + .fail(), + } + } + pub fn negative(mut self) -> Self { self.value = -self.value; self @@ -531,12 +556,8 @@ fn naive_datetime_to_timestamp( .context(ParseTimestampSnafu { raw: s }); }; - match datetime_to_utc(&datetime, timezone) { - LocalResult::None => ParseTimestampSnafu { raw: s }.fail(), - LocalResult::Single(utc) | LocalResult::Ambiguous(utc, _) => { - Timestamp::from_chrono_datetime(utc).context(ParseTimestampSnafu { raw: s }) - } - } + Timestamp::from_naive_datetime(datetime, timezone) + .map_err(|_| ParseTimestampSnafu { raw: s }.build()) } impl From for Timestamp { @@ -919,6 +940,48 @@ mod tests { ); } + #[test] + fn test_from_naive_datetime() { + let datetime = NaiveDate::from_ymd_opt(2026, 8, 13) + .unwrap() + .and_hms_opt(8, 0, 0) + .unwrap(); + + // A fixed-offset timezone shifts the datetime by a constant amount. + let shanghai = Timezone::from_tz_string("Asia/Shanghai").unwrap(); + assert_eq!( + "2026-08-13 00:00:00", + Timestamp::from_naive_datetime(datetime, &shanghai) + .unwrap() + .to_chrono_datetime() + .unwrap() + .to_string() + ); + + // 2026-03-08 02:30 does not exist in America/New_York (DST gap). + let new_york = Timezone::from_tz_string("America/New_York").unwrap(); + let gap = NaiveDate::from_ymd_opt(2026, 3, 8) + .unwrap() + .and_hms_opt(2, 30, 0) + .unwrap(); + assert!(Timestamp::from_naive_datetime(gap, &new_york).is_err()); + + // 2026-11-01 01:30 is ambiguous in America/New_York; picks the first + // instant (EDT, UTC-4). + let ambiguous = NaiveDate::from_ymd_opt(2026, 11, 1) + .unwrap() + .and_hms_opt(1, 30, 0) + .unwrap(); + assert_eq!( + "2026-11-01 05:30:00", + Timestamp::from_naive_datetime(ambiguous, &new_york) + .unwrap() + .to_chrono_datetime() + .unwrap() + .to_string() + ); + } + #[test] fn test_to_iso8601_string() { set_default_timezone(Some("Asia/Shanghai")).unwrap(); diff --git a/src/servers/src/mysql/handler.rs b/src/servers/src/mysql/handler.rs index b0eb33f8fd..30579d8846 100644 --- a/src/servers/src/mysql/handler.rs +++ b/src/servers/src/mysql/handler.rs @@ -25,6 +25,7 @@ use common_catalog::parse_optional_catalog_and_schema_from_db_string; use common_error::ext::ErrorExt; use common_query::Output; use common_telemetry::{debug, error, tracing, warn}; +use common_time::Timezone; use datafusion_common::ParamValues; use datafusion_expr::LogicalPlan; use datatypes::prelude::ConcreteDataType; @@ -289,12 +290,13 @@ impl MysqlInstanceShim { .fail(); } + let timezone = query_ctx.timezone(); let replaced_plan = match params { Params::ProtocolParams(params) => { - replace_params_with_values(&plan, param_types, ¶ms) + replace_params_with_values(&plan, param_types, ¶ms, &timezone) } Params::CliParams(params) => { - replace_params_with_exprs(&plan, param_types, ¶ms) + replace_params_with_exprs(&plan, param_types, ¶ms, &timezone) } }?; @@ -807,6 +809,7 @@ fn replace_params_with_values( plan: &LogicalPlan, param_types: HashMap>, params: &[ParamValue], + timezone: &Timezone, ) -> Result { debug_assert_eq!(param_types.len(), params.len()); @@ -824,7 +827,7 @@ fn replace_params_with_values( for (i, param) in params.iter().enumerate() { if let Some(Some(t)) = param_types.get(&format_placeholder(i + 1)) { - let value = helper::convert_value(param, t)?; + let value = helper::convert_value(param, t, timezone)?; values.push(value.into()); } @@ -839,6 +842,7 @@ fn replace_params_with_exprs( plan: &LogicalPlan, param_types: HashMap>, params: &[sql::ast::Expr], + timezone: &Timezone, ) -> Result { debug_assert_eq!(param_types.len(), params.len()); @@ -853,7 +857,7 @@ fn replace_params_with_exprs( for (i, param) in params.iter().enumerate() { if let Some(Some(t)) = param_types.get(&format_placeholder(i + 1)) { - let value = helper::convert_expr_to_scalar_value(param, t)?; + let value = helper::convert_expr_to_scalar_value(param, t, timezone)?; values.push(value.into()); } diff --git a/src/servers/src/mysql/helper.rs b/src/servers/src/mysql/helper.rs index a4f289cdfb..f9ab532592 100644 --- a/src/servers/src/mysql/helper.rs +++ b/src/servers/src/mysql/helper.rs @@ -15,10 +15,11 @@ use std::ops::ControlFlow; use std::time::Duration; -use chrono::NaiveDate; +use chrono::{NaiveDate, NaiveDateTime}; use common_query::prelude::ScalarValue; use common_sql::convert::sql_value_to_value; -use common_time::{Date, Timestamp}; +use common_time::timestamp::TimeUnit; +use common_time::{Date, Timestamp, Timezone}; use datatypes::prelude::{ConcreteDataType, DataType}; use datatypes::schema::ColumnSchema; use datatypes::types::TimestampType; @@ -137,11 +138,15 @@ where /// Convert [`ParamValue`] into [`Value`] according to param type. /// It will try it's best to do type conversions if possible -pub fn convert_value(param: &ParamValue, t: &ConcreteDataType) -> Result { +pub fn convert_value( + param: &ParamValue, + t: &ConcreteDataType, + timezone: &Timezone, +) -> Result { if let ConcreteDataType::Dictionary(dictionary) = t { return Ok(ScalarValue::Dictionary( Box::new(dictionary.key_type().as_arrow_type()), - Box::new(convert_value(param, dictionary.value_type())?), + Box::new(convert_value(param, dictionary.value_type(), timezone)?), )); } @@ -221,7 +226,9 @@ pub fn convert_value(param: &ParamValue, t: &ConcreteDataType) -> Result Ok(ScalarValue::Binary(Some(b.to_vec()))), - ConcreteDataType::Timestamp(ts_type) => convert_bytes_to_timestamp(b, ts_type), + ConcreteDataType::Timestamp(ts_type) => { + convert_bytes_to_timestamp(b, ts_type, timezone) + } ConcreteDataType::Date(_) => convert_bytes_to_date(b), _ => error::PreparedStmtTypeMismatchSnafu { expected: t, @@ -233,29 +240,26 @@ pub fn convert_value(param: &ParamValue, t: &ConcreteDataType) -> Result { - let timestamp_millis = to_naive_datetime(param.value) - .map_err(|e| { + ValueInner::Datetime(_) => match t { + ConcreteDataType::Timestamp(_) => { + let datetime = to_naive_datetime(param.value).map_err(|e| { error::MysqlValueConversionSnafu { err_msg: e.to_string(), } .build() - })? - .and_utc() - .timestamp_millis(); - - match t { - ConcreteDataType::Timestamp(_) => Ok(ScalarValue::TimestampMillisecond( + })?; + let timestamp_millis = datetime_to_timestamp_millis(datetime, timezone)?; + Ok(ScalarValue::TimestampMillisecond( Some(timestamp_millis), None, - )), - _ => error::PreparedStmtTypeMismatchSnafu { - expected: t, - actual: param.coltype, - } - .fail(), + )) } - } + _ => error::PreparedStmtTypeMismatchSnafu { + expected: t, + actual: param.coltype, + } + .fail(), + }, ValueInner::Time(_) => Ok(ScalarValue::Time64Nanosecond(Some( Duration::from(param.value).as_millis() as i64, ))), @@ -264,13 +268,18 @@ pub fn convert_value(param: &ParamValue, t: &ConcreteDataType) -> Result Result { +pub fn convert_expr_to_scalar_value( + param: &Expr, + t: &ConcreteDataType, + timezone: &Timezone, +) -> Result { if let ConcreteDataType::Dictionary(dictionary) = t { return Ok(ScalarValue::Dictionary( Box::new(dictionary.key_type().as_arrow_type()), Box::new(convert_expr_to_scalar_value( param, dictionary.value_type(), + timezone, )?), )); } @@ -278,7 +287,7 @@ pub fn convert_expr_to_scalar_value(param: &Expr, t: &ConcreteDataType) -> Resul let column_schema = ColumnSchema::new("", t.clone(), true); match param { Expr::Value(v) => { - let v = sql_value_to_value(&column_schema, &v.value, None, None, true); + let v = sql_value_to_value(&column_schema, &v.value, Some(timezone), None, true); match v { Ok(v) => v .try_to_scalar_value(t) @@ -290,7 +299,7 @@ pub fn convert_expr_to_scalar_value(param: &Expr, t: &ConcreteDataType) -> Resul } } Expr::UnaryOp { op, expr } if let Expr::Value(v) = &**expr => { - let v = sql_value_to_value(&column_schema, &v.value, None, Some(*op), true); + let v = sql_value_to_value(&column_schema, &v.value, Some(timezone), Some(*op), true); match v { Ok(v) => v .try_to_scalar_value(t) @@ -308,8 +317,32 @@ pub fn convert_expr_to_scalar_value(param: &Expr, t: &ConcreteDataType) -> Resul } } -fn convert_bytes_to_timestamp(bytes: &[u8], ts_type: &TimestampType) -> Result { - let ts = Timestamp::from_str_utc(&String::from_utf8_lossy(bytes)) +/// Interprets a timezone-less datetime in the given timezone and returns the +/// corresponding epoch timestamp in milliseconds. +fn datetime_to_timestamp_millis(datetime: NaiveDateTime, timezone: &Timezone) -> Result { + let ts = Timestamp::from_naive_datetime(datetime, timezone) + .map_err(|e| { + error::MysqlValueConversionSnafu { + err_msg: e.to_string(), + } + .build() + })? + .convert_to(TimeUnit::Millisecond) + .ok_or_else(|| { + error::MysqlValueConversionSnafu { + err_msg: "Overflow when converting datetime to milliseconds".to_string(), + } + .build() + })?; + Ok(ts.value()) +} + +fn convert_bytes_to_timestamp( + bytes: &[u8], + ts_type: &TimestampType, + timezone: &Timezone, +) -> Result { + let ts = Timestamp::from_str(&String::from_utf8_lossy(bytes), Some(timezone)) .map_err(|e| { error::MysqlValueConversionSnafu { err_msg: e.to_string(), @@ -446,19 +479,20 @@ mod tests { #[test] fn test_convert_expr_to_scalar_value() { + let utc = Timezone::from_tz_string("UTC").unwrap(); let expr = Expr::Value(ValueExpr::Number("123".to_string(), false).into()); let t = ConcreteDataType::int32_datatype(); - let v = convert_expr_to_scalar_value(&expr, &t).unwrap(); + let v = convert_expr_to_scalar_value(&expr, &t, &utc).unwrap(); assert_eq!(ScalarValue::Int32(Some(123)), v); let expr = Expr::Value(ValueExpr::Number("123.456789".to_string(), false).into()); let t = ConcreteDataType::float64_datatype(); - let v = convert_expr_to_scalar_value(&expr, &t).unwrap(); + let v = convert_expr_to_scalar_value(&expr, &t, &utc).unwrap(); assert_eq!(ScalarValue::Float64(Some(123.456789)), v); let expr = Expr::Value(ValueExpr::SingleQuotedString("2001-01-02".to_string()).into()); let t = ConcreteDataType::date_datatype(); - let v = convert_expr_to_scalar_value(&expr, &t).unwrap(); + let v = convert_expr_to_scalar_value(&expr, &t, &utc).unwrap(); let scalar_v = ScalarValue::Utf8(Some("2001-01-02".to_string())) .cast_to(&arrow_schema::DataType::Date32) .unwrap(); @@ -467,7 +501,7 @@ mod tests { let expr = Expr::Value(ValueExpr::SingleQuotedString("2001-01-02 03:04:05".to_string()).into()); let t = ConcreteDataType::timestamp_microsecond_datatype(); - let v = convert_expr_to_scalar_value(&expr, &t).unwrap(); + let v = convert_expr_to_scalar_value(&expr, &t, &utc).unwrap(); let scalar_v = ScalarValue::Utf8(Some("2001-01-02 03:04:05".to_string())) .cast_to(&arrow_schema::DataType::Timestamp( arrow_schema::TimeUnit::Microsecond, @@ -478,14 +512,14 @@ mod tests { let expr = Expr::Value(ValueExpr::SingleQuotedString("hello".to_string()).into()); let t = ConcreteDataType::string_datatype(); - let v = convert_expr_to_scalar_value(&expr, &t).unwrap(); + let v = convert_expr_to_scalar_value(&expr, &t, &utc).unwrap(); assert_eq!(ScalarValue::Utf8(Some("hello".to_string())), v); let t = ConcreteDataType::dictionary_datatype( ConcreteDataType::uint32_datatype(), ConcreteDataType::string_datatype(), ); - let v = convert_expr_to_scalar_value(&expr, &t).unwrap(); + let v = convert_expr_to_scalar_value(&expr, &t, &utc).unwrap(); assert_eq!( ScalarValue::Dictionary( Box::new(arrow_schema::DataType::UInt32), @@ -496,7 +530,7 @@ mod tests { let expr = Expr::Value(ValueExpr::Null.into()); let t = ConcreteDataType::time_microsecond_datatype(); - let v = convert_expr_to_scalar_value(&expr, &t).unwrap(); + let v = convert_expr_to_scalar_value(&expr, &t, &utc).unwrap(); assert_eq!(ScalarValue::Time64Microsecond(None), v); } @@ -577,12 +611,81 @@ mod tests { ), ]; + let utc = Timezone::from_tz_string("UTC").unwrap(); for (input, ts_type, expected) in test_cases { - let result = convert_bytes_to_timestamp(input.as_bytes(), &ts_type).unwrap(); + let result = convert_bytes_to_timestamp(input.as_bytes(), &ts_type, &utc).unwrap(); assert_eq!(result, expected); } } + fn utc_millis(year: i32, month: u32, day: u32, hour: u32, minute: u32, second: u32) -> i64 { + NaiveDate::from_ymd_opt(year, month, day) + .unwrap() + .and_hms_opt(hour, minute, second) + .unwrap() + .and_utc() + .timestamp_millis() + } + + #[test] + fn test_datetime_to_timestamp_millis() { + let datetime = NaiveDate::from_ymd_opt(2026, 8, 13) + .unwrap() + .and_hms_opt(8, 0, 0) + .unwrap(); + + let utc = Timezone::from_tz_string("UTC").unwrap(); + assert_eq!( + datetime_to_timestamp_millis(datetime, &utc).unwrap(), + utc_millis(2026, 8, 13, 8, 0, 0) + ); + + let shanghai = Timezone::from_tz_string("Asia/Shanghai").unwrap(); + assert_eq!( + datetime_to_timestamp_millis(datetime, &shanghai).unwrap(), + utc_millis(2026, 8, 13, 0, 0, 0) + ); + + let minus_7 = Timezone::from_tz_string("-07:00").unwrap(); + assert_eq!( + datetime_to_timestamp_millis(datetime, &minus_7).unwrap(), + utc_millis(2026, 8, 13, 15, 0, 0) + ); + + // 2026-03-08 02:30 does not exist in America/New_York (DST gap). + let new_york = Timezone::from_tz_string("America/New_York").unwrap(); + let gap = NaiveDate::from_ymd_opt(2026, 3, 8) + .unwrap() + .and_hms_opt(2, 30, 0) + .unwrap(); + assert!(datetime_to_timestamp_millis(gap, &new_york).is_err()); + + // 2026-11-01 01:30 is ambiguous in America/New_York; picks the first (EDT, UTC-4). + let ambiguous = NaiveDate::from_ymd_opt(2026, 11, 1) + .unwrap() + .and_hms_opt(1, 30, 0) + .unwrap(); + assert_eq!( + datetime_to_timestamp_millis(ambiguous, &new_york).unwrap(), + utc_millis(2026, 11, 1, 5, 30, 0) + ); + } + + #[test] + fn test_convert_bytes_to_timestamp_with_timezone() { + let shanghai = Timezone::from_tz_string("Asia/Shanghai").unwrap(); + let result = convert_bytes_to_timestamp( + "2026-08-13 08:00:00".as_bytes(), + &TimestampType::Millisecond(TimestampMillisecondType), + &shanghai, + ) + .unwrap(); + assert_eq!( + result, + ScalarValue::TimestampMillisecond(Some(utc_millis(2026, 8, 13, 0, 0, 0)), None) + ); + } + #[test] fn test_convert_bytes_to_date() { let test_cases = vec![ diff --git a/tests-integration/tests/sql.rs b/tests-integration/tests/sql.rs index e46a8f58d1..32a2881dd7 100644 --- a/tests-integration/tests/sql.rs +++ b/tests-integration/tests/sql.rs @@ -89,6 +89,7 @@ macro_rules! sql_tests { test_postgres_explain_bind_parameter, test_postgres_array_types, test_mysql_prepare_stmt_insert_timestamp, + test_mysql_prepare_stmt_timezone, test_mysql_federated_prepare_stmt, test_declare_fetch_close_cursor, test_alter_update_on, @@ -1779,6 +1780,77 @@ pub async fn test_mysql_prepare_stmt_insert_timestamp(store_type: StorageType) { guard.remove_all().await; } +pub async fn test_mysql_prepare_stmt_timezone(store_type: StorageType) { + let (mut guard, server) = + setup_mysql_server(store_type, "test_mysql_prepare_stmt_timezone").await; + let addr = server.bind_addr().unwrap().to_string(); + + let mut conn = MySqlConnection::connect(&format!("mysql://{addr}/public")) + .await + .unwrap(); + + conn.execute("create table demo(i bigint, ts timestamp time index)") + .await + .unwrap(); + conn.execute("SET time_zone = 'Asia/Shanghai'") + .await + .unwrap(); + + // Server-side prepared statement: the binary DATETIME parameter must be + // interpreted in the session timezone. + sqlx::query("insert into demo values(?, ?)") + .bind(1) + .bind( + NaiveDate::from_ymd_opt(2026, 8, 13) + .and_then(|x| x.and_hms_opt(8, 0, 0)) + .unwrap(), + ) + .execute(&mut conn) + .await + .unwrap(); + + // Text protocol with an equivalent timezone-less literal one hour later. + sqlx::query("insert into demo values(2, '2026-08-13 09:00:00')") + .execute(&mut conn) + .await + .unwrap(); + + // Timestamps are read back in the session timezone, and sqlx parses the + // timezone-less wire representation as UTC. + let rows = sqlx::query("select i, ts from demo order by i") + .fetch_all(&mut conn) + .await + .unwrap(); + assert_eq!(rows.len(), 2); + let ts: DateTime = rows[0].get("ts"); + assert_eq!(ts.to_string(), "2026-08-13 08:00:00 UTC"); + let ts: DateTime = rows[1].get("ts"); + assert_eq!(ts.to_string(), "2026-08-13 09:00:00 UTC"); + + // The prepared predicate compares the same instant as the text literal. + let rows = sqlx::query("select i from demo where ts = ? order by i") + .bind( + NaiveDate::from_ymd_opt(2026, 8, 13) + .and_then(|x| x.and_hms_opt(8, 0, 0)) + .unwrap(), + ) + .fetch_all(&mut conn) + .await + .unwrap(); + assert_eq!(rows.len(), 1); + assert_eq!(rows[0].get::("i"), 1); + + let rows = sqlx::query("select i from demo where ts = '2026-08-13 08:00:00'") + .fetch_all(&mut conn) + .await + .unwrap(); + assert_eq!(rows.len(), 1); + assert_eq!(rows[0].get::("i"), 1); + + let _ = server.shutdown().await; + guard.remove_all().await; +} + pub async fn test_mysql_federated_prepare_stmt(store_type: StorageType) { common_telemetry::init_default_ut_logging();