diff --git a/db/kv/kv_interface.go b/db/kv/kv_interface.go index e4b9427dcfb..7a4770b5efe 100644 --- a/db/kv/kv_interface.go +++ b/db/kv/kv_interface.go @@ -515,6 +515,10 @@ type TemporalDebugDB interface { Files() []string MergeLoop(ctx context.Context) error + // StartMergeLoopBackground spawns a background merge goroutine and registers + // it with the aggregator's WaitGroup before launching it, so that Close/ + // wg.Wait cannot observe a zero count before the goroutine has started. + StartMergeLoopBackground(ctx context.Context, onErr func(error)) } // FlushConfig holds optional behaviour for TemporalMemBatch.Flush, populated by diff --git a/db/kv/remotedb/kv_remote.go b/db/kv/remotedb/kv_remote.go index ba6ca9f1d32..2b193048a87 100644 --- a/db/kv/remotedb/kv_remote.go +++ b/db/kv/remotedb/kv_remote.go @@ -199,6 +199,9 @@ func (db *DB) EnableReadAhead() kv.TemporalDebugDB { panic("not func (db *DB) DisableReadAhead() { panic("not implemented") } func (db *DB) Files() []string { panic("not implemented") } func (db *DB) MergeLoop(ctx context.Context) error { panic("not implemented") } +func (db *DB) StartMergeLoopBackground(_ context.Context, _ func(error)) { + panic("not implemented") +} func (db *DB) BeginTemporalRo(ctx context.Context) (kv.TemporalTx, error) { t, err := db.BeginRo(ctx) //nolint:gocritic if err != nil { diff --git a/db/kv/temporal/kv_temporal.go b/db/kv/temporal/kv_temporal.go index ef4e421942d..d4eb72eb72d 100644 --- a/db/kv/temporal/kv_temporal.go +++ b/db/kv/temporal/kv_temporal.go @@ -718,6 +718,10 @@ func (db *DB) MergeLoop(ctx context.Context) error { return db.stateFiles.MergeLoop(ctx) } +func (db *DB) StartMergeLoopBackground(ctx context.Context, onErr func(error)) { + db.stateFiles.StartMergeLoopBackground(ctx, onErr) +} + func (tx *Tx) DomainFiles(domain ...kv.Domain) kv.VisibleFiles { return tx.aggtx.DomainFiles(domain...) } diff --git a/db/state/aggregator.go b/db/state/aggregator.go index 1225e215f05..85e999a06da 100644 --- a/db/state/aggregator.go +++ b/db/state/aggregator.go @@ -1230,6 +1230,21 @@ func (a *Aggregator) MergeLoop(ctx context.Context) error { return a.mergeLoop(ctx) } +// StartMergeLoopBackground spawns a background goroutine that runs MergeLoop. +// It calls wg.Add(1) synchronously before launching the goroutine, so that a +// concurrent Close/wg.Wait can never observe a zero count before the merge +// goroutine has registered itself. The caller should use this instead of +// manually wrapping MergeLoop in a goroutine to avoid the wg.Add/wg.Wait race. +func (a *Aggregator) StartMergeLoopBackground(ctx context.Context, onErr func(error)) { + a.wg.Add(1) + go func() { + defer a.wg.Done() + if err := a.mergeLoop(ctx); err != nil && onErr != nil { + onErr(err) + } + }() +} + func (a *Aggregator) mergeLoop(ctx context.Context) (err error) { if dbg.NoMerge() || !a.mergingFiles.CompareAndSwap(false, true) { return nil // currently merging or merge is prohibited diff --git a/node/eth/backend.go b/node/eth/backend.go index 6416014befd..ce0ed1adc44 100644 --- a/node/eth/backend.go +++ b/node/eth/backend.go @@ -1129,11 +1129,9 @@ func New(ctx context.Context, stack *node.Node, config *ethconfig.Config, logger } if !dbg.NoBackgroundMaintenance() { - go func() { - if err := temporalDb.Debug().MergeLoop(ctx); err != nil { - logger.Error("snapashot merge loop error", "err", err) - } - }() + temporalDb.Debug().StartMergeLoopBackground(ctx, func(err error) { + logger.Error("snapashot merge loop error", "err", err) + }) } return backend, nil