refactor(export): remove redundant trigger publication state

This commit is contained in:
Nicholas Dudfield
2026-07-16 16:37:49 +07:00
parent f5d9824e4a
commit 6b97686dde
13 changed files with 70 additions and 174 deletions

View File

@@ -174,8 +174,6 @@ async def scenario(ctx, log):
)
if event.get("origin_ledger_hash") != origin_hash:
raise AssertionError(f"Export stream origin hash mismatch: {event}")
if event.get("trigger_txid") != origin:
raise AssertionError(f"Export stream trigger mismatch: {event}")
position = int(event.get("committee_position", -1))
if position not in selected:
raise AssertionError(

View File

@@ -81,10 +81,10 @@ struct ExportLimits
// local tuning knobs and require measurement before activation; changing
// them does not change the canonical per-share format.
static constexpr std::size_t maxCanonicalExportSignatureBytes = 72;
// version + AccountID + 3 hashes + ledger sequence + committee position +
// version + AccountID + 2 hashes + ledger sequence + committee position +
// compressed public key + one-byte VL prefix + maximum signature.
static constexpr std::size_t maxSerializedExportShareBytes = 1 + 20 + 32 +
4 + 32 + 32 + 2 + 33 + 1 + maxCanonicalExportSignatureBytes;
4 + 32 + 2 + 33 + 1 + maxCanonicalExportSignatureBytes;
static constexpr std::size_t maxExportSharesPerRelay = 32;
static constexpr std::size_t maxExportShareRelayPayloadBytes =
maxSerializedExportShareBytes * maxExportSharesPerRelay;

View File

@@ -31,7 +31,6 @@ struct ExportShare
uint256 originTxn;
LedgerIndex originLedgerSeq{0};
uint256 originLedgerHash;
uint256 triggerTxn;
std::uint16_t committeePosition{0};
PublicKey signingKey;
Buffer signature;
@@ -61,7 +60,7 @@ struct ExportShare
{
return version == currentVersion && owner != beast::zero &&
!originTxn.isZero() && originLedgerSeq != 0 &&
!originLedgerHash.isZero() && !triggerTxn.isZero() &&
!originLedgerHash.isZero() &&
committeePosition < ExportLimits::maxCommitteeMembers &&
hasCanonicalSignature();
}
@@ -78,7 +77,6 @@ struct ExportShare
result.addBitString(originTxn);
result.add32(originLedgerSeq);
result.addBitString(originLedgerHash);
result.addBitString(triggerTxn);
result.add16(committeePosition);
result.addRaw(signingKey.slice());
result.addVL(signature);
@@ -88,9 +86,9 @@ struct ExportShare
uint256
wireHash() const
{
// Raw-wire suppression only. Routing context such as triggerTxn and
// committeePosition is checked against validated state before relay and
// is not authenticated by the destination multisignature.
// Raw-wire suppression only. Routing context such as committeePosition
// is checked against validated state before relay and is not
// authenticated by the destination multisignature.
auto const bytes = serialize();
return sha512Half(bytes.slice());
}
@@ -110,7 +108,6 @@ struct ExportShare
auto const originTxn = sit.get256();
auto const originLedgerSeq = sit.get32();
auto const originLedgerHash = sit.get256();
auto const triggerTxn = sit.get256();
auto const committeePosition = sit.get16();
auto const keySlice = sit.getSlice(33);
if (!publicKeyType(keySlice))
@@ -126,7 +123,6 @@ struct ExportShare
originTxn,
originLedgerSeq,
originLedgerHash,
triggerTxn,
committeePosition,
signingKey,
std::move(signature)};

View File

@@ -682,7 +682,6 @@ JSS(track); // out: PeerImp
JSS(traffic); // out: Overlay
JSS(trim); // in: get_aggregate_price
JSS(trimmed_set); // out: get_aggregate_price
JSS(trigger_txid); // out: NetworkOPs
JSS(total); // out: counters
JSS(total_bytes_recv); // out: Peers
JSS(total_bytes_sent); // out: Peers

View File

