Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 6 additions & 0 deletions cpp/src/common/tsfile_common.h
Original file line number Diff line number Diff line change
Expand Up @@ -78,6 +78,12 @@ struct PageHeader {
int ret = common::E_OK;
if (RET_FAIL(common::SerializationUtil::read_var_uint(
uncompressed_size_, in))) {
} else if (uncompressed_size_ == 0) {
// Empty page: the Java writer stores only a single zero varint and
// writes no compressed size or statistics (see
// ValueChunkWriter.writeEmptyPageToPageBuffer). Reading further
// would consume the following page's bytes as this page's header.
compressed_size_ = 0;
} else if (RET_FAIL(common::SerializationUtil::read_var_uint(
compressed_size_, in))) {
} else if (deserialize_stat) {
Expand Down
41 changes: 41 additions & 0 deletions cpp/test/common/tsfile_common_test.cc
Original file line number Diff line number Diff line change
Expand Up @@ -43,6 +43,47 @@ TEST(PageHeaderTest, Reset) {
EXPECT_EQ(header.compressed_size_, 0);
}

TEST(PageHeaderTest, DeserializeEmptyPageConsumesSingleVarint) {
// Java writes an empty value page as a single zero varint (see
// ValueChunkWriter.writeEmptyPageToPageBuffer): it carries no compressed
// size and no statistics. The C++ deserializer must not read past it.
common::ByteStream in(1024, common::MOD_DEFAULT);
ASSERT_EQ(common::E_OK, common::SerializationUtil::write_var_uint(0, in));
// A following non-empty page header (uncompressed=10, compressed=5).
ASSERT_EQ(common::E_OK, common::SerializationUtil::write_var_uint(10, in));
ASSERT_EQ(common::E_OK, common::SerializationUtil::write_var_uint(5, in));

PageHeader header;
ASSERT_EQ(common::E_OK,
header.deserialize_from(in, false, common::TSDataType::INT32));
EXPECT_EQ(header.uncompressed_size_, 0U);
EXPECT_EQ(header.compressed_size_, 0U);
EXPECT_EQ(header.statistic_, nullptr);
// The empty page header consumes exactly one byte (a single varint),
// leaving the following page header intact.
EXPECT_EQ(1U, in.read_pos());

PageHeader next;
ASSERT_EQ(common::E_OK,
next.deserialize_from(in, false, common::TSDataType::INT32));
EXPECT_EQ(next.uncompressed_size_, 10U);
EXPECT_EQ(next.compressed_size_, 5U);
EXPECT_EQ(next.statistic_, nullptr);
}

TEST(PageHeaderTest, DeserializeEmptyPageWithStatisticConsumesSingleVarint) {
common::ByteStream in(1024, common::MOD_DEFAULT);
ASSERT_EQ(common::E_OK, common::SerializationUtil::write_var_uint(0, in));

PageHeader header;
ASSERT_EQ(common::E_OK,
header.deserialize_from(in, true, common::TSDataType::INT32));
EXPECT_EQ(header.uncompressed_size_, 0U);
EXPECT_EQ(header.compressed_size_, 0U);
EXPECT_EQ(header.statistic_, nullptr);
EXPECT_EQ(1U, in.read_pos());
}

TEST(ChunkHeaderTest, DefaultConstructor) {
ChunkHeader header;
EXPECT_EQ(header.measurement_name_, "");
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -28,6 +28,7 @@
#include <cstring>
#include <fstream>
#include <limits>
#include <memory>
#include <sstream>
#include <string>
#include <unordered_set>
Expand All @@ -43,6 +44,7 @@
#endif
#include <fcntl.h>

#include "common/config/config.h"
#include "common/db_common.h"
#include "common/schema.h"
#include "common/tablet.h"
Expand All @@ -62,11 +64,14 @@ const char* const kGenerateDirEnv = "TSFILE_COMPAT_GENERATE_DIR";
const char* const kValidateDirEnv = "TSFILE_COMPAT_VALIDATE_DIR";
const char* const kManifestFile = "manifest.csv";
const char* const kManifestHeader =
"file,table,tagColumn,valueColumn,dataType,encoding,compression,rowCount";
"file,table,tagColumn,valueColumn,dataType,encoding,compression,rowCount,"
"nullPattern";
const char* const kTableName = "compat_table";
const char* const kTagColumn = "device";
const char* const kValueColumn = "value";
const char* const kTagValue = "compat_device";
const char* const kAnchorColumn = "anchor";
const int kNullPageRows = 64;
// Crosses the 129-value TS_2DIFF block boundary (300 -> 129+129+42), so
// page-wide bitmaps must survive block transitions. Applied to every case.
constexpr int kRowCount = 300;
Expand Down Expand Up @@ -181,6 +186,7 @@ struct FixtureCase {
TSEncoding encoding;
CompressionType compression;
int row_count;
std::string null_pattern = "NONE";

FixtureCase()
: table_name(kTableName),
Expand Down Expand Up @@ -221,7 +227,8 @@ struct FixtureCase {
out << file_name << "," << table_name << "," << tag_column << ","
<< value_column << "," << get_data_type_name(data_type) << ","
<< get_encoding_name(encoding) << ","
<< get_compression_name(compression) << "," << row_count;
<< get_compression_name(compression) << "," << row_count << ","
<< null_pattern;
return out.str();
}
};
Expand Down Expand Up @@ -264,7 +271,7 @@ bool ParseCompressionName(const std::string& value, CompressionType& out) {
FixtureCase ParseManifestLine(const std::string& line) {
FixtureCase fixture_case;
std::vector<std::string> parts = Split(line, ',');
if (parts.size() != 8) {
if (parts.size() != 9) {
ADD_FAILURE() << "Bad manifest line: " << line;
return fixture_case;
}
Expand All @@ -282,6 +289,10 @@ FixtureCase ParseManifestLine(const std::string& line) {
ADD_FAILURE() << "Bad compression in manifest line: " << line;
}
fixture_case.row_count = std::atoi(parts[7].c_str());
fixture_case.null_pattern = parts[8];
EXPECT_TRUE(fixture_case.null_pattern == "NONE" ||
fixture_case.null_pattern == "MIXED" ||
fixture_case.null_pattern == "ALL");
return fixture_case;
}

Expand Down Expand Up @@ -487,9 +498,55 @@ std::vector<FixtureCase> BuildMatrix() {
}
cases.emplace_back(DOUBLE, CAMEL, compression, kRowCount);
}
for (CompressionType compression : {UNCOMPRESSED, LZ4, ZSTD, LZMA2}) {
for (TSDataType data_type : data_types) {
for (const char* pattern : {"MIXED", "ALL"}) {
FixtureCase fixture(data_type, PLAIN, compression,
kNullPageRows * 8);
fixture.null_pattern = pattern;
fixture.file_name += "." + FixtureCase::Lower(pattern) + ".512";
cases.push_back(fixture);
}
FixtureCase fixture(data_type, PLAIN, compression, kNullPageRows);
fixture.null_pattern = "ALL";
fixture.file_name += ".all.64";
cases.push_back(fixture);
}
}
return cases;
}

// Shared contract with the Java test: empty first/middle/consecutive/trailing
// pages, non-empty pages 1/5, and alternating null values inside page 4.
bool IsNull(const FixtureCase& fixture_case, int row) {
if (fixture_case.null_pattern == "ALL") return true;
if (fixture_case.null_pattern == "NONE") return false;
int page = row / kNullPageRows;
return page == 0 || page == 2 || page == 3 || page >= 6 ||
(page == 4 && row % 2 == 0);
}

// Restore process-wide settings even when a fatal assertion returns early.
class ScopedNullPageSize {
public:
explicit ScopedNullPageSize(const FixtureCase& fixture_case)
: old_rows_(g_config_value_.page_writer_max_point_num_),
old_bytes_(g_config_value_.page_writer_max_memory_bytes_) {
if (fixture_case.null_pattern != "NONE") {
g_config_value_.page_writer_max_point_num_ = kNullPageRows;
g_config_value_.page_writer_max_memory_bytes_ = 1024 * 1024;
}
}
~ScopedNullPageSize() {
g_config_value_.page_writer_max_point_num_ = old_rows_;
g_config_value_.page_writer_max_memory_bytes_ = old_bytes_;
}

private:
uint32_t old_rows_;
uint32_t old_bytes_;
};

TableSchema* CreateTableSchema(const FixtureCase& fixture_case) {
std::vector<MeasurementSchema*> measurement_schemas;
std::vector<ColumnCategory> column_categories;
Expand All @@ -500,6 +557,11 @@ TableSchema* CreateTableSchema(const FixtureCase& fixture_case) {
new MeasurementSchema(fixture_case.value_column, fixture_case.data_type,
fixture_case.encoding, fixture_case.compression));
column_categories.emplace_back(ColumnCategory::FIELD);
if (fixture_case.null_pattern != "NONE") {
measurement_schemas.emplace_back(new MeasurementSchema(
kAnchorColumn, INT64, PLAIN, fixture_case.compression));
column_categories.emplace_back(ColumnCategory::FIELD);
}
return new TableSchema(fixture_case.table_name, measurement_schemas,
column_categories);
}
Expand Down Expand Up @@ -543,13 +605,21 @@ Tablet CreateTablet(TableSchema* table_schema,
for (int row = 0; row < fixture_case.row_count; ++row) {
EXPECT_EQ(E_OK, tablet.add_timestamp(row, row));
EXPECT_EQ(E_OK, tablet.add_value(row, kTagColumn, kTagValue));
AddValue(tablet, fixture_case, row);
// Tablet entries remain null until add_value is called.
if (!IsNull(fixture_case, row)) {
AddValue(tablet, fixture_case, row);
}
if (fixture_case.null_pattern != "NONE") {
EXPECT_EQ(E_OK, tablet.add_value(row, kAnchorColumn,
static_cast<int64_t>(row)));
}
}
return tablet;
}

void WriteFixture(const std::string& directory,
const FixtureCase& fixture_case) {
ScopedNullPageSize page_size(fixture_case);
std::string file_name = JoinPath(directory, fixture_case.file_name);
remove(file_name.c_str());
WriteFile write_file;
Expand Down Expand Up @@ -678,21 +748,30 @@ void AssertOnWireCodec(const std::string& directory,
}
}

void ValidateFixture(const std::string& directory,
const FixtureCase& fixture_case) {
SCOPED_TRACE(fixture_case.ToManifestLine());
AssertOnWireCodec(directory, fixture_case);
void ValidateQuery(const std::string& directory,
const FixtureCase& fixture_case, bool with_anchor, int start,
int end) {
SCOPED_TRACE(::testing::Message()
<< fixture_case.ToManifestLine() << " anchor=" << with_anchor
<< " range=" << start << ":" << end);
TsFileReader reader;
ASSERT_EQ(E_OK, reader.open(JoinPath(directory, fixture_case.file_name)));
std::vector<std::string> columns = {fixture_case.tag_column,
fixture_case.value_column};
if (with_anchor) columns.push_back(kAnchorColumn);
ResultSet* raw_result_set = nullptr;
ASSERT_EQ(E_OK,
reader.query(fixture_case.table_name,
{fixture_case.tag_column, fixture_case.value_column},
INT64_MIN, INT64_MAX, raw_result_set));
ASSERT_EQ(E_OK, reader.query(fixture_case.table_name, columns, start, end,
raw_result_set));
ASSERT_NE(nullptr, raw_result_set);
auto* result_set = static_cast<TableResultSet*>(raw_result_set);
// Assertion failures must release the query before its reader is destroyed.
auto cleanup = [&reader](TableResultSet* result) {
reader.destroy_query_data_set(result);
};
std::unique_ptr<TableResultSet, decltype(cleanup)> owned_result(result_set,
cleanup);
std::shared_ptr<ResultSetMetadata> metadata = result_set->get_metadata();
ASSERT_EQ(3U, metadata->get_column_count());
ASSERT_EQ(with_anchor ? 4U : 3U, metadata->get_column_count());
ASSERT_EQ("time", metadata->get_column_name(1));
ASSERT_EQ(INT64, metadata->get_column_type(1));
ASSERT_EQ(fixture_case.tag_column, metadata->get_column_name(2));
Expand All @@ -701,27 +780,51 @@ void ValidateFixture(const std::string& directory,
ASSERT_EQ(fixture_case.data_type, metadata->get_column_type(3));

bool has_next = false;
int row = 0;
int row = start;
while (true) {
ASSERT_EQ(E_OK, result_set->next(has_next));
if (!has_next) {
break;
}
if (!has_next) break;
SCOPED_TRACE(::testing::Message() << "row=" << row);
ASSERT_LE(row, end);
ASSERT_EQ(row, result_set->get_value<int64_t>(1));
ASSERT_FALSE(result_set->is_null(2));
ASSERT_EQ(kTagValue,
result_set->get_value<common::String*>(2)->to_std_string());
ASSERT_FALSE(result_set->is_null(3));
AssertValue(fixture_case, row, result_set);
ASSERT_EQ(IsNull(fixture_case, row), result_set->is_null(3));
if (!IsNull(fixture_case, row)) {
ASSERT_NO_FATAL_FAILURE(AssertValue(fixture_case, row, result_set));
}
if (with_anchor) {
ASSERT_FALSE(result_set->is_null(4));
ASSERT_EQ(row, result_set->get_value<int64_t>(4));
}
++row;
}
ASSERT_EQ(fixture_case.row_count, row);
result_set->close();
reader.destroy_query_data_set(result_set);
ASSERT_EQ(end + 1, row);
owned_result.reset();
ASSERT_EQ(E_OK, reader.close());
}

void ValidateFixture(const std::string& directory,
const FixtureCase& fixture_case) {
SCOPED_TRACE(fixture_case.ToManifestLine());
ASSERT_NO_FATAL_FAILURE(AssertOnWireCodec(directory, fixture_case));
ASSERT_NO_FATAL_FAILURE(ValidateQuery(directory, fixture_case, false, 0,
fixture_case.row_count - 1));
if (fixture_case.null_pattern != "NONE") {
ASSERT_NO_FATAL_FAILURE(ValidateQuery(directory, fixture_case, true, 0,
fixture_case.row_count - 1));
if (fixture_case.row_count > kNullPageRows) {
ASSERT_NO_FATAL_FAILURE(
ValidateQuery(directory, fixture_case, false, 63, 321));
ASSERT_NO_FATAL_FAILURE(
ValidateQuery(directory, fixture_case, true, 128, 255));
ASSERT_NO_FATAL_FAILURE(
ValidateQuery(directory, fixture_case, true, 320, 383));
}
}
}

} // namespace

class TableModelEncodingCompressionCompatibilityTest : public ::testing::Test {
Expand Down Expand Up @@ -751,6 +854,6 @@ TEST_F(TableModelEncodingCompressionCompatibilityTest, ValidateFixtures) {
}
std::vector<FixtureCase> cases = ReadManifest(input_dir);
for (const FixtureCase& fixture_case : cases) {
ValidateFixture(input_dir, fixture_case);
ASSERT_NO_FATAL_FAILURE(ValidateFixture(input_dir, fixture_case));
}
}
Loading
Loading