diff --git a/components/core/CMakeLists.txt b/components/core/CMakeLists.txt index a0889845b2..5d7cea3af0 100644 --- a/components/core/CMakeLists.txt +++ b/components/core/CMakeLists.txt @@ -679,6 +679,7 @@ set(SOURCE_FILES_unitTest tests/TestOutputCleaner.hpp tests/test-BoundedReader.cpp tests/test-BufferedFileReader.cpp + tests/test-clp_s-delta-encode-log-order.cpp tests/test-clp_s-end_to_end.cpp tests/test-clp_s-range_index.cpp tests/test-clp_s-search.cpp diff --git a/components/core/src/clp_s/ArchiveReader.cpp b/components/core/src/clp_s/ArchiveReader.cpp index f12bed273a..abf8605c1c 100644 --- a/components/core/src/clp_s/ArchiveReader.cpp +++ b/components/core/src/clp_s/ArchiveReader.cpp @@ -188,6 +188,9 @@ BaseColumnReader* ArchiveReader::append_reader_column(SchemaReader& reader, int3 case NodeType::Integer: column_reader = new Int64ColumnReader(column_id); break; + case NodeType::DeltaInteger: + column_reader = new DeltaEncodedInt64ColumnReader(column_id); + break; case NodeType::Float: column_reader = new FloatColumnReader(column_id); break; @@ -238,6 +241,9 @@ void ArchiveReader::append_unordered_reader_columns( case NodeType::Integer: column_reader = new Int64ColumnReader(column_id); break; + case NodeType::DeltaInteger: + column_reader = new DeltaEncodedInt64ColumnReader(column_id); + break; case NodeType::Float: column_reader = new FloatColumnReader(column_id); break; @@ -324,10 +330,8 @@ void ArchiveReader::initialize_schema_reader( } BaseColumnReader* column_reader = append_reader_column(reader, column_id); - if (column_id == m_log_event_idx_column_id - && nullptr != dynamic_cast(column_reader)) - { - reader.mark_column_as_log_event_idx(static_cast(column_reader)); + if (column_id == m_log_event_idx_column_id) { + reader.mark_column_as_log_event_idx(column_reader); } if (should_extract_timestamp && column_reader && timestamp_column_ids.count(column_id) > 0) diff --git a/components/core/src/clp_s/ArchiveWriter.cpp b/components/core/src/clp_s/ArchiveWriter.cpp index d152fd4ae7..a88a830ec7 100644 --- a/components/core/src/clp_s/ArchiveWriter.cpp +++ b/components/core/src/clp_s/ArchiveWriter.cpp @@ -331,6 +331,9 @@ void ArchiveWriter::initialize_schema_writer(SchemaWriter* writer, Schema const& case NodeType::DateString: writer->append_column(new DateStringColumnWriter(id)); break; + case NodeType::DeltaInteger: + writer->append_column(new DeltaEncodedInt64ColumnWriter(id)); + break; case NodeType::Metadata: case NodeType::NullValue: case NodeType::Object: diff --git a/components/core/src/clp_s/ColumnReader.cpp b/components/core/src/clp_s/ColumnReader.cpp index be914f4b47..493208de3e 100644 --- a/components/core/src/clp_s/ColumnReader.cpp +++ b/components/core/src/clp_s/ColumnReader.cpp @@ -16,6 +16,36 @@ std::variant Int64ColumnReader::extract_v return m_values[cur_message]; } +void DeltaEncodedInt64ColumnReader::load(BufferViewReader& reader, uint64_t num_messages) { + m_values = reader.read_unaligned_span(num_messages); + if (num_messages > 0) { + m_cur_idx = 0; + m_cur_value = m_values[0]; + } +} + +int64_t DeltaEncodedInt64ColumnReader::get_value_at_idx(size_t idx) { + if (m_cur_idx == idx) { + return m_cur_value; + } + if (idx > m_cur_idx) { + for (; m_cur_idx < idx; ++m_cur_idx) { + m_cur_value += m_values[m_cur_idx + 1]; + } + return m_cur_value; + } + for (; m_cur_idx > idx; --m_cur_idx) { + m_cur_value -= m_values[m_cur_idx]; + } + return m_cur_value; +} + +std::variant DeltaEncodedInt64ColumnReader::extract_value( + uint64_t cur_message +) { + return get_value_at_idx(cur_message); +} + void FloatColumnReader::load(BufferViewReader& reader, uint64_t num_messages) { m_values = reader.read_unaligned_span(num_messages); } @@ -25,6 +55,13 @@ Int64ColumnReader::extract_string_value_into_buffer(uint64_t cur_message, std::s buffer.append(std::to_string(m_values[cur_message])); } +void DeltaEncodedInt64ColumnReader::extract_string_value_into_buffer( + uint64_t cur_message, + std::string& buffer +) { + buffer.append(std::to_string(get_value_at_idx(cur_message))); +} + std::variant FloatColumnReader::extract_value( uint64_t cur_message ) { diff --git a/components/core/src/clp_s/ColumnReader.hpp b/components/core/src/clp_s/ColumnReader.hpp index 5800820042..d4b23cc54b 100644 --- a/components/core/src/clp_s/ColumnReader.hpp +++ b/components/core/src/clp_s/ColumnReader.hpp @@ -91,6 +91,39 @@ class Int64ColumnReader : public BaseColumnReader { UnalignedMemSpan m_values; }; +class DeltaEncodedInt64ColumnReader : public BaseColumnReader { +public: + // Constructor + explicit DeltaEncodedInt64ColumnReader(int32_t id) : BaseColumnReader(id) {} + + // Destructor + ~DeltaEncodedInt64ColumnReader() override = default; + + // Methods inherited from BaseColumnReader + void load(BufferViewReader& reader, uint64_t num_messages) override; + + NodeType get_type() override { return NodeType::DeltaInteger; } + + std::variant extract_value( + uint64_t cur_message + ) override; + + void extract_string_value_into_buffer(uint64_t cur_message, std::string& buffer) override; + +private: + /** + * Gets the value stored at a given index by summing up the stored deltas between the requested + * index and the last requested index. + * @param idx + * @return The value stored at the requested index. + */ + int64_t get_value_at_idx(size_t idx); + + UnalignedMemSpan m_values; + int64_t m_cur_value{}; + size_t m_cur_idx{}; +}; + class FloatColumnReader : public BaseColumnReader { public: // Constructor diff --git a/components/core/src/clp_s/ColumnWriter.cpp b/components/core/src/clp_s/ColumnWriter.cpp index 77fa51f155..384db0f5ec 100644 --- a/components/core/src/clp_s/ColumnWriter.cpp +++ b/components/core/src/clp_s/ColumnWriter.cpp @@ -11,6 +11,23 @@ void Int64ColumnWriter::store(ZstdCompressor& compressor) { compressor.write(reinterpret_cast(m_values.data()), size); } +size_t DeltaEncodedInt64ColumnWriter::add_value(ParsedMessage::variable_t& value) { + if (0 == m_values.size()) { + m_cur = std::get(value); + m_values.push_back(m_cur); + } else { + auto next = std::get(value); + m_values.push_back(next - m_cur); + m_cur = next; + } + return sizeof(int64_t); +} + +void DeltaEncodedInt64ColumnWriter::store(ZstdCompressor& compressor) { + size_t size = m_values.size() * sizeof(int64_t); + compressor.write(reinterpret_cast(m_values.data()), size); +} + size_t FloatColumnWriter::add_value(ParsedMessage::variable_t& value) { m_values.push_back(std::get(value)); return sizeof(double); diff --git a/components/core/src/clp_s/ColumnWriter.hpp b/components/core/src/clp_s/ColumnWriter.hpp index 64b79f960e..d126365f55 100644 --- a/components/core/src/clp_s/ColumnWriter.hpp +++ b/components/core/src/clp_s/ColumnWriter.hpp @@ -63,6 +63,24 @@ class Int64ColumnWriter : public BaseColumnWriter { std::vector m_values; }; +class DeltaEncodedInt64ColumnWriter : public BaseColumnWriter { +public: + // Constructor + explicit DeltaEncodedInt64ColumnWriter(int32_t id) : BaseColumnWriter(id) {} + + // Destructor + ~DeltaEncodedInt64ColumnWriter() override = default; + + // Methods inherited from BaseColumnWriter + size_t add_value(ParsedMessage::variable_t& value) override; + + void store(ZstdCompressor& compressor) override; + +private: + std::vector m_values; + int64_t m_cur{}; +}; + class FloatColumnWriter : public BaseColumnWriter { public: // Constructor diff --git a/components/core/src/clp_s/JsonParser.cpp b/components/core/src/clp_s/JsonParser.cpp index 142a1c3118..c30080ff63 100644 --- a/components/core/src/clp_s/JsonParser.cpp +++ b/components/core/src/clp_s/JsonParser.cpp @@ -513,7 +513,7 @@ bool JsonParser::parse() { auto initialize_fields_for_archive = [&]() -> bool { if (m_record_log_order) { log_event_idx_node_id - = add_metadata_field(constants::cLogEventIdxName, NodeType::Integer); + = add_metadata_field(constants::cLogEventIdxName, NodeType::DeltaInteger); } if (auto const rc = m_archive_writer->add_field_to_current_range( std::string{constants::range_index::cFilename}, @@ -982,7 +982,7 @@ auto JsonParser::parse_from_ir() -> bool { auto initialize_fields_for_archive = [&]() -> bool { if (m_record_log_order) { log_event_idx_node_id - = add_metadata_field(constants::cLogEventIdxName, NodeType::Integer); + = add_metadata_field(constants::cLogEventIdxName, NodeType::DeltaInteger); } if (auto const rc = m_archive_writer->add_field_to_current_range( std::string{constants::range_index::cFilename}, diff --git a/components/core/src/clp_s/SchemaReader.cpp b/components/core/src/clp_s/SchemaReader.cpp index 03bb9925c9..8fe842772f 100644 --- a/components/core/src/clp_s/SchemaReader.cpp +++ b/components/core/src/clp_s/SchemaReader.cpp @@ -29,6 +29,11 @@ void SchemaReader::mark_column_as_timestamp(BaseColumnReader* column_reader) { return std::get(static_cast(m_timestamp_column) ->extract_value(m_cur_message)); }; + } else if (m_timestamp_column->get_type() == NodeType::DeltaInteger) { + m_get_timestamp = [this]() { + return std::get(static_cast(m_timestamp_column) + ->extract_value(m_cur_message)); + }; } else if (m_timestamp_column->get_type() == NodeType::Float) { m_get_timestamp = [this]() { return static_cast( @@ -428,6 +433,7 @@ size_t SchemaReader::generate_structured_array_template( m_json_serializer.add_op(JsonSerializer::Op::EndArray); break; } + case NodeType::DeltaInteger: case NodeType::Integer: { m_json_serializer.add_op(JsonSerializer::Op::AddIntValue); m_reordered_columns.push_back(m_columns[column_idx++]); @@ -512,6 +518,7 @@ size_t SchemaReader::generate_structured_object_template( m_json_serializer.add_op(JsonSerializer::Op::EndArray); break; } + case NodeType::DeltaInteger: case NodeType::Integer: { m_json_serializer.add_op(JsonSerializer::Op::AddIntField); m_reordered_columns.push_back(m_columns[column_idx++]); @@ -620,6 +627,7 @@ void SchemaReader::generate_json_template(int32_t id) { m_json_serializer.add_op(JsonSerializer::Op::EndArray); break; } + case NodeType::DeltaInteger: case NodeType::Integer: { m_json_serializer.add_op(JsonSerializer::Op::AddIntField); m_reordered_columns.push_back(m_column_map[child_global_id]); diff --git a/components/core/src/clp_s/SchemaReader.hpp b/components/core/src/clp_s/SchemaReader.hpp index 6264707cae..64c397b593 100644 --- a/components/core/src/clp_s/SchemaReader.hpp +++ b/components/core/src/clp_s/SchemaReader.hpp @@ -206,7 +206,7 @@ class SchemaReader { /** * Marks a column as the log_event_idx column. */ - void mark_column_as_log_event_idx(Int64ColumnReader* column_reader) { + void mark_column_as_log_event_idx(BaseColumnReader* column_reader) { m_log_event_idx_column = column_reader; } @@ -321,7 +321,7 @@ class SchemaReader { BaseColumnReader* m_timestamp_column; std::function m_get_timestamp; - Int64ColumnReader* m_log_event_idx_column{nullptr}; + BaseColumnReader* m_log_event_idx_column{nullptr}; std::shared_ptr m_global_schema_tree; SchemaTree m_local_schema_tree; diff --git a/components/core/src/clp_s/SchemaTree.cpp b/components/core/src/clp_s/SchemaTree.cpp index 6e3019edb7..270a048745 100644 --- a/components/core/src/clp_s/SchemaTree.cpp +++ b/components/core/src/clp_s/SchemaTree.cpp @@ -11,6 +11,7 @@ auto node_to_literal_type(NodeType type) -> clp_s::search::ast::LiteralType { // type-per-token support. switch (type) { case NodeType::Integer: + case NodeType::DeltaInteger: return clp_s::search::ast::LiteralType::IntegerT; case NodeType::Float: return clp_s::search::ast::LiteralType::FloatT; diff --git a/components/core/src/clp_s/SchemaTree.hpp b/components/core/src/clp_s/SchemaTree.hpp index 55e2cb4ff8..f1289235be 100644 --- a/components/core/src/clp_s/SchemaTree.hpp +++ b/components/core/src/clp_s/SchemaTree.hpp @@ -41,6 +41,7 @@ enum class NodeType : uint8_t { DateString, StructuredArray, Metadata, + DeltaInteger, Unknown = std::underlying_type::type(~0ULL) }; diff --git a/components/core/src/clp_s/SingleFileArchiveDefs.hpp b/components/core/src/clp_s/SingleFileArchiveDefs.hpp index 9dea07188a..18bd3123d7 100644 --- a/components/core/src/clp_s/SingleFileArchiveDefs.hpp +++ b/components/core/src/clp_s/SingleFileArchiveDefs.hpp @@ -11,7 +11,7 @@ namespace clp_s { // define the version constexpr uint8_t cArchiveMajorVersion = 0; constexpr uint8_t cArchiveMinorVersion = 3; -constexpr uint16_t cArchivePatchVersion = 1; +constexpr uint16_t cArchivePatchVersion = 2; // define the magic number constexpr uint8_t cStructuredSFAMagicNumber[] = {0xFD, 0x2F, 0xC5, 0x30}; diff --git a/components/core/tests/test-clp_s-delta-encode-log-order.cpp b/components/core/tests/test-clp_s-delta-encode-log-order.cpp new file mode 100644 index 0000000000..89443f2795 --- /dev/null +++ b/components/core/tests/test-clp_s-delta-encode-log-order.cpp @@ -0,0 +1,118 @@ +#include +#include +#include +#include +#include +#include +#include + +#include + +#include "../src/clp_s/archive_constants.hpp" +#include "../src/clp_s/ArchiveReader.hpp" +#include "../src/clp_s/ColumnReader.hpp" +#include "../src/clp_s/InputConfig.hpp" +#include "../src/clp_s/SchemaReader.hpp" +#include "clp_s_test_utils.hpp" +#include "TestOutputCleaner.hpp" + +constexpr std::string_view cTestDeltaEncodeOrderArchiveDirectory{"test-delta-encode-order-archive"}; +constexpr std::string_view cTestDeltaEncodeOrderInputFileDirectory{"test_log_files"}; +constexpr std::string_view cTestDeltaEncodeOrderInputFile{"test_simple_order.jsonl"}; +constexpr size_t cNumEntries{3}; + +namespace { +/** + * A simple implementation of `clp_s::FilterClass` that allows us to grab the underlying + * `std::vector` from a `SchemaReader`. + */ +class SimpleFilterClass : public clp_s::FilterClass { +public: + void init( + clp_s::SchemaReader* reader, + std::vector const& column_readers + ) override { + m_column_readers = column_readers; + } + + auto filter(uint64_t cur_message) -> bool override { return true; } + + auto get_column_readers() -> std::vector const& { + return m_column_readers; + } + +private: + std::vector m_column_readers; +}; + +auto get_test_input_path_relative_to_tests_dir() -> std::filesystem::path; +auto get_test_input_local_path() -> std::string; + +auto get_test_input_path_relative_to_tests_dir() -> std::filesystem::path { + return std::filesystem::path{cTestDeltaEncodeOrderInputFileDirectory} + / cTestDeltaEncodeOrderInputFile; +} + +auto get_test_input_local_path() -> std::string { + std::filesystem::path const current_file_path{__FILE__}; + auto const tests_dir{current_file_path.parent_path()}; + return (tests_dir / get_test_input_path_relative_to_tests_dir()).string(); +} +} // namespace + +TEST_CASE("clp-s-delta-encode-log-order", "[clp-s][delta-encode-log-order]") { + auto start_index = GENERATE(0ULL, 1ULL, 2ULL); + TestOutputCleaner const test_cleanup{{std::string{cTestDeltaEncodeOrderArchiveDirectory}}}; + + REQUIRE_NOTHROW(compress_archive( + get_test_input_local_path(), + std::string{cTestDeltaEncodeOrderArchiveDirectory}, + true, + false, + clp_s::FileType::Json + )); + + std::vector archive_paths; + REQUIRE(clp_s::get_input_archives_for_raw_path( + std::string{cTestDeltaEncodeOrderArchiveDirectory}, + archive_paths + )); + REQUIRE(1 == archive_paths.size()); + + clp_s::ArchiveReader archive_reader; + REQUIRE_NOTHROW(archive_reader.open(archive_paths.back(), clp_s::NetworkAuthOption{})); + REQUIRE_NOTHROW(archive_reader.read_dictionaries_and_metadata()); + REQUIRE_NOTHROW(archive_reader.open_packed_streams()); + auto mpt = archive_reader.get_schema_tree(); + auto log_event_idx_node_id = mpt->get_metadata_field_id(clp_s::constants::cLogEventIdxName); + REQUIRE(-1 != log_event_idx_node_id); + + std::vector> schema_readers; + REQUIRE_NOTHROW(schema_readers = archive_reader.read_all_tables()); + REQUIRE(1 == schema_readers.size()); + auto schema_reader = schema_readers.back(); + REQUIRE(cNumEntries == schema_reader->get_num_messages()); + + SimpleFilterClass simple_filter_class; + schema_reader->initialize_filter(&simple_filter_class); + clp_s::BaseColumnReader* log_event_idx_reader{nullptr}; + for (auto* column_reader : simple_filter_class.get_column_readers()) { + if (log_event_idx_node_id == column_reader->get_id()) { + log_event_idx_reader = column_reader; + break; + } + } + REQUIRE(nullptr != log_event_idx_reader); + REQUIRE(clp_s::NodeType::DeltaInteger == log_event_idx_reader->get_type()); + REQUIRE(nullptr != dynamic_cast(log_event_idx_reader)); + + // Test forwards and backwards seeks on `DeltaEncodedInt64ColumnReader`. + size_t i{start_index}; + for (size_t num_iterations{0ULL}; num_iterations < cNumEntries; ++num_iterations) { + int64_t val{}; + REQUIRE_NOTHROW(val = std::get(log_event_idx_reader->extract_value(i))); + REQUIRE(val == static_cast(i)); + i = (i + 1) % cNumEntries; + } + REQUIRE_NOTHROW(archive_reader.close()); +} diff --git a/components/core/tests/test_log_files/test_simple_order.jsonl b/components/core/tests/test_log_files/test_simple_order.jsonl new file mode 100644 index 0000000000..e9eab76518 --- /dev/null +++ b/components/core/tests/test_log_files/test_simple_order.jsonl @@ -0,0 +1,3 @@ +{"idx": 0} +{"idx": 1} +{"idx": 2}