diff --git a/CHANGELOG.md b/CHANGELOG.md index 32e7c1131e6..51d479c14d2 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -17,6 +17,10 @@ and this project adheres to [Semantic Versioning](http://semver.org/spec/v2.0.0. [7.0.10]: https://github.com/microsoft/CCF/releases/tag/ccf-7.0.10 +### Added + +- Recovery can now use a COSE snapshot signed by an earlier service identity after one or more disaster recoveries. Before deserialising the snapshot, the node reads previous-service-identity endorsement candidates from the public ledger suffix, validates a complete chain against the operator-provided identity, and retains it only for the current recovery attempt. Invalid or incomplete endorsement chains fall back to full-ledger replay (#8092). + ### Changed - `ccf::http::ParsedQuery` (in `include/ccf/http_query.h`), returned by `ccf::http::parse_query()`, is now a `std::multimap>` that owns its decoded keys and values, rather than a `std::multimap` pointing into the source query string. Owned storage is required because each key and value is now URL-decoded individually after splitting, which produces bytes not present in the original query. Application code that consumed the previous `std::string_view` keys/values may need to be updated (#8024). diff --git a/CMakeLists.txt b/CMakeLists.txt index 5aa9b3450f5..4e3c19f0689 100644 --- a/CMakeLists.txt +++ b/CMakeLists.txt @@ -1109,6 +1109,12 @@ if(BUILD_TESTS) ADDITIONAL_ARGS --regex ^recovery_intermediate_snapshot_join$ ) + add_e2e_test( + NAME recovery_snapshot_endorsements_test + PYTHON_SCRIPT ${CMAKE_SOURCE_DIR}/tests/recovery_snapshot_endorsements.py + BUCKET bucket_b + ) + add_e2e_test( NAME recovery_test_suite PYTHON_SCRIPT ${CMAKE_SOURCE_DIR}/tests/e2e_suite.py diff --git a/src/kv/generic_serialise_wrapper.h b/src/kv/generic_serialise_wrapper.h index 13e5645fc16..b2a128f4d2a 100644 --- a/src/kv/generic_serialise_wrapper.h +++ b/src/kv/generic_serialise_wrapper.h @@ -347,7 +347,14 @@ namespace ccf::kv } serialized::skip(data_, size_, crypto_util->get_header_length()); - auto public_domain_length = serialized::read(data_, size_); + const auto public_domain_length = serialized::read(data_, size_); + if (public_domain_length > size_) + { + throw std::logic_error(fmt::format( + "Public domain length {} exceeds remaining entry size {}", + public_domain_length, + size_)); + } const auto* data_public = data_; public_reader.init(data_public, public_domain_length); diff --git a/src/kv/raw_serialise.h b/src/kv/raw_serialise.h index ef1b1104a2a..6ac43be675a 100644 --- a/src/kv/raw_serialise.h +++ b/src/kv/raw_serialise.h @@ -6,7 +6,9 @@ #include "generic_serialise_wrapper.h" #include +#include #include +#include #include namespace ccf::kv @@ -126,21 +128,69 @@ namespace ccf::kv class RawReader { - public: - const uint8_t* data_ptr; + private: + [[nodiscard]] size_t remaining_bytes() const + { + if (data_offset > data_size) + { + throw std::logic_error(fmt::format( + "Raw reader offset {} exceeds data size {}", data_offset, data_size)); + } + return data_size - data_offset; + } + + void require_bytes(size_t required, const char* description) const + { + const auto remaining = remaining_bytes(); + if (required > remaining) + { + throw std::runtime_error(fmt::format( + "Expected {} bytes for {}, found only {}", + required, + description, + remaining)); + } + } + + [[nodiscard]] const uint8_t* current_data() const + { + if (data_ptr == nullptr) + { + if (data_size != 0) + { + throw std::logic_error("Raw reader has non-zero size with null data"); + } + return nullptr; + } + return data_ptr + data_offset; + } + + void advance(size_t size) + { + require_bytes(size, "reader advance"); + data_offset += size; + } + + const uint8_t* data_ptr{nullptr}; size_t data_offset{0}; - size_t data_size; + size_t data_size{0}; + public: /** Reads the next entry, advancing data_offset */ template T read_entry() { - auto remainder = data_size - data_offset; - const auto* data = data_ptr + data_offset; + require_bytes(sizeof(T), "fixed-size entry"); + const auto before = remaining_bytes(); + auto remainder = before; + const auto* data = current_data(); const auto entry = serialized::read(data, remainder); - const auto bytes_read = data_size - data_offset - remainder; - data_offset += bytes_read; + if (remainder > before) + { + throw std::logic_error("Raw reader remaining size increased"); + } + advance(before - remainder); return entry; } @@ -148,17 +198,11 @@ namespace ccf::kv */ size_t read_size_prefixed_entry(size_t& start_offset) { - auto remainder = data_size - data_offset; - auto entry_size = read_entry(); - - if (remainder < entry_size) - { - throw std::runtime_error(fmt::format( - "Expected {} byte entry, found only {}", entry_size, remainder)); - } + const auto entry_size = read_entry(); + require_bytes(entry_size, "size-prefixed entry"); start_offset = data_offset; - data_offset += entry_size; + advance(entry_size); return entry_size; } @@ -166,13 +210,18 @@ namespace ccf::kv RawReader(const RawReader& other) = delete; RawReader& operator=(const RawReader& other) = delete; - RawReader(const uint8_t* data_in_ptr = nullptr, size_t data_in_size = 0) : - data_ptr(data_in_ptr), - data_size(data_in_size) - {} + RawReader(const uint8_t* data_in_ptr = nullptr, size_t data_in_size = 0) + { + init(data_in_ptr, data_in_size); + } void init(const uint8_t* data_in_ptr, size_t data_in_size) { + if (data_in_ptr == nullptr && data_in_size != 0) + { + throw std::invalid_argument( + "Cannot initialise raw reader with null data and non-zero size"); + } data_offset = 0; data_ptr = data_in_ptr; data_size = data_in_size; @@ -186,9 +235,30 @@ namespace ccf::kv std::is_same_v) { size_t entry_offset = 0; - size_t entry_size = read_size_prefixed_entry(entry_offset); + const auto entry_size = read_size_prefixed_entry(entry_offset); + using Element = typename T::value_type; + if (entry_size % sizeof(Element) != 0) + { + throw std::runtime_error(fmt::format( + "Size-prefixed entry of {} bytes is not divisible by element size " + "{}", + entry_size, + sizeof(Element))); + } - T ret(entry_size / sizeof(typename T::value_type)); + T ret; + const auto element_count = entry_size / sizeof(Element); + if (element_count > ret.max_size()) + { + throw std::length_error(fmt::format( + "Size-prefixed entry contains too many elements ({})", + element_count)); + } + ret.resize(element_count); + if (entry_size == 0) + { + return ret; + } auto* data_dest = reinterpret_cast(ret.data()); auto capacity = entry_size; // NOLINTNEXTLINE(readability-suspicious-call-argument) @@ -201,10 +271,15 @@ namespace ccf::kv { T ret{}; auto* data_ = reinterpret_cast(ret.data()); - constexpr size_t size = ret.size() * sizeof(typename T::value_type); + constexpr auto element_count = std::tuple_size_v; + static_assert( + element_count <= + std::numeric_limits::max() / sizeof(typename T::value_type)); + constexpr size_t size = element_count * sizeof(typename T::value_type); + require_bytes(size, "fixed-size array"); auto size_ = size; - serialized::write(data_, size_, data_ptr + data_offset, size); - data_offset += size; + serialized::write(data_, size_, current_data(), size); + advance(size); return ret; } @@ -240,7 +315,7 @@ namespace ccf::kv [[nodiscard]] bool is_eos() const { - return data_offset >= data_size; + return remaining_bytes() == 0; } }; diff --git a/src/kv/test/kv_serialisation.cpp b/src/kv/test/kv_serialisation.cpp index 92078bc2478..c41d1949482 100644 --- a/src/kv/test/kv_serialisation.cpp +++ b/src/kv/test/kv_serialisation.cpp @@ -2,12 +2,16 @@ // Licensed under the Apache 2.0 License. #include "ds/internal_logger.h" #include "kv/kv_serialiser.h" +#include "kv/raw_serialise.h" #include "kv/store.h" #include "kv/test/null_encryptor.h" #include "kv/test/stub_consensus.h" #include #undef FAIL +#include +#include +#include #include #include @@ -19,6 +23,138 @@ struct MapTypes using StringNum = ccf::kv::Map; }; +static std::vector make_size_prefixed_bytes( + size_t declared_size, size_t actual_size) +{ + std::vector bytes(sizeof(size_t) + actual_size); + std::memcpy(bytes.data(), &declared_size, sizeof(declared_size)); + return bytes; +} + +TEST_CASE( + "Raw reader rejects truncated entries" * doctest::test_suite("serialisation")) +{ + SUBCASE("Fixed-size integral") + { + std::vector bytes(sizeof(uint64_t) - 1); + ccf::kv::RawReader reader(bytes.data(), bytes.size()); + REQUIRE_THROWS(reader.read_next()); + } + + SUBCASE("Fixed-size array") + { + std::vector bytes(ccf::crypto::Sha256Hash::SIZE - 1); + ccf::kv::RawReader reader(bytes.data(), bytes.size()); + REQUIRE_THROWS(reader.read_next()); + } + + SUBCASE("Short size prefix") + { + std::vector bytes(sizeof(size_t) - 1); + ccf::kv::RawReader reader(bytes.data(), bytes.size()); + REQUIRE_THROWS(reader.read_next>()); + } + + SUBCASE("Prefix without payload") + { + auto bytes = make_size_prefixed_bytes(1, 0); + ccf::kv::RawReader reader(bytes.data(), bytes.size()); + REQUIRE_THROWS(reader.read_next>()); + } + + SUBCASE("Nine bytes remaining declare two payload bytes") + { + auto bytes = make_size_prefixed_bytes(2, 1); + ccf::kv::RawReader reader(bytes.data(), bytes.size()); + REQUIRE_THROWS(reader.read_next>()); + } + + SUBCASE("Impossible payload length") + { + auto bytes = + make_size_prefixed_bytes(std::numeric_limits::max(), 0); + ccf::kv::RawReader reader(bytes.data(), bytes.size()); + REQUIRE_THROWS(reader.read_next>()); + } + + SUBCASE("Payload is not a whole number of elements") + { + auto bytes = make_size_prefixed_bytes(1, 1); + ccf::kv::RawReader reader(bytes.data(), bytes.size()); + REQUIRE_THROWS(reader.read_next>()); + } + + SUBCASE("Null data with non-zero size") + { + REQUIRE_THROWS(ccf::kv::RawReader(nullptr, 1)); + } +} + +static std::vector make_public_domain_entry( + size_t declared_public_domain_size, const std::vector& public_domain) +{ + ccf::kv::SerialisedEntryHeader header; + header.set_size(sizeof(size_t) + public_domain.size()); + std::vector entry(sizeof(header) + header.size); + auto* data = entry.data(); + auto remaining = entry.size(); + serialized::write(data, remaining, header); + serialized::write(data, remaining, declared_public_domain_size); + serialized::write( + data, remaining, public_domain.data(), public_domain.size()); + return entry; +} + +TEST_CASE( + "KV deserialiser rejects invalid public domains" * + doctest::test_suite("serialisation")) +{ + auto encryptor = std::make_shared(); + + const auto initialise = [&](const std::vector& entry) { + ccf::kv::RawKvStoreDeserialiser deserialiser( + encryptor, ccf::kv::SecurityDomain::PUBLIC); + ccf::kv::Term term = 0; + ccf::kv::EntryFlags flags = {}; + return deserialiser.init(entry.data(), entry.size(), term, flags, false); + }; + + SUBCASE("Public domain exceeds remaining entry") + { + const auto entry = make_public_domain_entry(2, {0}); + REQUIRE_THROWS(initialise(entry)); + } + + SUBCASE("Public domain length prefix is truncated") + { + ccf::kv::SerialisedEntryHeader header; + header.set_size(sizeof(size_t) - 1); + std::vector entry(sizeof(header) + header.size); + std::memcpy(entry.data(), &header, sizeof(header)); + REQUIRE_THROWS(initialise(entry)); + } + + SUBCASE("Public domain length is impossible") + { + const auto entry = + make_public_domain_entry(std::numeric_limits::max(), {}); + REQUIRE_THROWS(initialise(entry)); + } + + SUBCASE("Claims digest is truncated") + { + std::vector public_domain( + sizeof(ccf::kv::EntryType) + sizeof(ccf::kv::Version)); + auto* data = public_domain.data(); + *data++ = static_cast(ccf::kv::EntryType::WriteSetWithClaims); + const ccf::kv::Version version = 1; + std::memcpy(data, &version, sizeof(version)); + const auto entry = + make_public_domain_entry(public_domain.size(), public_domain); + REQUIRE_THROWS(initialise(entry)); + } +} + TEST_CASE( "Serialise/deserialise public map only" * doctest::test_suite("serialisation")) diff --git a/src/node/node_state.h b/src/node/node_state.h index caf4357872a..e1bab416266 100644 --- a/src/node/node_state.h +++ b/src/node/node_state.h @@ -45,6 +45,7 @@ #include "node/node_inbound_message.h" #include "node/node_to_node_channel_manager.h" #include "node/recovery_decision_protocol.h" +#include "node/recovery_snapshot_ledger.h" #include "node/signature_cache_subsystem.h" #include "node/snapshotter.h" #include "node_to_node.h" @@ -527,8 +528,21 @@ namespace ccf snapshot_path, snapshot_data.size()); + // Structurally invalid snapshots remain fatal. const auto segments = separate_segments(snapshot_data); + if (start_type == StartType::Recover) + { + // Validate the receipt structure and claims now, but defer the + // previous-service-identity check to the public-ledger endorsement + // scan. The snapshot is not installed until that verification + // succeeds, so the Store is not mutated by an unverified snapshot. + verify_snapshot(segments); + startup_snapshot_info = std::make_unique( + snapshot_seqno, std::move(snapshot_data)); + return; + } + try { verify_snapshot(segments, config.recover.previous_service_identity); @@ -588,7 +602,18 @@ namespace ccf startup_snapshot_info = std::make_unique( snapshot_seqno, std::move(snapshot_data)); - LOG_INFO_FMT("Setting startup snapshot seqno to {}", snapshot_seqno); + install_startup_snapshot(); + } + + void install_startup_snapshot() + { + if (!startup_snapshot_info) + { + throw std::logic_error("No startup snapshot selected for installation"); + } + + LOG_INFO_FMT( + "Setting startup snapshot seqno to {}", startup_snapshot_info->seqno); startup_seqno = startup_snapshot_info->seqno; last_recovered_idx = startup_seqno; @@ -635,6 +660,113 @@ namespace ccf } } + void start_public_ledger_recovery_unsafe() + { + sm.advance(NodeStartupState::readingPublicLedger); + start_ledger_recovery_unsafe(); + } + + void fallback_from_recovery_snapshot_unsafe(const std::string& reason) + { + LOG_FAIL_FMT( + "Recovery snapshot cannot be verified under the configured previous " + "service identity: {}. Falling back to full-ledger recovery.", + reason); + + startup_snapshot_info.reset(); + startup_seqno = 0; + last_recovered_idx = 0; + last_recovered_signed_idx = 0; + view_history.clear(); + + start_public_ledger_recovery_unsafe(); + } + + void install_recovery_snapshot_and_start_unsafe() + { + install_startup_snapshot(); + start_public_ledger_recovery_unsafe(); + } + + void start_recovery_snapshot_verification_unsafe() + { + if ( + !startup_snapshot_info || + !config.recover.previous_service_identity.has_value()) + { + start_public_ledger_recovery_unsafe(); + return; + } + + const auto segments = separate_segments(startup_snapshot_info->raw); + const ccf::crypto::Pem target_identity( + *config.recover.previous_service_identity); + try + { + verify_snapshot_seqno( + segments, + network.tables->get_encryptor(), + startup_snapshot_info->seqno); + } + catch (const std::exception& e) + { + fallback_from_recovery_snapshot_unsafe(e.what()); + return; + } + + const auto direct_verification_error = + try_verify_and_install_recovery_snapshot( + [&]() { + verify_snapshot(segments, target_identity.raw()); + LOG_INFO_FMT( + "Recovery snapshot at {} is directly signed by the configured " + "previous service identity", + startup_snapshot_info->seqno); + }, + [&]() { install_recovery_snapshot_and_start_unsafe(); }); + if (!direct_verification_error.has_value()) + { + return; + } + + if (segments.receipt.empty() || segments.receipt[0] != 0xD2) + { + fallback_from_recovery_snapshot_unsafe(fmt::format( + "old-style snapshot receipt cannot use an endorsement chain: {}", + *direct_verification_error)); + return; + } + + LOG_INFO_FMT( + "Recovery snapshot at {} is not directly signed by the configured " + "previous service identity ({}); scanning the public ledger suffix " + "for COSE endorsements", + startup_snapshot_info->seqno, + *direct_verification_error); + + const auto chain_verification_error = + try_verify_and_install_recovery_snapshot( + [&]() { + const auto snapshot_seqno = startup_snapshot_info->seqno; + const auto scan = scan_recovery_snapshot_ledger_files( + config.ledger, network.tables->get_encryptor(), snapshot_seqno); + const auto target_key = ccf::crypto::public_key_der_from_cert( + ccf::crypto::cert_pem_to_der(target_identity)); + const auto snapshot_signer_key = + validate_recovery_snapshot_endorsement_chain( + scan.endorsements, target_key, snapshot_seqno); + verify_recovery_snapshot_receipt(segments, snapshot_signer_key); + LOG_INFO_FMT( + "Validated {} recovery snapshot endorsement(s) in memory", + scan.endorsements.size()); + }, + [&]() { install_recovery_snapshot_and_start_unsafe(); }); + if (chain_verification_error.has_value()) + { + fallback_from_recovery_snapshot_unsafe(*chain_verification_error); + } + } + RecoveryDecisionProtocolSubsystem recovery_decision_protocol; public: @@ -836,8 +968,14 @@ namespace ccf find_local_startup_snapshot(); - sm.advance(NodeStartupState::readingPublicLedger); - start_ledger_recovery_unsafe(); + if (startup_snapshot_info) + { + start_recovery_snapshot_verification_unsafe(); + } + else + { + start_public_ledger_recovery_unsafe(); + } return; } default: diff --git a/src/node/recovery_snapshot_ledger.h b/src/node/recovery_snapshot_ledger.h new file mode 100644 index 00000000000..70c14a8292f --- /dev/null +++ b/src/node/recovery_snapshot_ledger.h @@ -0,0 +1,453 @@ +// Copyright (c) Microsoft Corporation. All rights reserved. +// Licensed under the Apache 2.0 License. +#pragma once + +#include "ccf/node/startup_config.h" +#include "ds/internal_logger.h" +#include "host/ledger_filenames.h" +#include "kv/kv_serialiser.h" +#include "kv/serialised_entry_format.h" +#include "node/rpc/network_identity_chain_helpers.h" +#include "service/tables/previous_service_identity.h" + +#include +#include +#include +#include +#include + +namespace ccf +{ + static constexpr size_t MAX_RECOVERY_SNAPSHOT_ENDORSEMENTS_COUNT = 64; + static constexpr size_t MAX_RECOVERY_SNAPSHOT_ENDORSEMENT_SIZE = + size_t{1024} * 1024; + static constexpr size_t MAX_RECOVERY_SNAPSHOT_ENDORSEMENT_RECORD_SIZE = + 2 * MAX_RECOVERY_SNAPSHOT_ENDORSEMENT_SIZE; + static constexpr size_t MAX_RECOVERY_SNAPSHOT_ENDORSEMENTS_PAYLOAD_SIZE = + size_t{4} * 1024 * 1024; + static constexpr size_t MAX_RECOVERY_SNAPSHOT_LEDGER_ENTRY_SIZE = + size_t{16} * 1024 * 1024; + + struct RecoverySnapshotLedgerEntry + { + ccf::kv::Version version = 0; + std::optional endorsement = std::nullopt; + }; + + struct RecoverySnapshotLedgerScan + { + std::vector endorsements; + }; + + static RecoverySnapshotLedgerEntry parse_recovery_snapshot_ledger_entry( + const std::vector& entry, + const std::shared_ptr& encryptor) + { + auto deserialiser = ccf::kv::RawKvStoreDeserialiser( + encryptor, ccf::kv::SecurityDomain::PUBLIC); + ccf::kv::Term term = 0; + ccf::kv::EntryFlags flags = {}; + const auto version = + deserialiser.init(entry.data(), entry.size(), term, flags, false); + if (!version.has_value()) + { + throw std::logic_error( + "Failed to initialise public ledger entry deserialiser"); + } + + RecoverySnapshotLedgerEntry result; + result.version = *version; + + for (auto map_name = deserialiser.start_map(); map_name.has_value(); + map_name = deserialiser.start_map()) + { + std::ignore = deserialiser.deserialise_entry_version(); + + const auto read_count = deserialiser.deserialise_read_header(); + for (size_t i = 0; i < read_count; ++i) + { + std::ignore = deserialiser.deserialise_read(); + } + + const auto write_count = deserialiser.deserialise_write_header(); + for (size_t i = 0; i < write_count; ++i) + { + auto [key, value] = deserialiser.deserialise_write(); + if (*map_name != ccf::Tables::PREVIOUS_SERVICE_IDENTITY_ENDORSEMENT) + { + continue; + } + if ( + result.endorsement.has_value() || + key != ccf::PreviousServiceIdentityEndorsement::create_unit()) + { + throw std::logic_error( + "Invalid previous service identity endorsement table write"); + } + if (value.size() > MAX_RECOVERY_SNAPSHOT_ENDORSEMENT_RECORD_SIZE) + { + throw std::logic_error(fmt::format( + "Serialised previous service identity endorsement is too large " + "({} bytes; maximum {} bytes)", + value.size(), + MAX_RECOVERY_SNAPSHOT_ENDORSEMENT_RECORD_SIZE)); + } + result.endorsement = ccf::PreviousServiceIdentityEndorsement:: + ValueSerialiser::from_serialised(value); + } + + const auto remove_count = deserialiser.deserialise_remove_header(); + for (size_t i = 0; i < remove_count; ++i) + { + std::ignore = deserialiser.deserialise_remove(); + if (*map_name == ccf::Tables::PREVIOUS_SERVICE_IDENTITY_ENDORSEMENT) + { + throw std::logic_error( + "Unexpected removal from previous service identity endorsement " + "table"); + } + } + } + + if (!deserialiser.end()) + { + throw std::logic_error( + "Public ledger entry contains trailing serialised data"); + } + return result; + } + + struct RecoverySnapshotLedgerFile + { + std::filesystem::path path; + size_t start_idx; + std::optional end_idx; + bool committed; + }; + + static std::vector + find_recovery_snapshot_ledger_files(const ccf::CCFConfig::Ledger& config) + { + std::vector files; + + auto add_directory = + [&](const std::filesystem::path& directory, bool read_only) { + std::error_code ec; + const auto exists = std::filesystem::exists(directory, ec); + if (ec) + { + throw std::logic_error(fmt::format( + "Unable to inspect ledger directory {}: {}", + directory.string(), + ec.message())); + } + if (!exists) + { + return; + } + + for (std::filesystem::directory_iterator it(directory, ec), end; + it != end; + it.increment(ec)) + { + if (ec) + { + throw std::logic_error(fmt::format( + "Unable to iterate ledger directory {}: {}", + directory.string(), + ec.message())); + } + if (!it->is_regular_file()) + { + continue; + } + + const auto name = it->path().filename().string(); + if ( + !name.starts_with("ledger_") || + asynchost::is_ledger_file_ignored(name)) + { + continue; + } + + const auto committed = asynchost::is_ledger_file_name_committed(name); + if (read_only && !committed) + { + continue; + } + + try + { + files.push_back( + {it->path(), + asynchost::get_start_idx_from_file_name(name), + asynchost::get_last_idx_from_file_name(name), + committed}); + } + catch (const std::exception& e) + { + LOG_INFO_FMT( + "Ignoring invalid ledger file name {} while scanning recovery " + "snapshot endorsements: {}", + it->path().string(), + e.what()); + } + } + if (ec) + { + throw std::logic_error(fmt::format( + "Unable to iterate ledger directory {}: {}", + directory.string(), + ec.message())); + } + }; + + add_directory(config.directory, false); + for (const auto& directory : config.read_only_directories) + { + add_directory(directory, true); + } + + std::sort(files.begin(), files.end(), [](const auto& lhs, const auto& rhs) { + if (lhs.start_idx != rhs.start_idx) + { + return lhs.start_idx < rhs.start_idx; + } + return lhs.end_idx.has_value() && !rhs.end_idx.has_value(); + }); + return files; + } + + // NOLINTNEXTLINE(readability-function-cognitive-complexity) + static RecoverySnapshotLedgerScan scan_recovery_snapshot_ledger_files( + const ccf::CCFConfig::Ledger& ledger_config, + const std::shared_ptr& encryptor, + ccf::kv::Version snapshot_seqno) + { + if (snapshot_seqno == std::numeric_limits::max()) + { + throw std::logic_error( + "Snapshot seqno cannot be incremented for ledger scanning"); + } + + RecoverySnapshotLedgerScan scan; + auto expected_seqno = snapshot_seqno + 1; + size_t endorsements_payload_size = 0; + + for (const auto& ledger_file : + find_recovery_snapshot_ledger_files(ledger_config)) + { + if ( + ledger_file.end_idx.has_value() && + ledger_file.end_idx.value() < expected_seqno) + { + continue; + } + + std::ifstream file(ledger_file.path, std::ios::binary | std::ios::ate); + if (!file) + { + throw std::logic_error(fmt::format( + "Unable to open ledger file {}", ledger_file.path.string())); + } + + const auto file_size = static_cast(file.tellg()); + if (file_size < static_cast(sizeof(size_t))) + { + throw std::logic_error(fmt::format( + "Ledger file {} is too small", ledger_file.path.string())); + } + file.seekg(0); + + size_t positions_offset = 0; + file.read( + reinterpret_cast(&positions_offset), sizeof(positions_offset)); + if (!file) + { + throw std::logic_error(fmt::format( + "Unable to read ledger file header {}", ledger_file.path.string())); + } + if (ledger_file.committed && positions_offset == 0) + { + throw std::logic_error(fmt::format( + "Committed ledger file {} has no positions table", + ledger_file.path.string())); + } + + const auto entries_end = positions_offset == 0 ? + file_size : + static_cast(positions_offset); + if ( + entries_end < static_cast(sizeof(size_t)) || + entries_end > file_size) + { + throw std::logic_error(fmt::format( + "Ledger file {} has invalid positions table offset {}", + ledger_file.path.string(), + positions_offset)); + } + + std::optional first_file_version = std::nullopt; + std::optional previous_file_version = std::nullopt; + while (file.tellg() < entries_end) + { + const auto remaining = entries_end - file.tellg(); + if ( + remaining < + static_cast(sizeof(ccf::kv::SerialisedEntryHeader))) + { + if (positions_offset == 0) + { + break; + } + throw std::logic_error(fmt::format( + "Committed ledger file {} ends with a partial entry header", + ledger_file.path.string())); + } + + ccf::kv::SerialisedEntryHeader header{}; + file.read(reinterpret_cast(&header), sizeof(header)); + if (!file) + { + throw std::logic_error(fmt::format( + "Unable to read entry header from ledger file {}", + ledger_file.path.string())); + } + + const auto body_remaining = entries_end - file.tellg(); + if ( + header.size == 0 || header.size > static_cast(body_remaining)) + { + if (positions_offset == 0) + { + break; + } + throw std::logic_error(fmt::format( + "Committed ledger file {} contains a truncated entry", + ledger_file.path.string())); + } + if (header.size > MAX_RECOVERY_SNAPSHOT_LEDGER_ENTRY_SIZE) + { + throw std::logic_error(fmt::format( + "Ledger entry is too large ({} bytes; maximum {} bytes)", + static_cast(header.size), + MAX_RECOVERY_SNAPSHOT_LEDGER_ENTRY_SIZE)); + } + + std::vector entry(sizeof(header) + header.size); + std::memcpy(entry.data(), &header, sizeof(header)); + file.read( + reinterpret_cast(entry.data() + sizeof(header)), + static_cast(header.size)); + if (!file) + { + throw std::logic_error(fmt::format( + "Unable to read complete entry from ledger file {}", + ledger_file.path.string())); + } + + auto parsed = parse_recovery_snapshot_ledger_entry(entry, encryptor); + if (!first_file_version.has_value()) + { + first_file_version = parsed.version; + } + if (const auto previous_version_opt = previous_file_version) + { + const auto previous_version = *previous_version_opt; + if ( + previous_version == std::numeric_limits::max() || + parsed.version != previous_version + 1) + { + throw std::logic_error(fmt::format( + "Ledger file {} contains non-contiguous versions {} and {}", + ledger_file.path.string(), + previous_version, + parsed.version)); + } + } + previous_file_version = parsed.version; + + if (parsed.version < expected_seqno) + { + continue; + } + if (parsed.version > expected_seqno) + { + throw std::logic_error(fmt::format( + "Ledger suffix after snapshot is missing seqno {} (next entry is " + "{})", + expected_seqno, + parsed.version)); + } + + if (parsed.endorsement.has_value()) + { + const auto endorsement_size = parsed.endorsement->endorsement.size(); + if (endorsement_size > MAX_RECOVERY_SNAPSHOT_ENDORSEMENT_SIZE) + { + throw std::logic_error(fmt::format( + "Ledger endorsement at {} is too large ({} bytes; maximum {} " + "bytes)", + parsed.version, + endorsement_size, + MAX_RECOVERY_SNAPSHOT_ENDORSEMENT_SIZE)); + } + if ( + scan.endorsements.size() >= + MAX_RECOVERY_SNAPSHOT_ENDORSEMENTS_COUNT) + { + throw std::logic_error(fmt::format( + "Ledger suffix contains too many endorsements (maximum {})", + MAX_RECOVERY_SNAPSHOT_ENDORSEMENTS_COUNT)); + } + if ( + endorsement_size > MAX_RECOVERY_SNAPSHOT_ENDORSEMENTS_PAYLOAD_SIZE - + endorsements_payload_size) + { + throw std::logic_error(fmt::format( + "Ledger endorsements payload is too large (maximum {} bytes)", + MAX_RECOVERY_SNAPSHOT_ENDORSEMENTS_PAYLOAD_SIZE)); + } + endorsements_payload_size += endorsement_size; + scan.endorsements.push_back( + {parsed.version, std::move(*parsed.endorsement)}); + } + + if (expected_seqno == std::numeric_limits::max()) + { + throw std::logic_error( + "Ledger seqno overflow while scanning snapshot endorsements"); + } + ++expected_seqno; + } + + if (const auto first_version_opt = first_file_version) + { + const auto first_version = *first_version_opt; + if ( + first_version != static_cast(ledger_file.start_idx)) + { + throw std::logic_error(fmt::format( + "Ledger file {} does not start at its declared seqno {}", + ledger_file.path.string(), + ledger_file.start_idx)); + } + } + if (const auto end_idx_opt = ledger_file.end_idx) + { + const auto declared_end_idx = + static_cast(*end_idx_opt); + const auto previous_version_opt = previous_file_version; + if (!previous_version_opt || *previous_version_opt != declared_end_idx) + { + throw std::logic_error(fmt::format( + "Ledger file {} does not end at its declared seqno {}", + ledger_file.path.string(), + declared_end_idx)); + } + } + } + + return scan; + } +} diff --git a/src/node/rpc/network_identity_chain_helpers.h b/src/node/rpc/network_identity_chain_helpers.h index 8f1f154d842..0cac14d5844 100644 --- a/src/node/rpc/network_identity_chain_helpers.h +++ b/src/node/rpc/network_identity_chain_helpers.h @@ -2,8 +2,11 @@ // Licensed under the Apache 2.0 License. #pragma once +#include "ccf/crypto/cose_verifier.h" #include "ccf/tx_id.h" #include "consensus/aft/raft_types.h" +#include "crypto/cose.h" +#include "ds/internal_logger.h" #include "service/tables/previous_service_identity.h" #include @@ -11,6 +14,17 @@ namespace ccf { + struct CollectedCoseEndorsement + { + ccf::kv::Version write_version = ccf::kv::NoVersion; + ccf::CoseEndorsement endorsement; + }; + + inline std::string format_epoch(const std::optional& epoch_end) + { + return epoch_end.has_value() ? epoch_end->to_str() : "null"; + } + inline bool is_self_endorsement(const ccf::CoseEndorsement& endorsement) { return !endorsement.previous_version.has_value(); @@ -24,6 +38,74 @@ namespace ccf endorsement.endorsement_epoch_begin.seqno; } + inline void validate_fetched_endorsement( + const ccf::CoseEndorsement& endorsement) + { + LOG_INFO_FMT( + "Validating fetched endorsement from {} to {}", + endorsement.endorsement_epoch_begin.to_str(), + format_epoch(endorsement.endorsement_epoch_end)); + + if (!is_self_endorsement(endorsement)) + { + const auto [from, to] = + ccf::crypto::extract_cose_endorsement_validity(endorsement.endorsement); + + const auto from_txid = ccf::TxID::from_str(from); + if (!from_txid.has_value()) + { + throw std::logic_error(fmt::format( + "Cannot parse COSE endorsement header: {}", + ccf::cose::header::custom::TX_RANGE_BEGIN)); + } + + const auto to_txid = ccf::TxID::from_str(to); + if (!to_txid.has_value()) + { + throw std::logic_error(fmt::format( + "Cannot parse COSE endorsement header: {}", + ccf::cose::header::custom::TX_RANGE_END)); + } + + if (!endorsement.endorsement_epoch_end.has_value()) + { + throw std::logic_error( + "COSE endorsement does not contain epoch end in the table entry"); + } + if ( + endorsement.endorsement_epoch_begin != *from_txid || + *endorsement.endorsement_epoch_end != *to_txid) + { + throw std::logic_error(fmt::format( + "COSE endorsement fetched but range is invalid, epoch begin {}, " + "epoch end {}, header epoch begin: {}, header epoch end: {}", + endorsement.endorsement_epoch_begin.to_str(), + endorsement.endorsement_epoch_end->to_str(), + from, + to)); + } + } + } + + inline std::vector verify_cose_endorsement_signature( + std::span endorsement, + std::span endorsing_key) + { + auto verifier = ccf::crypto::make_cose_verifier_from_key(endorsing_key); + std::span endorsed_key; + if (!verifier->verify(endorsement, endorsed_key)) + { + throw std::logic_error("COSE endorsement failed signature verification"); + } + + if (endorsed_key.empty()) + { + throw std::logic_error("COSE endorsement contains an empty public key"); + } + + return {endorsed_key.begin(), endorsed_key.end()}; + } + inline void verify_endorsements_connected( const ccf::CoseEndorsement& newer, const ccf::CoseEndorsement& older) { diff --git a/src/node/rpc/network_identity_subsystem.h b/src/node/rpc/network_identity_subsystem.h index 60fe40d8e4f..bfd08a7409d 100644 --- a/src/node/rpc/network_identity_subsystem.h +++ b/src/node/rpc/network_identity_subsystem.h @@ -16,60 +16,6 @@ namespace ccf { - inline std::string format_epoch(const std::optional& epoch_end) - { - return epoch_end.has_value() ? epoch_end->to_str() : "null"; - } - - inline void validate_fetched_endorsement( - const ccf::CoseEndorsement& endorsement) - { - LOG_INFO_FMT( - "Validating fetched endorsement from {} to {}", - endorsement.endorsement_epoch_begin.to_str(), - format_epoch(endorsement.endorsement_epoch_end)); - - if (!is_self_endorsement(endorsement)) - { - const auto [from, to] = - ccf::crypto::extract_cose_endorsement_validity(endorsement.endorsement); - - const auto from_txid = ccf::TxID::from_str(from); - if (!from_txid.has_value()) - { - throw std::logic_error(fmt::format( - "Cannot parse COSE endorsement header: {}", - ccf::cose::header::custom::TX_RANGE_BEGIN)); - } - - const auto to_txid = ccf::TxID::from_str(to); - if (!to_txid.has_value()) - { - throw std::logic_error(fmt::format( - "Cannot parse COSE endorsement header: {}", - ccf::cose::header::custom::TX_RANGE_END)); - } - - if (!endorsement.endorsement_epoch_end.has_value()) - { - throw std::logic_error( - "COSE endorsement does not contain epoch end in the table entry"); - } - if ( - endorsement.endorsement_epoch_begin != *from_txid || - *endorsement.endorsement_epoch_end != *to_txid) - { - throw std::logic_error(fmt::format( - "COSE endorsement fetched but range is invalid, epoch begin {}, " - "epoch end {}, header epoch begin: {}, header epoch end: {}", - endorsement.endorsement_epoch_begin.to_str(), - endorsement.endorsement_epoch_end->to_str(), - from, - to)); - } - } - } - class NetworkIdentitySubsystem : public NetworkIdentitySubsystemInterface { protected: @@ -495,10 +441,13 @@ namespace ccf std::span previous_key_der{}; for (const auto& [seqno, endorsement] : endorsements) { - auto verifier = - ccf::crypto::make_cose_verifier_from_key(endorsement.endorsing_key); - std::span endorsed_key; - if (!verifier->verify(endorsement.endorsement, endorsed_key)) + std::vector endorsed_key; + try + { + endorsed_key = verify_cose_endorsement_signature( + endorsement.endorsement, endorsement.endorsing_key); + } + catch (const std::logic_error&) { throw std::logic_error(fmt::format( "COSE endorsement chain integrity is violated, endorsement from " diff --git a/src/node/snapshot_serdes.h b/src/node/snapshot_serdes.h index 0ac29955eab..c2b776da50f 100644 --- a/src/node/snapshot_serdes.h +++ b/src/node/snapshot_serdes.h @@ -15,9 +15,11 @@ #include "kv/serialised_entry_format.h" #include "node/cose_common.h" #include "node/history.h" +#include "node/rpc/network_identity_chain_helpers.h" #include "node/tx_receipt_impl.h" #include +#include namespace ccf { @@ -76,9 +78,176 @@ namespace ccf return SnapshotSegments{header_and_body, receipt}; } - static void verify_cose_snapshot_receipt( + static void verify_snapshot_seqno( const SnapshotSegments& segments, - const std::optional>& prev_service_identity) + const std::shared_ptr& encryptor, + ccf::kv::Version expected_seqno) + { + auto deserialiser = ccf::kv::RawKvStoreDeserialiser( + encryptor, ccf::kv::SecurityDomain::PUBLIC); + ccf::kv::Term term = 0; + ccf::kv::EntryFlags flags = {}; + const auto snapshot_seqno = deserialiser.init( + segments.header_and_body.data(), + segments.header_and_body.size(), + term, + flags, + false); + if (!snapshot_seqno.has_value()) + { + throw std::logic_error("Failed to read version from recovery snapshot"); + } + if (*snapshot_seqno != expected_seqno) + { + throw std::logic_error(fmt::format( + "Recovery snapshot body is at seqno {}, but its file name claims {}", + *snapshot_seqno, + expected_seqno)); + } + } + + template + static std::optional try_verify_and_install_recovery_snapshot( + Verify&& verify, Install&& install) + { + try + { + std::forward(verify)(); + } + catch (const std::exception& e) + { + return e.what(); + } + + std::forward(install)(); + return std::nullopt; + } + + // Validates the collected COSE endorsement chain against the configured + // previous service identity (target_key) and returns the snapshot signer + // key: the identity, active at snapshot_seqno, that signed the snapshot + // receipt. The oldest endorsement covers snapshot_seqno and endorses that + // key, so it is captured as the endorsed key of the first (oldest) link. + static std::vector validate_recovery_snapshot_endorsement_chain( + const std::vector& collected, + std::span target_key, + ccf::kv::Version snapshot_seqno) + { + if (collected.empty()) + { + throw std::logic_error( + "No previous service identity endorsements were found after the " + "snapshot"); + } + + std::vector snapshot_signer_key; + std::vector previous_endorsing_key; + for (size_t i = 0; i < collected.size(); ++i) + { + const auto& [write_version, endorsement] = collected[i]; + if (write_version <= snapshot_seqno) + { + throw std::logic_error(fmt::format( + "Collected endorsement write at {} is not after snapshot seqno {}", + write_version, + snapshot_seqno)); + } + if (is_self_endorsement(endorsement)) + { + throw std::logic_error(fmt::format( + "Unexpected self-endorsement after snapshot at {}", + endorsement.endorsement_epoch_begin.to_str())); + } + if (has_ill_formed_epoch_range(endorsement)) + { + throw std::logic_error(fmt::format( + "Collected endorsement has an ill-formed epoch range {} - {}", + endorsement.endorsement_epoch_begin.to_str(), + format_epoch(endorsement.endorsement_epoch_end))); + } + + validate_fetched_endorsement(endorsement); + + if (i == 0) + { + if ( + !endorsement.previous_version.has_value() || + endorsement.previous_version.value() > snapshot_seqno) + { + throw std::logic_error(fmt::format( + "Oldest collected endorsement at {} does not point to an " + "endorsement at or before snapshot seqno {}", + write_version, + snapshot_seqno)); + } + } + else + { + const auto& previous = collected[i - 1]; + if ( + !endorsement.previous_version.has_value() || + endorsement.previous_version.value() != previous.write_version) + { + throw std::logic_error(fmt::format( + "Collected endorsement at {} does not point to the previous " + "collected endorsement at {}", + write_version, + previous.write_version)); + } + verify_endorsements_connected(endorsement, previous.endorsement); + } + + const auto endorsed_key = verify_cose_endorsement_signature( + endorsement.endorsement, endorsement.endorsing_key); + if (i == 0) + { + // The oldest endorsement covers snapshot_seqno; the key it endorses is + // the identity that signed the snapshot receipt. + snapshot_signer_key = endorsed_key; + } + if ( + !previous_endorsing_key.empty() && + endorsed_key != previous_endorsing_key) + { + throw std::logic_error(fmt::format( + "Collected endorsement at {} does not endorse the preceding service " + "identity", + write_version)); + } + previous_endorsing_key = endorsement.endorsing_key; + } + + if ( + previous_endorsing_key.size() != target_key.size() || + !std::equal( + previous_endorsing_key.begin(), + previous_endorsing_key.end(), + target_key.begin())) + { + throw std::logic_error( + "Newest collected endorsement is not signed by the configured " + "previous service identity"); + } + + const auto& oldest = collected.front().endorsement; + if ( + !oldest.endorsement_epoch_end.has_value() || + oldest.endorsement_epoch_begin.seqno > snapshot_seqno || + oldest.endorsement_epoch_end->seqno < snapshot_seqno) + { + throw std::logic_error(fmt::format( + "Oldest collected endorsement range {} - {} does not cover snapshot " + "seqno {}", + oldest.endorsement_epoch_begin.to_str(), + format_epoch(oldest.endorsement_epoch_end), + snapshot_seqno)); + } + + return snapshot_signer_key; + } + + static auto decode_and_verify_cose_snapshot_receipt( + const SnapshotSegments& segments) { auto receipt = ccf::cose::decode_ccf_receipt( {segments.receipt.begin(), segments.receipt.end()}, @@ -99,6 +268,15 @@ namespace ccf ds::to_hex(receipt.claims_digest))); } + return receipt; + } + + static void verify_cose_snapshot_receipt( + const SnapshotSegments& segments, + const std::optional>& prev_service_identity) + { + const auto receipt = decode_and_verify_cose_snapshot_receipt(segments); + if (prev_service_identity) { auto verifier = ccf::crypto::make_cose_verifier_from_pem_cert( @@ -215,6 +393,29 @@ namespace ccf } } + // Verify the recovery snapshot receipt is signed by snapshot_signer_key, the + // snapshot service identity established by validating the endorsement chain. + static void verify_recovery_snapshot_receipt( + const SnapshotSegments& segments, + std::span snapshot_signer_key) + { + if (segments.receipt.empty() || segments.receipt[0] != 0xD2) + { + throw std::logic_error( + "Only snapshots with COSE receipts can use an endorsement chain"); + } + + const auto receipt = decode_and_verify_cose_snapshot_receipt(segments); + const auto verifier = + ccf::crypto::make_cose_verifier_from_key(snapshot_signer_key); + if (!verifier->verify_detached(segments.receipt, receipt.merkle_root)) + { + throw std::logic_error( + "Snapshot receipt signature verification failed under the endorsed " + "snapshot service identity"); + } + } + static void deserialise_snapshot( const std::shared_ptr& store, const SnapshotSegments& segments, diff --git a/src/node/test/network_identity_subsystem.cpp b/src/node/test/network_identity_subsystem.cpp index 0d6f0dd8ea1..305241b7840 100644 --- a/src/node/test/network_identity_subsystem.cpp +++ b/src/node/test/network_identity_subsystem.cpp @@ -7,6 +7,7 @@ #include "crypto/openssl/ec_key_pair.h" #include "node/rpc/network_identity_accessors.h" #include "node/rpc/network_identity_chain_helpers.h" +#include "node/snapshot_serdes.h" #define DOCTEST_CONFIG_IMPLEMENT_WITH_MAIN #include @@ -475,6 +476,69 @@ TEST_CASE( ccf::validate_chain_front_connection(bad, after), std::logic_error); } +TEST_CASE("Recovery snapshot endorsements are validated newest-to-oldest") +{ + ChainBuilder cb; + cb.add_self({2, 1}).add_next({2, 1}, {4, 100}).add_next({6, 101}, {8, 200}); + + std::vector collected{ + {cb.write_versions[1], cb.entries[1]}, + {cb.write_versions[2], cb.entries[2]}}; + const auto target_key = cb.current_pkey_der(); + + const auto snapshot_signer = + ccf::validate_recovery_snapshot_endorsement_chain(collected, target_key, 1); + REQUIRE(snapshot_signer == cb.service_keys.front()->public_key_der()); +} + +TEST_CASE("Recovery snapshot endorsements fail closed") +{ + ChainBuilder cb; + cb.add_self({2, 1}).add_next({2, 1}, {4, 100}).add_next({6, 101}, {8, 200}); + + std::vector collected{ + {cb.write_versions[1], cb.entries[1]}, + {cb.write_versions[2], cb.entries[2]}}; + const auto target_key = cb.current_pkey_der(); + + // Snapshot seqno inconsistent with the collected endorsement chain. + REQUIRE_THROWS_AS( + ccf::validate_recovery_snapshot_endorsement_chain( + collected, target_key, 1000), + std::logic_error); + + const auto wrong_target_key = + ccf::crypto::make_ec_key_pair()->public_key_der(); + REQUIRE_THROWS_AS( + ccf::validate_recovery_snapshot_endorsement_chain( + collected, wrong_target_key, 1), + std::logic_error); + + auto broken_back_pointer = collected; + broken_back_pointer.back().endorsement.previous_version = 999; + REQUIRE_THROWS_AS( + ccf::validate_recovery_snapshot_endorsement_chain( + broken_back_pointer, target_key, 1), + std::logic_error); + + auto epoch_mismatch = collected; + epoch_mismatch.front().endorsement.endorsement_epoch_end->seqno -= 1; + REQUIRE_THROWS_AS( + ccf::validate_recovery_snapshot_endorsement_chain( + epoch_mismatch, target_key, 1), + std::logic_error); + + auto injected = collected; + auto injected_endorsement = collected.back(); + ++injected_endorsement.write_version; + injected_endorsement.endorsement.previous_version = + collected.back().write_version; + injected.emplace_back(std::move(injected_endorsement)); + REQUIRE_THROWS_AS( + ccf::validate_recovery_snapshot_endorsement_chain(injected, target_key, 1), + std::logic_error); +} + // ------------------------------------------------------------------------ // State-machine tests using Mock{NodeState,HistoricalState}Accessor + // FakeTaskScheduler + ChainBuilder. Each test wires a synthetic ledger, diff --git a/src/node/test/snapshotter.cpp b/src/node/test/snapshotter.cpp index 4386522a8eb..9c6e7914cae 100644 --- a/src/node/test/snapshotter.cpp +++ b/src/node/test/snapshotter.cpp @@ -10,12 +10,15 @@ #include "kv/test/stub_consensus.h" #include "node/encryptor.h" #include "node/history.h" +#include "node/recovery_snapshot_ledger.h" +#include "node/snapshot_serdes.h" #include "snapshots/filenames.h" #define DOCTEST_CONFIG_IMPLEMENT #include #include #include +#include #include #include @@ -54,6 +57,226 @@ struct ScopedSnapshotDir } }; +void write_current_ledger_file( + const fs::path& path, const std::vector>& entries) +{ + std::ofstream ledger_file(path, std::ios::binary); + REQUIRE(ledger_file); + const size_t positions_offset = 0; + ledger_file.write( + reinterpret_cast(&positions_offset), sizeof(positions_offset)); + for (const auto& entry : entries) + { + ledger_file.write( + reinterpret_cast(entry.data()), + static_cast(entry.size())); + } + REQUIRE(ledger_file); +} + +TEST_CASE("Recovery snapshot installation failures do not request fallback") +{ + ccf::kv::Store store; + store.set_readiness(ccf::kv::StoreReadiness::InstallingSnapshot); + bool install_started = false; + bool fallback_requested = false; + + REQUIRE_THROWS_AS( + [&]() { + const auto verification_error = + ccf::try_verify_and_install_recovery_snapshot( + []() {}, + [&]() { + install_started = true; + store.set_readiness(ccf::kv::StoreReadiness::Failed); + throw std::logic_error("snapshot installation failed"); + }); + fallback_requested = verification_error.has_value(); + }(), + std::logic_error); + REQUIRE(install_started); + REQUIRE(store.get_readiness() == ccf::kv::StoreReadiness::Failed); + REQUIRE_FALSE(fallback_requested); + + bool install_called = false; + const auto verification_error = ccf::try_verify_and_install_recovery_snapshot( + []() { throw std::logic_error("snapshot verification failed"); }, + [&]() { install_called = true; }); + REQUIRE(verification_error.has_value()); + REQUIRE_FALSE(install_called); +} + +TEST_CASE("Recovery snapshot endorsement scan reads ledger files directly") +{ + ScopedSnapshotDir ledger_dir; + ccf::kv::Store source_store; + auto encryptor = std::make_shared(); + auto consensus = std::make_shared(); + source_store.set_encryptor(encryptor); + source_store.set_consensus(consensus); + source_store.initialise_term(2); + + std::vector> entries; + { + auto tx = source_store.create_tx(); + tx.rw("public:unrelated")->put("key", "value"); + REQUIRE(tx.commit() == ccf::kv::CommitResult::SUCCESS); + auto latest_entry = + consensus->get_latest_data().value_or(std::vector{}); + REQUIRE_FALSE(latest_entry.empty()); + entries.push_back(std::move(latest_entry)); + } + { + ccf::CoseEndorsement endorsement; + endorsement.endorsement = {0xd2, 0x01}; + endorsement.endorsing_key = {0x02, 0x03}; + endorsement.endorsement_epoch_begin = {2, 1}; + endorsement.endorsement_epoch_end = ccf::TxID{4, 1}; + endorsement.previous_version = 1; + + auto tx = source_store.create_tx(); + tx.rw( + ccf::Tables::PREVIOUS_SERVICE_IDENTITY_ENDORSEMENT) + ->put(endorsement); + REQUIRE(tx.commit() == ccf::kv::CommitResult::SUCCESS); + auto latest_entry = + consensus->get_latest_data().value_or(std::vector{}); + REQUIRE_FALSE(latest_entry.empty()); + entries.push_back(std::move(latest_entry)); + } + + const ccf::SnapshotSegments first_entry{ + std::span(entries.front()), {}}; + REQUIRE_NOTHROW(ccf::verify_snapshot_seqno(first_entry, encryptor, 1)); + REQUIRE_THROWS(ccf::verify_snapshot_seqno(first_entry, encryptor, 2)); + + auto malformed_entry = entries.front(); + const auto public_domain_size_offset = + sizeof(ccf::kv::SerialisedEntryHeader) + encryptor->get_header_length(); + const auto invalid_public_domain_size = malformed_entry.size(); + std::memcpy( + malformed_entry.data() + public_domain_size_offset, + &invalid_public_domain_size, + sizeof(invalid_public_domain_size)); + + ccf::kv::Store untouched_store; + const auto original_version = untouched_store.current_version(); + const auto original_readiness = untouched_store.get_readiness(); + + bool snapshot_install_called = false; + const ccf::SnapshotSegments malformed_snapshot{ + std::span(malformed_entry), {}}; + const auto snapshot_error = ccf::try_verify_and_install_recovery_snapshot( + [&]() { ccf::verify_snapshot_seqno(malformed_snapshot, encryptor, 1); }, + [&]() { snapshot_install_called = true; }); + REQUIRE(snapshot_error.has_value()); + REQUIRE_FALSE(snapshot_install_called); + REQUIRE(untouched_store.current_version() == original_version); + REQUIRE(untouched_store.get_readiness() == original_readiness); + + ScopedSnapshotDir malformed_ledger_dir; + write_current_ledger_file( + malformed_ledger_dir.path / "ledger_1", {malformed_entry}); + ccf::CCFConfig::Ledger malformed_ledger_config; + malformed_ledger_config.directory = malformed_ledger_dir.path.string(); + bool scanner_install_called = false; + const auto scanner_error = ccf::try_verify_and_install_recovery_snapshot( + [&]() { + std::ignore = ccf::scan_recovery_snapshot_ledger_files( + malformed_ledger_config, encryptor, 0); + }, + [&]() { scanner_install_called = true; }); + REQUIRE(scanner_error.has_value()); + REQUIRE_FALSE(scanner_install_called); + REQUIRE(untouched_store.current_version() == original_version); + REQUIRE(untouched_store.get_readiness() == original_readiness); + + write_current_ledger_file(ledger_dir.path / "ledger_1", entries); + + ccf::CCFConfig::Ledger ledger_config; + ledger_config.directory = ledger_dir.path.string(); + const auto scan = + ccf::scan_recovery_snapshot_ledger_files(ledger_config, encryptor, 1); + REQUIRE(scan.endorsements.size() == 1); + REQUIRE(scan.endorsements.front().write_version == 2); + + bool install_called = false; + const auto verification_error = ccf::try_verify_and_install_recovery_snapshot( + [&]() { + const auto target_key = ccf::crypto::make_ec_key_pair()->public_key_der(); + std::ignore = ccf::validate_recovery_snapshot_endorsement_chain( + scan.endorsements, target_key, 1); + }, + [&]() { install_called = true; }); + REQUIRE(verification_error.has_value()); + REQUIRE_FALSE(install_called); +} + +TEST_CASE("Recovery snapshot endorsement scan bounds candidate endorsements") +{ + ScopedSnapshotDir ledger_dir; + ccf::kv::Store source_store; + auto encryptor = std::make_shared(); + auto consensus = std::make_shared(); + source_store.set_encryptor(encryptor); + source_store.set_consensus(consensus); + source_store.initialise_term(2); + + std::vector> entries; + for (size_t i = 0; i < ccf::MAX_RECOVERY_SNAPSHOT_ENDORSEMENTS_COUNT + 1; ++i) + { + ccf::CoseEndorsement endorsement; + endorsement.endorsement = {0xd2, 0x01}; + endorsement.endorsing_key = {0x02, 0x03}; + endorsement.endorsement_epoch_begin = {2, i + 1}; + endorsement.endorsement_epoch_end = ccf::TxID{4, i + 1}; + endorsement.previous_version = 1; + + auto tx = source_store.create_tx(); + tx.rw( + ccf::Tables::PREVIOUS_SERVICE_IDENTITY_ENDORSEMENT) + ->put(endorsement); + REQUIRE(tx.commit() == ccf::kv::CommitResult::SUCCESS); + auto latest_entry = + consensus->get_latest_data().value_or(std::vector{}); + REQUIRE_FALSE(latest_entry.empty()); + entries.push_back(std::move(latest_entry)); + } + + write_current_ledger_file(ledger_dir.path / "ledger_1", entries); + + ccf::CCFConfig::Ledger ledger_config; + ledger_config.directory = ledger_dir.path.string(); + REQUIRE_THROWS( + ccf::scan_recovery_snapshot_ledger_files(ledger_config, encryptor, 0)); +} + +TEST_CASE("Recovery snapshot endorsement scan bounds ledger entry allocation") +{ + ScopedSnapshotDir ledger_dir; + const auto ledger_path = ledger_dir.path / "ledger_1"; + { + std::ofstream ledger_file(ledger_path, std::ios::binary); + REQUIRE(ledger_file); + const size_t positions_offset = 0; + ledger_file.write( + reinterpret_cast(&positions_offset), + sizeof(positions_offset)); + ccf::kv::SerialisedEntryHeader header{}; + header.size = ccf::MAX_RECOVERY_SNAPSHOT_LEDGER_ENTRY_SIZE + 1; + ledger_file.write(reinterpret_cast(&header), sizeof(header)); + ledger_file.seekp( + static_cast(header.size) - 1, std::ios::cur); + ledger_file.put(0); + REQUIRE(ledger_file); + } + + ccf::CCFConfig::Ledger ledger_config; + ledger_config.directory = ledger_dir.path.string(); + REQUIRE_THROWS(ccf::scan_recovery_snapshot_ledger_files( + ledger_config, std::make_shared(), 0)); +} + std::optional latest_committed_snapshot_path(const fs::path& dir) { return snapshots::find_latest_committed_snapshot_in_directory(dir); diff --git a/tests/ci-buckets.txt b/tests/ci-buckets.txt index 4ad9a4f5377..7ca731569c2 100644 --- a/tests/ci-buckets.txt +++ b/tests/ci-buckets.txt @@ -5,6 +5,7 @@ bucket_b: recovery_test recovery_stale_snapshot_join_test recovery_intermediate_snapshot_join_test + recovery_snapshot_endorsements_test schema_test nodes_test diff --git a/tests/recovery_snapshot_endorsements.py b/tests/recovery_snapshot_endorsements.py new file mode 100644 index 00000000000..16b71c3c62a --- /dev/null +++ b/tests/recovery_snapshot_endorsements.py @@ -0,0 +1,276 @@ +# Copyright (c) Microsoft Corporation. All rights reserved. +# Licensed under the Apache 2.0 License. + +import copy +import hashlib +import os +import shutil + +import ccf.ledger +import infra.e2e_args +import infra.logging_app as app +import infra.network +import infra.node +from infra.runner import ConcurrentRunner +from loguru import logger as LOG + + +def _logs(node): + out_path, _ = node.get_logs() + assert out_path is not None + with open(out_path, encoding="utf-8") as output: + return output.read() + + +def _stop_incomplete_recovery(network): + network.stop_all_nodes( + skip_verification=True, + skip_verify_chunking=True, + check_file_invariants=False, + ) + + +def _recover_and_open(network, args, label): + recovery_args = copy.deepcopy(args) + recovery_args.label = label + network.save_service_identity(recovery_args) + primary, _ = network.find_primary() + network.stop_all_nodes() + current_ledger_dir, committed_ledger_dirs = primary.get_ledger() + + recovered = infra.network.Network( + recovery_args.nodes, + recovery_args.binary_dir, + recovery_args.debug_nodes, + existing_network=network, + ) + recovered.start_in_recovery( + recovery_args, + ledger_dir=current_ledger_dir, + committed_ledger_dirs=committed_ledger_dirs, + ) + recovered.recover(recovery_args) + + app.LoggingTxs("user0").issue( + recovered, + number_txs=2, + send_private=False, + send_public=True, + wait_for_sync=True, + ) + recovered.get_latest_ledger_public_state() + return recovered, recovery_args + + +def _start_recovery_attempt( + base_network, + args, + label, + ledger_dir, + committed_ledger_dirs, + snapshots_dir, + previous_service_identity_file, + next_node_id, +): + attempt_args = copy.deepcopy(args) + attempt_args.label = label + attempt_args.previous_service_identity_file = previous_service_identity_file + attempt = infra.network.Network( + attempt_args.nodes, + attempt_args.binary_dir, + attempt_args.debug_nodes, + existing_network=base_network, + next_node_id=next_node_id, + ) + attempt.ignore_errors_on_shutdown() + attempt.start_in_recovery( + attempt_args, + ledger_dir=ledger_dir, + committed_ledger_dirs=committed_ledger_dirs, + snapshots_dir=snapshots_dir, + common_dir=base_network.common_dir, + ) + return attempt + + +def _copy_ledger_prefix(source_dirs, destination, first_excluded_seqno): + shutil.rmtree(destination, ignore_errors=True) + os.makedirs(destination) + copied = 0 + for source_dir in source_dirs: + for name in os.listdir(source_dir): + if not ccf.ledger.is_ledger_chunk_committed(name): + continue + _, end_seqno = ccf.ledger.range_from_filename(name) + if end_seqno is None or end_seqno >= first_excluded_seqno: + continue + destination_path = os.path.join(destination, name) + if not os.path.exists(destination_path): + shutil.copy(os.path.join(source_dir, name), destination_path) + copied += 1 + assert copied > 0 + + +def _assert_node_snapshot_unchanged( + network, node, snapshot_name, expected_snapshot_digest +): + snapshots_dir = network.get_committed_snapshots(node, force_txs=False) + snapshot_path = os.path.join(snapshots_dir, snapshot_name) + assert os.path.isfile(snapshot_path), snapshot_path + with open(snapshot_path, "rb") as snapshot_file: + assert hashlib.sha256(snapshot_file.read()).digest() == expected_snapshot_digest + + +def run_recovery_snapshot_endorsements(args): + with infra.network.network( + args.nodes, + args.binary_dir, + args.debug_nodes, + pdb=args.pdb, + ) as initial_network: + initial_network.start_and_open(args) + primary, _ = initial_network.find_primary() + + app.LoggingTxs("user0").issue( + initial_network, + number_txs=2, + send_private=False, + send_public=True, + wait_for_sync=True, + ) + snapshot_trigger = primary.trigger_snapshot() + committed_snapshots_dir = initial_network.get_committed_snapshots( + primary, + target_seqno=snapshot_trigger.seqno, + wait_for_target_seqno=True, + ) + snapshots = sorted( + ( + name + for name in os.listdir(committed_snapshots_dir) + if name.startswith("snapshot_") + and ccf.ledger.is_snapshot_file_committed(name) + ), + key=lambda name: infra.node.get_snapshot_seqnos(name)[0], + ) + assert snapshots + snapshot_name = snapshots[-1] + original_snapshot_path = os.path.join(committed_snapshots_dir, snapshot_name) + + source_snapshots_dir = os.path.join( + initial_network.common_dir, "recovery_snapshot_endorsements_source" + ) + shutil.rmtree(source_snapshots_dir, ignore_errors=True) + os.makedirs(source_snapshots_dir) + source_snapshot_path = shutil.copy(original_snapshot_path, source_snapshots_dir) + with open(source_snapshot_path, "rb") as snapshot_file: + source_snapshot_bytes = snapshot_file.read() + snapshot_digest = hashlib.sha256(source_snapshot_bytes).digest() + + first_recovery, first_args = _recover_and_open( + initial_network, args, f"{args.label}_identity_1" + ) + second_recovery, second_args = _recover_and_open( + first_recovery, first_args, f"{args.label}_identity_2" + ) + + second_recovery.save_service_identity(second_args) + target_identity_file = second_args.previous_service_identity_file + primary, _ = second_recovery.find_primary() + with primary.client() as client: + service_create_txid = client.get("/node/network").body.json()[ + "current_service_create_txid" + ] + service_create_seqno = int(service_create_txid.split(".")[1]) + second_recovery.stop_all_nodes() + current_ledger_dir, committed_ledger_dirs = primary.get_ledger() + + valid_attempt = _start_recovery_attempt( + second_recovery, + second_args, + f"{args.label}_in_memory_chain", + current_ledger_dir, + committed_ledger_dirs, + source_snapshots_dir, + target_identity_file, + 100, + ) + try: + valid_primary, _ = valid_attempt.find_primary() + logs = _logs(valid_primary) + scan_log = "scanning the public ledger suffix for COSE endorsements" + validated_log = "Validated 2 recovery snapshot endorsement(s) in memory" + snapshot_body_log = "Deserialising snapshot (size:" + public_recovery_log = "Starting to read public ledger" + assert ( + logs.index(scan_log) + < logs.index(validated_log) + < logs.index(snapshot_body_log) + < logs.index(public_recovery_log) + ) + _assert_node_snapshot_unchanged( + valid_attempt, valid_primary, snapshot_name, snapshot_digest + ) + finally: + _stop_incomplete_recovery(valid_attempt) + + incomplete_ledger_dir = os.path.join( + second_recovery.common_dir, "recovery_snapshot_incomplete_ledger" + ) + shutil.rmtree(incomplete_ledger_dir, ignore_errors=True) + os.makedirs(incomplete_ledger_dir) + incomplete_committed_ledger_dir = os.path.join( + second_recovery.common_dir, + "recovery_snapshot_incomplete_committed_ledger", + ) + _copy_ledger_prefix( + [current_ledger_dir, *committed_ledger_dirs], + incomplete_committed_ledger_dir, + service_create_seqno, + ) + fallback_attempt = _start_recovery_attempt( + second_recovery, + second_args, + f"{args.label}_incomplete_suffix", + incomplete_ledger_dir, + [incomplete_committed_ledger_dir], + source_snapshots_dir, + target_identity_file, + 101, + ) + try: + fallback_primary, _ = fallback_attempt.find_primary() + logs = _logs(fallback_primary) + assert "Falling back to full-ledger recovery" in logs + assert "Setting startup snapshot seqno" not in logs + _assert_node_snapshot_unchanged( + fallback_attempt, fallback_primary, snapshot_name, snapshot_digest + ) + finally: + _stop_incomplete_recovery(fallback_attempt) + + LOG.success( + "In-memory recovery snapshot endorsement validation and " + "incomplete-suffix fallback succeeded" + ) + + +if __name__ == "__main__": + + def add(parser): + parser.description = ( + "Verify in-memory recovery snapshot endorsement chains across multiple " + "disaster recoveries." + ) + + cr = ConcurrentRunner(add) + cr.add( + "recovery_snapshot_endorsements", + run_recovery_snapshot_endorsements, + package="samples/apps/logging/logging", + nodes=infra.e2e_args.min_nodes(cr.args, f=0), + ledger_chunk_bytes="50KB", + snapshot_tx_interval=10, + sig_tx_interval=1, + ) + cr.run()