fcb_core/http_reader/mod.rs
1//! HTTP reader for FlatCityBuf files
2//!
3//! This module contains HTTP range request patterns and streaming functionality
4//! derived from FlatGeobuf (https://github.com/flatgeobuf/flatgeobuf)
5//! Licensed under BSD 2-Clause License, Copyright (c) 2018-2024, Björn Harrtell and contributors
6
7use crate::deserializer::to_cj_feature;
8use crate::{add_indices_to_multi_memory_index, build_query, fb::*, AttrQuery};
9
10use crate::error::{Error, Result};
11use crate::packed_rtree::Query;
12use crate::reader::city_buffer::FcbBuffer;
13use crate::static_btree::{FixedStringKey, Float, KeyType, Operator};
14use crate::{
15 check_magic_bytes, size_prefixed_root_as_city_feature, HEADER_MAX_BUFFER_SIZE,
16 HEADER_SIZE_SIZE, MAGIC_BYTES_SIZE,
17};
18use byteorder::{ByteOrder, LittleEndian};
19use bytes::{BufMut, Bytes, BytesMut};
20use chrono::{DateTime, Utc};
21use cjseq::CityJSONFeature;
22use http_range_client::BufferedHttpRangeClient;
23use http_range_client::{AsyncBufferedHttpRangeClient, AsyncHttpRangeClient};
24use log::debug;
25use reqwest;
26
27use crate::packed_rtree::{http::HttpRange, http::HttpSearchResultItem, NodeItem, PackedRTree};
28use crate::static_btree::{
29 http::HttpRange as AttrHttpRange, http::HttpSearchResultItem as AttrHttpSearchResultItem,
30};
31use crate::static_btree::{HttpIndex, HttpMultiIndex};
32use std::collections::HashMap;
33use std::collections::VecDeque;
34use std::ops::Range;
35use tracing::trace;
36
37#[cfg(test)]
38mod mock_http_range_client;
39
40// The largest request we'll speculatively make.
41// If a single huge feature requires, we'll necessarily exceed this limit.
42const DEFAULT_HTTP_FETCH_SIZE: usize = 1_048_576; // 1MB
43
44/// FlatCityBuf dataset HTTP reader
45pub struct HttpFcbReader<T: AsyncHttpRangeClient + Send + Sync> {
46 client: AsyncBufferedHttpRangeClient<T>,
47 // feature reading requires header access, therefore
48 // header_buf is included in the FcbBuffer struct.
49 fbs: FcbBuffer,
50}
51
52pub struct AsyncFeatureIter<T: AsyncHttpRangeClient + Send + Sync> {
53 client: AsyncBufferedHttpRangeClient<T>,
54 // feature reading requires header access, therefore
55 // header_buf is included in the FcbBuffer struct.
56 fbs: FcbBuffer,
57 /// Which features to iterate
58 selection: FeatureSelection,
59 /// Number of selected features
60 count: usize,
61}
62
63impl HttpFcbReader<reqwest::Client> {
64 pub async fn open(url: &str) -> Result<HttpFcbReader<reqwest::Client>> {
65 let client = BufferedHttpRangeClient::new(url);
66 Self::_open(client).await
67 }
68}
69
70impl<T: AsyncHttpRangeClient + Send + Sync> HttpFcbReader<T> {
71 pub async fn new(client: AsyncBufferedHttpRangeClient<T>) -> Result<HttpFcbReader<T>> {
72 Self::_open(client).await
73 }
74
75 async fn _open(mut client: AsyncBufferedHttpRangeClient<T>) -> Result<HttpFcbReader<T>> {
76 // Because we use a buffered HTTP reader, anything extra we fetch here can
77 // be utilized to skip subsequent fetches.
78 // Immediately following the header is the optional spatial index, we deliberately fetch
79 // a small part of that to skip subsequent requests
80 let prefetch_index_bytes: usize = {
81 // The actual branching factor will be in the header, but since we don't have the header
82 // yet we guess. The consequence of getting this wrong isn't catastrophic, it just means
83 // we may be fetching slightly more than we need or that we make an extra request later.
84 let assumed_branching_factor = PackedRTree::DEFAULT_NODE_SIZE as usize;
85
86 // NOTE: each layer is exponentially larger
87 let prefetched_layers: u32 = 3;
88
89 (0..prefetched_layers)
90 .map(|i| assumed_branching_factor.pow(i) * std::mem::size_of::<NodeItem>())
91 .sum()
92 };
93
94 // In reality, the header is probably less than half this size, but better to overshoot and
95 // fetch an extra kb rather than have to issue a second request.
96 let assumed_header_size = 2024;
97 let min_req_size = assumed_header_size + prefetch_index_bytes;
98 client.set_min_req_size(min_req_size);
99 let mut read_bytes = 0;
100 let bytes = client.get_range(read_bytes, MAGIC_BYTES_SIZE).await?; // to get magic bytes
101 if !check_magic_bytes(bytes) {
102 return Err(Error::MissingMagicBytes);
103 }
104
105 read_bytes += MAGIC_BYTES_SIZE;
106 let mut bytes = BytesMut::from(client.get_range(read_bytes, HEADER_SIZE_SIZE).await?);
107 read_bytes += HEADER_SIZE_SIZE;
108
109 let header_size = LittleEndian::read_u32(&bytes) as usize;
110 if header_size > HEADER_MAX_BUFFER_SIZE || header_size < 8 {
111 // minimum size check avoids panic in FlatBuffers header decoding
112 return Err(Error::IllegalHeaderSize(header_size));
113 }
114
115 bytes.put(client.get_range(read_bytes, header_size).await?);
116 read_bytes += header_size;
117
118 let header_buf = bytes.to_vec();
119
120 // verify flatbuffer
121 let _header = size_prefixed_root_as_header(&header_buf)?;
122
123 Ok(HttpFcbReader {
124 client,
125 fbs: FcbBuffer {
126 header_buf,
127 features_buf: Vec::new(),
128 },
129 })
130 }
131
132 pub fn header(&self) -> Header {
133 self.fbs.header()
134 }
135 fn header_len(&self) -> usize {
136 MAGIC_BYTES_SIZE + self.fbs.header_buf.len()
137 }
138
139 fn rtree_index_size(&self) -> usize {
140 let header = self.fbs.header();
141 let feat_count = header.features_count() as usize;
142 if header.index_node_size() > 0 && feat_count > 0 {
143 PackedRTree::index_size(feat_count, header.index_node_size())
144 } else {
145 0
146 }
147 }
148
149 fn attr_index_size(&self) -> usize {
150 let header = self.fbs.header();
151 header
152 .attribute_index()
153 .map(|attr_index| {
154 attr_index
155 .iter()
156 .try_fold(0, |acc, ai| {
157 let len = ai.length() as usize;
158 if len > usize::MAX - acc {
159 Err(Error::AttributeIndexSizeOverflow)
160 } else {
161 Ok(acc + len)
162 }
163 }) // sum of all attribute index lengths
164 .unwrap_or(0)
165 })
166 .unwrap_or(0)
167 }
168
169 fn index_size(&self) -> usize {
170 self.rtree_index_size() + self.attr_index_size()
171 }
172
173 /// Select all features.
174 pub async fn select_all(self) -> Result<AsyncFeatureIter<T>> {
175 let header = self.fbs.header();
176 let count = header.features_count();
177 let index_size = self.index_size() as usize;
178 // Skip index
179 let feature_base = self.header_len() + index_size;
180 Ok(AsyncFeatureIter {
181 client: self.client,
182 fbs: self.fbs,
183 selection: FeatureSelection::SelectAll(SelectAll {
184 features_left: count,
185 pos: feature_base,
186 }),
187 count: count as usize,
188 })
189 }
190 /// Select features within a bounding box.
191 pub async fn select_query(mut self, query: Query) -> Result<AsyncFeatureIter<T>> {
192 self.select_query_paged(query, None, None).await
193 }
194
195 /// Select features within a bounding box with optional pagination.
196 /// If `limit`/`offset` are provided, only a page of features is returned while
197 /// `features_count()` on the returned iterator still reflects the total number of matches.
198 pub async fn select_query_paged(
199 mut self,
200 query: Query,
201 limit: Option<usize>,
202 offset: Option<usize>,
203 ) -> Result<AsyncFeatureIter<T>> {
204 // Read R-Tree index and build filter for features within bbox
205 let header = self.fbs.header();
206 if header.index_node_size() == 0 || header.features_count() == 0 {
207 return Err(Error::NoIndex);
208 }
209 let count = header.features_count() as usize;
210 // The R-tree branching factor is a per-file property: traversing with
211 // the compile-time default instead walks the wrong node ranges.
212 let index_node_size = header.index_node_size();
213 let header_len = self.header_len();
214
215 // request up to this many extra bytes if it means we can eliminate an extra request
216 let combine_request_threshold = 256 * 1024;
217 let attr_index_size = self.attr_index_size() as usize;
218 let list = PackedRTree::http_stream_search(
219 &mut self.client,
220 header_len,
221 attr_index_size,
222 count,
223 index_node_size,
224 query,
225 combine_request_threshold,
226 )
227 .await?;
228 debug_assert!(
229 list.windows(2)
230 .all(|w| w[0].range.start() < w[1].range.start()),
231 "Since the tree is traversed breadth first, list should be sorted by construction."
232 );
233
234 let total_count = list.len();
235
236 // Apply pagination
237 let start = offset.unwrap_or(0).min(total_count);
238 let end = match limit {
239 Some(l) => start.saturating_add(l).min(total_count),
240 None => total_count,
241 };
242 let page_list: Vec<_> = if start < end {
243 list.into_iter().skip(start).take(end - start).collect()
244 } else {
245 Vec::new()
246 };
247
248 let feature_batches =
249 FeatureBatch::make_batches(page_list, combine_request_threshold).await?;
250 let selection = FeatureSelection::SelectBbox(SelectBbox { feature_batches });
251 Ok(AsyncFeatureIter {
252 client: self.client,
253 fbs: self.fbs,
254 selection,
255 count: total_count,
256 })
257 }
258
259 /// This method uses the attribute index section to find matching feature offsets.
260 /// It then groups (batches) the remote feature ranges in order to reduce IO overhead.
261 pub async fn select_attr_query(mut self, query: &AttrQuery) -> Result<AsyncFeatureIter<T>> {
262 self.select_attr_query_paged(query, None, None).await
263 }
264
265 /// Attribute query with optional pagination where the iterator returns only the requested page,
266 /// while `features_count()` reflects the total number of matches.
267 pub async fn select_attr_query_paged(
268 mut self,
269 query: &AttrQuery,
270 limit: Option<usize>,
271 offset: Option<usize>,
272 ) -> Result<AsyncFeatureIter<T>> {
273 let header = self.fbs.header();
274 let header_len = self.header_len();
275 // Assume the header provides rtree and attribute index sizes.
276
277 // file structure:
278 // magic_bytes + header + rtree_index + attr_index1 + attr_index2 + ... + features
279 let rtree_index_size = self.rtree_index_size() as usize;
280 let attr_index_size = self.attr_index_size() as usize;
281 let attr_index_begin = header_len + rtree_index_size;
282 let feature_begin = header_len + rtree_index_size + attr_index_size;
283
284 let attr_index_entries = header
285 .attribute_index()
286 .ok_or_else(|| Error::AttributeIndexNotFound)?;
287 let mut attr_index_entries = attr_index_entries.iter().collect::<Vec<_>>();
288 let columns: Vec<Column> = header
289 .columns()
290 .ok_or_else(|| Error::NoColumnsInHeader)?
291 .iter()
292 .collect();
293 attr_index_entries.sort_by_key(|attr_info| attr_info.index());
294
295 // Build the query
296 let query = build_query(&query);
297
298 // Create a StreamableMultiIndex from HTTP range requests
299 let mut http_multi_index = HttpMultiIndex::new();
300
301 let mut current_index_begin = attr_index_begin;
302 for attr_info in attr_index_entries.iter() {
303 Self::add_indices_to_multi_http_index(
304 &mut http_multi_index,
305 &columns,
306 attr_info,
307 current_index_begin,
308 feature_begin,
309 )?;
310 current_index_begin += attr_info.length() as usize;
311 }
312
313 let result = http_multi_index
314 .query(&mut self.client, &query.conditions)
315 .await?;
316
317 let total_count = result.len();
318
319 // Apply pagination to attribute query results
320 let start = offset.unwrap_or(0).min(total_count);
321 let end = match limit {
322 Some(l) => start.saturating_add(l).min(total_count),
323 None => total_count,
324 };
325 let paged_iter: Vec<_> = if start < end {
326 result.into_iter().skip(start).take(end - start).collect()
327 } else {
328 Vec::new()
329 };
330
331 let http_ranges: Vec<HttpRange> = paged_iter
332 .into_iter()
333 .map(|item| match item.range {
334 AttrHttpRange::Range(range) => HttpRange::Range(range.start..range.end),
335 AttrHttpRange::RangeFrom(range) => HttpRange::RangeFrom(range.start..),
336 })
337 .collect();
338
339 Ok(AsyncFeatureIter {
340 client: self.client,
341 fbs: self.fbs,
342 selection: FeatureSelection::SelectAttr(SelectAttr {
343 ranges: http_ranges,
344 range_pos: 0,
345 }),
346 count: total_count,
347 })
348 }
349
350 pub fn add_indices_to_multi_http_index<C: AsyncHttpRangeClient + Send + Sync>(
351 multi_index: &mut HttpMultiIndex<C>,
352 columns: &[Column],
353 attr_info: &AttributeIndex,
354 index_begin: usize,
355 feature_begin: usize,
356 ) -> Result<()> {
357 if let Some(col) = columns.iter().find(|col| col.index() == attr_info.index()) {
358 // TODO: now it assuming to add all indices to the multi_index. However, we should only add the indices that are used in the query. To do that, we need to change the implementation of StreamMultiIndex. Current StreamMultiIndex's `add_index` method assumes that all indices are added to the multi_index. We'll change it to take Range<usize> as an argument.
359 match col.type_() {
360 ColumnType::Int => {
361 let index = HttpIndex::<i32>::new(
362 attr_info.num_unique_items() as usize,
363 attr_info.branching_factor(),
364 index_begin,
365 feature_begin,
366 1024 * 1024, // combine_request_threshold
367 );
368 multi_index.add_index(col.name().to_string(), index);
369 }
370 ColumnType::Float => {
371 let index = HttpIndex::<Float<f32>>::new(
372 attr_info.num_unique_items() as usize,
373 attr_info.branching_factor(),
374 index_begin,
375 feature_begin,
376 1024 * 1024, // combine_request_threshold
377 );
378 multi_index.add_index(col.name().to_string(), index);
379 }
380 ColumnType::Double => {
381 let index = HttpIndex::<Float<f64>>::new(
382 attr_info.num_unique_items() as usize,
383 attr_info.branching_factor(),
384 index_begin,
385 feature_begin,
386 1024 * 1024, // combine_request_threshold
387 );
388 multi_index.add_index(col.name().to_string(), index);
389 }
390 ColumnType::String => {
391 let index = HttpIndex::<FixedStringKey<50>>::new(
392 attr_info.num_unique_items() as usize,
393 attr_info.branching_factor(),
394 index_begin,
395 feature_begin,
396 1024 * 1024, // combine_request_threshold
397 );
398 multi_index.add_index(col.name().to_string(), index);
399 }
400
401 ColumnType::Bool => {
402 let index = HttpIndex::<bool>::new(
403 attr_info.num_unique_items() as usize,
404 attr_info.branching_factor(),
405 index_begin,
406 feature_begin,
407 1024 * 1024, // combine_request_threshold
408 );
409 multi_index.add_index(col.name().to_string(), index);
410 }
411 ColumnType::DateTime => {
412 let index = HttpIndex::<DateTime<Utc>>::new(
413 attr_info.num_unique_items() as usize,
414 attr_info.branching_factor(),
415 index_begin,
416 feature_begin,
417 1024 * 1024, // combine_request_threshold
418 );
419 multi_index.add_index(col.name().to_string(), index);
420 }
421 ColumnType::Short => {
422 let index = HttpIndex::<i16>::new(
423 attr_info.num_unique_items() as usize,
424 attr_info.branching_factor(),
425 index_begin,
426 feature_begin,
427 1024 * 1024, // combine_request_threshold
428 );
429 multi_index.add_index(col.name().to_string(), index);
430 }
431 ColumnType::UShort => {
432 let index = HttpIndex::<u16>::new(
433 attr_info.num_unique_items() as usize,
434 attr_info.branching_factor(),
435 index_begin,
436 feature_begin,
437 1024 * 1024, // combine_request_threshold
438 );
439 multi_index.add_index(col.name().to_string(), index);
440 }
441 ColumnType::UInt => {
442 let index = HttpIndex::<u32>::new(
443 attr_info.num_unique_items() as usize,
444 attr_info.branching_factor(),
445 index_begin,
446 feature_begin,
447 1024 * 1024, // combine_request_threshold
448 );
449 multi_index.add_index(col.name().to_string(), index);
450 }
451 ColumnType::Long => {
452 let index = HttpIndex::<i64>::new(
453 attr_info.num_unique_items() as usize,
454 attr_info.branching_factor(),
455 index_begin,
456 feature_begin,
457 1024 * 1024, // combine_request_threshold
458 );
459 multi_index.add_index(col.name().to_string(), index);
460 }
461 ColumnType::ULong => {
462 let index = HttpIndex::<u64>::new(
463 attr_info.num_unique_items() as usize,
464 attr_info.branching_factor(),
465 index_begin,
466 feature_begin,
467 1024 * 1024, // combine_request_threshold
468 );
469 multi_index.add_index(col.name().to_string(), index);
470 }
471 // Byte is stored as u8 by the writer; see reader/attr_query.rs.
472 ColumnType::Byte => {
473 let index = HttpIndex::<u8>::new(
474 attr_info.num_unique_items() as usize,
475 attr_info.branching_factor(),
476 index_begin,
477 feature_begin,
478 1024 * 1024, // combine_request_threshold
479 );
480 multi_index.add_index(col.name().to_string(), index);
481 }
482 ColumnType::UByte => {
483 let index = HttpIndex::<u8>::new(
484 attr_info.num_unique_items() as usize,
485 attr_info.branching_factor(),
486 index_begin,
487 feature_begin,
488 1024 * 1024, // combine_request_threshold
489 );
490 multi_index.add_index(col.name().to_string(), index);
491 }
492
493 _ => {
494 println!("Unsupported column type: {:?}", col.type_());
495 return Err(Error::UnsupportedColumnType(col.name().to_string()));
496 }
497 }
498 }
499 Ok(())
500 }
501}
502
503impl<T: AsyncHttpRangeClient + Send + Sync> AsyncFeatureIter<T> {
504 pub fn header(&self) -> Header {
505 self.fbs.header()
506 }
507 /// Number of selected features (might be unknown)
508 pub fn features_count(&self) -> Option<usize> {
509 if self.count > 0 {
510 Some(self.count)
511 } else {
512 None
513 }
514 }
515 /// Read next feature
516 pub async fn next(&mut self) -> Result<Option<&FcbBuffer>> {
517 let Some(buffer) = self.selection.next_feature_buffer(&mut self.client).await? else {
518 return Ok(None);
519 };
520
521 // Not zero-copy
522 self.fbs.features_buf = buffer.to_vec();
523 // verify flatbuffer
524 let _feature = size_prefixed_root_as_city_feature(&self.fbs.features_buf)?;
525 Ok(Some(&self.fbs))
526 }
527 /// Return current feature
528 pub fn cur_feature(&self) -> &FcbBuffer {
529 &self.fbs
530 }
531
532 pub fn cur_cj_feature(&self) -> Result<CityJSONFeature> {
533 let cj_feature = to_cj_feature(
534 self.cur_feature().feature(),
535 self.header().columns(),
536 self.header().semantic_columns(),
537 )?;
538 Ok(cj_feature)
539 }
540}
541
542enum FeatureSelection {
543 SelectAll(SelectAll),
544 SelectBbox(SelectBbox),
545 SelectAttr(SelectAttr),
546}
547
548impl FeatureSelection {
549 async fn next_feature_buffer<T: AsyncHttpRangeClient>(
550 &mut self,
551 client: &mut AsyncBufferedHttpRangeClient<T>,
552 ) -> Result<Option<Bytes>> {
553 match self {
554 FeatureSelection::SelectAll(select_all) => select_all.next_buffer(client).await,
555 FeatureSelection::SelectBbox(select_bbox) => select_bbox.next_buffer(client).await,
556 FeatureSelection::SelectAttr(select_attr) => select_attr.next_buffer(client).await,
557 }
558 }
559}
560
561struct SelectAll {
562 /// Features left
563 features_left: u64,
564
565 /// How many bytes into the file we've read so far
566 pos: usize,
567}
568
569impl SelectAll {
570 async fn next_buffer<T: AsyncHttpRangeClient>(
571 &mut self,
572 client: &mut AsyncBufferedHttpRangeClient<T>,
573 ) -> Result<Option<Bytes>> {
574 client.min_req_size(DEFAULT_HTTP_FETCH_SIZE);
575
576 if self.features_left == 0 {
577 return Ok(None);
578 }
579 self.features_left -= 1;
580
581 let mut feature_buffer = BytesMut::from(client.get_range(self.pos, 4).await?);
582 self.pos += 4;
583 let feature_size = LittleEndian::read_u32(&feature_buffer) as usize;
584 feature_buffer.put(client.get_range(self.pos, feature_size).await?);
585 self.pos += feature_size;
586
587 Ok(Some(feature_buffer.freeze()))
588 }
589}
590
591struct SelectBbox {
592 /// Selected features
593 feature_batches: Vec<FeatureBatch>,
594}
595
596impl SelectBbox {
597 async fn next_buffer<T: AsyncHttpRangeClient>(
598 &mut self,
599 client: &mut AsyncBufferedHttpRangeClient<T>,
600 ) -> Result<Option<Bytes>> {
601 let mut next_buffer = None;
602 while next_buffer.is_none() {
603 let Some(feature_batch) = self.feature_batches.last_mut() else {
604 break;
605 };
606 let Some(buffer) = feature_batch.next_buffer(client).await? else {
607 // done with this batch
608 self.feature_batches
609 .pop()
610 .expect("already asserted feature_batches was non-empty");
611 continue;
612 };
613 next_buffer = Some(buffer)
614 }
615
616 Ok(next_buffer)
617 }
618}
619
620struct FeatureBatch {
621 /// The byte location of each feature within the file
622 feature_ranges: VecDeque<HttpRange>,
623}
624
625impl FeatureBatch {
626 async fn make_batches(
627 feature_ranges: Vec<HttpSearchResultItem>,
628 combine_request_threshold: usize,
629 ) -> Result<Vec<Self>> {
630 let mut batched_ranges = vec![];
631
632 for search_result_item in feature_ranges.into_iter() {
633 let Some(latest_batch) = batched_ranges.last_mut() else {
634 let mut new_batch = VecDeque::new();
635 new_batch.push_back(search_result_item.range);
636 batched_ranges.push(new_batch);
637 continue;
638 };
639
640 let previous_item = latest_batch.back().expect("we never push an empty batch");
641
642 let HttpRange::Range(Range { end: prev_end, .. }) = previous_item else {
643 debug_assert!(false, "This shouldn't happen. Only the very last feature is expected to have an unknown length");
644 let mut new_batch = VecDeque::new();
645 new_batch.push_back(search_result_item.range);
646 batched_ranges.push(new_batch);
647 continue;
648 };
649
650 let wasted_bytes = search_result_item.range.start() - prev_end;
651 if wasted_bytes < combine_request_threshold {
652 latest_batch.push_back(search_result_item.range)
653 } else {
654 debug!("creating a new request for batch rather than wasting {wasted_bytes} bytes");
655 let mut new_batch = VecDeque::new();
656 new_batch.push_back(search_result_item.range);
657 batched_ranges.push(new_batch);
658 }
659 }
660
661 let mut batches: Vec<_> = batched_ranges.into_iter().map(FeatureBatch::new).collect();
662 batches.reverse();
663 Ok(batches)
664 }
665
666 fn new(feature_ranges: VecDeque<HttpRange>) -> Self {
667 Self { feature_ranges }
668 }
669
670 /// When fetching new data, how many bytes should we fetch at once.
671 /// It was computed based on the specific feature ranges of the batch
672 /// to optimize number of requests vs. wasted bytes vs. resident memory
673 fn request_size(&self) -> usize {
674 let Some(first) = self.feature_ranges.front() else {
675 return 0;
676 };
677 let Some(last) = self.feature_ranges.back() else {
678 return 0;
679 };
680
681 // `last.length()` should only be None if this batch includes the final feature
682 // in the dataset. Since we can't infer its actual length, we'll fetch only
683 // the first 4 bytes of that feature buffer, which will tell us the actual length
684 // of the feature buffer for the subsequent request.
685 let last_feature_length = last.length().unwrap_or(4);
686
687 let covering_range = first.start()..last.start() + last_feature_length;
688
689 covering_range
690 .len()
691 // Since it's all held in memory, don't fetch more than DEFAULT_HTTP_FETCH_SIZE at a time
692 // unless necessary.
693 .min(DEFAULT_HTTP_FETCH_SIZE)
694 }
695
696 async fn next_buffer<T: AsyncHttpRangeClient>(
697 &mut self,
698 client: &mut AsyncBufferedHttpRangeClient<T>,
699 ) -> Result<Option<Bytes>> {
700 let request_size = self.request_size();
701 client.set_min_req_size(request_size);
702 let Some(feature_range) = self.feature_ranges.pop_front() else {
703 return Ok(None);
704 };
705
706 let mut pos = feature_range.start();
707 let mut feature_buffer = BytesMut::from(client.get_range(pos, 4).await?);
708 pos += 4;
709 let feature_size = LittleEndian::read_u32(&feature_buffer) as usize;
710 feature_buffer.put(client.get_range(pos, feature_size).await?);
711
712 Ok(Some(feature_buffer.freeze()))
713 }
714}
715
716struct SelectAttr {
717 // TODO: change this implementation so it can batch features
718 ranges: Vec<HttpRange>,
719 range_pos: usize,
720}
721
722impl SelectAttr {
723 async fn next_buffer<T: AsyncHttpRangeClient>(
724 &mut self,
725 client: &mut AsyncBufferedHttpRangeClient<T>,
726 ) -> Result<Option<Bytes>> {
727 let Some(range) = self.ranges.get(self.range_pos) else {
728 return Ok(None);
729 };
730 let mut feature_buffer = BytesMut::from(client.get_range(range.start(), 4).await?);
731 let feature_size = LittleEndian::read_u32(&feature_buffer) as usize;
732 feature_buffer.put(client.get_range(range.start() + 4, feature_size).await?);
733 self.range_pos += 1;
734 Ok(Some(feature_buffer.freeze()))
735 }
736}
737
738#[cfg(test)]
739mod index_node_size_tests {
740 //! The HTTP spatial path must traverse the R-tree with the branching
741 //! factor recorded in the header, not the compile-time default.
742 //!
743 //! These tests live inside the crate because `mock_from_file` (and the
744 //! `MockHttpRangeClient` it returns a reader over) is `#[cfg(test)]`-only
745 //! and therefore unreachable from `tests/`.
746
747 use super::*;
748 use crate::http_reader::mock_http_range_client::MockHttpRangeClient;
749 use std::path::PathBuf;
750
751 /// A fixture from the shared conformance corpus, resolved from the crate
752 /// root so the test does not depend on the process working directory.
753 fn conformance_fixture(name: &str) -> String {
754 PathBuf::from(env!("CARGO_MANIFEST_DIR"))
755 .join("../../../conformance")
756 .join(name)
757 .to_str()
758 .expect("fixture path is valid UTF-8")
759 .to_string()
760 }
761
762 async fn feature_ids(iter: &mut AsyncFeatureIter<MockHttpRangeClient>) -> Result<Vec<String>> {
763 let mut ids = Vec::new();
764 while let Some(feature) = iter.next().await? {
765 ids.push(feature.cj_feature()?.id);
766 }
767 Ok(ids)
768 }
769
770 /// `appearance_depths_node8.fcb` holds 12 features at node size 8, so its
771 /// level bounds are `[12, 2, 1]` — three levels and two sibling leaf
772 /// ranges. Read with the default node size 16 the same 12 features imply
773 /// `[12, 1]`, a different node count and different level boundaries, so a
774 /// hardcoded 16 walks the wrong ranges.
775 ///
776 /// The oracle is `select_all`, which never touches the R-tree: a bbox
777 /// covering the whole extent must return exactly the same feature set.
778 #[tokio::test]
779 async fn http_spatial_query_honours_header_index_node_size() -> Result<()> {
780 let path = conformance_fixture("appearance_depths_node8.fcb");
781
782 let (reader, _stats) = HttpFcbReader::mock_from_file(&path).await?;
783 let (node_size, feature_count, minx, miny, maxx, maxy) = {
784 let header = reader.header();
785 let extent = header
786 .geographical_extent()
787 .expect("fixture carries a geographical extent");
788 (
789 header.index_node_size(),
790 header.features_count() as usize,
791 extent.min().x(),
792 extent.min().y(),
793 extent.max().x(),
794 extent.max().y(),
795 )
796 };
797 assert_eq!(
798 node_size, 8,
799 "fixture must have a NON-default node size, else a hardcoded 16 passes vacuously"
800 );
801 assert_eq!(feature_count, 12);
802
803 // Oracle: the full feature set, read sequentially without the R-tree.
804 let (all_reader, _stats) = HttpFcbReader::mock_from_file(&path).await?;
805 let mut all_iter = all_reader.select_all().await?;
806 let mut expected_ids = feature_ids(&mut all_iter).await?;
807 expected_ids.sort();
808 assert_eq!(expected_ids.len(), feature_count);
809
810 let mut query_iter = reader
811 .select_query(Query::BBox(minx, miny, maxx, maxy))
812 .await?;
813 assert_eq!(query_iter.features_count(), Some(feature_count));
814 let mut hit_ids = feature_ids(&mut query_iter).await?;
815 assert_eq!(
816 hit_ids
817 .iter()
818 .collect::<std::collections::HashSet<_>>()
819 .len(),
820 hit_ids.len(),
821 "a spatial query must return a set, not duplicates"
822 );
823 hit_ids.sort();
824 assert_eq!(hit_ids, expected_ids);
825 Ok(())
826 }
827}
828
829//TODO: Fix this test. It's failling bc of the mock client and payload cache.
830// #[cfg(test)]
831// mod tests {
832// use std::{path::PathBuf, str::FromStr};
833
834// use cjseq::CityJSONFeature;
835// use static_btree::{FixedStringKey, Float, KeyType, Operator};
836
837// use crate::error::Result;
838// use crate::HttpFcbReader;
839
840// #[tokio::test]
841// async fn fcb_http_reader_test() -> Result<()> {
842// #[derive(Debug)]
843// struct QueryTestCase {
844// test_name: &'static str,
845// query: Vec<(String, Operator, KeyType)>,
846// expected_count: usize,
847// validator: fn(&CityJSONFeature) -> bool,
848// }
849
850// let test_cases = vec![
851// // Test case: Expect one matching feature with b3_h_dak_50p > 2.0 and matching identificatie.
852// QueryTestCase {
853// test_name: "test_attr_index_multiple_queries: b3_h_dak_50p > 2.0 and identificatie == NL.IMBAG.Pand.0503100000012869",
854// query: vec![
855// (
856// "b3_h_dak_50p".to_string(),
857// Operator::Gt,
858// KeyType::Float64(Float::<f64>(2.0)),
859// ),
860// (
861// "identificatie".to_string(),
862// Operator::Eq,
863// KeyType::StringKey50(FixedStringKey::from_str(
864// "NL.IMBAG.Pand.0503100000012869",
865// )),
866// ),
867// ],
868// expected_count: 1,
869// validator: |feature: &CityJSONFeature| {
870// let mut valid_b3 = false;
871// let mut valid_ident = false;
872// for co in feature.city_objects.values() {
873// if let Some(attrs) = &co.attributes {
874// if let Some(val) = attrs.get("b3_h_dak_50p") {
875// if val.as_f64().unwrap() > 2.0 {
876// valid_b3 = true;
877// }
878// }
879// if let Some(ident) = attrs.get("identificatie") {
880// if ident.as_str().unwrap() == "NL.IMBAG.Pand.0503100000012869" {
881// valid_ident = true;
882// }
883// }
884// }
885// }
886// valid_b3 && valid_ident
887// },
888// },
889// // Test case: Expect zero features where tijdstipregistratie is before 2008-01-01.
890// QueryTestCase {
891// test_name: "test_attr_index_multiple_queries: tijdstipregistratie < 2008-01-01",
892// query: vec![(
893// "tijdstipregistratie".to_string(),
894// Operator::Lt,
895// KeyType::DateTime(chrono::DateTime::<chrono::Utc>::from_str(
896// "2008-01-01T00:00:00Z",
897// )
898// .unwrap()),
899// )],
900// expected_count: 0,
901// validator: |feature: &CityJSONFeature| {
902// let mut valid_tijdstip = true;
903// let query_tijdstip = chrono::NaiveDate::from_ymd(2008, 1, 1).and_hms(0, 0, 0);
904// for co in feature.city_objects.values() {
905// if let Some(attrs) = &co.attributes {
906// if let Some(val) = attrs.get("tijdstipregistratie") {
907// let val_tijdstip = chrono::NaiveDateTime::parse_from_str(
908// val.as_str().unwrap(),
909// "%Y-%m-%dT%H:%M:%S",
910// )
911// .unwrap();
912// if val_tijdstip < query_tijdstip {
913// valid_tijdstip = false;
914// }
915// }
916// }
917// }
918// valid_tijdstip
919// },
920// },
921// // Test case: Expect zero features where tijdstipregistratie is after 2008-01-01.
922// QueryTestCase {
923// test_name: "test_attr_index_multiple_queries: tijdstipregistratie > 2008-01-01",
924// query: vec![(
925// "tijdstipregistratie".to_string(),
926// Operator::Gt,
927// KeyType::DateTime(chrono::DateTime::<chrono::Utc>::from_utc(
928// chrono::NaiveDate::from_ymd(2008, 1, 1).and_hms(0, 0, 0),
929// chrono::Utc,
930// )),
931// )],
932// expected_count: 3,
933// validator: |feature: &CityJSONFeature| {
934// let mut valid_tijdstip = false;
935// let query_tijdstip = chrono::NaiveDate::from_ymd(2008, 1, 1).and_hms(0, 0, 0);
936// for co in feature.city_objects.values() {
937// if let Some(attrs) = &co.attributes {
938// if let Some(val) = attrs.get("tijdstipregistratie") {
939// let val_tijdstip =
940// chrono::DateTime::parse_from_rfc3339(val.as_str().unwrap())
941// .map_err(|e| eprintln!("Failed to parse datetime: {}", e))
942// .map(|dt| dt.naive_utc())
943// .unwrap_or_else(|_| {
944// chrono::NaiveDateTime::from_timestamp_opt(0, 0).unwrap()
945// });
946// if val_tijdstip > query_tijdstip {
947// valid_tijdstip = true;
948// }
949// }
950// }
951// }
952// valid_tijdstip
953// },
954// },
955// ];
956
957// for test_case in test_cases {
958// let manifest_dir = PathBuf::from(env!("CARGO_MANIFEST_DIR"));
959// let input_file_path = manifest_dir.join("tests/data/small.fcb");
960
961// let (fcb, stats) = HttpFcbReader::mock_from_file(&input_file_path.to_str().unwrap())
962// .await
963// .unwrap();
964
965// // {
966// // // The read guard needs to be in a scoped block, else we won't release the lock and the test will hang when
967// // // the actual FGB client code tries to update the stats.
968// // let stats = stats.read().unwrap();
969// // assert_eq!(stats.request_count, 1);
970// // // This number might change a little if the test data or logic changes, but they should be in the same ballpark.
971// // assert_eq!(stats.bytes_requested, 12944);
972// // }
973
974// let query = test_case.query;
975// let mut iter = fcb.select_attr_query(&query).await.unwrap();
976
977// let mut features = Vec::new();
978// while let Some(feat_buf) = iter.next().await.unwrap() {
979// let feature = feat_buf.cj_feature()?;
980// features.push(feature);
981// }
982// assert_eq!(features.len(), test_case.expected_count);
983
984// for feature in features {
985// assert!(
986// (test_case.validator)(&feature),
987// "Failed to validate feature in test case: {}",
988// test_case.test_name
989// );
990// }
991// }
992
993// // {
994// // let stats = stats.read().unwrap();
995
996// // assert_eq!(stats.request_count, 5);
997// // assert_eq!(stats.bytes_requested, 2131152);
998// // }
999// Ok(())
1000// }
1001// }