use arrow::array::{RecordBatch, RecordBatchOptions, UInt32Array};
use futures::StreamExt as _;
use re_log_types::{AbsoluteTimeRange, TimeInt, TimeType};
use re_protos::cloud::v1alpha1::ext::{
DataSource, Query, QueryDatasetDataframe, QueryDatasetRequest, QueryLatestAt, QueryRange,
};
use re_protos::cloud::v1alpha1::rerun_cloud_service_server::RerunCloudService;
use re_protos::headers::RerunHeadersInjectorExt as _;
use re_types_core::ChunkId;
use crate::tests::common::{
DataSourcesDefinition, LayerDefinition, RerunCloudServiceExt as _, concat_record_batches,
entry_name,
};
use crate::{FieldsTestExt as _, RecordBatchTestExt as _, TempPath};
pub async fn query_empty_dataset(service: impl RerunCloudService) {
let dataset_name = "dataset";
service.create_dataset_entry_with_name(dataset_name).await;
query_dataset_snapshot(
&service,
QueryDatasetRequest::default(),
&[],
dataset_name,
"empty_dataset",
)
.await;
}
pub async fn query_simple_dataset(service: impl RerunCloudService) {
let data_sources_def = DataSourcesDefinition::new_with_tuid_prefix(
1,
[
LayerDefinition::simple("my_segment_id1", &["my/entity", "my/other/entity"]),
LayerDefinition::simple("my_segment_id2", &["my/entity"]),
LayerDefinition::simple(
"my_segment_id3",
&["my/entity", "another/one", "yet/another/one"],
),
],
);
let dataset_name = "dataset";
service.create_dataset_entry_with_name(dataset_name).await;
service
.register_with_dataset_name_blocking(dataset_name, data_sources_def.to_data_sources())
.await;
let requests = vec![
(QueryDatasetRequest::default(), "default"),
(
QueryDatasetRequest {
segment_ids: vec!["my_segment_id3".into()],
..Default::default()
},
"single_segment",
),
(
QueryDatasetRequest {
entity_paths: vec!["/my/entity".into()],
select_all_entity_paths: false,
..Default::default()
},
"single_entity",
),
(
QueryDatasetRequest {
exclude_static_data: true,
..Default::default()
},
"exclude_static",
),
(
QueryDatasetRequest {
exclude_temporal_data: true,
..Default::default()
},
"exclude_temporal",
),
];
for (request, snapshot_name) in requests {
query_dataset_snapshot(
&service,
request,
&[],
dataset_name,
&format!("simple_dataset_{snapshot_name}"),
)
.await;
}
}
pub async fn query_simple_dataset_with_layers(service: impl RerunCloudService) {
let data_sources_def = DataSourcesDefinition::new_with_tuid_prefix(
1,
[
LayerDefinition::simple("partition1", &["my/entity"]),
LayerDefinition::simple("partition1", &["extra/entity"]).layer_name("extra"),
LayerDefinition::simple("partition2", &["another/one"]).layer_name("base"),
LayerDefinition::simple("partition2", &["extra/entity"]).layer_name("extra"),
LayerDefinition::simple("partition3", &["i/am/alone"]),
],
);
let dataset_name = "dataset_with_layers";
service.create_dataset_entry_with_name(dataset_name).await;
service
.register_with_dataset_name_blocking(dataset_name, data_sources_def.to_data_sources())
.await;
query_dataset_snapshot(
&service,
QueryDatasetRequest::default(),
&[],
dataset_name,
"simple_with_layer",
)
.await;
}
pub async fn query_dataset_unknown_segment_id_returns_empty(service: impl RerunCloudService) {
let data_sources_def = DataSourcesDefinition::new_with_tuid_prefix(
1,
[LayerDefinition::simple("real_segment", &["my/entity"])],
);
let dataset_name = "dataset";
service.create_dataset_entry_with_name(dataset_name).await;
service
.register_with_dataset_name_blocking(dataset_name, data_sources_def.to_data_sources())
.await;
let test_cases = [
("only unknown", vec!["doesnt_exist".into()], false),
(
"mixed real + unknown",
vec!["real_segment".into(), "doesnt_exist".into()],
true,
),
];
for (descr, segment_ids, expect_rows) in test_cases {
let chunk_info: Vec<RecordBatch> = service
.query_dataset(
tonic::Request::new(
QueryDatasetRequest {
segment_ids,
..Default::default()
}
.into(),
)
.with_entry_name(entry_name(dataset_name)),
)
.await
.unwrap_or_else(|err| panic!("query_dataset must succeed ({descr}): {err}"))
.into_inner()
.flat_map(|resp| {
futures::stream::iter(
resp.unwrap_or_else(|err| {
panic!("query_dataset stream must not error ({descr}): {err}")
})
.data,
)
})
.map(|dfp| dfp.try_into().unwrap())
.collect()
.await;
let merged = concat_record_batches(&chunk_info);
let has_rows = merged.num_rows() > 0;
assert_eq!(
has_rows,
expect_rows,
"unexpected row presence for {descr}: got {} rows",
merged.num_rows(),
);
}
}
pub async fn query_dataset_should_fail(service: impl RerunCloudService) {
let dataset_name = "dataset";
service.create_dataset_entry_with_name(dataset_name).await;
let test_cases = vec![
(
"cannot specify entity paths if `select_all_entity_paths` is true",
QueryDatasetRequest {
entity_paths: vec!["/entity/path".into()],
select_all_entity_paths: true,
..Default::default()
},
tonic::Code::InvalidArgument,
),
];
for (descr, request, expected_code) in test_cases {
let response = service
.query_dataset(tonic::Request::new(request.into()))
.await;
match response {
Ok(_) => {
panic!("expected failure with code {expected_code}, but got success ({descr})",);
}
Err(err) => {
assert_eq!(
err.code(),
expected_code,
"expected failure with code {expected_code}, but got {err} ({descr})"
);
}
}
}
}
fn create_recording_for_query_testing() -> anyhow::Result<TempPath> {
use re_chunk::{Chunk, TimePoint};
use re_log_types::example_components::{MyPoint, MyPoints};
use re_log_types::{EntityPath, TimeInt, build_frame_nr};
use re_sdk::RecordingStreamBuilder;
use crate::utils::rerun::{next_chunk_id_generator, next_row_id_generator};
let segment_id = "static_test_segment";
let tuid_prefix: u64 = 100;
let tmp_dir = tempfile::tempdir()?;
let tmp_path = tmp_dir.path().join(format!("{segment_id}.rrd"));
let rec = RecordingStreamBuilder::new(format!("rerun_example_{segment_id}"))
.recording_id(segment_id)
.send_properties(false)
.save(tmp_path.clone())?;
let mut next_chunk_id = next_chunk_id_generator(tuid_prefix);
let mut next_row_id = next_row_id_generator(tuid_prefix);
let frame0 = TimeInt::new_temporal(0);
let points = MyPoint::from_iter(0..1);
let static_only_chunk =
Chunk::builder_with_id(next_chunk_id(), EntityPath::from("/static_only"))
.with_sparse_component_batches(
next_row_id(),
TimePoint::default(),
[(MyPoints::descriptor_points(), Some(&points as _))],
)
.build()?;
rec.send_chunk(static_only_chunk);
let both_static_chunk = Chunk::builder_with_id(next_chunk_id(), EntityPath::from("/both"))
.with_sparse_component_batches(
next_row_id(),
TimePoint::default(),
[(MyPoints::descriptor_points(), Some(&points as _))],
)
.build()?;
rec.send_chunk(both_static_chunk);
let both_temporal_chunk = Chunk::builder_with_id(next_chunk_id(), EntityPath::from("/both"))
.with_sparse_component_batches(
next_row_id(),
[build_frame_nr(frame0)],
[(MyPoints::descriptor_points(), Some(&points as _))],
)
.build()?;
rec.send_chunk(both_temporal_chunk);
let temporal_only_chunk =
Chunk::builder_with_id(next_chunk_id(), EntityPath::from("/temporal_only"))
.with_sparse_component_batches(
next_row_id(),
[build_frame_nr(frame0)],
[(MyPoints::descriptor_points(), Some(&points as _))],
)
.build()?;
rec.send_chunk(temporal_only_chunk);
rec.flush_blocking()?;
Ok(crate::TempPath::new(tmp_dir, tmp_path))
}
pub async fn query_dataset_with_various_queries(service: impl RerunCloudService) {
let recording_path = create_recording_for_query_testing().unwrap();
let dataset_name = "dataset_with_layers";
service.create_dataset_entry_with_name(dataset_name).await;
service
.register_with_dataset_name_blocking(
dataset_name,
vec![
DataSource::new_rrd_url(
url::Url::from_file_path(recording_path.as_path()).unwrap(),
)
.into(),
],
)
.await;
let queries = [
(None, vec![], "none"),
(Some(Query::default()), vec![], "default"),
(
Some(Query {
latest_at: Some(QueryLatestAt::global(Some("frame_nr".into()), TimeInt::MAX)),
range: None,
..Default::default()
}),
vec![ChunkId::from_tuid(re_tuid::Tuid::from_nanos_and_inc(
100, 3,
))],
"latest_at_end",
),
(
Some(Query {
latest_at: None,
range: Some(QueryRange {
index: "frame_nr".into(),
index_range: AbsoluteTimeRange {
min: TimeInt::MIN,
max: TimeInt::MAX,
},
}),
..Default::default()
}),
vec![ChunkId::from_tuid(re_tuid::Tuid::from_nanos_and_inc(
100, 3,
))],
"range_all",
),
];
for (query, chunk_ids_to_remove, snapshot_name) in queries {
query_dataset_snapshot(
&service,
QueryDatasetRequest {
segment_ids: vec![],
chunk_ids: vec![],
entity_paths: vec![],
select_all_entity_paths: true,
fuzzy_descriptors: vec![],
exclude_static_data: false,
exclude_temporal_data: false,
scan_parameters: None,
query,
generate_direct_urls: false,
},
&chunk_ids_to_remove,
dataset_name,
&format!("with_query_{snapshot_name}"),
)
.await;
}
}
pub async fn query_dataset_has_uncompressed_sizes(service: impl RerunCloudService) {
let data_sources_def = DataSourcesDefinition::new_with_tuid_prefix(
1,
[LayerDefinition::simple("segment", &["my/entity"])],
);
let dataset_name = "dataset";
service.create_dataset_entry_with_name(dataset_name).await;
service
.register_with_dataset_name_blocking(dataset_name, data_sources_def.to_data_sources())
.await;
let chunk_info: Vec<RecordBatch> = service
.query_dataset(
tonic::Request::new(QueryDatasetRequest::default().into())
.with_entry_name(entry_name(dataset_name)),
)
.await
.unwrap()
.into_inner()
.flat_map(|resp| futures::stream::iter(resp.unwrap().data))
.map(|dfp| dfp.try_into().unwrap())
.collect()
.await;
let merged = concat_record_batches(&chunk_info);
assert!(
merged.num_rows() > 0,
"query should return at least one chunk"
);
let uncompressed = QueryDatasetDataframe::COLUMN_CHUNK_BYTE_SIZE_UNCOMPRESSED
.extract(&merged)
.expect("chunk_byte_size_uncompressed column must be present");
for (i, size) in uncompressed.iter().enumerate() {
let size = size.unwrap_or_else(|| panic!("row {i}: uncompressed size must not be null"));
assert!(0 < size, "row {i}: uncompressed size must be > 0");
}
}
pub async fn query_dataset_consistent_schema_across_timelines(service: impl RerunCloudService) {
let data_sources_def = DataSourcesDefinition::new_with_tuid_prefix(
1,
[
LayerDefinition::simple_with_time(
"segment_sequence",
&["my/entity"],
0,
TimeType::Sequence,
),
LayerDefinition::simple_with_time(
"segment_timestamp",
&["my/entity"],
0,
TimeType::TimestampNs,
),
],
);
let dataset_name = "dataset_mixed_timelines";
service.create_dataset_entry_with_name(dataset_name).await;
service
.register_with_dataset_name_blocking(dataset_name, data_sources_def.to_data_sources())
.await;
let request = QueryDatasetRequest {
query: Some(Query {
columns_always_include_global_indexes: true,
..Default::default()
}),
..Default::default()
};
let responses: Vec<RecordBatch> = service
.query_dataset(
tonic::Request::new(request.into()).with_entry_name(entry_name(dataset_name)),
)
.await
.unwrap()
.into_inner()
.flat_map(|resp| futures::stream::iter(resp.unwrap().data))
.map(|dfp| dfp.try_into().unwrap())
.collect()
.await;
assert!(
!responses.is_empty(),
"expected at least one response, got none",
);
let first_schema = responses[0].schema();
for (idx, rb) in responses.iter().enumerate() {
assert_eq!(
rb.schema(),
first_schema,
"response {idx} has a different schema than response 0 — client-side concatenation would fail",
);
}
for expected_col in ["frame_nr:start", "timestamp:start"] {
assert!(
first_schema.field_with_name(expected_col).is_ok(),
"expected `{expected_col}` in response schema, got: {:#?}",
first_schema.fields(),
);
}
let _ = concat_record_batches(&responses);
}
async fn query_dataset_snapshot(
service: &impl RerunCloudService,
query_dataset_request: QueryDatasetRequest,
chunk_ids_to_remove: &[ChunkId],
dataset_name: &str,
snapshot_name: &str,
) {
let chunk_info = service
.query_dataset(
tonic::Request::new(query_dataset_request.into())
.with_entry_name(entry_name(dataset_name)),
)
.await
.unwrap()
.into_inner()
.flat_map(|resp| futures::stream::iter(resp.unwrap().data))
.map(|dfp| dfp.try_into().unwrap())
.collect::<Vec<_>>()
.await;
let merged_chunk_info = concat_record_batches(&chunk_info);
let merged_chunk_info =
remove_rows_containing_chunk_id(&merged_chunk_info, chunk_ids_to_remove);
let required_field = QueryDatasetDataframe::min_schema().fields().to_vec();
assert!(
merged_chunk_info
.schema()
.fields()
.contains_unordered(&required_field),
"query dataset must return all guaranteed fields\nExpected: {:#?}\nGot: {:#?}",
required_field,
merged_chunk_info.schema().fields(),
);
let required_column_names = required_field
.iter()
.map(|f| f.name().as_str())
.collect::<Vec<_>>();
let required_chunk_info = merged_chunk_info.project_columns(&required_column_names);
insta::assert_snapshot!(
format!("{snapshot_name}_schema"),
required_chunk_info.format_schema_snapshot()
);
let filtered_chunk_info = required_chunk_info
.remove_columns(&[
QueryDatasetDataframe::COLUMN_CHUNK_KEY_NAME,
QueryDatasetDataframe::COLUMN_CHUNK_BYTE_LEN_NAME,
QueryDatasetDataframe::COLUMN_CHUNK_BYTE_SIZE_UNCOMPRESSED_NAME,
])
.auto_sort_rows()
.unwrap();
insta::assert_snapshot!(
format!("{snapshot_name}_data"),
filtered_chunk_info.format_snapshot(false)
);
}
fn remove_rows_containing_chunk_id(
rb: &RecordBatch,
chunk_ids: &[re_types_core::ChunkId],
) -> RecordBatch {
let chunk_id_col = QueryDatasetDataframe::COLUMN_CHUNK_ID
.extract(rb)
.expect("bad chunk_id column");
let mut indices_to_keep = Vec::new();
for (row_idx, chunk_id) in chunk_id_col.iter_owned().enumerate() {
if !chunk_ids.contains(&chunk_id) {
indices_to_keep.push(row_idx as u32);
}
}
let indices = UInt32Array::from(indices_to_keep);
let resultant_rows = arrow::compute::take_arrays(rb.columns(), &indices, None)
.expect("take_arrays should return arrays");
RecordBatch::try_new_with_options(rb.schema(), resultant_rows, &RecordBatchOptions::default())
.expect("should create record batch")
}