Skip to content

Commit 716025f

Browse files
committed
Make Soroban metric batching code lock-free.
The previous attempt at batching Soroban metrics still had some locks and was overall more convoluted than it needs to be. With this change we simply give every thread its own metrics batch to write to, and then cleanly move the batches to the publish step for the merge and Medida export. This is slightly more plumbing, but as a result we get much more straightforward code.
1 parent 7791bbe commit 716025f

27 files changed

Lines changed: 567 additions & 598 deletions

src/ledger/InMemorySorobanState.cpp

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -488,7 +488,7 @@ InMemorySorobanState::updateState(
488488
std::vector<LedgerEntry> const& liveEntries,
489489
std::vector<LedgerKey> const& deadEntries, LedgerHeader const& lh,
490490
std::optional<SorobanNetworkConfig const> const& sorobanConfig,
491-
SorobanMetrics& metrics)
491+
SorobanMetricsRegistry& metrics)
492492
{
493493
// After initialization, we must apply every ledger in order to the
494494
// in-memory state with no gaps.
@@ -575,7 +575,7 @@ InMemorySorobanState::getSize() const
575575
}
576576

577577
void
578-
InMemorySorobanState::reportMetrics(SorobanMetrics& metrics) const
578+
InMemorySorobanState::reportMetrics(SorobanMetricsRegistry& metrics) const
579579
{
580580
metrics.mContractCodeStateSize.set_count(mContractCodeStateSize);
581581
metrics.mContractDataStateSize.set_count(mContractDataStateSize);

src/ledger/InMemorySorobanState.h

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -20,7 +20,7 @@ namespace stellar
2020
class ApplyLedgerView;
2121

2222
class InvariantManagerImpl;
23-
class SorobanMetrics;
23+
class SorobanMetricsRegistry;
2424

2525
// TTLData stores both liveUntilLedgerSeq and lastModifiedLedgerSeq for TTL
2626
// entries. This allows us to construct a LedgerEntry for TTLs without having to
@@ -396,7 +396,7 @@ class InMemorySorobanState
396396
// CONTRACT_CODE.
397397
void deleteContractCode(LedgerKey const& ledgerKey);
398398

399-
void reportMetrics(SorobanMetrics& metrics) const;
399+
void reportMetrics(SorobanMetricsRegistry& metrics) const;
400400

401401
public:
402402
InMemorySorobanState() = default;
@@ -449,7 +449,7 @@ class InMemorySorobanState
449449
std::vector<LedgerKey> const& deadEntries,
450450
LedgerHeader const& lh,
451451
std::optional<SorobanNetworkConfig const> const& sorobanConfig,
452-
SorobanMetrics& metrics);
452+
SorobanMetricsRegistry& metrics);
453453

454454
// Should only be called in manual ledger close paths.
455455
void manuallyAdvanceLedgerHeader(LedgerHeader const& lh);

src/ledger/LedgerManager.h

Lines changed: 8 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -19,7 +19,8 @@ namespace stellar
1919

2020
class LedgerCloseData;
2121
class Database;
22-
class SorobanMetrics;
22+
class SorobanMetricsRegistry;
23+
struct SorobanApplyMetrics;
2324
class InMemorySorobanState;
2425

2526
// This diagram provides a schematic of the flow of (logical) ledgers coming in
@@ -373,11 +374,11 @@ class LedgerManager
373374
// upgradeApplied should be true if a protocol or network config setting
374375
// upgrade occurred during the ledger close. If inMemorySnapshotForInvariant
375376
// is not null, this will kick off a snapshot invariant check.
376-
virtual void completeLedgerClose(uint32_t ledgerSeq,
377-
bool calledViaExternalize,
378-
LedgerCloseData const& ledgerData,
379-
ImmutableLedgerDataPtr appliedLedgerState,
380-
bool upgradeApplied) = 0;
377+
virtual void completeLedgerClose(
378+
uint32_t ledgerSeq, bool calledViaExternalize,
379+
LedgerCloseData const& ledgerData,
380+
ImmutableLedgerDataPtr appliedLedgerState, bool upgradeApplied,
381+
std::vector<SorobanApplyMetrics>&& sorobanApplyMetrics) = 0;
381382

382383
virtual void assertSetupPhase() const = 0;
383384
#ifdef BUILD_TESTS
@@ -398,7 +399,7 @@ class LedgerManager
398399

