Skip to content
Draft
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
70 changes: 66 additions & 4 deletions runner/network/sequencer_benchmark.go
Original file line number Diff line number Diff line change
Expand Up @@ -28,20 +28,30 @@ import (
const gracefulWorkerShutdownTimeout = 90 * time.Second

type benchmarkRunController struct {
maxBlocks int
completion payloadworker.CompletionWorker
maxBlocks int
completion payloadworker.CompletionWorker
measurement payloadworker.MeasurementCompletionWorker
}

func newBenchmarkRunController(transactionWorker payloadworker.Worker, params benchtypes.RunParams) benchmarkRunController {
completion, ok := transactionWorker.(payloadworker.CompletionWorker)
if ok {
return benchmarkRunController{completion: completion}
measurement, _ := transactionWorker.(payloadworker.MeasurementCompletionWorker)
return benchmarkRunController{completion: completion, measurement: measurement}
}
return benchmarkRunController{maxBlocks: params.NumBlocks}
}

func (c benchmarkRunController) shouldStop(nextBlockIndex uint64) (bool, error) {
if c.completion != nil {
if c.measurement != nil {
select {
case <-c.measurement.MeasurementDone():
return true, c.completion.Err()
default:
return false, nil
}
}
select {
case <-c.completion.Done():
return true, c.completion.Err()
Expand All @@ -56,6 +66,10 @@ func (c benchmarkRunController) usesWorkerCompletion() bool {
return c.completion != nil
}

func (c benchmarkRunController) separatesMeasurementCompletion() bool {
return c.measurement != nil
}

type sequencerBenchmark struct {
log log.Logger
sequencerClient types.ExecutionClient
Expand Down Expand Up @@ -308,7 +322,12 @@ func (nb *sequencerBenchmark) Run(ctx context.Context, metricsCollector metrics.
blockIndex++
}

if !runController.usesWorkerCompletion() {
if runController.separatesMeasurementCompletion() {
if err := nb.settleCompletionWorker(benchmarkCtx, transactionWorker, consensusClient, pendingTxs, runController.completion); err != nil {
errChan <- err
return
}
} else if !runController.usesWorkerCompletion() {
if err := nb.settleGracefulWorkerShutdown(benchmarkCtx, transactionWorker, consensusClient, pendingTxs); err != nil {
errChan <- err
return
Expand Down Expand Up @@ -341,6 +360,49 @@ func (nb *sequencerBenchmark) Run(ctx context.Context, metricsCollector metrics.
}
}

func (nb *sequencerBenchmark) settleCompletionWorker(
ctx context.Context,
transactionWorker payloadworker.Worker,
consensusClient *consensus.SequencerConsensusClient,
pendingTxs int,
completion payloadworker.CompletionWorker,
) error {
settlementCtx, cancel := context.WithTimeout(ctx, gracefulWorkerShutdownTimeout)
defer cancel()

settlementBlock := 0
for {
select {
case <-completion.Done():
nb.log.Info("Load-test worker finished draining", "settlement_blocks", settlementBlock)
return completion.Err()
case <-settlementCtx.Done():
return errors.Wrapf(
settlementCtx.Err(),
"timed out waiting for load-test worker to drain after %d settlement blocks",
settlementBlock,
)
default:
}

var err error
_, pendingTxs, err = nb.proposeBlock(
settlementCtx,
transactionWorker,
consensusClient,
nil,
uint64(settlementBlock+1),
pendingTxs,
true,
false,
)
if err != nil {
return errors.Wrap(err, "failed to propose load-test settlement block")
}
settlementBlock++
}
}

func (nb *sequencerBenchmark) proposeBlock(
ctx context.Context,
transactionWorker payloadworker.Worker,
Expand Down
Loading
Loading