1use std::sync::OnceLock;
18
19use bytes::Bytes;
20use datatypes::prelude::ConcreteDataType;
21use mito_codec::row_converter::{DensePrimaryKeyCodec, PrimaryKeyCodec, SortField};
22use snafu::{ResultExt, ensure};
23use store_api::codec::PrimaryKeyEncoding;
24use store_api::metadata::{ColumnMetadata, RegionMetadataRef};
25use store_api::storage::RegionId;
26
27use crate::error::{DecodePrimaryKeyRangeSnafu, InvalidPrimaryKeyRangeSnafu, Result};
28
29#[derive(Debug)]
31pub(crate) struct PrimaryKeyRangeMapper {
32 metadata: RegionMetadataRef,
35 codec: DensePrimaryKeyCodec,
36 encoded_defaults: Vec<OnceLock<Option<Bytes>>>,
38 suffixes: Vec<OnceLock<Option<Bytes>>>,
39}
40
41impl PrimaryKeyRangeMapper {
42 pub(crate) fn new(metadata: RegionMetadataRef) -> Self {
43 let codec = DensePrimaryKeyCodec::new(&metadata);
44 let encoded_defaults = (0..codec.num_fields()).map(|_| OnceLock::new()).collect();
45 let suffixes = (0..codec.num_fields()).map(|_| OnceLock::new()).collect();
46 Self {
47 metadata,
48 codec,
49 encoded_defaults,
50 suffixes,
51 }
52 }
53
54 pub(crate) fn with_metadata(&self, metadata: RegionMetadataRef) -> Self {
56 let mut mapper = Self::new(metadata);
57 for (index, (old, new)) in self
58 .metadata
59 .primary_key_columns()
60 .zip(mapper.metadata.primary_key_columns())
61 .enumerate()
62 {
63 if same_pk_column(old, new) {
64 mapper.encoded_defaults[index] = self.encoded_defaults[index].clone();
65 }
66 }
67 mapper
68 }
69
70 pub(crate) fn schema_version(&self) -> u64 {
72 self.metadata.schema_version
73 }
74
75 pub(crate) fn region_id(&self) -> RegionId {
77 self.metadata.region_id
78 }
79
80 pub(crate) fn map(&self, (min, max): (Bytes, Bytes)) -> Result<Option<(Bytes, Bytes)>> {
85 ensure!(
86 min <= max,
87 InvalidPrimaryKeyRangeSnafu {
88 reason: "min is greater than max",
89 }
90 );
91 if self.metadata.primary_key_encoding == PrimaryKeyEncoding::Sparse {
93 return Ok(Some((min, max)));
94 }
95 let prefix_len = self.range_prefix_len(&min, &max)?;
98 if prefix_len == self.codec.num_fields() {
99 return Ok(Some((min, max)));
100 }
101 let Some(suffix) = self.suffixes[prefix_len]
102 .get_or_init(|| self.encode_suffix(prefix_len))
103 .as_ref()
104 else {
105 return Ok(None);
106 };
107 Ok(Some((
110 append_suffix(min, suffix),
111 append_suffix(max, suffix),
112 )))
113 }
114
115 fn range_prefix_len(&self, min: &[u8], max: &[u8]) -> Result<usize> {
116 let min_len = self
117 .codec
118 .decode_prefix_len(min)
119 .context(DecodePrimaryKeyRangeSnafu { endpoint: "min" })?;
120 let max_len = self
121 .codec
122 .decode_prefix_len(max)
123 .context(DecodePrimaryKeyRangeSnafu { endpoint: "max" })?;
124 ensure!(
125 min_len == max_len,
126 InvalidPrimaryKeyRangeSnafu {
127 reason: format!(
128 "endpoints have different field counts: min {min_len}, max {max_len}"
129 ),
130 }
131 );
132 Ok(min_len)
133 }
134
135 fn encode_suffix(&self, prefix_len: usize) -> Option<Bytes> {
136 let mut suffix = Vec::with_capacity(self.codec.num_fields() - prefix_len);
137 for (index, column) in self
138 .metadata
139 .primary_key_columns()
140 .enumerate()
141 .skip(prefix_len)
142 {
143 let encoded = self.encoded_defaults[index]
144 .get_or_init(|| encode_default(column))
145 .as_ref()?;
146 suffix.extend_from_slice(encoded);
147 }
148 Some(suffix.into())
149 }
150}
151
152fn same_pk_column(left: &ColumnMetadata, right: &ColumnMetadata) -> bool {
153 left.column_id == right.column_id
154 && left.column_schema.data_type == right.column_schema.data_type
155 && left.column_schema.is_nullable() == right.column_schema.is_nullable()
156 && left.column_schema.default_constraint() == right.column_schema.default_constraint()
157}
158
159fn encode_default(column: &ColumnMetadata) -> Option<Bytes> {
161 let schema = &column.column_schema;
162 if schema.is_default_impure() {
163 return None;
164 }
165 let default = schema.create_default().ok()??;
166 let field = SortField::new(schema.data_type.clone());
167 if default.is_null() {
168 return match field.encode_data_type() {
170 ConcreteDataType::Null(_)
171 | ConcreteDataType::List(_)
172 | ConcreteDataType::Struct(_)
173 | ConcreteDataType::Dictionary(_) => None,
174 _ => Some(Bytes::from_static(&[0])),
175 };
176 }
177 let mut encoded = Vec::new();
178 DensePrimaryKeyCodec::with_fields(vec![(column.column_id, field)])
179 .encode_values(&[(column.column_id, default)], &mut encoded)
180 .ok()?;
181 Some(encoded.into())
182}
183
184fn append_suffix(key: Bytes, suffix: &Bytes) -> Bytes {
185 let mut completed = Vec::with_capacity(key.len() + suffix.len());
186 completed.extend_from_slice(&key);
187 completed.extend_from_slice(suffix);
188 completed.into()
189}
190
191#[cfg(test)]
192mod tests {
193 use std::sync::Arc;
194
195 use api::v1::SemanticType;
196 use common_time::Timestamp;
197 use datatypes::prelude::{ConcreteDataType, Value};
198 use datatypes::schema::{ColumnDefaultConstraint, ColumnSchema};
199 use rstest::rstest;
200 use store_api::metadata::{ColumnMetadata, RegionMetadataBuilder};
201 use store_api::storage::FileId;
202
203 use super::*;
204 use crate::compaction::run::files_overlap_inclusive;
205 use crate::manifest::action::RegionEdit;
206 use crate::memtable::time_partition::TimePartitions;
207 use crate::memtable::time_series::TimeSeriesMemtableBuilder;
208 use crate::region::version::{VersionBuilder, VersionRef};
209 use crate::sst::file::{FileHandle, FileMeta};
210 use crate::test_util::new_noop_file_purger;
211
212 fn metadata(defaults: &[Value]) -> RegionMetadataRef {
213 let mut builder = RegionMetadataBuilder::new(RegionId::new(1, 1));
214 for (id, value) in defaults.iter().enumerate() {
215 let data_type = if value.is_null() {
216 ConcreteDataType::string_datatype()
217 } else {
218 value.data_type()
219 };
220 builder.push_column_metadata(ColumnMetadata {
221 column_id: id as u32,
222 semantic_type: SemanticType::Tag,
223 column_schema: ColumnSchema::new(format!("tag_{id}"), data_type, true)
224 .with_default_constraint(Some(ColumnDefaultConstraint::Value(value.clone())))
225 .unwrap(),
226 });
227 }
228 builder.push_column_metadata(ColumnMetadata {
229 column_id: 100,
230 semantic_type: SemanticType::Timestamp,
231 column_schema: ColumnSchema::new(
232 "ts",
233 ConcreteDataType::timestamp_millisecond_datatype(),
234 false,
235 ),
236 });
237 builder.primary_key((0..defaults.len() as u32).collect());
238 Arc::new(builder.build().unwrap())
239 }
240
241 fn metadata_at_version(defaults: &[Value], version: u64) -> RegionMetadataRef {
242 let mut metadata = (*metadata(defaults)).clone();
243 metadata.schema_version = version;
244 Arc::new(metadata)
245 }
246
247 fn encode(metadata: &RegionMetadataRef, values: &[Value]) -> Bytes {
248 let mut bytes = Vec::new();
249 let values: Vec<_> = values
250 .iter()
251 .enumerate()
252 .map(|(id, value)| (id as u32, value.clone()))
253 .collect();
254 DensePrimaryKeyCodec::new(metadata)
255 .encode_values(&values, &mut bytes)
256 .unwrap();
257 bytes.into()
258 }
259
260 fn file_meta(metadata: &RegionMetadataRef, values: &[Value]) -> FileMeta {
261 let key = encode(metadata, values);
262 FileMeta {
263 region_id: metadata.region_id,
264 file_id: FileId::random(),
265 time_range: (Timestamp::new_millisecond(0), Timestamp::new_millisecond(1)),
266 primary_key_min: Some(key.clone()),
267 primary_key_max: Some(key),
268 ..Default::default()
269 }
270 }
271
272 fn version_with_files(metadata: RegionMetadataRef, files: Vec<FileMeta>) -> VersionRef {
273 let mutable = Arc::new(TimePartitions::new(
274 metadata.clone(),
275 Arc::new(TimeSeriesMemtableBuilder::default()),
276 0,
277 None,
278 ));
279 Arc::new(
280 VersionBuilder::new(metadata, mutable)
281 .add_files(new_noop_file_purger(), files.into_iter())
282 .build(),
283 )
284 }
285
286 #[test]
287 fn test_field_addition_updates_cache_version_and_shares_sst_list() {
288 let metadata = metadata(&[Value::from(""), Value::from("default")]);
289 let raw = file_meta(&metadata, &[Value::from("a")]);
290 let file_id = raw.file_id;
291 let original = version_with_files(metadata.clone(), vec![raw]);
292 let mapper = original.ssts.primary_key_mapper();
293 let range = original.ssts.levels()[0].files[&file_id]
294 .primary_key_range(&mapper)
295 .unwrap();
296 let mut builder = RegionMetadataBuilder::from_existing((*metadata).clone());
297 builder
298 .push_column_metadata(ColumnMetadata {
299 column_id: 101,
300 semantic_type: SemanticType::Field,
301 column_schema: ColumnSchema::new("field", ConcreteDataType::int64_datatype(), true),
302 })
303 .bump_version();
304 let changed = VersionBuilder::from_version(original.clone())
305 .metadata(Arc::new(builder.build().unwrap()))
306 .build();
307
308 assert_eq!(1, changed.metadata.schema_version);
309 assert!(changed.metadata.column_by_id(101).is_some());
310 assert_eq!(
311 original.ssts.levels().as_ptr(),
312 changed.ssts.levels().as_ptr()
313 );
314 let changed_mapper = changed.ssts.primary_key_mapper();
315 assert_eq!(1, changed_mapper.schema_version());
316 let changed_range = changed.ssts.levels()[0].files[&file_id]
317 .primary_key_range(&changed_mapper)
318 .unwrap();
319 assert_eq!(range, changed_range);
320 let cached = changed.ssts.levels()[0].files[&file_id]
321 .primary_key_range(&changed_mapper)
322 .unwrap();
323 assert_eq!(cached.0.as_ptr(), changed_range.0.as_ptr());
324 assert_eq!(cached.1.as_ptr(), changed_range.1.as_ptr());
325 }
326
327 #[rstest]
329 #[case::set(Value::from("x"), Some(Value::from("y")))]
330 #[case::set_null(Value::from("x"), Some(Value::Null))]
331 #[case::drop(Value::from("x"), None)]
332 #[case::replace_null(Value::Null, Some(Value::from("y")))]
333 fn test_tag_default_changes_refresh_only_missing_values(
334 #[case] previous_default: Value,
335 #[case] new_default: Option<Value>,
336 ) {
337 let metadata = metadata(&[Value::from(""), previous_default.clone()]);
338 let missing = file_meta(&metadata, &[Value::from("a")]);
339 let complete = file_meta(&metadata, &[Value::from("b"), Value::from("stored")]);
340 let missing_id = missing.file_id;
341 let complete_id = complete.file_id;
342 let complete_range = complete.primary_key_range();
343 let original = version_with_files(metadata.clone(), vec![missing, complete]);
344 let old_mapper = original.ssts.primary_key_mapper();
345 let old_range = original.ssts.levels()[0].files[&missing_id].primary_key_range(&old_mapper);
346
347 let mut changed_metadata = (*metadata).clone();
348 changed_metadata.column_metadatas[1].column_schema = changed_metadata.column_metadatas[1]
349 .column_schema
350 .clone()
351 .with_default_constraint(new_default.clone().map(ColumnDefaultConstraint::Value))
352 .unwrap();
353 let mut builder = RegionMetadataBuilder::from_existing(changed_metadata);
354 builder.bump_version();
355 let changed_metadata = Arc::new(builder.build().unwrap());
356 let changed = VersionBuilder::from_version(original.clone())
357 .metadata(changed_metadata.clone())
358 .build();
359
360 let expected = encode(
361 &changed_metadata,
362 &[Value::from("a"), new_default.unwrap_or(Value::Null)],
363 );
364 let changed_mapper = changed.ssts.primary_key_mapper();
365 assert_eq!(
366 Some((expected.clone(), expected)),
367 changed.ssts.levels()[0].files[&missing_id].primary_key_range(&changed_mapper)
368 );
369 assert_eq!(
370 complete_range,
371 changed.ssts.levels()[0].files[&complete_id].primary_key_range(&changed_mapper)
372 );
373 let expected_old = encode(&metadata, &[Value::from("a"), previous_default]);
374 assert_eq!(Some((expected_old.clone(), expected_old)), old_range);
375 assert_eq!(
376 old_range,
377 original.ssts.levels()[0].files[&missing_id].primary_key_range(&old_mapper)
378 );
379 }
380
381 #[rstest]
382 #[case::mixed(vec![
383 Value::from("default"),
384 Value::Null,
385 Value::Int64(42),
386 Value::Binary(vec![0, 255].into()),
387 ])]
388 #[case::nulls(vec![Value::Null; 4])]
389 #[case::empty(vec![])]
390 fn test_dense_ranges_complete_every_historical_prefix(#[case] defaults: Vec<Value>) {
391 let metadata = metadata(&defaults);
392 let mapper = PrimaryKeyRangeMapper::new(metadata.clone());
393 let expected = encode(&metadata, &defaults);
394 for count in 0..=defaults.len() {
395 let prefix = encode(&metadata, &defaults[..count]);
396 assert_eq!(
397 Some((expected.clone(), expected.clone())),
398 mapper.map((prefix.clone(), prefix)).unwrap()
399 );
400 }
401 }
402
403 #[rstest]
404 #[case::cold(false)]
405 #[case::warm(true)]
406 fn test_successive_pk_appends_reuse_constants_and_complete_all_prefixes(#[case] warm: bool) {
407 let defaults = [
408 Value::from(""),
409 Value::from("constant"),
410 Value::Null,
411 Value::Int64(42),
412 Value::Null,
413 ];
414 let mut mapper = PrimaryKeyRangeMapper::new(metadata(&defaults[..2]));
415 let mut values = defaults.clone();
416 values[0] = Value::from("stored");
417 let raw = encode(&mapper.metadata, &values[..1]);
418 if warm {
419 assert!(mapper.map((raw.clone(), raw)).unwrap().is_some());
420 }
421
422 for count in 3..=defaults.len() {
423 let next_metadata = metadata(&defaults[..count]);
424 let next_mapper = mapper.with_metadata(next_metadata.clone());
425 if let Some(Some(encoded)) = mapper.encoded_defaults[1].get() {
426 let reused = next_mapper.encoded_defaults[1]
428 .get()
429 .unwrap()
430 .as_ref()
431 .unwrap();
432 assert_eq!(encoded.as_ptr(), reused.as_ptr());
433 }
434 let expected = encode(&next_metadata, &values[..count]);
435 for prefix_len in 1..=count {
436 let prefix = encode(&next_metadata, &values[..prefix_len]);
437 assert_eq!(
438 Some((expected.clone(), expected.clone())),
439 next_mapper.map((prefix.clone(), prefix)).unwrap()
440 );
441 }
442 let old_expected = encode(&mapper.metadata, &values[..count - 1]);
443 let raw = encode(&mapper.metadata, &values[..1]);
444 assert_eq!(
445 Some((old_expected.clone(), old_expected)),
446 mapper.map((raw.clone(), raw)).unwrap()
447 );
448 mapper = next_mapper;
449 }
450 }
451
452 #[rstest]
453 #[case::string(ConcreteDataType::string_datatype(), true, true)]
454 #[case::int(ConcreteDataType::int64_datatype(), true, true)]
455 #[case::binary(ConcreteDataType::binary_datatype(), true, true)]
456 #[case::timestamp(ConcreteDataType::timestamp_millisecond_datatype(), true, true)]
457 #[case::required(ConcreteDataType::string_datatype(), false, false)]
458 #[case::unsupported_list(
459 ConcreteDataType::list_datatype(Arc::new(ConcreteDataType::string_datatype())),
460 true,
461 false
462 )]
463 #[case::unsupported_null(ConcreteDataType::null_datatype(), true, false)]
464 fn test_null_padding_requires_a_nullable_encodable_column(
465 #[case] data_type: ConcreteDataType,
466 #[case] nullable: bool,
467 #[case] known: bool,
468 ) {
469 let mut metadata = (*metadata(&[Value::Null])).clone();
470 metadata.column_metadatas[0].column_schema =
471 ColumnSchema::new("tag_0", data_type, nullable);
472 let metadata = Arc::new(
473 RegionMetadataBuilder::from_existing(metadata)
474 .build()
475 .unwrap(),
476 );
477 let mapper = PrimaryKeyRangeMapper::new(metadata.clone());
478 let expected = known.then(|| (Bytes::from_static(&[0]), Bytes::from_static(&[0])));
479 assert_eq!(expected, mapper.map((Bytes::new(), Bytes::new())).unwrap());
480 if known {
481 let encoded_null = encode(&metadata, &[Value::Null]);
482 assert_eq!(Some((encoded_null.clone(), encoded_null)), expected);
483 }
484 }
485
486 #[rstest]
489 #[case::dense(PrimaryKeyEncoding::Dense)]
490 #[case::sparse(PrimaryKeyEncoding::Sparse)]
491 fn test_reversed_pk_bounds_return_error(#[case] encoding: PrimaryKeyEncoding) {
492 use common_error::ext::ErrorExt;
493 use common_error::status_code::StatusCode;
494 use mito_codec::row_converter::build_primary_key_codec;
495
496 let mut metadata = (*metadata(&[Value::from("")])).clone();
497 metadata.primary_key_encoding = encoding;
498 let metadata = Arc::new(metadata);
499 let codec = build_primary_key_codec(&metadata);
500 let mut a = Vec::new();
501 let mut b = Vec::new();
502 codec
503 .encode_values(&[(0, Value::from("a"))], &mut a)
504 .unwrap();
505 codec
506 .encode_values(&[(0, Value::from("b"))], &mut b)
507 .unwrap();
508 let mapper = PrimaryKeyRangeMapper::new(metadata);
509 let err = mapper.map((b.into(), a.into())).unwrap_err();
510 assert!(matches!(
511 err,
512 crate::error::Error::InvalidPrimaryKeyRange { .. }
513 ));
514 assert_eq!(StatusCode::Internal, err.status_code());
515 }
516
517 #[test]
518 fn test_different_pk_prefix_lengths_return_error() {
519 let metadata = metadata(&[Value::from(""), Value::Int64(42)]);
520 let mapper = PrimaryKeyRangeMapper::new(metadata.clone());
521 let min = encode(&metadata, &[Value::from("a")]);
522 let max = encode(&metadata, &[Value::from("b"), Value::Int64(42)]);
523 assert!(matches!(
524 mapper.map((min, max)),
525 Err(crate::error::Error::InvalidPrimaryKeyRange { .. })
526 ));
527 }
528
529 #[rstest]
530 #[case::min("min")]
531 #[case::max("max")]
532 fn test_truncated_pk_endpoint_preserves_decode_error(#[case] endpoint: &'static str) {
533 use common_error::ext::ErrorExt;
534 use common_error::status_code::StatusCode;
535
536 let metadata = metadata(&[Value::from("")]);
537 let mapper = PrimaryKeyRangeMapper::new(metadata.clone());
538 let mut min = encode(&metadata, &[Value::from("a")]);
539 let mut max = encode(&metadata, &[Value::from("b")]);
540 if endpoint == "min" {
541 min.truncate(min.len() - 1);
542 } else {
543 max.truncate(max.len() - 1);
544 }
545 let err = mapper.map((min, max)).unwrap_err();
546 assert_eq!(StatusCode::Internal, err.status_code());
547 assert!(matches!(err, crate::error::Error::DecodePrimaryKeyRange {
548 endpoint: actual,
549 source: mito_codec::error::Error::InvalidDensePrimaryKey { .. },
550 ..
551 } if actual == endpoint));
552 }
553
554 #[rstest]
555 #[case::source_first(false)]
556 #[case::destination_first(true)]
557 fn test_aligned_cache_uses_table_schema_across_region_migration(
558 #[case] destination_first: bool,
559 ) {
560 let source = metadata(&[Value::from("")]);
561 let file = FileHandle::new(
562 file_meta(&source, &[Value::from("a")]),
563 new_noop_file_purger(),
564 );
565 let target = metadata_at_version(&[Value::from(""), Value::Int64(42)], 1);
566 let expected = encode(&target, &[Value::from("a"), Value::Int64(42)]);
567 let mut regions = [RegionId::new(1, 1), RegionId::new(1, 2)];
568 if destination_first {
569 regions.reverse();
570 }
571 let mut previous: Option<(Bytes, Bytes)> = None;
572 for region_id in regions {
573 let mut metadata = (*target).clone();
574 metadata.region_id = region_id;
575 let mapper = PrimaryKeyRangeMapper::new(Arc::new(metadata));
576 let aligned = file.primary_key_range(&mapper).unwrap();
577 assert_eq!((expected.clone(), expected.clone()), aligned);
578 if let Some(previous) = previous {
579 assert_eq!(previous.0.as_ptr(), aligned.0.as_ptr());
581 assert_eq!(previous.1.as_ptr(), aligned.1.as_ptr());
582 }
583 previous = Some(aligned);
584 }
585 }
586
587 #[test]
588 fn test_invalid_file_ranges_log_once_and_keep_possible_overlap() {
589 use std::sync::atomic::{AtomicUsize, Ordering};
590
591 use common_telemetry::tracing_subscriber::Layer;
592 use common_telemetry::tracing_subscriber::layer::Context;
593 use common_telemetry::tracing_subscriber::prelude::*;
594 use tracing::{Event, Level, Subscriber};
595
596 struct WarningCounter(Arc<AtomicUsize>);
597 impl<S: Subscriber> Layer<S> for WarningCounter {
598 fn on_event(&self, event: &Event<'_>, _ctx: Context<'_, S>) {
599 if *event.metadata().level() == Level::WARN {
600 self.0.fetch_add(1, Ordering::Relaxed);
601 }
602 }
603 }
604
605 let warnings = Arc::new(AtomicUsize::new(0));
606 let subscriber =
607 common_telemetry::tracing_subscriber::registry().with(WarningCounter(warnings.clone()));
608 let metadata = metadata(&[Value::from("")]);
609 let mapper = Arc::new(PrimaryKeyRangeMapper::new(metadata.clone()));
610 let healthy_meta = file_meta(&metadata, &[Value::from("c")]);
611 let healthy = FileHandle::new(healthy_meta, new_noop_file_purger());
612 let a = encode(&metadata, &[Value::from("a")]);
613 let b = encode(&metadata, &[Value::from("b")]);
614
615 tracing::subscriber::with_default(subscriber, || {
616 for (min, max) in [(b.clone(), a.clone()), (a.slice(..a.len() - 1), b)] {
617 let mut meta = file_meta(&metadata, &[Value::from("a")]);
618 meta.primary_key_min = Some(min);
619 meta.primary_key_max = Some(max);
620 let file = FileHandle::new(meta, new_noop_file_purger());
621 assert_eq!(None, file.primary_key_range(&mapper));
622 assert_eq!(None, file.clone().primary_key_range(&mapper));
623 assert!(files_overlap_inclusive(&file, &healthy, &mapper));
625 assert!(files_overlap_inclusive(&healthy, &file, &mapper));
626 }
627 });
628 assert_eq!(2, warnings.load(Ordering::Relaxed));
629 }
630
631 #[test]
632 fn test_missing_impure_default_is_unknown_but_complete_key_is_usable() {
633 let mut metadata = (*metadata(&[Value::Timestamp(Timestamp::new_millisecond(0))])).clone();
634 metadata.column_metadatas[0].column_schema = metadata.column_metadatas[0]
635 .column_schema
636 .clone()
637 .with_default_constraint(Some(ColumnDefaultConstraint::Function(
638 "current_timestamp()".into(),
639 )))
640 .unwrap();
641 let metadata = Arc::new(metadata);
642 let mapper = PrimaryKeyRangeMapper::new(metadata.clone());
643 assert_eq!(None, mapper.map((Bytes::new(), Bytes::new())).unwrap());
644 let key = encode(
645 &metadata,
646 &[Value::Timestamp(Timestamp::new_millisecond(10))],
647 );
648 assert_eq!(
649 Some((key.clone(), key.clone())),
650 mapper.map((key.clone(), key)).unwrap()
651 );
652 }
653
654 #[test]
655 fn test_sparse_ranges_keep_their_original_encoding() {
656 use mito_codec::row_converter::SparsePrimaryKeyCodec;
657 let mut metadata = (*metadata(&[Value::from(""), Value::Null])).clone();
658 metadata.primary_key_encoding = PrimaryKeyEncoding::Sparse;
659 let metadata = Arc::new(metadata);
660 let mut bytes = Vec::new();
661 SparsePrimaryKeyCodec::new(&metadata)
662 .encode_values(&[(0, Value::from("a"))], &mut bytes)
663 .unwrap();
664 let key = Bytes::from(bytes);
665 let mapper = PrimaryKeyRangeMapper::new(metadata);
666 assert_eq!(
667 Some((key.clone(), key.clone())),
668 mapper.map((key.clone(), key)).unwrap()
669 );
670 }
671
672 #[test]
673 fn test_aligned_cache_isolates_schemas_and_shares_file_state() {
674 let old_metadata = metadata(&[Value::from("")]);
675 let new_metadata =
676 metadata_at_version(&[Value::from(""), Value::Null, Value::Int64(42)], 1);
677 let raw_meta = file_meta(&old_metadata, &[Value::from("b")]);
678 let original = version_with_files(old_metadata.clone(), vec![raw_meta.clone()]);
679 let old_file = original.ssts.levels()[0].files().next().unwrap();
680 let old_mapper = original.ssts.primary_key_mapper();
681 let changed = VersionBuilder::from_version(original.clone())
682 .metadata(new_metadata.clone())
683 .build();
684 let new_file = changed.ssts.levels()[0].files().next().unwrap();
685 let new_mapper = changed.ssts.primary_key_mapper();
686 let expected = encode(
687 &new_metadata,
688 &[Value::from("b"), Value::Null, Value::Int64(42)],
689 );
690 assert_eq!(
691 raw_meta.primary_key_range(),
692 old_file.primary_key_range(&old_mapper)
693 );
694 assert_eq!(
695 Some((expected.clone(), expected)),
696 new_file.primary_key_range(&new_mapper)
697 );
698 assert_eq!(raw_meta, *new_file.meta_ref());
699 assert_eq!(
700 raw_meta.primary_key_range(),
701 new_file.raw_primary_key_range()
702 );
703 assert_eq!(
704 old_file.primary_key_range(&old_mapper),
705 new_file.primary_key_range(&old_mapper)
706 );
707 assert_eq!(
708 new_file.primary_key_range(&new_mapper),
709 old_file.primary_key_range(&new_mapper)
710 );
711 old_file.set_compacting(true);
712 assert!(new_file.compacting());
713 new_file.mark_deleted();
714 assert!(old_file.is_deleted());
715
716 let added_meta = file_meta(&old_metadata, &[Value::from("a")]);
718 let added_id = added_meta.file_id;
719 let changed = VersionBuilder::from_version(Arc::new(changed))
720 .apply_edit(
721 RegionEdit {
722 files_to_add: vec![added_meta],
723 files_to_remove: vec![],
724 timestamp_ms: None,
725 compaction_time_window: None,
726 flushed_entry_id: None,
727 flushed_sequence: None,
728 committed_sequence: None,
729 },
730 new_noop_file_purger(),
731 )
732 .build();
733 let added = &changed.ssts.levels()[0].files[&added_id];
734 let expected = encode(
735 &new_metadata,
736 &[Value::from("a"), Value::Null, Value::Int64(42)],
737 );
738 assert_eq!(
739 Some((expected.clone(), expected)),
740 added.primary_key_range(&new_mapper)
741 );
742 }
743
744 #[test]
745 fn test_late_statistics_and_changed_defaults_use_each_pinned_schema() {
746 let metadata_v1 = metadata(&[Value::from(""), Value::Int64(42)]);
747 let metadata_v2 = metadata_at_version(&[Value::from(""), Value::Int64(100)], 1);
748 let mut meta = file_meta(&metadata_v1, &[Value::from("a")]);
749 let raw = meta.primary_key_range().unwrap();
750 meta.primary_key_min = None;
751 meta.primary_key_max = None;
752 let file = FileHandle::new(meta, new_noop_file_purger());
753 let v1 = PrimaryKeyRangeMapper::new(metadata_v1.clone());
754 let v2 = PrimaryKeyRangeMapper::new(metadata_v2.clone());
755 assert_eq!(None, file.primary_key_range(&v1));
756 assert_eq!(None, file.primary_key_range(&v2));
757 file.set_primary_key_range(raw.clone());
758 for (mapper, metadata, default) in [
759 (&v2, &metadata_v2, 100),
760 (&v1, &metadata_v1, 42),
761 (&v2, &metadata_v2, 100),
762 ] {
763 let expected = encode(metadata, &[Value::from("a"), Value::Int64(default)]);
764 assert_eq!(
765 Some((expected.clone(), expected)),
766 file.primary_key_range(mapper)
767 );
768 }
769 assert_eq!(Some(raw), file.raw_primary_key_range());
770 }
771
772 #[test]
773 fn test_concurrent_snapshots_do_not_return_each_others_aligned_bounds() {
774 let old_metadata = metadata(&[Value::from(""), Value::Int64(42)]);
775 let new_metadata = metadata_at_version(&[Value::from(""), Value::Int64(100)], 1);
776 let raw = file_meta(&old_metadata, &[Value::from("a")]);
777 let original_bounds = raw.primary_key_range();
778 let file = FileHandle::new(raw, new_noop_file_purger());
779 let barrier = std::sync::Barrier::new(2);
780 std::thread::scope(|scope| {
781 for (metadata, default) in [(old_metadata, 42), (new_metadata, 100)] {
782 let file = &file;
783 let barrier = &barrier;
784 scope.spawn(move || {
785 let expected = encode(&metadata, &[Value::from("a"), Value::Int64(default)]);
786 let mapper = PrimaryKeyRangeMapper::new(metadata);
787 let mut all_matched = true;
788 for _ in 0..64 {
789 barrier.wait();
790 all_matched &= file.primary_key_range(&mapper)
791 == Some((expected.clone(), expected.clone()));
792 }
793 assert!(
795 all_matched,
796 "wrong bounds for schema version {}",
797 mapper.schema_version()
798 );
799 });
800 }
801 });
802 assert_eq!(original_bounds, file.raw_primary_key_range());
803 }
804
805 #[test]
806 fn test_deserialized_compaction_files_use_the_comparison_schema() {
807 use crate::compaction::CompactionOutput;
808 use crate::compaction::picker::{PickerOutput, SerializedPickerOutput};
809
810 let metadata = metadata(&[Value::from(""), Value::Int64(42)]);
811 let raw = file_meta(&metadata, &[Value::from("a")]);
812 let picked = PickerOutput {
813 outputs: vec![CompactionOutput {
814 output_level: 1,
815 inputs: vec![FileHandle::new(raw.clone(), new_noop_file_purger())],
816 filter_deleted: false,
817 output_time_range: None,
818 }],
819 ..Default::default()
820 };
821 let serialized = SerializedPickerOutput::from(&picked);
822 let output = PickerOutput::from_serialized(serialized, new_noop_file_purger());
823 let mapper = PrimaryKeyRangeMapper::new(metadata.clone());
824 let file = &output.outputs[0].inputs[0];
825 let completed = encode(&metadata, &[Value::from("a"), Value::Int64(42)]);
826 assert_eq!(
827 Some((completed.clone(), completed)),
828 file.primary_key_range(&mapper)
829 );
830 assert_eq!(raw, *file.meta_ref());
831 }
832
833 #[test]
834 fn test_compaction_overlap_uses_completed_endpoints() {
835 let metadata = metadata(&[Value::from(""), Value::Null]);
836 let mapper = Arc::new(PrimaryKeyRangeMapper::new(metadata.clone()));
837 let mut old = file_meta(&metadata, &[Value::from("a")]);
838 old.primary_key_max = Some(encode(&metadata, &[Value::from("b")]));
839 let mut new = file_meta(&metadata, &[Value::from("b"), Value::Null]);
840 new.primary_key_max = Some(encode(&metadata, &[Value::from("c"), Value::Null]));
841 assert!(old.primary_key_max < new.primary_key_min);
842 let old = FileHandle::new(old, new_noop_file_purger());
843 let new = FileHandle::new(new, new_noop_file_purger());
844 assert!(files_overlap_inclusive(&old, &new, &mapper));
845 assert!(files_overlap_inclusive(&new, &old, &mapper));
846 }
847}