399400
virtual void manuallyAdvanceLedgerHeader(LedgerHeader const& header) = 0;
400401

401-
virtual SorobanMetrics& getSorobanMetrics() = 0;
402+
virtual SorobanMetricsRegistry& getSorobanMetrics() = 0;
402403
virtual ::rust::Box<rust_bridge::SorobanModuleCache> getModuleCache() = 0;
403404

404405
virtual ~LedgerManager()

src/ledger/LedgerManagerImpl.cpp

Lines changed: 74 additions & 40 deletions
Original file line numberDiff line numberDiff line change
@@ -958,7 +958,7 @@ LedgerManagerImpl::getSorobanInMemoryStateSizeForTesting()
958958
}
959959
#endif
960960

961-
SorobanMetrics&
961+
SorobanMetricsRegistry&
962962
LedgerManagerImpl::getSorobanMetrics()
963963
{
964964
return mApplyState.getMetrics().mSorobanMetrics;
@@ -1198,7 +1198,8 @@ LedgerManagerImpl::ApplyState::maybeRebuildModuleCache(
11981198
}
11991199

12001200
void
1201-
LedgerManagerImpl::publishSorobanMetrics()
1201+
LedgerManagerImpl::publishSorobanMetrics(
1202+
std::vector<SorobanApplyMetrics>& sorobanApplyMetricsPerThread)
12021203
{
12031204
if (!hasLastClosedSorobanNetworkConfig())
12041205
{
@@ -1235,8 +1236,13 @@ LedgerManagerImpl::publishSorobanMetrics()
12351236
conf.sorobanStateTargetSizeBytes());
12361237
m.mConfigFeeWrite1KB.set_count(conf.feeRent1KB());
12371238

1238-
// then publish the actual ledger usage
1239-
m.publishAndResetLedgerWideMetrics();
1239+
// then publish the actual ledger usage, merged into a single instance
1240+
auto& totalSorobanMetrics = sorobanApplyMetricsPerThread[0];
1241+
for (size_t i = 1; i < sorobanApplyMetricsPerThread.size(); ++i)
1242+
{
1243+
totalSorobanMetrics.merge(std::move(sorobanApplyMetricsPerThread[i]));
1244+
}
1245+
m.recordApplyMetrics(totalSorobanMetrics);
12401246
}
12411247

12421248
// called by txherder
@@ -1575,7 +1581,8 @@ void
15751581
LedgerManagerImpl::completeLedgerClose(
15761582
uint32_t ledgerSeq, bool calledViaExternalize,
15771583
LedgerCloseData const& ledgerData,
1578-
ImmutableLedgerDataPtr appliedLedgerState, bool upgradeApplied)
1584+
ImmutableLedgerDataPtr appliedLedgerState, bool upgradeApplied,
1585+
std::vector<SorobanApplyMetrics>&& sorobanApplyMetricsPerThread)
15791586
{
15801587
#ifdef BUILD_TESTS
15811588
if (mCompleteLedgerCloseOverride)
@@ -1598,7 +1605,7 @@ LedgerManagerImpl::completeLedgerClose(
15981605

15991606
// We can publish Soroban metrics at any point after advancing the LCL
16001607
// state.
1601-
publishSorobanMetrics();
1608+
publishSorobanMetrics(sorobanApplyMetricsPerThread);
16021609

16031610
// Maybe kick off publishing on complete checkpoint files
16041611
auto& hm = mApp.getHistoryManager();
@@ -1813,6 +1820,15 @@ LedgerManagerImpl::applyLedger(LedgerCloseData const& ledgerData,
18131820
#endif
18141821

18151822
TransactionResultSet txResultSet;
1823+
// Soroban apply metrics are collected by different worker threads, so
1824+
// in order to avoid synchronization we give every thread its own metrics
1825+
// container. At least one thread (the apply thread) is going to exist, so
1826+
// this starts at 1, and is extended on-demand before spawning more workers.
1827+
// The metrics set in each thread's slot should be considered to be
1828+
// arbitrary by the consumers - the metrics must commute, and in the end
1829+
// we merge all the metric containers together without worrying about their
1830+
// order or origin.
1831+
std::vector<SorobanApplyMetrics> sorobanApplyMetricsPerThread(1);
18161832
#ifdef BUILD_TESTS
18171833
if (mApp.getRunInOverlayOnlyMode())
18181834
{
@@ -1854,8 +1870,9 @@ LedgerManagerImpl::applyLedger(LedgerCloseData const& ledgerData,
18541870
// to use
18551871
auto const mutableTxResults = processFeesSeqNums(
18561872
*applicableTxSet, ltx, ledgerCloseMeta, ledgerData);
1857-
txResultSet = applyTransactions(*applicableTxSet, mutableTxResults, ltx,
1858-
ledgerCloseMeta);
1873+
txResultSet =
1874+
applyTransactions(*applicableTxSet, mutableTxResults, ltx,
1875+
ledgerCloseMeta, sorobanApplyMetricsPerThread);
18591876
}
18601877

18611878
auto ledgerSeq = ltx.loadHeader().current().ledgerSeq;
@@ -2081,15 +2098,19 @@ LedgerManagerImpl::applyLedger(LedgerCloseData const& ledgerData,
20812098
if (threadIsMain())
20822099
{
20832100
completeLedgerClose(ledgerSeq, calledViaExternalize, ledgerData,
2084-
std::move(appliedLedgerState), upgradeApplied);
2101+
std::move(appliedLedgerState), upgradeApplied,
2102+
std::move(sorobanApplyMetricsPerThread));
20852103
}
20862104
else
20872105
{
20882106
auto cb = [this, ledgerSeq, calledViaExternalize, ledgerData,
20892107
appliedLedgerState = std::move(appliedLedgerState),
2090-
upgradeApplied]() mutable {
2108+
upgradeApplied,
2109+
sorobanApplyMetricsPerThread =
2110+
std::move(sorobanApplyMetricsPerThread)]() mutable {
20912111
completeLedgerClose(ledgerSeq, calledViaExternalize, ledgerData,
2092-
std::move(appliedLedgerState), upgradeApplied);
2112+
std::move(appliedLedgerState), upgradeApplied,
2113+
std::move(sorobanApplyMetricsPerThread));
20932114
};
20942115
mApp.postOnMainThread(std::move(cb), "completeLedgerClose");
20952116
}
@@ -2604,26 +2625,18 @@ LedgerManagerImpl::applyThread(
26042625
AppConnector& app,
26052626
std::unique_ptr<ThreadParallelApplyLedgerState> threadState,
26062627
Cluster const& cluster, Config const& config, ParallelLedgerInfo ledgerInfo,
2607-
Hash sorobanBasePrngSeed)
2628+
Hash sorobanBasePrngSeed, SorobanApplyMetrics& sorobanMetrics)
26082629
{
26092630
for (auto const& txBundle : cluster)
26102631
{
2611-
// Apply timer; samples go into the thread's metrics batch and are
2612-
// published at ledger close.
2613-
std::optional<BatchedTimerScope> txTime;
2614-
if (!mApp.getConfig().DISABLE_SOROBAN_METRICS_FOR_TESTING)
2615-
{
2616-
txTime.emplace(getSorobanMetrics(),
2617-
&SorobanMetrics::ApplyMetricsBatch::mTxApplyNsecs);
2618-
}
2619-
2632+
auto applyStart = std::chrono::steady_clock::now();
26202633
Hash txSubSeed = subSha256(sorobanBasePrngSeed, txBundle.getTxNum());
26212634

26222635
threadState->flushRoTTLBumpsInTxWriteFootprint(txBundle);
26232636

26242637
auto res = txBundle.getTx()->parallelApply(
26252638
app, *threadState, config, ledgerInfo, txBundle.getResPayload(),
2626-
getSorobanMetrics(), txSubSeed, txBundle.getEffects());
2639+
sorobanMetrics, txSubSeed, txBundle.getEffects());
26272640

26282641
if (res)
26292642
{
@@ -2633,6 +2646,10 @@ LedgerManagerImpl::applyThread(
26332646
{
26342647
releaseAssert(!txBundle.getResPayload().isSuccess());
26352648
}
2649+
sorobanMetrics.mTxApplyNsecs.push_back(
2650+
std::chrono::duration_cast<std::chrono::nanoseconds>(
2651+
std::chrono::steady_clock::now() - applyStart)
2652+
.count());
26362653
}
26372654

26382655
threadState->flushRemainingRoTTLBumps();
@@ -2652,11 +2669,18 @@ LedgerManagerImpl::applySorobanStageClustersInParallel(
26522669
AppConnector& app, ApplyStage const& stage,
26532670
GlobalParallelApplyLedgerState const& globalState,
26542671
Hash const& sorobanBasePrngSeed, Config const& config,
2655-
ParallelLedgerInfo const& ledgerInfo)
2672+
ParallelLedgerInfo const& ledgerInfo,
2673+
std::vector<SorobanApplyMetrics>& sorobanApplyMetricsPerThread)
26562674
{
26572675
ZoneScoped;
26582676

26592677
DeactivateScopeGuard globalStateDeactivateGuard(globalState);
2678+
// Stages may contain a different number of clusters, so we ensure that
2679+
// there is a corresponding metrics entry for each cluster.
2680+
if (sorobanApplyMetricsPerThread.size() < stage.numClusters())
2681+
{
2682+
sorobanApplyMetricsPerThread.resize(stage.numClusters());
2683+
}
26602684

26612685
std::vector<
26622686
std::function<std::unique_ptr<ThreadParallelApplyLedgerState>()>>
@@ -2665,13 +2689,16 @@ LedgerManagerImpl::applySorobanStageClustersInParallel(
26652689
for (size_t i = 0; i < stage.numClusters(); ++i)
26662690
{
26672691
tasks.emplace_back([this, &app, &globalState, &stage, i, &config,
2668-
&ledgerInfo, &sorobanBasePrngSeed]() {
2692+
&ledgerInfo, &sorobanBasePrngSeed,
2693+
&sorobanApplyMetricsPerThread]() {
26692694
auto const& cluster = stage.getCluster(i);
26702695
auto threadStatePtr =
26712696
std::make_unique<ThreadParallelApplyLedgerState>(
26722697
app, globalState, cluster, i);
2698+
// Give every thread its own metrics entry to write to.
26732699
return applyThread(app, std::move(threadStatePtr), cluster, config,
2674-
ledgerInfo, sorobanBasePrngSeed);
2700+
ledgerInfo, sorobanBasePrngSeed,
2701+
sorobanApplyMetricsPerThread[i]);
26752702
});
26762703
}
26772704

@@ -2728,14 +2755,16 @@ void
27282755
LedgerManagerImpl::applySorobanStage(
27292756
AppConnector& app, LedgerHeader const& header,
27302757
GlobalParallelApplyLedgerState& globalParState, ApplyStage const& stage,
2731-
Hash const& sorobanBasePrngSeed)
2758+
Hash const& sorobanBasePrngSeed,
2759+
std::vector<SorobanApplyMetrics>& sorobanApplyMetricsPerThread)
27322760
{
27332761
ZoneScoped;
27342762
auto const& config = app.getConfig();
27352763
auto ledgerInfo = getParallelLedgerInfo(app, header);
27362764

27372765
auto threadStates = applySorobanStageClustersInParallel(
2738-
app, stage, globalParState, sorobanBasePrngSeed, config, ledgerInfo);
2766+
app, stage, globalParState, sorobanBasePrngSeed, config, ledgerInfo,
2767+
sorobanApplyMetricsPerThread);
27392768

27402769
if (config.invariantsEnabled())
27412770
{
@@ -2754,22 +2783,24 @@ LedgerManagerImpl::applySorobanStage(
27542783
}
27552784

27562785
void
2757-
LedgerManagerImpl::applySorobanStages(AppConnector& app, AbstractLedgerTxn& ltx,
2758-
std::vector<ApplyStage> const& stages,
2759-
SorobanNetworkConfig const& sorobanConfig,
2760-
Hash const& sorobanBasePrngSeed)
2786+
LedgerManagerImpl::applySorobanStages(
2787+
AppConnector& app, AbstractLedgerTxn& ltx,
2788+
std::vector<ApplyStage> const& stages,
2789+
SorobanNetworkConfig const& sorobanConfig, Hash const& sorobanBasePrngSeed,
2790+
std::vector<SorobanApplyMetrics>& sorobanApplyMetricsPerThread)
27612791
{
27622792
ZoneScoped;
27632793
GlobalParallelApplyLedgerState globalParState(
27642794
app, mApplyState.copyApplyLedgerView(), ltx, stages,
2765-
mApplyState.getInMemorySorobanState(), sorobanConfig);
2795+
mApplyState.getInMemorySorobanState(), sorobanConfig,
2796+
sorobanApplyMetricsPerThread[0]);
27662797
// LedgerTxn is not passed into applySorobanStage, so there's no risk
27672798
// of the header being updated while we apply the stages.
27682799
auto const& header = ltx.loadHeader().current();
27692800
for (auto const& stage : stages)
27702801
{
27712802
applySorobanStage(app, header, globalParState, stage,
2772-
sorobanBasePrngSeed);
2803+
sorobanBasePrngSeed, sorobanApplyMetricsPerThread);
27732804
}
27742805
globalParState.commitChangesToLedgerTxn(ltx);
27752806
}
@@ -2837,7 +2868,8 @@ LedgerManagerImpl::applyTransactions(
28372868
ApplicableTxSetFrame const& txSet,
28382869
std::vector<MutableTxResultPtr> const& mutableTxResults,
28392870
AbstractLedgerTxn& ltx,
2840-
std::unique_ptr<LedgerCloseMetaFrame> const& ledgerCloseMeta)
2871+
std::unique_ptr<LedgerCloseMetaFrame> const& ledgerCloseMeta,
2872+
std::vector<SorobanApplyMetrics>& sorobanApplyMetricsPerThread)
28412873
{
28422874
ZoneNamedN(txsZone, "applyTransactions", true);
28432875
size_t numTxs = txSet.sizeTxTotal();
@@ -2896,7 +2928,8 @@ LedgerManagerImpl::applyTransactions(
28962928
releaseAssert(sorobanConfig.has_value());
28972929
applyParallelPhase(phase, applyStages, mutableTxResults, index,
28982930
ltx, enableTxMeta, *sorobanConfig,
2899-
sorobanBasePrngSeed);
2931+
sorobanBasePrngSeed,
2932+
sorobanApplyMetricsPerThread);
29002933
}
29012934
catch (std::exception const& e)
29022935
{
@@ -2914,7 +2947,7 @@ LedgerManagerImpl::applyTransactions(
29142947
applySequentialPhase(phase, mutableTxResults, index, ltx,
29152948
enableTxMeta, sorobanConfig,
29162949
sorobanBasePrngSeed, ledgerCloseMeta,
2917-
txResultSet);
2950+
txResultSet, sorobanApplyMetricsPerThread[0]);
29182951
}
29192952
}
29202953

