From c9c44a09979de2f518575c2745e07441b65a9954 Mon Sep 17 00:00:00 2001 From: jeffyanta Date: Fri, 24 Jul 2026 14:27:27 -0400 Subject: [PATCH] Implement max batch size for fee burn txn --- ocp/worker/currency/feeburner/config.go | 13 +++++++++---- ocp/worker/currency/feeburner/worker.go | 13 ++++++++++--- ocp/worker/currency/feeburner/worker_test.go | 12 ++++++++---- 3 files changed, 27 insertions(+), 11 deletions(-) diff --git a/ocp/worker/currency/feeburner/config.go b/ocp/worker/currency/feeburner/config.go index a1af67c..80b73e1 100644 --- a/ocp/worker/currency/feeburner/config.go +++ b/ocp/worker/currency/feeburner/config.go @@ -13,11 +13,15 @@ const ( BatchSizeConfigEnvName = envConfigPrefix + "WORKER_BATCH_SIZE" defaultBatchSize = 100 + + MaxBurnsPerBatchConfigEnvName = envConfigPrefix + "MAX_BURNS_PER_BATCH" + defaultMaxBurnsPerBatch = 10 ) type conf struct { - subsidizer config.String - batchSize config.Uint64 + subsidizer config.String + batchSize config.Uint64 + maxBurnsPerBatch config.Uint64 } // ConfigProvider defines how config values are pulled @@ -27,8 +31,9 @@ type ConfigProvider func() *conf func WithEnvConfigs() ConfigProvider { return func() *conf { return &conf{ - subsidizer: env.NewStringConfig(SubsidizerConfigEnvName, defaultSubsidizer), - batchSize: env.NewUint64Config(BatchSizeConfigEnvName, defaultBatchSize), + subsidizer: env.NewStringConfig(SubsidizerConfigEnvName, defaultSubsidizer), + batchSize: env.NewUint64Config(BatchSizeConfigEnvName, defaultBatchSize), + maxBurnsPerBatch: env.NewUint64Config(MaxBurnsPerBatchConfigEnvName, defaultMaxBurnsPerBatch), } } } diff --git a/ocp/worker/currency/feeburner/worker.go b/ocp/worker/currency/feeburner/worker.go index 5a49b6e..fc0e514 100644 --- a/ocp/worker/currency/feeburner/worker.go +++ b/ocp/worker/currency/feeburner/worker.go @@ -34,6 +34,8 @@ func (p *runtime) sweep(runtimeCtx context.Context) error { defer trace.End() tracedCtx := metrics.NewContext(runtimeCtx, trace) + maxBurnsPerBatch := max(1, int(p.conf.maxBurnsPerBatch.Get(tracedCtx))) + var cursor query.Cursor for { items, err := p.data.GetAllCurrencyMetadataByState( @@ -77,7 +79,7 @@ func (p *runtime) sweep(runtimeCtx context.Context) error { targets = append(targets, target) } - for _, batch := range p.packBurnBatches(targets) { + for _, batch := range p.packBurnBatches(targets, maxBurnsPerBatch) { err := p.burnFeesForBatch(tracedCtx, batch) if err != nil { trace.OnError(err) @@ -135,11 +137,16 @@ func (p *runtime) hasFeesToBurn(ctx context.Context, record *currency.MetadataRe } // packBurnBatches greedily packs burn targets into the fewest transactions -// that fit within the transaction size limit. -func (p *runtime) packBurnBatches(targets []*burnTarget) [][]*burnTarget { +// that fit within the transaction size limit and maxBurnsPerBatch. +func (p *runtime) packBurnBatches(targets []*burnTarget, maxBurnsPerBatch int) [][]*burnTarget { var batches [][]*burnTarget var current []*burnTarget for _, target := range targets { + if len(current) >= maxBurnsPerBatch { + batches = append(batches, current) + current = nil + } + candidate := append(current, target) txn := p.makeBurnTransaction(candidate) if len(txn.Marshal()) > solana.MaxTransactionSize { diff --git a/ocp/worker/currency/feeburner/worker_test.go b/ocp/worker/currency/feeburner/worker_test.go index af0522f..de0f728 100644 --- a/ocp/worker/currency/feeburner/worker_test.go +++ b/ocp/worker/currency/feeburner/worker_test.go @@ -32,7 +32,7 @@ func TestPackBurnBatches(t *testing.T) { targets = append(targets, target) } - batches := p.packBurnBatches(targets) + batches := p.packBurnBatches(targets, defaultMaxBurnsPerBatch) var flattened []*burnTarget for _, batch := range batches { @@ -46,11 +46,15 @@ func TestPackBurnBatches(t *testing.T) { for i, batch := range batches { txn := p.makeBurnTransaction(batch) assert.LessOrEqual(t, len(txn.Marshal()), solana.MaxTransactionSize, fmt.Sprintf("batch %d exceeds size limit", i)) + assert.LessOrEqual(t, len(batch), defaultMaxBurnsPerBatch, fmt.Sprintf("batch %d exceeds max burns", i)) } - // Every batch except the last must be full: adding the next target would - // exceed the transaction size limit + // Every batch except the last must be full: it hit the max burn count, or + // adding the next target would exceed the transaction size limit for i := 0; i < len(batches)-1; i++ { + if len(batches[i]) == defaultMaxBurnsPerBatch { + continue + } overfilled := append(append([]*burnTarget{}, batches[i]...), batches[i+1][0]) txn := p.makeBurnTransaction(overfilled) assert.Greater(t, len(txn.Marshal()), solana.MaxTransactionSize, fmt.Sprintf("batch %d is not fully packed", i)) @@ -65,5 +69,5 @@ func TestPackBurnBatches_Empty(t *testing.T) { subsidizer: testutil.NewRandomAccount(t), } - assert.Empty(t, p.packBurnBatches(nil)) + assert.Empty(t, p.packBurnBatches(nil, defaultMaxBurnsPerBatch)) }