From 01af38cf456e1ebf2b85f78d347fb402b6adad78 Mon Sep 17 00:00:00 2001 From: Eddy Ashton Date: Fri, 21 Aug 2026 14:49:07 +0000 Subject: [PATCH 01/12] history doesn't need to know about term_of_next_version --- src/kv/committable_tx.h | 3 +-- src/kv/kv_types.h | 10 +++------- src/kv/store.h | 8 +++----- src/node/history.h | 38 ++++++++------------------------------ src/node/node_state.h | 2 +- src/node/rpc/frontend.h | 4 ++-- 6 files changed, 18 insertions(+), 47 deletions(-) diff --git a/src/kv/committable_tx.h b/src/kv/committable_tx.h index f383563bf14a..8c49d469069d 100644 --- a/src/kv/committable_tx.h +++ b/src/kv/committable_tx.h @@ -328,14 +328,13 @@ namespace ccf::kv return TxID(pimpl->commit_view, version); } - void set_read_txid(const TxID& tx_id, Term commit_view_) + void set_read_txid(const TxID& tx_id) { if (pimpl->read_txid.has_value()) { throw std::logic_error("Read TxID already set"); } pimpl->read_txid = tx_id; - pimpl->commit_view = commit_view_; } void set_root_at_read_version(const ccf::crypto::Sha256Hash& r) diff --git a/src/kv/kv_types.h b/src/kv/kv_types.h index 7a85981ad16d..ea0663308f46 100644 --- a/src/kv/kv_types.h +++ b/src/kv/kv_types.h @@ -393,8 +393,7 @@ namespace ccf::kv virtual ccf::crypto::Sha256Hash get_replicated_state_root() = 0; virtual std::tuple< ccf::TxID /* TxID of last transaction seen by history */, - ccf::crypto::Sha256Hash /* root as of TxID */, - ccf::kv::Term /* term_of_next_version */> + ccf::crypto::Sha256Hash /* root as of TxID */> get_replicated_state_txid_and_root() = 0; virtual std::vector get_proof(Version v) = 0; virtual bool verify_proof(const std::vector& proof) = 0; @@ -402,11 +401,8 @@ namespace ccf::kv const std::vector& hash_at_snapshot) = 0; virtual std::vector get_raw_leaf(uint64_t index) = 0; virtual void append(const std::vector& data) = 0; - virtual void append_entry( - const ccf::crypto::Sha256Hash& digest, - std::optional expected_term = std::nullopt) = 0; - virtual void rollback( - const ccf::TxID& tx_id, ccf::kv::Term term_of_next_version_) = 0; + virtual void append_entry(const ccf::crypto::Sha256Hash& digest) = 0; + virtual void rollback(const ccf::TxID& tx_id) = 0; virtual void compact(Version v) = 0; virtual void set_term(ccf::kv::Term) = 0; virtual std::vector serialise_tree(size_t to) = 0; diff --git a/src/kv/store.h b/src/kv/store.h index dce78e55f41d..8155b9dfba7a 100644 --- a/src/kv/store.h +++ b/src/kv/store.h @@ -665,7 +665,7 @@ namespace ccf::kv auto h = get_history(); if (h) { - h->rollback(tx_id, term_of_next_version); + h->rollback(tx_id); } if (tx_id.seqno >= version) @@ -1049,10 +1049,8 @@ namespace ccf::kv if (h) { - h->append_entry( - ccf::entry_leaf( - *data_shared, commit_evidence_digest_, claims_digest_), - replication_view); + h->append_entry(ccf::entry_leaf( + *data_shared, commit_evidence_digest_, claims_digest_)); } if (chunker) diff --git a/src/node/history.h b/src/node/history.h index db80959d33e6..a18791bcbfd0 100644 --- a/src/node/history.h +++ b/src/node/history.h @@ -119,7 +119,6 @@ namespace ccf protected: ccf::kv::Version version = 0; ccf::kv::Term term_of_last_version = 0; - ccf::kv::Term term_of_next_version = 0; public: NullTxHistory( @@ -133,10 +132,7 @@ namespace ccf version++; } - void append_entry( - const ccf::crypto::Sha256Hash& /*digest*/, - std::optional /*term_of_next_version_*/ = - std::nullopt) override + void append_entry(const ccf::crypto::Sha256Hash& /*digest*/) override { version++; } @@ -149,14 +145,12 @@ namespace ccf void set_term(ccf::kv::Term t) override { term_of_last_version = t; - term_of_next_version = t; } - void rollback(const ccf::TxID& tx_id, ccf::kv::Term commit_term_) override + void rollback(const ccf::TxID& tx_id) override { version = tx_id.seqno; term_of_last_version = tx_id.view; - term_of_next_version = commit_term_; } void compact(ccf::kv::Version /*v*/) override {} @@ -201,13 +195,12 @@ namespace ccf return ccf::crypto::Sha256Hash(std::to_string(version)); } - std::tuple + std::tuple get_replicated_state_txid_and_root() override { return { {term_of_last_version, version}, - ccf::crypto::Sha256Hash(std::to_string(version)), - term_of_next_version}; + ccf::crypto::Sha256Hash(std::to_string(version))}; } std::vector get_proof(ccf::kv::Version /*v*/) override @@ -575,7 +568,6 @@ namespace ccf ccf::pal::Mutex state_lock; ccf::kv::Term term_of_last_version = 0; - ccf::kv::Term term_of_next_version{}; std::optional endorsed_cert = std::nullopt; @@ -757,15 +749,14 @@ namespace ccf return replicated_state_tree.get_root(); } - std::tuple + std::tuple get_replicated_state_txid_and_root() override { std::lock_guard guard(state_lock); return { {term_of_last_version, static_cast(replicated_state_tree.end_index())}, - replicated_state_tree.get_root(), - term_of_next_version}; + replicated_state_tree.get_root()}; } bool verify_root_signatures(ccf::kv::Version version) override @@ -891,16 +882,13 @@ namespace ccf // term std::lock_guard guard(state_lock); term_of_last_version = t; - term_of_next_version = t; } - void rollback( - const ccf::TxID& tx_id, ccf::kv::Term term_of_next_version_) override + void rollback(const ccf::TxID& tx_id) override { std::lock_guard guard(state_lock); LOG_TRACE_FMT("Rollback to {}.{}", tx_id.view, tx_id.seqno); term_of_last_version = tx_id.view; - term_of_next_version = term_of_next_version_; replicated_state_tree.retract(tx_id.seqno); log_hash(replicated_state_tree.get_root(), ROLLBACK); } @@ -1003,20 +991,10 @@ namespace ccf replicated_state_tree.append(rh); } - void append_entry( - const ccf::crypto::Sha256Hash& digest, - std::optional expected_term_of_next_version = - std::nullopt) override + void append_entry(const ccf::crypto::Sha256Hash& digest) override { log_hash(digest, APPEND); std::lock_guard guard(state_lock); - if (expected_term_of_next_version.has_value()) - { - if (expected_term_of_next_version.value() != term_of_next_version) - { - return; - } - } replicated_state_tree.append(digest); } diff --git a/src/node/node_state.h b/src/node/node_state.h index eb4f64e7a911..38eff1274625 100644 --- a/src/node/node_state.h +++ b/src/node/node_state.h @@ -2432,7 +2432,7 @@ namespace ccf // version can advance before its Merkle history is updated during commit, // so these must be captured together from the history. auto* h = dynamic_cast(history.get()); - const auto& [txid, root, _] = h->get_replicated_state_txid_and_root(); + const auto& [txid, root] = h->get_replicated_state_txid_and_root(); recovery_v = txid.seqno; recovery_root = root; diff --git a/src/node/rpc/frontend.h b/src/node/rpc/frontend.h index ec65f95ab601..9bc5506b74cf 100644 --- a/src/node/rpc/frontend.h +++ b/src/node/rpc/frontend.h @@ -1113,9 +1113,9 @@ namespace ccf // should only ever be used for the proposal creation endpoint and // nothing else. Many bad things could happen otherwise (e.g. breaking // session consistency). - const auto& [txid, root, term_of_next_version] = + const auto& [txid, root] = current_history->get_replicated_state_txid_and_root(); - tx.set_read_txid(txid, term_of_next_version); + tx.set_read_txid(txid); tx.set_root_at_read_version(root); } } From f03647e95d72078e711f8e403641fcac6f0f4cd5 Mon Sep 17 00:00:00 2001 From: Eddy Ashton Date: Mon, 24 Aug 2026 13:14:36 +0000 Subject: [PATCH 02/12] VersionResolver becomes TxIDResolver, remove idea of TxID having a separate "commit_view" --- src/kv/apply_changes.h | 25 ++++++------- src/kv/committable_tx.h | 79 ++++++++++++++++++++--------------------- src/kv/kv_types.h | 3 -- src/kv/store.h | 33 ++--------------- src/kv/tx.cpp | 4 +-- src/kv/tx_pimpl.h | 1 - src/node/rpc/frontend.h | 3 +- src/node/snapshotter.h | 2 +- 8 files changed, 54 insertions(+), 96 deletions(-) diff --git a/src/kv/apply_changes.h b/src/kv/apply_changes.h index dc64c5dd3dc1..4b1d0ff6cc38 100644 --- a/src/kv/apply_changes.h +++ b/src/kv/apply_changes.h @@ -16,17 +16,14 @@ namespace ccf::kv using MapCollection = std::map>; // Atomically checks for conflicts then applies the writes in the given change - // sets to their underlying Maps. Calls f() at most once, iff the writes are - // applied, to retrieve a unique Version for the write set and return the max - // version which can have a conflict with the transaction. + // sets to their underlying Maps. Calls tx_id_resolver() at most once, iff the + // writes are applied, to retrieve a unique TxID for the write set. - using VersionLastNewMap = Version; - using VersionResolver = std::function( - bool tx_contains_new_map)>; + using TxIDResolver = std::function; - static inline std::optional apply_changes( + static inline std::optional apply_changes( OrderedChanges& changes, - VersionResolver version_resolver_fn, + TxIDResolver tx_id_resolver, ccf::kv::ConsensusHookPtrs& hooks, const MapCollection& new_maps, const std::optional& new_maps_conflict_version, @@ -37,7 +34,7 @@ namespace ccf::kv // and possibly committed, and then all maps with pending writes are // unlocked. This is to prevent transactions from being committed in an // interleaved fashion. - Version version = NoVersion; + ccf::TxID tx_id; bool has_writes = false; std::map> views; @@ -117,9 +114,7 @@ namespace ccf::kv if (ok && has_writes) { // Get the version number to be used for this commit. - ccf::kv::Version version_last_new_map = 0; - std::tie(version, version_last_new_map) = - version_resolver_fn(!new_maps.empty()); + tx_id = tx_id_resolver(); // Transfer ownership of these new maps to their target stores, iff we // have writes to them @@ -128,13 +123,13 @@ namespace ccf::kv const auto it = views.find(map_name); if (it != views.end() && it->second->has_writes()) { - map_ptr->get_store()->add_dynamic_map(version, map_ptr); + map_ptr->get_store()->add_dynamic_map(tx_id.seqno, map_ptr); } } for (auto& [view_name, view_ptr] : views) { - view_ptr->commit(version, track_deletes_on_missing_keys); + view_ptr->commit(tx_id.seqno, track_deletes_on_missing_keys); } // Collect ConsensusHooks @@ -158,6 +153,6 @@ namespace ccf::kv return std::nullopt; } - return version; + return tx_id; } } diff --git a/src/kv/committable_tx.h b/src/kv/committable_tx.h index 8c49d469069d..a7e99a4f6e48 100644 --- a/src/kv/committable_tx.h +++ b/src/kv/committable_tx.h @@ -28,10 +28,15 @@ namespace ccf::kv }; protected: + // Indicates that this transaction's changes have been applied to the local + // KV. Monotonic and gating, to form a crude linear type - some functions + // are available only pre-commit, some only post-commit. bool committed = false; - bool success = false; - Version version = NoVersion; + // Populated only after commit() has been called, and only if the + // transaction was successful. This is the version at which the transaction + // was applied to the local KV. + std::optional applied_txid = std::nullopt; TxFlags flags = 0; SerialisedEntryFlags entry_flags = 0; @@ -47,7 +52,7 @@ namespace ccf::kv throw std::logic_error("Transaction not yet committed"); } - if (!success) + if (!applied_txid.has_value()) { throw std::logic_error("Transaction aborted"); } @@ -74,7 +79,7 @@ namespace ccf::kv throw KvSerialiserException("No encryptor set"); } - commit_evidence = e->get_commit_evidence({pimpl->commit_view, version}); + commit_evidence = e->get_commit_evidence(*applied_txid); LOG_TRACE_FMT("Commit evidence: {}", commit_evidence); ccf::crypto::Sha256Hash tx_commit_evidence_digest(commit_evidence); commit_evidence_digest = tx_commit_evidence_digest; @@ -87,7 +92,7 @@ namespace ccf::kv RawKvStoreSerialiser replicated_serialiser( e, - {pimpl->commit_view, version}, + *applied_txid, entry_type, entry_flags, tx_commit_evidence_digest, @@ -134,8 +139,6 @@ namespace ccf::kv */ CommitResult commit( const ccf::ClaimsDigest& claims = ccf::empty_claims(), - std::function(bool has_new_map)> - version_resolver = nullptr, WriteSetObserver write_set_observer = nullptr) { if (committed) @@ -146,7 +149,7 @@ namespace ccf::kv if (all_changes.empty()) { committed = true; - success = true; + applied_txid = pimpl->read_txid; return CommitResult::SUCCESS; } @@ -163,13 +166,11 @@ namespace ccf::kv std::optional new_maps_conflict_version = std::nullopt; bool track_deletes_on_missing_keys = false; - auto c = apply_changes( + auto txid_resolver = [&]() { return this->pimpl->store->next_txid(); }; + + applied_txid = apply_changes( all_changes, - version_resolver == nullptr ? - [&](bool has_new_map) { - return pimpl->store->next_version(has_new_map); - } : - version_resolver, + txid_resolver, hooks, pimpl->created_maps, new_maps_conflict_version, @@ -180,9 +181,7 @@ namespace ccf::kv this->pimpl->store->unlock_map_set(); } - success = c.has_value(); - - if (!success) + if (!applied_txid.has_value()) { // This Tx is now in a dead state. Caller should create a new Tx and try // again. @@ -191,14 +190,13 @@ namespace ccf::kv } committed = true; - version = c.value(); if (tx_flag_enabled(TxFlag::LEDGER_CHUNK_AT_NEXT_SIGNATURE)) { auto chunker = pimpl->store->get_chunker(); if (chunker) { - chunker->force_end_of_chunk(version); + chunker->force_end_of_chunk(applied_txid->seqno); } } @@ -209,7 +207,7 @@ namespace ccf::kv unset_tx_flag(TxFlag::SNAPSHOT_AT_NEXT_SIGNATURE); } - if (version == NoVersion) + if (applied_txid->seqno == NoVersion) { // Read-only transaction return CommitResult::SUCCESS; @@ -238,7 +236,7 @@ namespace ccf::kv auto claims_ = claims; return pimpl->store->commit( - {pimpl->commit_view, version}, + *applied_txid, std::make_unique( std::move(data), std::move(claims_), @@ -273,12 +271,12 @@ namespace ccf::kv throw std::logic_error("Transaction not yet committed"); } - if (!success) + if (!applied_txid.has_value()) { throw std::logic_error("Transaction aborted"); } - return version; + return applied_txid->seqno; } /** Get term in which this transaction was committed. @@ -295,12 +293,12 @@ namespace ccf::kv throw std::logic_error("Transaction not yet committed"); } - if (!success) + if (!applied_txid.has_value()) { throw std::logic_error("Transaction aborted"); } - return pimpl->commit_view; + return applied_txid->view; } [[nodiscard]] std::optional get_txid() const @@ -318,14 +316,15 @@ namespace ccf::kv // A committed tx is read-only (i.e. no write to any map) if it was not // assigned a version when it was committed - if (version == NoVersion) + // TODO: This is no longer true. Just return applied_txid, if possible + if (!applied_txid.has_value() || applied_txid->seqno == NoVersion) { // Read-only transaction return pimpl->read_txid; } // Write transaction - return TxID(pimpl->commit_view, version); + return applied_txid; } void set_read_txid(const TxID& tx_id) @@ -369,19 +368,19 @@ namespace ccf::kv { private: Version rollback_count = 0; + const TxID reserved_txid; public: ReservedTx( AbstractStore* _store, Term read_term, - const TxID& reserved_tx_id, + const TxID& reserved_txid_, Version rollback_count_) : CommittableTx(_store), - rollback_count(rollback_count_) + rollback_count(rollback_count_), + reserved_txid(reserved_txid_) { - version = reserved_tx_id.seqno; - pimpl->commit_view = reserved_tx_id.view; - pimpl->read_txid = TxID(read_term, reserved_tx_id.seqno - 1); + pimpl->read_txid = TxID(read_term, reserved_txid.seqno - 1); } // Used by frontend to commit reserved transactions @@ -399,15 +398,15 @@ namespace ccf::kv std::vector hooks; bool track_deletes_on_missing_keys = false; - auto c = apply_changes( + applied_txid = apply_changes( all_changes, - [this](bool) { return std::make_tuple(version, version - 1); }, + [this]() { return reserved_txid; }, hooks, pimpl->created_maps, - version, + reserved_txid.seqno, track_deletes_on_missing_keys, rollback_count); - success = c.has_value(); + const auto success = applied_txid.has_value(); if (!success) { @@ -427,18 +426,18 @@ namespace ccf::kv // This is a signature and, if the ledger chunking or snapshot flags are // enabled, we want the host to create a chunk when it sees this entry. // version_lock held by Store::commit - if (pimpl->store->should_create_ledger_chunk_unsafe(version)) + if (pimpl->store->should_create_ledger_chunk_unsafe(applied_txid->seqno)) { entry_flags |= EntryFlags::FORCE_LEDGER_CHUNK_AFTER; LOG_DEBUG_FMT( "Ending ledger chunk with signature at {}.{}", - pimpl->commit_view, - version); + applied_txid->view, + applied_txid->seqno); auto chunker = pimpl->store->get_chunker(); if (chunker) { - chunker->produced_chunk_at(version); + chunker->produced_chunk_at(applied_txid->seqno); } } diff --git a/src/kv/kv_types.h b/src/kv/kv_types.h index ea0663308f46..e48d239702f8 100644 --- a/src/kv/kv_types.h +++ b/src/kv/kv_types.h @@ -696,13 +696,10 @@ namespace ccf::kv virtual void lock_map_set() = 0; virtual void unlock_map_set() = 0; - virtual Version next_version() = 0; - virtual std::tuple next_version(bool commit_new_map) = 0; virtual ccf::TxID next_txid() = 0; virtual Version current_version() = 0; virtual ccf::TxID current_txid() = 0; - virtual std::pair current_txid_and_commit_term() = 0; virtual Version compacted_version() = 0; virtual Term commit_view() = 0; diff --git a/src/kv/store.h b/src/kv/store.h index 8155b9dfba7a..4a63afda46ea 100644 --- a/src/kv/store.h +++ b/src/kv/store.h @@ -43,7 +43,6 @@ namespace ccf::kv ccf::pal::Mutex version_lock; std::atomic version = 0; - Version last_new_map = ccf::kv::NoVersion; std::atomic compacted = 0; // Calls to Store::commit are made atomic by taking this lock. @@ -78,7 +77,6 @@ namespace ccf::kv pending_txs.clear(); version = 0; - last_new_map = ccf::kv::NoVersion; compacted = 0; term_of_next_version = 0; term_of_last_version = 0; @@ -139,7 +137,7 @@ namespace ccf::kv auto c = apply_changes( changes, - [v](bool) { return std::make_tuple(v, v - 1); }, + [term, v]() { return ccf::TxID(term, v); }, hooks, new_maps, std::nullopt, @@ -530,7 +528,7 @@ namespace ccf::kv bool track_deletes_on_missing_keys = false; auto r = apply_changes( changes, - [](bool) { return std::make_tuple(NoVersion, NoVersion); }, + [term, v]() { return ccf::TxID(term, v); }, hooks, new_maps, std::nullopt, @@ -913,13 +911,6 @@ namespace ccf::kv return current_txid_unsafe(); } - std::pair current_txid_and_commit_term() override - { - // Must lock in case the version or commit term is being incremented. - std::lock_guard vguard(version_lock); - return {current_txid_unsafe(), term_of_next_version}; - } - Version compacted_version() override { return compacted; @@ -1148,26 +1139,6 @@ namespace ccf::kv return rollback_count == count; } - std::tuple next_version(bool commit_new_map) override - { - std::lock_guard vguard(version_lock); - Version v = next_version_unsafe(); - - auto previous_last_new_map = last_new_map; - if (commit_new_map) - { - last_new_map = v; - } - - return std::make_tuple(v, previous_last_new_map); - } - - Version next_version() override - { - std::lock_guard vguard(version_lock); - return next_version_unsafe(); - } - TxID next_txid() override { std::lock_guard vguard(version_lock); diff --git a/src/kv/tx.cpp b/src/kv/tx.cpp index 3432868cc2dc..9ecbba0b453e 100644 --- a/src/kv/tx.cpp +++ b/src/kv/tx.cpp @@ -58,9 +58,7 @@ namespace ccf::kv // rather than earlier, at Tx construction. This is to minimise the // window during which concurrent transactions can write to the same map // and cause this transaction to conflict on commit. - auto p = pimpl->store->current_txid_and_commit_term(); - read_txid = p.first; - pimpl->commit_view = p.second; + read_txid = pimpl->store->current_txid(); } auto abstract_map = pimpl->store->get_map(read_txid->seqno, map_name); diff --git a/src/kv/tx_pimpl.h b/src/kv/tx_pimpl.h index 1f14e91297e4..e42f866be060 100644 --- a/src/kv/tx_pimpl.h +++ b/src/kv/tx_pimpl.h @@ -21,7 +21,6 @@ namespace ccf::kv // Note: read_txid version is set to NoVersion for the first transaction in // the service, before anything has been applied to the KV. std::optional read_txid = std::nullopt; - ccf::View commit_view = ccf::VIEW_UNKNOWN; std::map> created_maps; }; diff --git a/src/node/rpc/frontend.h b/src/node/rpc/frontend.h index 9bc5506b74cf..953c753e74cf 100644 --- a/src/node/rpc/frontend.h +++ b/src/node/rpc/frontend.h @@ -900,8 +900,7 @@ namespace ccf }; } - ccf::kv::CommitResult result = - tx.commit(ctx->claims, nullptr, ws_observer); + ccf::kv::CommitResult result = tx.commit(ctx->claims, ws_observer); switch (result) { diff --git a/src/node/snapshotter.h b/src/node/snapshotter.h index 5b437c2b1901..03a0d93703bc 100644 --- a/src/node/snapshotter.h +++ b/src/node/snapshotter.h @@ -277,7 +277,7 @@ namespace ccf commit_evidence = commit_evidence_; }; - auto rc = tx.commit(cd, nullptr, capture_ws_digest_and_commit_evidence); + auto rc = tx.commit(cd, capture_ws_digest_and_commit_evidence); if (rc != ccf::kv::CommitResult::SUCCESS) { LOG_FAIL_FMT( From 98e3461ded6fd659ff92fad38a89db067641c3e6 Mon Sep 17 00:00:00 2001 From: Eddy Ashton Date: Mon, 24 Aug 2026 14:02:14 +0000 Subject: [PATCH 03/12] API shape change - BatchVector holds TxID --- src/consensus/aft/raft.h | 9 ++- src/consensus/aft/test/committable_suffix.cpp | 56 ++++++++++++------ src/consensus/aft/test/driver.h | 4 +- src/consensus/aft/test/main.cpp | 59 ++++++++++--------- src/indexing/test/common.h | 2 +- src/kv/kv_types.h | 2 +- src/kv/store.h | 30 +++------- src/kv/test/stub_consensus.h | 6 +- src/node/test/historical_queries.cpp | 12 ++-- src/node/test/history.cpp | 8 +-- tests/schema.py | 17 +++--- 11 files changed, 112 insertions(+), 93 deletions(-) diff --git a/src/consensus/aft/raft.h b/src/consensus/aft/raft.h index 6842946a80e4..466f61a34fb6 100644 --- a/src/consensus/aft/raft.h +++ b/src/consensus/aft/raft.h @@ -621,6 +621,7 @@ namespace aft return details; } + // TODO TODO: Mid-refactor. Term argument should be removed bool replicate(const ccf::kv::BatchVector& entries, Term term) override { std::lock_guard guard(state->lock); @@ -654,12 +655,18 @@ namespace aft RAFT_DEBUG_FMT("Replicating {} entries", entries.size()); - for (const auto& [index, data, is_globally_committable, hooks] : entries) + for (const auto& [tx_id, data, is_globally_committable, hooks] : entries) { + // TODO TODO: This is a temporary hack to allow the new TxID type to be + // used in the raft code. Once the raft code is fully migrated to use + // TxID, this can be removed. + const auto index = tx_id.seqno; bool globally_committable = is_globally_committable; if (index != state->last_idx + 1) { + LOG_INFO_FMT( + "!!!! Not contiguous ({} != {} + 1)", index, state->last_idx); return false; } diff --git a/src/consensus/aft/test/committable_suffix.cpp b/src/consensus/aft/test/committable_suffix.cpp index 3571bedea2e5..35b60e249b5b 100644 --- a/src/consensus/aft/test/committable_suffix.cpp +++ b/src/consensus/aft/test/committable_suffix.cpp @@ -214,7 +214,8 @@ DOCTEST_TEST_CASE("Retention of dead leader's commit") DOCTEST_INFO("Entry at 1.1 is received by all nodes"); { auto entry = make_ledger_entry(1, 1); - rA.replicate(ccf::kv::BatchVector{{1, entry, true, hooks}}, 1); + rA.replicate( + ccf::kv::BatchVector{{ccf::TxID{1, 1}, entry, true, hooks}}, 1); DOCTEST_REQUIRE(rA.get_last_idx() == 1); DOCTEST_REQUIRE(rA.get_committed_seqno() == 0); // Size limit was reached, so periodic is not needed @@ -244,21 +245,24 @@ DOCTEST_TEST_CASE("Retention of dead leader's commit") "committed"); { auto entry = make_ledger_entry(1, 2); - rA.replicate(ccf::kv::BatchVector{{2, entry, true, hooks}}, 1); + rA.replicate( + ccf::kv::BatchVector{{ccf::TxID{1, 2}, entry, true, hooks}}, 1); DOCTEST_REQUIRE(rA.get_last_idx() == 2); DOCTEST_REQUIRE(rA.get_committed_seqno() == 1); // Size limit was reached, so periodic is not needed // rA.periodic(request_timeout); entry = make_ledger_entry(1, 3); - rA.replicate(ccf::kv::BatchVector{{3, entry, true, hooks}}, 1); + rA.replicate( + ccf::kv::BatchVector{{ccf::TxID{1, 3}, entry, true, hooks}}, 1); DOCTEST_REQUIRE(rA.get_last_idx() == 3); DOCTEST_REQUIRE(rA.get_committed_seqno() == 1); // Size limit was reached, so periodic is not needed // rA.periodic(request_timeout); entry = make_ledger_entry(1, 4); - rA.replicate(ccf::kv::BatchVector{{4, entry, true, hooks}}, 1); + rA.replicate( + ccf::kv::BatchVector{{ccf::TxID{1, 4}, entry, true, hooks}}, 1); DOCTEST_REQUIRE(rA.get_last_idx() == 4); DOCTEST_REQUIRE(rA.get_committed_seqno() == 1); // Size limit was reached, so periodic is not needed @@ -292,7 +296,8 @@ DOCTEST_TEST_CASE("Retention of dead leader's commit") "committed"); { auto entry = make_ledger_entry(1, 5); - rA.replicate(ccf::kv::BatchVector{{5, entry, true, hooks}}, 1); + rA.replicate( + ccf::kv::BatchVector{{ccf::TxID{1, 5}, entry, true, hooks}}, 1); DOCTEST_REQUIRE(rA.get_last_idx() == 5); // Size limit was reached, so periodic is not needed // rB.periodic(request_timeout); @@ -367,11 +372,13 @@ DOCTEST_TEST_CASE("Retention of dead leader's commit") DOCTEST_INFO("Node B writes some entries, though they are lost"); { auto entry = make_ledger_entry(2, 6); - rB.replicate(ccf::kv::BatchVector{{6, entry, true, hooks}}, 2); + rB.replicate( + ccf::kv::BatchVector{{ccf::TxID{2, 6}, entry, true, hooks}}, 2); DOCTEST_REQUIRE(rB.get_last_idx() == 6); entry = make_ledger_entry(2, 7); - rB.replicate(ccf::kv::BatchVector{{7, entry, true, hooks}}, 2); + rB.replicate( + ccf::kv::BatchVector{{ccf::TxID{2, 7}, entry, true, hooks}}, 2); DOCTEST_REQUIRE(rB.get_last_idx() == 7); // Size limit was reached, so periodic is not needed @@ -426,15 +433,18 @@ DOCTEST_TEST_CASE("Retention of dead leader's commit") DOCTEST_REQUIRE("Node C produces 3.5, 3.6, and 3.7"); { auto entry = make_ledger_entry(3, 5); - rC.replicate(ccf::kv::BatchVector{{5, entry, true, hooks}}, 3); + rC.replicate( + ccf::kv::BatchVector{{ccf::TxID{3, 5}, entry, true, hooks}}, 3); DOCTEST_REQUIRE(rC.get_last_idx() == 5); entry = make_ledger_entry(3, 6); - rC.replicate(ccf::kv::BatchVector{{6, entry, true, hooks}}, 3); + rC.replicate( + ccf::kv::BatchVector{{ccf::TxID{3, 6}, entry, true, hooks}}, 3); DOCTEST_REQUIRE(rC.get_last_idx() == 6); entry = make_ledger_entry(3, 7); - rC.replicate(ccf::kv::BatchVector{{7, entry, true, hooks}}, 3); + rC.replicate( + ccf::kv::BatchVector{{ccf::TxID{3, 7}, entry, true, hooks}}, 3); DOCTEST_REQUIRE(rC.get_last_idx() == 7); // The early AppendEntries that describe this are lost @@ -639,7 +649,9 @@ DOCTEST_TEST_CASE_TEMPLATE("Multi-term divergence", T, WorstCase, RandomCase) { auto entry = make_ledger_entry(primary.get_view(), idx); primary.replicate( - ccf::kv::BatchVector{{idx, entry, true, hooks}}, primary.get_view()); + ccf::kv::BatchVector{ + {ccf::TxID{primary.get_view(), idx}, entry, true, hooks}}, + primary.get_view()); } // All related AppendEntries are lost @@ -668,9 +680,11 @@ DOCTEST_TEST_CASE_TEMPLATE("Multi-term divergence", T, WorstCase, RandomCase) // be committed auto entry = make_ledger_entry(1, 1); - rA.replicate(ccf::kv::BatchVector{{1, entry, true, hooks}}, 1); + rA.replicate( + ccf::kv::BatchVector{{ccf::TxID{1, 1}, entry, true, hooks}}, 1); entry = make_ledger_entry(1, 2); - rA.replicate(ccf::kv::BatchVector{{2, entry, true, hooks}}, 1); + rA.replicate( + ccf::kv::BatchVector{{ccf::TxID{1, 2}, entry, true, hooks}}, 1); DOCTEST_REQUIRE(rA.get_last_idx() == 2); DOCTEST_REQUIRE(rA.get_committed_seqno() == 0); // Size limit was reached, so periodic is not needed @@ -703,18 +717,22 @@ DOCTEST_TEST_CASE_TEMPLATE("Multi-term divergence", T, WorstCase, RandomCase) // Node A produces 2 additional entries that A and B have, and 2 additional // entries that are only present on A entry = make_ledger_entry(1, 3); - rA.replicate(ccf::kv::BatchVector{{3, entry, true, hooks}}, 1); + rA.replicate( + ccf::kv::BatchVector{{ccf::TxID{1, 3}, entry, true, hooks}}, 1); entry = make_ledger_entry(1, 4); - rA.replicate(ccf::kv::BatchVector{{4, entry, true, hooks}}, 1); + rA.replicate( + ccf::kv::BatchVector{{ccf::TxID{1, 4}, entry, true, hooks}}, 1); keep_messages_for(node_idB, channelsA->messages); DOCTEST_REQUIRE(2 == dispatch_all(nodes, node_idA)); entry = make_ledger_entry(1, 5); - rA.replicate(ccf::kv::BatchVector{{5, entry, true, hooks}}, 1); + rA.replicate( + ccf::kv::BatchVector{{ccf::TxID{1, 5}, entry, true, hooks}}, 1); entry = make_ledger_entry(1, 6); - rA.replicate(ccf::kv::BatchVector{{6, entry, true, hooks}}, 1); + rA.replicate( + ccf::kv::BatchVector{{ccf::TxID{1, 6}, entry, true, hooks}}, 1); channelsA->messages.clear(); channelsB->messages.clear(); @@ -991,7 +1009,9 @@ DOCTEST_TEST_CASE_TEMPLATE("Multi-term divergence", T, WorstCase, RandomCase) const auto seqno = rPrimary.get_last_idx() + 1; auto final_entry = make_ledger_entry(view, seqno); rPrimary.replicate( - ccf::kv::BatchVector{{seqno, final_entry, true, hooks}}, view); + ccf::kv::BatchVector{ + {ccf::TxID{view, seqno}, final_entry, true, hooks}}, + view); rPrimary.periodic(request_timeout); keep_earliest_append_entries_for_each_target(channelsPrimary->messages); diff --git a/src/consensus/aft/test/driver.h b/src/consensus/aft/test/driver.h index aa2693bef970..1a993eed8afd 100644 --- a/src/consensus/aft/test/driver.h +++ b/src/consensus/aft/test/driver.h @@ -188,7 +188,9 @@ class RaftDriver auto s = nlohmann::json(aft::ReplicatedData{type, data}).dump(); auto d = std::make_shared>(s.begin(), s.end()); - raft->replicate(ccf::kv::BatchVector{{idx, d, committable, hooks}}, term); + raft->replicate( + ccf::kv::BatchVector{{ccf::TxID{term, idx}, d, committable, hooks}}, + term); } void add_node(ccf::NodeId node_id) diff --git a/src/consensus/aft/test/main.cpp b/src/consensus/aft/test/main.cpp index 317ec0668bac..c1eda8233095 100644 --- a/src/consensus/aft/test/main.cpp +++ b/src/consensus/aft/test/main.cpp @@ -77,7 +77,8 @@ DOCTEST_TEST_CASE("Single node commit" * doctest::test_suite("single")) entry->push_back(2); entry->push_back(3); - r0.replicate(ccf::kv::BatchVector{{i, entry, true, hooks}}, 1); + r0.replicate( + ccf::kv::BatchVector{{ccf::TxID{1, i}, entry, true, hooks}}, 1); DOCTEST_REQUIRE(r0.get_last_idx() == i); DOCTEST_REQUIRE(r0.get_committed_seqno() == i); } @@ -428,12 +429,12 @@ DOCTEST_TEST_CASE( DOCTEST_INFO("Try to replicate on a follower, and fail"); std::vector entry = {1, 2, 3}; auto data = std::make_shared>(entry); - DOCTEST_REQUIRE_FALSE( - r1.replicate(ccf::kv::BatchVector{{1, data, true, hooks}}, 1)); + DOCTEST_REQUIRE_FALSE(r1.replicate( + ccf::kv::BatchVector{{ccf::TxID{1, 1}, data, true, hooks}}, 1)); DOCTEST_INFO("Tell the leader to replicate a message"); - DOCTEST_REQUIRE( - r0.replicate(ccf::kv::BatchVector{{1, data, true, hooks}}, 1)); + DOCTEST_REQUIRE(r0.replicate( + ccf::kv::BatchVector{{ccf::TxID{1, 1}, data, true, hooks}}, 1)); DOCTEST_REQUIRE(r0.ledger->ledger.size() == 1); // The test ledger adds its own header. Confirm that the expected data is @@ -548,8 +549,8 @@ DOCTEST_TEST_CASE("Multiple nodes late join" * doctest::test_suite("multiple")) std::vector first_entry = {1, 2, 3}; auto data = std::make_shared>(first_entry); - DOCTEST_REQUIRE( - r0.replicate(ccf::kv::BatchVector{{1, data, true, hooks}}, 1)); + DOCTEST_REQUIRE(r0.replicate( + ccf::kv::BatchVector{{ccf::TxID{1, 1}, data, true, hooks}}, 1)); r0.periodic(request_timeout); DOCTEST_REQUIRE( @@ -662,10 +663,10 @@ DOCTEST_TEST_CASE("Recv append entries logic" * doctest::test_suite("multiple")) std::vector second_entry = {2, 2, 2}; auto data_2 = std::make_shared>(second_entry); - DOCTEST_REQUIRE( - r0.replicate(ccf::kv::BatchVector{{1, data_1, true, hooks}}, 1)); - DOCTEST_REQUIRE( - r0.replicate(ccf::kv::BatchVector{{2, data_2, true, hooks}}, 1)); + DOCTEST_REQUIRE(r0.replicate( + ccf::kv::BatchVector{{ccf::TxID{1, 1}, data_1, true, hooks}}, 1)); + DOCTEST_REQUIRE(r0.replicate( + ccf::kv::BatchVector{{ccf::TxID{1, 2}, data_2, true, hooks}}, 1)); DOCTEST_REQUIRE(r0.ledger->ledger.size() == 2); r0.periodic(request_timeout); DOCTEST_REQUIRE(r0c->messages.size() == 1); @@ -686,8 +687,8 @@ DOCTEST_TEST_CASE("Recv append entries logic" * doctest::test_suite("multiple")) { std::vector third_entry = {3, 3, 3}; auto data = std::make_shared>(third_entry); - DOCTEST_REQUIRE( - r0.replicate(ccf::kv::BatchVector{{3, data, true, hooks}}, 1)); + DOCTEST_REQUIRE(r0.replicate( + ccf::kv::BatchVector{{ccf::TxID{1, 3}, data, true, hooks}}, 1)); DOCTEST_REQUIRE(r0.ledger->ledger.size() == 3); // Simulate that the append entries was not deserialised successfully @@ -719,8 +720,8 @@ DOCTEST_TEST_CASE("Recv append entries logic" * doctest::test_suite("multiple")) { std::vector fourth_entry = {4, 4, 4}; auto data = std::make_shared>(fourth_entry); - DOCTEST_REQUIRE( - r0.replicate(ccf::kv::BatchVector{{4, data, true, hooks}}, 1)); + DOCTEST_REQUIRE(r0.replicate( + ccf::kv::BatchVector{{ccf::TxID{1, 4}, data, true, hooks}}, 1)); DOCTEST_REQUIRE(r0.ledger->ledger.size() == 4); r0.periodic(request_timeout); DOCTEST_REQUIRE(r0c->messages.size() == 1); @@ -733,8 +734,8 @@ DOCTEST_TEST_CASE("Recv append entries logic" * doctest::test_suite("multiple")) { std::vector fifth_entry = {5, 5, 5}; auto data = std::make_shared>(fifth_entry); - DOCTEST_REQUIRE( - r0.replicate(ccf::kv::BatchVector{{5, data, true, hooks}}, 1)); + DOCTEST_REQUIRE(r0.replicate( + ccf::kv::BatchVector{{ccf::TxID{1, 5}, data, true, hooks}}, 1)); DOCTEST_REQUIRE(r0.ledger->ledger.size() == 5); r0.periodic(request_timeout); DOCTEST_REQUIRE(r0c->messages.size() == 1); @@ -763,8 +764,8 @@ DOCTEST_TEST_CASE("Recv append entries logic" * doctest::test_suite("multiple")) { std::vector entry_6 = {6, 6, 6}; auto data = std::make_shared>(entry_6); - DOCTEST_REQUIRE( - r0.replicate(ccf::kv::BatchVector{{6, data, true, hooks}}, 1)); + DOCTEST_REQUIRE(r0.replicate( + ccf::kv::BatchVector{{ccf::TxID{1, 6}, data, true, hooks}}, 1)); DOCTEST_REQUIRE(r0.ledger->ledger.size() == 6); } const auto last_correct_version = r0.ledger->ledger.size(); @@ -773,8 +774,8 @@ DOCTEST_TEST_CASE("Recv append entries logic" * doctest::test_suite("multiple")) { std::vector entry_7 = {7, 7, 7}; auto data = std::make_shared>(entry_7); - DOCTEST_REQUIRE( - r0.replicate(ccf::kv::BatchVector{{7, data, true, hooks}}, 1)); + DOCTEST_REQUIRE(r0.replicate( + ccf::kv::BatchVector{{ccf::TxID{1, 7}, data, true, hooks}}, 1)); DOCTEST_REQUIRE(r0.ledger->ledger.size() == 7); dead_branch = r0.ledger->ledger.back(); } @@ -795,8 +796,8 @@ DOCTEST_TEST_CASE("Recv append entries logic" * doctest::test_suite("multiple")) { std::vector entry_7b = {7, 7, 'b'}; auto data = std::make_shared>(entry_7b); - DOCTEST_REQUIRE( - r0.replicate(ccf::kv::BatchVector{{7, data, true, hooks}}, 4)); + DOCTEST_REQUIRE(r0.replicate( + ccf::kv::BatchVector{{ccf::TxID{4, 7}, data, true, hooks}}, 4)); DOCTEST_REQUIRE(r0.ledger->ledger.size() == 7); live_branch = r0.ledger->ledger.back(); } @@ -804,8 +805,8 @@ DOCTEST_TEST_CASE("Recv append entries logic" * doctest::test_suite("multiple")) { std::vector entry_8 = {8, 8, 8}; auto data = std::make_shared>(entry_8); - DOCTEST_REQUIRE( - r0.replicate(ccf::kv::BatchVector{{8, data, true, hooks}}, 4)); + DOCTEST_REQUIRE(r0.replicate( + ccf::kv::BatchVector{{ccf::TxID{4, 8}, data, true, hooks}}, 4)); DOCTEST_REQUIRE(r0.ledger->ledger.size() == 8); DOCTEST_REQUIRE(r0.ledger->ledger.size() > last_correct_version); } @@ -936,8 +937,8 @@ DOCTEST_TEST_CASE("Exceed append entries limit") for (size_t i = 1; i <= static_cast(num_big_entries); ++i) { - DOCTEST_REQUIRE( - r0.replicate(ccf::kv::BatchVector{{i, data, true, hooks}}, 1)); + DOCTEST_REQUIRE(r0.replicate( + ccf::kv::BatchVector{{ccf::TxID{1, i}, data, true, hooks}}, 1)); const auto received_ae = dispatch_all_and_DOCTEST_CHECK( nodes, node_id0, r0c->messages, [](const auto& msg) { @@ -955,8 +956,8 @@ DOCTEST_TEST_CASE("Exceed append entries limit") i <= static_cast(individual_entries); ++i) { - DOCTEST_REQUIRE( - r0.replicate(ccf::kv::BatchVector{{i, smaller_data, true, hooks}}, 1)); + DOCTEST_REQUIRE(r0.replicate( + ccf::kv::BatchVector{{ccf::TxID{1, i}, smaller_data, true, hooks}}, 1)); dispatch_all(nodes, node_id0, r0c->messages); } diff --git a/src/indexing/test/common.h b/src/indexing/test/common.h index 9c65339b4458..e77f6515a360 100644 --- a/src/indexing/test/common.h +++ b/src/indexing/test/common.h @@ -92,7 +92,7 @@ class AllCommittableWrapper : public TConsensus // Rather than building a history that produces real signatures, we just // overwrite the entries here to say that everything is committable ccf::kv::BatchVector entries(entries_); - for (auto& [seqno, data, committable, hooks] : entries) + for (auto& [tx_id, data, committable, hooks] : entries) { committable = true; } diff --git a/src/kv/kv_types.h b/src/kv/kv_types.h index e48d239702f8..0d6b0a98c4f9 100644 --- a/src/kv/kv_types.h +++ b/src/kv/kv_types.h @@ -206,7 +206,7 @@ namespace ccf::kv }; using BatchVector = std::vector>, bool, std::shared_ptr>>; diff --git a/src/kv/store.h b/src/kv/store.h index 4a63afda46ea..180474008bf6 100644 --- a/src/kv/store.h +++ b/src/kv/store.h @@ -64,7 +64,9 @@ namespace ccf::kv Version rollback_count = 0; - std::unordered_map, bool>> + std::unordered_map< + Version, + std::tuple, bool>> pending_txs; public: @@ -944,25 +946,13 @@ namespace ccf::kv Version previous_last_replicated = 0; Version next_last_replicated = 0; Version previous_rollback_count = 0; - ccf::View replication_view = 0; + ccf::View replication_view = 0; // TODO: Remove - std::vector, bool>> - contiguous_pending_txs; + std::vector contiguous_pending_txs; auto h = get_history(); { std::lock_guard vguard(version_lock); - if (txid.view != term_of_next_version && get_consensus()->is_primary()) - { - // This can happen when a transaction started before a view change, - // but tries to commit after the view change is complete. - LOG_DEBUG_FMT( - "Want to commit for term {} but term is {}", - txid.view, - term_of_next_version); - - return CommitResult::FAIL_NO_REPLICATE; - } if (globally_committable && txid.seqno > last_committable) { @@ -971,7 +961,7 @@ namespace ccf::kv pending_txs.insert( {txid.seqno, - std::make_tuple(std::move(pending_tx), globally_committable)}); + std::make_tuple(txid, std::move(pending_tx), globally_committable)}); LOG_TRACE_FMT("Inserting pending tx at {}", txid.seqno); @@ -1009,7 +999,8 @@ namespace ccf::kv } size_t offset = 1; - for (auto& [pending_tx_, committable_] : contiguous_pending_txs) + for (auto& [pending_txid_, pending_tx_, committable_] : + contiguous_pending_txs) { auto [success_, data_, claims_digest_, commit_evidence_digest_, hooks_] = @@ -1057,10 +1048,7 @@ namespace ccf::kv txid.seqno); batch.emplace_back( - previous_last_replicated + offset, - data_shared, - committable_, - hooks_shared); + pending_txid_, data_shared, committable_, hooks_shared); offset++; } diff --git a/src/kv/test/stub_consensus.h b/src/kv/test/stub_consensus.h index 5d40a1558183..6f5220e5a874 100644 --- a/src/kv/test/stub_consensus.h +++ b/src/kv/test/stub_consensus.h @@ -111,15 +111,15 @@ namespace ccf::kv::test { replica.push_back(entry); - const auto& [v, data, committable, hooks] = entry; + const auto& [tx_id, data, committable, hooks] = entry; // Simplification: all entries are replicated in the same term - view_history.update(v, view); + view_history.update(tx_id.seqno, view); if (committable) { // All committable indices are instantly committed - committed_txid = {view, v}; + committed_txid = tx_id; } } current_view = view; diff --git a/src/node/test/historical_queries.cpp b/src/node/test/historical_queries.cpp index 47a8561da949..4c0e62ad6ef7 100644 --- a/src/node/test/historical_queries.cpp +++ b/src/node/test/historical_queries.cpp @@ -223,18 +223,18 @@ std::map> construct_host_ledger( std::map> ledger; auto next_ledger_entry = consensus->pop_oldest_entry(); - auto version = std::get<0>(next_ledger_entry.value()); + auto version = std::get<0>(next_ledger_entry.value()).seqno; while (next_ledger_entry.has_value()) { - const auto ib = ledger.insert(std::make_pair( - std::get<0>(next_ledger_entry.value()), - *std::get<1>(next_ledger_entry.value()))); + const auto next_version = std::get<0>(next_ledger_entry.value()).seqno; + const auto ib = ledger.insert( + std::make_pair(next_version, *std::get<1>(next_ledger_entry.value()))); REQUIRE(ib.second); next_ledger_entry = consensus->pop_oldest_entry(); if (next_ledger_entry.has_value()) { - REQUIRE(version + 1 == std::get<0>(next_ledger_entry.value())); - version = std::get<0>(next_ledger_entry.value()); + REQUIRE(version + 1 == next_version); + version = next_version; } } diff --git a/src/node/test/history.cpp b/src/node/test/history.cpp index 39c955488238..c5a9789db4c8 100644 --- a/src/node/test/history.cpp +++ b/src/node/test/history.cpp @@ -258,11 +258,11 @@ class CompactingConsensus : public ccf::kv::test::StubConsensus bool replicate(const ccf::kv::BatchVector& entries, ccf::View view) override { - for (auto& [version, data, committable, hooks] : entries) + for (auto& [tx_id, data, committable, hooks] : entries) { count++; if (committable) - store->compact(version); + store->compact(tx_id.seqno); } return true; } @@ -431,10 +431,10 @@ class RollbackConsensus : public ccf::kv::test::StubConsensus bool replicate(const ccf::kv::BatchVector& entries, ccf::View view) override { - for (auto& [version, data, committable, hook] : entries) + for (auto& [tx_id, data, committable, hook] : entries) { count++; - if (version == rollback_at) + if (tx_id.seqno == rollback_at) store->rollback({view, rollback_to}, store->commit_view()); } return true; diff --git a/tests/schema.py b/tests/schema.py index 331a64c5cd53..5a35f8856c17 100644 --- a/tests/schema.py +++ b/tests/schema.py @@ -264,14 +264,15 @@ def add(parser): initial_member_count=1, ) - cr.add( - "operations", - e2e_operations.run, - package="samples/apps/logging/logging", - nodes=infra.e2e_args.min_nodes(cr.args, f=0), - initial_user_count=1, - ledger_chunk_bytes="1B", # Chunk ledger at every signature transaction - ) + ## TODO: Too slow, temporarily disabled for faster smoke testing + # cr.add( + # "operations", + # e2e_operations.run, + # package="samples/apps/logging/logging", + # nodes=infra.e2e_args.min_nodes(cr.args, f=0), + # initial_user_count=1, + # ledger_chunk_bytes="1B", # Chunk ledger at every signature transaction + # ) cr.add( "download", From f32ab6dbec06efd665e485948c9eef5bf3582690 Mon Sep 17 00:00:00 2001 From: Eddy Ashton Date: Mon, 24 Aug 2026 14:47:45 +0000 Subject: [PATCH 04/12] Behaviour change - Raft considers each transaction's actual term, in order --- src/consensus/aft/raft.h | 47 +++++++++++++++++++--------------------- 1 file changed, 22 insertions(+), 25 deletions(-) diff --git a/src/consensus/aft/raft.h b/src/consensus/aft/raft.h index 466f61a34fb6..fea12a04cef7 100644 --- a/src/consensus/aft/raft.h +++ b/src/consensus/aft/raft.h @@ -634,16 +634,6 @@ namespace aft return false; } - if (term != state->current_view) - { - RAFT_DEBUG_FMT( - "Failed to replicate {} items at term {}, current term is {}", - entries.size(), - term, - state->current_view); - return false; - } - if (is_retired_committed()) { RAFT_DEBUG_FMT( @@ -657,23 +647,31 @@ namespace aft for (const auto& [tx_id, data, is_globally_committable, hooks] : entries) { - // TODO TODO: This is a temporary hack to allow the new TxID type to be - // used in the raft code. Once the raft code is fully migrated to use - // TxID, this can be removed. - const auto index = tx_id.seqno; bool globally_committable = is_globally_committable; - if (index != state->last_idx + 1) + if (tx_id.seqno != state->last_idx + 1) { - LOG_INFO_FMT( - "!!!! Not contiguous ({} != {} + 1)", index, state->last_idx); + RAFT_DEBUG_FMT( + "Received non-contiguous batch: {} != {} + 1", + tx_id.seqno, + state->last_idx); + return false; + } + + if (tx_id.view != state->current_view) + { + RAFT_DEBUG_FMT( + "Failed to replicate item at {}.{}, current term is {}", + tx_id.view, + tx_id.seqno, + state->current_view); return false; } RAFT_DEBUG_FMT( "Replicated on leader {}: {}{} ({} hooks)", state->node_id, - index, + tx_id.seqno, (globally_committable ? " committable" : ""), hooks->size()); @@ -683,7 +681,7 @@ namespace aft j["state"] = *state; COMMITTABLE_INDICES(j["state"], state); j["view"] = term; - j["seqno"] = index; + j["seqno"] = tx_id.seqno; j["globally_committable"] = globally_committable; RAFT_TRACE_JSON_OUT(j); #endif @@ -703,9 +701,9 @@ namespace aft state->membership_state == ccf::kv::MembershipState::Retired && state->retirement_phase == ccf::kv::RetirementPhase::Ordered) { - become_retired(index, ccf::kv::RetirementPhase::Signed); + become_retired(tx_id.seqno, ccf::kv::RetirementPhase::Signed); } - state->committable_indices.push_back(index); + state->committable_indices.push_back(tx_id.seqno); start_ticking_if_necessary(); // Reset should_sign here - whenever we see a committable entry we @@ -713,13 +711,12 @@ namespace aft should_sign = false; } - state->last_idx = index; - ledger->put_entry( - *data, globally_committable, state->current_view, index); + state->last_idx = tx_id.seqno; + ledger->put_entry(*data, globally_committable, tx_id.view, tx_id.seqno); entry_size_not_limited += data->size(); entry_count++; - state->view_history.update(index, state->current_view); + state->view_history.update(tx_id.seqno, state->current_view); if (entry_size_not_limited >= append_entries_size_limit) { update_batch_size(); From 0ae1e06d64ac93e02cdf66386703580a62f36d83 Mon Sep 17 00:00:00 2001 From: Eddy Ashton Date: Mon, 24 Aug 2026 14:48:33 +0000 Subject: [PATCH 05/12] API change - Stop passing replicate a standalone "replication_view" --- src/consensus/aft/raft.h | 3 +- src/consensus/aft/test/committable_suffix.cpp | 60 +++++++------------ src/consensus/aft/test/driver.h | 3 +- src/consensus/aft/test/main.cpp | 45 +++++++------- src/indexing/test/common.h | 4 +- src/kv/kv_types.h | 2 +- src/kv/store.h | 5 +- src/kv/test/kv_contention.cpp | 4 +- src/kv/test/stub_consensus.h | 9 ++- src/node/test/history.cpp | 8 +-- 10 files changed, 58 insertions(+), 85 deletions(-) diff --git a/src/consensus/aft/raft.h b/src/consensus/aft/raft.h index fea12a04cef7..0d31175dd78e 100644 --- a/src/consensus/aft/raft.h +++ b/src/consensus/aft/raft.h @@ -621,8 +621,7 @@ namespace aft return details; } - // TODO TODO: Mid-refactor. Term argument should be removed - bool replicate(const ccf::kv::BatchVector& entries, Term term) override + bool replicate(const ccf::kv::BatchVector& entries) override { std::lock_guard guard(state->lock); diff --git a/src/consensus/aft/test/committable_suffix.cpp b/src/consensus/aft/test/committable_suffix.cpp index 35b60e249b5b..6a35c2f8888a 100644 --- a/src/consensus/aft/test/committable_suffix.cpp +++ b/src/consensus/aft/test/committable_suffix.cpp @@ -214,8 +214,7 @@ DOCTEST_TEST_CASE("Retention of dead leader's commit") DOCTEST_INFO("Entry at 1.1 is received by all nodes"); { auto entry = make_ledger_entry(1, 1); - rA.replicate( - ccf::kv::BatchVector{{ccf::TxID{1, 1}, entry, true, hooks}}, 1); + rA.replicate(ccf::kv::BatchVector{{ccf::TxID{1, 1}, entry, true, hooks}}); DOCTEST_REQUIRE(rA.get_last_idx() == 1); DOCTEST_REQUIRE(rA.get_committed_seqno() == 0); // Size limit was reached, so periodic is not needed @@ -245,24 +244,21 @@ DOCTEST_TEST_CASE("Retention of dead leader's commit") "committed"); { auto entry = make_ledger_entry(1, 2); - rA.replicate( - ccf::kv::BatchVector{{ccf::TxID{1, 2}, entry, true, hooks}}, 1); + rA.replicate(ccf::kv::BatchVector{{ccf::TxID{1, 2}, entry, true, hooks}}); DOCTEST_REQUIRE(rA.get_last_idx() == 2); DOCTEST_REQUIRE(rA.get_committed_seqno() == 1); // Size limit was reached, so periodic is not needed // rA.periodic(request_timeout); entry = make_ledger_entry(1, 3); - rA.replicate( - ccf::kv::BatchVector{{ccf::TxID{1, 3}, entry, true, hooks}}, 1); + rA.replicate(ccf::kv::BatchVector{{ccf::TxID{1, 3}, entry, true, hooks}}); DOCTEST_REQUIRE(rA.get_last_idx() == 3); DOCTEST_REQUIRE(rA.get_committed_seqno() == 1); // Size limit was reached, so periodic is not needed // rA.periodic(request_timeout); entry = make_ledger_entry(1, 4); - rA.replicate( - ccf::kv::BatchVector{{ccf::TxID{1, 4}, entry, true, hooks}}, 1); + rA.replicate(ccf::kv::BatchVector{{ccf::TxID{1, 4}, entry, true, hooks}}); DOCTEST_REQUIRE(rA.get_last_idx() == 4); DOCTEST_REQUIRE(rA.get_committed_seqno() == 1); // Size limit was reached, so periodic is not needed @@ -296,8 +292,7 @@ DOCTEST_TEST_CASE("Retention of dead leader's commit") "committed"); { auto entry = make_ledger_entry(1, 5); - rA.replicate( - ccf::kv::BatchVector{{ccf::TxID{1, 5}, entry, true, hooks}}, 1); + rA.replicate(ccf::kv::BatchVector{{ccf::TxID{1, 5}, entry, true, hooks}}); DOCTEST_REQUIRE(rA.get_last_idx() == 5); // Size limit was reached, so periodic is not needed // rB.periodic(request_timeout); @@ -372,13 +367,11 @@ DOCTEST_TEST_CASE("Retention of dead leader's commit") DOCTEST_INFO("Node B writes some entries, though they are lost"); { auto entry = make_ledger_entry(2, 6); - rB.replicate( - ccf::kv::BatchVector{{ccf::TxID{2, 6}, entry, true, hooks}}, 2); + rB.replicate(ccf::kv::BatchVector{{ccf::TxID{2, 6}, entry, true, hooks}}); DOCTEST_REQUIRE(rB.get_last_idx() == 6); entry = make_ledger_entry(2, 7); - rB.replicate( - ccf::kv::BatchVector{{ccf::TxID{2, 7}, entry, true, hooks}}, 2); + rB.replicate(ccf::kv::BatchVector{{ccf::TxID{2, 7}, entry, true, hooks}}); DOCTEST_REQUIRE(rB.get_last_idx() == 7); // Size limit was reached, so periodic is not needed @@ -433,18 +426,15 @@ DOCTEST_TEST_CASE("Retention of dead leader's commit") DOCTEST_REQUIRE("Node C produces 3.5, 3.6, and 3.7"); { auto entry = make_ledger_entry(3, 5); - rC.replicate( - ccf::kv::BatchVector{{ccf::TxID{3, 5}, entry, true, hooks}}, 3); + rC.replicate(ccf::kv::BatchVector{{ccf::TxID{3, 5}, entry, true, hooks}}); DOCTEST_REQUIRE(rC.get_last_idx() == 5); entry = make_ledger_entry(3, 6); - rC.replicate( - ccf::kv::BatchVector{{ccf::TxID{3, 6}, entry, true, hooks}}, 3); + rC.replicate(ccf::kv::BatchVector{{ccf::TxID{3, 6}, entry, true, hooks}}); DOCTEST_REQUIRE(rC.get_last_idx() == 6); entry = make_ledger_entry(3, 7); - rC.replicate( - ccf::kv::BatchVector{{ccf::TxID{3, 7}, entry, true, hooks}}, 3); + rC.replicate(ccf::kv::BatchVector{{ccf::TxID{3, 7}, entry, true, hooks}}); DOCTEST_REQUIRE(rC.get_last_idx() == 7); // The early AppendEntries that describe this are lost @@ -648,10 +638,8 @@ DOCTEST_TEST_CASE_TEMPLATE("Multi-term divergence", T, WorstCase, RandomCase) for (auto idx = start_idx + 1; idx <= start_idx + num_entries; ++idx) { auto entry = make_ledger_entry(primary.get_view(), idx); - primary.replicate( - ccf::kv::BatchVector{ - {ccf::TxID{primary.get_view(), idx}, entry, true, hooks}}, - primary.get_view()); + primary.replicate(ccf::kv::BatchVector{ + {ccf::TxID{primary.get_view(), idx}, entry, true, hooks}}); } // All related AppendEntries are lost @@ -680,11 +668,9 @@ DOCTEST_TEST_CASE_TEMPLATE("Multi-term divergence", T, WorstCase, RandomCase) // be committed auto entry = make_ledger_entry(1, 1); - rA.replicate( - ccf::kv::BatchVector{{ccf::TxID{1, 1}, entry, true, hooks}}, 1); + rA.replicate(ccf::kv::BatchVector{{ccf::TxID{1, 1}, entry, true, hooks}}); entry = make_ledger_entry(1, 2); - rA.replicate( - ccf::kv::BatchVector{{ccf::TxID{1, 2}, entry, true, hooks}}, 1); + rA.replicate(ccf::kv::BatchVector{{ccf::TxID{1, 2}, entry, true, hooks}}); DOCTEST_REQUIRE(rA.get_last_idx() == 2); DOCTEST_REQUIRE(rA.get_committed_seqno() == 0); // Size limit was reached, so periodic is not needed @@ -717,22 +703,18 @@ DOCTEST_TEST_CASE_TEMPLATE("Multi-term divergence", T, WorstCase, RandomCase) // Node A produces 2 additional entries that A and B have, and 2 additional // entries that are only present on A entry = make_ledger_entry(1, 3); - rA.replicate( - ccf::kv::BatchVector{{ccf::TxID{1, 3}, entry, true, hooks}}, 1); + rA.replicate(ccf::kv::BatchVector{{ccf::TxID{1, 3}, entry, true, hooks}}); entry = make_ledger_entry(1, 4); - rA.replicate( - ccf::kv::BatchVector{{ccf::TxID{1, 4}, entry, true, hooks}}, 1); + rA.replicate(ccf::kv::BatchVector{{ccf::TxID{1, 4}, entry, true, hooks}}); keep_messages_for(node_idB, channelsA->messages); DOCTEST_REQUIRE(2 == dispatch_all(nodes, node_idA)); entry = make_ledger_entry(1, 5); - rA.replicate( - ccf::kv::BatchVector{{ccf::TxID{1, 5}, entry, true, hooks}}, 1); + rA.replicate(ccf::kv::BatchVector{{ccf::TxID{1, 5}, entry, true, hooks}}); entry = make_ledger_entry(1, 6); - rA.replicate( - ccf::kv::BatchVector{{ccf::TxID{1, 6}, entry, true, hooks}}, 1); + rA.replicate(ccf::kv::BatchVector{{ccf::TxID{1, 6}, entry, true, hooks}}); channelsA->messages.clear(); channelsB->messages.clear(); @@ -1008,10 +990,8 @@ DOCTEST_TEST_CASE_TEMPLATE("Multi-term divergence", T, WorstCase, RandomCase) const auto view = rPrimary.get_view(); const auto seqno = rPrimary.get_last_idx() + 1; auto final_entry = make_ledger_entry(view, seqno); - rPrimary.replicate( - ccf::kv::BatchVector{ - {ccf::TxID{view, seqno}, final_entry, true, hooks}}, - view); + rPrimary.replicate(ccf::kv::BatchVector{ + {ccf::TxID{view, seqno}, final_entry, true, hooks}}); rPrimary.periodic(request_timeout); keep_earliest_append_entries_for_each_target(channelsPrimary->messages); diff --git a/src/consensus/aft/test/driver.h b/src/consensus/aft/test/driver.h index 1a993eed8afd..e556d64d1fa2 100644 --- a/src/consensus/aft/test/driver.h +++ b/src/consensus/aft/test/driver.h @@ -189,8 +189,7 @@ class RaftDriver auto s = nlohmann::json(aft::ReplicatedData{type, data}).dump(); auto d = std::make_shared>(s.begin(), s.end()); raft->replicate( - ccf::kv::BatchVector{{ccf::TxID{term, idx}, d, committable, hooks}}, - term); + ccf::kv::BatchVector{{ccf::TxID{term, idx}, d, committable, hooks}}); } void add_node(ccf::NodeId node_id) diff --git a/src/consensus/aft/test/main.cpp b/src/consensus/aft/test/main.cpp index c1eda8233095..6b1d50841439 100644 --- a/src/consensus/aft/test/main.cpp +++ b/src/consensus/aft/test/main.cpp @@ -77,8 +77,7 @@ DOCTEST_TEST_CASE("Single node commit" * doctest::test_suite("single")) entry->push_back(2); entry->push_back(3); - r0.replicate( - ccf::kv::BatchVector{{ccf::TxID{1, i}, entry, true, hooks}}, 1); + r0.replicate(ccf::kv::BatchVector{{ccf::TxID{1, i}, entry, true, hooks}}); DOCTEST_REQUIRE(r0.get_last_idx() == i); DOCTEST_REQUIRE(r0.get_committed_seqno() == i); } @@ -429,12 +428,12 @@ DOCTEST_TEST_CASE( DOCTEST_INFO("Try to replicate on a follower, and fail"); std::vector entry = {1, 2, 3}; auto data = std::make_shared>(entry); - DOCTEST_REQUIRE_FALSE(r1.replicate( - ccf::kv::BatchVector{{ccf::TxID{1, 1}, data, true, hooks}}, 1)); + DOCTEST_REQUIRE_FALSE( + r1.replicate(ccf::kv::BatchVector{{ccf::TxID{1, 1}, data, true, hooks}})); DOCTEST_INFO("Tell the leader to replicate a message"); - DOCTEST_REQUIRE(r0.replicate( - ccf::kv::BatchVector{{ccf::TxID{1, 1}, data, true, hooks}}, 1)); + DOCTEST_REQUIRE( + r0.replicate(ccf::kv::BatchVector{{ccf::TxID{1, 1}, data, true, hooks}})); DOCTEST_REQUIRE(r0.ledger->ledger.size() == 1); // The test ledger adds its own header. Confirm that the expected data is @@ -549,8 +548,8 @@ DOCTEST_TEST_CASE("Multiple nodes late join" * doctest::test_suite("multiple")) std::vector first_entry = {1, 2, 3}; auto data = std::make_shared>(first_entry); - DOCTEST_REQUIRE(r0.replicate( - ccf::kv::BatchVector{{ccf::TxID{1, 1}, data, true, hooks}}, 1)); + DOCTEST_REQUIRE( + r0.replicate(ccf::kv::BatchVector{{ccf::TxID{1, 1}, data, true, hooks}})); r0.periodic(request_timeout); DOCTEST_REQUIRE( @@ -664,9 +663,9 @@ DOCTEST_TEST_CASE("Recv append entries logic" * doctest::test_suite("multiple")) auto data_2 = std::make_shared>(second_entry); DOCTEST_REQUIRE(r0.replicate( - ccf::kv::BatchVector{{ccf::TxID{1, 1}, data_1, true, hooks}}, 1)); + ccf::kv::BatchVector{{ccf::TxID{1, 1}, data_1, true, hooks}})); DOCTEST_REQUIRE(r0.replicate( - ccf::kv::BatchVector{{ccf::TxID{1, 2}, data_2, true, hooks}}, 1)); + ccf::kv::BatchVector{{ccf::TxID{1, 2}, data_2, true, hooks}})); DOCTEST_REQUIRE(r0.ledger->ledger.size() == 2); r0.periodic(request_timeout); DOCTEST_REQUIRE(r0c->messages.size() == 1); @@ -687,8 +686,8 @@ DOCTEST_TEST_CASE("Recv append entries logic" * doctest::test_suite("multiple")) { std::vector third_entry = {3, 3, 3}; auto data = std::make_shared>(third_entry); - DOCTEST_REQUIRE(r0.replicate( - ccf::kv::BatchVector{{ccf::TxID{1, 3}, data, true, hooks}}, 1)); + DOCTEST_REQUIRE( + r0.replicate(ccf::kv::BatchVector{{ccf::TxID{1, 3}, data, true, hooks}})); DOCTEST_REQUIRE(r0.ledger->ledger.size() == 3); // Simulate that the append entries was not deserialised successfully @@ -720,8 +719,8 @@ DOCTEST_TEST_CASE("Recv append entries logic" * doctest::test_suite("multiple")) { std::vector fourth_entry = {4, 4, 4}; auto data = std::make_shared>(fourth_entry); - DOCTEST_REQUIRE(r0.replicate( - ccf::kv::BatchVector{{ccf::TxID{1, 4}, data, true, hooks}}, 1)); + DOCTEST_REQUIRE( + r0.replicate(ccf::kv::BatchVector{{ccf::TxID{1, 4}, data, true, hooks}})); DOCTEST_REQUIRE(r0.ledger->ledger.size() == 4); r0.periodic(request_timeout); DOCTEST_REQUIRE(r0c->messages.size() == 1); @@ -734,8 +733,8 @@ DOCTEST_TEST_CASE("Recv append entries logic" * doctest::test_suite("multiple")) { std::vector fifth_entry = {5, 5, 5}; auto data = std::make_shared>(fifth_entry); - DOCTEST_REQUIRE(r0.replicate( - ccf::kv::BatchVector{{ccf::TxID{1, 5}, data, true, hooks}}, 1)); + DOCTEST_REQUIRE( + r0.replicate(ccf::kv::BatchVector{{ccf::TxID{1, 5}, data, true, hooks}})); DOCTEST_REQUIRE(r0.ledger->ledger.size() == 5); r0.periodic(request_timeout); DOCTEST_REQUIRE(r0c->messages.size() == 1); @@ -765,7 +764,7 @@ DOCTEST_TEST_CASE("Recv append entries logic" * doctest::test_suite("multiple")) std::vector entry_6 = {6, 6, 6}; auto data = std::make_shared>(entry_6); DOCTEST_REQUIRE(r0.replicate( - ccf::kv::BatchVector{{ccf::TxID{1, 6}, data, true, hooks}}, 1)); + ccf::kv::BatchVector{{ccf::TxID{1, 6}, data, true, hooks}})); DOCTEST_REQUIRE(r0.ledger->ledger.size() == 6); } const auto last_correct_version = r0.ledger->ledger.size(); @@ -775,7 +774,7 @@ DOCTEST_TEST_CASE("Recv append entries logic" * doctest::test_suite("multiple")) std::vector entry_7 = {7, 7, 7}; auto data = std::make_shared>(entry_7); DOCTEST_REQUIRE(r0.replicate( - ccf::kv::BatchVector{{ccf::TxID{1, 7}, data, true, hooks}}, 1)); + ccf::kv::BatchVector{{ccf::TxID{1, 7}, data, true, hooks}})); DOCTEST_REQUIRE(r0.ledger->ledger.size() == 7); dead_branch = r0.ledger->ledger.back(); } @@ -797,7 +796,7 @@ DOCTEST_TEST_CASE("Recv append entries logic" * doctest::test_suite("multiple")) std::vector entry_7b = {7, 7, 'b'}; auto data = std::make_shared>(entry_7b); DOCTEST_REQUIRE(r0.replicate( - ccf::kv::BatchVector{{ccf::TxID{4, 7}, data, true, hooks}}, 4)); + ccf::kv::BatchVector{{ccf::TxID{4, 7}, data, true, hooks}})); DOCTEST_REQUIRE(r0.ledger->ledger.size() == 7); live_branch = r0.ledger->ledger.back(); } @@ -806,7 +805,7 @@ DOCTEST_TEST_CASE("Recv append entries logic" * doctest::test_suite("multiple")) std::vector entry_8 = {8, 8, 8}; auto data = std::make_shared>(entry_8); DOCTEST_REQUIRE(r0.replicate( - ccf::kv::BatchVector{{ccf::TxID{4, 8}, data, true, hooks}}, 4)); + ccf::kv::BatchVector{{ccf::TxID{4, 8}, data, true, hooks}})); DOCTEST_REQUIRE(r0.ledger->ledger.size() == 8); DOCTEST_REQUIRE(r0.ledger->ledger.size() > last_correct_version); } @@ -937,8 +936,8 @@ DOCTEST_TEST_CASE("Exceed append entries limit") for (size_t i = 1; i <= static_cast(num_big_entries); ++i) { - DOCTEST_REQUIRE(r0.replicate( - ccf::kv::BatchVector{{ccf::TxID{1, i}, data, true, hooks}}, 1)); + DOCTEST_REQUIRE( + r0.replicate(ccf::kv::BatchVector{{ccf::TxID{1, i}, data, true, hooks}})); const auto received_ae = dispatch_all_and_DOCTEST_CHECK( nodes, node_id0, r0c->messages, [](const auto& msg) { @@ -957,7 +956,7 @@ DOCTEST_TEST_CASE("Exceed append entries limit") ++i) { DOCTEST_REQUIRE(r0.replicate( - ccf::kv::BatchVector{{ccf::TxID{1, i}, smaller_data, true, hooks}}, 1)); + ccf::kv::BatchVector{{ccf::TxID{1, i}, smaller_data, true, hooks}})); dispatch_all(nodes, node_id0, r0c->messages); } diff --git a/src/indexing/test/common.h b/src/indexing/test/common.h index e77f6515a360..458bc6ea75a8 100644 --- a/src/indexing/test/common.h +++ b/src/indexing/test/common.h @@ -87,7 +87,7 @@ class AllCommittableWrapper : public TConsensus public: using TConsensus::TConsensus; - bool replicate(const ccf::kv::BatchVector& entries_, ccf::View view) override + bool replicate(const ccf::kv::BatchVector& entries_) override { // Rather than building a history that produces real signatures, we just // overwrite the entries here to say that everything is committable @@ -97,7 +97,7 @@ class AllCommittableWrapper : public TConsensus committable = true; } - return TConsensus::replicate(entries, view); + return TConsensus::replicate(entries); } }; diff --git a/src/kv/kv_types.h b/src/kv/kv_types.h index 0d6b0a98c4f9..f49c1f87ab8c 100644 --- a/src/kv/kv_types.h +++ b/src/kv/kv_types.h @@ -440,7 +440,7 @@ namespace ccf::kv virtual void init_as_backup( ccf::SeqNo, ccf::View, const std::vector&, ccf::SeqNo) = 0; - virtual bool replicate(const BatchVector& entries, ccf::View view) = 0; + virtual bool replicate(const BatchVector& entries) = 0; virtual std::pair get_committed_txid() = 0; virtual ccf::View get_view(ccf::SeqNo seqno) = 0; diff --git a/src/kv/store.h b/src/kv/store.h index 180474008bf6..f6580accf231 100644 --- a/src/kv/store.h +++ b/src/kv/store.h @@ -946,7 +946,6 @@ namespace ccf::kv Version previous_last_replicated = 0; Version next_last_replicated = 0; Version previous_rollback_count = 0; - ccf::View replication_view = 0; // TODO: Remove std::vector contiguous_pending_txs; auto h = get_history(); @@ -988,8 +987,6 @@ namespace ccf::kv previous_rollback_count = rollback_count; previous_last_replicated = last_replicated; next_last_replicated = last_replicated + contiguous_pending_txs.size(); - - replication_view = term_of_next_version; } // Release version lock @@ -1053,7 +1050,7 @@ namespace ccf::kv offset++; } - if (c->replicate(batch, replication_view)) + if (c->replicate(batch)) { std::lock_guard vguard(version_lock); if ( diff --git a/src/kv/test/kv_contention.cpp b/src/kv/test/kv_contention.cpp index 3bdde2a260f7..ff4cae578896 100644 --- a/src/kv/test/kv_contention.cpp +++ b/src/kv/test/kv_contention.cpp @@ -22,7 +22,7 @@ class SlowStubConsensus : public ccf::kv::test::StubConsensus public: using ccf::kv::test::StubConsensus::StubConsensus; - bool replicate(const ccf::kv::BatchVector& entries, ccf::View view) override + bool replicate(const ccf::kv::BatchVector& entries) override { if (rand() % 2 == 0) { @@ -30,7 +30,7 @@ class SlowStubConsensus : public ccf::kv::test::StubConsensus std::this_thread::sleep_for(std::chrono::milliseconds(delay)); } - return ccf::kv::test::StubConsensus::replicate(entries, view); + return ccf::kv::test::StubConsensus::replicate(entries); } }; diff --git a/src/kv/test/stub_consensus.h b/src/kv/test/stub_consensus.h index 6f5220e5a874..a7510aa6eb5c 100644 --- a/src/kv/test/stub_consensus.h +++ b/src/kv/test/stub_consensus.h @@ -105,7 +105,7 @@ namespace ccf::kv::test state = Backup; } - bool replicate(const BatchVector& entries, ccf::View view) override + bool replicate(const BatchVector& entries) override { for (const auto& entry : entries) { @@ -113,16 +113,15 @@ namespace ccf::kv::test const auto& [tx_id, data, committable, hooks] = entry; - // Simplification: all entries are replicated in the same term - view_history.update(tx_id.seqno, view); + view_history.update(tx_id.seqno, tx_id.view); if (committable) { // All committable indices are instantly committed committed_txid = tx_id; } + current_view = tx_id.view; } - current_view = view; return true; } @@ -236,7 +235,7 @@ namespace ccf::kv::test return false; } - bool replicate(const BatchVector& entries, ccf::View view) override + bool replicate(const BatchVector& entries) override { return false; } diff --git a/src/node/test/history.cpp b/src/node/test/history.cpp index c5a9789db4c8..84f3f490fe65 100644 --- a/src/node/test/history.cpp +++ b/src/node/test/history.cpp @@ -40,7 +40,7 @@ class DummyConsensus : public ccf::kv::test::StubConsensus DummyConsensus(ccf::kv::Store* store_) : store(store_) {} - bool replicate(const ccf::kv::BatchVector& entries, ccf::View view) override + bool replicate(const ccf::kv::BatchVector& entries) override { if (store) { @@ -256,7 +256,7 @@ class CompactingConsensus : public ccf::kv::test::StubConsensus CompactingConsensus(ccf::kv::Store* store_) : store(store_) {} - bool replicate(const ccf::kv::BatchVector& entries, ccf::View view) override + bool replicate(const ccf::kv::BatchVector& entries) override { for (auto& [tx_id, data, committable, hooks] : entries) { @@ -429,13 +429,13 @@ class RollbackConsensus : public ccf::kv::test::StubConsensus rollback_to(rollback_to_) {} - bool replicate(const ccf::kv::BatchVector& entries, ccf::View view) override + bool replicate(const ccf::kv::BatchVector& entries) override { for (auto& [tx_id, data, committable, hook] : entries) { count++; if (tx_id.seqno == rollback_at) - store->rollback({view, rollback_to}, store->commit_view()); + store->rollback(tx_id, store->commit_view()); } return true; } From 5e8c7645f1c63d976c7cc143c986748998db95cc Mon Sep 17 00:00:00 2001 From: Eddy Ashton Date: Tue, 25 Aug 2026 14:02:17 +0000 Subject: [PATCH 06/12] More API tweaks for partial application, a few error cases handled early, and unit tests --- CMakeLists.txt | 3 +- src/consensus/aft/raft.h | 166 +++--- .../aft/test/view_straddling_transactions.cpp | 491 ++++++++++++++++++ src/kv/kv_types.h | 2 +- src/kv/store.h | 58 ++- src/kv/test/stub_consensus.h | 8 +- 6 files changed, 626 insertions(+), 102 deletions(-) create mode 100644 src/consensus/aft/test/view_straddling_transactions.cpp diff --git a/CMakeLists.txt b/CMakeLists.txt index e136af2d2920..6c3d888d8f8c 100644 --- a/CMakeLists.txt +++ b/CMakeLists.txt @@ -716,8 +716,9 @@ if(BUILD_TESTS) ${CMAKE_CURRENT_SOURCE_DIR}/src/consensus/aft/test/main.cpp ${CMAKE_CURRENT_SOURCE_DIR}/src/consensus/aft/test/view_history.cpp ${CMAKE_CURRENT_SOURCE_DIR}/src/consensus/aft/test/committable_suffix.cpp + ${CMAKE_CURRENT_SOURCE_DIR}/src/consensus/aft/test/view_straddling_transactions.cpp ) - target_link_libraries(raft_test PRIVATE ccfcrypto ccf_tasks) + target_link_libraries(raft_test PRIVATE ccfcrypto ccf_kv ccf_tasks) add_unit_test( raft_enclave_test diff --git a/src/consensus/aft/raft.h b/src/consensus/aft/raft.h index 0d31175dd78e..57ef71a32ad5 100644 --- a/src/consensus/aft/raft.h +++ b/src/consensus/aft/raft.h @@ -621,121 +621,129 @@ namespace aft return details; } - bool replicate(const ccf::kv::BatchVector& entries) override + size_t replicate(const ccf::kv::BatchVector& entries) override { std::lock_guard guard(state->lock); + size_t replicated_count = 0; + if (state->leadership_state != ccf::kv::LeadershipState::Leader) { RAFT_DEBUG_FMT( "Failed to replicate {} items: not leader", entries.size()); - rollback(state->last_idx); - return false; } - - if (is_retired_committed()) + else if (is_retired_committed()) { RAFT_DEBUG_FMT( "Failed to replicate {} items: node retirement is complete", entries.size()); - rollback(state->last_idx); - return false; } - - RAFT_DEBUG_FMT("Replicating {} entries", entries.size()); - - for (const auto& [tx_id, data, is_globally_committable, hooks] : entries) + else { - bool globally_committable = is_globally_committable; + RAFT_DEBUG_FMT("Replicating {} entries", entries.size()); - if (tx_id.seqno != state->last_idx + 1) + for (const auto& [tx_id, data, is_globally_committable, hooks] : + entries) { - RAFT_DEBUG_FMT( - "Received non-contiguous batch: {} != {} + 1", - tx_id.seqno, - state->last_idx); - return false; - } + bool globally_committable = is_globally_committable; + + if (tx_id.seqno != state->last_idx + 1) + { + RAFT_DEBUG_FMT( + "Received non-contiguous batch: {} != {} + 1", + tx_id.seqno, + state->last_idx); + break; + } + + if (tx_id.view != state->current_view) + { + RAFT_DEBUG_FMT( + "Failed to replicate item at {}.{}, current term is {}", + tx_id.view, + tx_id.seqno, + state->current_view); + break; + } - if (tx_id.view != state->current_view) - { RAFT_DEBUG_FMT( - "Failed to replicate item at {}.{}, current term is {}", - tx_id.view, + "Replicated on leader {}: {}{} ({} hooks)", + state->node_id, tx_id.seqno, - state->current_view); - return false; - } - - RAFT_DEBUG_FMT( - "Replicated on leader {}: {}{} ({} hooks)", - state->node_id, - tx_id.seqno, - (globally_committable ? " committable" : ""), - hooks->size()); + (globally_committable ? " committable" : ""), + hooks->size()); #ifdef CCF_RAFT_TRACING - nlohmann::json j = {}; - j["function"] = "replicate"; - j["state"] = *state; - COMMITTABLE_INDICES(j["state"], state); - j["view"] = term; - j["seqno"] = tx_id.seqno; - j["globally_committable"] = globally_committable; - RAFT_TRACE_JSON_OUT(j); + nlohmann::json j = {}; + j["function"] = "replicate"; + j["state"] = *state; + COMMITTABLE_INDICES(j["state"], state); + j["view"] = term; // TODO + j["seqno"] = tx_id.seqno; + j["globally_committable"] = globally_committable; + RAFT_TRACE_JSON_OUT(j); #endif - for (auto& hook : *hooks) - { - hook->call(this); - } - - if (globally_committable) - { - RAFT_DEBUG_FMT( - "membership: {} leadership: {}", - state->membership_state, - state->leadership_state); - if ( - state->membership_state == ccf::kv::MembershipState::Retired && - state->retirement_phase == ccf::kv::RetirementPhase::Ordered) + for (auto& hook : *hooks) { - become_retired(tx_id.seqno, ccf::kv::RetirementPhase::Signed); + hook->call(this); } - state->committable_indices.push_back(tx_id.seqno); - start_ticking_if_necessary(); - // Reset should_sign here - whenever we see a committable entry we - // don't need to produce _another_ signature - should_sign = false; - } + if (globally_committable) + { + RAFT_DEBUG_FMT( + "membership: {} leadership: {}", + state->membership_state, + state->leadership_state); + if ( + state->membership_state == ccf::kv::MembershipState::Retired && + state->retirement_phase == ccf::kv::RetirementPhase::Ordered) + { + become_retired(tx_id.seqno, ccf::kv::RetirementPhase::Signed); + } + state->committable_indices.push_back(tx_id.seqno); + start_ticking_if_necessary(); - state->last_idx = tx_id.seqno; - ledger->put_entry(*data, globally_committable, tx_id.view, tx_id.seqno); - entry_size_not_limited += data->size(); - entry_count++; + // Reset should_sign here - whenever we see a committable entry we + // don't need to produce _another_ signature + should_sign = false; + } - state->view_history.update(tx_id.seqno, state->current_view); - if (entry_size_not_limited >= append_entries_size_limit) - { - update_batch_size(); - entry_count = 0; - entry_size_not_limited = 0; - for (const auto& it : all_other_nodes) + state->last_idx = tx_id.seqno; + ledger->put_entry( + *data, globally_committable, tx_id.view, tx_id.seqno); + entry_size_not_limited += data->size(); + entry_count++; + + state->view_history.update(tx_id.seqno, state->current_view); + if (entry_size_not_limited >= append_entries_size_limit) { - RAFT_DEBUG_FMT("Sending updates to follower {}", it.first); - send_append_entries(it.first, it.second.sent_idx + 1); + update_batch_size(); + entry_count = 0; + entry_size_not_limited = 0; + for (const auto& it : all_other_nodes) + { + RAFT_DEBUG_FMT("Sending updates to follower {}", it.first); + send_append_entries(it.first, it.second.sent_idx + 1); + } } + + replicated_count++; + } + + // Try to advance commit at once if there are no other nodes. + if (other_nodes_in_active_configs().size() == 0) + { + update_commit(); } } - // Try to advance commit at once if there are no other nodes. - if (other_nodes_in_active_configs().size() == 0) + if (replicated_count != entries.size()) { - update_commit(); + rollback(state->last_idx); } - return true; + return replicated_count; } void recv_message( diff --git a/src/consensus/aft/test/view_straddling_transactions.cpp b/src/consensus/aft/test/view_straddling_transactions.cpp new file mode 100644 index 000000000000..4f5aec05ea1b --- /dev/null +++ b/src/consensus/aft/test/view_straddling_transactions.cpp @@ -0,0 +1,491 @@ +// Copyright (c) Microsoft Corporation. All rights reserved. +// Licensed under the Apache 2.0 License. + +#include "kv/store.h" +#include "kv/test/null_encryptor.h" +#include "test_common.h" + +#include +#include +#include +#include + +namespace +{ + using TestMap = ccf::kv::Map; + using Raft = aft::Aft; + + class BaselinePendingTx : public ccf::kv::PendingTx + { + ccf::TxID txid; + ccf::kv::Store& store; + TestMap& table; + + public: + BaselinePendingTx( + ccf::TxID txid_, ccf::kv::Store& store_, TestMap& table_) : + txid(txid_), + store(store_), + table(table_) + {} + + ccf::kv::PendingTxInfo call() override + { + auto tx = store.create_reserved_tx(txid); + tx.rw(table)->put(0, 1); + return tx.commit_reserved(); + } + }; + + struct CommitPause + { + std::mutex lock; + std::condition_variable paused_cv; + std::condition_variable resume_cv; + bool paused = false; + bool resume = false; + + void pause() + { + { + std::lock_guard guard(lock); + paused = true; + } + paused_cv.notify_one(); + + std::unique_lock guard(lock); + resume_cv.wait(guard, [this]() { return resume; }); + } + + void wait_until_paused() + { + std::unique_lock guard(lock); + paused_cv.wait(guard, [this]() { return paused; }); + } + + void release() + { + { + std::lock_guard guard(lock); + resume = true; + } + resume_cv.notify_one(); + } + }; + + static std::optional read_value( + ccf::kv::Store& store, TestMap& table, size_t key) + { + auto tx = store.create_read_only_tx(); + return tx.ro(table)->get(key); + } + + struct Fixture + { + const ccf::NodeId node_id = ccf::kv::test::PrimaryNodeId; + std::shared_ptr store = std::make_shared(); + TestMap table{"public:table"}; + std::shared_ptr raft; + ccf::View initial_view = 0; + + Fixture() + { + store->set_encryptor(std::make_shared()); + raft = std::make_shared( + raft_settings, + std::make_unique>(store), + std::make_unique(node_id), + std::make_shared(), + std::make_shared(node_id), + nullptr); + store->set_consensus(raft); + + ccf::kv::Configuration::Nodes configuration; + configuration.try_emplace(node_id); + raft->add_configuration(0, configuration); + raft->force_become_primary(); + initial_view = raft->get_view(); + + const auto baseline_txid = store->next_txid(); + REQUIRE( + store->commit( + baseline_txid, + std::make_unique(baseline_txid, *store, table), + true) == ccf::kv::CommitResult::SUCCESS); + REQUIRE(store->current_txid() == ccf::TxID(initial_view, 1)); + REQUIRE(raft->get_committed_seqno() == 1); + REQUIRE(raft->ledger->ledger.size() == 1); + } + + void step_down() + { + const auto next_view = raft->get_view() + 1; + raft->become_aware_of_new_term(next_view); + } + + ccf::View reelect() + { + step_down(); + raft->force_become_primary(); + return raft->get_view(); + } + }; + + static ccf::kv::BatchVector::value_type make_entry( + const ccf::TxID& tx_id, bool globally_committable = false) + { + return { + tx_id, + std::make_shared>(16, tx_id.seqno), + globally_committable, + std::make_shared()}; + } +} + +TEST_CASE( + "Long-lived transaction is rolled back after leadership loss" * + doctest::test_suite("view_straddling_transactions")) +{ + Fixture fixture; + + INFO("Start applying a local transaction in the initial view"); + auto stale_tx = fixture.store->create_tx(); + stale_tx.rw(fixture.table)->put(1, 2); + + CommitPause pause; + std::optional stale_result; + std::thread stale_worker([&]() { + stale_result = + stale_tx.commit(ccf::empty_claims(), [&pause](const auto&, const auto&) { + pause.pause(); + }); + }); + pause.wait_until_paused(); + REQUIRE(stale_tx.get_txid() == ccf::TxID(fixture.initial_view, 2)); + + INFO("Step down after the transaction has been assigned its old-view TxID"); + fixture.step_down(); + + INFO("AFT rejects the transaction and rolls Store back to the baseline"); + pause.release(); + stale_worker.join(); + REQUIRE(stale_result.has_value()); + REQUIRE(stale_result.value() == ccf::kv::CommitResult::FAIL_NO_REPLICATE); + REQUIRE(fixture.store->current_txid() == ccf::TxID(fixture.initial_view, 1)); + REQUIRE_FALSE(read_value(*fixture.store, fixture.table, 1).has_value()); + REQUIRE(fixture.raft->get_last_idx() == 1); + REQUIRE(fixture.raft->ledger->ledger.size() == 1); + + INFO("Win a later election and replicate the next transaction normally"); + fixture.raft->force_become_primary(); + const auto fresh_view = fixture.raft->get_view(); + auto fresh_tx = fixture.store->create_tx(); + fresh_tx.rw(fixture.table)->put(2, 3); + REQUIRE(fresh_tx.commit() == ccf::kv::CommitResult::SUCCESS); + REQUIRE(fixture.store->current_txid() == ccf::TxID(fresh_view, 2)); + REQUIRE(read_value(*fixture.store, fixture.table, 2) == 3); + REQUIRE(fixture.raft->get_last_idx() == 2); + REQUIRE(fixture.raft->ledger->ledger.size() == 2); +} + +TEST_CASE( + "Read-only transaction can finish after re-election" * + doctest::test_suite("view_straddling_transactions")) +{ + Fixture fixture; + + INFO("Read the baseline in the initial view"); + auto read_tx = fixture.store->create_tx(); + REQUIRE(read_tx.ro(fixture.table)->get(0) == 1); + const auto read_txid = fixture.store->current_txid(); + + INFO("Lose leadership and win a later election before finishing the read"); + fixture.reelect(); + + INFO("The read remains valid at the TxID where it observed state"); + REQUIRE(read_tx.commit() == ccf::kv::CommitResult::SUCCESS); + REQUIRE(read_tx.get_txid() == read_txid); + CHECK(fixture.store->current_txid() == read_txid); + CHECK(fixture.raft->get_last_idx() == 1); + CHECK(fixture.raft->ledger->ledger.size() == 1); +} + +TEST_CASE( + "Transaction begun before re-election can commit in the new view" * + doctest::test_suite("view_straddling_transactions")) +{ + Fixture fixture; + + INFO("Read state and prepare writes in the initial view"); + auto tx = fixture.store->create_tx(); + auto handle = tx.rw(fixture.table); + REQUIRE(handle->get(0) == 1); + handle->put(1, 2); + + INFO("Win a later election before assigning the transaction a TxID"); + const auto reelection_view = fixture.reelect(); + + INFO("Revalidate the read and assign the write a TxID in the new view"); + REQUIRE(tx.commit() == ccf::kv::CommitResult::SUCCESS); + REQUIRE(tx.get_txid() == ccf::TxID(reelection_view, 2)); + CHECK(fixture.store->current_txid() == ccf::TxID(reelection_view, 2)); + CHECK(read_value(*fixture.store, fixture.table, 1) == 2); + CHECK(fixture.raft->get_last_idx() == 2); + CHECK(fixture.raft->ledger->ledger.size() == 2); +} + +TEST_CASE( + "Transaction conflicts when re-election rolls back its read snapshot" * + doctest::test_suite("view_straddling_transactions")) +{ + Fixture fixture; + + INFO("Read state and prepare writes in the initial view"); + auto tx = fixture.store->create_tx(); + auto handle = tx.rw(fixture.table); + REQUIRE(handle->get(0) == 1); + handle->put(1, 2); + + INFO("Create an unreplicated local suffix after the transaction's read"); + REQUIRE(fixture.store->next_txid() == ccf::TxID(fixture.initial_view, 2)); + auto suffix_tx = fixture.store->create_tx(); + suffix_tx.rw(fixture.table)->put(0, 9); + REQUIRE(suffix_tx.commit() == ccf::kv::CommitResult::SUCCESS); + REQUIRE(suffix_tx.get_txid() == ccf::TxID(fixture.initial_view, 3)); + REQUIRE(fixture.raft->get_last_idx() == 1); + + INFO("Win a later election, rolling the local suffix back"); + fixture.reelect(); + REQUIRE(read_value(*fixture.store, fixture.table, 0) == 1); + + INFO("The transaction's pre-rollback change set is no longer valid"); + CHECK(tx.commit() == ccf::kv::CommitResult::FAIL_CONFLICT); + CHECK(fixture.store->current_txid() == ccf::TxID(fixture.initial_view, 1)); + CHECK_FALSE(read_value(*fixture.store, fixture.table, 1).has_value()); + CHECK(fixture.raft->get_last_idx() == 1); + CHECK(fixture.raft->ledger->ledger.size() == 1); +} + +TEST_CASE( + "Assigned old-view transaction is rolled back after re-election" * + doctest::test_suite("view_straddling_transactions")) +{ + Fixture fixture; + + INFO("Apply a transaction and assign its TxID in the initial view"); + auto stale_tx = fixture.store->create_tx(); + stale_tx.rw(fixture.table)->put(1, 2); + CommitPause pause; + std::optional stale_result; + std::thread stale_worker([&]() { + stale_result = + stale_tx.commit(ccf::empty_claims(), [&pause](const auto&, const auto&) { + pause.pause(); + }); + }); + pause.wait_until_paused(); + REQUIRE(stale_tx.get_txid() == ccf::TxID(fixture.initial_view, 2)); + + INFO("Lose leadership and win a later election before replication"); + const auto reelection_view = fixture.reelect(); + + INFO("AFT rejects the assigned old-view transaction"); + pause.release(); + stale_worker.join(); + REQUIRE(stale_result.has_value()); + REQUIRE(stale_result.value() == ccf::kv::CommitResult::FAIL_NO_REPLICATE); + CHECK(fixture.store->current_txid() == ccf::TxID(fixture.initial_view, 1)); + CHECK_FALSE(read_value(*fixture.store, fixture.table, 1).has_value()); + CHECK(fixture.raft->get_last_idx() == 1); + CHECK(fixture.raft->ledger->ledger.size() == 1); + + INFO("Replicate a fresh transaction at the next index in the new view"); + auto fresh_tx = fixture.store->create_tx(); + fresh_tx.rw(fixture.table)->put(2, 3); + CHECK(fresh_tx.commit() == ccf::kv::CommitResult::SUCCESS); + CHECK(fixture.store->current_txid() == ccf::TxID(reelection_view, 2)); + CHECK(read_value(*fixture.store, fixture.table, 2) == 3); + CHECK(fixture.raft->get_last_idx() == 2); + CHECK(fixture.raft->ledger->ledger.size() == 2); +} + +TEST_CASE( + "Rolled-back stale transaction cannot invalidate current-view work" * + doctest::test_suite("view_straddling_transactions")) +{ + Fixture fixture; + + INFO( + "Assign an old-view transaction seqno 2, then pause before Store::commit"); + auto stale_tx = fixture.store->create_tx(); + stale_tx.rw(fixture.table)->put(1, 2); + CommitPause stale_pause; + std::optional stale_result; + std::thread stale_worker([&]() { + stale_result = stale_tx.commit( + ccf::empty_claims(), + [&stale_pause](const auto&, const auto&) { stale_pause.pause(); }); + }); + stale_pause.wait_until_paused(); + REQUIRE(stale_tx.get_txid() == ccf::TxID(fixture.initial_view, 2)); + + INFO("Win a later election, reclaiming seqno 2"); + const auto reelection_view = fixture.reelect(); + + INFO("Assign current-view seqno 2, then pause before Store::commit"); + auto current_head = fixture.store->create_tx(); + current_head.rw(fixture.table)->put(2, 3); + CommitPause current_pause; + std::optional current_result; + std::thread current_worker([&]() { + current_result = current_head.commit( + ccf::empty_claims(), + [¤t_pause](const auto&, const auto&) { current_pause.pause(); }); + }); + current_pause.wait_until_paused(); + REQUIRE(current_head.get_txid() == ccf::TxID(reelection_view, 2)); + + INFO("Queue current-view seqno 3 behind the missing seqno 2"); + auto current_suffix = fixture.store->create_tx(); + current_suffix.rw(fixture.table)->put(3, 4); + REQUIRE(current_suffix.commit() == ccf::kv::CommitResult::SUCCESS); + REQUIRE(current_suffix.get_txid() == ccf::TxID(reelection_view, 3)); + + INFO("Resume old 2.2; Store must reject its invalidated local application"); + stale_pause.release(); + stale_worker.join(); + REQUIRE(stale_result.has_value()); + CHECK(stale_result.value() == ccf::kv::CommitResult::FAIL_NO_REPLICATE); + + INFO("Resume current 3.2, which can now replicate with pending 3.3"); + current_pause.release(); + current_worker.join(); + REQUIRE(current_result.has_value()); + CHECK(current_result.value() == ccf::kv::CommitResult::SUCCESS); + + const auto a = fixture.store->current_txid(); + const auto b = ccf::TxID(reelection_view, 3); + std::cout << "Current TxID: " << a.to_str() << ", expected: " << b.to_str() + << std::endl; + CHECK(fixture.store->current_txid() == ccf::TxID(reelection_view, 3)); + CHECK_FALSE(read_value(*fixture.store, fixture.table, 1).has_value()); + CHECK(read_value(*fixture.store, fixture.table, 2) == 3); + CHECK(read_value(*fixture.store, fixture.table, 3) == 4); + CHECK(fixture.raft->get_last_idx() == 3); + CHECK(fixture.raft->ledger->ledger.size() == 3); + + INFO("Replicate a fresh new-view transaction at seqno 4"); + auto fresh_tx = fixture.store->create_tx(); + fresh_tx.rw(fixture.table)->put(4, 5); + CHECK(fresh_tx.commit() == ccf::kv::CommitResult::SUCCESS); + CHECK(fixture.store->current_txid() == ccf::TxID(reelection_view, 4)); + CHECK(read_value(*fixture.store, fixture.table, 4) == 5); + CHECK(fixture.raft->get_last_idx() == 4); + CHECK(fixture.raft->ledger->ledger.size() == 4); +} + +TEST_CASE( + "Rolled-back stale suffix cannot block a current-view prefix" * + doctest::test_suite("view_straddling_transactions")) +{ + Fixture fixture; + + INFO("Reserve old-view seqno 2, leaving a hole before stale seqno 3"); + REQUIRE(fixture.store->next_txid() == ccf::TxID(fixture.initial_view, 2)); + + auto stale_suffix = fixture.store->create_tx(); + stale_suffix.rw(fixture.table)->put(1, 2); + CommitPause stale_pause; + std::optional stale_result; + std::thread stale_worker([&]() { + stale_result = stale_suffix.commit( + ccf::empty_claims(), + [&stale_pause](const auto&, const auto&) { stale_pause.pause(); }); + }); + stale_pause.wait_until_paused(); + REQUIRE(stale_suffix.get_txid() == ccf::TxID(fixture.initial_view, 3)); + + INFO("Win a later election, rolling back the old reservation and write"); + const auto reelection_view = fixture.reelect(); + + INFO( + "Apply the current-view transaction at seqno 2, then pause before " + "Store::commit"); + auto current_tx = fixture.store->create_tx(); + current_tx.rw(fixture.table)->put(2, 3); + CommitPause current_pause; + std::optional current_result; + std::thread current_worker([&]() { + current_result = current_tx.commit( + ccf::empty_claims(), + [¤t_pause](const auto&, const auto&) { current_pause.pause(); }); + }); + current_pause.wait_until_paused(); + REQUIRE(current_tx.get_txid() == ccf::TxID(reelection_view, 2)); + + INFO("Resume stale 2.3; Store must reject its invalidated local application"); + stale_pause.release(); + stale_worker.join(); + REQUIRE(stale_result.has_value()); + CHECK(stale_result.value() == ccf::kv::CommitResult::FAIL_NO_REPLICATE); + + INFO("Resume seqno 2; AFT can accept the current-view transaction"); + current_pause.release(); + current_worker.join(); + + REQUIRE(current_result.has_value()); + CHECK(current_result.value() == ccf::kv::CommitResult::SUCCESS); + CHECK(fixture.store->current_txid() == ccf::TxID(reelection_view, 2)); + CHECK_FALSE(read_value(*fixture.store, fixture.table, 1).has_value()); + CHECK(read_value(*fixture.store, fixture.table, 2) == 3); + CHECK(fixture.raft->get_last_idx() == 2); + CHECK(fixture.raft->ledger->ledger.size() == 2); + + INFO("Replicate the next new-view transaction at seqno 3"); + auto fresh_tx = fixture.store->create_tx(); + fresh_tx.rw(fixture.table)->put(3, 4); + CHECK(fresh_tx.commit() == ccf::kv::CommitResult::SUCCESS); + CHECK(fixture.store->current_txid() == ccf::TxID(reelection_view, 3)); + CHECK(read_value(*fixture.store, fixture.table, 3) == 4); + CHECK(fixture.raft->get_last_idx() == 3); + CHECK(fixture.raft->ledger->ledger.size() == 3); +} + +TEST_CASE( + "AFT accepts a current-view prefix before a stale suffix" * + doctest::test_suite("view_straddling_transactions")) +{ + const ccf::NodeId node_id = ccf::kv::test::PrimaryNodeId; + auto store = std::make_shared(node_id); + Raft raft( + raft_settings, + std::make_unique>(store), + std::make_unique(node_id), + std::make_shared(), + std::make_shared(node_id), + nullptr); + + ccf::kv::Configuration::Nodes configuration; + configuration.try_emplace(node_id); + raft.add_configuration(0, configuration); + raft.force_become_primary(); + + const auto initial_view = raft.get_view(); + REQUIRE(raft.replicate({make_entry({initial_view, 1}, true)})); + REQUIRE(raft.get_committed_seqno() == 1); + + raft.become_aware_of_new_term(initial_view + 1); + raft.force_become_primary(); + const auto current_view = raft.get_view(); + + INFO("Submit current 3.2 followed by stale 2.3 in one candidate batch"); + CHECK(raft.replicate( + {make_entry({current_view, 2}), make_entry({initial_view, 3})})); + CHECK(raft.get_last_idx() == 2); + CHECK(raft.ledger->ledger.size() == 2); + + INFO("The next current-view entry can reuse seqno 3"); + CHECK(raft.replicate({make_entry({current_view, 3})})); + CHECK(raft.get_last_idx() == 3); + CHECK(raft.ledger->ledger.size() == 3); +} \ No newline at end of file diff --git a/src/kv/kv_types.h b/src/kv/kv_types.h index f49c1f87ab8c..8c2bede0790a 100644 --- a/src/kv/kv_types.h +++ b/src/kv/kv_types.h @@ -440,7 +440,7 @@ namespace ccf::kv virtual void init_as_backup( ccf::SeqNo, ccf::View, const std::vector&, ccf::SeqNo) = 0; - virtual bool replicate(const BatchVector& entries) = 0; + virtual size_t replicate(const BatchVector& entries) = 0; virtual std::pair get_committed_txid() = 0; virtual ccf::View get_view(ccf::SeqNo seqno) = 0; diff --git a/src/kv/store.h b/src/kv/store.h index f6580accf231..4a6cf116a011 100644 --- a/src/kv/store.h +++ b/src/kv/store.h @@ -944,7 +944,6 @@ namespace ccf::kv BatchVector batch; Version previous_last_replicated = 0; - Version next_last_replicated = 0; Version previous_rollback_count = 0; std::vector contiguous_pending_txs; @@ -953,14 +952,38 @@ namespace ccf::kv { std::lock_guard vguard(version_lock); + if (txid.view != term_of_next_version) + { + LOG_DEBUG_FMT( + "Discarding transaction {} after Store moved to view {}", + txid.to_str(), + term_of_next_version); + return CommitResult::FAIL_NO_REPLICATE; + } + if (globally_committable && txid.seqno > last_committable) { last_committable = txid.seqno; } - pending_txs.insert( - {txid.seqno, - std::make_tuple(txid, std::move(pending_tx), globally_committable)}); + auto [it, inserted] = pending_txs.try_emplace( + txid.seqno, txid, std::move(pending_tx), globally_committable); + + if (!inserted) + { + // Extremely unexpected case: Something went very wrong with TxID + // assignment, but still fail report here rather than persisting the + // confusion + const auto& existing_txid = std::get<0>(it->second); + + LOG_FAIL_FMT( + "Conflicting pending transactions at seqno {}: {} and {}", + txid.seqno, + existing_txid.to_str(), + txid.to_str()); + + return CommitResult::FAIL_NO_REPLICATE; + } LOG_TRACE_FMT("Inserting pending tx at {}", txid.seqno); @@ -971,12 +994,11 @@ namespace ccf::kv { LOG_TRACE_FMT( "Couldn't find {} = {} + {}, giving up on batch while committing " - "{}.{}", + "{}", last_replicated + offset, last_replicated, offset, - txid.view, - txid.seqno); + txid.to_str()); break; } @@ -986,7 +1008,6 @@ namespace ccf::kv previous_rollback_count = rollback_count; previous_last_replicated = last_replicated; - next_last_replicated = last_replicated + contiguous_pending_txs.size(); } // Release version lock @@ -1020,10 +1041,9 @@ namespace ccf::kv if (success_ != CommitResult::SUCCESS) { LOG_FAIL_FMT( - "Unexpected failure reason {} during commit of {}.{}", + "Unexpected failure reason {} during commit of {}", static_cast(success_), - txid.view, - txid.seqno); + txid.to_str()); } if (h) @@ -1038,11 +1058,10 @@ namespace ccf::kv } LOG_DEBUG_FMT( - "Batching {} ({}) during commit of {}.{}", + "Batching {} ({}) during commit of {}", previous_last_replicated + offset, data_shared->size(), - txid.view, - txid.seqno); + txid.to_str()); batch.emplace_back( pending_txid_, data_shared, committable_, hooks_shared); @@ -1050,16 +1069,21 @@ namespace ccf::kv offset++; } - if (c->replicate(batch)) + const auto replicated_count = c->replicate(batch); + if (replicated_count > 0) { std::lock_guard vguard(version_lock); if ( last_replicated == previous_last_replicated && previous_rollback_count == rollback_count) { - last_replicated = next_last_replicated; + last_replicated += replicated_count; + } + + if (last_replicated >= txid.seqno) + { + return CommitResult::SUCCESS; } - return CommitResult::SUCCESS; } LOG_DEBUG_FMT("Failed to replicate"); diff --git a/src/kv/test/stub_consensus.h b/src/kv/test/stub_consensus.h index a7510aa6eb5c..579630789c55 100644 --- a/src/kv/test/stub_consensus.h +++ b/src/kv/test/stub_consensus.h @@ -105,7 +105,7 @@ namespace ccf::kv::test state = Backup; } - bool replicate(const BatchVector& entries) override + size_t replicate(const BatchVector& entries) override { for (const auto& entry : entries) { @@ -122,7 +122,7 @@ namespace ccf::kv::test } current_view = tx_id.view; } - return true; + return entries.size(); } std::optional> get_latest_data() @@ -235,9 +235,9 @@ namespace ccf::kv::test return false; } - bool replicate(const BatchVector& entries) override + size_t replicate(const BatchVector& entries) override { - return false; + return 0; } bool can_replicate() override From 12c708f4e0fe38ee6e8a141c55c8cbdd6c6351af Mon Sep 17 00:00:00 2001 From: Eddy Ashton Date: Tue, 25 Aug 2026 15:01:33 +0000 Subject: [PATCH 07/12] TODOnes --- src/consensus/aft/raft.h | 1 - src/indexing/test/common.h | 2 +- src/kv/apply_changes.h | 2 ++ src/kv/committable_tx.h | 17 +---------------- src/kv/test/kv_contention.cpp | 2 +- src/node/test/history.cpp | 20 ++++++++++++-------- tests/schema.py | 17 ++++++++--------- 7 files changed, 25 insertions(+), 36 deletions(-) diff --git a/src/consensus/aft/raft.h b/src/consensus/aft/raft.h index 57ef71a32ad5..c8333928a5fe 100644 --- a/src/consensus/aft/raft.h +++ b/src/consensus/aft/raft.h @@ -678,7 +678,6 @@ namespace aft j["function"] = "replicate"; j["state"] = *state; COMMITTABLE_INDICES(j["state"], state); - j["view"] = term; // TODO j["seqno"] = tx_id.seqno; j["globally_committable"] = globally_committable; RAFT_TRACE_JSON_OUT(j); diff --git a/src/indexing/test/common.h b/src/indexing/test/common.h index 458bc6ea75a8..fc14d0830122 100644 --- a/src/indexing/test/common.h +++ b/src/indexing/test/common.h @@ -87,7 +87,7 @@ class AllCommittableWrapper : public TConsensus public: using TConsensus::TConsensus; - bool replicate(const ccf::kv::BatchVector& entries_) override + size_t replicate(const ccf::kv::BatchVector& entries_) override { // Rather than building a history that produces real signatures, we just // overwrite the entries here to say that everything is committable diff --git a/src/kv/apply_changes.h b/src/kv/apply_changes.h index 4b1d0ff6cc38..a387b8dbb4ab 100644 --- a/src/kv/apply_changes.h +++ b/src/kv/apply_changes.h @@ -18,6 +18,8 @@ namespace ccf::kv // Atomically checks for conflicts then applies the writes in the given change // sets to their underlying Maps. Calls tx_id_resolver() at most once, iff the // writes are applied, to retrieve a unique TxID for the write set. + // Returns std::nullopt on conflict, a default TxID on successful application + // with no writes, or the assigned TxID on successful application with writes. using TxIDResolver = std::function; diff --git a/src/kv/committable_tx.h b/src/kv/committable_tx.h index a7e99a4f6e48..a8fc55769ab9 100644 --- a/src/kv/committable_tx.h +++ b/src/kv/committable_tx.h @@ -210,6 +210,7 @@ namespace ccf::kv if (applied_txid->seqno == NoVersion) { // Read-only transaction + applied_txid = pimpl->read_txid; return CommitResult::SUCCESS; } @@ -308,22 +309,6 @@ namespace ccf::kv throw std::logic_error("Transaction not yet committed"); } - if (!pimpl->read_txid.has_value()) - { - // Transaction did not get a handle on any map. - return std::nullopt; - } - - // A committed tx is read-only (i.e. no write to any map) if it was not - // assigned a version when it was committed - // TODO: This is no longer true. Just return applied_txid, if possible - if (!applied_txid.has_value() || applied_txid->seqno == NoVersion) - { - // Read-only transaction - return pimpl->read_txid; - } - - // Write transaction return applied_txid; } diff --git a/src/kv/test/kv_contention.cpp b/src/kv/test/kv_contention.cpp index ff4cae578896..690a1462b62b 100644 --- a/src/kv/test/kv_contention.cpp +++ b/src/kv/test/kv_contention.cpp @@ -22,7 +22,7 @@ class SlowStubConsensus : public ccf::kv::test::StubConsensus public: using ccf::kv::test::StubConsensus::StubConsensus; - bool replicate(const ccf::kv::BatchVector& entries) override + size_t replicate(const ccf::kv::BatchVector& entries) override { if (rand() % 2 == 0) { diff --git a/src/node/test/history.cpp b/src/node/test/history.cpp index 84f3f490fe65..70484bab0920 100644 --- a/src/node/test/history.cpp +++ b/src/node/test/history.cpp @@ -40,15 +40,19 @@ class DummyConsensus : public ccf::kv::test::StubConsensus DummyConsensus(ccf::kv::Store* store_) : store(store_) {} - bool replicate(const ccf::kv::BatchVector& entries) override + size_t replicate(const ccf::kv::BatchVector& entries) override { if (store) { REQUIRE(entries.size() == 1); - return store->deserialize(*std::get<1>(entries[0]))->apply() != - ccf::kv::ApplyResult::FAIL; + if ( + store->deserialize(*std::get<1>(entries[0]))->apply() != + ccf::kv::ApplyResult::FAIL) + { + return 1; + } } - return true; + return 0; } std::pair get_committed_txid() override @@ -256,7 +260,7 @@ class CompactingConsensus : public ccf::kv::test::StubConsensus CompactingConsensus(ccf::kv::Store* store_) : store(store_) {} - bool replicate(const ccf::kv::BatchVector& entries) override + size_t replicate(const ccf::kv::BatchVector& entries) override { for (auto& [tx_id, data, committable, hooks] : entries) { @@ -264,7 +268,7 @@ class CompactingConsensus : public ccf::kv::test::StubConsensus if (committable) store->compact(tx_id.seqno); } - return true; + return entries.size(); } std::pair get_committed_txid() override @@ -429,7 +433,7 @@ class RollbackConsensus : public ccf::kv::test::StubConsensus rollback_to(rollback_to_) {} - bool replicate(const ccf::kv::BatchVector& entries) override + size_t replicate(const ccf::kv::BatchVector& entries) override { for (auto& [tx_id, data, committable, hook] : entries) { @@ -437,7 +441,7 @@ class RollbackConsensus : public ccf::kv::test::StubConsensus if (tx_id.seqno == rollback_at) store->rollback(tx_id, store->commit_view()); } - return true; + return entries.size(); } std::pair get_committed_txid() override diff --git a/tests/schema.py b/tests/schema.py index 5a35f8856c17..331a64c5cd53 100644 --- a/tests/schema.py +++ b/tests/schema.py @@ -264,15 +264,14 @@ def add(parser): initial_member_count=1, ) - ## TODO: Too slow, temporarily disabled for faster smoke testing - # cr.add( - # "operations", - # e2e_operations.run, - # package="samples/apps/logging/logging", - # nodes=infra.e2e_args.min_nodes(cr.args, f=0), - # initial_user_count=1, - # ledger_chunk_bytes="1B", # Chunk ledger at every signature transaction - # ) + cr.add( + "operations", + e2e_operations.run, + package="samples/apps/logging/logging", + nodes=infra.e2e_args.min_nodes(cr.args, f=0), + initial_user_count=1, + ledger_chunk_bytes="1B", # Chunk ledger at every signature transaction + ) cr.add( "download", From f99c5e34a24c5399c2ad0f9b1b166006c98331b6 Mon Sep 17 00:00:00 2001 From: Eddy Ashton Date: Tue, 25 Aug 2026 15:52:58 +0000 Subject: [PATCH 08/12] Test fixup --- src/node/test/historical_queries.cpp | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/src/node/test/historical_queries.cpp b/src/node/test/historical_queries.cpp index 4c0e62ad6ef7..db0d252f0fd3 100644 --- a/src/node/test/historical_queries.cpp +++ b/src/node/test/historical_queries.cpp @@ -226,13 +226,13 @@ std::map> construct_host_ledger( auto version = std::get<0>(next_ledger_entry.value()).seqno; while (next_ledger_entry.has_value()) { - const auto next_version = std::get<0>(next_ledger_entry.value()).seqno; const auto ib = ledger.insert( - std::make_pair(next_version, *std::get<1>(next_ledger_entry.value()))); + std::make_pair(version, *std::get<1>(next_ledger_entry.value()))); REQUIRE(ib.second); next_ledger_entry = consensus->pop_oldest_entry(); if (next_ledger_entry.has_value()) { + const auto next_version = std::get<0>(next_ledger_entry.value()).seqno; REQUIRE(version + 1 == next_version); version = next_version; } From a601c010010a48ba826cfb87f78b31ec8acd29ed Mon Sep 17 00:00:00 2001 From: Eddy Ashton Date: Wed, 26 Aug 2026 14:37:45 +0000 Subject: [PATCH 09/12] Merge fixup --- src/kv/committable_tx.h | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/src/kv/committable_tx.h b/src/kv/committable_tx.h index acb01ffa8c65..2121b287d7ca 100644 --- a/src/kv/committable_tx.h +++ b/src/kv/committable_tx.h @@ -83,7 +83,9 @@ namespace ccf::kv SizeKvStoreSerialiser size_serialiser( e, - TxID{pimpl->commit_view, NoVersion}, + // Used as IV for encrypted serialisation, but does not affect the + // projected size. + ccf::TxID{0, 0}, EntryType::WriteSetWithCommitEvidenceAndClaims, entry_flags, // Both digests are fixed-size, so their values do not affect the From 252b8c8b54c866351ba7616c8d8a5b4c9c0a3bcd Mon Sep 17 00:00:00 2001 From: Eddy Ashton Date: Wed, 26 Aug 2026 14:48:07 +0000 Subject: [PATCH 10/12] Remove debug logging --- src/consensus/aft/test/view_straddling_transactions.cpp | 4 ---- 1 file changed, 4 deletions(-) diff --git a/src/consensus/aft/test/view_straddling_transactions.cpp b/src/consensus/aft/test/view_straddling_transactions.cpp index 4f5aec05ea1b..a60b13326163 100644 --- a/src/consensus/aft/test/view_straddling_transactions.cpp +++ b/src/consensus/aft/test/view_straddling_transactions.cpp @@ -363,10 +363,6 @@ TEST_CASE( REQUIRE(current_result.has_value()); CHECK(current_result.value() == ccf::kv::CommitResult::SUCCESS); - const auto a = fixture.store->current_txid(); - const auto b = ccf::TxID(reelection_view, 3); - std::cout << "Current TxID: " << a.to_str() << ", expected: " << b.to_str() - << std::endl; CHECK(fixture.store->current_txid() == ccf::TxID(reelection_view, 3)); CHECK_FALSE(read_value(*fixture.store, fixture.table, 1).has_value()); CHECK(read_value(*fixture.store, fixture.table, 2) == 3); From 76b84eb3468bb8e6d2e8d567fa249d1417151d5a Mon Sep 17 00:00:00 2001 From: Eddy Ashton Date: Wed, 26 Aug 2026 15:16:39 +0000 Subject: [PATCH 11/12] Legitimate complaint: Clean up type confusion --- src/kv/committable_tx.h | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/src/kv/committable_tx.h b/src/kv/committable_tx.h index 2121b287d7ca..91abaa56ac73 100644 --- a/src/kv/committable_tx.h +++ b/src/kv/committable_tx.h @@ -347,7 +347,7 @@ namespace ccf::kv * * @return Commit term */ - [[nodiscard]] Version commit_term() const + [[nodiscard]] Term commit_term() const { if (!committed) { From 98428b608bd6880e8f5f93d681f3ddf00c188232 Mon Sep 17 00:00:00 2001 From: Eddy Ashton Date: Wed, 26 Aug 2026 16:29:47 +0000 Subject: [PATCH 12/12] Slightly improve post-commit metdata-access semantics? --- src/kv/committable_tx.h | 19 +++++++++++++++---- src/kv/test/kv_test.cpp | 2 ++ 2 files changed, 17 insertions(+), 4 deletions(-) diff --git a/src/kv/committable_tx.h b/src/kv/committable_tx.h index 91abaa56ac73..bf50ec6d09a3 100644 --- a/src/kv/committable_tx.h +++ b/src/kv/committable_tx.h @@ -33,9 +33,11 @@ namespace ccf::kv // are available only pre-commit, some only post-commit. bool committed = false; - // Populated only after commit() has been called, and only if the - // transaction was successful. This is the version at which the transaction - // was applied to the local KV. + // The TxID at which this transaction was applied to the local KV. A + // successful transaction that never acquired a map handle has no changes + // and no TxID; its legacy commit metadata is represented by NoVersion and + // VIEW_UNKNOWN. A committed transaction with changes but no applied TxID + // was aborted. std::optional applied_txid = std::nullopt; TxFlags flags = 0; @@ -187,7 +189,6 @@ namespace ccf::kv if (all_changes.empty()) { committed = true; - applied_txid = pimpl->read_txid; return CommitResult::SUCCESS; } @@ -334,6 +335,11 @@ namespace ccf::kv if (!applied_txid.has_value()) { + if (all_changes.empty()) + { + return NoVersion; + } + throw std::logic_error("Transaction aborted"); } @@ -356,6 +362,11 @@ namespace ccf::kv if (!applied_txid.has_value()) { + if (all_changes.empty()) + { + return ccf::VIEW_UNKNOWN; + } + throw std::logic_error("Transaction aborted"); } diff --git a/src/kv/test/kv_test.cpp b/src/kv/test/kv_test.cpp index a356b968e6d1..5df739621b04 100644 --- a/src/kv/test/kv_test.cpp +++ b/src/kv/test/kv_test.cpp @@ -3027,6 +3027,8 @@ TEST_CASE("Reported TxID after commit") // Committed transaction was not assigned a TxID because it was empty REQUIRE_FALSE(tx.get_txid().has_value()); + REQUIRE_EQ(tx.commit_version(), ccf::kv::NoVersion); + REQUIRE_EQ(tx.commit_term(), ccf::VIEW_UNKNOWN); } INFO("Simple read-only tx");