use std::collections::{HashMap, HashSet};
use std::sync::Arc;
use std::sync::atomic::Ordering;
use std::time::Duration;
use arrow::array::{RecordBatch, StringArray};
use futures::StreamExt as _;
use re_dataframe::TimelineName;
use re_log_types::TimeInt;
use re_protos::cloud::v1alpha1::ext::QueryDatasetDataframe;
use re_redap_client::{ApiError, ApiResult};
use re_types_core::SegmentId;
use tokio::sync::mpsc::Sender;
use tracing::Instrument as _;
use crate::analytics::{QueryErrorKind, TaskFetchStats};
use crate::chunk_fetcher::{
ChunksWithSegment, SortedChunksWithSegment, batch_byte_size, batch_byte_size_uncompressed,
batch_has_any_direct_urls, fetch_batch_direct, fetch_batch_group_via_grpc,
split_batch_by_direct_url,
};
use crate::dataframe_query_common::{DataframeClientAPI, force_grpc};
use crate::metrics_capture::QueryMetrics;
use crate::pipeline_budget::{MAX_CONCURRENT_SEGMENTS, PipelineBudget};
use crate::segment_chunk_manifest::SegmentChunkManifest;
use re_dataframe::external::re_chunk::{Chunk, TimeColumn};
use super::cpu_worker::CpuWorkerMsg;
const TARGET_BATCH_SIZE_BYTES: usize = 8 * 1024 * 1024;
const GRPC_BATCH_SIZE: usize = 12;
const IO_PIPELINE_BUFFER: usize = 24;
const DIRECT_FETCH_CONNECT_TIMEOUT: Duration = Duration::from_secs(10);
const DIRECT_FETCH_READ_TIMEOUT: Duration = Duration::from_secs(30);
fn build_segment_manifests(
chunk_infos: &[RecordBatch],
filtered_timeline: Option<TimelineName>,
) -> ApiResult<HashMap<SegmentId, SegmentChunkManifest>> {
let mut manifests: HashMap<SegmentId, SegmentChunkManifest> = HashMap::new();
let Some(timeline_name) = filtered_timeline else {
return Ok(manifests);
};
let start_col_name = format!("{timeline_name}:start");
if chunk_infos
.iter()
.any(|rb| rb.column_by_name(&start_col_name).is_none())
{
return Ok(manifests);
}
for rb in chunk_infos {
let start_col = rb
.column_by_name(&start_col_name)
.expect("pre-check above guarantees presence on every batch");
let (start_values, start_nulls) = TimeColumn::read_nullable_array(start_col.as_ref())
.map_err(|err| {
ApiError::internal(format!(
"`{start_col_name}` column has unsupported type: {err}"
))
})?;
let segment_ids = QueryDatasetDataframe::COLUMN_CHUNK_SEGMENT_ID
.extract(rb)
.map_err(ApiError::internal_quiver)?;
let entity_paths = QueryDatasetDataframe::COLUMN_CHUNK_ENTITY_PATH
.extract(rb)
.map_err(ApiError::internal_quiver)?;
let is_statics = QueryDatasetDataframe::COLUMN_CHUNK_IS_STATIC
.extract(rb)
.map_err(ApiError::internal_quiver)?;
for i in 0..rb.num_rows() {
if start_nulls.as_ref().is_some_and(|n| n.is_null(i)) {
continue;
}
if is_statics.value(i) {
continue;
}
let time_min = TimeInt::saturated_temporal_i64(start_values[i]);
if time_min.is_static() {
continue;
}
let seg = segment_ids.value_owned(i);
let entity = entity_paths.value_owned(i);
manifests
.entry(seg)
.or_default()
.expect_chunk(entity, time_min);
}
}
for m in manifests.values_mut() {
m.lock();
}
Ok(manifests)
}
fn count_chunks_per_segment(chunk_infos: &[RecordBatch]) -> ApiResult<Vec<(SegmentId, usize)>> {
let mut counts: HashMap<SegmentId, usize> = HashMap::new();
let mut order: Vec<SegmentId> = Vec::new();
for rb in chunk_infos {
let segment_ids = QueryDatasetDataframe::COLUMN_CHUNK_SEGMENT_ID
.extract(rb)
.map_err(ApiError::internal_quiver)?;
for seg in &segment_ids {
if let Some(c) = counts.get_mut(seg) {
*c += 1;
} else {
let seg = SegmentId::from(seg);
order.push(seg.clone());
counts.insert(seg, 1);
}
}
}
Ok(order
.into_iter()
.map(|s| {
let c = counts.remove(&s).unwrap_or(0);
(s, c)
})
.collect())
}
fn extend_distinct_segment_ids(
batch: &RecordBatch,
seen: &mut HashSet<String>,
order: &mut Vec<String>,
) -> ApiResult<()> {
let seg_col = batch
.column_by_name(QueryDatasetDataframe::COLUMN_CHUNK_SEGMENT_ID_NAME)
.ok_or_else(|| ApiError::internal("missing segment_id column in fetch batch"))?
.as_any()
.downcast_ref::<StringArray>()
.ok_or_else(|| ApiError::internal("segment_id column is not a string array"))?;
for i in 0..batch.num_rows() {
let s = seg_col.value(i);
if !seen.contains(s) {
seen.insert(s.to_owned());
order.push(s.to_owned());
}
}
Ok(())
}
fn extract_task_time_min(batch: &RecordBatch, filtered_timeline: Option<&str>) -> TimeInt {
let Some(timeline) = filtered_timeline else {
return TimeInt::MAX;
};
let col_name = format!("{timeline}:start");
let Some(start_col) = batch.column_by_name(&col_name) else {
return TimeInt::MAX;
};
let Ok((start_values, start_nulls)) = TimeColumn::read_nullable_array(start_col.as_ref())
else {
return TimeInt::MAX;
};
let mut min_seen: Option<i64> = None;
for i in 0..start_values.len() {
if start_nulls.as_ref().is_some_and(|n| n.is_null(i)) {
continue;
}
let v = start_values[i];
min_seen = Some(min_seen.map_or(v, |m| m.min(v)));
}
min_seen.map_or(TimeInt::MAX, TimeInt::saturated_temporal_i64)
}
fn extract_segment_id(chunk_info: &RecordBatch) -> ApiResult<SegmentId> {
let segment_ids = QueryDatasetDataframe::COLUMN_CHUNK_SEGMENT_ID
.extract(chunk_info)
.map_err(ApiError::internal_quiver)?;
Ok(segment_ids.value_owned(0))
}
fn extract_chunk_sizes(chunk_info: &RecordBatch) -> ApiResult<quiver::Column<u64>> {
QueryDatasetDataframe::COLUMN_CHUNK_BYTE_LEN
.extract(chunk_info)
.map_err(ApiError::internal_quiver)
}
type BatchingResult = (Vec<RecordBatch>, Vec<SegmentId>);
enum FetchTask {
Direct(RecordBatch),
Grpc(RecordBatch),
}
fn split_batches_into_fetch_tasks(batches: &[RecordBatch]) -> (Vec<FetchTask>, usize, usize) {
let mut work_items: Vec<FetchTask> = Vec::new();
let mut n_direct = 0usize;
let mut n_grpc = 0usize;
for batch in batches {
let (direct_batch, grpc_batch) = split_batch_by_direct_url(batch);
if let Some(b) = direct_batch {
n_direct += 1;
work_items.push(FetchTask::Direct(b));
}
if let Some(b) = grpc_batch {
n_grpc += 1;
work_items.push(FetchTask::Grpc(b));
}
}
(work_items, n_direct, n_grpc)
}
#[tracing::instrument(
level = "info",
skip_all,
fields(
num_chunk_infos = chunk_infos.len(),
target_size_bytes,
output_batches,
byte_target_flushes,
segment_limit_flushes,
large_segment_batches,
end_of_input_batches,
)
)]
fn create_request_batches(
chunk_infos: Vec<RecordBatch>,
target_size_bytes: u64,
) -> ApiResult<BatchingResult> {
re_tracing::profile_function!();
let merge_err = |err: arrow::error::ArrowError, ctx: &'static str| {
ApiError::deserialization_with_source(None, err, ctx)
};
let mut request_batches = Vec::new();
let mut current_batch = Vec::new();
let mut current_batch_size = 0u64;
let mut current_batch_segments: HashSet<SegmentId> = HashSet::new();
let mut segment_order = Vec::new();
let mut segment_seen = HashSet::new();
let mut byte_target_flushes = 0usize;
let mut segment_limit_flushes = 0usize;
let mut large_segment_batches = 0usize;
let mut end_of_input_batches = 0usize;
for chunk_info in chunk_infos {
let segment_id = extract_segment_id(&chunk_info)?;
let chunk_sizes = extract_chunk_sizes(&chunk_info)?;
let segment_size: u64 = chunk_sizes.iter().sum();
if segment_seen.insert(segment_id.clone()) {
segment_order.push(segment_id.clone());
}
let adds_new_segment = !current_batch_segments.contains(&segment_id);
let would_exceed_size = current_batch_size + segment_size > target_size_bytes;
let would_exceed_segments =
adds_new_segment && current_batch_segments.len() >= MAX_CONCURRENT_SEGMENTS;
if !current_batch.is_empty() && (would_exceed_size || would_exceed_segments) {
if would_exceed_segments {
segment_limit_flushes += 1;
} else {
byte_target_flushes += 1;
}
let merged_batch = re_arrow_util::concat_polymorphic_batches(¤t_batch)
.map_err(|err| merge_err(err, "merging chunk-info batches"))?;
request_batches.push(merged_batch);
current_batch = Vec::new();
current_batch_size = 0;
current_batch_segments.clear();
}
if segment_size > target_size_bytes {
if !current_batch.is_empty() {
byte_target_flushes += 1;
let merged_batch = re_arrow_util::concat_polymorphic_batches(¤t_batch)
.map_err(|err| merge_err(err, "merging chunk-info batches"))?;
request_batches.push(merged_batch);
current_batch = Vec::new();
current_batch_size = 0;
current_batch_segments.clear();
}
let split_batches =
split_large_segments(&segment_id, &chunk_info, target_size_bytes, &chunk_sizes)?;
large_segment_batches += split_batches.len();
for split_batch in split_batches {
request_batches.push(split_batch);
}
} else {
current_batch.push(chunk_info);
current_batch_size += segment_size;
current_batch_segments.insert(segment_id);
}
}
if !current_batch.is_empty() {
let merged_batch = re_arrow_util::concat_polymorphic_batches(¤t_batch)
.map_err(|err| merge_err(err, "merging final chunk-info batch"))?;
request_batches.push(merged_batch);
end_of_input_batches += 1;
}
re_log::debug_assert_eq!(
request_batches.len(),
byte_target_flushes + segment_limit_flushes + large_segment_batches + end_of_input_batches,
"every planned batch must have exactly one flush reason"
);
let span = tracing::Span::current();
span.record("output_batches", request_batches.len());
span.record("byte_target_flushes", byte_target_flushes);
span.record("segment_limit_flushes", segment_limit_flushes);
span.record("large_segment_batches", large_segment_batches);
span.record("end_of_input_batches", end_of_input_batches);
tracing::debug!(
"Batching complete: {} segments → {} batches (target_size={}KB)",
segment_order.len(),
request_batches.len(),
target_size_bytes / 1024
);
Ok((request_batches, segment_order))
}
fn split_large_segments(
segment_id: &SegmentId,
chunk_info: &RecordBatch,
target_size: u64,
chunk_sizes: &quiver::Column<u64>,
) -> ApiResult<Vec<RecordBatch>> {
re_tracing::profile_function!();
let take_err = |err: arrow::error::ArrowError| {
ApiError::deserialization_with_source(None, err, "slicing large segment into sub-batches")
};
let mut result_batches = Vec::new();
let mut current_indices = Vec::new();
let mut current_size = 0u64;
for row_idx in 0..chunk_info.num_rows() {
let chunk_size = chunk_sizes[row_idx];
if current_indices.is_empty() || current_size + chunk_size <= target_size {
current_indices.push(row_idx);
current_size += chunk_size;
} else {
let batch =
re_arrow_util::take_record_batch(chunk_info, ¤t_indices).map_err(take_err)?;
result_batches.push(batch);
current_indices = vec![row_idx];
current_size = chunk_size;
}
}
if !current_indices.is_empty() {
let batch =
re_arrow_util::take_record_batch(chunk_info, ¤t_indices).map_err(take_err)?;
result_batches.push(batch);
}
tracing::debug!(
"Split large segment '{}' ({}) into {} requests",
segment_id,
re_format::format_bytes(chunk_sizes.iter().sum::<u64>() as _),
result_batches.len()
);
Ok(result_batches)
}
fn sort_chunks_by_segment_order(
chunks: Vec<ChunksWithSegment>,
segment_order: &[SegmentId],
) -> Vec<SortedChunksWithSegment> {
let mut segment_groups: HashMap<SegmentId, Vec<Chunk>> = HashMap::default();
for chunks_with_segment in chunks {
for (chunk, segment_id_opt) in chunks_with_segment {
let Some(segment_id) = segment_id_opt else {
continue;
};
segment_groups.entry(segment_id).or_default().push(chunk);
}
}
segment_order
.iter()
.filter_map(|segment_id| segment_groups.remove_entry(segment_id))
.collect()
}
async fn send_sorted_chunks(
chunks: Vec<ChunksWithSegment>,
global_segment_order: &[SegmentId],
output_channel: &Sender<ApiResult<CpuWorkerMsg>>,
) -> bool {
let sorted = {
let _span = tracing::info_span!("sort_chunks").entered();
sort_chunks_by_segment_order(chunks, global_segment_order)
};
let n_sorted = sorted.len();
async {
for chunk in sorted {
if output_channel
.send(Ok(CpuWorkerMsg::Chunks(chunk)))
.await
.is_err()
{
return false;
}
}
true
}
.instrument(tracing::info_span!("send_chunks", n = n_sorted))
.await
}
fn create_grpc_batch_groups(batches: &[RecordBatch]) -> ApiResult<Vec<&[RecordBatch]>> {
let mut groups = Vec::new();
let mut group_start = 0usize;
let mut group_segments: HashSet<String> = HashSet::new();
for (idx, batch) in batches.iter().enumerate() {
let mut batch_seen = HashSet::new();
let mut batch_segments = Vec::new();
extend_distinct_segment_ids(batch, &mut batch_seen, &mut batch_segments)?;
if batch_segments.len() > MAX_CONCURRENT_SEGMENTS {
return Err(ApiError::internal(format!(
"single gRPC fetch batch spans {} distinct segments, exceeding the cap of {}",
batch_segments.len(),
MAX_CONCURRENT_SEGMENTS,
)));
}
let would_exceed_batch_count = idx - group_start >= GRPC_BATCH_SIZE;
let new_segments = batch_segments
.iter()
.filter(|segment_id| !group_segments.contains(*segment_id))
.count();
let would_exceed_segments = group_segments.len() + new_segments > MAX_CONCURRENT_SEGMENTS;
if idx > group_start && (would_exceed_batch_count || would_exceed_segments) {
groups.push(&batches[group_start..idx]);
group_start = idx;
group_segments.clear();
}
group_segments.extend(batch_segments);
}
if group_start < batches.len() {
groups.push(&batches[group_start..]);
}
Ok(groups)
}
fn distinct_segment_ids_for_batch(batch: &RecordBatch) -> ApiResult<Vec<String>> {
distinct_segment_ids_for_batches(std::iter::once(batch))
}
fn distinct_segment_ids_for_batches<'a>(
batches: impl IntoIterator<Item = &'a RecordBatch>,
) -> ApiResult<Vec<String>> {
let mut seen = HashSet::new();
let mut segment_ids = Vec::new();
for batch in batches {
extend_distinct_segment_ids(batch, &mut seen, &mut segment_ids)?;
}
Ok(segment_ids)
}
fn segment_wave_index(
segment_id: &str,
segment_to_wave: &HashMap<String, usize>,
) -> ApiResult<usize> {
segment_to_wave.get(segment_id).copied().ok_or_else(|| {
ApiError::internal(format!(
"fetch batch references segment {segment_id} that is not present in global segment order"
))
})
}
fn batches_by_segment_wave(
batches: &[RecordBatch],
global_segment_order: &[SegmentId],
) -> ApiResult<(Vec<Vec<RecordBatch>>, usize, usize)> {
let n_waves = global_segment_order.len().div_ceil(MAX_CONCURRENT_SEGMENTS);
let mut segment_to_wave = HashMap::new();
for (idx, segment_id) in global_segment_order.iter().enumerate() {
segment_to_wave.insert(
segment_id.as_ref().to_owned(),
idx / MAX_CONCURRENT_SEGMENTS,
);
}
let mut waves = vec![Vec::new(); n_waves];
let mut wave_segment_ids = vec![HashSet::new(); n_waves];
let mut max_segments_per_batch = 0;
for batch in batches {
let segment_ids = distinct_segment_ids_for_batch(batch)?;
let Some(first_segment) = segment_ids.first() else {
continue;
};
max_segments_per_batch = max_segments_per_batch.max(segment_ids.len());
let wave_idx = segment_wave_index(first_segment, &segment_to_wave)?;
for segment_id in &segment_ids[1..] {
let other_wave_idx = segment_wave_index(segment_id, &segment_to_wave)?;
if other_wave_idx != wave_idx {
return Err(ApiError::internal(format!(
"fetch batch spans segment waves: first segment {first_segment} is in wave \
{wave_idx}, but segment {segment_id} is in wave {other_wave_idx}"
)));
}
}
wave_segment_ids[wave_idx].extend(segment_ids);
waves[wave_idx].push(batch.clone());
}
let max_segments_per_wave = wave_segment_ids
.into_iter()
.map(|segment_ids| segment_ids.len())
.max()
.unwrap_or(0);
Ok((waves, max_segments_per_batch, max_segments_per_wave))
}
async fn admit_segment_wave(
wave_segments: &[SegmentId],
wave_batches: &[RecordBatch],
segment_chunk_counts: &HashMap<SegmentId, usize>,
manifests: &mut HashMap<SegmentId, SegmentChunkManifest>,
output_channel: &Sender<ApiResult<CpuWorkerMsg>>,
pipeline_budget: &PipelineBudget,
filtered_index_timeline: Option<&str>,
) -> ApiResult<bool> {
if wave_segments.is_empty() {
return Ok(true);
}
let task_time_min = wave_batches
.iter()
.map(|b| extract_task_time_min(b, filtered_index_timeline))
.min()
.unwrap_or(TimeInt::MAX);
let segment_ids = wave_segments
.iter()
.map(|segment_id| segment_id.as_ref().to_owned())
.collect();
let admission_guard = pipeline_budget
.reserve_guarded_with_priority(0, task_time_min, segment_ids)
.await;
for segment_id in wave_segments {
let count = segment_chunk_counts
.get(segment_id)
.copied()
.ok_or_else(|| {
ApiError::internal(format!(
"missing chunk count for admitted segment {segment_id}"
))
})?;
if output_channel
.send(Ok(CpuWorkerMsg::SegmentChunkCount {
segment_id: segment_id.clone(),
count,
}))
.await
.is_err()
{
return Ok(false);
}
if let Some(manifest) = manifests.remove(segment_id)
&& output_channel
.send(Ok(CpuWorkerMsg::SegmentManifest {
segment_id: segment_id.clone(),
manifest: Box::new(manifest),
}))
.await
.is_err()
{
return Ok(false);
}
}
admission_guard.commit(0);
Ok(true)
}
async fn fetch_remaining_via_grpc<T: DataframeClientAPI>(
batches: &[RecordBatch],
client: &T,
global_segment_order: &[SegmentId],
filtered_index_timeline: Option<&str>,
output_channel: &Sender<ApiResult<CpuWorkerMsg>>,
pipeline_budget: &PipelineBudget,
metrics: &QueryMetrics,
) -> ApiResult<()> {
let total_batches = batches.len();
let mut batches_completed = 0usize;
for batch_group in create_grpc_batch_groups(batches)? {
#[cfg(not(target_arch = "wasm32"))]
{
let bytes: u64 = batch_group.iter().map(batch_byte_size).sum();
crate::chunk_fetcher::metrics::record_grpc_no_direct_urls(bytes);
}
let estimated = batch_group
.iter()
.map(|b| batch_byte_size_uncompressed(b).unwrap_or_else(|| batch_byte_size(b)))
.sum::<u64>() as usize;
let task_time_min = batch_group
.iter()
.map(|b| extract_task_time_min(b, filtered_index_timeline))
.min()
.unwrap_or(TimeInt::MAX);
let mut segment_ids: Vec<String> = Vec::new();
{
let mut seen: HashSet<String> = HashSet::new();
for b in batch_group {
extend_distinct_segment_ids(b, &mut seen, &mut segment_ids)?;
}
}
re_log::debug_assert!(
segment_ids.len() <= MAX_CONCURRENT_SEGMENTS,
"gRPC batch group exceeded segment gate cap"
);
let guard = pipeline_budget
.reserve_guarded_with_priority(estimated, task_time_min, segment_ids)
.await;
let mut stats = TaskFetchStats::default();
let all_chunks = fetch_batch_group_via_grpc(
batch_group,
client,
&metrics.fetch_grpc_requests,
&mut stats,
)
.await?;
let actual: usize = all_chunks
.iter()
.flat_map(|segment_chunks| {
segment_chunks
.iter()
.map(|(chunk, _)| re_byte_size::SizeBytes::total_size_bytes(chunk) as usize)
})
.sum();
guard.commit(actual);
stats.flush_into(metrics);
batches_completed += batch_group.len();
if !send_sorted_chunks(all_chunks, global_segment_order, output_channel).await {
tracing::info!(
total_batches,
batches_completed,
batches_skipped = total_batches.saturating_sub(batches_completed),
"FetchChunks IO loop short-circuited: downstream consumer closed (likely LIMIT or plan cancellation)"
);
return Ok(());
}
}
Ok(())
}
#[tracing::instrument(
level = "info",
skip_all,
fields(
n_chunks,
n_batches,
n_segments,
n_waves,
segment_admission_limit,
max_segments_per_batch,
max_segments_per_wave,
fetch_strategy,
)
)]
pub(super) async fn chunk_stream_io_loop<T: DataframeClientAPI>(
client: T,
chunk_infos: Vec<RecordBatch>,
filtered_index_timeline: Option<TimelineName>,
output_channel: Sender<ApiResult<CpuWorkerMsg>>,
pending_analytics: crate::PendingQueryAnalytics,
pipeline_budget: Arc<PipelineBudget>,
) -> ApiResult<()> {
let target_size_bytes = TARGET_BATCH_SIZE_BYTES as u64;
let metrics = Arc::clone(pending_analytics.metrics());
let n_chunks: usize = chunk_infos.iter().map(|rb| rb.num_rows()).sum();
let segment_chunk_counts: HashMap<SegmentId, usize> = count_chunks_per_segment(&chunk_infos)?
.into_iter()
.collect();
let filtered_index_timeline_str: Option<Arc<str>> = filtered_index_timeline
.as_ref()
.map(|t| Arc::<str>::from(t.as_str()));
let mut manifests = build_segment_manifests(&chunk_infos, filtered_index_timeline)?;
let (request_batches, global_segment_order) =
create_request_batches(chunk_infos, target_size_bytes)?;
let (request_batches_by_wave, max_segments_per_batch, max_segments_per_wave) =
batches_by_segment_wave(&request_batches, &global_segment_order)?;
metrics
.planned_fetch_batches
.fetch_add(request_batches.len() as u64, Ordering::Relaxed);
metrics
.planned_segment_waves
.fetch_add(request_batches_by_wave.len() as u64, Ordering::Relaxed);
metrics
.segment_admission_limit
.fetch_max(MAX_CONCURRENT_SEGMENTS as u64, Ordering::Relaxed);
metrics
.max_segments_per_fetch_batch
.fetch_max(max_segments_per_batch as u64, Ordering::Relaxed);
metrics
.max_segments_per_wave
.fetch_max(max_segments_per_wave as u64, Ordering::Relaxed);
let span = tracing::Span::current();
span.record("n_chunks", n_chunks);
span.record("n_batches", request_batches.len());
span.record("n_segments", global_segment_order.len());
span.record("n_waves", request_batches_by_wave.len());
span.record("segment_admission_limit", MAX_CONCURRENT_SEGMENTS);
span.record("max_segments_per_batch", max_segments_per_batch);
span.record("max_segments_per_wave", max_segments_per_wave);
re_log::debug!(
"Fetching {n_chunks} chunks in {} batches ({} segments)",
request_batches.len(),
global_segment_order.len()
);
let force_grpc = force_grpc();
let grpc_only = force_grpc || !request_batches.iter().any(batch_has_any_direct_urls);
if grpc_only {
let reason = if force_grpc {
"grpc_forced"
} else {
"no_direct_urls"
};
span.record("fetch_strategy", reason);
re_log::debug!(
"{reason}, fetching all {} chunks via FetchChunks gRPC in segment waves",
request_batches.len()
);
let result = async {
for (wave_idx, wave_batches) in request_batches_by_wave.iter().enumerate() {
let start = wave_idx * MAX_CONCURRENT_SEGMENTS;
let end = (start + MAX_CONCURRENT_SEGMENTS).min(global_segment_order.len());
let wave_segments = &global_segment_order[start..end];
if !admit_segment_wave(
wave_segments,
wave_batches,
&segment_chunk_counts,
&mut manifests,
&output_channel,
&pipeline_budget,
filtered_index_timeline_str.as_deref(),
)
.await?
{
return Ok(());
}
fetch_remaining_via_grpc(
wave_batches,
&client,
&global_segment_order,
filtered_index_timeline_str.as_deref(),
&output_channel,
&pipeline_budget,
&metrics,
)
.await?;
}
Ok(())
}
.await;
if result.is_err() {
pending_analytics.record_error(QueryErrorKind::GrpcFetch);
}
return result;
}
let total_n_direct: usize = request_batches
.iter()
.filter_map(|batch| split_batch_by_direct_url(batch).0)
.count();
let total_n_grpc: usize = request_batches
.iter()
.filter_map(|batch| split_batch_by_direct_url(batch).1)
.count();
if total_n_grpc == 0 {
span.record("fetch_strategy", "direct");
} else {
span.record(
"fetch_strategy",
format!("hybrid(direct={total_n_direct},grpc={total_n_grpc})"),
);
}
re_log::debug!("Fetch tasks: {total_n_direct} direct, {total_n_grpc} gRPC fallback");
let http_client = reqwest::Client::builder()
.connect_timeout(DIRECT_FETCH_CONNECT_TIMEOUT)
.read_timeout(DIRECT_FETCH_READ_TIMEOUT)
.build()
.expect("static reqwest client config is valid");
let total_tasks = total_n_direct + total_n_grpc;
let mut tasks_completed: usize = 0;
for (wave_idx, wave_batches) in request_batches_by_wave.into_iter().enumerate() {
let start = wave_idx * MAX_CONCURRENT_SEGMENTS;
let end = (start + MAX_CONCURRENT_SEGMENTS).min(global_segment_order.len());
let wave_segments = &global_segment_order[start..end];
if !admit_segment_wave(
wave_segments,
&wave_batches,
&segment_chunk_counts,
&mut manifests,
&output_channel,
&pipeline_budget,
filtered_index_timeline_str.as_deref(),
)
.await?
{
return Ok(());
}
let (work_items, _n_direct, _n_grpc) = split_batches_into_fetch_tasks(&wave_batches);
let fetch_stream = futures::stream::iter(work_items.into_iter().enumerate())
.map(|(task_idx, task)| {
let http_client = http_client.clone();
let client = client.clone();
let pending_analytics = pending_analytics.clone();
let pipeline_budget = Arc::clone(&pipeline_budget);
let metrics = Arc::clone(&metrics);
let filtered_index_timeline = filtered_index_timeline_str.as_ref().map(Arc::clone);
async move {
let mut stats = TaskFetchStats::default();
let batch_ref = match &task {
FetchTask::Direct(b) | FetchTask::Grpc(b) => b,
};
let estimated = batch_byte_size_uncompressed(batch_ref)
.unwrap_or_else(|| batch_byte_size(batch_ref))
as usize;
let task_time_min =
extract_task_time_min(batch_ref, filtered_index_timeline.as_deref());
let mut seen: HashSet<String> = HashSet::new();
let mut segment_ids: Vec<String> = Vec::new();
extend_distinct_segment_ids(batch_ref, &mut seen, &mut segment_ids)?;
let guard = pipeline_budget
.reserve_guarded_with_priority(estimated, task_time_min, segment_ids)
.await;
let chunks = match task {
FetchTask::Direct(batch) => {
match fetch_batch_direct(
&batch,
&http_client,
&metrics.fetch_direct_requests,
&mut stats,
&pending_analytics,
)
.await
{
Ok(chunks) => chunks,
Err(err) => {
stats.try_flush_into(
&pending_analytics,
Err(QueryErrorKind::DirectFetch),
);
return Err(err);
}
}
}
FetchTask::Grpc(batch) => {
let bytes = batch_byte_size(&batch);
#[cfg(not(target_arch = "wasm32"))]
crate::chunk_fetcher::metrics::record_grpc_no_direct_urls(bytes);
match fetch_batch_group_via_grpc(
std::slice::from_ref(&batch),
&client,
&metrics.fetch_grpc_requests,
&mut stats,
)
.await
{
Ok(chunks) => chunks,
Err(err) => {
stats.try_flush_into(
&pending_analytics,
Err(QueryErrorKind::GrpcFetch),
);
return Err(err);
}
}
}
};
stats.try_flush_into(&pending_analytics, Ok(()));
let actual: usize = chunks
.iter()
.flat_map(|seg| {
seg.iter()
.map(|(c, _)| re_byte_size::SizeBytes::total_size_bytes(c) as usize)
})
.sum();
guard.commit(actual);
Ok::<_, ApiError>(chunks)
}
.instrument(tracing::info_span!("fetch_task", wave_idx, task_idx))
})
.buffer_unordered(IO_PIPELINE_BUFFER);
tokio::pin!(fetch_stream);
while let Some(result) = fetch_stream.next().await {
let chunks = result?;
tasks_completed += 1;
if !send_sorted_chunks(chunks, &global_segment_order, &output_channel).await {
tracing::info!(
total_tasks,
tasks_completed,
in_flight_or_pending = total_tasks.saturating_sub(tasks_completed),
"FetchChunks IO loop short-circuited (hybrid path): downstream consumer closed (likely LIMIT or plan cancellation)"
);
return Ok(());
}
}
}
Ok(())
}
#[cfg(test)]
mod tests {
use std::collections::HashMap;
use arrow::array::{
Array as _, BooleanArray, FixedSizeBinaryBuilder, Int64Array, RecordBatchOptions,
StringArray, UInt64Array,
};
use arrow::datatypes::{Field, Schema};
use re_log_types::EntityPath;
use super::*;
fn extract_segment_id_from_chunk((segment_id, _chunks): &SortedChunksWithSegment) -> &str {
segment_id.as_ref()
}
fn create_test_chunk_info(segment_id: &str, chunk_sizes: &[u64]) -> RecordBatch {
let num_chunks = chunk_sizes.len();
let segment_ids = StringArray::from(vec![segment_id; num_chunks]);
let sizes = UInt64Array::from(chunk_sizes.to_vec());
let mut chunk_id_builder = FixedSizeBinaryBuilder::with_capacity(num_chunks, 16);
for i in 0..num_chunks {
let mut id_bytes = [0u8; 16];
id_bytes[0..4].copy_from_slice(&(i as u32).to_le_bytes());
chunk_id_builder.append_value(id_bytes).unwrap();
}
let chunk_ids = chunk_id_builder.finish();
let schema = Arc::new(Schema::new_with_metadata(
vec![
QueryDatasetDataframe::COLUMN_CHUNK_SEGMENT_ID.arrow_field(),
Field::new(
QueryDatasetDataframe::COLUMN_CHUNK_BYTE_LEN_NAME,
arrow::datatypes::DataType::UInt64,
false,
),
QueryDatasetDataframe::COLUMN_CHUNK_ID.arrow_field(),
],
HashMap::default(),
));
RecordBatch::try_new_with_options(
schema,
vec![Arc::new(segment_ids), Arc::new(sizes), Arc::new(chunk_ids)],
&RecordBatchOptions::new().with_row_count(Some(num_chunks)),
)
.unwrap()
}
fn segment_order_as_strs(segment_order: &[SegmentId]) -> Vec<&str> {
segment_order.iter().map(SegmentId::as_ref).collect()
}
#[test]
fn test_create_request_batches_single_small_segment() {
let chunk_info = create_test_chunk_info("seg1", &[100, 200, 300]); let target_size = 1000;
let (batches, segment_order) =
create_request_batches(vec![chunk_info], target_size).unwrap();
assert_eq!(batches.len(), 1);
assert_eq!(batches[0].num_rows(), 3);
assert_eq!(segment_order_as_strs(&segment_order), vec!["seg1"]);
}
#[test]
fn test_create_request_batches_single_large_segment() {
let chunk_info = create_test_chunk_info("seg1", &[300, 400, 500, 600]); let target_size = 1000;
let (batches, segment_order) =
create_request_batches(vec![chunk_info], target_size).unwrap();
assert_eq!(batches.len(), 3);
assert_eq!(segment_order_as_strs(&segment_order), vec!["seg1"]);
}
#[test]
fn test_create_request_batches_multiple_small_segments() {
let chunk_infos = vec![
create_test_chunk_info("seg1", &[100, 150]), create_test_chunk_info("seg2", &[200, 250]), create_test_chunk_info("seg3", &[300]), create_test_chunk_info("seg4", &[100]), ];
let target_size = 800;
let (batches, segment_order) = create_request_batches(chunk_infos, target_size).unwrap();
assert_eq!(batches.len(), 2);
assert_eq!(batches[0].num_rows(), 4);
assert_eq!(batches[1].num_rows(), 2);
assert_eq!(
segment_order_as_strs(&segment_order),
vec!["seg1", "seg2", "seg3", "seg4"]
);
}
#[test]
fn test_create_request_batches_caps_segments_per_batch() {
let chunk_infos = vec![
create_test_chunk_info("seg1", &[10]),
create_test_chunk_info("seg2", &[10]),
create_test_chunk_info("seg3", &[10]),
create_test_chunk_info("seg4", &[10]),
create_test_chunk_info("seg5", &[10]),
create_test_chunk_info("seg6", &[10]),
];
let target_size = 100_000;
let (batches, segment_order) = create_request_batches(chunk_infos, target_size).unwrap();
assert_eq!(batches.len(), 2);
for batch in &batches {
let seg_col = batch
.column_by_name(QueryDatasetDataframe::COLUMN_CHUNK_SEGMENT_ID_NAME)
.unwrap()
.as_any()
.downcast_ref::<StringArray>()
.unwrap();
let distinct: HashSet<&str> = (0..seg_col.len()).map(|i| seg_col.value(i)).collect();
assert!(
distinct.len() <= MAX_CONCURRENT_SEGMENTS,
"batch has {} distinct segments, cap is {}",
distinct.len(),
MAX_CONCURRENT_SEGMENTS,
);
}
assert_eq!(
segment_order_as_strs(&segment_order),
vec!["seg1", "seg2", "seg3", "seg4", "seg5", "seg6"]
);
}
#[test]
fn test_many_tiny_segments_expose_serial_wave_shape() {
const NUM_SEGMENTS: usize = 3_990;
let chunk_infos = (0..NUM_SEGMENTS)
.map(|idx| create_test_chunk_info(&format!("seg{idx:04}"), &[10]))
.collect();
let (batches, segment_order) =
create_request_batches(chunk_infos, TARGET_BATCH_SIZE_BYTES as u64).unwrap();
let (waves, max_segments_per_batch, max_segments_per_wave) =
batches_by_segment_wave(&batches, &segment_order).unwrap();
let expected = NUM_SEGMENTS.div_ceil(MAX_CONCURRENT_SEGMENTS);
assert_eq!(segment_order.len(), NUM_SEGMENTS);
assert_eq!(batches.len(), expected);
assert_eq!(waves.len(), expected);
assert_eq!(max_segments_per_batch, MAX_CONCURRENT_SEGMENTS);
assert_eq!(max_segments_per_wave, MAX_CONCURRENT_SEGMENTS);
assert!(waves.iter().all(|wave| wave.len() == 1));
assert!(
batches
.iter()
.all(|batch| batch.num_rows() == MAX_CONCURRENT_SEGMENTS)
);
}
fn distinct_segments_in_batches<'a>(
batches: impl IntoIterator<Item = &'a RecordBatch>,
) -> HashSet<String> {
let mut seen = HashSet::new();
let mut order = Vec::new();
for batch in batches {
extend_distinct_segment_ids(batch, &mut seen, &mut order).unwrap();
}
seen
}
#[test]
fn test_batches_by_segment_wave_keeps_each_wave_within_segment_cap() {
let chunk_infos = vec![
create_test_chunk_info("seg1", &[100]),
create_test_chunk_info("seg2", &[100]),
create_test_chunk_info("seg3", &[100]),
create_test_chunk_info("seg4", &[100]),
create_test_chunk_info("seg5", &[100]),
];
let (batches, segment_order) = create_request_batches(chunk_infos, 10_000).unwrap();
let (waves, _, _) = batches_by_segment_wave(&batches, &segment_order).unwrap();
assert_eq!(waves.len(), 2);
assert_eq!(waves[0].len(), 1);
assert_eq!(waves[1].len(), 1);
assert_eq!(
distinct_segments_in_batches(&waves[0]),
HashSet::from_iter(["seg1".to_owned(), "seg2".to_owned(), "seg3".to_owned()])
);
assert_eq!(
distinct_segments_in_batches(&waves[1]),
HashSet::from_iter(["seg4".to_owned(), "seg5".to_owned()])
);
}
#[test]
fn test_create_grpc_batch_groups_preserves_segment_cap() {
let chunk_infos = vec![
create_test_chunk_info("seg1", &[10]),
create_test_chunk_info("seg2", &[10]),
create_test_chunk_info("seg3", &[10]),
create_test_chunk_info("seg4", &[10]),
create_test_chunk_info("seg5", &[10]),
create_test_chunk_info("seg6", &[10]),
];
let (batches, _) = create_request_batches(chunk_infos, 100_000).unwrap();
assert_eq!(batches.len(), 2);
let groups = create_grpc_batch_groups(&batches).unwrap();
assert_eq!(groups.len(), 2);
for group in groups {
assert!(group.len() <= GRPC_BATCH_SIZE);
assert!(distinct_segments_in_batches(group.iter()).len() <= MAX_CONCURRENT_SEGMENTS);
}
}
#[test]
fn test_create_grpc_batch_groups_keeps_same_segment_batches_together() {
let chunk_info = create_test_chunk_info("seg1", &[10; GRPC_BATCH_SIZE + 1]);
let (batches, _) = create_request_batches(vec![chunk_info], 10).unwrap();
assert_eq!(batches.len(), GRPC_BATCH_SIZE + 1);
let groups = create_grpc_batch_groups(&batches).unwrap();
assert_eq!(groups.len(), 2);
assert_eq!(groups[0].len(), GRPC_BATCH_SIZE);
assert_eq!(groups[1].len(), 1);
assert_eq!(distinct_segments_in_batches(groups[0].iter()).len(), 1);
assert_eq!(distinct_segments_in_batches(groups[1].iter()).len(), 1);
}
#[test]
fn test_create_request_batches_mixed_small_and_large() {
let chunk_infos = vec![
create_test_chunk_info("seg1", &[100, 200]), create_test_chunk_info("seg2", &[800, 900, 700]), create_test_chunk_info("seg3", &[150]), ];
let target_size = 1000;
let (batches, segment_order) = create_request_batches(chunk_infos, target_size).unwrap();
assert_eq!(batches.len(), 5);
assert_eq!(
segment_order_as_strs(&segment_order),
vec!["seg1", "seg2", "seg3"]
);
}
#[test]
fn test_skewed_segment_sizes_preserve_rows_order_and_caps() {
let mut chunk_infos: Vec<_> = (0..30)
.map(|idx| create_test_chunk_info(&format!("tiny{idx:02}"), &[1]))
.collect();
chunk_infos.push(create_test_chunk_info("large", &[600, 600, 600]));
let expected_rows: usize = chunk_infos.iter().map(RecordBatch::num_rows).sum();
let (batches, segment_order) = create_request_batches(chunk_infos, 1_000).unwrap();
let (waves, max_segments_per_batch, max_segments_per_wave) =
batches_by_segment_wave(&batches, &segment_order).unwrap();
assert_eq!(segment_order.len(), 31);
assert_eq!(segment_order.last().unwrap().as_ref(), "large");
assert_eq!(batches.len(), 13);
assert_eq!(waves.len(), 11);
assert_eq!(max_segments_per_batch, MAX_CONCURRENT_SEGMENTS);
assert_eq!(max_segments_per_wave, MAX_CONCURRENT_SEGMENTS);
assert_eq!(
batches.iter().map(RecordBatch::num_rows).sum::<usize>(),
expected_rows
);
assert!(batches.iter().all(|batch| {
distinct_segment_ids_for_batch(batch).unwrap().len() <= MAX_CONCURRENT_SEGMENTS
}));
assert!(waves.iter().all(|wave| {
distinct_segments_in_batches(wave.iter()).len() <= MAX_CONCURRENT_SEGMENTS
}));
}
#[test]
fn test_segment_order_within_batches_is_preserved() {
let chunk_infos = vec![
create_test_chunk_info("segA", &[100]), create_test_chunk_info("segB", &[200]), create_test_chunk_info("segC", &[300]), ];
let target_size = 1000;
let (batches, segment_order) = create_request_batches(chunk_infos, target_size).unwrap();
assert_eq!(batches.len(), 1);
assert_eq!(batches[0].num_rows(), 3);
assert_eq!(
segment_order_as_strs(&segment_order),
vec!["segA", "segB", "segC"]
);
let segment_id_column = QueryDatasetDataframe::COLUMN_CHUNK_SEGMENT_ID
.extract(&batches[0])
.unwrap();
let batch_segment_ids: Vec<&str> = segment_id_column.iter().collect();
assert_eq!(batch_segment_ids, ["segA", "segB", "segC"]);
}
#[test]
fn test_sort_chunks_by_segment_order_simple_case() {
use re_dataframe::external::re_chunk::Chunk;
use re_log_types::EntityPath;
let empty_chunk = Chunk::builder(EntityPath::root()).build().unwrap();
let segment_order: Vec<SegmentId> = vec!["segA".into(), "segB".into(), "segC".into()];
let chunks: Vec<ChunksWithSegment> = vec![
vec![(empty_chunk.clone(), Some("segC".into()))],
vec![(empty_chunk.clone(), Some("segA".into()))],
vec![(empty_chunk.clone(), Some("segB".into()))],
];
let sorted_chunks = sort_chunks_by_segment_order(chunks, &segment_order);
let sorted_segments: Vec<&str> = sorted_chunks
.iter()
.map(extract_segment_id_from_chunk)
.collect();
assert_eq!(sorted_segments, vec!["segA", "segB", "segC"]);
}
#[test]
fn test_sort_chunks_by_segment_order_multi_segment_response() {
use re_dataframe::external::re_chunk::Chunk;
use re_log_types::EntityPath;
let empty_chunk = Chunk::builder(EntityPath::root()).build().unwrap();
let segment_order: Vec<SegmentId> = vec!["segA".into(), "segB".into(), "segC".into()];
let chunks: Vec<ChunksWithSegment> = vec![
vec![
(empty_chunk.clone(), Some("segC".into())),
(empty_chunk.clone(), Some("segC".into())), (empty_chunk.clone(), Some("segA".into())),
(empty_chunk.clone(), Some("segB".into())),
(empty_chunk.clone(), Some("segB".into())), (empty_chunk.clone(), Some("segA".into())), (empty_chunk.clone(), Some("segB".into())), ],
];
let sorted_chunks = sort_chunks_by_segment_order(chunks, &segment_order);
assert_eq!(sorted_chunks.len(), 3);
let sorted_segments: Vec<&str> = sorted_chunks
.iter()
.map(extract_segment_id_from_chunk)
.collect();
assert_eq!(sorted_segments, vec!["segA", "segB", "segC"]);
let seg_a_chunks = sorted_chunks[0].1.len();
let seg_b_chunks = sorted_chunks[1].1.len();
let seg_c_chunks = sorted_chunks[2].1.len();
assert_eq!(seg_a_chunks, 2);
assert_eq!(seg_b_chunks, 3);
assert_eq!(seg_c_chunks, 2);
}
#[test]
fn test_sort_chunks_by_segment_order_mixed_responses() {
use re_dataframe::external::re_chunk::Chunk;
use re_log_types::EntityPath;
let empty_chunk = Chunk::builder(EntityPath::root()).build().unwrap();
let segment_order: Vec<SegmentId> = vec!["segA".into(), "segB".into(), "segC".into()];
let chunks: Vec<ChunksWithSegment> = vec![
vec![(empty_chunk.clone(), Some("segC".into()))],
vec![
(empty_chunk.clone(), Some("segB".into())),
(empty_chunk.clone(), Some("segA".into())),
],
vec![(empty_chunk.clone(), Some("segB".into()))],
];
let sorted_chunks = sort_chunks_by_segment_order(chunks, &segment_order);
assert_eq!(sorted_chunks.len(), 3);
let sorted_segments: Vec<&str> = sorted_chunks
.iter()
.map(extract_segment_id_from_chunk)
.collect();
assert_eq!(sorted_segments, vec!["segA", "segB", "segC"]);
let seg_b_chunks = sorted_chunks[1].1.len();
assert_eq!(seg_b_chunks, 2);
}
#[test]
fn test_count_chunks_per_segment_basic() {
let chunk_infos = vec![
create_test_chunk_info("segA", &[10, 20, 30]),
create_test_chunk_info("segB", &[40]),
create_test_chunk_info("segC", &[50, 60]),
];
let counts = count_chunks_per_segment(&chunk_infos).unwrap();
assert_eq!(
counts,
vec![
(SegmentId::from("segA"), 3),
(SegmentId::from("segB"), 1),
(SegmentId::from("segC"), 2),
]
);
}
#[test]
fn test_count_chunks_per_segment_sums_across_batches() {
let chunk_infos = vec![
create_test_chunk_info("segA", &[1, 2]),
create_test_chunk_info("segB", &[3]),
create_test_chunk_info("segA", &[4, 5, 6]),
];
let counts = count_chunks_per_segment(&chunk_infos).unwrap();
assert_eq!(
counts,
vec![(SegmentId::from("segA"), 5), (SegmentId::from("segB"), 1),]
);
}
fn create_chunk_info_with_starts(
timeline_name: &str,
rows: &[(&str, &str, bool, Option<i64>)],
) -> RecordBatch {
let num_rows = rows.len();
let segment_ids = StringArray::from(rows.iter().map(|r| r.0).collect::<Vec<_>>());
let entity_paths = StringArray::from(rows.iter().map(|r| r.1).collect::<Vec<_>>());
let is_static = BooleanArray::from(rows.iter().map(|r| r.2).collect::<Vec<_>>());
let starts = Int64Array::from(rows.iter().map(|r| r.3).collect::<Vec<_>>());
let schema = Arc::new(Schema::new_with_metadata(
vec![
QueryDatasetDataframe::COLUMN_CHUNK_SEGMENT_ID.arrow_field(),
QueryDatasetDataframe::COLUMN_CHUNK_ENTITY_PATH.arrow_field(),
QueryDatasetDataframe::COLUMN_CHUNK_IS_STATIC.arrow_field(),
re_protos::cloud::v1alpha1::QueryDatasetResponse::field_timeline_start(
timeline_name,
)
.as_ref()
.clone(),
],
HashMap::default(),
));
RecordBatch::try_new_with_options(
schema,
vec![
Arc::new(segment_ids),
Arc::new(entity_paths),
Arc::new(is_static),
Arc::new(starts),
],
&RecordBatchOptions::new().with_row_count(Some(num_rows)),
)
.unwrap()
}
#[test]
fn test_build_segment_manifests_no_timeline_yields_empty() {
let rb = create_chunk_info_with_starts("time", &[("seg1", "/a", false, Some(10))]);
let manifests = build_segment_manifests(&[rb], None).unwrap();
assert!(manifests.is_empty());
}
#[test]
fn test_build_segment_manifests_missing_start_column_yields_empty() {
let rb = create_test_chunk_info("seg1", &[100]);
let manifests = build_segment_manifests(&[rb], Some(TimelineName::from("time"))).unwrap();
assert!(manifests.is_empty());
}
#[test]
fn test_build_segment_manifests_mixed_batches_yields_empty() {
let with_start = create_chunk_info_with_starts("time", &[("seg1", "/a", false, Some(10))]);
let without_start = create_test_chunk_info("seg1", &[20]);
let manifests = build_segment_manifests(
&[with_start, without_start],
Some(TimelineName::from("time")),
)
.unwrap();
assert!(
manifests.is_empty(),
"any batch missing `:start` must trigger full fallback, even if other batches carry it",
);
}
#[test]
fn test_build_segment_manifests_filters_null_and_static() {
let rb = create_chunk_info_with_starts(
"time",
&[
("seg1", "/a", false, None), ("seg1", "/b", true, Some(10)), ("seg1", "/c", false, Some(20)), ("seg1", "/d", false, Some(i64::MIN)), ],
);
let manifests = build_segment_manifests(&[rb], Some(TimelineName::from("time"))).unwrap();
let m = manifests
.get(&SegmentId::from("seg1"))
.expect("seg1 has at least one temporal chunk");
assert_eq!(m.outstanding_count(), 2);
assert!(m.is_locked());
assert_eq!(m.safe_horizon(), Some(TimeInt::MIN));
}
#[test]
fn test_build_segment_manifests_entity_path_keys_roundtrip() {
let rb = create_chunk_info_with_starts("time", &[("seg1", "/foo/bar", false, Some(42))]);
let mut manifests =
build_segment_manifests(&[rb], Some(TimelineName::from("time"))).unwrap();
let m = manifests.get_mut(&SegmentId::from("seg1")).unwrap();
assert!(m.record_arrival(
&EntityPath::from("/foo/bar"),
TimeInt::saturated_temporal_i64(42)
));
assert!(m.is_complete());
}
#[test]
fn test_build_segment_manifests_all_static_segment_absent() {
let rb = create_chunk_info_with_starts(
"time",
&[
("segA", "/a", true, Some(10)),
("segA", "/b", true, Some(20)),
("segB", "/c", false, Some(30)),
],
);
let manifests = build_segment_manifests(&[rb], Some(TimelineName::from("time"))).unwrap();
assert!(!manifests.contains_key(&SegmentId::from("segA")));
let b = &manifests[&SegmentId::from("segB")];
assert_eq!(b.outstanding_count(), 1);
}
#[test]
fn test_build_segment_manifests_accepts_timestamp_ns_start_column() {
use arrow::array::TimestampNanosecondArray;
let num_rows = 2;
let segment_ids = StringArray::from(vec!["seg1", "seg1"]);
let entity_paths = StringArray::from(vec!["/a", "/b"]);
let is_static = BooleanArray::from(vec![false, false]);
let starts = TimestampNanosecondArray::from(vec![Some(100_i64), Some(200_i64)]);
let schema = Arc::new(Schema::new_with_metadata(
vec![
QueryDatasetDataframe::COLUMN_CHUNK_SEGMENT_ID.arrow_field(),
QueryDatasetDataframe::COLUMN_CHUNK_ENTITY_PATH.arrow_field(),
QueryDatasetDataframe::COLUMN_CHUNK_IS_STATIC.arrow_field(),
Field::new(
"time:start",
arrow::datatypes::DataType::Timestamp(
arrow::datatypes::TimeUnit::Nanosecond,
None,
),
true,
),
],
HashMap::default(),
));
let rb = RecordBatch::try_new_with_options(
schema,
vec![
Arc::new(segment_ids),
Arc::new(entity_paths),
Arc::new(is_static),
Arc::new(starts),
],
&RecordBatchOptions::new().with_row_count(Some(num_rows)),
)
.unwrap();
let manifests = build_segment_manifests(&[rb], Some(TimelineName::from("time"))).unwrap();
let m = manifests
.get(&SegmentId::from("seg1"))
.expect("seg1 has temporal chunks");
assert_eq!(m.outstanding_count(), 2);
assert_eq!(m.safe_horizon(), Some(TimeInt::new_temporal(99)));
}
}