Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
13 changes: 9 additions & 4 deletions ocp/worker/currency/feeburner/config.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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),
}
}
}
13 changes: 10 additions & 3 deletions ocp/worker/currency/feeburner/worker.go
Original file line number Diff line number Diff line change
Expand Up @@ -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(
Expand Down Expand Up @@ -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)
Expand Down Expand Up @@ -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 {
Expand Down
12 changes: 8 additions & 4 deletions ocp/worker/currency/feeburner/worker_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand All @@ -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))
Expand All @@ -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))
}
Loading