diff --git a/Cargo.lock b/Cargo.lock index e0276a71a9..ac7e49d22a 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -574,11 +574,10 @@ checksum = "96d30a06541fbafbc7f82ed10c06164cfbd2c401138f6addd8404629c4b16711" [[package]] name = "arrow" -version = "49.0.0" +version = "50.0.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "5bc25126d18a012146a888a0298f2c22e1150327bd2765fc76d710a556b2d614" +checksum = "aa285343fba4d829d49985bdc541e3789cf6000ed0e84be7c039438df4a4e78c" dependencies = [ - "ahash 0.8.7", "arrow-arith", "arrow-array", "arrow-buffer", @@ -596,9 +595,9 @@ dependencies = [ [[package]] name = "arrow-arith" -version = "49.0.0" +version = "50.0.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "34ccd45e217ffa6e53bbb0080990e77113bdd4e91ddb84e97b77649810bcf1a7" +checksum = "753abd0a5290c1bcade7c6623a556f7d1659c5f4148b140b5b63ce7bd1a45705" dependencies = [ "arrow-array", "arrow-buffer", @@ -611,9 +610,9 @@ dependencies = [ [[package]] name = "arrow-array" -version = "49.0.0" +version = "50.0.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "6bda9acea48b25123c08340f3a8ac361aa0f74469bb36f5ee9acf923fce23e9d" +checksum = "d390feeb7f21b78ec997a4081a025baef1e2e0d6069e181939b61864c9779609" dependencies = [ "ahash 0.8.7", "arrow-buffer", @@ -624,14 +623,13 @@ dependencies = [ "half 2.3.1", "hashbrown 0.14.3", "num", - "packed_simd", ] [[package]] name = "arrow-buffer" -version = "49.0.0" +version = "50.0.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "01a0fc21915b00fc6c2667b069c1b64bdd920982f426079bc4a7cab86822886c" +checksum = "69615b061701bcdffbc62756bc7e85c827d5290b472b580c972ebbbf690f5aa4" dependencies = [ "bytes", "half 2.3.1", @@ -640,9 +638,9 @@ dependencies = [ [[package]] name = "arrow-cast" -version = "49.0.0" +version = "50.0.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "5dc0368ed618d509636c1e3cc20db1281148190a78f43519487b2daf07b63b4a" +checksum = "e448e5dd2f4113bf5b74a1f26531708f5edcacc77335b7066f9398f4bcf4cdef" dependencies = [ "arrow-array", "arrow-buffer", @@ -659,9 +657,9 @@ dependencies = [ [[package]] name = "arrow-csv" -version = "49.0.0" +version = "50.0.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "2e09aa6246a1d6459b3f14baeaa49606cfdbca34435c46320e14054d244987ca" +checksum = "46af72211f0712612f5b18325530b9ad1bfbdc87290d5fbfd32a7da128983781" dependencies = [ "arrow-array", "arrow-buffer", @@ -678,9 +676,9 @@ dependencies = [ [[package]] name = "arrow-data" -version = "49.0.0" +version = "50.0.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "907fafe280a3874474678c1858b9ca4cb7fd83fb8034ff5b6d6376205a08c634" +checksum = "67d644b91a162f3ad3135ce1184d0a31c28b816a581e08f29e8e9277a574c64e" dependencies = [ "arrow-buffer", "arrow-schema", @@ -690,9 +688,9 @@ dependencies = [ [[package]] name = "arrow-ipc" -version = "49.0.0" +version = "50.0.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "79a43d6808411886b8c7d4f6f7dd477029c1e77ffffffb7923555cc6579639cd" +checksum = "03dea5e79b48de6c2e04f03f62b0afea7105be7b77d134f6c5414868feefb80d" dependencies = [ "arrow-array", "arrow-buffer", @@ -706,9 +704,9 @@ dependencies = [ [[package]] name = "arrow-json" -version = "49.0.0" +version = "50.0.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "d82565c91fd627922ebfe2810ee4e8346841b6f9361b87505a9acea38b614fee" +checksum = "8950719280397a47d37ac01492e3506a8a724b3fb81001900b866637a829ee0f" dependencies = [ "arrow-array", "arrow-buffer", @@ -726,9 +724,9 @@ dependencies = [ [[package]] name = "arrow-ord" -version = "49.0.0" +version = "50.0.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "9b23b0e53c0db57c6749997fd343d4c0354c994be7eca67152dd2bdb9a3e1bb4" +checksum = "1ed9630979034077982d8e74a942b7ac228f33dd93a93b615b4d02ad60c260be" dependencies = [ "arrow-array", "arrow-buffer", @@ -741,9 +739,9 @@ dependencies = [ [[package]] name = "arrow-row" -version = "49.0.0" +version = "50.0.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "361249898d2d6d4a6eeb7484be6ac74977e48da12a4dd81a708d620cc558117a" +checksum = "007035e17ae09c4e8993e4cb8b5b96edf0afb927cd38e2dff27189b274d83dcf" dependencies = [ "ahash 0.8.7", "arrow-array", @@ -756,18 +754,18 @@ dependencies = [ [[package]] name = "arrow-schema" -version = "49.0.0" +version = "50.0.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "09e28a5e781bf1b0f981333684ad13f5901f4cd2f20589eab7cf1797da8fc167" +checksum = "0ff3e9c01f7cd169379d269f926892d0e622a704960350d09d331be3ec9e0029" dependencies = [ "serde", ] [[package]] name = "arrow-select" -version = "49.0.0" +version = "50.0.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "4f6208466590960efc1d2a7172bc4ff18a67d6e25c529381d7f96ddaf0dc4036" +checksum = "1ce20973c1912de6514348e064829e50947e35977bb9d7fb637dc99ea9ffd78c" dependencies = [ "ahash 0.8.7", "arrow-array", @@ -779,9 +777,9 @@ dependencies = [ [[package]] name = "arrow-string" -version = "49.0.0" +version = "50.0.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "a4a48149c63c11c9ff571e50ab8f017d2a7cb71037a882b42f6354ed2da9acc7" +checksum = "00f3b37f2aeece31a2636d1b037dabb69ef590e03bdc7eb68519b51ec86932a7" dependencies = [ "arrow-array", "arrow-buffer", @@ -2443,12 +2441,14 @@ checksum = "7e962a19be5cfc3f3bf6dd8f61eb50107f356ad6270fbb3ed41476571db78be5" [[package]] name = "datafusion" -version = "34.0.0" -source = "git+https://github.com/openobserve/arrow-datafusion.git?rev=45e5537ca43d2c2a6e55b9804073b191b337b9e5#45e5537ca43d2c2a6e55b9804073b191b337b9e5" +version = "35.0.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "4328f5467f76d890fe3f924362dbc3a838c6a733f762b32d87f9e0b7bef5fb49" dependencies = [ "ahash 0.8.7", "arrow", "arrow-array", + "arrow-ipc", "arrow-schema", "async-compression", "async-trait", @@ -2489,8 +2489,9 @@ dependencies = [ [[package]] name = "datafusion-common" -version = "34.0.0" -source = "git+https://github.com/openobserve/arrow-datafusion.git?rev=45e5537ca43d2c2a6e55b9804073b191b337b9e5#45e5537ca43d2c2a6e55b9804073b191b337b9e5" +version = "35.0.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d29a7752143b446db4a2cccd9a6517293c6b97e8c39e520ca43ccd07135a4f7e" dependencies = [ "ahash 0.8.7", "arrow", @@ -2508,8 +2509,9 @@ dependencies = [ [[package]] name = "datafusion-execution" -version = "34.0.0" -source = "git+https://github.com/openobserve/arrow-datafusion.git?rev=45e5537ca43d2c2a6e55b9804073b191b337b9e5#45e5537ca43d2c2a6e55b9804073b191b337b9e5" +version = "35.0.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "2d447650af16e138c31237f53ddaef6dd4f92f0e2d3f2f35d190e16c214ca496" dependencies = [ "arrow", "chrono", @@ -2528,8 +2530,9 @@ dependencies = [ [[package]] name = "datafusion-expr" -version = "34.0.0" -source = "git+https://github.com/openobserve/arrow-datafusion.git?rev=45e5537ca43d2c2a6e55b9804073b191b337b9e5#45e5537ca43d2c2a6e55b9804073b191b337b9e5" +version = "35.0.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d8d19598e48a498850fb79f97a9719b1f95e7deb64a7a06f93f313e8fa1d524b" dependencies = [ "ahash 0.8.7", "arrow", @@ -2543,8 +2546,9 @@ dependencies = [ [[package]] name = "datafusion-optimizer" -version = "34.0.0" -source = "git+https://github.com/openobserve/arrow-datafusion.git?rev=45e5537ca43d2c2a6e55b9804073b191b337b9e5#45e5537ca43d2c2a6e55b9804073b191b337b9e5" +version = "35.0.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "8b7feb0391f1fc75575acb95b74bfd276903dc37a5409fcebe160bc7ddff2010" dependencies = [ "arrow", "async-trait", @@ -2560,8 +2564,9 @@ dependencies = [ [[package]] name = "datafusion-physical-expr" -version = "34.0.0" -source = "git+https://github.com/openobserve/arrow-datafusion.git?rev=45e5537ca43d2c2a6e55b9804073b191b337b9e5#45e5537ca43d2c2a6e55b9804073b191b337b9e5" +version = "35.0.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "e911bca609c89a54e8f014777449d8290327414d3e10c57a3e3c2122e38878d0" dependencies = [ "ahash 0.8.7", "arrow", @@ -2593,8 +2598,9 @@ dependencies = [ [[package]] name = "datafusion-physical-plan" -version = "34.0.0" -source = "git+https://github.com/openobserve/arrow-datafusion.git?rev=45e5537ca43d2c2a6e55b9804073b191b337b9e5#45e5537ca43d2c2a6e55b9804073b191b337b9e5" +version = "35.0.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "e96b546b8a02e9c2ab35ac6420d511f12a4701950c1eb2e568c122b4fefb0be3" dependencies = [ "ahash 0.8.7", "arrow", @@ -2623,8 +2629,9 @@ dependencies = [ [[package]] name = "datafusion-sql" -version = "34.0.0" -source = "git+https://github.com/openobserve/arrow-datafusion.git?rev=45e5537ca43d2c2a6e55b9804073b191b337b9e5#45e5537ca43d2c2a6e55b9804073b191b337b9e5" +version = "35.0.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "2d18d36f260bbbd63aafdb55339213a23d540d3419810575850ef0a798a6b768" dependencies = [ "arrow", "arrow-schema", @@ -4635,9 +4642,9 @@ dependencies = [ [[package]] name = "object_store" -version = "0.8.0" +version = "0.9.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "2524735495ea1268be33d200e1ee97455096a0846295a21548cd2f3541de7050" +checksum = "d139f545f64630e2e3688fd9f81c470888ab01edeb72d13b4e86c566f1130000" dependencies = [ "async-trait", "base64 0.21.7", @@ -4646,14 +4653,14 @@ dependencies = [ "futures", "humantime", "hyper", - "itertools 0.11.0", + "itertools 0.12.0", "parking_lot 0.12.1", "percent-encoding", "quick-xml", "rand", "reqwest", "ring 0.17.7", - "rustls-pemfile", + "rustls-pemfile 2.0.0", "serde", "serde_json", "snafu 0.7.5", @@ -4830,6 +4837,7 @@ dependencies = [ "serde_json", "sha256", "sled", + "snafu 0.7.5", "snap", "sqlparser", "sqlx", @@ -5113,16 +5121,6 @@ dependencies = [ "sha2", ] -[[package]] -name = "packed_simd" -version = "0.3.9" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "1f9f08af0c877571712e2e3e686ad79efad9657dbf0f7c3c8ba943ff6c38932d" -dependencies = [ - "cfg-if 1.0.0", - "num-traits", -] - [[package]] name = "packedvec" version = "1.2.4" @@ -5189,9 +5187,9 @@ dependencies = [ [[package]] name = "parquet" -version = "49.0.0" +version = "50.0.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "af88740a842787da39b3d69ce5fbf6fce97d20211d3b299fee0a0da6430c74d4" +checksum = "547b92ebf0c1177e3892f44c8f79757ee62e678d564a9834189725f2c5b7a750" dependencies = [ "ahash 0.8.7", "arrow-array", @@ -5207,6 +5205,7 @@ dependencies = [ "chrono", "flate2", "futures", + "half 2.3.1", "hashbrown 0.14.3", "lz4_flex", "num", @@ -6029,7 +6028,8 @@ dependencies = [ "percent-encoding", "pin-project-lite", "rustls", - "rustls-pemfile", + "rustls-native-certs", + "rustls-pemfile 1.0.4", "serde", "serde_json", "serde_urlencoded", @@ -6304,7 +6304,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "a9aace74cb666635c918e9c12bc0d348266037aa8eb599b5cba565709a8dff00" dependencies = [ "openssl-probe", - "rustls-pemfile", + "rustls-pemfile 1.0.4", "schannel", "security-framework", ] @@ -6318,6 +6318,22 @@ dependencies = [ "base64 0.21.7", ] +[[package]] +name = "rustls-pemfile" +version = "2.0.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "35e4980fa29e4c4b212ffb3db068a564cbf560e51d3944b7c88bd8bf5bec64f4" +dependencies = [ + "base64 0.21.7", + "rustls-pki-types", +] + +[[package]] +name = "rustls-pki-types" +version = "1.1.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "9e9d979b3ce68192e42760c7810125eb6cf2ea10efae545a156063e61f314e2a" + [[package]] name = "rustls-webpki" version = "0.101.7" @@ -6855,9 +6871,9 @@ dependencies = [ [[package]] name = "sqlparser" -version = "0.40.0" +version = "0.41.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "7c80afe31cdb649e56c0d9bb5503be9166600d68a852c38dd445636d126858e5" +checksum = "5cc2c25a6c66789625ef164b4c7d2e548d627902280c13710d33da8222169964" dependencies = [ "log", "serde", @@ -6918,7 +6934,7 @@ dependencies = [ "paste", "percent-encoding", "rustls", - "rustls-pemfile", + "rustls-pemfile 1.0.4", "serde", "serde_json", "sha2", @@ -7621,7 +7637,7 @@ dependencies = [ "pin-project", "prost 0.12.3", "rustls", - "rustls-pemfile", + "rustls-pemfile 1.0.4", "tokio", "tokio-rustls", "tokio-stream", diff --git a/Cargo.toml b/Cargo.toml index d92d3818b0..ca86d783de 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -76,14 +76,12 @@ clap = { version = "4.1", default-features = false, features = [ cloudevents-sdk = { version = "0.7.0", features = ["actix"] } csv = "1.2.1" dashmap = { version = "5.4", features = ["serde"] } -datafusion = { git = "https://github.com/openobserve/arrow-datafusion.git", rev = "45e5537ca43d2c2a6e55b9804073b191b337b9e5", version = "34", features = [ - "simd", -] } -datafusion-expr = { git = "https://github.com/openobserve/arrow-datafusion.git", rev = "45e5537ca43d2c2a6e55b9804073b191b337b9e5", version = "34" } -arrow = { version = "49", features = ["simd", "ipc_compression"] } -arrow-schema = { version = "49", features = ["serde"] } -parquet = { version = "49", features = ["arrow", "async"] } -object_store = { version = "0.8", features = ["aws", "azure", "gcp"] } +datafusion = "35" +datafusion-expr = "35" +arrow = { version = "50.0.0", features = ["ipc_compression"] } +arrow-schema = { version = "50.0.0", features = ["serde"] } +parquet = { version = "50.0.0", features = ["arrow", "async", "object_store"] } +object_store = { version = "0.9", features = ["aws", "azure", "gcp"] } dotenv_config = "0.1.7" dotenvy = "0.15" env_logger = "0.10" @@ -141,9 +139,10 @@ segment = "0.2" serde = { version = "1", features = ["derive"] } serde_json = "1" sha256 = "1.4.0" +snafu = "0.7.5" sled = "0.34" snap = "1" -sqlparser = { version = "0.40", features = ["serde", "visitor"] } +sqlparser = { version = "0.41", features = ["serde", "visitor"] } sqlx = { version = "0.7", features = [ "runtime-tokio-rustls", "postgres", @@ -214,10 +213,10 @@ bytes = "1.4" byteorder = "1.4.3" chrono = { version = "0.4", default-features = false, features = ["clock"] } dashmap = { version = "5.4", features = ["serde"] } -arrow = { version = "49", features = ["simd", "ipc_compression"] } -arrow-json = "49" -arrow-schema = { version = "49", features = ["serde"] } -parquet = { version = "49", features = ["arrow", "async"] } +arrow = { version = "50.0.0", features = ["ipc_compression"] } +arrow-json = "50.0.0" +arrow-schema = { version = "50.0.0", features = ["serde"] } +parquet = { version = "50.0.0", features = ["arrow", "async", "object_store"] } dotenv_config = "0.1.7" dotenvy = "0.15" faststr = "0.2" diff --git a/src/service/metrics/otlp_grpc.rs b/src/service/metrics/otlp_grpc.rs index 47709abb1d..9cddaa48c4 100644 --- a/src/service/metrics/otlp_grpc.rs +++ b/src/service/metrics/otlp_grpc.rs @@ -302,7 +302,11 @@ pub async fn handle_grpc_request( let buf = metric_data_map .entry(local_metric_name.to_owned()) .or_default(); - let schema = metric_schema_map.get(local_metric_name).unwrap().clone(); + let schema = metric_schema_map + .get(local_metric_name) + .unwrap() + .clone() + .with_metadata(HashMap::new()); let schema_key = schema.hash_key(); // get hour key let hour_key = crate::service::ingestion::get_wal_time_key( diff --git a/src/service/search/datafusion/date_format_udf.rs b/src/service/search/datafusion/date_format_udf.rs index 2cd8e103c4..a8911da4c0 100644 --- a/src/service/search/datafusion/date_format_udf.rs +++ b/src/service/search/datafusion/date_format_udf.rs @@ -51,9 +51,12 @@ pub(crate) static DATE_FORMAT_UDF: Lazy = Lazy::new(|| { pub fn date_format_expr_impl() -> ScalarFunctionImplementation { let func = move |args: &[ArrayRef]| -> datafusion::error::Result { if args.len() != 3 { - return Err(DataFusionError::SQL(ParserError::ParserError( - "UDF params should be: date_format(field, format, zone)".to_string(), - ))); + return Err(DataFusionError::SQL( + ParserError::ParserError( + "UDF params should be: date_format(field, format, zone)".to_string(), + ), + None, + )); } // 1. cast both arguments to Union. These casts MUST be aligned with the signature or this diff --git a/src/service/search/datafusion/exec.rs b/src/service/search/datafusion/exec.rs index fe46705b02..d60b7b73c0 100644 --- a/src/service/search/datafusion/exec.rs +++ b/src/service/search/datafusion/exec.rs @@ -1002,6 +1002,10 @@ pub fn create_session_config(search_type: &SearchType) -> Result let mut config = SessionConfig::from_env()? .with_batch_size(PARQUET_BATCH_SIZE) .with_information_schema(true); + config = config.set_bool( + "datafusion.execution.listing_table_ignore_subdirectory", + false, + ); if search_type == &SearchType::Normal { config = config.set_bool("datafusion.execution.parquet.pushdown_filters", true); config = config.set_bool("datafusion.execution.parquet.reorder_filters", true); diff --git a/src/service/search/datafusion/match_udf.rs b/src/service/search/datafusion/match_udf.rs index 44da27152c..33dce650a0 100644 --- a/src/service/search/datafusion/match_udf.rs +++ b/src/service/search/datafusion/match_udf.rs @@ -60,9 +60,10 @@ pub(crate) static MATCH_IGNORE_CASE_UDF: Lazy = Lazy::new(|| { pub fn match_expr_impl(case_insensitive: bool) -> ScalarFunctionImplementation { let func = move |args: &[ArrayRef]| -> datafusion::error::Result { if args.len() != 2 { - return Err(DataFusionError::SQL(ParserError::ParserError( - "match UDF expects two string".to_string(), - ))); + return Err(DataFusionError::SQL( + ParserError::ParserError("match UDF expects two string".to_string()), + None, + )); } // 1. cast both arguments to string. These casts MUST be aligned with the signature or this diff --git a/src/service/search/datafusion/storage/memory.rs b/src/service/search/datafusion/storage/memory.rs index a6d6b5af08..e84d995dc4 100644 --- a/src/service/search/datafusion/storage/memory.rs +++ b/src/service/search/datafusion/storage/memory.rs @@ -24,6 +24,7 @@ use object_store::{ }; use tokio::io::AsyncWrite; +use super::GetRangeExt; use crate::common::{ infra::{cache::file_data, storage}, utils::time::BASE_TIME, @@ -105,7 +106,12 @@ impl ObjectStore for FS { version: None, }; let (range, data) = match options.range { - Some(range) => (range.clone(), data.slice(range)), + Some(range) => { + let r = range + .as_range(data.len()) + .map_err(|e| super::Error::BadRange(e.to_string()))?; + (r.clone(), data.slice(r)) + } None => (0..data.len(), data), }; Ok(GetResult { diff --git a/src/service/search/datafusion/storage/mod.rs b/src/service/search/datafusion/storage/mod.rs index 7d68a403ec..51cc7cbc11 100644 --- a/src/service/search/datafusion/storage/mod.rs +++ b/src/service/search/datafusion/storage/mod.rs @@ -13,6 +13,10 @@ // You should have received a copy of the GNU Affero General Public License // along with this program. If not, see . +use std::ops::Range; + +use object_store::GetRange; +use snafu::Snafu; use thiserror::Error as ThisError; pub mod file_list; @@ -28,7 +32,7 @@ pub enum StorageType { /// A specialized `Error` for in-memory object store-related errors #[derive(ThisError, Debug)] #[allow(missing_docs)] -enum Error { +pub(crate) enum Error { #[error("Out of range")] OutOfRange(String), #[error("Bad range")] @@ -43,3 +47,65 @@ impl From for object_store::Error { } } } + +#[derive(Debug, Snafu)] +pub(crate) enum InvalidGetRange { + #[snafu(display( + "Wanted range starting at {requested}, but object was only {length} bytes long" + ))] + StartTooLarge { requested: usize, length: usize }, + + #[snafu(display("Range started at {start} and ended at {end}"))] + Inconsistent { start: usize, end: usize }, +} + +pub(crate) trait GetRangeExt { + fn is_valid(&self) -> Result<(), InvalidGetRange>; + /// Convert to a [`Range`] if valid. + fn as_range(&self, len: usize) -> Result, InvalidGetRange>; +} + +impl GetRangeExt for GetRange { + fn is_valid(&self) -> Result<(), InvalidGetRange> { + match self { + Self::Bounded(r) if r.end <= r.start => { + return Err(InvalidGetRange::Inconsistent { + start: r.start, + end: r.end, + }); + } + _ => (), + }; + Ok(()) + } + + /// Convert to a [`Range`] if valid. + fn as_range(&self, len: usize) -> Result, InvalidGetRange> { + self.is_valid()?; + match self { + Self::Bounded(r) => { + if r.start >= len { + Err(InvalidGetRange::StartTooLarge { + requested: r.start, + length: len, + }) + } else if r.end > len { + Ok(r.start..len) + } else { + Ok(r.clone()) + } + } + Self::Offset(o) => { + if *o >= len { + Err(InvalidGetRange::StartTooLarge { + requested: *o, + length: len, + }) + } else { + Ok(*o..len) + } + } + Self::Suffix(n) => Ok(len.saturating_sub(*n)..len), + } + } +} diff --git a/src/service/search/datafusion/storage/tmpfs.rs b/src/service/search/datafusion/storage/tmpfs.rs index eb1c87b273..d162287d47 100644 --- a/src/service/search/datafusion/storage/tmpfs.rs +++ b/src/service/search/datafusion/storage/tmpfs.rs @@ -26,6 +26,7 @@ use object_store::{ use thiserror::Error as ThisError; use tokio::io::AsyncWrite; +use super::GetRangeExt; use crate::common::{infra::cache::tmpfs, utils::time::BASE_TIME}; /// A specialized `Error` for in-memory object store-related errors @@ -94,7 +95,12 @@ impl ObjectStore for Tmpfs { version: None, }; let (range, data) = match options.range { - Some(range) => (range.clone(), data.slice(range)), + Some(range) => { + let r = range + .as_range(data.len()) + .map_err(|e| super::Error::BadRange(e.to_string()))?; + (r.clone(), data.slice(r)) + } None => (0..data.len(), data), }; Ok(GetResult { diff --git a/src/service/search/datafusion/time_range_udf.rs b/src/service/search/datafusion/time_range_udf.rs index 2e396f28ea..376ad44d66 100644 --- a/src/service/search/datafusion/time_range_udf.rs +++ b/src/service/search/datafusion/time_range_udf.rs @@ -50,9 +50,12 @@ pub(crate) static TIME_RANGE_UDF: Lazy = Lazy::new(|| { pub fn time_range_expr_impl() -> ScalarFunctionImplementation { let func = move |args: &[ArrayRef]| -> datafusion::error::Result { if args.len() != 3 { - return Err(DataFusionError::SQL(ParserError::ParserError( - "UDF params should be: time_range(field, start, end)".to_string(), - ))); + return Err(DataFusionError::SQL( + ParserError::ParserError( + "UDF params should be: time_range(field, start, end)".to_string(), + ), + None, + )); } // 1. cast both arguments to Union. These casts MUST be aligned with the signature or this diff --git a/src/service/search/grpc/errors.rs b/src/service/search/grpc/errors.rs index 589362bd69..8bff9d4da2 100644 --- a/src/service/search/grpc/errors.rs +++ b/src/service/search/grpc/errors.rs @@ -36,10 +36,13 @@ fn get_key_from_error(err: &str, pos: usize) -> Option { impl From for Error { fn from(err: DataFusionError) -> Self { - if let DataFusionError::SchemaError(SchemaError::FieldNotFound { - field, - valid_fields: _, - }) = err + if let DataFusionError::SchemaError( + SchemaError::FieldNotFound { + field, + valid_fields: _, + }, + _, + ) = err { return Error::ErrorCode(ErrorCodes::SearchFieldNotFound(field.name)); }