mirror of
https://github.com/lightninglabs/loop.git
synced 2026-08-13 12:33:03 +02:00
Merge pull request #889 from starius/fix-unit-test-races
sweepbatcher: fix race conditions in unit tests
This commit is contained in:
commit
800f0e0fba
8 changed files with 612 additions and 289 deletions
2
.github/workflows/main.yml
vendored
2
.github/workflows/main.yml
vendored
|
|
@ -20,7 +20,7 @@ env:
|
||||||
|
|
||||||
# If you change this value, please change it in the following files as well:
|
# If you change this value, please change it in the following files as well:
|
||||||
# /Dockerfile
|
# /Dockerfile
|
||||||
GO_VERSION: 1.21.10
|
GO_VERSION: 1.24.0
|
||||||
|
|
||||||
jobs:
|
jobs:
|
||||||
########################
|
########################
|
||||||
|
|
|
||||||
|
|
@ -1,4 +1,4 @@
|
||||||
FROM --platform=${BUILDPLATFORM} golang:1.22-alpine as builder
|
FROM --platform=${BUILDPLATFORM} golang:1.24-alpine as builder
|
||||||
|
|
||||||
# Copy in the local repository to build from.
|
# Copy in the local repository to build from.
|
||||||
COPY . /go/src/github.com/lightningnetwork/loop
|
COPY . /go/src/github.com/lightningnetwork/loop
|
||||||
|
|
|
||||||
|
|
@ -92,8 +92,8 @@ func (b *Batcher) greedyAddSweep(ctx context.Context, sweep *sweep) error {
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
log.Debugf("Batch selection algorithm returned batch id %d for"+
|
debugf("Batch selection algorithm returned batch id %d "+
|
||||||
" sweep %x, but acceptance failed.", batchId,
|
"for sweep %x, but acceptance failed.", batchId,
|
||||||
sweep.swapHash[:6])
|
sweep.swapHash[:6])
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -2,15 +2,21 @@ package sweepbatcher
|
||||||
|
|
||||||
import (
|
import (
|
||||||
"fmt"
|
"fmt"
|
||||||
|
"sync/atomic"
|
||||||
|
|
||||||
"github.com/btcsuite/btclog"
|
"github.com/btcsuite/btclog"
|
||||||
"github.com/lightningnetwork/lnd/build"
|
"github.com/lightningnetwork/lnd/build"
|
||||||
)
|
)
|
||||||
|
|
||||||
// log is a logger that is initialized with no output filters. This
|
// log_ is a logger that is initialized with no output filters. This
|
||||||
// means the package will not perform any logging by default until the
|
// means the package will not perform any logging by default until the
|
||||||
// caller requests it.
|
// caller requests it.
|
||||||
var log btclog.Logger
|
var log_ atomic.Pointer[btclog.Logger]
|
||||||
|
|
||||||
|
// log returns active logger.
|
||||||
|
func log() btclog.Logger {
|
||||||
|
return *log_.Load()
|
||||||
|
}
|
||||||
|
|
||||||
// The default amount of logging is none.
|
// The default amount of logging is none.
|
||||||
func init() {
|
func init() {
|
||||||
|
|
@ -20,12 +26,32 @@ func init() {
|
||||||
// batchPrefixLogger returns a logger that prefixes all log messages with
|
// batchPrefixLogger returns a logger that prefixes all log messages with
|
||||||
// the ID.
|
// the ID.
|
||||||
func batchPrefixLogger(batchID string) btclog.Logger {
|
func batchPrefixLogger(batchID string) btclog.Logger {
|
||||||
return build.NewPrefixLog(fmt.Sprintf("[Batch %s]", batchID), log)
|
return build.NewPrefixLog(fmt.Sprintf("[Batch %s]", batchID), log())
|
||||||
}
|
}
|
||||||
|
|
||||||
// UseLogger uses a specified Logger to output package logging info.
|
// UseLogger uses a specified Logger to output package logging info.
|
||||||
// This should be used in preference to SetLogWriter if the caller is also
|
// This should be used in preference to SetLogWriter if the caller is also
|
||||||
// using btclog.
|
// using btclog.
|
||||||
func UseLogger(logger btclog.Logger) {
|
func UseLogger(logger btclog.Logger) {
|
||||||
log = logger
|
log_.Store(&logger)
|
||||||
|
}
|
||||||
|
|
||||||
|
// debugf logs a message with level DEBUG.
|
||||||
|
func debugf(format string, params ...interface{}) {
|
||||||
|
log().Debugf(format, params...)
|
||||||
|
}
|
||||||
|
|
||||||
|
// infof logs a message with level INFO.
|
||||||
|
func infof(format string, params ...interface{}) {
|
||||||
|
log().Infof(format, params...)
|
||||||
|
}
|
||||||
|
|
||||||
|
// warnf logs a message with level WARN.
|
||||||
|
func warnf(format string, params ...interface{}) {
|
||||||
|
log().Warnf(format, params...)
|
||||||
|
}
|
||||||
|
|
||||||
|
// errorf logs a message with level ERROR.
|
||||||
|
func errorf(format string, params ...interface{}) {
|
||||||
|
log().Errorf(format, params...)
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -4,6 +4,7 @@ import (
|
||||||
"context"
|
"context"
|
||||||
"errors"
|
"errors"
|
||||||
"sort"
|
"sort"
|
||||||
|
"sync"
|
||||||
|
|
||||||
"github.com/btcsuite/btcd/btcutil"
|
"github.com/btcsuite/btcd/btcutil"
|
||||||
"github.com/lightningnetwork/lnd/lntypes"
|
"github.com/lightningnetwork/lnd/lntypes"
|
||||||
|
|
@ -13,6 +14,7 @@ import (
|
||||||
type StoreMock struct {
|
type StoreMock struct {
|
||||||
batches map[int32]dbBatch
|
batches map[int32]dbBatch
|
||||||
sweeps map[lntypes.Hash]dbSweep
|
sweeps map[lntypes.Hash]dbSweep
|
||||||
|
mu sync.Mutex
|
||||||
}
|
}
|
||||||
|
|
||||||
// NewStoreMock instantiates a new mock store.
|
// NewStoreMock instantiates a new mock store.
|
||||||
|
|
@ -28,6 +30,9 @@ func NewStoreMock() *StoreMock {
|
||||||
func (s *StoreMock) FetchUnconfirmedSweepBatches(ctx context.Context) (
|
func (s *StoreMock) FetchUnconfirmedSweepBatches(ctx context.Context) (
|
||||||
[]*dbBatch, error) {
|
[]*dbBatch, error) {
|
||||||
|
|
||||||
|
s.mu.Lock()
|
||||||
|
defer s.mu.Unlock()
|
||||||
|
|
||||||
result := []*dbBatch{}
|
result := []*dbBatch{}
|
||||||
for _, batch := range s.batches {
|
for _, batch := range s.batches {
|
||||||
batch := batch
|
batch := batch
|
||||||
|
|
@ -44,6 +49,9 @@ func (s *StoreMock) FetchUnconfirmedSweepBatches(ctx context.Context) (
|
||||||
func (s *StoreMock) InsertSweepBatch(ctx context.Context,
|
func (s *StoreMock) InsertSweepBatch(ctx context.Context,
|
||||||
batch *dbBatch) (int32, error) {
|
batch *dbBatch) (int32, error) {
|
||||||
|
|
||||||
|
s.mu.Lock()
|
||||||
|
defer s.mu.Unlock()
|
||||||
|
|
||||||
var id int32
|
var id int32
|
||||||
|
|
||||||
if len(s.batches) == 0 {
|
if len(s.batches) == 0 {
|
||||||
|
|
@ -66,12 +74,18 @@ func (s *StoreMock) DropBatch(ctx context.Context, id int32) error {
|
||||||
func (s *StoreMock) UpdateSweepBatch(ctx context.Context,
|
func (s *StoreMock) UpdateSweepBatch(ctx context.Context,
|
||||||
batch *dbBatch) error {
|
batch *dbBatch) error {
|
||||||
|
|
||||||
|
s.mu.Lock()
|
||||||
|
defer s.mu.Unlock()
|
||||||
|
|
||||||
s.batches[batch.ID] = *batch
|
s.batches[batch.ID] = *batch
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
// ConfirmBatch confirms a batch.
|
// ConfirmBatch confirms a batch.
|
||||||
func (s *StoreMock) ConfirmBatch(ctx context.Context, id int32) error {
|
func (s *StoreMock) ConfirmBatch(ctx context.Context, id int32) error {
|
||||||
|
s.mu.Lock()
|
||||||
|
defer s.mu.Unlock()
|
||||||
|
|
||||||
batch, ok := s.batches[id]
|
batch, ok := s.batches[id]
|
||||||
if !ok {
|
if !ok {
|
||||||
return errors.New("batch not found")
|
return errors.New("batch not found")
|
||||||
|
|
@ -87,6 +101,9 @@ func (s *StoreMock) ConfirmBatch(ctx context.Context, id int32) error {
|
||||||
func (s *StoreMock) FetchBatchSweeps(ctx context.Context,
|
func (s *StoreMock) FetchBatchSweeps(ctx context.Context,
|
||||||
id int32) ([]*dbSweep, error) {
|
id int32) ([]*dbSweep, error) {
|
||||||
|
|
||||||
|
s.mu.Lock()
|
||||||
|
defer s.mu.Unlock()
|
||||||
|
|
||||||
result := []*dbSweep{}
|
result := []*dbSweep{}
|
||||||
for _, sweep := range s.sweeps {
|
for _, sweep := range s.sweeps {
|
||||||
sweep := sweep
|
sweep := sweep
|
||||||
|
|
@ -104,7 +121,11 @@ func (s *StoreMock) FetchBatchSweeps(ctx context.Context,
|
||||||
|
|
||||||
// UpsertSweep inserts a sweep into the database, or updates an existing sweep.
|
// UpsertSweep inserts a sweep into the database, or updates an existing sweep.
|
||||||
func (s *StoreMock) UpsertSweep(ctx context.Context, sweep *dbSweep) error {
|
func (s *StoreMock) UpsertSweep(ctx context.Context, sweep *dbSweep) error {
|
||||||
|
s.mu.Lock()
|
||||||
|
defer s.mu.Unlock()
|
||||||
|
|
||||||
s.sweeps[sweep.SwapHash] = *sweep
|
s.sweeps[sweep.SwapHash] = *sweep
|
||||||
|
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -112,6 +133,9 @@ func (s *StoreMock) UpsertSweep(ctx context.Context, sweep *dbSweep) error {
|
||||||
func (s *StoreMock) GetSweepStatus(ctx context.Context,
|
func (s *StoreMock) GetSweepStatus(ctx context.Context,
|
||||||
swapHash lntypes.Hash) (bool, error) {
|
swapHash lntypes.Hash) (bool, error) {
|
||||||
|
|
||||||
|
s.mu.Lock()
|
||||||
|
defer s.mu.Unlock()
|
||||||
|
|
||||||
sweep, ok := s.sweeps[swapHash]
|
sweep, ok := s.sweeps[swapHash]
|
||||||
if !ok {
|
if !ok {
|
||||||
return false, nil
|
return false, nil
|
||||||
|
|
@ -127,6 +151,9 @@ func (s *StoreMock) Close() error {
|
||||||
|
|
||||||
// AssertSweepStored asserts that a sweep is stored.
|
// AssertSweepStored asserts that a sweep is stored.
|
||||||
func (s *StoreMock) AssertSweepStored(id lntypes.Hash) bool {
|
func (s *StoreMock) AssertSweepStored(id lntypes.Hash) bool {
|
||||||
|
s.mu.Lock()
|
||||||
|
defer s.mu.Unlock()
|
||||||
|
|
||||||
_, ok := s.sweeps[id]
|
_, ok := s.sweeps[id]
|
||||||
return ok
|
return ok
|
||||||
}
|
}
|
||||||
|
|
@ -135,6 +162,9 @@ func (s *StoreMock) AssertSweepStored(id lntypes.Hash) bool {
|
||||||
func (s *StoreMock) GetParentBatch(ctx context.Context, swapHash lntypes.Hash) (
|
func (s *StoreMock) GetParentBatch(ctx context.Context, swapHash lntypes.Hash) (
|
||||||
*dbBatch, error) {
|
*dbBatch, error) {
|
||||||
|
|
||||||
|
s.mu.Lock()
|
||||||
|
defer s.mu.Unlock()
|
||||||
|
|
||||||
for _, sweep := range s.sweeps {
|
for _, sweep := range s.sweeps {
|
||||||
if sweep.SwapHash == swapHash {
|
if sweep.SwapHash == swapHash {
|
||||||
batch, ok := s.batches[sweep.BatchID]
|
batch, ok := s.batches[sweep.BatchID]
|
||||||
|
|
@ -153,6 +183,9 @@ func (s *StoreMock) GetParentBatch(ctx context.Context, swapHash lntypes.Hash) (
|
||||||
func (s *StoreMock) TotalSweptAmount(ctx context.Context, batchID int32) (
|
func (s *StoreMock) TotalSweptAmount(ctx context.Context, batchID int32) (
|
||||||
btcutil.Amount, error) {
|
btcutil.Amount, error) {
|
||||||
|
|
||||||
|
s.mu.Lock()
|
||||||
|
defer s.mu.Unlock()
|
||||||
|
|
||||||
batch, ok := s.batches[batchID]
|
batch, ok := s.batches[batchID]
|
||||||
if !ok {
|
if !ok {
|
||||||
return 0, errors.New("batch not found")
|
return 0, errors.New("batch not found")
|
||||||
|
|
|
||||||
|
|
@ -9,6 +9,7 @@ import (
|
||||||
"math"
|
"math"
|
||||||
"strings"
|
"strings"
|
||||||
"sync"
|
"sync"
|
||||||
|
"sync/atomic"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
"github.com/btcsuite/btcd/blockchain"
|
"github.com/btcsuite/btcd/blockchain"
|
||||||
|
|
@ -214,6 +215,12 @@ type batch struct {
|
||||||
// reorgChan is the channel over which reorg notifications are received.
|
// reorgChan is the channel over which reorg notifications are received.
|
||||||
reorgChan chan struct{}
|
reorgChan chan struct{}
|
||||||
|
|
||||||
|
// testReqs is a channel where test requests are received.
|
||||||
|
// This is used only in unit tests! The reason to have this is to
|
||||||
|
// avoid data races in require.Eventually calls running in parallel
|
||||||
|
// to the event loop. See method testRunInEventLoop().
|
||||||
|
testReqs chan *testRequest
|
||||||
|
|
||||||
// errChan is the channel over which errors are received.
|
// errChan is the channel over which errors are received.
|
||||||
errChan chan error
|
errChan chan error
|
||||||
|
|
||||||
|
|
@ -284,8 +291,8 @@ type batch struct {
|
||||||
// cfg is the configuration for this batch.
|
// cfg is the configuration for this batch.
|
||||||
cfg *batchConfig
|
cfg *batchConfig
|
||||||
|
|
||||||
// log is the logger for this batch.
|
// log_ is the logger for this batch.
|
||||||
log btclog.Logger
|
log_ atomic.Pointer[btclog.Logger]
|
||||||
|
|
||||||
wg sync.WaitGroup
|
wg sync.WaitGroup
|
||||||
}
|
}
|
||||||
|
|
@ -351,6 +358,7 @@ func NewBatch(cfg batchConfig, bk batchKit) *batch {
|
||||||
spendChan: make(chan *chainntnfs.SpendDetail),
|
spendChan: make(chan *chainntnfs.SpendDetail),
|
||||||
confChan: make(chan *chainntnfs.TxConfirmation, 1),
|
confChan: make(chan *chainntnfs.TxConfirmation, 1),
|
||||||
reorgChan: make(chan struct{}, 1),
|
reorgChan: make(chan struct{}, 1),
|
||||||
|
testReqs: make(chan *testRequest),
|
||||||
errChan: make(chan error, 1),
|
errChan: make(chan error, 1),
|
||||||
callEnter: make(chan struct{}),
|
callEnter: make(chan struct{}),
|
||||||
callLeave: make(chan struct{}),
|
callLeave: make(chan struct{}),
|
||||||
|
|
@ -387,7 +395,7 @@ func NewBatchFromDB(cfg batchConfig, bk batchKit) (*batch, error) {
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
return &batch{
|
b := &batch{
|
||||||
id: bk.id,
|
id: bk.id,
|
||||||
state: bk.state,
|
state: bk.state,
|
||||||
primarySweepID: bk.primaryID,
|
primarySweepID: bk.primaryID,
|
||||||
|
|
@ -395,6 +403,7 @@ func NewBatchFromDB(cfg batchConfig, bk batchKit) (*batch, error) {
|
||||||
spendChan: make(chan *chainntnfs.SpendDetail),
|
spendChan: make(chan *chainntnfs.SpendDetail),
|
||||||
confChan: make(chan *chainntnfs.TxConfirmation, 1),
|
confChan: make(chan *chainntnfs.TxConfirmation, 1),
|
||||||
reorgChan: make(chan struct{}, 1),
|
reorgChan: make(chan struct{}, 1),
|
||||||
|
testReqs: make(chan *testRequest),
|
||||||
errChan: make(chan error, 1),
|
errChan: make(chan error, 1),
|
||||||
callEnter: make(chan struct{}),
|
callEnter: make(chan struct{}),
|
||||||
callLeave: make(chan struct{}),
|
callLeave: make(chan struct{}),
|
||||||
|
|
@ -412,9 +421,42 @@ func NewBatchFromDB(cfg batchConfig, bk batchKit) (*batch, error) {
|
||||||
publishErrorHandler: bk.publishErrorHandler,
|
publishErrorHandler: bk.publishErrorHandler,
|
||||||
purger: bk.purger,
|
purger: bk.purger,
|
||||||
store: bk.store,
|
store: bk.store,
|
||||||
log: bk.log,
|
|
||||||
cfg: &cfg,
|
cfg: &cfg,
|
||||||
}, nil
|
}
|
||||||
|
|
||||||
|
b.setLog(bk.log)
|
||||||
|
|
||||||
|
return b, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// log returns current logger.
|
||||||
|
func (b *batch) log() btclog.Logger {
|
||||||
|
return *b.log_.Load()
|
||||||
|
}
|
||||||
|
|
||||||
|
// setLog atomically replaces the logger.
|
||||||
|
func (b *batch) setLog(logger btclog.Logger) {
|
||||||
|
b.log_.Store(&logger)
|
||||||
|
}
|
||||||
|
|
||||||
|
// Debugf logs a message with level DEBUG.
|
||||||
|
func (b *batch) Debugf(format string, params ...interface{}) {
|
||||||
|
b.log().Debugf(format, params...)
|
||||||
|
}
|
||||||
|
|
||||||
|
// Infof logs a message with level INFO.
|
||||||
|
func (b *batch) Infof(format string, params ...interface{}) {
|
||||||
|
b.log().Infof(format, params...)
|
||||||
|
}
|
||||||
|
|
||||||
|
// Warnf logs a message with level WARN.
|
||||||
|
func (b *batch) Warnf(format string, params ...interface{}) {
|
||||||
|
b.log().Warnf(format, params...)
|
||||||
|
}
|
||||||
|
|
||||||
|
// Errorf logs a message with level ERROR.
|
||||||
|
func (b *batch) Errorf(format string, params ...interface{}) {
|
||||||
|
b.log().Errorf(format, params...)
|
||||||
}
|
}
|
||||||
|
|
||||||
// addSweep tries to add a sweep to the batch. If this is the first sweep being
|
// addSweep tries to add a sweep to the batch. If this is the first sweep being
|
||||||
|
|
@ -430,7 +472,7 @@ func (b *batch) addSweep(ctx context.Context, sweep *sweep) (bool, error) {
|
||||||
// If the provided sweep is nil, we can't proceed with any checks, so
|
// If the provided sweep is nil, we can't proceed with any checks, so
|
||||||
// we just return early.
|
// we just return early.
|
||||||
if sweep == nil {
|
if sweep == nil {
|
||||||
b.log.Infof("the sweep is nil")
|
b.Infof("the sweep is nil")
|
||||||
|
|
||||||
return false, nil
|
return false, nil
|
||||||
}
|
}
|
||||||
|
|
@ -473,7 +515,7 @@ func (b *batch) addSweep(ctx context.Context, sweep *sweep) (bool, error) {
|
||||||
// the batch, do not add another sweep to prevent the tx from becoming
|
// the batch, do not add another sweep to prevent the tx from becoming
|
||||||
// non-standard.
|
// non-standard.
|
||||||
if len(b.sweeps) >= MaxSweepsPerBatch {
|
if len(b.sweeps) >= MaxSweepsPerBatch {
|
||||||
b.log.Infof("the batch has already too many sweeps (%d >= %d)",
|
b.Infof("the batch has already too many sweeps %d >= %d",
|
||||||
len(b.sweeps), MaxSweepsPerBatch)
|
len(b.sweeps), MaxSweepsPerBatch)
|
||||||
|
|
||||||
return false, nil
|
return false, nil
|
||||||
|
|
@ -483,7 +525,7 @@ func (b *batch) addSweep(ctx context.Context, sweep *sweep) (bool, error) {
|
||||||
// arrive here after the batch got closed because of a spend. In this
|
// arrive here after the batch got closed because of a spend. In this
|
||||||
// case we cannot add the sweep to this batch.
|
// case we cannot add the sweep to this batch.
|
||||||
if b.state != Open {
|
if b.state != Open {
|
||||||
b.log.Infof("the batch state (%v) is not open", b.state)
|
b.Infof("the batch state (%v) is not open", b.state)
|
||||||
|
|
||||||
return false, nil
|
return false, nil
|
||||||
}
|
}
|
||||||
|
|
@ -493,15 +535,15 @@ func (b *batch) addSweep(ctx context.Context, sweep *sweep) (bool, error) {
|
||||||
// we cannot add this sweep to the batch.
|
// we cannot add this sweep to the batch.
|
||||||
for _, s := range b.sweeps {
|
for _, s := range b.sweeps {
|
||||||
if s.isExternalAddr {
|
if s.isExternalAddr {
|
||||||
b.log.Infof("the batch already has a sweep (%x) with "+
|
b.Infof("the batch already has a sweep %x with "+
|
||||||
"an external address", s.swapHash[:6])
|
"an external address", s.swapHash[:6])
|
||||||
|
|
||||||
return false, nil
|
return false, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
if sweep.isExternalAddr {
|
if sweep.isExternalAddr {
|
||||||
b.log.Infof("the batch is not empty and new sweep (%x)"+
|
b.Infof("the batch is not empty and new sweep %x "+
|
||||||
" has an external address", sweep.swapHash[:6])
|
"has an external address", sweep.swapHash[:6])
|
||||||
|
|
||||||
return false, nil
|
return false, nil
|
||||||
}
|
}
|
||||||
|
|
@ -515,7 +557,7 @@ func (b *batch) addSweep(ctx context.Context, sweep *sweep) (bool, error) {
|
||||||
int32(math.Abs(float64(sweep.timeout - s.timeout)))
|
int32(math.Abs(float64(sweep.timeout - s.timeout)))
|
||||||
|
|
||||||
if timeoutDistance > b.cfg.maxTimeoutDistance {
|
if timeoutDistance > b.cfg.maxTimeoutDistance {
|
||||||
b.log.Infof("too long timeout distance between the "+
|
b.Infof("too long timeout distance between the "+
|
||||||
"batch and sweep %x: %d > %d",
|
"batch and sweep %x: %d > %d",
|
||||||
sweep.swapHash[:6], timeoutDistance,
|
sweep.swapHash[:6], timeoutDistance,
|
||||||
b.cfg.maxTimeoutDistance)
|
b.cfg.maxTimeoutDistance)
|
||||||
|
|
@ -544,7 +586,7 @@ func (b *batch) addSweep(ctx context.Context, sweep *sweep) (bool, error) {
|
||||||
}
|
}
|
||||||
|
|
||||||
// Add the sweep to the batch's sweeps.
|
// Add the sweep to the batch's sweeps.
|
||||||
b.log.Infof("adding sweep %x", sweep.swapHash[:6])
|
b.Infof("adding sweep %x", sweep.swapHash[:6])
|
||||||
b.sweeps[sweep.swapHash] = *sweep
|
b.sweeps[sweep.swapHash] = *sweep
|
||||||
|
|
||||||
// Update FeeRate. Max(sweep.minFeeRate) for all the sweeps of
|
// Update FeeRate. Max(sweep.minFeeRate) for all the sweeps of
|
||||||
|
|
@ -572,7 +614,7 @@ func (b *batch) sweepExists(hash lntypes.Hash) bool {
|
||||||
|
|
||||||
// Wait waits for the batch to gracefully stop.
|
// Wait waits for the batch to gracefully stop.
|
||||||
func (b *batch) Wait() {
|
func (b *batch) Wait() {
|
||||||
b.log.Infof("Stopping")
|
b.Infof("Stopping")
|
||||||
<-b.finished
|
<-b.finished
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -613,8 +655,7 @@ func (b *batch) Run(ctx context.Context) error {
|
||||||
// Set currentHeight here, because it may be needed in monitorSpend.
|
// Set currentHeight here, because it may be needed in monitorSpend.
|
||||||
select {
|
select {
|
||||||
case b.currentHeight = <-blockChan:
|
case b.currentHeight = <-blockChan:
|
||||||
b.log.Debugf("initial height for the batch is %v",
|
b.Debugf("initial height for the batch is %v", b.currentHeight)
|
||||||
b.currentHeight)
|
|
||||||
|
|
||||||
case <-runCtx.Done():
|
case <-runCtx.Done():
|
||||||
return runCtx.Err()
|
return runCtx.Err()
|
||||||
|
|
@ -652,7 +693,7 @@ func (b *batch) Run(ctx context.Context) error {
|
||||||
// completes.
|
// completes.
|
||||||
timerChan := clock.TickAfter(b.cfg.batchPublishDelay)
|
timerChan := clock.TickAfter(b.cfg.batchPublishDelay)
|
||||||
|
|
||||||
b.log.Infof("started, primary %x, total sweeps %v",
|
b.Infof("started, primary %x, total sweeps %v",
|
||||||
b.primarySweepID[0:6], len(b.sweeps))
|
b.primarySweepID[0:6], len(b.sweeps))
|
||||||
|
|
||||||
for {
|
for {
|
||||||
|
|
@ -662,7 +703,7 @@ func (b *batch) Run(ctx context.Context) error {
|
||||||
|
|
||||||
// blockChan provides immediately the current tip.
|
// blockChan provides immediately the current tip.
|
||||||
case height := <-blockChan:
|
case height := <-blockChan:
|
||||||
b.log.Debugf("received block %v", height)
|
b.Debugf("received block %v", height)
|
||||||
|
|
||||||
// Set the timer to publish the batch transaction after
|
// Set the timer to publish the batch transaction after
|
||||||
// the configured delay.
|
// the configured delay.
|
||||||
|
|
@ -670,7 +711,7 @@ func (b *batch) Run(ctx context.Context) error {
|
||||||
b.currentHeight = height
|
b.currentHeight = height
|
||||||
|
|
||||||
case <-initialDelayChan:
|
case <-initialDelayChan:
|
||||||
b.log.Debugf("initial delay of duration %v has ended",
|
b.Debugf("initial delay of duration %v has ended",
|
||||||
b.cfg.initialDelay)
|
b.cfg.initialDelay)
|
||||||
|
|
||||||
// Set the timer to publish the batch transaction after
|
// Set the timer to publish the batch transaction after
|
||||||
|
|
@ -680,8 +721,8 @@ func (b *batch) Run(ctx context.Context) error {
|
||||||
case <-timerChan:
|
case <-timerChan:
|
||||||
// Check that batch is still open.
|
// Check that batch is still open.
|
||||||
if b.state != Open {
|
if b.state != Open {
|
||||||
b.log.Debugf("Skipping publishing, because the"+
|
b.Debugf("Skipping publishing, because "+
|
||||||
" batch is not open (%v).", b.state)
|
"the batch is not open (%v).", b.state)
|
||||||
continue
|
continue
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -695,7 +736,7 @@ func (b *batch) Run(ctx context.Context) error {
|
||||||
// initialDelayChan has just fired, this check passes.
|
// initialDelayChan has just fired, this check passes.
|
||||||
now := clock.Now()
|
now := clock.Now()
|
||||||
if skipBefore.After(now) {
|
if skipBefore.After(now) {
|
||||||
b.log.Debugf(stillWaitingMsg, skipBefore, now)
|
b.Debugf(stillWaitingMsg, skipBefore, now)
|
||||||
continue
|
continue
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -715,14 +756,18 @@ func (b *batch) Run(ctx context.Context) error {
|
||||||
|
|
||||||
case <-b.reorgChan:
|
case <-b.reorgChan:
|
||||||
b.state = Open
|
b.state = Open
|
||||||
b.log.Warnf("reorg detected, batch is able to accept " +
|
b.Warnf("reorg detected, batch is able to " +
|
||||||
"new sweeps")
|
"accept new sweeps")
|
||||||
|
|
||||||
err := b.monitorSpend(ctx, b.sweeps[b.primarySweepID])
|
err := b.monitorSpend(ctx, b.sweeps[b.primarySweepID])
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
|
|
||||||
|
case testReq := <-b.testReqs:
|
||||||
|
testReq.handler()
|
||||||
|
close(testReq.quit)
|
||||||
|
|
||||||
case err := <-blockErrChan:
|
case err := <-blockErrChan:
|
||||||
return err
|
return err
|
||||||
|
|
||||||
|
|
@ -735,6 +780,36 @@ func (b *batch) Run(ctx context.Context) error {
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// testRunInEventLoop runs a function in the event loop blocking until
|
||||||
|
// the function returns. For unit tests only!
|
||||||
|
func (b *batch) testRunInEventLoop(ctx context.Context, handler func()) {
|
||||||
|
// If the event loop is finished, run the function.
|
||||||
|
select {
|
||||||
|
case <-b.stopping:
|
||||||
|
handler()
|
||||||
|
|
||||||
|
return
|
||||||
|
default:
|
||||||
|
}
|
||||||
|
|
||||||
|
quit := make(chan struct{})
|
||||||
|
req := &testRequest{
|
||||||
|
handler: handler,
|
||||||
|
quit: quit,
|
||||||
|
}
|
||||||
|
|
||||||
|
select {
|
||||||
|
case b.testReqs <- req:
|
||||||
|
case <-ctx.Done():
|
||||||
|
return
|
||||||
|
}
|
||||||
|
|
||||||
|
select {
|
||||||
|
case <-quit:
|
||||||
|
case <-ctx.Done():
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
// timeout returns minimum timeout as block height among sweeps of the batch.
|
// timeout returns minimum timeout as block height among sweeps of the batch.
|
||||||
// If the batch is empty, return -1.
|
// If the batch is empty, return -1.
|
||||||
func (b *batch) timeout() int32 {
|
func (b *batch) timeout() int32 {
|
||||||
|
|
@ -755,8 +830,10 @@ func (b *batch) timeout() int32 {
|
||||||
func (b *batch) isUrgent(skipBefore time.Time) bool {
|
func (b *batch) isUrgent(skipBefore time.Time) bool {
|
||||||
timeout := b.timeout()
|
timeout := b.timeout()
|
||||||
if timeout <= 0 {
|
if timeout <= 0 {
|
||||||
b.log.Warnf("Method timeout() returned %v. Number of"+
|
// This may happen if the batch is empty or if SweepInfo.Timeout
|
||||||
" sweeps: %d. It may be an empty batch.",
|
// is not set, may be possible in tests or if there is a bug.
|
||||||
|
b.Warnf("Method timeout() returned %v. Number of "+
|
||||||
|
"sweeps: %d. It may be an empty batch.",
|
||||||
timeout, len(b.sweeps))
|
timeout, len(b.sweeps))
|
||||||
return false
|
return false
|
||||||
}
|
}
|
||||||
|
|
@ -779,7 +856,7 @@ func (b *batch) isUrgent(skipBefore time.Time) bool {
|
||||||
return false
|
return false
|
||||||
}
|
}
|
||||||
|
|
||||||
b.log.Debugf("cancelling waiting for urgent sweep (timeBank is %v, "+
|
b.Debugf("cancelling waiting for urgent sweep (timeBank is %v, "+
|
||||||
"remainingWaiting is %v)", timeBank, remainingWaiting)
|
"remainingWaiting is %v)", timeBank, remainingWaiting)
|
||||||
|
|
||||||
// Signal to the caller to cancel initialDelay.
|
// Signal to the caller to cancel initialDelay.
|
||||||
|
|
@ -795,7 +872,7 @@ func (b *batch) publish(ctx context.Context) error {
|
||||||
)
|
)
|
||||||
|
|
||||||
if len(b.sweeps) == 0 {
|
if len(b.sweeps) == 0 {
|
||||||
b.log.Debugf("skipping publish: no sweeps in the batch")
|
b.Debugf("skipping publish: no sweeps in the batch")
|
||||||
|
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
@ -808,7 +885,7 @@ func (b *batch) publish(ctx context.Context) error {
|
||||||
|
|
||||||
// logPublishError is a function which logs publish errors.
|
// logPublishError is a function which logs publish errors.
|
||||||
logPublishError := func(errMsg string, err error) {
|
logPublishError := func(errMsg string, err error) {
|
||||||
b.publishErrorHandler(err, errMsg, b.log)
|
b.publishErrorHandler(err, errMsg, b.log())
|
||||||
}
|
}
|
||||||
|
|
||||||
fee, err, signSuccess = b.publishMixedBatch(ctx)
|
fee, err, signSuccess = b.publishMixedBatch(ctx)
|
||||||
|
|
@ -830,9 +907,9 @@ func (b *batch) publish(ctx context.Context) error {
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
b.log.Infof("published, total sweeps: %v, fees: %v", len(b.sweeps), fee)
|
b.Infof("published, total sweeps: %v, fees: %v", len(b.sweeps), fee)
|
||||||
for _, sweep := range b.sweeps {
|
for _, sweep := range b.sweeps {
|
||||||
b.log.Infof("published sweep %x, value: %v",
|
b.Infof("published sweep %x, value: %v",
|
||||||
sweep.swapHash[:6], sweep.value)
|
sweep.swapHash[:6], sweep.value)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -1026,7 +1103,7 @@ func (b *batch) publishMixedBatch(ctx context.Context) (btcutil.Amount, error,
|
||||||
coopInputs int
|
coopInputs int
|
||||||
)
|
)
|
||||||
for attempt := 1; ; attempt++ {
|
for attempt := 1; ; attempt++ {
|
||||||
b.log.Infof("Attempt %d of collecting cooperative signatures.",
|
b.Infof("Attempt %d of collecting cooperative signatures.",
|
||||||
attempt)
|
attempt)
|
||||||
|
|
||||||
// Construct unsigned batch transaction.
|
// Construct unsigned batch transaction.
|
||||||
|
|
@ -1062,7 +1139,7 @@ func (b *batch) publishMixedBatch(ctx context.Context) (btcutil.Amount, error,
|
||||||
ctx, i, sweep, tx, prevOutsMap, psbtBytes,
|
ctx, i, sweep, tx, prevOutsMap, psbtBytes,
|
||||||
)
|
)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
b.log.Infof("cooperative signing failed for "+
|
b.Infof("cooperative signing failed for "+
|
||||||
"sweep %x: %v", sweep.swapHash[:6], err)
|
"sweep %x: %v", sweep.swapHash[:6], err)
|
||||||
|
|
||||||
// Set coopFailed flag for this sweep in all the
|
// Set coopFailed flag for this sweep in all the
|
||||||
|
|
@ -1201,7 +1278,7 @@ func (b *batch) publishMixedBatch(ctx context.Context) (btcutil.Amount, error,
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
txHash := tx.TxHash()
|
txHash := tx.TxHash()
|
||||||
b.log.Infof("attempting to publish batch tx=%v with feerate=%v, "+
|
b.Infof("attempting to publish batch tx=%v with feerate=%v, "+
|
||||||
"weight=%v, feeForWeight=%v, fee=%v, sweeps=%d, "+
|
"weight=%v, feeForWeight=%v, fee=%v, sweeps=%d, "+
|
||||||
"%d cooperative: (%s) and %d non-cooperative (%s), destAddr=%s",
|
"%d cooperative: (%s) and %d non-cooperative (%s), destAddr=%s",
|
||||||
txHash, b.rbfCache.FeeRate, weight, feeForWeight, fee,
|
txHash, b.rbfCache.FeeRate, weight, feeForWeight, fee,
|
||||||
|
|
@ -1215,7 +1292,7 @@ func (b *batch) publishMixedBatch(ctx context.Context) (btcutil.Amount, error,
|
||||||
blockchain.GetTransactionWeight(btcutil.NewTx(tx)),
|
blockchain.GetTransactionWeight(btcutil.NewTx(tx)),
|
||||||
)
|
)
|
||||||
if realWeight != weight {
|
if realWeight != weight {
|
||||||
b.log.Warnf("actual weight of tx %v is %v, estimated as %d",
|
b.Warnf("actual weight of tx %v is %v, estimated as %d",
|
||||||
txHash, realWeight, weight)
|
txHash, realWeight, weight)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -1239,11 +1316,11 @@ func (b *batch) debugLogTx(msg string, tx *wire.MsgTx) {
|
||||||
// Serialize the transaction and convert to hex string.
|
// Serialize the transaction and convert to hex string.
|
||||||
buf := bytes.NewBuffer(make([]byte, 0, tx.SerializeSize()))
|
buf := bytes.NewBuffer(make([]byte, 0, tx.SerializeSize()))
|
||||||
if err := tx.Serialize(buf); err != nil {
|
if err := tx.Serialize(buf); err != nil {
|
||||||
b.log.Errorf("failed to serialize tx for debug log: %v", err)
|
b.Errorf("failed to serialize tx for debug log: %v", err)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
b.log.Debugf("%s: %s", msg, hex.EncodeToString(buf.Bytes()))
|
b.Debugf("%s: %s", msg, hex.EncodeToString(buf.Bytes()))
|
||||||
}
|
}
|
||||||
|
|
||||||
// musig2sign signs one sweep using musig2.
|
// musig2sign signs one sweep using musig2.
|
||||||
|
|
@ -1405,15 +1482,16 @@ func (b *batch) updateRbfRate(ctx context.Context) error {
|
||||||
if b.rbfCache.FeeRate == 0 {
|
if b.rbfCache.FeeRate == 0 {
|
||||||
// We set minFeeRate in each sweep, so fee rate is expected to
|
// We set minFeeRate in each sweep, so fee rate is expected to
|
||||||
// be initiated here.
|
// be initiated here.
|
||||||
b.log.Warnf("rbfCache.FeeRate is 0, which must not happen.")
|
b.Warnf("rbfCache.FeeRate is 0, which must not happen.")
|
||||||
|
|
||||||
if b.cfg.batchConfTarget == 0 {
|
if b.cfg.batchConfTarget == 0 {
|
||||||
b.log.Warnf("updateRbfRate called with zero " +
|
b.Warnf("updateRbfRate called with zero " +
|
||||||
"batchConfTarget")
|
"batchConfTarget")
|
||||||
}
|
}
|
||||||
|
|
||||||
b.log.Infof("initializing rbf fee rate for conf target=%v",
|
b.Infof("initializing rbf fee rate for conf target=%v",
|
||||||
b.cfg.batchConfTarget)
|
b.cfg.batchConfTarget)
|
||||||
|
|
||||||
rate, err := b.wallet.EstimateFeeRate(
|
rate, err := b.wallet.EstimateFeeRate(
|
||||||
ctx, b.cfg.batchConfTarget,
|
ctx, b.cfg.batchConfTarget,
|
||||||
)
|
)
|
||||||
|
|
@ -1453,6 +1531,7 @@ func (b *batch) monitorSpend(ctx context.Context, primarySweep sweep) error {
|
||||||
)
|
)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
cancel()
|
cancel()
|
||||||
|
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -1461,7 +1540,7 @@ func (b *batch) monitorSpend(ctx context.Context, primarySweep sweep) error {
|
||||||
defer cancel()
|
defer cancel()
|
||||||
defer b.wg.Done()
|
defer b.wg.Done()
|
||||||
|
|
||||||
b.log.Infof("monitoring spend for outpoint %s",
|
b.Infof("monitoring spend for outpoint %s",
|
||||||
primarySweep.outpoint.String())
|
primarySweep.outpoint.String())
|
||||||
|
|
||||||
for {
|
for {
|
||||||
|
|
@ -1584,7 +1663,7 @@ func (b *batch) handleSpend(ctx context.Context, spendTx *wire.MsgTx) error {
|
||||||
if len(spendTx.TxOut) > 0 {
|
if len(spendTx.TxOut) > 0 {
|
||||||
b.batchPkScript = spendTx.TxOut[0].PkScript
|
b.batchPkScript = spendTx.TxOut[0].PkScript
|
||||||
} else {
|
} else {
|
||||||
b.log.Warnf("transaction %v has no outputs", txHash)
|
b.Warnf("transaction %v has no outputs", txHash)
|
||||||
}
|
}
|
||||||
|
|
||||||
// As a previous version of the batch transaction may get confirmed,
|
// As a previous version of the batch transaction may get confirmed,
|
||||||
|
|
@ -1666,13 +1745,13 @@ func (b *batch) handleSpend(ctx context.Context, spendTx *wire.MsgTx) error {
|
||||||
|
|
||||||
err := b.purger(&sweep)
|
err := b.purger(&sweep)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
b.log.Errorf("unable to purge sweep %x: %v",
|
b.Errorf("unable to purge sweep %x: %v",
|
||||||
sweep.SwapHash[:6], err)
|
sweep.SwapHash[:6], err)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}()
|
}()
|
||||||
|
|
||||||
b.log.Infof("spent, total sweeps: %v, purged sweeps: %v",
|
b.Infof("spent, total sweeps: %v, purged sweeps: %v",
|
||||||
len(notifyList), len(purgeList))
|
len(notifyList), len(purgeList))
|
||||||
|
|
||||||
err := b.monitorConfirmations(ctx)
|
err := b.monitorConfirmations(ctx)
|
||||||
|
|
@ -1690,7 +1769,7 @@ func (b *batch) handleSpend(ctx context.Context, spendTx *wire.MsgTx) error {
|
||||||
// handleConf handles a confirmation notification. This is the final step of the
|
// handleConf handles a confirmation notification. This is the final step of the
|
||||||
// batch. Here we signal to the batcher that this batch was completed.
|
// batch. Here we signal to the batcher that this batch was completed.
|
||||||
func (b *batch) handleConf(ctx context.Context) error {
|
func (b *batch) handleConf(ctx context.Context) error {
|
||||||
b.log.Infof("confirmed in txid %s", b.batchTxid)
|
b.Infof("confirmed in txid %s", b.batchTxid)
|
||||||
b.state = Confirmed
|
b.state = Confirmed
|
||||||
|
|
||||||
return b.store.ConfirmBatch(ctx, b.id)
|
return b.store.ConfirmBatch(ctx, b.id)
|
||||||
|
|
@ -1769,7 +1848,7 @@ func (b *batch) insertAndAcquireID(ctx context.Context) (int32, error) {
|
||||||
}
|
}
|
||||||
|
|
||||||
b.id = id
|
b.id = id
|
||||||
b.log = batchPrefixLogger(fmt.Sprintf("%d", b.id))
|
b.setLog(batchPrefixLogger(fmt.Sprintf("%d", b.id)))
|
||||||
|
|
||||||
return id, nil
|
return id, nil
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -225,6 +225,16 @@ var (
|
||||||
ErrBatcherShuttingDown = errors.New("batcher shutting down")
|
ErrBatcherShuttingDown = errors.New("batcher shutting down")
|
||||||
)
|
)
|
||||||
|
|
||||||
|
// testRequest is a function passed to an event loop and a channel used to
|
||||||
|
// wait until the function is executed. This is used in unit tests only!
|
||||||
|
type testRequest struct {
|
||||||
|
// handler is the function to an event loop.
|
||||||
|
handler func()
|
||||||
|
|
||||||
|
// quit is closed when the handler completes.
|
||||||
|
quit chan struct{}
|
||||||
|
}
|
||||||
|
|
||||||
// Batcher is a system that is responsible for accepting sweep requests and
|
// Batcher is a system that is responsible for accepting sweep requests and
|
||||||
// placing them in appropriate batches. It will spin up new batches as needed.
|
// placing them in appropriate batches. It will spin up new batches as needed.
|
||||||
type Batcher struct {
|
type Batcher struct {
|
||||||
|
|
@ -234,6 +244,12 @@ type Batcher struct {
|
||||||
// sweepReqs is a channel where sweep requests are received.
|
// sweepReqs is a channel where sweep requests are received.
|
||||||
sweepReqs chan SweepRequest
|
sweepReqs chan SweepRequest
|
||||||
|
|
||||||
|
// testReqs is a channel where test requests are received.
|
||||||
|
// This is used only in unit tests! The reason to have this is to
|
||||||
|
// avoid data races in require.Eventually calls running in parallel
|
||||||
|
// to the event loop. See method testRunInEventLoop().
|
||||||
|
testReqs chan *testRequest
|
||||||
|
|
||||||
// errChan is a channel where errors are received.
|
// errChan is a channel where errors are received.
|
||||||
errChan chan error
|
errChan chan error
|
||||||
|
|
||||||
|
|
@ -461,6 +477,7 @@ func NewBatcher(wallet lndclient.WalletKitClient,
|
||||||
return &Batcher{
|
return &Batcher{
|
||||||
batches: make(map[int32]*batch),
|
batches: make(map[int32]*batch),
|
||||||
sweepReqs: make(chan SweepRequest),
|
sweepReqs: make(chan SweepRequest),
|
||||||
|
testReqs: make(chan *testRequest),
|
||||||
errChan: make(chan error, 1),
|
errChan: make(chan error, 1),
|
||||||
quit: make(chan struct{}),
|
quit: make(chan struct{}),
|
||||||
initDone: make(chan struct{}),
|
initDone: make(chan struct{}),
|
||||||
|
|
@ -518,22 +535,30 @@ func (b *Batcher) Run(ctx context.Context) error {
|
||||||
case sweepReq := <-b.sweepReqs:
|
case sweepReq := <-b.sweepReqs:
|
||||||
sweep, err := b.fetchSweep(runCtx, sweepReq)
|
sweep, err := b.fetchSweep(runCtx, sweepReq)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
log.Warnf("fetchSweep failed: %v.", err)
|
warnf("fetchSweep failed: %v.", err)
|
||||||
|
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
|
|
||||||
err = b.handleSweep(runCtx, sweep, sweepReq.Notifier)
|
err = b.handleSweep(runCtx, sweep, sweepReq.Notifier)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
log.Warnf("handleSweep failed: %v.", err)
|
warnf("handleSweep failed: %v.", err)
|
||||||
|
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
|
|
||||||
|
case testReq := <-b.testReqs:
|
||||||
|
testReq.handler()
|
||||||
|
close(testReq.quit)
|
||||||
|
|
||||||
case err := <-b.errChan:
|
case err := <-b.errChan:
|
||||||
log.Warnf("Batcher received an error: %v.", err)
|
warnf("Batcher received an error: %v.", err)
|
||||||
|
|
||||||
return err
|
return err
|
||||||
|
|
||||||
case <-runCtx.Done():
|
case <-runCtx.Done():
|
||||||
log.Infof("Stopping Batcher: run context cancelled.")
|
infof("Stopping Batcher: run context cancelled.")
|
||||||
|
|
||||||
return runCtx.Err()
|
return runCtx.Err()
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
@ -551,6 +576,36 @@ func (b *Batcher) AddSweep(sweepReq *SweepRequest) error {
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// testRunInEventLoop runs a function in the event loop blocking until
|
||||||
|
// the function returns. For unit tests only!
|
||||||
|
func (b *Batcher) testRunInEventLoop(ctx context.Context, handler func()) {
|
||||||
|
// If the event loop is finished, run the function.
|
||||||
|
select {
|
||||||
|
case <-b.quit:
|
||||||
|
handler()
|
||||||
|
|
||||||
|
return
|
||||||
|
default:
|
||||||
|
}
|
||||||
|
|
||||||
|
quit := make(chan struct{})
|
||||||
|
req := &testRequest{
|
||||||
|
handler: handler,
|
||||||
|
quit: quit,
|
||||||
|
}
|
||||||
|
|
||||||
|
select {
|
||||||
|
case b.testReqs <- req:
|
||||||
|
case <-ctx.Done():
|
||||||
|
return
|
||||||
|
}
|
||||||
|
|
||||||
|
select {
|
||||||
|
case <-quit:
|
||||||
|
case <-ctx.Done():
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
// handleSweep handles a sweep request by either placing it in an existing
|
// handleSweep handles a sweep request by either placing it in an existing
|
||||||
// batch, or by spinning up a new batch for it.
|
// batch, or by spinning up a new batch for it.
|
||||||
func (b *Batcher) handleSweep(ctx context.Context, sweep *sweep,
|
func (b *Batcher) handleSweep(ctx context.Context, sweep *sweep,
|
||||||
|
|
@ -561,8 +616,8 @@ func (b *Batcher) handleSweep(ctx context.Context, sweep *sweep,
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
|
|
||||||
log.Infof("Batcher handling sweep %x, completed=%v", sweep.swapHash[:6],
|
infof("Batcher handling sweep %x, completed=%v",
|
||||||
completed)
|
sweep.swapHash[:6], completed)
|
||||||
|
|
||||||
// If the sweep has already been completed in a confirmed batch then we
|
// If the sweep has already been completed in a confirmed batch then we
|
||||||
// can't attach its notifier to the batch as that is no longer running.
|
// can't attach its notifier to the batch as that is no longer running.
|
||||||
|
|
@ -573,8 +628,8 @@ func (b *Batcher) handleSweep(ctx context.Context, sweep *sweep,
|
||||||
// on-chain confirmations to prevent issues caused by reorgs.
|
// on-chain confirmations to prevent issues caused by reorgs.
|
||||||
parentBatch, err := b.store.GetParentBatch(ctx, sweep.swapHash)
|
parentBatch, err := b.store.GetParentBatch(ctx, sweep.swapHash)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
log.Errorf("unable to get parent batch for sweep %x: "+
|
errorf("unable to get parent batch for sweep %x:"+
|
||||||
"%v", sweep.swapHash[:6], err)
|
" %v", sweep.swapHash[:6], err)
|
||||||
|
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
|
|
@ -590,16 +645,17 @@ func (b *Batcher) handleSweep(ctx context.Context, sweep *sweep,
|
||||||
|
|
||||||
sweep.notifier = notifier
|
sweep.notifier = notifier
|
||||||
|
|
||||||
|
// This is a check to see if a batch is completed. In that case we just
|
||||||
|
// lazily delete it.
|
||||||
|
for _, batch := range b.batches {
|
||||||
|
if batch.isComplete() {
|
||||||
|
delete(b.batches, batch.id)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
// Check if the sweep is already in a batch. If that is the case, we
|
// Check if the sweep is already in a batch. If that is the case, we
|
||||||
// provide the sweep to that batch and return.
|
// provide the sweep to that batch and return.
|
||||||
for _, batch := range b.batches {
|
for _, batch := range b.batches {
|
||||||
// This is a check to see if a batch is completed. In that case
|
|
||||||
// we just lazily delete it and continue our scan.
|
|
||||||
if batch.isComplete() {
|
|
||||||
delete(b.batches, batch.id)
|
|
||||||
continue
|
|
||||||
}
|
|
||||||
|
|
||||||
if batch.sweepExists(sweep.swapHash) {
|
if batch.sweepExists(sweep.swapHash) {
|
||||||
accepted, err := batch.addSweep(ctx, sweep)
|
accepted, err := batch.addSweep(ctx, sweep)
|
||||||
if err != nil && !errors.Is(err, ErrBatchShuttingDown) {
|
if err != nil && !errors.Is(err, ErrBatchShuttingDown) {
|
||||||
|
|
@ -624,8 +680,8 @@ func (b *Batcher) handleSweep(ctx context.Context, sweep *sweep,
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
log.Warnf("Greedy batch selection algorithm failed for sweep %x: %v. "+
|
warnf("Greedy batch selection algorithm failed for sweep %x: %v."+
|
||||||
"Falling back to old approach.", sweep.swapHash[:6], err)
|
" Falling back to old approach.", sweep.swapHash[:6], err)
|
||||||
|
|
||||||
// If one of the batches accepts the sweep, we provide it to that batch.
|
// If one of the batches accepts the sweep, we provide it to that batch.
|
||||||
for _, batch := range b.batches {
|
for _, batch := range b.batches {
|
||||||
|
|
@ -730,13 +786,13 @@ func (b *Batcher) spinUpBatchFromDB(ctx context.Context, batch *batch) error {
|
||||||
}
|
}
|
||||||
|
|
||||||
if len(dbSweeps) == 0 {
|
if len(dbSweeps) == 0 {
|
||||||
log.Infof("skipping restored batch %d as it has no sweeps",
|
infof("skipping restored batch %d as it has no sweeps",
|
||||||
batch.id)
|
batch.id)
|
||||||
|
|
||||||
// It is safe to drop this empty batch as it has no sweeps.
|
// It is safe to drop this empty batch as it has no sweeps.
|
||||||
err := b.store.DropBatch(ctx, batch.id)
|
err := b.store.DropBatch(ctx, batch.id)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
log.Warnf("unable to drop empty batch %d: %v",
|
warnf("unable to drop empty batch %d: %v",
|
||||||
batch.id, err)
|
batch.id, err)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -878,7 +934,7 @@ func (b *Batcher) monitorSpendAndNotify(ctx context.Context, sweep *sweep,
|
||||||
b.wg.Add(1)
|
b.wg.Add(1)
|
||||||
go func() {
|
go func() {
|
||||||
defer b.wg.Done()
|
defer b.wg.Done()
|
||||||
log.Infof("Batcher monitoring spend for swap %x",
|
infof("Batcher monitoring spend for swap %x",
|
||||||
sweep.swapHash[:6])
|
sweep.swapHash[:6])
|
||||||
|
|
||||||
for {
|
for {
|
||||||
|
|
@ -1057,7 +1113,7 @@ func (b *Batcher) loadSweep(ctx context.Context, swapHash lntypes.Hash,
|
||||||
}
|
}
|
||||||
} else {
|
} else {
|
||||||
if s.ConfTarget == 0 {
|
if s.ConfTarget == 0 {
|
||||||
log.Warnf("Fee estimation was requested for zero "+
|
warnf("Fee estimation was requested for zero "+
|
||||||
"confTarget for sweep %x.", swapHash[:6])
|
"confTarget for sweep %x.", swapHash[:6])
|
||||||
}
|
}
|
||||||
minFeeRate, err = b.wallet.EstimateFeeRate(ctx, s.ConfTarget)
|
minFeeRate, err = b.wallet.EstimateFeeRate(ctx, s.ConfTarget)
|
||||||
|
|
|
||||||
File diff suppressed because it is too large
Load diff
Loading…
Add table
Add a link
Reference in a new issue