diff --git a/Cargo.lock b/Cargo.lock index ec477e5c9a..d77ec340a8 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -4628,10 +4628,23 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "e7c1832837b905bbfb5101e07cc24c8deddf52f93225eee6ead5f4d63d53ddcb" dependencies = [ "const-oid 0.9.6", + "der_derive", + "flagset", "pem-rfc7468", "zeroize", ] +[[package]] +name = "der_derive" +version = "0.7.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "8034092389675178f570469e6c3b0465d3d30b4505c294a6550db47f3c17ad18" +dependencies = [ + "proc-macro2", + "quote", + "syn 2.0.117", +] + [[package]] name = "deranged" version = "0.5.8" @@ -5334,6 +5347,12 @@ version = "0.5.7" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "1d674e81391d1e1ab681a28d99df07927c6d4aa5b027d7da16ba32d1d21ecd99" +[[package]] +name = "flagset" +version = "0.4.7" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "b7ac824320a75a52197e8f2d787f6a38b6718bb6897a35142d749af3c0e8f4fe" + [[package]] name = "flatbuffers" version = "25.2.10" @@ -5973,7 +5992,7 @@ dependencies = [ "cfg-if", "js-sys", "libc", - "wasi", + "wasi 0.11.1+wasi-snapshot-preview1", "wasm-bindgen", ] @@ -7866,13 +7885,14 @@ checksum = "f9fbbcab51052fe104eb5e5d351cf728d30a5be1fe14d9be8a3b097481fb97de" [[package]] name = "libredox" -version = "0.1.4" +version = "0.1.20" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "1580801010e535496706ba011c15f8532df6b42297d2e471fec38ceadd8c0638" +checksum = "28d0a00925a9f930d679b6789b721e3a7f9ed110f41b86d2497caa780c3a070a" dependencies = [ "bitflags 2.12.1", "libc", - "redox_syscall 0.5.13", + "plain", + "redox_syscall 0.9.3", ] [[package]] @@ -8543,7 +8563,7 @@ checksum = "02bd0af71c67b473010cbbc60715ee815645a4dc942899111f494b4b737d6fda" dependencies = [ "libc", "log", - "wasi", + "wasi 0.11.1+wasi-snapshot-preview1", "windows-sys 0.61.2", ] @@ -9358,6 +9378,15 @@ dependencies = [ "objc2-encode", ] +[[package]] +name = "objc2-core-foundation" +version = "0.3.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "2a180dd8642fa45cdb7dd721cd4c11b1cadd4929ce112ebd8b9f5803cc79d536" +dependencies = [ + "bitflags 2.12.1", +] + [[package]] name = "objc2-encode" version = "4.1.0" @@ -9374,6 +9403,15 @@ dependencies = [ "objc2", ] +[[package]] +name = "objc2-system-configuration" +version = "0.3.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7216bd11cbda54ccabcab84d523dc93b858ec75ecfb3a7d89513fa22464da396" +dependencies = [ + "objc2-core-foundation", +] + [[package]] name = "object" version = "0.36.7" @@ -10552,16 +10590,7 @@ dependencies = [ "tokio", "tokio-rustls", "tokio-util", - "x509-certificate 0.25.0", -] - -[[package]] -name = "phf" -version = "0.11.3" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "1fd6780a80ae0c52cc120a26a1a42c1ae51b247a253e4e06113d23d2c2edd078" -dependencies = [ - "phf_shared 0.11.3", + "x509-certificate", ] [[package]] @@ -10814,6 +10843,12 @@ version = "0.3.32" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "7edddbd0b52d732b21ad9a5fab5c704c14cd949e5e9a1ec5929a24fded1b904c" +[[package]] +name = "plain" +version = "0.2.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "b4596b6d070b27117e987119b4dac604f3c58cfb0b191112e24771b2faeac1a6" + [[package]] name = "plist" version = "1.7.2" @@ -10924,16 +10959,16 @@ dependencies = [ [[package]] name = "postgres-types" -version = "0.2.9" +version = "0.2.14" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "613283563cd90e1dfc3518d548caee47e0e725455ed619881f5cf21f36de4b48" +checksum = "851ca9db4932932d69f3ea811b1abe63087a0f740a47692619dd40d4899b68be" dependencies = [ "array-init", "bytes", "chrono", "fallible-iterator", "postgres-protocol", - "serde", + "serde_core", "serde_json", ] @@ -11299,7 +11334,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "22505a5c94da8e3b7c2996394d1c933236c4d743e81a410bcca4e6989fc066a4" dependencies = [ "bytes", - "heck 0.4.1", + "heck 0.5.0", "itertools 0.12.1", "log", "multimap", @@ -11319,7 +11354,7 @@ version = "0.14.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "ac6c3320f9abac597dcbc668774ef006702672474aad53c6d596b62e487b40b1" dependencies = [ - "heck 0.4.1", + "heck 0.5.0", "itertools 0.14.0", "log", "multimap", @@ -12087,6 +12122,15 @@ dependencies = [ "bitflags 2.12.1", ] +[[package]] +name = "redox_syscall" +version = "0.9.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d678d17679829e73d371e96880897e98fee2ded7acc0a50bdf8af2affa4b2fe5" +dependencies = [ + "bitflags 2.12.1", +] + [[package]] name = "redox_users" version = "0.4.6" @@ -13785,7 +13829,7 @@ version = "0.8.6" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "1961e2ef424c1424204d3a5d6975f934f56b6d50ff5732382d84ebf460e147f7" dependencies = [ - "heck 0.4.1", + "heck 0.5.0", "proc-macro2", "quote", "syn 2.0.117", @@ -14117,7 +14161,7 @@ dependencies = [ "stringprep", "thiserror 2.0.17", "tracing", - "whoami", + "whoami 1.6.0", ] [[package]] @@ -14156,7 +14200,7 @@ dependencies = [ "stringprep", "thiserror 2.0.17", "tracing", - "whoami", + "whoami 1.6.0", ] [[package]] @@ -15175,6 +15219,27 @@ version = "0.1.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "1f3ccbac311fea05f86f61904b462b55fb3df8837a366dfc601a0161d0532f20" +[[package]] +name = "tls_codec" +version = "0.4.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "0de2e01245e2bb89d6f05801c564fa27624dbd7b1846859876c7dad82e90bf6b" +dependencies = [ + "tls_codec_derive", + "zeroize", +] + +[[package]] +name = "tls_codec_derive" +version = "0.4.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "2d2e76690929402faae40aebdda620a2c0e25dd6d3b9afe48867dfd95991f4bd" +dependencies = [ + "proc-macro2", + "quote", + "syn 2.0.117", +] + [[package]] name = "tokio" version = "1.52.3" @@ -15241,9 +15306,9 @@ dependencies = [ [[package]] name = "tokio-postgres" -version = "0.7.13" +version = "0.7.18" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "6c95d533c83082bb6490e0189acaa0bbeef9084e60471b696ca6988cd0541fb0" +checksum = "a528f7d280f6d5b9cd149635c8705b0dd049754bc67d81d31fa25169a93809d3" dependencies = [ "async-trait", "byteorder", @@ -15254,29 +15319,29 @@ dependencies = [ "log", "parking_lot 0.12.4", "percent-encoding", - "phf 0.11.3", + "phf 0.13.1", "pin-project-lite", "postgres-protocol", "postgres-types", - "rand 0.9.4", - "socket2 0.5.10", + "rand 0.10.1", + "socket2 0.6.4", "tokio", "tokio-util", - "whoami", + "whoami 2.1.3", ] [[package]] name = "tokio-postgres-rustls" -version = "0.12.0" +version = "0.14.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "04fb792ccd6bbcd4bba408eb8a292f70fc4a3589e5d793626f45190e6454b6ab" +checksum = "4c2ad44aa0ae96db89c4742212ed41645b2f597311ff6e1945542a4d9fadc2fb" dependencies = [ - "ring", "rustls", + "sha2 0.11.0", "tokio", "tokio-postgres", "tokio-rustls", - "x509-certificate 0.23.1", + "x509-cert", ] [[package]] @@ -16284,6 +16349,15 @@ version = "0.11.1+wasi-snapshot-preview1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "ccf3ec651a847eb01de73ccad15eb7d99f80485de043efb2f370cd654f4ea44b" +[[package]] +name = "wasi" +version = "0.14.7+wasi-0.2.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "883478de20367e224c0090af9cf5f9fa85bed63a95c1abf3afc5c083ebc06e8c" +dependencies = [ + "wasip2", +] + [[package]] name = "wasip2" version = "1.0.2+wasi-0.2.9" @@ -16308,6 +16382,15 @@ version = "0.1.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "b8dad83b4f25e74f184f64c43b150b91efe7647395b42289f38e50566d82855b" +[[package]] +name = "wasite" +version = "1.0.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "66fe902b4a6b8028a753d5424909b764ccf79b7a209eac9bf97e59cda9f71a42" +dependencies = [ + "wasi 0.14.7+wasi-0.2.4", +] + [[package]] name = "wasm-bindgen" version = "0.2.122" @@ -16530,7 +16613,19 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "6994d13118ab492c3c80c1f81928718159254c53c472bf9ce36f8dae4add02a7" dependencies = [ "redox_syscall 0.5.13", - "wasite", + "wasite 0.1.0", +] + +[[package]] +name = "whoami" +version = "2.1.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "626c4bac6755d76ffc12cb01b2eac751db1996b9e0041de9aa02c8c211ddc82c" +dependencies = [ + "libc", + "libredox", + "objc2-system-configuration", + "wasite 1.0.2", "web-sys", ] @@ -17139,22 +17234,15 @@ dependencies = [ ] [[package]] -name = "x509-certificate" -version = "0.23.1" +name = "x509-cert" +version = "0.2.5" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "66534846dec7a11d7c50a74b7cdb208b9a581cad890b7866430d438455847c85" +checksum = "1301e935010a701ae5f8655edc0ad17c44bad3ac5ce8c39185f75453b720ae94" dependencies = [ - "bcder", - "bytes", - "chrono", + "const-oid 0.9.6", "der", - "hex", - "pem", - "ring", - "signature", "spki", - "thiserror 1.0.69", - "zeroize", + "tls_codec", ] [[package]] diff --git a/Cargo.toml b/Cargo.toml index 7c3fccfbbf..fab7b2aa73 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -255,7 +255,7 @@ strum = { version = "0.27", features = ["derive"] } sysinfo = "0.33" tempfile = "3" tokio = { version = "1.47", features = ["full"] } -tokio-postgres = "0.7" +tokio-postgres = "0.7.18" tokio-rustls = { version = "0.26.2", default-features = false } tokio-stream = "0.1" tokio-util = { version = "0.7", features = ["io-util", "compat"] } diff --git a/src/common/meta/Cargo.toml b/src/common/meta/Cargo.toml index 8e0cbbdf35..8deb3d0081 100644 --- a/src/common/meta/Cargo.toml +++ b/src/common/meta/Cargo.toml @@ -89,7 +89,7 @@ strum.workspace = true table = { workspace = true, features = ["testing"] } tokio.workspace = true tokio-postgres = { workspace = true, optional = true } -tokio-postgres-rustls = { version = "0.12", optional = true } +tokio-postgres-rustls = { version = "0.14", optional = true } tonic.workspace = true tracing.workspace = true typetag.workspace = true diff --git a/src/common/recordbatch/src/cursor.rs b/src/common/recordbatch/src/cursor.rs index a741953ccc..60ae9d54fa 100644 --- a/src/common/recordbatch/src/cursor.rs +++ b/src/common/recordbatch/src/cursor.rs @@ -12,6 +12,7 @@ // See the License for the specific language governing permissions and // limitations under the License. +use datatypes::schema::SchemaRef; use futures::StreamExt; use tokio::sync::Mutex; @@ -28,12 +29,15 @@ struct Inner { /// A cursor on RecordBatchStream that fetches data batch by batch pub struct RecordBatchStreamCursor { + schema: SchemaRef, inner: Mutex, } impl RecordBatchStreamCursor { pub fn new(stream: SendableRecordBatchStream) -> RecordBatchStreamCursor { + let schema = stream.schema(); Self { + schema, inner: Mutex::new(Inner { stream, current_row_index: 0, @@ -43,6 +47,11 @@ impl RecordBatchStreamCursor { } } + /// Returns the schema of the underlying record batch stream. + pub fn schema(&self) -> SchemaRef { + self.schema.clone() + } + /// Take `size` of row from the `RecordBatchStream` and create a new /// `RecordBatch` for these rows. pub async fn take(&self, size: usize) -> Result { diff --git a/src/query/src/analyze.rs b/src/query/src/analyze.rs index ec0abf25ba..b1522fc05d 100644 --- a/src/query/src/analyze.rs +++ b/src/query/src/analyze.rs @@ -46,6 +46,21 @@ const STAGE: &str = "stage"; const NODE: &str = "node"; const PLAN: &str = "plan"; +/// Returns the fixed output schema of [`DistAnalyzeExec`]: +/// (`stage`: UInt32, `node`: UInt32, `plan`: Utf8). +/// +/// Exposed so protocol-level `Describe` handlers can advertise the schema +/// clients will actually receive (the physical plan is rewritten in +/// `optimize_physical_plan` and its schema differs from the DataFusion +/// `Analyze` logical plan schema). +pub fn dist_analyze_output_schema() -> SchemaRef { + SchemaRef::new(Schema::new(vec![ + Field::new(STAGE, DataType::UInt32, true), + Field::new(NODE, DataType::UInt32, true), + Field::new(PLAN, DataType::Utf8, true), + ])) +} + #[derive(Debug)] pub struct DistAnalyzeExec { input: Arc, @@ -58,11 +73,7 @@ pub struct DistAnalyzeExec { impl DistAnalyzeExec { /// Create a new DistAnalyzeExec pub fn new(input: Arc, verbose: bool, format: AnalyzeFormat) -> Self { - let schema = SchemaRef::new(Schema::new(vec![ - Field::new(STAGE, DataType::UInt32, true), - Field::new(NODE, DataType::UInt32, true), - Field::new(PLAN, DataType::Utf8, true), - ])); + let schema = dist_analyze_output_schema(); let properties = Arc::new(Self::compute_properties(&input, schema.clone())); Self { input, diff --git a/src/query/src/lib.rs b/src/query/src/lib.rs index 68d8ff3a9e..4b2a2d5264 100644 --- a/src/query/src/lib.rs +++ b/src/query/src/lib.rs @@ -45,7 +45,7 @@ pub(crate) mod test_util; #[cfg(test)] mod tests; -pub use crate::analyze::analyze_plan_metrics_to_json_value; +pub use crate::analyze::{analyze_plan_metrics_to_json_value, dist_analyze_output_schema}; pub use crate::datafusion::DfContextProviderAdapter; pub use crate::query_engine::{ QueryEngine, QueryEngineContext, QueryEngineFactory, QueryEngineRef, diff --git a/src/servers/Cargo.toml b/src/servers/Cargo.toml index 7936c96fbd..bb358ea76a 100644 --- a/src/servers/Cargo.toml +++ b/src/servers/Cargo.toml @@ -167,8 +167,8 @@ servers = { workspace = true, features = ["testing"] } session = { workspace = true, features = ["testing"] } table.workspace = true tempfile = "3.0.0" -tokio-postgres = "0.7" -tokio-postgres-rustls = "0.12" +tokio-postgres.workspace = true +tokio-postgres-rustls = "0.14" [target.'cfg(unix)'.dev-dependencies] pprof = { version = "0.14", features = ["criterion", "flamegraph"] } diff --git a/src/servers/src/postgres/handler.rs b/src/servers/src/postgres/handler.rs index 484bb6a1f1..433eb11cda 100644 --- a/src/servers/src/postgres/handler.rs +++ b/src/servers/src/postgres/handler.rs @@ -40,6 +40,7 @@ use pgwire::error::{ErrorInfo, PgWireError, PgWireResult}; use pgwire::messages::PgWireBackendMessage; use pgwire::messages::copy::CopyData; use pgwire::messages::data::DataRow; +use query::dist_analyze_output_schema; use query::planner::DfLogicalPlanner; use query::query_engine::DescribeResult; use session::Session; @@ -554,6 +555,16 @@ fn describe_fields( session: &Arc, ) -> PgWireResult> { match sql_plan { + // EXPLAIN ANALYZE: at execution time the physical plan is replaced with + // GreptimeDB's `DistAnalyzeExec` (see `optimize_physical_plan`), whose + // output schema (stage/node/plan) differs from the DataFusion `Analyze` + // logical plan schema (plan_type/plan). Describe with the schema the + // client will actually receive so the DataRow field count matches. + SqlPlan::Plan(LogicalPlan::Analyze(_), _) => { + let schema: Schema = + Schema::try_from(dist_analyze_output_schema()).map_err(convert_err)?; + schema_to_pg(&schema, format, None).map_err(convert_err) + } // query SqlPlan::Plan(plan, _) if !matches!(plan, LogicalPlan::Dml(_) | LogicalPlan::Ddl(_)) => { let schema: Schema = plan.schema().clone().try_into().map_err(convert_err)?; @@ -679,6 +690,24 @@ fn describe_fields( Ok(vec![]) } } + // FETCH cursor: return the cursor's schema so the RowDescription + // matches the DataRow field count sent during Execute. + SqlPlan::Statement(Statement::FetchCursor(fetch), _) => { + let cursor_name = fetch.cursor_name.to_string(); + match session.get_cursor(&cursor_name) { + Some(cursor) => { + // `cursor.schema()` is the GreptimeDB `SchemaRef` captured + // when the cursor was declared; `schema_to_pg` accepts it + // directly (no DataFusion -> GreptimeDB conversion needed). + schema_to_pg(&cursor.schema(), format, None).map_err(convert_err) + } + None => { + // Cursor not found (e.g. DECLARE hasn't executed yet in + // an extended-protocol batch). Return NoData as a fallback. + Ok(vec![]) + } + } + } _ => { // NoData Ok(vec![]) diff --git a/src/session/src/lib.rs b/src/session/src/lib.rs index 110560b324..44f3ee6689 100644 --- a/src/session/src/lib.rs +++ b/src/session/src/lib.rs @@ -111,6 +111,15 @@ impl Session { .into() } + /// Returns the cursor with the given name, if it exists. + /// + /// Cursors are stored in the session's mutable inner data and shared across + /// all query contexts created from this session. + pub fn get_cursor(&self, name: &str) -> Option> { + let guard = self.mutable_inner.read().unwrap(); + guard.cursors.get(name).cloned() + } + pub fn conn_info(&self) -> &ConnInfo { &self.conn_info }