diff --git a/quickwit/quickwit-doc-mapper/src/doc_mapper/borrowed_doc_tests.rs b/quickwit/quickwit-doc-mapper/src/doc_mapper/borrowed_doc_tests.rs new file mode 100644 index 00000000000..fdd4ba2447c --- /dev/null +++ b/quickwit/quickwit-doc-mapper/src/doc_mapper/borrowed_doc_tests.rs @@ -0,0 +1,228 @@ +// 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. + +//! Differential tests: `DocMapper::doc_from_borrowed_json` must produce the same partitions, +//! documents and errors as `DocMapper::doc_from_json_obj`. + +use tantivy::TantivyDocument as Document; +use tantivy::schema::OwnedValue; + +use crate::doc_mapper::{BorrowedJsonDoc, DocMapper, JsonObject, RandomJsonDocs}; + +/// A field mapping exercising every leaf type, arrays, objects, concatenate fields, coercion and +/// the partition key. `{mode}` is replaced by each mode. +const DOC_MAPPING_TEMPLATE: &str = r#"{ + "mode": "{mode}", + "store_source": {store_source}, + "store_document_size": true, + "index_field_presence": true, + "timestamp_field": "timestamp", + "partition_key": "service,hash_mod(resource.host, 7)", + "field_mappings": [ + {"name": "timestamp", "type": "datetime", "fast": true, + "input_formats": ["rfc3339", "unix_timestamp", "%Y-%m-%d %H:%M:%S"]}, + {"name": "service", "type": "text", "tokenizer": "raw", "fast": true}, + {"name": "body", "type": "text"}, + {"name": "count", "type": "u64"}, + {"name": "delta", "type": "i64", "coerce": false}, + {"name": "ratio", "type": "f64"}, + {"name": "flag", "type": "bool"}, + {"name": "ip", "type": "ip"}, + {"name": "payload", "type": "bytes"}, + {"name": "hex_payload", "type": "bytes", "input_format": "hex"}, + {"name": "tags", "type": "array", "tokenizer": "raw"}, + {"name": "values", "type": "array"}, + {"name": "attributes", "type": "json"}, + {"name": "events", "type": "array"}, + {"name": "resource", "type": "object", "field_mappings": [ + {"name": "host", "type": "text", "tokenizer": "raw"}, + {"name": "pid", "type": "u64"}, + {"name": "inner", "type": "object", "field_mappings": [ + {"name": "zone", "type": "text"} + ]} + ]}, + {"name": "all_text", "type": "concatenate", + "concatenate_fields": ["body", "service", "attributes", "count", "flag"], + "include_dynamic_fields": {include_dynamic}} + ] +}"#; + +fn build_doc_mappers() -> Vec<(String, DocMapper)> { + let mut doc_mappers = Vec::new(); + for mode in ["dynamic", "lenient", "strict"] { + for store_source in ["true", "false"] { + let include_dynamic = if mode == "dynamic" { "true" } else { "false" }; + let doc_mapping_json = DOC_MAPPING_TEMPLATE + .replace("{mode}", mode) + .replace("{store_source}", store_source) + .replace("{include_dynamic}", include_dynamic); + let doc_mapper: DocMapper = serde_json::from_str(&doc_mapping_json).unwrap(); + doc_mappers.push((format!("{mode}/store_source={store_source}"), doc_mapper)); + } + } + // The default dynamic mapping, and the dynamic mapping used by OTel logs. + let default_doc_mapper: DocMapper = serde_json::from_str("{}").unwrap(); + doc_mappers.push(("default".to_string(), default_doc_mapper)); + let otel_like_doc_mapper: DocMapper = serde_json::from_str( + r#"{ + "mode": "dynamic", + "dynamic_mapping": {"indexed": false, "stored": true, "expand_dots": true}, + "field_mappings": [ + {"name": "timestamp_nanos", "type": "datetime", "input_formats": ["unix_timestamp"], + "fast": true}, + {"name": "service_name", "type": "text", "tokenizer": "raw", "fast": true}, + {"name": "severity_text", "type": "text", "tokenizer": "raw"}, + {"name": "body", "type": "json"}, + {"name": "attributes", "type": "json", "tokenizer": "raw"}, + {"name": "resource", "type": "object", "field_mappings": [ + {"name": "host", "type": "text"} + ]} + ] + }"#, + ) + .unwrap(); + doc_mappers.push(("otel_like".to_string(), otel_like_doc_mapper)); + doc_mappers +} + +/// Returns the field values in insertion order. Tantivy's `PartialEq` on documents ignores the +/// order of the values, which is not strict enough here. +fn ordered_field_values(document: &Document) -> Vec<(u32, OwnedValue)> { + document + .field_values() + .map(|(field, value)| (field.field_id(), OwnedValue::from(value))) + .collect() +} + +fn assert_same_conversion(doc_mapper_name: &str, doc_mapper: &DocMapper, json_doc: &str) { + let document_len = json_doc.len() as u64; + let owned_result = serde_json::from_str::(json_doc) + .map_err(|error| error.to_string()) + .map(|json_obj| doc_mapper.doc_from_json_obj(json_obj, document_len)); + let borrowed_doc_result = BorrowedJsonDoc::parse(json_doc.as_bytes()); + let borrowed_result = borrowed_doc_result + .as_ref() + .map_err(|error| error.to_string()) + .map(|borrowed_doc| doc_mapper.doc_from_borrowed_json(borrowed_doc, document_len)); + let context = || format!("doc mapper: {doc_mapper_name}, doc: {json_doc}"); + match (owned_result, borrowed_result) { + (Ok(Ok((owned_partition, owned_document))), Ok(Ok((partition, document)))) => { + assert_eq!(owned_partition, partition, "{}", context()); + assert_eq!( + ordered_field_values(&owned_document), + ordered_field_values(&document), + "{}", + context() + ); + } + (Ok(Err(owned_error)), Ok(Err(error))) => { + assert_eq!(owned_error, error, "{}", context()); + } + (Err(owned_parse_error), Err(parse_error)) => { + assert_eq!(owned_parse_error, parse_error, "{}", context()); + } + (owned_result, borrowed_result) => { + panic!( + "{}: owned: {:?}, borrowed: {:?}", + context(), + owned_result.map(|result| result.map(|(_, doc)| ordered_field_values(&doc))), + borrowed_result.map(|result| result.map(|(_, doc)| ordered_field_values(&doc))) + ); + } + } +} + +const HAND_WRITTEN_DOCS: &[&str] = &[ + r#"{}"#, + r#"{"timestamp": "2024-01-02T03:04:05Z", "service": "api", "body": "hello world"}"#, + r#"{"timestamp": 1704164645, "count": 3, "delta": -4, "ratio": 0.5, "flag": true}"#, + r#"{"timestamp": 1704164645.123, "ratio": 3, "count": 18446744073709551615}"#, + r#"{"timestamp": "2024-01-02 03:04:05", "ip": "192.168.0.1", "payload": "aGVsbG8="}"#, + r#"{"ip": "::1", "hex_payload": "deadbeef", "tags": ["a", null, "b"], "values": [1, -2]}"#, + r#"{"attributes": {"k": "v", "n": 1, "d": "2024-01-02T03:04:05+01:00", "z": [null, {}]}}"#, + r#"{"events": [{"name": "e1"}, null, {"name": "e2", "at": "2024-01-02T03:04:05Z"}]}"#, + r#"{"resource": {"host": "h1", "pid": 12, "inner": {"zone": "z1"}}}"#, + r#"{"resource": {"host": "h1", "unmapped": 1, "inner": {"zone": "z", "other": [1, "x"]}}}"#, + r#"{"resource": {"inner": {"other": null}}, "unmapped_root": {"a": {"b": "2024-01-02T03:04:05Z"}}}"#, + r#"{"service": "a", "service": "b", "count": 1, "count": 2, "x": 1, "x": {"y": 2}}"#, + r#"{"body": "esc\"aped \u00e9 \ud83d\ude00\n", "service": "\t"}"#, + r#"{"unmapped": "1 not a date", "unmapped2": "2024-13-45T00:00:00Z", "u3": ""}"#, + r#"{"unmapped": [1, -1, 1.5, 18446744073709551615, true, null, "s", [], {}]}"#, + // Errors. + r#"{"count": -1}"#, + r#"{"count": "12"}"#, + r#"{"delta": "12"}"#, + r#"{"count": 1.5}"#, + r#"{"ratio": "not a float"}"#, + r#"{"flag": "true"}"#, + r#"{"ip": "not an ip"}"#, + r#"{"ip": 1}"#, + r#"{"payload": "not base64!"}"#, + r#"{"payload": 12}"#, + r#"{"hex_payload": "xyz"}"#, + r#"{"timestamp": "yesterday"}"#, + r#"{"timestamp": true}"#, + r#"{"timestamp": [1, 2]}"#, + r#"{"body": 12}"#, + r#"{"body": ["a", "b"]}"#, + r#"{"tags": [["nested"]]}"#, + r#"{"tags": "single"}"#, + r#"{"attributes": "not an object"}"#, + r#"{"events": [1]}"#, + r#"{"resource": "not an object"}"#, + r#"{"resource": null}"#, + r#"{"resource": {"inner": 3}}"#, + r#"{"resource": {"pid": "x", "host": 1}}"#, + r#"{"zzz": 1, "count": "bad"}"#, + r#"{"aaa": 1, "count": "bad"}"#, + // Parse errors. + r#"[]"#, + r#"{"a": }"#, + r#"{"a": 1} {"b": 2}"#, +]; + +#[test] +fn test_borrowed_doc_same_as_owned_doc_hand_written() { + for (doc_mapper_name, doc_mapper) in build_doc_mappers() { + for json_doc in HAND_WRITTEN_DOCS { + assert_same_conversion(&doc_mapper_name, &doc_mapper, json_doc); + } + } +} + +#[test] +fn test_borrowed_doc_same_as_owned_doc_random() { + let doc_mappers = build_doc_mappers(); + let mut random_json_docs = RandomJsonDocs::new(0x2545_f491_4f6c_dd1d); + let mut num_successes = 0; + for _ in 0..20_000 { + let json_doc = random_json_docs.next_doc(); + for (doc_mapper_name, doc_mapper) in &doc_mappers { + assert_same_conversion(doc_mapper_name, doc_mapper, &json_doc); + } + let json_obj: JsonObject = serde_json::from_str(&json_doc).unwrap(); + if doc_mappers[0] + .1 + .doc_from_json_obj(json_obj, json_doc.len() as u64) + .is_ok() + { + num_successes += 1; + } + } + // Make sure the generator does not only produce invalid documents. + assert!( + num_successes > 1_000, + "only {num_successes} valid documents" + ); +} diff --git a/quickwit/quickwit-doc-mapper/src/doc_mapper/borrowed_json.rs b/quickwit/quickwit-doc-mapper/src/doc_mapper/borrowed_json.rs new file mode 100644 index 00000000000..f47fa1027fc --- /dev/null +++ b/quickwit/quickwit-doc-mapper/src/doc_mapper/borrowed_json.rs @@ -0,0 +1,480 @@ +// 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. + +//! A JSON tree borrowing its strings from the input buffer. +//! +//! It is the input of the allocation-light document conversion path +//! ([`crate::DocMapper::doc_from_borrowed_json`]). That path must produce exactly the same tantivy +//! documents, partitions and errors as [`crate::DocMapper::doc_from_json_obj`] applied to a +//! [`serde_json::Map`]. For this reason the tree reproduces the semantics of `serde_json::Value` +//! as built by `serde_json` (without the `arbitrary_precision` feature): +//! - object entries are ordered like in `serde_json::Map`: sorted by key (byte-wise, like +//! `BTreeMap`) by default, or in insertion order when the `preserve_order` feature of +//! `serde_json` is enabled (it is enabled by some test dependencies); +//! - duplicate keys are resolved by keeping the last value (at the position of the first occurrence +//! with `preserve_order`); +//! - numbers are stored as `serde_json::Number`, built from the same deserializer events. + +use std::borrow::Cow; +use std::fmt; +use std::sync::LazyLock; + +use serde::de::{self, Deserialize, DeserializeSeed, Deserializer, MapAccess, SeqAccess, Visitor}; +use serde_json::{Number, Value as JsonValue}; + +/// Key `serde_json` reserves to deserialize `RawValue`s when its `raw_value` feature is enabled. +/// See `BorrowedValue` deserialization. +const SERDE_JSON_RAW_VALUE_TOKEN: &str = "$serde_json::private::RawValue"; + +/// Whether `serde_json::Map` preserves insertion order, i.e. whether the `preserve_order` +/// feature of `serde_json` is enabled. Cargo feature unification can enable it from any crate of +/// the build, so we check the actual behavior. +static SERDE_JSON_PRESERVES_ORDER: LazyLock = LazyLock::new(|| { + let mut json_obj = serde_json::Map::new(); + json_obj.insert("b".to_string(), JsonValue::Null); + json_obj.insert("a".to_string(), JsonValue::Null); + json_obj.keys().next().map(String::as_str) == Some("b") +}); + +pub(crate) fn serde_json_preserves_order() -> bool { + *SERDE_JSON_PRESERVES_ORDER +} + +/// A JSON object whose entries have unique keys and are ordered like in `serde_json::Map`. +pub type BorrowedObject<'a> = Vec<(Cow<'a, str>, BorrowedValue<'a>)>; + +/// A JSON value borrowing strings from the input when they do not contain escape sequences. +#[derive(Debug, Clone, PartialEq)] +pub enum BorrowedValue<'a> { + /// JSON `null`. + Null, + /// JSON boolean. + Bool(bool), + /// JSON number, as built by `serde_json`. + Number(Number), + /// JSON string, borrowed from the input unless it contains escape sequences. + Str(Cow<'a, str>), + /// JSON array. + Array(Vec>), + /// Entries have unique keys and are ordered like in `serde_json::Map`. + Object(BorrowedObject<'a>), +} + +impl<'a> BorrowedValue<'a> { + /// Returns the value associated with `key` if `self` is an object. + pub fn get(&self, key: &str) -> Option<&BorrowedValue<'a>> { + let BorrowedValue::Object(entries) = self else { + return None; + }; + get_in_object(entries, key) + } + + /// Returns true if the value is JSON `null`. + pub fn is_null(&self) -> bool { + matches!(self, BorrowedValue::Null) + } + + /// Converts the value into an owned `serde_json::Value`. + /// + /// This allocates and is only meant for error messages and tests. + pub fn to_serde_json(&self) -> JsonValue { + match self { + BorrowedValue::Null => JsonValue::Null, + BorrowedValue::Bool(bool_val) => JsonValue::Bool(*bool_val), + BorrowedValue::Number(number) => JsonValue::Number(number.clone()), + BorrowedValue::Str(text) => JsonValue::String(text.to_string()), + BorrowedValue::Array(elements) => { + JsonValue::Array(elements.iter().map(BorrowedValue::to_serde_json).collect()) + } + BorrowedValue::Object(entries) => JsonValue::Object( + entries + .iter() + .map(|(key, value)| (key.to_string(), value.to_serde_json())) + .collect(), + ), + } + } +} + +/// Formats the value exactly like the equivalent `serde_json::Value`, so that error messages are +/// the same as the ones of the owned conversion path. +impl fmt::Display for BorrowedValue<'_> { + fn fmt(&self, formatter: &mut fmt::Formatter) -> fmt::Result { + self.to_serde_json().fmt(formatter) + } +} + +/// Looks up a key in an object. Relies on the keys being unique. +pub(crate) fn get_in_object<'b, 'a>( + object: &'b BorrowedObject<'a>, + key: &str, +) -> Option<&'b BorrowedValue<'a>> { + if *SERDE_JSON_PRESERVES_ORDER { + return object + .iter() + .find(|(entry_key, _)| entry_key.as_ref() == key) + .map(|(_, value)| value); + } + let position = object + .binary_search_by(|(entry_key, _)| entry_key.as_ref().cmp(key)) + .ok()?; + Some(&object[position].1) +} + +/// A parsed JSON document whose root is an object. +#[derive(Debug, Clone, PartialEq)] +pub struct BorrowedJsonDoc<'a> { + root: BorrowedObject<'a>, +} + +impl<'a> BorrowedJsonDoc<'a> { + /// Parses a JSON object. + /// + /// Accepts and rejects exactly the same inputs as + /// `serde_json::from_slice::>`, with one documented + /// exception: nested objects whose first key is `serde_json`'s private raw value token are + /// rejected (`serde_json` would parse the associated string as embedded JSON). + /// Error messages may differ for invalid UTF-8. + pub fn parse(json_bytes: &'a [u8]) -> Result { + // Validating UTF-8 once for the whole document is cheaper than validating each string. + let json_str = std::str::from_utf8(json_bytes).map_err(de::Error::custom)?; + let mut deserializer = serde_json::Deserializer::from_str(json_str); + let root = deserializer.deserialize_map(ObjectVisitor { + reject_raw_value_token: false, + })?; + deserializer.end()?; + Ok(BorrowedJsonDoc { root }) + } + + /// Returns the entries of the root object. They have unique keys and are ordered like in + /// `serde_json::Map`. + pub fn root(&self) -> &BorrowedObject<'a> { + &self.root + } + + /// Returns the value associated with `key` in the root object. + pub fn get(&self, key: &str) -> Option<&BorrowedValue<'a>> { + get_in_object(&self.root, key) + } +} + +impl<'de> Deserialize<'de> for BorrowedValue<'de> { + fn deserialize(deserializer: D) -> Result + where D: Deserializer<'de> { + deserializer.deserialize_any(ValueVisitor) + } +} + +struct ValueVisitor; + +impl<'de> Visitor<'de> for ValueVisitor { + type Value = BorrowedValue<'de>; + + fn expecting(&self, formatter: &mut fmt::Formatter) -> fmt::Result { + formatter.write_str("any valid JSON value") + } + + fn visit_bool(self, value: bool) -> Result { + Ok(BorrowedValue::Bool(value)) + } + + fn visit_i64(self, value: i64) -> Result { + Ok(BorrowedValue::Number(value.into())) + } + + fn visit_u64(self, value: u64) -> Result { + Ok(BorrowedValue::Number(value.into())) + } + + fn visit_f64(self, value: f64) -> Result { + // Same as `serde_json::Value`: non-finite floats become `null`. + match Number::from_f64(value) { + Some(number) => Ok(BorrowedValue::Number(number)), + None => Ok(BorrowedValue::Null), + } + } + + fn visit_borrowed_str(self, value: &'de str) -> Result { + Ok(BorrowedValue::Str(Cow::Borrowed(value))) + } + + fn visit_str(self, value: &str) -> Result { + Ok(BorrowedValue::Str(Cow::Owned(value.to_string()))) + } + + fn visit_string(self, value: String) -> Result { + Ok(BorrowedValue::Str(Cow::Owned(value))) + } + + fn visit_none(self) -> Result { + Ok(BorrowedValue::Null) + } + + fn visit_some(self, deserializer: D) -> Result + where D: Deserializer<'de> { + Deserialize::deserialize(deserializer) + } + + fn visit_unit(self) -> Result { + Ok(BorrowedValue::Null) + } + + fn visit_seq(self, mut seq: V) -> Result + where V: SeqAccess<'de> { + let mut elements = Vec::with_capacity(seq.size_hint().unwrap_or(0)); + while let Some(element) = seq.next_element()? { + elements.push(element); + } + Ok(BorrowedValue::Array(elements)) + } + + fn visit_map(self, map: V) -> Result + where V: MapAccess<'de> { + let object_visitor = ObjectVisitor { + reject_raw_value_token: true, + }; + object_visitor.visit_map(map).map(BorrowedValue::Object) + } +} + +struct ObjectVisitor { + /// `serde_json::Value` (but not `serde_json::Map`) gives a special meaning to objects whose + /// first key is the raw value token. We do not support it and reject such documents. + reject_raw_value_token: bool, +} + +impl<'de> Visitor<'de> for ObjectVisitor { + type Value = BorrowedObject<'de>; + + fn expecting(&self, formatter: &mut fmt::Formatter) -> fmt::Result { + // Same as the visitor of `serde_json::Map`. + formatter.write_str("a map") + } + + fn visit_unit(self) -> Result { + // Same as the visitor of `serde_json::Map`. + Ok(Vec::new()) + } + + fn visit_map(self, mut map: V) -> Result + where V: MapAccess<'de> { + let mut entries: BorrowedObject<'de> = Vec::with_capacity(map.size_hint().unwrap_or(0)); + while let Some(key) = map.next_key_seed(KeySeed)? { + if self.reject_raw_value_token + && entries.is_empty() + && key.as_ref() == SERDE_JSON_RAW_VALUE_TOKEN + { + return Err(de::Error::custom(format!( + "unsupported object key `{SERDE_JSON_RAW_VALUE_TOKEN}`" + ))); + } + let value: BorrowedValue<'de> = map.next_value()?; + entries.push((key, value)); + } + if *SERDE_JSON_PRESERVES_ORDER { + dedup_in_place_keeping_last(&mut entries); + } else { + sort_and_dedup_keeping_last(&mut entries); + } + Ok(entries) + } +} + +/// Removes duplicate keys without reordering, keeping the value of the last occurrence at the +/// position of the first one, which is the behavior of `IndexMap::insert`. +fn dedup_in_place_keeping_last(entries: &mut BorrowedObject) { + let mut positions: Vec = (0..entries.len()).collect(); + // Stable sort: positions of duplicate keys stay in increasing order. + positions.sort_by(|left, right| entries[*left].0.cmp(&entries[*right].0)); + let mut is_removed = vec![false; entries.len()]; + // (first position, last position) of each duplicated key. + let mut swaps: Vec<(usize, usize)> = Vec::new(); + for group in positions.chunk_by(|left, right| entries[*left].0 == entries[*right].0) { + let [first_position, .., last_position] = group else { + continue; + }; + swaps.push((*first_position, *last_position)); + for duplicate_position in &group[1..] { + is_removed[*duplicate_position] = true; + } + } + if swaps.is_empty() { + return; + } + // The first occurrence of each key gets the last value. The last occurrence gets the first + // value but is removed. + for (first_position, last_position) in swaps { + entries.swap(first_position, last_position); + } + let mut position = 0; + entries.retain(|_| { + let keep = !is_removed[position]; + position += 1; + keep + }); +} + +/// Sorts the entries by key and removes duplicate keys, keeping the value of the last occurrence, +/// which is the behavior of `BTreeMap::insert`. +fn sort_and_dedup_keeping_last(entries: &mut BorrowedObject) { + if entries.is_sorted_by(|left, right| left.0 < right.0) { + // Strictly sorted: no duplicates. + return; + } + // The sort must be stable for the last occurrence of a key to stay last among its duplicates. + entries.sort_by(|left, right| left.0.cmp(&right.0)); + // `dedup_by` calls the closure with `(later, retained)` and removes `later` when it returns + // true. Swapping first moves the later value into the retained entry. + entries.dedup_by(|later, retained| { + if later.0 != retained.0 { + return false; + } + std::mem::swap(&mut later.1, &mut retained.1); + true + }); +} + +struct KeySeed; + +impl<'de> DeserializeSeed<'de> for KeySeed { + type Value = Cow<'de, str>; + + fn deserialize(self, deserializer: D) -> Result + where D: Deserializer<'de> { + deserializer.deserialize_str(KeyVisitor) + } +} + +struct KeyVisitor; + +impl<'de> Visitor<'de> for KeyVisitor { + type Value = Cow<'de, str>; + + fn expecting(&self, formatter: &mut fmt::Formatter) -> fmt::Result { + formatter.write_str("a string key") + } + + fn visit_borrowed_str(self, value: &'de str) -> Result { + Ok(Cow::Borrowed(value)) + } + + fn visit_str(self, value: &str) -> Result { + Ok(Cow::Owned(value.to_string())) + } + + fn visit_string(self, value: String) -> Result { + Ok(Cow::Owned(value)) + } +} + +#[cfg(test)] +mod tests { + use super::*; + + fn assert_same_as_serde_json(json: &str) { + let expected: Result, _> = serde_json::from_str(json); + let parsed = BorrowedJsonDoc::parse(json.as_bytes()); + match (expected, parsed) { + (Ok(expected_obj), Ok(parsed_doc)) => { + let parsed_value = BorrowedValue::Object(parsed_doc.root).to_serde_json(); + assert_eq!( + parsed_value, + JsonValue::Object(expected_obj), + "input: {json}" + ); + } + (Err(expected_error), Err(parsed_error)) => { + assert_eq!( + expected_error.to_string(), + parsed_error.to_string(), + "input: {json}" + ); + } + (expected, parsed) => { + panic!("input: {json}, serde_json: {expected:?}, borrowed: {parsed:?}") + } + } + } + + #[test] + fn test_borrowed_json_same_as_serde_json() { + let inputs = [ + r#"{}"#, + r#"{"b": 1, "a": 2, "c": {"z": null, "y": [1, -2, 3.5, true, "x"]}}"#, + r#"{"a": 1, "a": 2, "b": 3, "a": 4}"#, + r#"{"nested": {"k": 1, "k": {"x": 1}, "j": 0}}"#, + r#"{"escaped\"key": "line\nbreak \u00e9 \ud83d\ude00", "plain": "abc"}"#, + r#"{"big": 18446744073709551615, "bigger": 18446744073709551616, "neg": -9223372036854775808}"#, + r#"{"float": 1e300, "neg_zero": -0.0, "exp": 2E-5, "int_float": 3.0}"#, + r#"{"é": 1, "e": 2, "Z": 3, "": 4}"#, + r#"{"a": [[], [{}], [[null]]]}"#, + " {\"a\" : 1 } \n", + r#"{"a": 1"#, + r#"{"a": 1} trailing"#, + r#"[1, 2]"#, + r#""string""#, + r#"null"#, + r#"42"#, + r#"{"a": 1e400}"#, + r#"{"a": "\ud800"}"#, + r#"{"a": nul}"#, + r#"{1: 2}"#, + r#"{"$serde_json::private::RawValue": 1}"#, + ]; + for input in inputs { + assert_same_as_serde_json(input); + } + } + + #[test] + fn test_borrowed_json_borrows_unescaped_strings() { + let json = r#"{"key": "value", "escaped": "va\"lue"}"#; + let doc = BorrowedJsonDoc::parse(json.as_bytes()).unwrap(); + let root = doc.root(); + assert!(matches!( + get_in_object(root, "key"), + Some(BorrowedValue::Str(Cow::Borrowed("value"))) + )); + assert!(matches!( + get_in_object(root, "escaped"), + Some(BorrowedValue::Str(Cow::Owned(_))) + )); + assert!(get_in_object(root, "missing").is_none()); + } + + #[test] + fn test_borrowed_json_invalid_utf8() { + let invalid_utf8: &[u8] = b"{\"a\": \"\xff\"}"; + assert!( + serde_json::from_slice::>(invalid_utf8).is_err() + ); + assert!(BorrowedJsonDoc::parse(invalid_utf8).is_err()); + } + + #[test] + fn test_borrowed_json_rejects_nested_raw_value_token() { + let json = r#"{"a": {"$serde_json::private::RawValue": "1"}}"#; + let error = BorrowedJsonDoc::parse(json.as_bytes()).unwrap_err(); + assert!(error.to_string().contains("unsupported object key")); + } + + #[test] + fn test_borrowed_json_display_same_as_serde_json() { + let json = r#"{"a": [1, "two", {"b": null}], "c": 1.5}"#; + let doc = BorrowedJsonDoc::parse(json.as_bytes()).unwrap(); + let expected: JsonValue = serde_json::from_str(json).unwrap(); + let value = BorrowedValue::Object(doc.root().clone()); + assert_eq!(value.to_string(), expected.to_string()); + } +} diff --git a/quickwit/quickwit-doc-mapper/src/doc_mapper/borrowed_value_view.rs b/quickwit/quickwit-doc-mapper/src/doc_mapper/borrowed_value_view.rs new file mode 100644 index 00000000000..04866428a4a --- /dev/null +++ b/quickwit/quickwit-doc-mapper/src/doc_mapper/borrowed_value_view.rs @@ -0,0 +1,267 @@ +// 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. + +//! Exposes borrowed JSON values to tantivy through its [`Value`] trait, so that objects are +//! written into the tantivy document without building intermediate `OwnedValue`s. + +use std::borrow::Cow; +use std::slice; + +use serde_json::Number; +use tantivy::schema::document::{ReferenceValue, ReferenceValueLeaf}; +use tantivy::schema::{Field, Value}; +use tantivy::time::format_description::well_known::Rfc3339; +use tantivy::time::{OffsetDateTime, UtcOffset}; +use tantivy::{DateTime, TantivyDocument as Document}; + +use super::borrowed_json::{BorrowedObject, BorrowedValue, serde_json_preserves_order}; +use super::mapping_tree::NumVal; + +/// Unmapped fields collected while walking the mapping tree in dynamic mode, mirroring the +/// `dynamic_json_obj` map of the owned path. Entries are pushed while iterating over sorted +/// objects, so they are sorted by key too. +pub(crate) type DynamicObject<'b, 'a> = Vec<(&'b str, DynamicEntry<'b, 'a>)>; + +#[derive(Debug)] +pub(crate) enum DynamicEntry<'b, 'a> { + /// An unmapped field, with its whole JSON value. + Value(&'b BorrowedValue<'a>), + /// The unmapped fields of an object mapped with an `object` field mapping. Never empty. + Object(DynamicObject<'b, 'a>), +} + +/// Adds a JSON object to a document, like `document.add_object(field, object.into())` does in +/// the owned path (string values that look like RFC 3339 dates become dates, and the top-level +/// keys are sorted because `add_object` takes a `BTreeMap`). +pub(crate) fn add_borrowed_object( + document: &mut Document, + field: Field, + json_obj: &BorrowedObject, +) { + let value_view = TantivyValueView::JsonObject { + entries: json_obj, + sort_keys: true, + }; + document.add_field_value(field, value_view); +} + +/// Adds a JSON value to a document. See [`add_borrowed_object`]. +pub(crate) fn add_borrowed_value(document: &mut Document, field: Field, json_val: &BorrowedValue) { + document.add_field_value(field, TantivyValueView::Json(json_val)); +} + +/// Adds the unmapped fields to the dynamic field. See [`add_borrowed_object`]. +pub(crate) fn add_dynamic_object( + document: &mut Document, + field: Field, + dynamic_obj: &DynamicObject, +) { + let value_view = TantivyValueView::DynamicObject { + entries: dynamic_obj, + sort_keys: true, + }; + document.add_field_value(field, value_view); +} + +/// Adds every primitive value of the unmapped fields to the concatenate fields, mirroring +/// `JsonValueIterator` followed by `map_primitive_json_to_concatenate_value` in the owned path. +pub(crate) fn add_dynamic_concatenate_values( + document: &mut Document, + concatenate_fields: &[Field], + dynamic_obj: &DynamicObject, +) { + for (_key, dynamic_entry) in dynamic_obj { + match dynamic_entry { + DynamicEntry::Value(json_val) => { + add_concatenate_leaves(document, concatenate_fields, json_val); + } + DynamicEntry::Object(child_dynamic_obj) => { + add_dynamic_concatenate_values(document, concatenate_fields, child_dynamic_obj); + } + } + } +} + +/// Adds every primitive value of `json_val`, in depth-first order, to the concatenate fields. +/// Nulls are skipped and strings are never interpreted as dates. +pub(crate) fn add_concatenate_leaves( + document: &mut Document, + concatenate_fields: &[Field], + json_val: &BorrowedValue, +) { + let leaf: ReferenceValueLeaf = match json_val { + BorrowedValue::Null => return, + BorrowedValue::Array(elements) => { + for element in elements { + add_concatenate_leaves(document, concatenate_fields, element); + } + return; + } + BorrowedValue::Object(entries) => { + for (_key, child_json_val) in entries { + add_concatenate_leaves(document, concatenate_fields, child_json_val); + } + return; + } + BorrowedValue::Str(text) => ReferenceValueLeaf::Str(text), + BorrowedValue::Bool(bool_val) => (*bool_val).into(), + BorrowedValue::Number(number) => { + if let Some(i64_val) = i64::from_json_number(number) { + i64_val.into() + } else if let Some(u64_val) = u64::from_json_number(number) { + u64_val.into() + } else if let Some(f64_val) = f64::from_json_number(number) { + f64_val.into() + } else { + return; + } + } + }; + for field in concatenate_fields { + document.add_leaf_field_value(*field, leaf.clone()); + } +} + +/// Mirrors `From for OwnedValue` for numbers. +fn number_to_leaf(number: &Number) -> ReferenceValueLeaf<'static> { + if let Some(i64_val) = number.as_i64() { + ReferenceValueLeaf::I64(i64_val) + } else if let Some(u64_val) = number.as_u64() { + ReferenceValueLeaf::U64(u64_val) + } else if let Some(f64_val) = number.as_f64() { + ReferenceValueLeaf::F64(f64_val) + } else { + // Without the `arbitrary_precision` feature, a number is always an i64, a u64 or a f64. + panic!("unsupported serde_json number `{number}`") + } +} + +/// Mirrors `From for OwnedValue` for strings: strings that start with a digit +/// and parse as RFC 3339 are converted into dates. +fn str_to_leaf(text: &str) -> ReferenceValueLeaf<'_> { + // Same pre-check as tantivy's `can_be_rfc3339_date_time`. + let Some(first_byte) = text.as_bytes().first() else { + return ReferenceValueLeaf::Str(text); + }; + if !first_byte.is_ascii_digit() { + return ReferenceValueLeaf::Str(text); + } + match OffsetDateTime::parse(text, &Rfc3339) { + Ok(date_time) => { + ReferenceValueLeaf::Date(DateTime::from_utc(date_time.to_offset(UtcOffset::UTC))) + } + Err(_) => ReferenceValueLeaf::Str(text), + } +} + +#[derive(Debug, Clone, Copy)] +enum TantivyValueView<'b, 'a> { + Json(&'b BorrowedValue<'a>), + /// `sort_keys` mirrors the `BTreeMap` used by `add_object` for top-level objects. Objects are + /// already sorted unless the `preserve_order` feature of `serde_json` is enabled. + JsonObject { + entries: &'b BorrowedObject<'a>, + sort_keys: bool, + }, + DynamicObject { + entries: &'b DynamicObject<'b, 'a>, + sort_keys: bool, + }, +} + +impl<'b, 'a: 'b> TantivyValueView<'b, 'a> { + fn object_iter(&self) -> Option> { + let (object_iter, sort_keys) = match *self { + TantivyValueView::Json(BorrowedValue::Object(entries)) => { + (ObjectIter::Json(entries.iter()), false) + } + TantivyValueView::Json(_) => return None, + TantivyValueView::JsonObject { entries, sort_keys } => { + (ObjectIter::Json(entries.iter()), sort_keys) + } + TantivyValueView::DynamicObject { entries, sort_keys } => { + (ObjectIter::Dynamic(entries.iter()), sort_keys) + } + }; + if !sort_keys || !serde_json_preserves_order() { + return Some(object_iter); + } + let mut sorted_entries: Vec<(&'b str, TantivyValueView<'b, 'a>)> = object_iter.collect(); + sorted_entries.sort_by(|left, right| left.0.cmp(right.0)); + Some(ObjectIter::Sorted(sorted_entries.into_iter())) + } +} + +impl<'b, 'a: 'b> Value<'b> for TantivyValueView<'b, 'a> { + type ArrayIter = ArrayIter<'b, 'a>; + type ObjectIter = ObjectIter<'b, 'a>; + + fn as_value(&self) -> ReferenceValue<'b, Self> { + if let Some(object_iter) = self.object_iter() { + return ReferenceValue::Object(object_iter); + } + let TantivyValueView::Json(json_val) = *self else { + unreachable!("objects are handled above") + }; + match json_val { + BorrowedValue::Null => ReferenceValue::Leaf(ReferenceValueLeaf::Null), + BorrowedValue::Bool(bool_val) => ReferenceValue::Leaf((*bool_val).into()), + BorrowedValue::Number(number) => ReferenceValue::Leaf(number_to_leaf(number)), + BorrowedValue::Str(text) => ReferenceValue::Leaf(str_to_leaf(text)), + BorrowedValue::Array(elements) => ReferenceValue::Array(ArrayIter(elements.iter())), + BorrowedValue::Object(_) => unreachable!("objects are handled above"), + } + } +} + +struct ArrayIter<'b, 'a>(slice::Iter<'b, BorrowedValue<'a>>); + +impl<'b, 'a: 'b> Iterator for ArrayIter<'b, 'a> { + type Item = TantivyValueView<'b, 'a>; + + fn next(&mut self) -> Option { + self.0.next().map(TantivyValueView::Json) + } +} + +enum ObjectIter<'b, 'a> { + Json(slice::Iter<'b, (Cow<'a, str>, BorrowedValue<'a>)>), + Dynamic(slice::Iter<'b, (&'b str, DynamicEntry<'b, 'a>)>), + Sorted(std::vec::IntoIter<(&'b str, TantivyValueView<'b, 'a>)>), +} + +impl<'b, 'a: 'b> Iterator for ObjectIter<'b, 'a> { + type Item = (&'b str, TantivyValueView<'b, 'a>); + + fn next(&mut self) -> Option { + match self { + ObjectIter::Json(entries) => { + let (key, json_val) = entries.next()?; + Some((key.as_ref(), TantivyValueView::Json(json_val))) + } + ObjectIter::Dynamic(entries) => { + let (key, dynamic_entry) = entries.next()?; + let value_view = match dynamic_entry { + DynamicEntry::Value(json_val) => TantivyValueView::Json(json_val), + DynamicEntry::Object(child_obj) => TantivyValueView::DynamicObject { + entries: child_obj, + sort_keys: false, + }, + }; + Some((key, value_view)) + } + ObjectIter::Sorted(entries) => entries.next(), + } + } +} diff --git a/quickwit/quickwit-doc-mapper/src/doc_mapper/date_time_type.rs b/quickwit/quickwit-doc-mapper/src/doc_mapper/date_time_type.rs index 772819d6dc2..55d5ea42e1d 100644 --- a/quickwit/quickwit-doc-mapper/src/doc_mapper/date_time_type.rs +++ b/quickwit/quickwit-doc-mapper/src/doc_mapper/date_time_type.rs @@ -19,6 +19,8 @@ use serde::{Deserialize, Deserializer, Serialize}; use serde_json::Value as JsonValue; use tantivy::schema::{DateTimePrecision, OwnedValue as TantivyValue}; +use super::borrowed_json::BorrowedValue; + /// A struct holding DateTime field options. #[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)] #[serde(deny_unknown_fields)] @@ -125,6 +127,35 @@ impl QuickwitDateTimeOptions { Ok(TantivyValue::Date(date_time)) } + /// Same as [`Self::parse_json`] for a borrowed JSON value. + pub(crate) fn parse_borrowed_json( + &self, + json_value: &BorrowedValue, + ) -> Result { + match json_value { + BorrowedValue::Number(timestamp) => { + // `.as_f64()` actually converts floats to integers, so we must check for integers + // first. + if let Some(timestamp_i64) = timestamp.as_i64() { + quickwit_datetime::parse_timestamp_int(timestamp_i64, &self.input_formats.0) + } else if let Some(timestamp_f64) = timestamp.as_f64() { + quickwit_datetime::parse_timestamp_float(timestamp_f64, &self.input_formats.0) + } else { + Err(format!( + "failed to parse datetime `{timestamp:?}`: value is larger than i64::MAX", + )) + } + } + BorrowedValue::Str(date_time_str) => { + quickwit_datetime::parse_date_time_str(date_time_str, &self.input_formats.0) + } + _ => Err(format!( + "failed to parse datetime: expected a float, integer, or string, got \ + `{json_value}`" + )), + } + } + pub(crate) fn reparse_tantivy_value( &self, tantivy_value: &TantivyValue, diff --git a/quickwit/quickwit-doc-mapper/src/doc_mapper/doc_mapper_impl.rs b/quickwit/quickwit-doc-mapper/src/doc_mapper/doc_mapper_impl.rs index 087fe2839ab..1118ca38162 100644 --- a/quickwit/quickwit-doc-mapper/src/doc_mapper/doc_mapper_impl.rs +++ b/quickwit/quickwit-doc-mapper/src/doc_mapper/doc_mapper_impl.rs @@ -30,6 +30,10 @@ use tantivy::query::Query; use tantivy::schema::{Field, FieldType, INDEXED, OwnedValue as TantivyValue, STORED, Schema}; use super::DocMapperBuilder; +use super::borrowed_json::BorrowedJsonDoc; +use super::borrowed_value_view::{ + DynamicObject, add_borrowed_object, add_dynamic_concatenate_values, add_dynamic_object, +}; use super::field_mapping_entry::RAW_TOKENIZER_NAME; use super::field_presence::populate_field_presence; use super::tantivy_val_to_json::tantivy_value_to_json; @@ -547,18 +551,68 @@ impl DocMapper { ); } + self.add_document_size_and_field_presence(&mut document, document_len); + Ok((partition, document)) + } + + /// Same as [`Self::doc_from_json_obj`] for a borrowed JSON document, without allocating + /// intermediate JSON or tantivy values. + /// + /// Produces exactly the same partition, document (same field values in the same order) and + /// errors as `doc_from_json_obj` applied to the equivalent `serde_json::Map`. + pub fn doc_from_borrowed_json( + &self, + json_doc: &BorrowedJsonDoc, + document_len: u64, + ) -> Result<(Partition, Document), DocParsingError> { + let partition: Partition = self.partition_key.eval_hash(json_doc); + + let mut dynamic_obj = DynamicObject::new(); + let mut field_path = Vec::new(); + let mut document = Document::default(); + + if let Some(source_field) = self.source_field { + add_borrowed_object(&mut document, source_field, json_doc.root()); + } + + let mode = self.mode.mode_type(); + self.field_mappings.doc_from_borrowed_json( + json_doc.root(), + mode, + &mut document, + &mut field_path, + &mut dynamic_obj, + )?; + + if let Some(dynamic_field) = self.dynamic_field + && !dynamic_obj.is_empty() + { + if !self.concatenate_dynamic_fields.is_empty() { + add_dynamic_concatenate_values( + &mut document, + &self.concatenate_dynamic_fields, + &dynamic_obj, + ); + } + add_dynamic_object(&mut document, dynamic_field, &dynamic_obj); + } + + self.add_document_size_and_field_presence(&mut document, document_len); + Ok((partition, document)) + } + + fn add_document_size_and_field_presence(&self, document: &mut Document, document_len: u64) { if let Some(document_size_field) = self.document_size_field { document.add_u64(document_size_field, document_len); } if self.index_field_presence { let field_presence_hashes: FnvHashSet = - populate_field_presence(&document, &self.schema, true); + populate_field_presence(document, &self.schema, true); for field_presence_hash in field_presence_hashes { document.add_field_value(FIELD_PRESENCE_FIELD, &field_presence_hash); } } - Ok((partition, document)) } /// Converts a tantivy named Document to the json format. diff --git a/quickwit/quickwit-doc-mapper/src/doc_mapper/mapping_tree.rs b/quickwit/quickwit-doc-mapper/src/doc_mapper/mapping_tree.rs index 30b16fea4c8..fa023f77350 100644 --- a/quickwit/quickwit-doc-mapper/src/doc_mapper/mapping_tree.rs +++ b/quickwit/quickwit-doc-mapper/src/doc_mapper/mapping_tree.rs @@ -37,6 +37,8 @@ use crate::doc_mapper::field_mapping_entry::{ use crate::doc_mapper::{FieldMappingType, QuickwitJsonOptions}; use crate::{Cardinality, DocParsingError, FieldMappingEntry, ModeType}; +mod borrowed_doc; + #[derive(Clone, Debug)] pub enum LeafType { Bool(QuickwitBoolOptions), diff --git a/quickwit/quickwit-doc-mapper/src/doc_mapper/mapping_tree/borrowed_doc.rs b/quickwit/quickwit-doc-mapper/src/doc_mapper/mapping_tree/borrowed_doc.rs new file mode 100644 index 00000000000..939fb1b1398 --- /dev/null +++ b/quickwit/quickwit-doc-mapper/src/doc_mapper/mapping_tree/borrowed_doc.rs @@ -0,0 +1,292 @@ +// 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. + +//! Conversion of a borrowed JSON object into a tantivy document, following the mapping tree. +//! +//! Every function here mirrors a function of the owned conversion path (`doc_from_json` and +//! friends in the parent module) and must produce the same field values, in the same order, and +//! the same errors. The equivalence is checked by the differential tests in +//! `doc_mapper/borrowed_doc_tests.rs`. Any change to one path must be reflected in the other. + +use std::any::type_name; +use std::net::IpAddr; +use std::str::FromStr; + +use tantivy::TantivyDocument as Document; +use tantivy::schema::document::ReferenceValueLeaf; +use tantivy::schema::{Field, IntoIpv6Addr}; + +use super::{LeafType, MappingLeaf, MappingNode, MappingTree, NumVal}; +use crate::doc_mapper::borrowed_json::{BorrowedObject, BorrowedValue}; +use crate::doc_mapper::borrowed_value_view::{ + DynamicEntry, DynamicObject, add_borrowed_value, add_concatenate_leaves, +}; +use crate::{Cardinality, DocParsingError, ModeType}; + +impl MappingNode { + /// Mirrors [`MappingNode::doc_from_json`]. + pub(crate) fn doc_from_borrowed_json<'b, 'a>( + &self, + json_obj: &'b BorrowedObject<'a>, + mode: ModeType, + document: &mut Document, + path: &mut Vec<&'b str>, + dynamic_obj: &mut DynamicObject<'b, 'a>, + ) -> Result<(), DocParsingError> { + for (field_name, json_val) in json_obj { + let field_name: &'b str = field_name; + let Some(child_tree) = self.branches.get(field_name) else { + match mode { + ModeType::Lenient => { + // In lenient mode we simply ignore these unmapped fields. + } + ModeType::Dynamic => { + dynamic_obj.push((field_name, DynamicEntry::Value(json_val))); + } + ModeType::Strict => { + path.push(field_name); + return Err(DocParsingError::NoSuchFieldInSchema(path.join("."))); + } + } + continue; + }; + path.push(field_name); + match child_tree { + MappingTree::Leaf(mapping_leaf) => { + mapping_leaf.doc_from_borrowed_json(json_val, document, path)?; + } + MappingTree::Node(mapping_node) => { + let BorrowedValue::Object(child_obj) = json_val else { + return Err(DocParsingError::ValueError( + path.join("."), + format!("expected an JSON object, got {json_val}"), + )); + }; + let mut child_dynamic_obj = DynamicObject::new(); + mapping_node.doc_from_borrowed_json( + child_obj, + mode, + document, + path, + &mut child_dynamic_obj, + )?; + if !child_dynamic_obj.is_empty() { + dynamic_obj.push((field_name, DynamicEntry::Object(child_dynamic_obj))); + } + } + } + path.pop(); + } + Ok(()) + } +} + +impl MappingLeaf { + /// Mirrors [`MappingLeaf::doc_from_json`]. + fn doc_from_borrowed_json( + &self, + json_val: &BorrowedValue, + document: &mut Document, + path: &[&str], + ) -> Result<(), DocParsingError> { + if json_val.is_null() { + // We just ignore `null`. + return Ok(()); + } + let BorrowedValue::Array(elements) = json_val else { + return self.add_borrowed_value(json_val, document, path); + }; + if self.cardinality == Cardinality::SingleValued { + return Err(DocParsingError::MultiValuesNotSupported(path.join("."))); + } + for element in elements { + if element.is_null() { + // We just ignore `null`. + continue; + } + self.add_borrowed_value(element, document, path)?; + } + Ok(()) + } + + fn add_borrowed_value( + &self, + json_val: &BorrowedValue, + document: &mut Document, + path: &[&str], + ) -> Result<(), DocParsingError> { + let to_value_error = |err_msg| DocParsingError::ValueError(path.join("."), err_msg); + if !self.concatenate.is_empty() { + self.typ + .add_borrowed_concatenate_values(json_val, &self.concatenate, document) + .map_err(to_value_error)?; + } + self.typ + .add_borrowed_value(self.field, json_val, document) + .map_err(to_value_error) + } +} + +impl LeafType { + /// Mirrors [`LeafType::value_from_json`] followed by adding the value to the document. + fn add_borrowed_value( + &self, + field: Field, + json_val: &BorrowedValue, + document: &mut Document, + ) -> Result<(), String> { + match self { + LeafType::Text(_) => { + let BorrowedValue::Str(text) = json_val else { + return Err(format!("expected string, got `{json_val}`")); + }; + document.add_text(field, text); + } + LeafType::I64(numeric_options) => { + let value = num_from_borrowed_json::(json_val, numeric_options.coerce)?; + document.add_i64(field, value); + } + LeafType::U64(numeric_options) => { + let value = num_from_borrowed_json::(json_val, numeric_options.coerce)?; + document.add_u64(field, value); + } + LeafType::F64(numeric_options) => { + let value = num_from_borrowed_json::(json_val, numeric_options.coerce)?; + document.add_f64(field, value); + } + LeafType::Bool(_) => { + let BorrowedValue::Bool(value) = json_val else { + return Err(format!("expected boolean, got `{json_val}`")); + }; + document.add_bool(field, *value); + } + LeafType::IpAddr(_) => { + let BorrowedValue::Str(ip_address) = json_val else { + return Err(format!("expected string, got `{json_val}`")); + }; + let ipv6_value = IpAddr::from_str(ip_address) + .map_err(|err| format!("failed to parse IP address `{ip_address}`: {err}"))? + .into_ipv6_addr(); + document.add_ip_addr(field, ipv6_value); + } + LeafType::DateTime(date_time_options) => { + let date_time = date_time_options.parse_borrowed_json(json_val)?; + document.add_date(field, date_time); + } + LeafType::Bytes(binary_options) => { + let BorrowedValue::Str(byte_str) = json_val else { + return Err(format!( + "expected {} string, got `{json_val}`", + binary_options.input_format.as_str() + )); + }; + let payload = binary_options.input_format.parse_str(byte_str)?; + document.add_bytes(field, &payload); + } + LeafType::Json(_) => { + let BorrowedValue::Object(_) = json_val else { + return Err(format!("expected object, got `{json_val}`")); + }; + add_borrowed_value(document, field, json_val); + } + } + Ok(()) + } + + /// Mirrors [`LeafType::concatenate_values_from_json`] followed by adding the values to each + /// of the `concatenate_fields`. + /// + /// Errors are only returned before any value is added. + fn add_borrowed_concatenate_values( + &self, + json_val: &BorrowedValue, + concatenate_fields: &[Field], + document: &mut Document, + ) -> Result<(), String> { + let leaf: ReferenceValueLeaf = match self { + LeafType::Text(_) => { + let BorrowedValue::Str(text) = json_val else { + return Err(format!("expected string, got `{json_val}`")); + }; + ReferenceValueLeaf::Str(text) + } + LeafType::I64(numeric_options) => { + num_from_borrowed_json::(json_val, numeric_options.coerce)?.into() + } + LeafType::U64(numeric_options) => { + num_from_borrowed_json::(json_val, numeric_options.coerce)?.into() + } + LeafType::F64(numeric_options) => { + num_from_borrowed_json::(json_val, numeric_options.coerce)?.into() + } + LeafType::Bool(_) => { + let BorrowedValue::Bool(value) = json_val else { + return Err(format!("expected boolean, got `{json_val}`")); + }; + (*value).into() + } + LeafType::IpAddr(_) => return Err("unsupported concat type: IpAddr".to_string()), + LeafType::DateTime(_) => return Err("unsupported concat type: DateTime".to_string()), + LeafType::Bytes(_) => return Err("unsupported concat type: Bytes".to_string()), + LeafType::Json(_) => { + let BorrowedValue::Object(json_obj) = json_val else { + return Err(format!("expected object, got `{json_val}`")); + }; + for (_key, child_json_val) in json_obj { + add_concatenate_leaves(document, concatenate_fields, child_json_val); + } + return Ok(()); + } + }; + for field in concatenate_fields { + document.add_leaf_field_value(*field, leaf.clone()); + } + Ok(()) + } +} + +/// Mirrors [`NumVal::from_json_to_self`]. +fn num_from_borrowed_json(json_val: &BorrowedValue, coerce: bool) -> Result { + match json_val { + BorrowedValue::Number(num_val) => T::from_json_number(num_val).ok_or_else(|| { + format!( + "expected {}, got inconvertible JSON number `{}`", + type_name::(), + num_val + ) + }), + BorrowedValue::Str(str_val) => { + if !coerce { + return Err(format!( + "expected JSON number, got string `\"{str_val}\"`. enable coercion to {} with \ + the `coerce` parameter in the field mapping", + type_name::() + )); + } + str_val.parse::().map_err(|_| { + format!( + "failed to coerce JSON string `\"{str_val}\"` to {}", + type_name::() + ) + }) + } + _ => { + if coerce { + Err(format!("expected JSON number or string, got `{json_val}`")) + } else { + Err(format!("expected JSON number, got `{json_val}`")) + } + } + } +} diff --git a/quickwit/quickwit-doc-mapper/src/doc_mapper/mod.rs b/quickwit/quickwit-doc-mapper/src/doc_mapper/mod.rs index e9f199fc417..2bbcc853bb5 100644 --- a/quickwit/quickwit-doc-mapper/src/doc_mapper/mod.rs +++ b/quickwit/quickwit-doc-mapper/src/doc_mapper/mod.rs @@ -12,6 +12,10 @@ // See the License for the specific language governing permissions and // limitations under the License. +#[cfg(test)] +mod borrowed_doc_tests; +pub(crate) mod borrowed_json; +mod borrowed_value_view; mod date_time_type; mod doc_mapper_builder; mod doc_mapper_impl; @@ -19,6 +23,8 @@ mod field_mapping_entry; mod field_mapping_type; mod field_presence; mod mapping_tree; +#[cfg(any(test, feature = "testsuite"))] +mod random_json; mod tantivy_val_to_json; mod tokenizer_entry; @@ -26,6 +32,7 @@ use std::collections::{HashMap, HashSet}; use std::fmt::Debug; use std::ops::Bound; +pub use borrowed_json::{BorrowedJsonDoc, BorrowedObject, BorrowedValue}; pub use doc_mapper_builder::DocMapperBuilder; pub use doc_mapper_impl::DocMapper; pub use field_mapping_entry::{ @@ -38,6 +45,8 @@ pub(crate) use field_mapping_entry::{ #[cfg(test)] pub(crate) use field_mapping_entry::{QuickwitNumericOptions, QuickwitTextOptions}; pub use field_mapping_type::FieldMappingType; +#[cfg(any(test, feature = "testsuite"))] +pub use random_json::RandomJsonDocs; use serde_json::Value as JsonValue; use tantivy::Term; use tantivy::schema::{Field, FieldType}; diff --git a/quickwit/quickwit-doc-mapper/src/doc_mapper/random_json.rs b/quickwit/quickwit-doc-mapper/src/doc_mapper/random_json.rs new file mode 100644 index 00000000000..a37e35473a0 --- /dev/null +++ b/quickwit/quickwit-doc-mapper/src/doc_mapper/random_json.rs @@ -0,0 +1,175 @@ +// 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. + +//! Generator of random JSON objects, used by differential tests comparing code paths that consume +//! JSON documents (for instance owned and borrowed JSON trees). +//! +//! Keys are drawn from a small vocabulary so objects regularly contain duplicate keys, and values +//! cover escapes, dates, IPs, numbers at the integer type boundaries, nested arrays and objects. + +/// Deterministic generator of random JSON object documents. +pub struct RandomJsonDocs { + rng: Rng, +} + +impl RandomJsonDocs { + /// The same seed always produces the same sequence of documents. `seed` must not be zero. + pub fn new(seed: u64) -> Self { + assert_ne!(seed, 0, "xorshift seed must not be zero"); + RandomJsonDocs { rng: Rng(seed) } + } + + /// Returns a random JSON object, serialized. Keys may be duplicated. + pub fn next_doc(&mut self) -> String { + let mut json_doc = String::new(); + write_random_object(&mut self.rng, 0, &mut json_doc); + json_doc + } +} + +/// Deterministic xorshift generator, to keep failures reproducible without a new dependency. +struct Rng(u64); + +impl Rng { + fn next(&mut self) -> u64 { + self.0 ^= self.0 << 13; + self.0 ^= self.0 >> 7; + self.0 ^= self.0 << 17; + self.0 + } + + fn below(&mut self, bound: u64) -> u64 { + self.next() % bound + } + + fn pick<'a>(&mut self, items: &[&'a str]) -> &'a str { + items[self.below(items.len() as u64) as usize] + } +} + +const KEYS: &[&str] = &[ + "timestamp", + "service", + "body", + "count", + "delta", + "ratio", + "flag", + "ip", + "payload", + "hex_payload", + "tags", + "values", + "attributes", + "events", + "resource", + "host", + "pid", + "inner", + "zone", + "all_text", + "unmapped", + "a", + "b", + "é", + "", + "timestamp_nanos", + "service_name", + "severity_text", + "k.with.dots", +]; + +const STRINGS: &[&str] = &[ + "", + "text", + "2024-01-02T03:04:05Z", + "2024-01-02T03:04:05.123+02:00", + "2024-01-02 03:04:05", + "1704164645", + "-12", + "1.5", + "192.168.1.1", + "::ffff:10.0.0.1", + "aGVsbG8=", + "deadbeef", + "true", + "9 lives", + "esc\\\"aped\\n", + "\\u00e9t\\u00e9", + "\\ud83d\\ude00", +]; + +const NUMBERS: &[&str] = &[ + "0", + "1", + "-1", + "42", + "1704164645", + "1704164645123", + "-9223372036854775808", + "9223372036854775808", + "18446744073709551615", + "18446744073709551616", + "0.5", + "-0.0", + "1e3", + "1.7976931348623157e308", + "3.0", +]; + +fn write_random_value(rng: &mut Rng, depth: usize, output: &mut String) { + let kind = if depth >= 3 { + rng.below(5) + } else { + rng.below(8) + }; + match kind { + 0 => output.push_str("null"), + 1 => output.push_str(if rng.below(2) == 0 { "true" } else { "false" }), + 2 | 3 => output.push_str(rng.pick(NUMBERS)), + 4 => { + output.push('"'); + output.push_str(rng.pick(STRINGS)); + output.push('"'); + } + 5 => { + output.push('['); + let num_elements = rng.below(4); + for i in 0..num_elements { + if i > 0 { + output.push(','); + } + write_random_value(rng, depth + 1, output); + } + output.push(']'); + } + _ => write_random_object(rng, depth + 1, output), + } +} + +/// Keys are drawn from a small vocabulary, so objects regularly contain duplicate keys. +fn write_random_object(rng: &mut Rng, depth: usize, output: &mut String) { + output.push('{'); + let num_entries = rng.below(6); + for i in 0..num_entries { + if i > 0 { + output.push(','); + } + output.push('"'); + output.push_str(rng.pick(KEYS)); + output.push_str("\":"); + write_random_value(rng, depth, output); + } + output.push('}'); +} diff --git a/quickwit/quickwit-doc-mapper/src/lib.rs b/quickwit/quickwit-doc-mapper/src/lib.rs index 8dee8d700ed..1f6232776ca 100644 --- a/quickwit/quickwit-doc-mapper/src/lib.rs +++ b/quickwit/quickwit-doc-mapper/src/lib.rs @@ -29,10 +29,13 @@ mod routing_expression; /// Pruning tags manipulation. pub mod tag_pruning; +#[cfg(any(test, feature = "testsuite"))] +pub use doc_mapper::RandomJsonDocs; pub use doc_mapper::{ - Automaton, BinaryFormat, DocMapper, DocMapperBuilder, FastFieldWarmupInfo, FieldMappingEntry, - FieldMappingType, JsonObject, NamedField, QuickwitBytesOptions, QuickwitJsonOptions, TermRange, - TokenizerConfig, TokenizerEntry, WarmupInfo, analyze_text, + Automaton, BinaryFormat, BorrowedJsonDoc, BorrowedObject, BorrowedValue, DocMapper, + DocMapperBuilder, FastFieldWarmupInfo, FieldMappingEntry, FieldMappingType, JsonObject, + NamedField, QuickwitBytesOptions, QuickwitJsonOptions, TermRange, TokenizerConfig, + TokenizerEntry, WarmupInfo, analyze_text, }; use doc_mapper::{ FastFieldOptions, FieldMappingEntryForSerialization, IndexRecordOptionSchema, diff --git a/quickwit/quickwit-doc-mapper/src/routing_expression/borrowed.rs b/quickwit/quickwit-doc-mapper/src/routing_expression/borrowed.rs new file mode 100644 index 00000000000..e6c174d9976 --- /dev/null +++ b/quickwit/quickwit-doc-mapper/src/routing_expression/borrowed.rs @@ -0,0 +1,130 @@ +// 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. + +//! Routing expression evaluation on borrowed JSON documents. +//! +//! The partition of a document must not depend on the conversion path, so this mirrors exactly the +//! `serde_json::Map` implementation of the parent module. + +use std::hash::{Hash, Hasher}; + +use super::RoutingExprContext; +use crate::doc_mapper::BorrowedJsonDoc; +use crate::doc_mapper::borrowed_json::{BorrowedValue, get_in_object}; + +/// Mirrors `hash_json_val`. Relies on objects being sorted by key with unique keys, like +/// `serde_json::Map`. +fn hash_borrowed_json_val(json_val: &BorrowedValue, hasher: &mut H) { + match json_val { + BorrowedValue::Null => { + hasher.write_u8(0u8); + } + BorrowedValue::Bool(bool_val) => { + hasher.write_u8(1u8); + bool_val.hash(hasher); + } + BorrowedValue::Number(num) => { + hasher.write_u8(2u8); + num.hash(hasher); + } + BorrowedValue::Str(text) => { + hasher.write_u8(3u8); + hasher.write_u64(text.len() as u64); + hasher.write(text.as_bytes()); + } + BorrowedValue::Array(elements) => { + hasher.write_u8(4u8); + hasher.write_u64(elements.len() as u64); + for element in elements { + hash_borrowed_json_val(element, hasher); + } + } + BorrowedValue::Object(entries) => { + hasher.write_u8(5u8); + hasher.write_u64(entries.len() as u64); + for (key, value) in entries { + hasher.write_u64(key.len() as u64); + hasher.write(key.as_bytes()); + hash_borrowed_json_val(value, hasher); + } + } + } +} + +/// Mirrors `find_value_in_map`. `keys` is never empty. +fn find_value_in_doc<'b, 'a>( + json_doc: &'b BorrowedJsonDoc<'a>, + keys: &[String], +) -> Option<&'b BorrowedValue<'a>> { + let mut value = get_in_object(json_doc.root(), &keys[0])?; + for key in &keys[1..] { + let BorrowedValue::Object(entries) = value else { + return None; + }; + value = get_in_object(entries, key)?; + } + Some(value) +} + +impl RoutingExprContext for BorrowedJsonDoc<'_> { + fn hash_attribute(&self, attr_name: &[String], hasher: &mut H) { + if let Some(json_val) = find_value_in_doc(self, attr_name) { + hasher.write_u8(1u8); + hash_borrowed_json_val(json_val, hasher); + } else { + hasher.write_u8(0u8); + } + } +} + +#[cfg(test)] +mod tests { + use crate::RoutingExpr; + use crate::doc_mapper::BorrowedJsonDoc; + + #[test] + fn test_borrowed_routing_hash_same_as_serde_json() { + let expressions = [ + "service", + "service,tenant_id", + "resource.attributes.host", + "hash_mod(tenant_id, 10)", + "resource", + "missing,service", + "tags", + ]; + let docs = [ + r#"{"service": "api", "tenant_id": 12}"#, + r#"{"tenant_id": -3, "service": "api", "service": "web"}"#, + r#"{"resource": {"attributes": {"host": "h1", "zone": 1.5}, "a": null}}"#, + r#"{"resource": {"attributes": "not an object"}}"#, + r#"{"tags": ["a", 1, true, null, {"z": 1, "a": 2, "z": 3}]}"#, + r#"{"service": "esc\"aped", "tenant_id": 18446744073709551615}"#, + r#"{}"#, + ]; + for expression in expressions { + let routing_expr = RoutingExpr::new(expression).unwrap(); + for doc in docs { + let json_obj: serde_json::Map = + serde_json::from_str(doc).unwrap(); + let borrowed_doc = BorrowedJsonDoc::parse(doc.as_bytes()).unwrap(); + assert_eq!( + routing_expr.eval_hash(&json_obj), + routing_expr.eval_hash(&borrowed_doc), + "expression: {expression}, doc: {doc}" + ); + } + } + } +} diff --git a/quickwit/quickwit-doc-mapper/src/routing_expression/mod.rs b/quickwit/quickwit-doc-mapper/src/routing_expression/mod.rs index 69af8751be5..c63003dd786 100644 --- a/quickwit/quickwit-doc-mapper/src/routing_expression/mod.rs +++ b/quickwit/quickwit-doc-mapper/src/routing_expression/mod.rs @@ -22,6 +22,8 @@ pub(crate) use expression_dsl::parse_field_name; use serde_json::Value as JsonValue; use siphasher::sip::SipHasher; +mod borrowed; + pub trait RoutingExprContext { fn hash_attribute(&self, attr_name: &[String], hasher: &mut H); } @@ -29,6 +31,8 @@ pub trait RoutingExprContext { /// This is a bit overkill but this function has the merit of /// ensuring that the data that is sent to the hasher is unique /// to the value, so we do not lose injectivity there. +/// +/// `borrowed::hash_borrowed_json_val` must feed the hasher exactly the same data. fn hash_json_val(json_val: &JsonValue, hasher: &mut H) { match json_val { JsonValue::Null => { diff --git a/quickwit/quickwit-indexing/src/actors/doc_processor.rs b/quickwit/quickwit-indexing/src/actors/doc_processor.rs index 9c6294bea6d..ce03aa81319 100644 --- a/quickwit/quickwit-indexing/src/actors/doc_processor.rs +++ b/quickwit/quickwit-indexing/src/actors/doc_processor.rs @@ -23,7 +23,7 @@ use quickwit_actors::{Actor, ActorContext, ActorExitStatus, Handler, Mailbox, Qu use quickwit_common::rate_limited_tracing::rate_limited_warn; use quickwit_common::runtimes::RuntimeType; use quickwit_config::{SourceInputFormat, TransformConfig}; -use quickwit_doc_mapper::{DocMapper, DocParsingError, JsonObject}; +use quickwit_doc_mapper::{BorrowedJsonDoc, DocMapper, DocParsingError, JsonObject}; use quickwit_metrics::{Counter, counter, labels}; use quickwit_opentelemetry::otlp::{ JsonLogIterator, JsonSpanIterator, OtlpLogsError, OtlpTracesError, parse_otlp_logs_json, @@ -467,6 +467,12 @@ impl DocProcessor { fn process_raw_doc(&mut self, raw_doc: Bytes, processed_docs: &mut Vec) { let num_bytes = raw_doc.len(); + if self.can_process_borrowed_json() { + let processed_doc_result = self.process_borrowed_json_doc(&raw_doc); + self.record_processed_doc(processed_doc_result, num_bytes, processed_docs); + return; + } + #[cfg(feature = "vrl")] let transform_opt = self.transform_opt.as_mut(); #[cfg(not(feature = "vrl"))] @@ -475,25 +481,68 @@ impl DocProcessor { for json_doc_result in parse_raw_doc(self.input_format, raw_doc, num_bytes, transform_opt) { let processed_doc_result = json_doc_result.and_then(|json_doc| self.process_json_doc(json_doc)); + self.record_processed_doc(processed_doc_result, num_bytes, processed_docs); + } + } - match processed_doc_result { - Ok(processed_doc) => { - self.counters.record_valid(processed_doc.num_bytes as u64); - processed_docs.push(processed_doc); - } - Err(error) => { - rate_limited_warn!( - limit_per_min = 10, - index_id = self.counters.index_id, - source_id = self.counters.source_id, - "{error}", - ); - self.counters.record_error(error, num_bytes as u64); - } + fn record_processed_doc( + &self, + processed_doc_result: Result, + num_bytes: usize, + processed_docs: &mut Vec, + ) { + match processed_doc_result { + Ok(processed_doc) => { + self.counters.record_valid(processed_doc.num_bytes as u64); + processed_docs.push(processed_doc); + } + Err(error) => { + rate_limited_warn!( + limit_per_min = 10, + index_id = self.counters.index_id, + source_id = self.counters.source_id, + "{error}", + ); + self.counters.record_error(error, num_bytes as u64); } } } + /// Returns true if raw documents can be converted with + /// [`DocMapper::doc_from_borrowed_json`], which avoids building an owned JSON object. + /// + /// This requires JSON input and no VRL transform (it operates on owned values). + fn can_process_borrowed_json(&self) -> bool { + #[cfg(feature = "vrl")] + let has_transform = self.transform_opt.is_some(); + #[cfg(not(feature = "vrl"))] + let has_transform = false; + + self.input_format == SourceInputFormat::Json && !has_transform + } + + /// Same as `try_into_json_docs` followed by [`Self::process_json_doc`] for JSON input, without + /// a transform. + fn process_borrowed_json_doc(&self, raw_doc: &[u8]) -> Result { + let num_bytes = raw_doc.len(); + let json_doc = BorrowedJsonDoc::parse(raw_doc)?; + let (partition, doc) = self + .doc_mapper + .doc_from_borrowed_json(&json_doc, num_bytes as u64)?; + let timestamp_opt = self.extract_timestamp(&doc)?; + let fingerprint_opt = self + .fingerprinter_opt + .as_ref() + .map(|fingerprinter| fingerprinter.fingerprint_borrowed(&json_doc)); + Ok(ProcessedDoc { + doc, + fingerprint_opt, + timestamp_opt, + partition, + num_bytes, + }) + } + fn process_json_doc(&self, json_doc: JsonDoc) -> Result { let num_bytes = json_doc.num_bytes; let json_value = JsonValue::Object(json_doc.json_obj); @@ -762,11 +811,10 @@ mod tests { .unwrap(); let (doc_processor_mailbox, doc_processor_handle) = universe.spawn_builder().spawn(doc_processor); + // The duplicate key makes sure the fingerprint is computed on the deduplicated document. + let raw_doc = br#"{"body":"sad 1","timestamp":1628837062,"body":"happy 2"}"#; doc_processor_mailbox - .send_message(RawDocBatch::for_test( - &[br#"{"body":"happy","timestamp":1628837062}"#], - 0..1, - )) + .send_message(RawDocBatch::for_test(&[raw_doc], 0..1)) .await .unwrap(); doc_processor_handle.process_pending_and_observe().await; @@ -774,6 +822,9 @@ mod tests { let output_messages: Vec = indexer_inbox.drain_for_test_typed(); assert_eq!(output_messages.len(), 1); let fingerprint = output_messages[0].docs[0].fingerprint_opt.as_ref().unwrap(); + let json_value: JsonValue = serde_json::from_slice(raw_doc).unwrap(); + let expected_fingerprint = raw_body_fingerprinter_for_test().fingerprint(&json_value); + assert_eq!(fingerprint, &expected_fingerprint); assert_eq!(fingerprint.len(), 3); universe.assert_quit().await; } diff --git a/quickwit/quickwit-indexing/src/docs_clustering/fingerprinter.rs b/quickwit/quickwit-indexing/src/docs_clustering/fingerprinter.rs index d7bae17a397..3afa34210d6 100644 --- a/quickwit/quickwit-indexing/src/docs_clustering/fingerprinter.rs +++ b/quickwit/quickwit-indexing/src/docs_clustering/fingerprinter.rs @@ -65,9 +65,11 @@ use fnv::FnvHasher; use quickwit_config::{ ClusteringMethod, ClusteringPolicy, DocsClusteringConfig, FingerprintPolicy, JsonPath, }; +use quickwit_doc_mapper::BorrowedJsonDoc; use serde_json::Value as JsonValue; use smallvec::SmallVec; +use super::json_view::{BorrowedJsonNode, JsonView, JsonViewKind}; use super::tokenize; const PATH_COMPONENT_SEPARATOR: u8 = 0xF9; @@ -127,7 +129,18 @@ impl Fingerprinter { &self.config } + /// Computes the fingerprint of a document parsed as an owned JSON value. pub fn fingerprint(&self, json_value: &JsonValue) -> Fingerprint { + self.fingerprint_view(json_value) + } + + /// Computes the fingerprint of a document parsed as a [`BorrowedJsonDoc`]. Returns the same + /// fingerprint as [`Self::fingerprint`] for the same JSON input. + pub fn fingerprint_borrowed(&self, json_doc: &BorrowedJsonDoc) -> Fingerprint { + self.fingerprint_view(BorrowedJsonNode::Root(json_doc)) + } + + fn fingerprint_view<'a>(&self, json_view: impl JsonView<'a>) -> Fingerprint { let mut fingerprint = SmallVec::new(); for policy in self.policies.iter() { let mut hasher = FnvHasher::default(); @@ -135,13 +148,13 @@ impl Fingerprinter { for method in policy.fingerprint.iter() { match method { ClusteringMethod::Structure { exclude } => { - self.hash_structure(json_value, exclude, &mut hasher); + hash_structure(json_view, exclude, &mut hasher); } ClusteringMethod::Raw { path } => { - self.hash_raw_value(json_value, path, &mut hasher); + hash_raw_value(json_view, path, &mut hasher); } ClusteringMethod::Tokenized { path, max_tokens } => { - self.hash_string_tokenized(json_value, path, *max_tokens, &mut hasher); + hash_string_tokenized(json_view, path, *max_tokens, &mut hasher); } } } @@ -151,404 +164,136 @@ impl Fingerprinter { Fingerprint(fingerprint) } +} - fn hash_structure(&self, value: &JsonValue, exclude: &[JsonPath], hasher: &mut FnvHasher) { - fn walk<'a>( - json_value: &'a JsonValue, - exclude: &[JsonPath], - current: &mut Vec<&'a str>, - paths: &mut Vec>, - ) { - fn is_excluded(exclude: &[JsonPath], path: &[&str]) -> bool { - exclude.iter().any(|excluded_path| { - excluded_path.len() == path.len() - && excluded_path.iter().zip(path.iter()).all( - |(excluded_component, component)| { - excluded_component.as_str() == *component - }, - ) - }) - } +fn is_excluded(exclude: &[JsonPath], path: &[&str]) -> bool { + exclude.iter().any(|excluded_path| { + excluded_path.len() == path.len() + && excluded_path + .iter() + .zip(path.iter()) + .all(|(excluded_component, component)| excluded_component.as_str() == *component) + }) +} - match json_value { - JsonValue::Object(obj) => { - for (key, child_value) in obj.iter() { - current.push(key.as_str()); - if !is_excluded(exclude, current) { - walk(child_value, exclude, current, paths); - } - current.pop(); - } - } - _ => paths.push(current.clone()), - } +fn collect_leaf_paths<'a, V: JsonView<'a>>( + json_view: V, + exclude: &[JsonPath], + current: &mut Vec<&'a str>, + paths: &mut Vec>, +) { + if !matches!(json_view.kind(), JsonViewKind::Object { .. }) { + paths.push(current.clone()); + return; + } + json_view.for_each_entry(|key, child_view| { + current.push(key); + if !is_excluded(exclude, current) { + collect_leaf_paths(child_view, exclude, current, paths); } + current.pop(); + }); +} - let mut current = Vec::with_capacity(16); - let mut paths = Vec::with_capacity(32); - walk(value, exclude, &mut current, &mut paths); - paths.sort_unstable(); +fn hash_structure<'a>(json_view: impl JsonView<'a>, exclude: &[JsonPath], hasher: &mut FnvHasher) { + let mut current = Vec::with_capacity(16); + let mut paths = Vec::with_capacity(32); + collect_leaf_paths(json_view, exclude, &mut current, &mut paths); + paths.sort_unstable(); - for path in paths { - for component in path { - hasher.write(component.as_bytes()); - hasher.write_u8(PATH_COMPONENT_SEPARATOR); - } - hasher.write_u8(PATH_SEPARATOR); + for path in paths { + for component in path { + hasher.write(component.as_bytes()); + hasher.write_u8(PATH_COMPONENT_SEPARATOR); } + hasher.write_u8(PATH_SEPARATOR); } +} - fn hash_raw_value(&self, json_value: &JsonValue, path: &JsonPath, hasher: &mut FnvHasher) { - fn hash_raw_value_inner(json_value: &JsonValue, hasher: &mut FnvHasher) { - const RAW_NULL: u8 = 0; - const RAW_BOOL: u8 = 1; - const RAW_NUMBER: u8 = 2; - const RAW_STRING: u8 = 3; - const RAW_ARRAY: u8 = 4; - const RAW_OBJECT: u8 = 5; - - match json_value { - JsonValue::Null => hasher.write_u8(RAW_NULL), - JsonValue::Bool(value) => { - hasher.write_u8(RAW_BOOL); - hasher.write_u8(*value as u8); - } - JsonValue::Number(value) => { - hasher.write_u8(RAW_NUMBER); - let number = value.to_string(); - hasher.write(number.as_bytes()); - } - JsonValue::String(value) => { - hasher.write_u8(RAW_STRING); - hasher.write_usize(value.len()); - hasher.write(value.as_bytes()); - } - JsonValue::Array(values) => { - hasher.write_u8(RAW_ARRAY); - hasher.write_usize(values.len()); - for value in values { - hash_raw_value_inner(value, hasher); - } - } - JsonValue::Object(map) => { - hasher.write_u8(RAW_OBJECT); - hasher.write_usize(map.len()); - for (key, value) in map.iter() { - hasher.write_usize(key.len()); - hasher.write(key.as_bytes()); - hash_raw_value_inner(value, hasher); - } - } - } +fn hash_raw_value_inner<'a, V: JsonView<'a>>(json_view: V, hasher: &mut FnvHasher) { + const RAW_NULL: u8 = 0; + const RAW_BOOL: u8 = 1; + const RAW_NUMBER: u8 = 2; + const RAW_STRING: u8 = 3; + const RAW_ARRAY: u8 = 4; + const RAW_OBJECT: u8 = 5; + + match json_view.kind() { + JsonViewKind::Null => hasher.write_u8(RAW_NULL), + JsonViewKind::Bool(value) => { + hasher.write_u8(RAW_BOOL); + hasher.write_u8(value as u8); } - - let Some(json_value) = get_leaf_json_value(json_value, path) else { - hasher.write_u8(FIELD_ABSENT); - hasher.write_u8(FIELD_BOUNDARY); - return; - }; - hasher.write_u8(FIELD_PRESENT); - hash_raw_value_inner(json_value, hasher); - hasher.write_u8(FIELD_BOUNDARY); - } - - fn hash_string_tokenized( - &self, - json_value: &JsonValue, - path: &JsonPath, - max_tokens: Option, - hasher: &mut FnvHasher, - ) { - let Some(value) = get_leaf_string(json_value, path) else { - hasher.write_u8(FIELD_ABSENT); - hasher.write_u8(FIELD_BOUNDARY); - return; - }; - hasher.write_u8(FIELD_PRESENT); - - let max_tokens = max_tokens.unwrap_or(DEFAULT_MAX_GROUPING_TOKENS); - for span in tokenize(value).take(max_tokens) { - hasher.write_u8(span.token_type as u8); - hasher.write_u8(TOKENIZED_TOKEN_SEPARATOR); + JsonViewKind::Number(value) => { + hasher.write_u8(RAW_NUMBER); + let number = value.to_string(); + hasher.write(number.as_bytes()); + } + JsonViewKind::Str(value) => { + hasher.write_u8(RAW_STRING); + hasher.write_usize(value.len()); + hasher.write(value.as_bytes()); + } + JsonViewKind::Array { len } => { + hasher.write_u8(RAW_ARRAY); + hasher.write_usize(len); + json_view.for_each_element(|element_view| hash_raw_value_inner(element_view, hasher)); + } + JsonViewKind::Object { len } => { + hasher.write_u8(RAW_OBJECT); + hasher.write_usize(len); + json_view.for_each_entry(|key, child_view| { + hasher.write_usize(key.len()); + hasher.write(key.as_bytes()); + hash_raw_value_inner(child_view, hasher); + }); } - - hasher.write_u8(FIELD_BOUNDARY); } } -fn get_leaf_json_value<'a>(json_value: &'a JsonValue, path: &[String]) -> Option<&'a JsonValue> { - if path.is_empty() { - return Some(json_value); - } - let JsonValue::Object(obj) = json_value else { - return None; +fn hash_raw_value<'a>(json_view: impl JsonView<'a>, path: &JsonPath, hasher: &mut FnvHasher) { + let Some(leaf_view) = get_leaf_json_view(json_view, path) else { + hasher.write_u8(FIELD_ABSENT); + hasher.write_u8(FIELD_BOUNDARY); + return; }; - get_leaf_json_value(obj.get(path.first()?)?, &path[1..]) + hasher.write_u8(FIELD_PRESENT); + hash_raw_value_inner(leaf_view, hasher); + hasher.write_u8(FIELD_BOUNDARY); } -fn get_leaf_string<'a>(json_value: &'a JsonValue, path: &[String]) -> Option<&'a str> { - get_leaf_json_value(json_value, path)?.as_str() -} - -#[cfg(test)] -mod tests { - use quickwit_config::DocsClusteringConfig; - use serde_json::Value as JsonValue; - - use super::Fingerprinter; - - fn parse(s: &str) -> JsonValue { - serde_json::from_str(s).unwrap() - } - - fn test_docs_clustering_config() -> DocsClusteringConfig { - docs_clustering_config(serde_json::json!([ - { - "fingerprint": [{ - "kind": "structure", - "exclude": ["tag", "custom"] - }] - }, - { - "fingerprint": [{ - "path": "message", - "kind": "tokenized" - }, - { - "path": "service", - "kind": "raw" - }] - } - ])) - } - - fn test_fingerprinter() -> Fingerprinter { - let docs_clustering_config = test_docs_clustering_config(); - Fingerprinter::new(&docs_clustering_config) - } - - fn docs_clustering_config(json_value: JsonValue) -> DocsClusteringConfig { - serde_json::from_value(json_value).unwrap() - } - - #[test] - fn configured_fingerprinter_returns_config() { - let docs_clustering_config = test_docs_clustering_config(); - let fingerprinter = test_fingerprinter(); - assert_eq!(fingerprinter.config(), &docs_clustering_config); - } - - #[test] - fn identical_logs_have_equal_fingerprint() { - let fingerprinter = test_fingerprinter(); - let doc = parse(r#"{"message":"server started at 8080","service":"api"}"#); - assert_eq!( - fingerprinter.fingerprint(&doc), - fingerprinter.fingerprint(&doc) - ); - } - - #[test] - fn dotted_key_and_nested_path_have_different_schema_fingerprints() { - let fingerprinter = test_fingerprinter(); - let dotted_key_doc = parse(r#"{"a.b":1}"#); - let nested_path_doc = parse(r#"{"a":{"b":1}}"#); - let dotted_key_fingerprint = fingerprinter.fingerprint(&dotted_key_doc); - let nested_path_fingerprint = fingerprinter.fingerprint(&nested_path_doc); - - assert_ne!(dotted_key_fingerprint[0], nested_path_fingerprint[0]); - assert_eq!(dotted_key_fingerprint[1], nested_path_fingerprint[1]); - } - - #[test] - fn same_message_template_has_equal_fingerprint() { - let fingerprinter = test_fingerprinter(); - let doc1 = parse(r#"{"message":"server started at 8080","service":"api"}"#); - let doc2 = parse(r#"{"message":"server started at 9090","service":"api"}"#); - assert_eq!( - fingerprinter.fingerprint(&doc1), - fingerprinter.fingerprint(&doc2) - ); - } - - #[test] - fn different_message_template_changes_grouping_fingerprint_only() { - let fingerprinter = test_fingerprinter(); - let doc1 = parse(r#"{"message":"server started at 8080","service":"api"}"#); - let doc2 = parse(r#"{"message":"connection from 1.2.3.4","service":"api"}"#); - let doc1_fingerprint = fingerprinter.fingerprint(&doc1); - let doc2_fingerprint = fingerprinter.fingerprint(&doc2); - assert_eq!(doc1_fingerprint[0], doc2_fingerprint[0]); - assert_ne!(doc1_fingerprint[1], doc2_fingerprint[1]); - } - - #[test] - fn tokenized_field_respects_hardcoded_grouping_token_limit() { - let docs_clustering_config = docs_clustering_config(serde_json::json!([ - { - "fingerprint": [{ - "kind": "structure" - }] - }, - { - "fingerprint": [{ - "path": "message", - "kind": "tokenized" - }] - } - ])); - let fingerprinter = Fingerprinter::new(&docs_clustering_config); - let prefix = "alpha ".repeat(25); - let doc1 = parse(&format!(r#"{{"message":"{prefix}123"}}"#)); - let doc2 = parse(&format!(r#"{{"message":"{prefix}beta"}}"#)); - assert_eq!( - fingerprinter.fingerprint(&doc1)[1], - fingerprinter.fingerprint(&doc2)[1] - ); - } - - #[test] - fn different_service_changes_grouping_fingerprint_only() { - let fingerprinter = test_fingerprinter(); - let doc1 = parse(r#"{"message":"server started at 8080","service":"api"}"#); - let doc2 = parse(r#"{"message":"server started at 8080","service":"worker"}"#); - let doc1_fingerprint = fingerprinter.fingerprint(&doc1); - let doc2_fingerprint = fingerprinter.fingerprint(&doc2); - assert_eq!(doc1_fingerprint[0], doc2_fingerprint[0]); - assert_ne!(doc1_fingerprint[1], doc2_fingerprint[1]); - } - - #[test] - fn ignored_custom_shape_does_not_change_fingerprint() { - let fingerprinter = test_fingerprinter(); - let doc1 = - parse(r#"{"message":"server started at 8080","service":"api","custom":{"a":1}}"#); - let doc2 = - parse(r#"{"message":"server started at 8080","service":"api","custom":{"b":2}}"#); - assert_eq!( - fingerprinter.fingerprint(&doc1), - fingerprinter.fingerprint(&doc2) - ); - } - - #[test] - fn extra_non_ignored_shape_changes_schema_fingerprint_only() { - let fingerprinter = test_fingerprinter(); - let doc1 = parse(r#"{"message":"server started at 8080","service":"api"}"#); - let doc2 = parse(r#"{"message":"server started at 8080","service":"api","host":"web-1"}"#); - let doc1_fingerprint = fingerprinter.fingerprint(&doc1); - let doc2_fingerprint = fingerprinter.fingerprint(&doc2); - assert_ne!(doc1_fingerprint[0], doc2_fingerprint[0]); - assert_eq!(doc1_fingerprint[1], doc2_fingerprint[1]); - } - - #[test] - fn configured_raw_field_changes_grouping_fingerprint_only() { - let docs_clustering_config = docs_clustering_config(serde_json::json!([ - { - "fingerprint": [{ - "kind": "structure" - }] - }, - { - "fingerprint": [{ - "path": "host", - "kind": "raw" - }] - } - ])); - let fingerprinter = Fingerprinter::new(&docs_clustering_config); - let doc1 = parse(r#"{"message":"same","host":"web-1"}"#); - let doc2 = parse(r#"{"message":"same","host":"web-2"}"#); - let doc1_fingerprint = fingerprinter.fingerprint(&doc1); - let doc2_fingerprint = fingerprinter.fingerprint(&doc2); - assert_eq!(doc1_fingerprint[0], doc2_fingerprint[0]); - assert_ne!(doc1_fingerprint[1], doc2_fingerprint[1]); - } +fn hash_string_tokenized<'a>( + json_view: impl JsonView<'a>, + path: &JsonPath, + max_tokens: Option, + hasher: &mut FnvHasher, +) { + let Some(value) = get_leaf_string(json_view, path) else { + hasher.write_u8(FIELD_ABSENT); + hasher.write_u8(FIELD_BOUNDARY); + return; + }; + hasher.write_u8(FIELD_PRESENT); - #[test] - fn raw_field_hashes_non_string_values() { - let docs_clustering_config = docs_clustering_config(serde_json::json!([ - { - "fingerprint": [{ - "kind": "structure" - }] - }, - { - "fingerprint": [{ - "path": "status", - "kind": "raw" - }] - } - ])); - let fingerprinter = Fingerprinter::new(&docs_clustering_config); - let doc1 = parse(r#"{"status":200}"#); - let doc2 = parse(r#"{"status":500}"#); - let doc1_fingerprint = fingerprinter.fingerprint(&doc1); - let doc2_fingerprint = fingerprinter.fingerprint(&doc2); - assert_eq!(doc1_fingerprint[0], doc2_fingerprint[0]); - assert_ne!(doc1_fingerprint[1], doc2_fingerprint[1]); + let max_tokens = max_tokens.unwrap_or(DEFAULT_MAX_GROUPING_TOKENS); + for span in tokenize(value).take(max_tokens) { + hasher.write_u8(span.token_type as u8); + hasher.write_u8(TOKENIZED_TOKEN_SEPARATOR); } - #[test] - fn raw_field_preserves_json_type_and_structure_boundaries() { - let docs_clustering_config = docs_clustering_config(serde_json::json!([ - { - "fingerprint": [{ - "path": "value", - "kind": "raw" - }] - } - ])); - let fingerprinter = Fingerprinter::new(&docs_clustering_config); - let docs = [ - parse(r#"{"value":"null"}"#), - parse(r#"{"value":null}"#), - parse(r#"{"value":""}"#), - parse(r#"{"value":[]}"#), - parse(r#"{"value":{}}"#), - parse(r#"{"value":false}"#), - parse(r#"{"value":-1}"#), - parse(r#"{"value":18446744073709551615}"#), - ]; - let fingerprints: Vec<_> = docs - .iter() - .map(|doc| fingerprinter.fingerprint(doc)) - .collect(); - - for (left_idx, left_fingerprint) in fingerprints.iter().enumerate() { - for right_fingerprint in &fingerprints[left_idx + 1..] { - assert_ne!(left_fingerprint, right_fingerprint); - } - } - } + hasher.write_u8(FIELD_BOUNDARY); +} - #[test] - fn absent_grouping_values_preserve_field_position() { - let docs_clustering_config = docs_clustering_config(serde_json::json!([ - { - "fingerprint": [{ - "kind": "structure" - }] - }, - { - "fingerprint": [{ - "path": "a", - "kind": "raw" - }, - { - "path": "b", - "kind": "raw" - }] - } - ])); - let fingerprinter = Fingerprinter::new(&docs_clustering_config); - let doc1 = parse(r#"{"a":"x","b":null}"#); - let doc2 = parse(r#"{"a":null,"b":"x"}"#); - let doc1_fingerprint = fingerprinter.fingerprint(&doc1); - let doc2_fingerprint = fingerprinter.fingerprint(&doc2); +fn get_leaf_json_view<'a, V: JsonView<'a>>(json_view: V, path: &[String]) -> Option { + let Some((first_component, remaining_components)) = path.split_first() else { + return Some(json_view); + }; + get_leaf_json_view(json_view.get(first_component)?, remaining_components) +} - assert_eq!(doc1_fingerprint[0], doc2_fingerprint[0]); - assert_ne!(doc1_fingerprint[1], doc2_fingerprint[1]); - } +fn get_leaf_string<'a>(json_view: impl JsonView<'a>, path: &[String]) -> Option<&'a str> { + let JsonViewKind::Str(value) = get_leaf_json_view(json_view, path)?.kind() else { + return None; + }; + Some(value) } diff --git a/quickwit/quickwit-indexing/src/docs_clustering/fingerprinter_tests.rs b/quickwit/quickwit-indexing/src/docs_clustering/fingerprinter_tests.rs new file mode 100644 index 00000000000..142940ea11d --- /dev/null +++ b/quickwit/quickwit-indexing/src/docs_clustering/fingerprinter_tests.rs @@ -0,0 +1,335 @@ +// 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 quickwit_config::DocsClusteringConfig; +use quickwit_doc_mapper::{BorrowedJsonDoc, RandomJsonDocs}; +use serde_json::Value as JsonValue; + +use super::fingerprinter::Fingerprinter; + +fn parse(s: &str) -> JsonValue { + serde_json::from_str(s).unwrap() +} + +fn test_docs_clustering_config() -> DocsClusteringConfig { + docs_clustering_config(serde_json::json!([ + { + "fingerprint": [{ + "kind": "structure", + "exclude": ["tag", "custom"] + }] + }, + { + "fingerprint": [{ + "path": "message", + "kind": "tokenized" + }, + { + "path": "service", + "kind": "raw" + }] + } + ])) +} + +fn test_fingerprinter() -> Fingerprinter { + let docs_clustering_config = test_docs_clustering_config(); + Fingerprinter::new(&docs_clustering_config) +} + +fn docs_clustering_config(json_value: JsonValue) -> DocsClusteringConfig { + serde_json::from_value(json_value).unwrap() +} + +#[test] +fn configured_fingerprinter_returns_config() { + let docs_clustering_config = test_docs_clustering_config(); + let fingerprinter = test_fingerprinter(); + assert_eq!(fingerprinter.config(), &docs_clustering_config); +} + +#[test] +fn identical_logs_have_equal_fingerprint() { + let fingerprinter = test_fingerprinter(); + let doc = parse(r#"{"message":"server started at 8080","service":"api"}"#); + assert_eq!( + fingerprinter.fingerprint(&doc), + fingerprinter.fingerprint(&doc) + ); +} + +#[test] +fn dotted_key_and_nested_path_have_different_schema_fingerprints() { + let fingerprinter = test_fingerprinter(); + let dotted_key_doc = parse(r#"{"a.b":1}"#); + let nested_path_doc = parse(r#"{"a":{"b":1}}"#); + let dotted_key_fingerprint = fingerprinter.fingerprint(&dotted_key_doc); + let nested_path_fingerprint = fingerprinter.fingerprint(&nested_path_doc); + + assert_ne!(dotted_key_fingerprint[0], nested_path_fingerprint[0]); + assert_eq!(dotted_key_fingerprint[1], nested_path_fingerprint[1]); +} + +#[test] +fn same_message_template_has_equal_fingerprint() { + let fingerprinter = test_fingerprinter(); + let doc1 = parse(r#"{"message":"server started at 8080","service":"api"}"#); + let doc2 = parse(r#"{"message":"server started at 9090","service":"api"}"#); + assert_eq!( + fingerprinter.fingerprint(&doc1), + fingerprinter.fingerprint(&doc2) + ); +} + +#[test] +fn different_message_template_changes_grouping_fingerprint_only() { + let fingerprinter = test_fingerprinter(); + let doc1 = parse(r#"{"message":"server started at 8080","service":"api"}"#); + let doc2 = parse(r#"{"message":"connection from 1.2.3.4","service":"api"}"#); + let doc1_fingerprint = fingerprinter.fingerprint(&doc1); + let doc2_fingerprint = fingerprinter.fingerprint(&doc2); + assert_eq!(doc1_fingerprint[0], doc2_fingerprint[0]); + assert_ne!(doc1_fingerprint[1], doc2_fingerprint[1]); +} + +#[test] +fn tokenized_field_respects_hardcoded_grouping_token_limit() { + let docs_clustering_config = docs_clustering_config(serde_json::json!([ + { + "fingerprint": [{ + "kind": "structure" + }] + }, + { + "fingerprint": [{ + "path": "message", + "kind": "tokenized" + }] + } + ])); + let fingerprinter = Fingerprinter::new(&docs_clustering_config); + let prefix = "alpha ".repeat(25); + let doc1 = parse(&format!(r#"{{"message":"{prefix}123"}}"#)); + let doc2 = parse(&format!(r#"{{"message":"{prefix}beta"}}"#)); + assert_eq!( + fingerprinter.fingerprint(&doc1)[1], + fingerprinter.fingerprint(&doc2)[1] + ); +} + +#[test] +fn different_service_changes_grouping_fingerprint_only() { + let fingerprinter = test_fingerprinter(); + let doc1 = parse(r#"{"message":"server started at 8080","service":"api"}"#); + let doc2 = parse(r#"{"message":"server started at 8080","service":"worker"}"#); + let doc1_fingerprint = fingerprinter.fingerprint(&doc1); + let doc2_fingerprint = fingerprinter.fingerprint(&doc2); + assert_eq!(doc1_fingerprint[0], doc2_fingerprint[0]); + assert_ne!(doc1_fingerprint[1], doc2_fingerprint[1]); +} + +#[test] +fn ignored_custom_shape_does_not_change_fingerprint() { + let fingerprinter = test_fingerprinter(); + let doc1 = parse(r#"{"message":"server started at 8080","service":"api","custom":{"a":1}}"#); + let doc2 = parse(r#"{"message":"server started at 8080","service":"api","custom":{"b":2}}"#); + assert_eq!( + fingerprinter.fingerprint(&doc1), + fingerprinter.fingerprint(&doc2) + ); +} + +#[test] +fn extra_non_ignored_shape_changes_schema_fingerprint_only() { + let fingerprinter = test_fingerprinter(); + let doc1 = parse(r#"{"message":"server started at 8080","service":"api"}"#); + let doc2 = parse(r#"{"message":"server started at 8080","service":"api","host":"web-1"}"#); + let doc1_fingerprint = fingerprinter.fingerprint(&doc1); + let doc2_fingerprint = fingerprinter.fingerprint(&doc2); + assert_ne!(doc1_fingerprint[0], doc2_fingerprint[0]); + assert_eq!(doc1_fingerprint[1], doc2_fingerprint[1]); +} + +#[test] +fn configured_raw_field_changes_grouping_fingerprint_only() { + let docs_clustering_config = docs_clustering_config(serde_json::json!([ + { + "fingerprint": [{ + "kind": "structure" + }] + }, + { + "fingerprint": [{ + "path": "host", + "kind": "raw" + }] + } + ])); + let fingerprinter = Fingerprinter::new(&docs_clustering_config); + let doc1 = parse(r#"{"message":"same","host":"web-1"}"#); + let doc2 = parse(r#"{"message":"same","host":"web-2"}"#); + let doc1_fingerprint = fingerprinter.fingerprint(&doc1); + let doc2_fingerprint = fingerprinter.fingerprint(&doc2); + assert_eq!(doc1_fingerprint[0], doc2_fingerprint[0]); + assert_ne!(doc1_fingerprint[1], doc2_fingerprint[1]); +} + +#[test] +fn raw_field_hashes_non_string_values() { + let docs_clustering_config = docs_clustering_config(serde_json::json!([ + { + "fingerprint": [{ + "kind": "structure" + }] + }, + { + "fingerprint": [{ + "path": "status", + "kind": "raw" + }] + } + ])); + let fingerprinter = Fingerprinter::new(&docs_clustering_config); + let doc1 = parse(r#"{"status":200}"#); + let doc2 = parse(r#"{"status":500}"#); + let doc1_fingerprint = fingerprinter.fingerprint(&doc1); + let doc2_fingerprint = fingerprinter.fingerprint(&doc2); + assert_eq!(doc1_fingerprint[0], doc2_fingerprint[0]); + assert_ne!(doc1_fingerprint[1], doc2_fingerprint[1]); +} + +#[test] +fn raw_field_preserves_json_type_and_structure_boundaries() { + let docs_clustering_config = docs_clustering_config(serde_json::json!([ + { + "fingerprint": [{ + "path": "value", + "kind": "raw" + }] + } + ])); + let fingerprinter = Fingerprinter::new(&docs_clustering_config); + let docs = [ + parse(r#"{"value":"null"}"#), + parse(r#"{"value":null}"#), + parse(r#"{"value":""}"#), + parse(r#"{"value":[]}"#), + parse(r#"{"value":{}}"#), + parse(r#"{"value":false}"#), + parse(r#"{"value":-1}"#), + parse(r#"{"value":18446744073709551615}"#), + ]; + let fingerprints: Vec<_> = docs + .iter() + .map(|doc| fingerprinter.fingerprint(doc)) + .collect(); + + for (left_idx, left_fingerprint) in fingerprints.iter().enumerate() { + for right_fingerprint in &fingerprints[left_idx + 1..] { + assert_ne!(left_fingerprint, right_fingerprint); + } + } +} + +#[test] +fn absent_grouping_values_preserve_field_position() { + let docs_clustering_config = docs_clustering_config(serde_json::json!([ + { + "fingerprint": [{ + "kind": "structure" + }] + }, + { + "fingerprint": [{ + "path": "a", + "kind": "raw" + }, + { + "path": "b", + "kind": "raw" + }] + } + ])); + let fingerprinter = Fingerprinter::new(&docs_clustering_config); + let doc1 = parse(r#"{"a":"x","b":null}"#); + let doc2 = parse(r#"{"a":null,"b":"x"}"#); + let doc1_fingerprint = fingerprinter.fingerprint(&doc1); + let doc2_fingerprint = fingerprinter.fingerprint(&doc2); + + assert_eq!(doc1_fingerprint[0], doc2_fingerprint[0]); + assert_ne!(doc1_fingerprint[1], doc2_fingerprint[1]); +} + +/// Exercises every clustering method, nested paths and exclusions. +fn differential_docs_clustering_config() -> DocsClusteringConfig { + docs_clustering_config(serde_json::json!([ + {"fingerprint": [{"kind": "structure"}]}, + {"fingerprint": [{"kind": "structure", "exclude": ["attributes", "inner.body"]}]}, + {"fingerprint": [ + {"kind": "raw", "path": "service"}, + {"kind": "raw", "path": "count"}, + {"kind": "tokenized", "path": "body"} + ]}, + {"fingerprint": [ + {"kind": "raw", "path": "attributes"}, + {"kind": "raw", "path": "inner.zone"}, + {"kind": "tokenized", "path": "inner.body", "max_tokens": 3}, + {"kind": "raw", "path": "tags"}, + {"kind": "raw", "path": "values"} + ]} + ])) +} + +fn assert_same_fingerprint(fingerprinter: &Fingerprinter, json_doc: &str) { + let json_value = JsonValue::Object(serde_json::from_str(json_doc).unwrap()); + let borrowed_json_doc = BorrowedJsonDoc::parse(json_doc.as_bytes()).unwrap(); + assert_eq!( + fingerprinter.fingerprint_borrowed(&borrowed_json_doc), + fingerprinter.fingerprint(&json_value), + "doc: {json_doc}" + ); +} + +#[test] +fn borrowed_fingerprint_same_as_owned_fingerprint_hand_written() { + let fingerprinter = Fingerprinter::new(&differential_docs_clustering_config()); + let json_docs = [ + r#"{}"#, + r#"{"service": "api", "body": "server started at 8080", "count": 3}"#, + r#"{"body": "first", "service": "api", "body": "connection from 1.2.3.4"}"#, + r#"{"service": "api", "count": 3, "count": -3.5e10, "count": 18446744073709551615}"#, + r#"{"service": 12, "body": 12, "count": "12"}"#, + r#"{"service": null, "body": null, "count": [1, 2.0, -0.0, 1e3]}"#, + r#"{"attributes": {"z": 1, "a": {"y": [true, null], "b": "x"}, "z": 2}}"#, + r#"{"inner": {"zone": {"b": 1, "a": 2}, "body": "job 123 finished in 42ms"}}"#, + r#"{"inner": {"body": "a b c d e f", "body": "x"}, "inner": {"zone": "eu"}}"#, + r#"{"body": "esc\"aped\nline \u00e9t\u00e9 \ud83d\ude00", "service": "s\u0000"}"#, + r#"{"tags": ["a", {"k": "v", "j": [1]}, []], "values": {}}"#, + r#"{"k.with.dots": {"x": 1}, "k": {"with": {"dots": 2}}, "": {"": null}}"#, + r#"{"b": 1, "a": 2, "é": 3, "c": {"y": 1, "x": 2}}"#, + ]; + for json_doc in json_docs { + assert_same_fingerprint(&fingerprinter, json_doc); + } +} + +#[test] +fn borrowed_fingerprint_same_as_owned_fingerprint_random() { + let fingerprinter = Fingerprinter::new(&differential_docs_clustering_config()); + let mut random_json_docs = RandomJsonDocs::new(0x9e37_79b9_7f4a_7c15); + for _ in 0..20_000 { + let json_doc = random_json_docs.next_doc(); + assert_same_fingerprint(&fingerprinter, &json_doc); + } +} diff --git a/quickwit/quickwit-indexing/src/docs_clustering/json_view.rs b/quickwit/quickwit-indexing/src/docs_clustering/json_view.rs new file mode 100644 index 00000000000..8f73cbe8957 --- /dev/null +++ b/quickwit/quickwit-indexing/src/docs_clustering/json_view.rs @@ -0,0 +1,142 @@ +// 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. + +//! Read-only view over a JSON tree, so the [`Fingerprinter`](super::Fingerprinter) can hash both +//! owned `serde_json` values and the borrowed tree built by [`BorrowedJsonDoc`]. +//! +//! Both implementations must expose the same object entries in the same order for the same JSON +//! input: fingerprints are only comparable if they are computed identically. `BorrowedJsonDoc` +//! guarantees this by mirroring the key ordering and duplicate key resolution of +//! `serde_json::Map`. + +use quickwit_doc_mapper::{BorrowedJsonDoc, BorrowedObject, BorrowedValue}; +use serde_json::{Number, Value as JsonValue}; + +/// Kind of a JSON value, with access to its scalar content. +pub(crate) enum JsonViewKind<'a> { + Null, + Bool(bool), + Number(&'a Number), + Str(&'a str), + Array { len: usize }, + Object { len: usize }, +} + +/// Read-only access to a JSON value. Iteration uses callbacks to avoid allocations. +pub(crate) trait JsonView<'a>: Copy { + fn kind(self) -> JsonViewKind<'a>; + + /// Calls `callback` on each element if `self` is an array. + fn for_each_element(self, callback: impl FnMut(Self)); + + /// Calls `callback` on each entry, in the iteration order of `serde_json::Map`, if `self` is + /// an object. + fn for_each_entry(self, callback: impl FnMut(&'a str, Self)); + + /// Returns the value associated with `key` if `self` is an object. + fn get(self, key: &str) -> Option; +} + +impl<'a> JsonView<'a> for &'a JsonValue { + fn kind(self) -> JsonViewKind<'a> { + match self { + JsonValue::Null => JsonViewKind::Null, + JsonValue::Bool(value) => JsonViewKind::Bool(*value), + JsonValue::Number(number) => JsonViewKind::Number(number), + JsonValue::String(value) => JsonViewKind::Str(value), + JsonValue::Array(values) => JsonViewKind::Array { len: values.len() }, + JsonValue::Object(map) => JsonViewKind::Object { len: map.len() }, + } + } + + fn for_each_element(self, callback: impl FnMut(Self)) { + if let JsonValue::Array(values) = self { + values.iter().for_each(callback); + } + } + + fn for_each_entry(self, mut callback: impl FnMut(&'a str, Self)) { + if let JsonValue::Object(map) = self { + for (key, value) in map { + callback(key, value); + } + } + } + + fn get(self, key: &str) -> Option { + self.as_object()?.get(key) + } +} + +/// A node of a [`BorrowedJsonDoc`]: the root object is not stored as a [`BorrowedValue`]. +#[derive(Clone, Copy)] +pub(crate) enum BorrowedJsonNode<'b, 'a> { + Root(&'b BorrowedJsonDoc<'a>), + Value(&'b BorrowedValue<'a>), +} + +impl<'b, 'a> BorrowedJsonNode<'b, 'a> { + fn object_entries(self) -> Option<&'b BorrowedObject<'a>> { + match self { + BorrowedJsonNode::Root(json_doc) => Some(json_doc.root()), + BorrowedJsonNode::Value(BorrowedValue::Object(entries)) => Some(entries), + BorrowedJsonNode::Value(_) => None, + } + } +} + +impl<'b, 'a: 'b> JsonView<'b> for BorrowedJsonNode<'b, 'a> { + fn kind(self) -> JsonViewKind<'b> { + if let Some(entries) = self.object_entries() { + return JsonViewKind::Object { len: entries.len() }; + } + let BorrowedJsonNode::Value(json_value) = self else { + unreachable!("the root is an object") + }; + match json_value { + BorrowedValue::Null => JsonViewKind::Null, + BorrowedValue::Bool(value) => JsonViewKind::Bool(*value), + BorrowedValue::Number(number) => JsonViewKind::Number(number), + BorrowedValue::Str(value) => JsonViewKind::Str(value), + BorrowedValue::Array(values) => JsonViewKind::Array { len: values.len() }, + BorrowedValue::Object(_) => unreachable!("objects are handled above"), + } + } + + fn for_each_element(self, callback: impl FnMut(Self)) { + if let BorrowedJsonNode::Value(BorrowedValue::Array(values)) = self { + values + .iter() + .map(BorrowedJsonNode::Value) + .for_each(callback); + } + } + + fn for_each_entry(self, mut callback: impl FnMut(&'b str, Self)) { + let Some(entries) = self.object_entries() else { + return; + }; + for (key, value) in entries { + callback(key, BorrowedJsonNode::Value(value)); + } + } + + fn get(self, key: &str) -> Option { + let child_value = match self { + BorrowedJsonNode::Root(json_doc) => json_doc.get(key)?, + BorrowedJsonNode::Value(json_value) => json_value.get(key)?, + }; + Some(BorrowedJsonNode::Value(child_value)) + } +} diff --git a/quickwit/quickwit-indexing/src/docs_clustering/mod.rs b/quickwit/quickwit-indexing/src/docs_clustering/mod.rs index 0250c3a4b2f..98e0b1eb059 100644 --- a/quickwit/quickwit-indexing/src/docs_clustering/mod.rs +++ b/quickwit/quickwit-indexing/src/docs_clustering/mod.rs @@ -112,6 +112,9 @@ mod clusterer; mod fingerprinter; +#[cfg(test)] +mod fingerprinter_tests; +mod json_view; mod tokenizer; pub use clusterer::DocIdClusterer;