@@ -2942,7 +2975,8 @@ LedgerManagerImpl::applyParallelPhase(
29422975
TxSetPhaseFrame const& phase, std::vector<stellar::ApplyStage>& applyStages,
29432976
std::vector<stellar::MutableTxResultPtr> const& mutableTxResults,
29442977
uint32_t& index, stellar::AbstractLedgerTxn& ltx, bool enableTxMeta,
2945-
SorobanNetworkConfig const& sorobanConfig, Hash const& sorobanBasePrngSeed)
2978+
SorobanNetworkConfig const& sorobanConfig, Hash const& sorobanBasePrngSeed,
2979+
std::vector<SorobanApplyMetrics>& sorobanApplyMetricsPerThread)
29462980
{
29472981
ZoneScoped;
29482982

@@ -2991,7 +3025,7 @@ LedgerManagerImpl::applyParallelPhase(
29913025
}
29923026

29933027
applySorobanStages(mApp.getAppConnector(), ltx, applyStages, sorobanConfig,
2994-
sorobanBasePrngSeed);
3028+
sorobanBasePrngSeed, sorobanApplyMetricsPerThread);
29953029

29963030
// meta will be processed in processPostTxSetApply
29973031
}
@@ -3004,7 +3038,7 @@ LedgerManagerImpl::applySequentialPhase(
30043038
std::optional<SorobanNetworkConfig const> const& sorobanConfig,
30053039
Hash const& sorobanBasePrngSeed,
30063040
std::unique_ptr<LedgerCloseMetaFrame> const& ledgerCloseMeta,
3007-
TransactionResultSet& txResultSet)
3041+
TransactionResultSet& txResultSet, SorobanApplyMetrics& sorobanMetrics)
30083042
{
30093043
for (auto const& tx : phase)
30103044
{
@@ -3040,7 +3074,7 @@ LedgerManagerImpl::applySequentialPhase(
30403074
}
30413075

30423076
tx->apply(mApp.getAppConnector(), ltx, tm, mutableTxResult,
3043-
sorobanConfig, subSeed);
3077+
sorobanConfig, subSeed, sorobanMetrics);
30443078
tx->processPostApply(mApp.getAppConnector(), ltx, tm, mutableTxResult);
30453079

30463080
tm.maybeSetRefundableFeeMeta(mutableTxResult.getRefundableFeeTracker());

0 commit comments

Comments
 (0)