From 421a8270b5dbc2c5f55c06ab25959c3d684f9fb5 Mon Sep 17 00:00:00 2001 From: ColinLee Date: Sat, 10 Oct 2026 16:42:35 +0800 Subject: [PATCH] fix(cpp): apply residual row offsets in non-aligned page decoding --- cpp/src/reader/chunk_reader.cc | 156 +++++++++++--------- cpp/src/reader/chunk_reader.h | 19 +-- cpp/test/reader/prepared_series_test.cc | 180 +++++++++++++++++++++++- python/tests/test_tsfile_dataset.py | 32 +++++ 4 files changed, 313 insertions(+), 74 deletions(-) diff --git a/cpp/src/reader/chunk_reader.cc b/cpp/src/reader/chunk_reader.cc index abeff8375..d713570e7 100644 --- a/cpp/src/reader/chunk_reader.cc +++ b/cpp/src/reader/chunk_reader.cc @@ -315,7 +315,7 @@ int ChunkReader::skip_cur_page() { } int ChunkReader::decode_cur_page_data(TsBlock*& ret_tsblock, Filter* filter, - PageArena& pa) { + PageArena& pa, int* row_offset) { int ret = E_OK; // Step 1: make sure we load the whole page data in @in_stream_ @@ -396,8 +396,8 @@ int ChunkReader::decode_cur_page_data(TsBlock*& ret_tsblock, Filter* filter, // ret = decode_tv_buf_into_tsblock(time_buf, value_buf, time_buf_size, // value_buf_size, ret_tsblock, // filter); - ret = decode_tv_buf_into_tsblock_by_datatype(time_in_, value_in_, - ret_tsblock, filter, &pa); + ret = decode_tv_buf_into_tsblock_by_datatype( + time_in_, value_in_, ret_tsblock, filter, &pa, row_offset); // if we return during @decode_tv_buf_into_tsblock, we should keep // @uncompressed_buf_ valid until all TV pairs are decoded. if (ret != E_OVERFLOW) { @@ -414,29 +414,32 @@ int ChunkReader::decode_cur_page_data(TsBlock*& ret_tsblock, Filter* filter, return ret; } -#define DECODE_TYPED_TV_INTO_TSBLOCK(CppType, ReadType, time_in, value_in, \ - row_appender) \ - do { \ - int64_t time = 0; \ - CppType value; \ - while (time_decoder_->has_remaining(time_in)) { \ - ASSERT(value_decoder_->has_remaining(value_in)); \ - if (UNLIKELY(!row_appender.add_row())) { \ - ret = E_OVERFLOW; \ - break; \ - } else if (RET_FAIL(time_decoder_->read_int64(time, time_in))) { \ - } else if (RET_FAIL(value_decoder_->read_##ReadType(value, \ - value_in))) { \ - } else if (filter != nullptr && !filter->satisfy(time, value)) { \ - row_appender.backoff_add_row(); \ - continue; \ - } else { \ - /*std::cout << "decoder: time=" << time << ", value=" << value \ - * << std::endl;*/ \ - row_appender.append(0, (char*)&time, sizeof(time)); \ - row_appender.append(1, (char*)&value, sizeof(value)); \ - } \ - } \ +#define DECODE_TYPED_TV_INTO_TSBLOCK(CppType, ReadType, time_in, value_in, \ + row_appender) \ + do { \ + int64_t time = 0; \ + CppType value; \ + while (time_decoder_->has_remaining(time_in)) { \ + ASSERT(value_decoder_->has_remaining(value_in)); \ + if (UNLIKELY(!row_appender.add_row())) { \ + ret = E_OVERFLOW; \ + break; \ + } else if (RET_FAIL(time_decoder_->read_int64(time, time_in))) { \ + } else if (RET_FAIL(value_decoder_->read_##ReadType(value, \ + value_in))) { \ + } else if (filter != nullptr && !filter->satisfy(time, value)) { \ + row_appender.backoff_add_row(); \ + continue; \ + } else { \ + if (row_offset > 0) { \ + --row_offset; \ + row_appender.backoff_add_row(); \ + continue; \ + } \ + row_appender.append(0, (char*)&time, sizeof(time)); \ + row_appender.append(1, (char*)&value, sizeof(value)); \ + } \ + } \ } while (false) int ChunkReader::i32_DECODE_TYPED_TV_INTO_TSBLOCK(ByteStream& time_in, @@ -467,19 +470,19 @@ int ChunkReader::i32_DECODE_TYPED_TV_INTO_TSBLOCK(ByteStream& time_in, } int ChunkReader::i32_DECODE_TV_BATCH(ByteStream& time_in, ByteStream& value_in, - RowAppender& row_appender, - Filter* filter) { + RowAppender& row_appender, Filter* filter, + int& row_offset) { int ret = E_OK; const int BATCH = 129; int64_t times[BATCH]; int32_t values[BATCH]; while (time_decoder_->has_remaining(time_in)) { - // Cap each pass to what the appender can still hold; the old - // "remaining < BATCH → OVERFLOW" check made progress impossible on - // TsBlocks with capacity below BATCH. - int eff_batch = - std::min(BATCH, static_cast(row_appender.remaining())); + // Skipped rows consume no output capacity. Bound the batch so every + // accepted row fits after the remaining offset has been consumed. + int eff_batch = std::min( + BATCH, + std::max(static_cast(row_appender.remaining()), row_offset)); if (eff_batch <= 0) { ret = E_OVERFLOW; break; @@ -543,6 +546,10 @@ int ChunkReader::i32_DECODE_TV_BATCH(ByteStream& time_in, ByteStream& value_in, !filter->satisfy(times[i], (int64_t)values[i])) { continue; } + if (row_offset > 0) { + --row_offset; + continue; + } if (UNLIKELY(!row_appender.add_row())) { ret = E_OVERFLOW; break; @@ -556,16 +563,17 @@ int ChunkReader::i32_DECODE_TV_BATCH(ByteStream& time_in, ByteStream& value_in, } int ChunkReader::i64_DECODE_TV_BATCH(ByteStream& time_in, ByteStream& value_in, - RowAppender& row_appender, - Filter* filter) { + RowAppender& row_appender, Filter* filter, + int& row_offset) { int ret = E_OK; const int BATCH = 129; int64_t times[BATCH]; int64_t values[BATCH]; while (time_decoder_->has_remaining(time_in)) { - int eff_batch = - std::min(BATCH, static_cast(row_appender.remaining())); + int eff_batch = std::min( + BATCH, + std::max(static_cast(row_appender.remaining()), row_offset)); if (eff_batch <= 0) { ret = E_OVERFLOW; break; @@ -629,6 +637,10 @@ int ChunkReader::i64_DECODE_TV_BATCH(ByteStream& time_in, ByteStream& value_in, !filter->satisfy(times[i], values[i])) { continue; } + if (row_offset > 0) { + --row_offset; + continue; + } if (UNLIKELY(!row_appender.add_row())) { ret = E_OVERFLOW; break; @@ -644,15 +656,16 @@ int ChunkReader::i64_DECODE_TV_BATCH(ByteStream& time_in, ByteStream& value_in, int ChunkReader::float_DECODE_TV_BATCH(ByteStream& time_in, ByteStream& value_in, RowAppender& row_appender, - Filter* filter) { + Filter* filter, int& row_offset) { int ret = E_OK; const int BATCH = 129; int64_t times[BATCH]; float values[BATCH]; while (time_decoder_->has_remaining(time_in)) { - int eff_batch = - std::min(BATCH, static_cast(row_appender.remaining())); + int eff_batch = std::min( + BATCH, + std::max(static_cast(row_appender.remaining()), row_offset)); if (eff_batch <= 0) { ret = E_OVERFLOW; break; @@ -712,6 +725,10 @@ int ChunkReader::float_DECODE_TV_BATCH(ByteStream& time_in, if (filter != nullptr && !block_all_pass && !time_mask[i]) { continue; } + if (row_offset > 0) { + --row_offset; + continue; + } if (UNLIKELY(!row_appender.add_row())) { ret = E_OVERFLOW; break; @@ -727,15 +744,16 @@ int ChunkReader::float_DECODE_TV_BATCH(ByteStream& time_in, int ChunkReader::double_DECODE_TV_BATCH(ByteStream& time_in, ByteStream& value_in, RowAppender& row_appender, - Filter* filter) { + Filter* filter, int& row_offset) { int ret = E_OK; const int BATCH = 129; int64_t times[BATCH]; double values[BATCH]; while (time_decoder_->has_remaining(time_in)) { - int eff_batch = - std::min(BATCH, static_cast(row_appender.remaining())); + int eff_batch = std::min( + BATCH, + std::max(static_cast(row_appender.remaining()), row_offset)); if (eff_batch <= 0) { ret = E_OVERFLOW; break; @@ -795,6 +813,10 @@ int ChunkReader::double_DECODE_TV_BATCH(ByteStream& time_in, if (filter != nullptr && !block_all_pass && !time_mask[i]) { continue; } + if (row_offset > 0) { + --row_offset; + continue; + } if (UNLIKELY(!row_appender.add_row())) { ret = E_OVERFLOW; break; @@ -807,11 +829,9 @@ int ChunkReader::double_DECODE_TV_BATCH(ByteStream& time_in, return ret; } -int ChunkReader::STRING_DECODE_TYPED_TV_INTO_TSBLOCK(ByteStream& time_in, - ByteStream& value_in, - RowAppender& row_appender, - PageArena& pa, - Filter* filter) { +int ChunkReader::STRING_DECODE_TYPED_TV_INTO_TSBLOCK( + ByteStream& time_in, ByteStream& value_in, RowAppender& row_appender, + PageArena& pa, Filter* filter, int& row_offset) { int ret = E_OK; int64_t time = 0; common::String value; @@ -825,6 +845,10 @@ int ChunkReader::STRING_DECODE_TYPED_TV_INTO_TSBLOCK(ByteStream& time_in, } else if (filter != nullptr && !filter->satisfy(time, value)) { row_appender.backoff_add_row(); continue; + } else if (row_offset > 0) { + --row_offset; + row_appender.backoff_add_row(); + continue; } else { row_appender.append(0, (char*)&time, sizeof(time)); row_appender.append(1, value.buf_, value.len_); @@ -833,12 +857,13 @@ int ChunkReader::STRING_DECODE_TYPED_TV_INTO_TSBLOCK(ByteStream& time_in, return ret; } -int ChunkReader::decode_tv_buf_into_tsblock_by_datatype(ByteStream& time_in, - ByteStream& value_in, - TsBlock* ret_tsblock, - Filter* filter, - common::PageArena* pa) { +int ChunkReader::decode_tv_buf_into_tsblock_by_datatype( + ByteStream& time_in, ByteStream& value_in, TsBlock* ret_tsblock, + Filter* filter, common::PageArena* pa, int* remaining_offset) { int ret = E_OK; + int unused_offset = 0; + int& row_offset = + remaining_offset == nullptr ? unused_offset : *remaining_offset; RowAppender row_appender(ret_tsblock); switch (chunk_header_.data_type_) { case common::BOOLEAN: @@ -847,33 +872,36 @@ int ChunkReader::decode_tv_buf_into_tsblock_by_datatype(ByteStream& time_in, break; case common::DATE: case common::INT32: - ret = - i32_DECODE_TV_BATCH(time_in_, value_in_, row_appender, filter); + ret = i32_DECODE_TV_BATCH(time_in_, value_in_, row_appender, filter, + row_offset); break; case TIMESTAMP: case common::INT64: - ret = - i64_DECODE_TV_BATCH(time_in_, value_in_, row_appender, filter); + ret = i64_DECODE_TV_BATCH(time_in_, value_in_, row_appender, filter, + row_offset); break; case common::FLOAT: ret = float_DECODE_TV_BATCH(time_in_, value_in_, row_appender, - filter); + filter, row_offset); break; case common::DOUBLE: ret = double_DECODE_TV_BATCH(time_in_, value_in_, row_appender, - filter); + filter, row_offset); break; case common::TEXT: case common::BLOB: case common::STRING: ret = STRING_DECODE_TYPED_TV_INTO_TSBLOCK( - time_in, value_in, row_appender, *pa, filter); + time_in, value_in, row_appender, *pa, filter, row_offset); break; default: ret = E_NOT_SUPPORT; ASSERT(false); } - if (ret_tsblock->get_row_count() == 0 && ret == E_OK) { + // The offset-aware iterator must keep scanning after a page whose + // matching rows were all consumed by the offset (or time filter). + if (remaining_offset == nullptr && ret_tsblock->get_row_count() == 0 && + ret == E_OK) { ret = E_NO_MORE_DATA; } return ret; @@ -917,8 +945,8 @@ int ChunkReader::get_next_page(TsBlock* ret_tsblock, Filter* oneshoot_filter, } if (prev_page_not_finish()) { - ret = decode_tv_buf_into_tsblock_by_datatype(time_in_, value_in_, - ret_tsblock, filter, &pa); + ret = decode_tv_buf_into_tsblock_by_datatype( + time_in_, value_in_, ret_tsblock, filter, &pa, &row_offset); if (ret == E_OVERFLOW) { ret = E_OK; } else { @@ -953,7 +981,7 @@ int ChunkReader::get_next_page(TsBlock* ret_tsblock, Filter* oneshoot_filter, } if (IS_SUCC(ret)) { - ret = decode_cur_page_data(ret_tsblock, filter, pa); + ret = decode_cur_page_data(ret_tsblock, filter, pa, &row_offset); } return ret; } diff --git a/cpp/src/reader/chunk_reader.h b/cpp/src/reader/chunk_reader.h index b54301199..e49f727a5 100644 --- a/cpp/src/reader/chunk_reader.h +++ b/cpp/src/reader/chunk_reader.h @@ -91,7 +91,7 @@ class ChunkReader : public IChunkReader { bool cur_page_fully_satisfies_filter(Filter* filter); int skip_cur_page(); int decode_cur_page_data(common::TsBlock*& ret_tsblock, Filter* filter, - common::PageArena& pa); + common::PageArena& pa, int* row_offset = nullptr); bool prev_page_not_finish() const { return (time_decoder_ && time_decoder_->has_remaining(time_in_)) || time_in_.has_remaining(); @@ -101,30 +101,33 @@ class ChunkReader : public IChunkReader { common::ByteStream& value_in, common::TsBlock* ret_tsblock, Filter* filter, - common::PageArena* pa = nullptr); + common::PageArena* pa = nullptr, + int* remaining_offset = nullptr); int i32_DECODE_TYPED_TV_INTO_TSBLOCK(common::ByteStream& time_in, common::ByteStream& value_in, common::RowAppender& row_appender, Filter* filter); int i32_DECODE_TV_BATCH(common::ByteStream& time_in, common::ByteStream& value_in, - common::RowAppender& row_appender, Filter* filter); + common::RowAppender& row_appender, Filter* filter, + int& row_offset); int i64_DECODE_TV_BATCH(common::ByteStream& time_in, common::ByteStream& value_in, - common::RowAppender& row_appender, Filter* filter); + common::RowAppender& row_appender, Filter* filter, + int& row_offset); int float_DECODE_TV_BATCH(common::ByteStream& time_in, common::ByteStream& value_in, - common::RowAppender& row_appender, - Filter* filter); + common::RowAppender& row_appender, Filter* filter, + int& row_offset); int double_DECODE_TV_BATCH(common::ByteStream& time_in, common::ByteStream& value_in, common::RowAppender& row_appender, - Filter* filter); + Filter* filter, int& row_offset); int STRING_DECODE_TYPED_TV_INTO_TSBLOCK(common::ByteStream& time_in, common::ByteStream& value_in, common::RowAppender& row_appender, common::PageArena& pa, - Filter* filter); + Filter* filter, int& row_offset); private: RandomAccessReadFile* read_file_; diff --git a/cpp/test/reader/prepared_series_test.cc b/cpp/test/reader/prepared_series_test.cc index cd1389fbb..d90101fd9 100644 --- a/cpp/test/reader/prepared_series_test.cc +++ b/cpp/test/reader/prepared_series_test.cc @@ -22,6 +22,7 @@ #include #include +#include #include #include @@ -32,22 +33,29 @@ #include "reader/table_result_set.h" #include "reader/tsfile_reader.h" #include "writer/tsfile_table_writer.h" +#include "writer/tsfile_writer.h" namespace storage { namespace { class PagePointGuard { public: - explicit PagePointGuard(uint32_t page_points) - : saved_(common::g_config_value_.page_writer_max_point_num_) { + explicit PagePointGuard(uint32_t page_points, uint32_t page_bytes = 0) + : saved_(common::g_config_value_.page_writer_max_point_num_), + saved_bytes_(common::g_config_value_.page_writer_max_memory_bytes_) { common::g_config_value_.page_writer_max_point_num_ = page_points; + if (page_bytes != 0) { + common::g_config_value_.page_writer_max_memory_bytes_ = page_bytes; + } } ~PagePointGuard() { common::g_config_value_.page_writer_max_point_num_ = saved_; + common::g_config_value_.page_writer_max_memory_bytes_ = saved_bytes_; } private: uint32_t saved_; + uint32_t saved_bytes_; }; class PreparedSeriesBatchTest : public ::testing::Test { @@ -59,6 +67,7 @@ class PreparedSeriesBatchTest : public ::testing::Test { file_name_ = std::string("prepared_series_batch_test_") + (test_info == nullptr ? "unknown" : test_info->name()) + ".tsfile"; + std::replace(file_name_.begin(), file_name_.end(), '/', '_'); std::remove(file_name_.c_str()); } @@ -111,6 +120,173 @@ class PreparedSeriesBatchTest : public ::testing::Test { std::string file_name_ = "prepared_series_batch_test.tsfile"; }; +class NonAlignedPreparedSeriesOffsetTest + : public PreparedSeriesBatchTest, + public ::testing::WithParamInterface {}; + +TEST_P(NonAlignedPreparedSeriesOffsetTest, AppliesRowOffset) { + // Include pages larger than the 65536-row output block so the decoder + // resumes an unfinished page on the next read. + PagePointGuard guard(GetParam(), 16 * 1024 * 1024); + const std::string device = "root.offset"; + const std::vector names = {"boolean", "int32", "int64", + "float", "double", "string"}; + const std::vector types = { + common::BOOLEAN, common::INT32, common::INT64, + common::FLOAT, common::DOUBLE, common::STRING}; + auto schemas = std::make_shared>(); + TsFileWriter writer; + ASSERT_EQ(common::E_OK, writer.open(file_name_)); + for (size_t i = 0; i < names.size(); ++i) { + schemas->emplace_back(names[i], types[i], common::PLAIN, + common::UNCOMPRESSED); + ASSERT_EQ(common::E_OK, + writer.register_timeseries(device, schemas->back())); + } + // Flush halfway through to exercise both page and chunk skipping. + for (int start : {0, 70000}) { + Tablet tablet(device, schemas, 70000); + for (int row = 0; row < 70000; ++row) { + const int value = start + row; + ASSERT_EQ(common::E_OK, tablet.add_timestamp(row, value)); + ASSERT_EQ(common::E_OK, tablet.add_value(row, 0u, value % 2 == 0)); + ASSERT_EQ(common::E_OK, tablet.add_value(row, 1u, int32_t(value))); + ASSERT_EQ(common::E_OK, tablet.add_value(row, 2u, int64_t(value))); + ASSERT_EQ(common::E_OK, + tablet.add_value(row, 3u, float(value) + 0.5f)); + ASSERT_EQ(common::E_OK, + tablet.add_value(row, 4u, double(value) + 0.5)); + const std::string text = std::to_string(value); + ASSERT_EQ(common::E_OK, tablet.add_value(row, 5u, text.c_str())); + } + ASSERT_EQ(common::E_OK, writer.write_tablet(tablet)); + ASSERT_EQ(common::E_OK, writer.flush()); + } + ASSERT_EQ(common::E_OK, writer.close()); + + TsFileReader reader; + ASSERT_EQ(common::E_OK, reader.open(file_name_)); + FileGeneration generation; + generation.mapped_index_identity = 1; + generation.file_id = 0; + struct stat file_stat {}; + ASSERT_EQ(0, stat(file_name_.c_str(), &file_stat)); + generation.file_size = static_cast(file_stat.st_size); + generation.file_fingerprint = 0; + + struct Window { + int64_t start; + int64_t end; + int offset; + int limit; + }; + const std::vector windows = { + {0, 139999, 30, 10}, // Offset inside the first page. + {0, 139999, 10003, 1}, // Output capacity below the decode batch size. + {0, 139999, 70003, 7}, // Whole chunk plus a residual. + {0, 139999, 13, 65537}, // Resumes decoding a large page. + {0, 139999, 69995, 10}, // Window crossing a chunk boundary. + {0, 139999, 139990, -1}, // Unlimited tail. + {0, 139999, 140001, 5}, // Offset beyond the data. + {99, 150, 5, 7}, // Count only rows passing the time filter. + {99, 150, 60, 7}, // Offset beyond the filtered rows. + {9998, 10005, 3, 4}, // Filter and offset crossing a page boundary. + {0, 139999, 30, 0}, + }; + auto metadata = reader.get_timeseries_metadata(); + ASSERT_EQ(1U, metadata.size()); + ASSERT_EQ(names.size(), metadata.begin()->second.size()); + for (const auto& device_entry : metadata) { + for (const auto& index : device_entry.second) { + auto* series_index = dynamic_cast(index.get()); + ASSERT_NE(nullptr, series_index); + PreparedLocator locator; + locator.layout = 0; + locator.value_metadata_offset = series_index->get_metadata_offset(); + locator.value_metadata_length = series_index->get_metadata_length(); + std::shared_ptr prepared; + ASSERT_EQ(common::E_OK, + reader.prepare_series(generation, locator, prepared)); + for (const auto& window : windows) { + SCOPED_TRACE(index->get_measurement_name().to_std_string()); + SCOPED_TRACE(window.offset); + ResultSet* result = nullptr; + ASSERT_EQ( + common::E_OK, + reader.query_prepared(prepared, window.start, window.end, + window.offset, window.limit, result)); + auto* table = dynamic_cast(result); + ASSERT_NE(nullptr, table); + const int available = std::max( + 0, int(window.end - window.start + 1) - window.offset); + const int expected_count = + window.limit < 0 ? available + : std::min(available, window.limit); + int count = 0; + common::TsBlock* block = nullptr; + int ret = common::E_OK; + while ((ret = table->get_next_tsblock(block)) == common::E_OK) { + common::RowIterator rows(block); + while (rows.has_next()) { + const int64_t expected = + window.start + window.offset + count; + uint32_t len = 0; + bool is_null = false; + const char* time = rows.read(0, &len, &is_null); + ASSERT_FALSE(is_null); + ASSERT_EQ(expected, + *reinterpret_cast(time)); + const char* value = rows.read(1, &len, &is_null); + ASSERT_FALSE(is_null); + switch (index->get_data_type()) { + case common::BOOLEAN: + EXPECT_EQ( + expected % 2 == 0, + *reinterpret_cast(value)); + break; + case common::INT32: + EXPECT_EQ( + expected, + *reinterpret_cast(value)); + break; + case common::INT64: + EXPECT_EQ( + expected, + *reinterpret_cast(value)); + break; + case common::FLOAT: + EXPECT_FLOAT_EQ( + expected + 0.5f, + *reinterpret_cast(value)); + break; + case common::DOUBLE: + EXPECT_DOUBLE_EQ( + expected + 0.5, + *reinterpret_cast(value)); + break; + case common::STRING: + EXPECT_EQ(std::to_string(expected), + std::string(value, len)); + break; + default: + FAIL() << "Unexpected type"; + } + ++count; + rows.next(); + } + } + EXPECT_EQ(common::E_NO_MORE_DATA, ret); + EXPECT_EQ(expected_count, count); + reader.destroy_query_data_set(result); + } + } + } + EXPECT_EQ(common::E_OK, reader.close()); +} + +INSTANTIATE_TEST_SUITE_P(PageSizes, NonAlignedPreparedSeriesOffsetTest, + ::testing::Values(10000U, 100000U)); + TEST_F(PreparedSeriesBatchTest, PreparedQueryReturnsDirectTableResultSetBatches) { write_nullable_table(); diff --git a/python/tests/test_tsfile_dataset.py b/python/tests/test_tsfile_dataset.py index 02baa9bec..964dde6f8 100644 --- a/python/tests/test_tsfile_dataset.py +++ b/python/tests/test_tsfile_dataset.py @@ -2329,6 +2329,38 @@ def _write_tree_rows(path, device_measurements, t_start=0, t_count=3): writer.close() +@pytest.mark.parametrize("dtype", [TSDataType.INT32, TSDataType.DOUBLE]) +def test_dataset_tree_model_row_slices_apply_page_offset( + tmp_path, dataframe_use_index, dtype +): + paths = [tmp_path / "part0.tsfile", tmp_path / "part1.tsfile"] + for start, path in zip((0, 40), paths): + _write_tree_rows( + path, {"root.offset": [("value", dtype)]}, t_start=start, t_count=40 + ) + expected = np.arange(80, dtype=np.float64) + if dtype == TSDataType.DOUBLE: + expected += 0.5 + + # Check both the initial build and reopening the persisted index. + for _ in range(2): + with TsFileDataFrame( + [str(path) for path in paths], + show_progress=False, + use_index=dataframe_use_index, + ) as dataframe: + series = dataframe["root.offset.value"] + assert series[30] == expected[30] + for window in ( + slice(30, 40), + slice(35, 55), + slice(43, 47), + slice(-10, None), + slice(30, 40, 3), + ): + np.testing.assert_array_equal(series[window], expected[window]) + + def test_dataset_tree_model_merges_identical_structure_across_files( tmp_path, dataframe_use_index ):