@@ -76,8 +76,7 @@ public:
ExportSigCollector collector;
auto const w = origin(1);
auto const publication = collector.reopenPublication(w, w, 9);
BEAST_EXPECT(publication.has_value());
BEAST_EXPECT(collector.registerOrigin(w, 9));
auto first = collector.beginAttributedAdmission(
w, contribution(7, keyA_, 1), 10);
BEAST_EXPECT(first.result == ExportSigCollector::BeginResult::verify);
@@ -124,14 +123,13 @@ public:
}
void
testReservationAndPublicationLifecycle()
testReservationAndOriginLifecycle()
{
testcase("verification reservations and publication reopening");
testcase("verification reservations and origin registration");
ExportSigCollector collector;
auto const w = origin(2);
auto publication = collector.reopenPublication(w, w, 19);
BEAST_EXPECT(publication.has_value());
BEAST_EXPECT(collector.registerOrigin(w, 19));
auto invalid = collector.beginAttributedAdmission(
w, contribution(3, keyA_, 10), 20);
BEAST_EXPECT(invalid.ticket.has_value());
@@ -159,38 +157,16 @@ public:
collector.admitContribution(std::move(*other.ticket), true, 20)
.result == ExportSigCollector::AdmitResult::accepted);
BEAST_EXPECT(publication.has_value());
if (!publication)
return;
BEAST_EXPECT(collector.claimPublication(*publication, 3, 2));
BEAST_EXPECT(collector.publicationGeneration(w) == 1);
BEAST_EXPECT(!collector.claimPublication(*publication, 3, 2));
// Reopening the same trigger is idempotent and does not reset slots.
auto samePublication = collector.reopenPublication(w, w, 21);
BEAST_EXPECT(samePublication.has_value());
if (!samePublication)
return;
BEAST_EXPECT(collector.publicationGeneration(w) == 1);
BEAST_EXPECT(!collector.claimPublication(*samePublication, 3, 2));
auto const nextTrigger = origin(22);
auto nextPublication = collector.reopenPublication(w, nextTrigger, 22);
BEAST_EXPECT(nextPublication.has_value());
if (!nextPublication)
return;
BEAST_EXPECT(collector.publicationGeneration(w) == 2);
BEAST_EXPECT(!collector.claimPublication(*publication, 4, 2));
BEAST_EXPECT(collector.claimPublication(*nextPublication, 4, 1));
BEAST_EXPECT(!collector.reopenPublication(w, w, 23));
BEAST_EXPECT(collector.registerOrigin(w, 21));
BEAST_EXPECT(
collector.positionStatus(w, 3) ==
ExportSigCollector::PositionStatus::unique);
BEAST_EXPECT(collector.fullUnionSnapshot().at(w).size() == 2);
collector.cleanupStale(278);
// Idempotent registration refreshes the stale-cleanup cursor.
collector.cleanupStale(277);
BEAST_EXPECT(!collector.fullUnionSnapshot().empty());
collector.cleanupStale(279);
collector.cleanupStale(278);
BEAST_EXPECT(collector.fullUnionSnapshot().empty());
}
@@ -209,7 +185,7 @@ public:
BEAST_EXPECT(
collector.beginAttributedAdmission(w, good, 1).result ==
ExportSigCollector::BeginResult::unknownOrigin);
BEAST_EXPECT(collector.reopenPublication(w, w, 1).has_value());
BEAST_EXPECT(collector.registerOrigin(w, 1));
good.position = ExportLimits::maxCommitteeMembers;
BEAST_EXPECT(
@@ -246,21 +222,40 @@ public:
collector.beginAttributedAdmission(w, contribution(6, keyB_, 7), 5)
.result == ExportSigCollector::BeginResult::verify);
auto oldToken = collector.reopenPublication(w, origin(30), 6);
BEAST_EXPECT(oldToken.has_value());
collector.clear(w);
auto newToken = collector.reopenPublication(w, origin(31), 7);
BEAST_EXPECT(newToken.has_value());
if (oldToken)
BEAST_EXPECT(!collector.claimPublication(*oldToken, 0, 1));
BEAST_EXPECT(
collector.beginAttributedAdmission(w, good, 7).result ==
ExportSigCollector::BeginResult::unknownOrigin);
BEAST_EXPECT(collector.registerOrigin(w, 7));
}
void
testOriginRegistrationBounds()
{
testcase("origin registration bounds");
ExportSigCollector collector;
BEAST_EXPECT(!collector.registerOrigin(uint256{}, 1));
BEAST_EXPECT(!collector.registerOrigin(origin(100), 0));
bool registeredAll = true;
for (std::size_t i = 0; i < ExportSigCollector::maxTrackedOrigins; ++i)
registeredAll =
collector.registerOrigin(
origin(static_cast<std::uint32_t>(i + 1'000)), 1) &&
registeredAll;
BEAST_EXPECT(registeredAll);
BEAST_EXPECT(!collector.registerOrigin(origin(99'999), 1));
BEAST_EXPECT(collector.registerOrigin(origin(1'000), 2));
}
void
run() override
{
testAdmissionAndConflict();
testReservationAndPublicationLifecycle();
testReservationAndOriginLifecycle();
testMalformedBoundaries();
testOriginRegistrationBounds();
}
};

View File

@@ -1390,7 +1390,6 @@ class ConsensusExtensions_test : public beast::unit_test::suite
ConsensusExtensions ce{enabledEnv.app(), activeNoopJournal()};
auto const origin = makeHash("on-round-start-export-latch");
auto const trigger = makeHash("on-round-start-export-trigger");
auto const pk = makeValidatorKeys().front();
std::uint8_t const sigBytes[] = {1, 2, 3};
Buffer const sig{sigBytes, sizeof(sigBytes)};
@@ -1401,8 +1400,8 @@ class ConsensusExtensions_test : public beast::unit_test::suite
BEAST_EXPECT(ce.rngEnabled());
BEAST_EXPECT(ce.exportEnabled());
BEAST_EXPECT(ce.postValidationExportSigCollector().reopenPublication(
origin, trigger, 10));
BEAST_EXPECT(
ce.postValidationExportSigCollector().registerOrigin(origin, 10));
auto admission =
ce.postValidationExportSigCollector().beginAttributedAdmission(
origin, ExportSigCollector::Contribution{0, pk, sig}, 10);
@@ -2149,7 +2148,6 @@ class ConsensusExtensions_test : public beast::unit_test::suite
makeHash("precheck-export-origin"),
10,
makeHash("precheck-export-origin-ledger"),
makeHash("precheck-export-trigger"),
3,
signer.first,
std::move(signature)};
@@ -2324,7 +2322,6 @@ class ConsensusExtensions_test : public beast::unit_test::suite
ConsensusExtensions ce{env.app(), activeNoopJournal()};
auto& collector = ce.postValidationExportSigCollector();
auto const origin = makeHash("attributed-origin");
auto const trigger = makeHash("attributed-trigger");
auto const signerA = randomKeyPair(KeyType::secp256k1).first;
auto const signerB = randomKeyPair(KeyType::secp256k1).first;
std::uint8_t const signatureABytes[] = {1, 2};
@@ -2332,7 +2329,7 @@ class ConsensusExtensions_test : public beast::unit_test::suite
Buffer const signatureA{signatureABytes, sizeof(signatureABytes)};
Buffer const signatureB{signatureBBytes, sizeof(signatureBBytes)};
BEAST_EXPECT(collector.reopenPublication(origin, trigger, 10));
BEAST_EXPECT(collector.registerOrigin(origin, 10));
auto admit = [&](ExportSigCollector::Position position,
PublicKey const& key,
Buffer const& signature) {
@@ -2424,7 +2421,7 @@ class ConsensusExtensions_test : public beast::unit_test::suite
auto const signer = randomKeyPair(KeyType::secp256k1).first;
std::uint8_t const signatureBytes[] = {1, 2, 3};
Buffer const signature{signatureBytes, sizeof(signatureBytes)};
BEAST_EXPECT(collector.reopenPublication(origin, origin, deadline));
BEAST_EXPECT(collector.registerOrigin(origin, deadline));
auto admission = collector.beginAttributedAdmission(
origin,
ExportSigCollector::Contribution{0, signer, signature},
@@ -3628,14 +3625,13 @@ class ConsensusExtensions_test : public beast::unit_test::suite
Env env{*this, envconfig(), supported_amendments(), nullptr};
ConsensusExtensions ce{env.app(), env.journal};
auto const origin = makeHash("export-disabled-clears-collector");
auto const trigger = makeHash("export-disabled-trigger");
auto const pk = makeValidatorKeys().front();
std::uint8_t const sigBytes[] = {1, 2, 3};
Buffer const sig{sigBytes, sizeof(sigBytes)};
ce.setExportEnabledThisRound(true);
BEAST_EXPECT(ce.postValidationExportSigCollector().reopenPublication(
origin, trigger, 10));
BEAST_EXPECT(
ce.postValidationExportSigCollector().registerOrigin(origin, 10));
auto admission =
ce.postValidationExportSigCollector().beginAttributedAdmission(
origin, ExportSigCollector::Contribution{0, pk, sig}, 10);
@@ -3822,7 +3818,6 @@ class ConsensusExtensions_test : public beast::unit_test::suite
origin,
originLedger->info().seq,
originLedger->info().hash,
origin,
0,
valKeys.keys->publicKey,
signature};

View File

@@ -287,7 +287,6 @@ public:
makeHash("export-origin"),
42,
makeHash("export-ledger"),
makeHash("export-trigger"),
3,
key,
signature}

View File

@@ -48,7 +48,6 @@ class ExportShareTransport_test : public beast::unit_test::suite
uint256{1},
4'200'000,
uint256{2},
uint256{3},
17,
key,
signature};

View File

@@ -37,7 +37,6 @@ class ExportShare_test : public beast::unit_test::suite
uint256{1},
4'200'000,
uint256{2},
uint256{3},
17,
key,
signature};
@@ -63,7 +62,6 @@ public:
BEAST_EXPECT(parsed->originTxn == share.originTxn);
BEAST_EXPECT(parsed->originLedgerSeq == share.originLedgerSeq);
BEAST_EXPECT(parsed->originLedgerHash == share.originLedgerHash);
BEAST_EXPECT(parsed->triggerTxn == share.triggerTxn);
BEAST_EXPECT(parsed->committeePosition == share.committeePosition);
BEAST_EXPECT(parsed->signingKey == share.signingKey);
BEAST_EXPECT(parsed->signature == share.signature);

View File

@@ -450,7 +450,6 @@ public:
uint256{1},
4'200'000,
uint256{2},
uint256{3},
17,
key,
signature};
@@ -471,7 +470,6 @@ public:
share.originLedgerSeq &&
event[jss::origin_ledger_hash] ==
to_string(share.originLedgerHash) &&
event[jss::trigger_txid] == to_string(share.triggerTxn) &&
event[jss::committee_position].asUInt() ==
share.committeePosition &&
event[jss::signing_key] ==

View File

@@ -125,8 +125,7 @@ resolveExportShare(
{
if (!validated)
return {ExportShareResolutionStatus::deferred, std::nullopt};
if (!validated->rules().enabled(featureExport) ||
share.triggerTxn != share.originTxn)
if (!validated->rules().enabled(featureExport))
return {ExportShareResolutionStatus::invalid, std::nullopt};
if (share.originLedgerSeq > validated->info().seq)
return {ExportShareResolutionStatus::deferred, std::nullopt};
@@ -245,7 +244,6 @@ withContribution(
context.originTxn,
context.originLedgerSeq,
context.originLedgerHash,
context.triggerTxn,
contribution.position,
contribution.signingKey,
contribution.signature};
@@ -465,8 +463,8 @@ ConsensusExtensions::admitExportShare(
return {
ExportShareDisposition::invalid, ExportShareCharge::invalidData};
if (!postValidationExportSigCollector_.reopenPublication(
share.originTxn, share.triggerTxn, validated->info().seq))
if (!postValidationExportSigCollector_.registerOrigin(
share.originTxn, validated->info().seq))
return {ExportShareDisposition::deferred, ExportShareCharge::none};
ExportSigCollector::Contribution contribution{
@@ -605,7 +603,6 @@ ConsensusExtensions::onValidatedLedger(
origin,
originSeq,
*originHash,
origin,
contribution.position,
contribution.signingKey,
contribution.signature},
@@ -722,7 +719,6 @@ ConsensusExtensions::onValidatedLedger(
origin,
originSeq,
*originHash,
origin,
*position,
keys.keys->publicKey,
signature};
@@ -3247,7 +3243,6 @@ ConsensusExtensions::attachExportSignatures(
origin,
originSeq,
*originHash,
origin,
contribution.position,
contribution.signingKey,
contribution.signature};

View File

@@ -9,7 +9,6 @@
#include <map>
#include <mutex>
#include <optional>
#include <set>
#include <unordered_map>
#include <utility>
#include <vector>
@@ -19,8 +18,9 @@ namespace ripple {
/** Post-validation Export contribution collector.
Contribution identity is the immutable Export origin plus the position in
that origin's immutable validator committee. Publication attempts may
reopen, but never reset admitted contributions or conflict state.
that origin's immutable validator committee. Callers register an origin
only after resolving it against validated state; registration enables
attributed admission without resetting contributions or conflict state.
*/
class ExportSigCollector
{
@@ -93,19 +93,6 @@ public:
std::optional<AdmissionTicket> ticket;
};
class PublicationToken
{
friend class ExportSigCollector;
uint256 origin_;
std::uint64_t generation_;
PublicationToken(uint256 const& origin, std::uint64_t generation)
: origin_(origin), generation_(generation)
{
}
};
struct AdmitOutcome
{
AdmitResult result;
@@ -121,7 +108,6 @@ public:
static constexpr std::size_t maxTrackedOrigins = 4096;
static constexpr std::uint32_t maxStaleLedgers = 256;
static constexpr std::uint32_t maxReservationLedgers = 1;
static constexpr std::size_t maxPublicationTriggers = 64;
static constexpr std::size_t maxSignatureBytes = 72;
private:
@@ -141,17 +127,12 @@ private:
struct OriginEntry
{
std::map<Position, PositionEntry> positions;
std::set<Position> published;
std::set<uint256> publicationTriggers;
uint256 publicationTrigger;
std::uint64_t publicationGeneration{0};
std::uint32_t lastTouchedSeq{0};
};
mutable std::mutex mutex_;
std::unordered_map<uint256, OriginEntry> origins_;
std::uint64_t nextReservation_{1};
std::uint64_t nextPublicationGeneration_{1};
static bool
sameEncoding(Contribution const& lhs, Contribution const& rhs)
@@ -200,16 +181,26 @@ private:
return result;
}
std::uint64_t
nextPublicationGeneration()
public:
/** Register an origin after resolving it against validated state. */
bool
registerOrigin(uint256 const& origin, std::uint32_t currentSeq)
{
auto const result = nextPublicationGeneration_++;
if (nextPublicationGeneration_ == 0)
nextPublicationGeneration_ = 1;
return result;
if (origin.isZero() || currentSeq == 0)
return false;
std::lock_guard lock(mutex_);
auto it = origins_.find(origin);
if (it == origins_.end())
{
if (origins_.size() >= maxTrackedOrigins)
return false;
it = origins_.emplace(origin, OriginEntry{}).first;
}
touch(it->second, currentSeq);
return true;
}
public:
/** Reserve one contribution encoding for verification.
The caller must first attribute the signing key to this exact selected
@@ -337,64 +328,6 @@ public:
AdmitResult::conflicted, std::move(prior), std::move(conflicting)};
}
/** Reopen per-attempt publication without changing contribution state. */
std::optional<PublicationToken>
reopenPublication(
uint256 const& origin,
uint256 const& trigger,
std::uint32_t currentSeq)
{
if (origin.isZero() || trigger.isZero() || currentSeq == 0)
return std::nullopt;
std::lock_guard lock(mutex_);
auto it = origins_.find(origin);
if (it == origins_.end())
{
if (origins_.size() >= maxTrackedOrigins)
return std::nullopt;
it = origins_.emplace(origin, OriginEntry{}).first;
}
auto& entry = it->second;
if (entry.publicationGeneration != 0 &&
entry.publicationTrigger == trigger)
return PublicationToken{origin, entry.publicationGeneration};
if (entry.publicationTriggers.count(trigger) != 0 ||
entry.publicationTriggers.size() >= maxPublicationTriggers)
return std::nullopt;
entry.published.clear();
entry.publicationTrigger = trigger;
entry.publicationTriggers.insert(trigger);
entry.publicationGeneration = nextPublicationGeneration();
touch(entry, currentSeq);
return PublicationToken{origin, entry.publicationGeneration};
}
/** Claim one position's publication slot in the current attempt. */
bool
claimPublication(
PublicationToken const& token,
Position position,
std::size_t maxDistinct)
{
std::lock_guard lock(mutex_);
auto it = origins_.find(token.origin_);
if (it == origins_.end() ||
it->second.publicationGeneration != token.generation_)
return false;
auto const positionIt = it->second.positions.find(position);
if (positionIt == it->second.positions.end() ||
!positionIt->second.unique || positionIt->second.conflicted)
return false;
if (it->second.published.count(position) != 0 ||
it->second.published.size() >= maxDistinct)
return false;
it->second.published.insert(position);
return true;
}
PositionStatus
positionStatus(uint256 const& origin, Position position) const
{
@@ -411,14 +344,6 @@ public:
: PositionStatus::empty;
}
std::uint64_t
publicationGeneration(uint256 const& origin) const
{
std::lock_guard lock(mutex_);
auto const it = origins_.find(origin);
return it == origins_.end() ? 0 : it->second.publicationGeneration;
}
/** Deterministic full union of every unique, verified contribution. */
UnionSnapshot
fullUnionSnapshot() const

View File

@@ -2407,7 +2407,6 @@ NetworkOPsImp::pubExportSignature(
event[jss::origin_txid] = to_string(share.originTxn);
event[jss::origin_ledger_seq] = Json::UInt(share.originLedgerSeq);
event[jss::origin_ledger_hash] = to_string(share.originLedgerHash);
event[jss::trigger_txid] = to_string(share.triggerTxn);
event[jss::committee_position] = Json::UInt(share.committeePosition);
event[jss::signing_key] = toBase58(TokenType::NodePublic, share.signingKey);
event[jss::signature] =