diff --git a/quickwit/Cargo.lock b/quickwit/Cargo.lock index caede75b90f..6007e9212a4 100644 --- a/quickwit/Cargo.lock +++ b/quickwit/Cargo.lock @@ -14,7 +14,7 @@ version = "0.25.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "1b5d307320b3181d6d7954e663bd7c774a838b8220fe0593c86d9fb09f498b4b" dependencies = [ - "gimli", + "gimli 0.32.3", ] [[package]] @@ -228,6 +228,12 @@ dependencies = [ "object", ] +[[package]] +name = "arbitrary" +version = "1.4.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "c3d036a3c4ab069c7b410a2ce876bd74808d2d0888a82667669f8e783a898bf1" + [[package]] name = "arc-swap" version = "1.9.2" @@ -1765,6 +1771,9 @@ name = "bumpalo" version = "3.20.3" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "72f5acc6cb2ba439de613abc23857ec3d78374d8ed5ac84e9d11336e87da8649" +dependencies = [ + "allocator-api2", +] [[package]] name = "bytecount" @@ -1998,7 +2007,7 @@ dependencies = [ "async-trait", "bytes", "itertools 0.15.0", - "lru 0.18.0", + "lru 0.18.2", "rand 0.10.2", "serde", "tokio", @@ -2434,6 +2443,186 @@ dependencies = [ "libc", ] +[[package]] +name = "cranelift" +version = "0.134.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "4ae1a7acf94cec79fd145e25ef76d437630f7598f3ed4962edbc3feb6079eae3" +dependencies = [ + "cranelift-codegen", + "cranelift-frontend", + "cranelift-module", +] + +[[package]] +name = "cranelift-assembler-x64" +version = "0.134.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "16c273ba1bfc3a1cb27cb4f83df27134f3cb638a61f5ed79ab22cf86eb13da0d" +dependencies = [ + "cranelift-assembler-x64-meta", +] + +[[package]] +name = "cranelift-assembler-x64-meta" +version = "0.134.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "cbf8909326a466739e83ffe11e74c0c8c6b4bfa911b612aece915daa746cff58" +dependencies = [ + "cranelift-srcgen", +] + +[[package]] +name = "cranelift-bforest" +version = "0.134.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "3cb26b06b54d8b2f8cdf4d440a32fc3b3b91c71c162333a69a70a66fa265fc45" +dependencies = [ + "cranelift-entity", + "wasmtime-internal-core", +] + +[[package]] +name = "cranelift-bitset" +version = "0.134.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7ce4e57dbb7f73d08011808a19787c3a03c3e7854b2a479ed0b788915fb671be" +dependencies = [ + "wasmtime-internal-core", +] + +[[package]] +name = "cranelift-codegen" +version = "0.134.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "74b0a8968f33d0bc13bc86fe70dd3947e93acc9f358ed9848e9988448972a9be" +dependencies = [ + "bumpalo", + "cranelift-assembler-x64", + "cranelift-bforest", + "cranelift-bitset", + "cranelift-codegen-meta", + "cranelift-codegen-shared", + "cranelift-control", + "cranelift-entity", + "cranelift-isle", + "gimli 0.33.0", + "hashbrown 0.17.1", + "libm", + "log", + "regalloc2", + "rustc-hash", + "serde", + "smallvec", + "target-lexicon", + "wasmtime-internal-core", +] + +[[package]] +name = "cranelift-codegen-meta" +version = "0.134.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f69e8e2fe9deb48ef7b8d0c61612f8c1e2277d0421c201e30a4c844a4ebac698" +dependencies = [ + "cranelift-assembler-x64-meta", + "cranelift-codegen-shared", + "cranelift-srcgen", + "heck 0.5.0", +] + +[[package]] +name = "cranelift-codegen-shared" +version = "0.134.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "ad4db1e87a65e9c5a832e3ede160830d4e34df938ac38b5604819f8825687e73" + +[[package]] +name = "cranelift-control" +version = "0.134.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "9248cd37bb6ec460981f57ead368e32b1eca7a207ad50e3e9c7adc3b161f20e9" +dependencies = [ + "arbitrary", +] + +[[package]] +name = "cranelift-entity" +version = "0.134.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "26b8bc91c0a3530d132acbf118965dabe9a084029eb0fa74b64760360dec3316" +dependencies = [ + "cranelift-bitset", + "wasmtime-internal-core", +] + +[[package]] +name = "cranelift-frontend" +version = "0.134.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "5086f2ebc084a98b387e522680f62ef65a512444d8382151e03088716d3a1714" +dependencies = [ + "cranelift-codegen", + "hashbrown 0.17.1", + "log", + "smallvec", + "target-lexicon", +] + +[[package]] +name = "cranelift-isle" +version = "0.134.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "b6a658670f779afc2df083be201ebe54805f69b7b9ffa8dab04142b72ed751c9" + +[[package]] +name = "cranelift-jit" +version = "0.134.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "08e06a19e861cb46dafe15e5602fcd8a8aeb302e82df10f7886c73410aed1f8b" +dependencies = [ + "anyhow", + "cranelift-codegen", + "cranelift-control", + "cranelift-entity", + "cranelift-module", + "cranelift-native", + "libc", + "log", + "memmap2", + "region", + "target-lexicon", + "wasmtime-internal-jit-icache-coherence", + "windows-sys 0.61.2", +] + +[[package]] +name = "cranelift-module" +version = "0.134.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "75623dde4c952712d12e3281a1b7bf7587bca23139ac0d0e3d45de1f99abd772" +dependencies = [ + "anyhow", + "cranelift-codegen", + "cranelift-control", +] + +[[package]] +name = "cranelift-native" +version = "0.134.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "16a41083b1fd953debd6f32d21fcda765b258f323ce907fdd6b4715f2769c478" +dependencies = [ + "cranelift-codegen", + "libc", + "target-lexicon", +] + +[[package]] +name = "cranelift-srcgen" +version = "0.134.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "244d9cd0759b0a3ccdd587cafb13434f9ea43b5c9d07701d277daf5cd0f7b7c1" + [[package]] name = "crc" version = "3.4.0" @@ -4545,6 +4734,18 @@ version = "0.32.3" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "e629b9b98ef3dd8afe6ca2bd0f89306cec16d43d907889945bc5d6687f2f13c7" +[[package]] +name = "gimli" +version = "0.33.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "0bf7f043f89559805f8c7cacc432749b2fa0d0a0a9ee46ce47164ed5ba7f126c" +dependencies = [ + "fnv", + "hashbrown 0.16.1", + "indexmap 2.14.0", + "stable_deref_trait", +] + [[package]] name = "glob" version = "0.3.3" @@ -5526,6 +5727,19 @@ dependencies = [ "jiff-tzdb", ] +[[package]] +name = "jitexpr" +version = "0.1.0" +source = "git+https://github.com/quickwit-oss/tantivy/?rev=c661f6a#c661f6a740fac429a7538bbf697942d7c7e3e5f4" +dependencies = [ + "cranelift", + "cranelift-jit", + "cranelift-module", + "cranelift-native", + "regex", + "thiserror 2.0.18", +] + [[package]] name = "jni" version = "0.22.4" @@ -5980,9 +6194,9 @@ dependencies = [ [[package]] name = "lru" -version = "0.18.0" +version = "0.18.2" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "8a860605968fce16869fd239cf4237a82f3ac470723415db603b0e8b6c8d4fb9" +checksum = "5d2f2f9b4ba7e6b24d95e7e899329d35be83bcded72c8540cdd5368932d1d90a" dependencies = [ "hashbrown 0.17.1", ] @@ -6036,6 +6250,15 @@ version = "0.14.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "ecbdfe44b1bd960b68170b417450a628c43f7cf56bb3c5317e61cb230ee7f226" +[[package]] +name = "mach2" +version = "0.4.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d640282b302c0bb0a2a8e0233ead9035e3bed871f0b7e81fe4a1ec829765db44" +dependencies = [ + "libc", +] + [[package]] name = "matchers" version = "0.2.0" @@ -7087,7 +7310,7 @@ checksum = "1a80800c0488c3a21695ea981a54918fbb37abf04f4d0720c453632255e2ff0e" [[package]] name = "ownedbytes" version = "0.9.0" -source = "git+https://github.com/quickwit-oss/tantivy/?rev=86641f7#86641f72df7f79867486aa4645b66767a80cc4e8" +source = "git+https://github.com/quickwit-oss/tantivy/?rev=c661f6a#c661f6a740fac429a7538bbf697942d7c7e3e5f4" dependencies = [ "stable_deref_trait", ] @@ -8597,7 +8820,7 @@ dependencies = [ "fnv", "futures", "itertools 0.15.0", - "lru 0.18.0", + "lru 0.18.2", "mockall", "proptest", "quickwit-actors", @@ -9415,7 +9638,7 @@ dependencies = [ "http 1.4.2", "http-body-util", "hyper 1.10.1", - "lru 0.18.0", + "lru 0.18.2", "md5", "mini-moka", "mockall", @@ -9886,6 +10109,20 @@ dependencies = [ "serde_json", ] +[[package]] +name = "regalloc2" +version = "0.15.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "757712e8e61590d6d4f5d563483755538b5aa13467837a3b41cd9832509a7f85" +dependencies = [ + "allocator-api2", + "bumpalo", + "hashbrown 0.17.1", + "log", + "rustc-hash", + "smallvec", +] + [[package]] name = "regex" version = "1.13.0" @@ -9933,6 +10170,18 @@ version = "0.8.11" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "d6f6ff9a378485b298a5286656da665ba74413d36db0979633275d2e708145d4" +[[package]] +name = "region" +version = "3.0.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "e6b6ebd13bc009aef9cd476c1310d49ac354d36e240cf1bd753290f3dc7199a7" +dependencies = [ + "bitflags 1.3.2", + "libc", + "mach2", + "windows-sys 0.52.0", +] + [[package]] name = "regress" version = "0.10.5" @@ -11765,11 +12014,11 @@ checksum = "7b2093cf4c8eb1e67749a6762251bc9cd836b6fc171623bd0a9d324d37af2417" [[package]] name = "tantivy" version = "0.27.0" -source = "git+https://github.com/quickwit-oss/tantivy/?rev=86641f7#86641f72df7f79867486aa4645b66767a80cc4e8" +source = "git+https://github.com/quickwit-oss/tantivy/?rev=c661f6a#c661f6a740fac429a7538bbf697942d7c7e3e5f4" dependencies = [ "aho-corasick", "arc-swap", - "base64 0.22.1", + "base64 0.23.0", "bitpacking", "bon", "byteorder", @@ -11785,9 +12034,10 @@ dependencies = [ "futures-util", "htmlescape", "itertools 0.14.0", + "jitexpr", "levenshtein_automata", "log", - "lru 0.16.4", + "lru 0.18.2", "lz4_flex 0.14.0", "measure_time", "memmap2", @@ -11804,7 +12054,7 @@ dependencies = [ "tantivy-columnar", "tantivy-common", "tantivy-fst", - "tantivy-query-grammar 0.26.0 (git+https://github.com/quickwit-oss/tantivy/?rev=86641f7)", + "tantivy-query-grammar 0.26.0 (git+https://github.com/quickwit-oss/tantivy/?rev=c661f6a)", "tantivy-sstable", "tantivy-stacker", "tantivy-tokenizer-api", @@ -11821,7 +12071,7 @@ dependencies = [ [[package]] name = "tantivy-bitpacker" version = "0.10.0" -source = "git+https://github.com/quickwit-oss/tantivy/?rev=86641f7#86641f72df7f79867486aa4645b66767a80cc4e8" +source = "git+https://github.com/quickwit-oss/tantivy/?rev=c661f6a#c661f6a740fac429a7538bbf697942d7c7e3e5f4" dependencies = [ "bitpacking", ] @@ -11829,7 +12079,7 @@ dependencies = [ [[package]] name = "tantivy-columnar" version = "0.7.0" -source = "git+https://github.com/quickwit-oss/tantivy/?rev=86641f7#86641f72df7f79867486aa4645b66767a80cc4e8" +source = "git+https://github.com/quickwit-oss/tantivy/?rev=c661f6a#c661f6a740fac429a7538bbf697942d7c7e3e5f4" dependencies = [ "downcast-rs", "fastdivide", @@ -11844,7 +12094,7 @@ dependencies = [ [[package]] name = "tantivy-common" version = "0.11.0" -source = "git+https://github.com/quickwit-oss/tantivy/?rev=86641f7#86641f72df7f79867486aa4645b66767a80cc4e8" +source = "git+https://github.com/quickwit-oss/tantivy/?rev=c661f6a#c661f6a740fac429a7538bbf697942d7c7e3e5f4" dependencies = [ "async-trait", "byteorder", @@ -11880,7 +12130,7 @@ dependencies = [ [[package]] name = "tantivy-query-grammar" version = "0.26.0" -source = "git+https://github.com/quickwit-oss/tantivy/?rev=86641f7#86641f72df7f79867486aa4645b66767a80cc4e8" +source = "git+https://github.com/quickwit-oss/tantivy/?rev=c661f6a#c661f6a740fac429a7538bbf697942d7c7e3e5f4" dependencies = [ "fnv", "nom 7.1.3", @@ -11892,7 +12142,7 @@ dependencies = [ [[package]] name = "tantivy-sstable" version = "0.7.0" -source = "git+https://github.com/quickwit-oss/tantivy/?rev=86641f7#86641f72df7f79867486aa4645b66767a80cc4e8" +source = "git+https://github.com/quickwit-oss/tantivy/?rev=c661f6a#c661f6a740fac429a7538bbf697942d7c7e3e5f4" dependencies = [ "futures-util", "itertools 0.14.0", @@ -11905,7 +12155,7 @@ dependencies = [ [[package]] name = "tantivy-stacker" version = "0.7.0" -source = "git+https://github.com/quickwit-oss/tantivy/?rev=86641f7#86641f72df7f79867486aa4645b66767a80cc4e8" +source = "git+https://github.com/quickwit-oss/tantivy/?rev=c661f6a#c661f6a740fac429a7538bbf697942d7c7e3e5f4" dependencies = [ "murmurhash32", "tantivy-common", @@ -11914,11 +12164,17 @@ dependencies = [ [[package]] name = "tantivy-tokenizer-api" version = "0.7.0" -source = "git+https://github.com/quickwit-oss/tantivy/?rev=86641f7#86641f72df7f79867486aa4645b66767a80cc4e8" +source = "git+https://github.com/quickwit-oss/tantivy/?rev=c661f6a#c661f6a740fac429a7538bbf697942d7c7e3e5f4" dependencies = [ "serde", ] +[[package]] +name = "target-lexicon" +version = "0.13.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "adb6935a6f5c20170eeceb1a3835a49e12e19d792f6dd344ccc76a985ca5a6ca" + [[package]] name = "tempfile" version = "3.27.0" @@ -13301,6 +13557,28 @@ dependencies = [ "web-sys", ] +[[package]] +name = "wasmtime-internal-core" +version = "47.0.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f41d3b2cb9c3de7b696690af070d408b5e357011832f8b8bb7fbb315dbb620ed" +dependencies = [ + "hashbrown 0.17.1", + "libm", +] + +[[package]] +name = "wasmtime-internal-jit-icache-coherence" +version = "47.0.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "4a52de8b6afbfcd618073f592c1450cb661574e7061f3f4a44f1f7ec0c7f7909" +dependencies = [ + "cfg-if", + "libc", + "wasmtime-internal-core", + "windows-sys 0.61.2", +] + [[package]] name = "wasmtimer" version = "0.4.3" diff --git a/quickwit/Cargo.toml b/quickwit/Cargo.toml index 8904ed1dd3c..6a148207198 100644 --- a/quickwit/Cargo.toml +++ b/quickwit/Cargo.toml @@ -426,7 +426,8 @@ quickwit-storage = { path = "quickwit-storage" } quickwit-telemetry-exporters = { path = "quickwit-telemetry-exporters" } quickwit-transport = { path = "quickwit-transport" } -tantivy = { git = "https://github.com/quickwit-oss/tantivy/", rev = "86641f7", default-features = false, features = [ +tantivy = { git = "https://github.com/quickwit-oss/tantivy/", rev = "c661f6a", default-features = false, features = [ + "jitexpr", "lz4-compression", "mmap", "quickwit", diff --git a/quickwit/quickwit-directories/src/hot_directory.rs b/quickwit/quickwit-directories/src/hot_directory.rs index 2019505b542..0cc6ee00f50 100644 --- a/quickwit/quickwit-directories/src/hot_directory.rs +++ b/quickwit/quickwit-directories/src/hot_directory.rs @@ -25,6 +25,7 @@ use serde::{Deserialize, Serialize}; use tantivy::directory::error::OpenReadError; use tantivy::directory::{FileHandle, FileSlice, OwnedBytes}; use tantivy::error::DataCorruption; +use tantivy::index::SegmentComponent; use tantivy::{Directory, HasLen, Index, IndexReader, ReloadPolicy, TantivyError}; use crate::{CachingDirectory, DebugProxyDirectory}; @@ -459,11 +460,24 @@ impl Directory for HotDirectory { fn list_index_files(index: &Index) -> tantivy::Result> { let index_meta = index.load_metas()?; - let mut files: HashSet = index_meta - .segments - .into_iter() - .flat_map(|segment_meta| segment_meta.list_files()) - .collect(); + let segment_components = [ + SegmentComponent::Postings, + SegmentComponent::Positions, + SegmentComponent::Terms, + SegmentComponent::Store, + SegmentComponent::FastFields, + SegmentComponent::FieldNorms, + SegmentComponent::Delete, + ]; + let mut files = HashSet::new(); + for segment_meta in index_meta.segments { + for segment_component in &segment_components { + let path = segment_meta.relative_path(segment_component.clone()); + if index.directory().exists(&path)? { + files.insert(path); + } + } + } files.insert(Path::new("meta.json").to_path_buf()); files.insert(Path::new(".managed.json").to_path_buf()); Ok(files) diff --git a/quickwit/quickwit-doc-mapper/src/doc_mapper/field_mapping_entry.rs b/quickwit/quickwit-doc-mapper/src/doc_mapper/field_mapping_entry.rs index 2f3e755348e..dcd1c9f6d89 100644 --- a/quickwit/quickwit-doc-mapper/src/doc_mapper/field_mapping_entry.rs +++ b/quickwit/quickwit-doc-mapper/src/doc_mapper/field_mapping_entry.rs @@ -783,6 +783,7 @@ fn deserialize_mapping_type( let json_options: QuickwitJsonOptions = serde_json::from_value(json)?; Ok(FieldMappingType::Json(json_options, cardinality)) } + Type::Custom => bail!("custom fields are not supported in quickwit yet."), } } diff --git a/quickwit/quickwit-doc-mapper/src/doc_mapper/field_mapping_type.rs b/quickwit/quickwit-doc-mapper/src/doc_mapper/field_mapping_type.rs index 9e258a3141b..45b957340cd 100644 --- a/quickwit/quickwit-doc-mapper/src/doc_mapper/field_mapping_type.rs +++ b/quickwit/quickwit-doc-mapper/src/doc_mapper/field_mapping_type.rs @@ -138,6 +138,7 @@ fn primitive_type_to_str(primitive_type: &Type) -> &'static str { Type::Facet => { unimplemented!("Facets are not supported by quickwit at the moment.") } + Type::Custom => "custom", } } diff --git a/quickwit/quickwit-doc-mapper/src/doc_mapper/tantivy_val_to_json.rs b/quickwit/quickwit-doc-mapper/src/doc_mapper/tantivy_val_to_json.rs index 56b63dd71d7..c4e13862352 100644 --- a/quickwit/quickwit-doc-mapper/src/doc_mapper/tantivy_val_to_json.rs +++ b/quickwit/quickwit-doc-mapper/src/doc_mapper/tantivy_val_to_json.rs @@ -212,6 +212,7 @@ pub fn tantivy_value_to_json(value: TantivyValue) -> JsonValue { .expect("Invalid datetime is not allowed."), TantivyValue::Facet(facet) => JsonValue::String(facet.to_string()), TantivyValue::Bytes(bytes) => BinaryFormat::Base64.format_to_json(&bytes), + TantivyValue::Custom(bytes) => BinaryFormat::Base64.format_to_json(&bytes), TantivyValue::IpAddr(ip_v6) => { let ip_str = if let Some(ip_v4) = ip_v6.to_ipv4_mapped() { ip_v4.to_string() diff --git a/quickwit/quickwit-indexing/src/actors/packager.rs b/quickwit/quickwit-indexing/src/actors/packager.rs index 69c6d13a5d7..c9a058004bb 100644 --- a/quickwit/quickwit-indexing/src/actors/packager.rs +++ b/quickwit/quickwit-indexing/src/actors/packager.rs @@ -28,7 +28,7 @@ use quickwit_directories::write_hotcache; use quickwit_doc_mapper::NamedField; use quickwit_doc_mapper::tag_pruning::append_to_tag_set; use quickwit_proto::search::{ListFieldsEntry, ListFieldsMetadata, ListFieldsType}; -use tantivy::index::FieldMetadata; +use tantivy::index::{FieldMetadata, SegmentComponent}; use tantivy::schema::{FieldType, Type}; use tantivy::{InvertedIndexReader, ReloadPolicy, SegmentMeta}; use tokio::runtime::Handle; @@ -187,15 +187,24 @@ fn list_split_files( scratch_directory: &TempDirectory, ) -> io::Result> { let mut split_files = vec![scratch_directory.path().join("meta.json")]; + let segment_components = [ + SegmentComponent::Postings, + SegmentComponent::Positions, + SegmentComponent::Terms, + SegmentComponent::Store, + SegmentComponent::FastFields, + SegmentComponent::FieldNorms, + SegmentComponent::Delete, + ]; // list the segment files for segment_meta in segment_metas { - for relative_path in segment_meta.list_files() { + for segment_component in &segment_components { + let relative_path = segment_meta.relative_path(segment_component.clone()); let filepath = scratch_directory.path().join(relative_path); if filepath.try_exists()? { // If the file is missing, this is fine. - // segment_meta.list_files() may actually returns files that - // may not exist. + // Segment metadata may reference optional files that do not exist. split_files.push(filepath); } } @@ -361,6 +370,7 @@ fn tantivy_type_to_list_field_type(typ: Type) -> ListFieldsType { Type::Json => ListFieldsType::Json, Type::Str => ListFieldsType::Str, Type::U64 => ListFieldsType::U64, + Type::Custom => ListFieldsType::Custom, } } diff --git a/quickwit/quickwit-proto/build.rs b/quickwit/quickwit-proto/build.rs index bf7977e00dd..ee4ec457daa 100644 --- a/quickwit/quickwit-proto/build.rs +++ b/quickwit/quickwit-proto/build.rs @@ -235,6 +235,10 @@ fn main() -> Result<(), Box> { ) .type_attribute("PartialHit.sort_value", "#[derive(Copy)]") .type_attribute("SortByValue", "#[derive(Ord, PartialOrd)]") + .type_attribute("CalculatedPredicate", "#[derive(Hash, Eq)]") + .type_attribute("CalculatedPredicateExpr", "#[derive(Hash, Eq)]") + .type_attribute("CalculatedPredicateExpr.node", "#[derive(Hash, Eq)]") + .type_attribute("CalculatedPredicateFuncCall", "#[derive(Hash, Eq)]") .type_attribute("SearchRequest", "#[derive(Hash, Eq)]") .type_attribute("PartialHit", "#[derive(Hash, Eq)]") .out_dir("src/codegen/quickwit") diff --git a/quickwit/quickwit-proto/protos/quickwit/search.proto b/quickwit/quickwit-proto/protos/quickwit/search.proto index d266f889508..68d4b38ad4c 100644 --- a/quickwit/quickwit-proto/protos/quickwit/search.proto +++ b/quickwit/quickwit-proto/protos/quickwit/search.proto @@ -200,10 +200,85 @@ enum ListFieldsType { BYTES = 7; IP_ADDR = 8; JSON = 9; + CUSTOM = 10; } // -- Search ------------------- +// A predicate expression evaluated against fast-field values on each leaf +// segment. The expression is lowered to tantivy::query::CalculatedPredicateQuery +// by the search leaf. +message CalculatedPredicate { + CalculatedPredicateExpr expr = 1; +} + +message CalculatedPredicateExpr { + oneof node { + CalculatedPredicateLiteral literal = 1; + // Fast-field name read by the calculated predicate query. + string variable = 2; + CalculatedPredicateFuncCall func_call = 3; + } +} + +message CalculatedPredicateLiteral { + oneof value { + int64 int_value = 1; + uint64 uint_value = 2; + // IEEE-754 bits for an f64 literal. This keeps SearchRequest Eq/Hash-safe. + fixed64 double_value_bits = 3; + string string_value = 4; + bool bool_value = 5; + } +} + +message CalculatedPredicateFuncCall { + enum Function { + FUNCTION_UNSPECIFIED = 0; + FUNCTION_ABS = 1; + FUNCTION_AND = 2; + FUNCTION_CEIL = 3; + FUNCTION_CONCAT = 4; + FUNCTION_ADD = 5; + FUNCTION_DIVIDE = 6; + FUNCTION_EQ = 7; + FUNCTION_FLOOR = 8; + FUNCTION_GT = 9; + FUNCTION_GT_EQ = 10; + FUNCTION_IF = 11; + FUNCTION_INT_MOD = 12; + FUNCTION_LEFT = 13; + FUNCTION_LT = 14; + FUNCTION_LT_EQ = 15; + FUNCTION_IS_NOT_NULL = 16; + FUNCTION_IS_NULL = 17; + FUNCTION_LOWER = 18; + FUNCTION_MAX = 19; + FUNCTION_MIN = 20; + FUNCTION_MULTIPLY = 21; + FUNCTION_NEQ = 22; + FUNCTION_NOT = 23; + FUNCTION_OR = 24; + FUNCTION_POW = 25; + FUNCTION_SQRT = 26; + FUNCTION_REGEXP_EXTRACT = 27; + FUNCTION_REGEXP_LIKE = 28; + FUNCTION_RIGHT = 29; + FUNCTION_ROUND = 30; + FUNCTION_SPLIT_AFTER = 31; + FUNCTION_SPLIT_BEFORE = 32; + FUNCTION_SUBTRACT = 33; + FUNCTION_SUBSTRING = 34; + FUNCTION_SUBSTRING_COUNT = 35; + FUNCTION_TEXT_JOIN = 36; + FUNCTION_TRIM = 37; + FUNCTION_UPPER = 38; + } + + Function function = 1; + repeated CalculatedPredicateExpr args = 2; +} + message SearchRequest { // Index ID patterns repeated string index_id_patterns = 1; @@ -280,6 +355,9 @@ message SearchRequest { // Scheduling priority for leaf search execution. Negative values are allowed, // and lower values have higher priority. Callers that omit it get priority 0. int32 priority = 21; + + // Predicate expression evaluated by Tantivy against fast fields on each leaf. + optional CalculatedPredicate calculated_predicate = 22; } enum CountHits { diff --git a/quickwit/quickwit-proto/src/codegen/quickwit/quickwit.search.rs b/quickwit/quickwit-proto/src/codegen/quickwit/quickwit.search.rs index 347892fb124..c21482516fa 100644 --- a/quickwit/quickwit-proto/src/codegen/quickwit/quickwit.search.rs +++ b/quickwit/quickwit-proto/src/codegen/quickwit/quickwit.search.rs @@ -152,6 +152,225 @@ pub struct ListFieldsEntry { #[prost(uint64, tag = "8")] pub num_splits: u64, } +/// A predicate expression evaluated against fast-field values on each leaf +/// segment. The expression is lowered to tantivy::query::CalculatedPredicateQuery +/// by the search leaf. +#[derive(serde::Serialize, serde::Deserialize, utoipa::ToSchema)] +#[derive(Hash, Eq)] +#[derive(Clone, PartialEq, ::prost::Message)] +pub struct CalculatedPredicate { + #[prost(message, optional, tag = "1")] + pub expr: ::core::option::Option, +} +#[derive(serde::Serialize, serde::Deserialize, utoipa::ToSchema)] +#[derive(Hash, Eq)] +#[derive(Clone, PartialEq, ::prost::Message)] +pub struct CalculatedPredicateExpr { + #[prost(oneof = "calculated_predicate_expr::Node", tags = "1, 2, 3")] + pub node: ::core::option::Option, +} +/// Nested message and enum types in `CalculatedPredicateExpr`. +pub mod calculated_predicate_expr { + #[derive(serde::Serialize, serde::Deserialize, utoipa::ToSchema)] + #[derive(Hash, Eq)] + #[serde(rename_all = "snake_case")] + #[derive(Clone, PartialEq, ::prost::Oneof)] + pub enum Node { + #[prost(message, tag = "1")] + Literal(super::CalculatedPredicateLiteral), + /// Fast-field name read by the calculated predicate query. + #[prost(string, tag = "2")] + Variable(::prost::alloc::string::String), + #[prost(message, tag = "3")] + FuncCall(super::CalculatedPredicateFuncCall), + } +} +#[derive(serde::Serialize, serde::Deserialize, utoipa::ToSchema)] +#[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)] +pub struct CalculatedPredicateLiteral { + #[prost(oneof = "calculated_predicate_literal::Value", tags = "1, 2, 3, 4, 5")] + pub value: ::core::option::Option, +} +/// Nested message and enum types in `CalculatedPredicateLiteral`. +pub mod calculated_predicate_literal { + #[derive(serde::Serialize, serde::Deserialize, utoipa::ToSchema)] + #[serde(rename_all = "snake_case")] + #[derive(Clone, PartialEq, Eq, Hash, ::prost::Oneof)] + pub enum Value { + #[prost(int64, tag = "1")] + IntValue(i64), + #[prost(uint64, tag = "2")] + UintValue(u64), + /// IEEE-754 bits for an f64 literal. This keeps SearchRequest Eq/Hash-safe. + #[prost(fixed64, tag = "3")] + DoubleValueBits(u64), + #[prost(string, tag = "4")] + StringValue(::prost::alloc::string::String), + #[prost(bool, tag = "5")] + BoolValue(bool), + } +} +#[derive(serde::Serialize, serde::Deserialize, utoipa::ToSchema)] +#[derive(Hash, Eq)] +#[derive(Clone, PartialEq, ::prost::Message)] +pub struct CalculatedPredicateFuncCall { + #[prost(enumeration = "calculated_predicate_func_call::Function", tag = "1")] + pub function: i32, + #[prost(message, repeated, tag = "2")] + pub args: ::prost::alloc::vec::Vec, +} +/// Nested message and enum types in `CalculatedPredicateFuncCall`. +pub mod calculated_predicate_func_call { + #[derive(serde::Serialize, serde::Deserialize, utoipa::ToSchema)] + #[serde(rename_all = "snake_case")] + #[derive( + Clone, + Copy, + Debug, + PartialEq, + Eq, + Hash, + PartialOrd, + Ord, + ::prost::Enumeration + )] + #[repr(i32)] + pub enum Function { + Unspecified = 0, + Abs = 1, + And = 2, + Ceil = 3, + Concat = 4, + Add = 5, + Divide = 6, + Eq = 7, + Floor = 8, + Gt = 9, + GtEq = 10, + If = 11, + IntMod = 12, + Left = 13, + Lt = 14, + LtEq = 15, + IsNotNull = 16, + IsNull = 17, + Lower = 18, + Max = 19, + Min = 20, + Multiply = 21, + Neq = 22, + Not = 23, + Or = 24, + Pow = 25, + Sqrt = 26, + RegexpExtract = 27, + RegexpLike = 28, + Right = 29, + Round = 30, + SplitAfter = 31, + SplitBefore = 32, + Subtract = 33, + Substring = 34, + SubstringCount = 35, + TextJoin = 36, + Trim = 37, + Upper = 38, + } + impl Function { + /// String value of the enum field names used in the ProtoBuf definition. + /// + /// The values are not transformed in any way and thus are considered stable + /// (if the ProtoBuf definition does not change) and safe for programmatic use. + pub fn as_str_name(&self) -> &'static str { + match self { + Self::Unspecified => "FUNCTION_UNSPECIFIED", + Self::Abs => "FUNCTION_ABS", + Self::And => "FUNCTION_AND", + Self::Ceil => "FUNCTION_CEIL", + Self::Concat => "FUNCTION_CONCAT", + Self::Add => "FUNCTION_ADD", + Self::Divide => "FUNCTION_DIVIDE", + Self::Eq => "FUNCTION_EQ", + Self::Floor => "FUNCTION_FLOOR", + Self::Gt => "FUNCTION_GT", + Self::GtEq => "FUNCTION_GT_EQ", + Self::If => "FUNCTION_IF", + Self::IntMod => "FUNCTION_INT_MOD", + Self::Left => "FUNCTION_LEFT", + Self::Lt => "FUNCTION_LT", + Self::LtEq => "FUNCTION_LT_EQ", + Self::IsNotNull => "FUNCTION_IS_NOT_NULL", + Self::IsNull => "FUNCTION_IS_NULL", + Self::Lower => "FUNCTION_LOWER", + Self::Max => "FUNCTION_MAX", + Self::Min => "FUNCTION_MIN", + Self::Multiply => "FUNCTION_MULTIPLY", + Self::Neq => "FUNCTION_NEQ", + Self::Not => "FUNCTION_NOT", + Self::Or => "FUNCTION_OR", + Self::Pow => "FUNCTION_POW", + Self::Sqrt => "FUNCTION_SQRT", + Self::RegexpExtract => "FUNCTION_REGEXP_EXTRACT", + Self::RegexpLike => "FUNCTION_REGEXP_LIKE", + Self::Right => "FUNCTION_RIGHT", + Self::Round => "FUNCTION_ROUND", + Self::SplitAfter => "FUNCTION_SPLIT_AFTER", + Self::SplitBefore => "FUNCTION_SPLIT_BEFORE", + Self::Subtract => "FUNCTION_SUBTRACT", + Self::Substring => "FUNCTION_SUBSTRING", + Self::SubstringCount => "FUNCTION_SUBSTRING_COUNT", + Self::TextJoin => "FUNCTION_TEXT_JOIN", + Self::Trim => "FUNCTION_TRIM", + Self::Upper => "FUNCTION_UPPER", + } + } + /// Creates an enum from field names used in the ProtoBuf definition. + pub fn from_str_name(value: &str) -> ::core::option::Option { + match value { + "FUNCTION_UNSPECIFIED" => Some(Self::Unspecified), + "FUNCTION_ABS" => Some(Self::Abs), + "FUNCTION_AND" => Some(Self::And), + "FUNCTION_CEIL" => Some(Self::Ceil), + "FUNCTION_CONCAT" => Some(Self::Concat), + "FUNCTION_ADD" => Some(Self::Add), + "FUNCTION_DIVIDE" => Some(Self::Divide), + "FUNCTION_EQ" => Some(Self::Eq), + "FUNCTION_FLOOR" => Some(Self::Floor), + "FUNCTION_GT" => Some(Self::Gt), + "FUNCTION_GT_EQ" => Some(Self::GtEq), + "FUNCTION_IF" => Some(Self::If), + "FUNCTION_INT_MOD" => Some(Self::IntMod), + "FUNCTION_LEFT" => Some(Self::Left), + "FUNCTION_LT" => Some(Self::Lt), + "FUNCTION_LT_EQ" => Some(Self::LtEq), + "FUNCTION_IS_NOT_NULL" => Some(Self::IsNotNull), + "FUNCTION_IS_NULL" => Some(Self::IsNull), + "FUNCTION_LOWER" => Some(Self::Lower), + "FUNCTION_MAX" => Some(Self::Max), + "FUNCTION_MIN" => Some(Self::Min), + "FUNCTION_MULTIPLY" => Some(Self::Multiply), + "FUNCTION_NEQ" => Some(Self::Neq), + "FUNCTION_NOT" => Some(Self::Not), + "FUNCTION_OR" => Some(Self::Or), + "FUNCTION_POW" => Some(Self::Pow), + "FUNCTION_SQRT" => Some(Self::Sqrt), + "FUNCTION_REGEXP_EXTRACT" => Some(Self::RegexpExtract), + "FUNCTION_REGEXP_LIKE" => Some(Self::RegexpLike), + "FUNCTION_RIGHT" => Some(Self::Right), + "FUNCTION_ROUND" => Some(Self::Round), + "FUNCTION_SPLIT_AFTER" => Some(Self::SplitAfter), + "FUNCTION_SPLIT_BEFORE" => Some(Self::SplitBefore), + "FUNCTION_SUBTRACT" => Some(Self::Subtract), + "FUNCTION_SUBSTRING" => Some(Self::Substring), + "FUNCTION_SUBSTRING_COUNT" => Some(Self::SubstringCount), + "FUNCTION_TEXT_JOIN" => Some(Self::TextJoin), + "FUNCTION_TRIM" => Some(Self::Trim), + "FUNCTION_UPPER" => Some(Self::Upper), + _ => None, + } + } + } +} #[derive(serde::Serialize, serde::Deserialize, utoipa::ToSchema)] #[derive(Hash, Eq)] #[derive(Clone, PartialEq, ::prost::Message)] @@ -216,6 +435,9 @@ pub struct SearchRequest { #[prost(int32, tag = "21")] #[serde(default)] pub priority: i32, + /// Predicate expression evaluated by Tantivy against fast fields on each leaf. + #[prost(message, optional, tag = "22")] + pub calculated_predicate: ::core::option::Option, } #[derive(serde::Serialize, serde::Deserialize, utoipa::ToSchema)] #[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)] @@ -789,6 +1011,7 @@ pub enum ListFieldsType { Bytes = 7, IpAddr = 8, Json = 9, + Custom = 10, } impl ListFieldsType { /// String value of the enum field names used in the ProtoBuf definition. @@ -807,6 +1030,7 @@ impl ListFieldsType { Self::Bytes => "BYTES", Self::IpAddr => "IP_ADDR", Self::Json => "JSON", + Self::Custom => "CUSTOM", } } /// Creates an enum from field names used in the ProtoBuf definition. @@ -822,6 +1046,7 @@ impl ListFieldsType { "BYTES" => Some(Self::Bytes), "IP_ADDR" => Some(Self::IpAddr), "JSON" => Some(Self::Json), + "CUSTOM" => Some(Self::Custom), _ => None, } } diff --git a/quickwit/quickwit-query/src/query_ast/range_query.rs b/quickwit/quickwit-query/src/query_ast/range_query.rs index 346f02feed3..7212a816935 100644 --- a/quickwit/quickwit-query/src/query_ast/range_query.rs +++ b/quickwit/quickwit-query/src/query_ast/range_query.rs @@ -273,6 +273,12 @@ impl BuildTantivyAst for RangeQuery { ) .into() } + tantivy::schema::FieldType::Custom(_) => { + return Err(InvalidQuery::RangeQueryNotSupportedForField { + value_type: "custom", + field_name: field_entry.name().to_string(), + }); + } }) } } diff --git a/quickwit/quickwit-query/src/query_ast/utils.rs b/quickwit/quickwit-query/src/query_ast/utils.rs index 71f3d4e5303..a5d6fd15bec 100644 --- a/quickwit/quickwit-query/src/query_ast/utils.rs +++ b/quickwit/quickwit-query/src/query_ast/utils.rs @@ -192,6 +192,9 @@ fn compute_query_with_field( let term = Term::from_field_bytes(field, &buffer[..]); Ok(make_term_query(term)) } + FieldType::Custom(_) => Err(InvalidQuery::SchemaError( + "custom fields are not supported in Quickwit queries".to_string(), + )), } } diff --git a/quickwit/quickwit-search/src/calculated_predicate.rs b/quickwit/quickwit-search/src/calculated_predicate.rs new file mode 100644 index 00000000000..647d95b3490 --- /dev/null +++ b/quickwit/quickwit-search/src/calculated_predicate.rs @@ -0,0 +1,232 @@ +// Copyright 2021-Present Datadog, Inc. +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +use std::collections::HashSet; + +use anyhow::{Context, bail}; +use quickwit_proto::search::calculated_predicate_expr::Node; +use quickwit_proto::search::calculated_predicate_func_call::Function as ProtoFunction; +use quickwit_proto::search::calculated_predicate_literal::Value; +use quickwit_proto::search::{ + CalculatedPredicate, CalculatedPredicateExpr, CalculatedPredicateFuncCall, + CalculatedPredicateLiteral, +}; +use tantivy::jitexpr::ast::{Function, Literal, UntypedExpr}; +use tantivy::query::CalculatedPredicateQuery; + +pub(crate) fn calculated_predicate_query( + predicate: &CalculatedPredicate, +) -> anyhow::Result<(CalculatedPredicateQuery, HashSet)> { + let expression = predicate + .expr + .as_ref() + .context("calculated predicate is missing its root expression")?; + let mut variables = HashSet::new(); + let expression = to_untyped_expr(expression, &mut variables)?; + let query = + CalculatedPredicateQuery::new(expression).context("invalid calculated predicate")?; + Ok((query, variables)) +} + +fn to_untyped_expr( + expression: &CalculatedPredicateExpr, + variables: &mut HashSet, +) -> anyhow::Result { + match expression.node.as_ref() { + Some(Node::Literal(literal)) => literal_to_untyped_expr(literal), + Some(Node::Variable(variable)) => { + variables.insert(variable.clone()); + Ok(UntypedExpr::variable(variable)) + } + Some(Node::FuncCall(func_call)) => func_call_to_untyped_expr(func_call, variables), + None => bail!("calculated predicate expression is missing a node"), + } +} + +fn literal_to_untyped_expr(literal: &CalculatedPredicateLiteral) -> anyhow::Result { + let literal = match literal.value.as_ref() { + Some(Value::IntValue(value)) => Literal::I64(*value), + Some(Value::UintValue(value)) => Literal::U64(*value), + Some(Value::DoubleValueBits(value)) => Literal::F64(f64::from_bits(*value)), + Some(Value::StringValue(value)) => Literal::from(value.as_str()), + Some(Value::BoolValue(value)) => Literal::Bool(*value), + None => bail!("calculated predicate literal is missing a value"), + }; + Ok(UntypedExpr::Literal(literal)) +} + +fn func_call_to_untyped_expr( + func_call: &CalculatedPredicateFuncCall, + variables: &mut HashSet, +) -> anyhow::Result { + let function = ProtoFunction::try_from(func_call.function) + .context("unknown calculated predicate function")?; + let function = to_jitexpr_function(function)?; + let mut args = Vec::with_capacity(func_call.args.len()); + for arg in &func_call.args { + args.push(to_untyped_expr(arg, variables)?); + } + Ok(UntypedExpr::Call { function, args }) +} + +fn to_jitexpr_function(function: ProtoFunction) -> anyhow::Result { + match function { + ProtoFunction::Unspecified => bail!("calculated predicate function is unspecified"), + ProtoFunction::Abs => Ok(Function::Abs), + ProtoFunction::And => Ok(Function::And), + ProtoFunction::Ceil => Ok(Function::Ceil), + ProtoFunction::Concat => Ok(Function::Concat), + ProtoFunction::Add => Ok(Function::Add), + ProtoFunction::Divide => Ok(Function::Divide), + ProtoFunction::Eq => Ok(Function::Eq), + ProtoFunction::Floor => Ok(Function::Floor), + ProtoFunction::Gt => Ok(Function::Gt), + ProtoFunction::GtEq => Ok(Function::GtEq), + ProtoFunction::If => Ok(Function::If), + ProtoFunction::IntMod => Ok(Function::IntMod), + ProtoFunction::Left => Ok(Function::Left), + ProtoFunction::Lt => Ok(Function::Lt), + ProtoFunction::LtEq => Ok(Function::LtEq), + ProtoFunction::IsNotNull => Ok(Function::IsNotNull), + ProtoFunction::IsNull => Ok(Function::IsNull), + ProtoFunction::Lower => Ok(Function::Lower), + ProtoFunction::Max => Ok(Function::Max), + ProtoFunction::Min => Ok(Function::Min), + ProtoFunction::Multiply => Ok(Function::Multiply), + ProtoFunction::Neq => Ok(Function::Neq), + ProtoFunction::Not => Ok(Function::Not), + ProtoFunction::Or => Ok(Function::Or), + ProtoFunction::Pow => Ok(Function::Pow), + ProtoFunction::Sqrt => Ok(Function::Sqrt), + ProtoFunction::RegexpExtract => Ok(Function::RegexpExtract), + ProtoFunction::RegexpLike => Ok(Function::RegexpLike), + ProtoFunction::Right => Ok(Function::Right), + ProtoFunction::Round => Ok(Function::Round), + ProtoFunction::SplitAfter => Ok(Function::SplitAfter), + ProtoFunction::SplitBefore => Ok(Function::SplitBefore), + ProtoFunction::Subtract => Ok(Function::Subtract), + ProtoFunction::Substring => Ok(Function::Substring), + ProtoFunction::SubstringCount => Ok(Function::SubstringCount), + ProtoFunction::TextJoin => Ok(Function::TextJoin), + ProtoFunction::Trim => Ok(Function::Trim), + ProtoFunction::Upper => Ok(Function::Upper), + } +} + +#[cfg(test)] +mod tests { + use quickwit_proto::search::calculated_predicate_expr::Node; + use quickwit_proto::search::calculated_predicate_func_call::Function as ProtoFunction; + use quickwit_proto::search::calculated_predicate_literal::Value; + use quickwit_proto::search::{ + CalculatedPredicate, CalculatedPredicateExpr, CalculatedPredicateFuncCall, + CalculatedPredicateLiteral, + }; + + use super::calculated_predicate_query; + + fn variable(name: &str) -> CalculatedPredicateExpr { + CalculatedPredicateExpr { + node: Some(Node::Variable(name.to_string())), + } + } + + fn literal(value: Value) -> CalculatedPredicateExpr { + CalculatedPredicateExpr { + node: Some(Node::Literal(CalculatedPredicateLiteral { + value: Some(value), + })), + } + } + + fn call( + function: ProtoFunction, + args: Vec, + ) -> CalculatedPredicateExpr { + CalculatedPredicateExpr { + node: Some(Node::FuncCall(CalculatedPredicateFuncCall { + function: function as i32, + args, + })), + } + } + + #[test] + fn test_calculated_predicate_query_accepts_boolean_expression() { + let predicate = CalculatedPredicate { + expr: Some(call( + ProtoFunction::Gt, + vec![variable("latency_ms"), literal(Value::IntValue(100))], + )), + }; + + let (_query, variables) = calculated_predicate_query(&predicate).unwrap(); + assert_eq!(variables, ["latency_ms".to_string()].into_iter().collect()); + } + + #[test] + fn test_calculated_predicate_query_rejects_missing_root() { + let predicate = CalculatedPredicate { expr: None }; + + assert!( + calculated_predicate_query(&predicate) + .unwrap_err() + .to_string() + .contains("root expression") + ); + } + + #[test] + fn test_calculated_predicate_query_rejects_non_boolean_expression() { + let predicate = CalculatedPredicate { + expr: Some(literal(Value::StringValue("hello".to_string()))), + }; + + assert!( + calculated_predicate_query(&predicate) + .unwrap_err() + .to_string() + .contains("invalid calculated predicate") + ); + } + + #[test] + fn test_calculated_predicate_query_rejects_unspecified_function() { + let predicate = CalculatedPredicate { + expr: Some(call(ProtoFunction::Unspecified, Vec::new())), + }; + + assert!( + calculated_predicate_query(&predicate) + .unwrap_err() + .to_string() + .contains("unspecified") + ); + } + + #[test] + fn test_double_literals_are_encoded_as_bits() { + let predicate = CalculatedPredicate { + expr: Some(call( + ProtoFunction::LtEq, + vec![ + variable("score"), + literal(Value::DoubleValueBits(12.5f64.to_bits())), + ], + )), + }; + + calculated_predicate_query(&predicate).unwrap(); + } +} diff --git a/quickwit/quickwit-search/src/leaf.rs b/quickwit/quickwit-search/src/leaf.rs index 245a18db75a..e088225faea 100644 --- a/quickwit/quickwit-search/src/leaf.rs +++ b/quickwit/quickwit-search/src/leaf.rs @@ -52,12 +52,14 @@ use tantivy::aggregation::agg_req::{AggregationVariants, Aggregations}; use tantivy::collector::Collector; use tantivy::fastfield::FastFieldReaders; use tantivy::index::SegmentId; +use tantivy::query::{BooleanQuery, Occur, Query}; use tantivy::schema::Field; use tantivy::{DateTime, Index, ReloadPolicy, Searcher, TantivyError, Term}; use tokio::task::{JoinError, JoinSet}; use tokio_util::sync::CancellationToken; use tracing::*; +use crate::calculated_predicate::calculated_predicate_query; use crate::collector::{IncrementalCollector, make_collector_for_split, make_merge_collector}; use crate::leaf_cache::LeafSearchCache; use crate::metrics::{ @@ -739,12 +741,31 @@ async fn leaf_search_single_split( )) }; let split_schema = index.schema(); - let (query, mut warmup_info) = ctx.doc_mapper.query( + let (mut query, mut warmup_info) = ctx.doc_mapper.query( split_schema.clone(), query_ast.clone(), false, predicate_cache, )?; + if let Some(calculated_predicate) = search_request.calculated_predicate.as_ref() { + let (predicate_query, predicate_variables) = + calculated_predicate_query(calculated_predicate) + .map_err(|err| SearchError::InvalidQuery(err.to_string()))?; + warmup_info + .fast_fields + .extend( + predicate_variables + .into_iter() + .map(|name| FastFieldWarmupInfo { + name, + with_subfields: false, + }), + ); + query = Box::new(BooleanQuery::new(vec![ + (Occur::Must, query), + (Occur::Must, Box::new(predicate_query) as Box), + ])); + } let collector_warmup_info = collector.warmup_info(); warmup_info.merge(collector_warmup_info); diff --git a/quickwit/quickwit-search/src/lib.rs b/quickwit/quickwit-search/src/lib.rs index 2d891dbfa65..f3d7128cccb 100644 --- a/quickwit/quickwit-search/src/lib.rs +++ b/quickwit/quickwit-search/src/lib.rs @@ -17,6 +17,7 @@ #![allow(clippy::bool_assert_comparison)] #![deny(clippy::disallowed_methods)] +mod calculated_predicate; mod client; mod cluster_client; mod collector; diff --git a/quickwit/quickwit-search/src/root.rs b/quickwit/quickwit-search/src/root.rs index c31a848f8ea..7242398ed13 100644 --- a/quickwit/quickwit-search/src/root.rs +++ b/quickwit/quickwit-search/src/root.rs @@ -377,6 +377,7 @@ fn simplify_search_request_for_scroll_api(req: &SearchRequest) -> crate::Result< ignore_missing_indexes: req.ignore_missing_indexes, skip_aggregation_finalization: false, priority: req.priority, + calculated_predicate: req.calculated_predicate.clone(), }) } @@ -685,6 +686,9 @@ pub fn is_metadata_count_request_with_ast(query_ast: &QueryAst, request: &Search if request.aggregation_request.is_some() || !request.snippet_fields.is_empty() { return false; } + if request.calculated_predicate.is_some() { + return false; + } true } @@ -1919,7 +1923,7 @@ mod tests { MockMetastoreService, }; use quickwit_proto::search::{ - ScrollRequest, SortByValue, SortOrder, SortValue, SplitSearchError, + CalculatedPredicate, ScrollRequest, SortByValue, SortOrder, SortValue, SplitSearchError, }; use quickwit_query::query_ast::{qast_helper, qast_json_helper, query_ast_from_user_text}; use tantivy::schema::{FAST, STORED, TEXT}; @@ -1927,6 +1931,27 @@ mod tests { use super::*; use crate::{MockSearchService, searcher_pool_for_test}; + #[test] + fn test_metadata_count_request_with_calculated_predicate() { + let query_ast = QueryAst::MatchAll; + let mut search_request = SearchRequest { + query_ast: serde_json::to_string(&query_ast).unwrap(), + ..Default::default() + }; + assert!(is_metadata_count_request_with_ast( + &query_ast, + &search_request + )); + + search_request.calculated_predicate = Some(CalculatedPredicate { expr: None }); + + assert!(!is_metadata_count_request_with_ast( + &query_ast, + &search_request + )); + assert!(!is_metadata_count_request(&search_request)); + } + #[track_caller] fn check_snippet_fields_validation(snippet_fields: &[String]) -> anyhow::Result<()> { let mut schema_builder = Schema::builder(); diff --git a/quickwit/quickwit-serve/src/elasticsearch_api/model/field_capability.rs b/quickwit/quickwit-serve/src/elasticsearch_api/model/field_capability.rs index 292f7202644..2c281a290ca 100644 --- a/quickwit/quickwit-serve/src/elasticsearch_api/model/field_capability.rs +++ b/quickwit/quickwit-serve/src/elasticsearch_api/model/field_capability.rs @@ -160,6 +160,7 @@ pub fn convert_to_es_field_capabilities_response( ListFieldsType::Date => vec![FieldCapabilityEntryType::DateNanos], ListFieldsType::Facet => continue, ListFieldsType::Json => continue, + ListFieldsType::Custom => continue, ListFieldsType::Bytes => vec![FieldCapabilityEntryType::Binary], ListFieldsType::IpAddr => vec![FieldCapabilityEntryType::Ip], }; diff --git a/quickwit/quickwit-serve/src/elasticsearch_api/model/mappings.rs b/quickwit/quickwit-serve/src/elasticsearch_api/model/mappings.rs index 15bcc16f060..c1f98e02442 100644 --- a/quickwit/quickwit-serve/src/elasticsearch_api/model/mappings.rs +++ b/quickwit/quickwit-serve/src/elasticsearch_api/model/mappings.rs @@ -159,7 +159,7 @@ fn es_type_from_list_field_type(field_type: ListFieldsType) -> Option<&'static s ListFieldsType::Date => Some("date"), ListFieldsType::Bytes => Some("binary"), ListFieldsType::IpAddr => Some("ip"), - ListFieldsType::Facet | ListFieldsType::Json => None, + ListFieldsType::Facet | ListFieldsType::Json | ListFieldsType::Custom => None, } }