sweepbatcher: re-add sweeps after fully confirmed

In case of a reorg sweeps should not go to another batch but stay in the current
batch until it is reorg-safely confirmed. Only after that the remaining sweeps
are re-added to another batch.

Field sweep.completed is now set to true only for reorg-safely confirmed sweeps.

In handleConf we now use batch.persist() (i.e. store.UpdateSweepBatch) instead
of ConfirmBatch, because we set not only Confirmed flag, but also batchTxid.
This commit is contained in:
Boris Nagaev 2025-04-26 01:01:39 -03:00
parent d2c168945b
commit 650cf20fe9
No known key found for this signature in database
5 changed files with 249 additions and 159 deletions

View file

@ -201,7 +201,7 @@ type dbBatch struct {
// ID is the unique identifier of the batch.
ID int32
// Confirmed is set when the batch is fully confirmed.
// Confirmed is set when the batch is reorg-safely confirmed.
Confirmed bool
// BatchTxid is the txid of the batch transaction.
@ -236,7 +236,7 @@ type dbSweep struct {
// Amount is the amount of the sweep.
Amount btcutil.Amount
// Completed indicates whether this sweep is completed.
// Completed indicates whether this sweep is fully-confirmed.
Completed bool
}

View file

@ -1943,7 +1943,6 @@ func getFeePortionPaidBySweep(spendTx *wire.MsgTx, feePortionPerSweep,
func (b *batch) handleSpend(ctx context.Context, spendTx *wire.MsgTx) error {
var (
txHash = spendTx.TxHash()
purgeList = make([]SweepRequest, 0, len(b.sweeps))
notifyList = make([]sweep, 0, len(b.sweeps))
)
b.batchTxid = &txHash
@ -1953,7 +1952,105 @@ func (b *batch) handleSpend(ctx context.Context, spendTx *wire.MsgTx) error {
b.Warnf("transaction %v has no outputs", txHash)
}
// Determine if we should use presigned mode for the batch.
// Make a set of confirmed sweeps.
confirmedSet := make(map[wire.OutPoint]struct{}, len(spendTx.TxIn))
for _, txIn := range spendTx.TxIn {
confirmedSet[txIn.PreviousOutPoint] = struct{}{}
}
// As a previous version of the batch transaction may get confirmed,
// which does not contain the latest sweeps, we need to detect which
// sweeps are in the transaction to correctly calculate fee portions
// and notify proper sweeps.
var (
totalSweptAmt btcutil.Amount
confirmedSweeps = []wire.OutPoint{}
)
for _, sweep := range b.sweeps {
// Skip sweeps that were not included into the confirmed tx.
_, found := confirmedSet[sweep.outpoint]
if !found {
continue
}
totalSweptAmt += sweep.value
notifyList = append(notifyList, sweep)
confirmedSweeps = append(confirmedSweeps, sweep.outpoint)
}
// Calculate the fee portion that each sweep should pay for the batch.
feePortionPaidPerSweep, roundingDifference := getFeePortionForSweep(
spendTx, len(notifyList), totalSweptAmt,
)
for _, sweep := range notifyList {
// If the sweep's notifier is empty then this means that a swap
// is not waiting to read an update from it, so we can skip
// the notification part.
if sweep.notifier == nil ||
*sweep.notifier == (SpendNotifier{}) {
continue
}
spendDetail := SpendDetail{
Tx: spendTx,
OnChainFeePortion: getFeePortionPaidBySweep(
spendTx, feePortionPaidPerSweep,
roundingDifference, &sweep,
),
}
// Dispatch the sweep notifier, we don't care about the outcome
// of this action so we don't wait for it.
go func() {
// Make sure this context doesn't expire so we
// successfully notify the caller.
ctx := context.WithoutCancel(ctx)
sweep.notifySweepSpend(ctx, &spendDetail)
}()
}
b.Infof("spent, confirmed sweeps: %v", confirmedSweeps)
// We are no longer able to accept new sweeps, so we mark the batch as
// closed and persist on storage.
b.state = Closed
if err := b.persist(ctx); err != nil {
return fmt.Errorf("saving batch failed: %w", err)
}
if err := b.monitorConfirmations(ctx); err != nil {
return fmt.Errorf("monitorConfirmations failed: %w", err)
}
return nil
}
// handleConf handles a confirmation notification. This is the final step of the
// batch. Here we signal to the batcher that this batch was completed.
func (b *batch) handleConf(ctx context.Context,
conf *chainntnfs.TxConfirmation) error {
spendTx := conf.Tx
txHash := spendTx.TxHash()
if b.batchTxid == nil || *b.batchTxid != txHash {
b.Warnf("Mismatch of batch txid: tx in spend notification had "+
"txid %v, but confirmation notification has txif %v. "+
"Using the later.", b.batchTxid, txHash)
}
b.batchTxid = &txHash
b.Infof("confirmed in txid %s", b.batchTxid)
b.state = Confirmed
if err := b.persist(ctx); err != nil {
return fmt.Errorf("saving batch failed: %w", err)
}
// If the batch is in presigned mode, cleanup presignedHelper.
presigned, err := b.isPresigned()
if err != nil {
return fmt.Errorf("failed to determine if the batch %d uses "+
@ -1971,40 +2068,46 @@ func (b *batch) handleSpend(ctx context.Context, spendTx *wire.MsgTx) error {
b.id, err)
}
// Make a set of confirmed sweeps.
confirmedSet := make(map[wire.OutPoint]struct{}, len(spendTx.TxIn))
for _, txIn := range spendTx.TxIn {
confirmedSet[txIn.PreviousOutPoint] = struct{}{}
}
// As a previous version of the batch transaction may get confirmed,
// which does not contain the latest sweeps, we need to detect the
// sweeps that did not make it to the confirmed transaction and feed
// them back to the batcher. This will ensure that the sweeps will enter
// a new batch instead of remaining dangling.
var (
totalSweptAmt btcutil.Amount
confirmedSweeps = []wire.OutPoint{}
purgedSweeps = []wire.OutPoint{}
purgedSwaps = []lntypes.Hash{}
purgeList = make([]SweepRequest, 0, len(b.sweeps))
totalSweptAmt btcutil.Amount
)
for _, sweep := range allSweeps {
found := false
for _, txIn := range spendTx.TxIn {
if txIn.PreviousOutPoint == sweep.outpoint {
found = true
totalSweptAmt += sweep.value
notifyList = append(notifyList, sweep)
confirmedSweeps = append(
confirmedSweeps, sweep.outpoint,
)
break
_, found := confirmedSet[sweep.outpoint]
if found {
// Save the sweep as completed. Note that sweeps are
// marked completed after the batch is marked confirmed
// because the check in handleSweeps checks sweep's
// status first and then checks the batch status.
err := b.persistSweep(ctx, sweep, true)
if err != nil {
return err
}
confirmedSweeps = append(
confirmedSweeps, sweep.outpoint,
)
totalSweptAmt += sweep.value
continue
}
// If the sweep's outpoint was not found in the transaction's
// inputs this means it was left out. So we delete it from this
// batch and feed it back to the batcher.
if found {
continue
}
newSweep := sweep
delete(b.sweeps, sweep.outpoint)
@ -2036,6 +2139,10 @@ func (b *batch) handleSpend(ctx context.Context, spendTx *wire.MsgTx) error {
})
}
}
var (
purgedSweeps = []wire.OutPoint{}
purgedSwaps = []lntypes.Hash{}
)
for _, sweepReq := range purgeList {
purgedSwaps = append(purgedSwaps, sweepReq.SwapHash)
for _, input := range sweepReq.Inputs {
@ -2043,45 +2150,8 @@ func (b *batch) handleSpend(ctx context.Context, spendTx *wire.MsgTx) error {
}
}
// Calculate the fee portion that each sweep should pay for the batch.
feePortionPaidPerSweep, roundingDifference := getFeePortionForSweep(
spendTx, len(notifyList), totalSweptAmt,
)
for _, sweep := range notifyList {
// Save the sweep as completed.
err := b.persistSweep(ctx, sweep, true)
if err != nil {
return err
}
// If the sweep's notifier is empty then this means that a swap
// is not waiting to read an update from it, so we can skip
// the notification part.
if sweep.notifier == nil ||
*sweep.notifier == (SpendNotifier{}) {
continue
}
spendDetail := SpendDetail{
Tx: spendTx,
OnChainFeePortion: getFeePortionPaidBySweep(
spendTx, feePortionPaidPerSweep,
roundingDifference, &sweep,
),
}
// Dispatch the sweep notifier, we don't care about the outcome
// of this action so we don't wait for it.
go func() {
// Make sure this context doesn't expire so we
// successfully notify the caller.
ctx := context.WithoutCancel(ctx)
sweep.notifySweepSpend(ctx, &spendDetail)
}()
}
b.Infof("fully confirmed sweeps: %v, purged sweeps: %v, "+
"purged swaps: %v", confirmedSweeps, purgedSweeps, purgedSwaps)
// Proceed with purging the sweeps. This will feed the sweeps that
// didn't make it to the confirmed batch transaction back to the batcher
@ -2103,49 +2173,6 @@ func (b *batch) handleSpend(ctx context.Context, spendTx *wire.MsgTx) error {
}
}()
b.Infof("spent, confirmed sweeps: %v, purged sweeps: %v, "+
"purged swaps: %v, purged groups: %v", confirmedSweeps,
purgedSweeps, purgedSwaps, len(purgeList))
// We are no longer able to accept new sweeps, so we mark the batch as
// closed and persist on storage.
b.state = Closed
if err = b.persist(ctx); err != nil {
return fmt.Errorf("saving batch failed: %w", err)
}
if err = b.monitorConfirmations(ctx); err != nil {
return fmt.Errorf("monitorConfirmations failed: %w", err)
}
return nil
}
// handleConf handles a confirmation notification. This is the final step of the
// batch. Here we signal to the batcher that this batch was completed. We also
// cleanup up presigned transactions whose primarySweepID is one of the sweeps
// that were spent and fully confirmed: such a transaction can't be broadcasted
// since it is either in a block or double-spends one of spent outputs.
func (b *batch) handleConf(ctx context.Context,
conf *chainntnfs.TxConfirmation) error {
spendTx := conf.Tx
txHash := spendTx.TxHash()
if b.batchTxid == nil || *b.batchTxid != txHash {
b.Warnf("Mismatch of batch txid: tx in spend notification had "+
"txid %v, but confirmation notification has txif %v. "+
"Using the later.", b.batchTxid, txHash)
}
b.batchTxid = &txHash
// If the batch is in presigned mode, cleanup presignedHelper.
presigned, err := b.isPresigned()
if err != nil {
return fmt.Errorf("failed to determine if the batch %d uses "+
"presigned mode: %w", b.id, err)
}
if presigned {
b.Infof("Cleaning up presigned store")
@ -2161,19 +2188,7 @@ func (b *batch) handleConf(ctx context.Context,
}
}
b.Infof("confirmed in txid %s", b.batchTxid)
b.state = Confirmed
if err := b.store.ConfirmBatch(ctx, b.id); err != nil {
return fmt.Errorf("failed to store confirmed state: %w", err)
}
// Calculate the fee portion that each sweep should pay for the batch.
// TODO: make sure spendTx matches b.sweeps.
var totalSweptAmt btcutil.Amount
for _, s := range b.sweeps {
totalSweptAmt += s.value
}
feePortionPaidPerSweep, roundingDifference := getFeePortionForSweep(
spendTx, len(b.sweeps), totalSweptAmt,
)

View file

@ -192,7 +192,8 @@ type PresignedHelper interface {
loadOnly bool) (*wire.MsgTx, error)
// CleanupTransactions removes all transactions related to any of the
// outpoints. Should be called after sweep batch tx is fully confirmed.
// outpoints. Should be called after sweep batch tx is reorg-safely
// confirmed.
CleanupTransactions(ctx context.Context, inputs []wire.OutPoint) error
}
@ -274,7 +275,7 @@ type addSweepsRequest struct {
notifier *SpendNotifier
// completed is set if the sweep is spent and the spending transaction
// is confirmed.
// is reorg-safely confirmed.
completed bool
// parentBatch is the parent batch of this sweep. It is loaded ony if
@ -296,7 +297,7 @@ type SpendDetail struct {
}
// ConfDetail is a notification that is send to the user of sweepbatcher when
// a batch is fully confirmed, i.e. gets batchConfHeight confirmations.
// a batch is reorg-safely confirmed, i.e. gets batchConfHeight confirmations.
type ConfDetail struct {
// TxConfirmation has data about the confirmation of the transaction.
*chainntnfs.TxConfirmation
@ -808,8 +809,8 @@ func (b *Batcher) AddSweep(ctx context.Context, sweepReq *SweepRequest) error {
}
// If this is a presigned mode, make sure PresignSweepsGroup was called.
// We skip the check for fully confirmed sweeps, because their presigned
// transactions were already cleaned up from the store.
// We skip the check for reorg-safely confirmed sweeps, because their
// presigned transactions were already cleaned up from the store.
if sweep.presigned && !fullyConfirmed {
err := ensurePresigned(
ctx, sweeps, b.presignedHelper, b.chainParams,
@ -822,8 +823,8 @@ func (b *Batcher) AddSweep(ctx context.Context, sweepReq *SweepRequest) error {
}
infof("Batcher adding sweep group of %d sweeps with primarySweep %x, "+
"presigned=%v, completed=%v", len(sweeps), sweep.swapHash[:6],
sweep.presigned, completed)
"presigned=%v, fully_confirmed=%v", len(sweeps),
sweep.swapHash[:6], sweep.presigned, completed)
req := &addSweepsRequest{
sweeps: sweeps,
@ -883,14 +884,10 @@ func (b *Batcher) handleSweeps(ctx context.Context, sweeps []*sweep,
// 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.
// Instead we directly detect and return the spend here.
if completed && *notifier != (SpendNotifier{}) {
// The parent batch is indeed confirmed, meaning it is complete
// and we won't be able to attach this sweep to it.
if parentBatch.Confirmed {
return b.monitorSpendAndNotify(
ctx, sweep, parentBatch.ID, notifier,
)
}
if completed && parentBatch.Confirmed {
return b.monitorSpendAndNotify(
ctx, sweep, parentBatch.ID, notifier,
)
}
sweep.notifier = notifier
@ -1158,11 +1155,11 @@ func (b *Batcher) FetchUnconfirmedBatches(ctx context.Context) ([]*batch,
} else {
// We don't store Closed state separately in DB.
// If the batch is closed (included into a block, but
// not fully confirmed), it is now considered Open
// again. It will receive a spending notification as
// soon as it starts, so it is not an issue. If a sweep
// manages to be added during this time, it will be
// detected as missing when analyzing the spend
// not reorg-safely confirmed), it is now considered
// Open again. It will receive a spending notification
// as soon as it starts, so it is not an issue. If a
// sweep manages to be added during this time, it will
// be detected as missing when analyzing the spend
// notification and will be added to new batch.
batch.state = Open
}
@ -1192,6 +1189,11 @@ func (b *Batcher) FetchUnconfirmedBatches(ctx context.Context) ([]*batch,
func (b *Batcher) monitorSpendAndNotify(ctx context.Context, sweep *sweep,
parentBatchID int32, notifier *SpendNotifier) error {
// If the caller has not provided a notifier, stop.
if notifier == nil || *notifier == (SpendNotifier{}) {
return nil
}
spendCtx, cancel := context.WithCancel(ctx)
// Then we get the total amount that was swept by the batch.
@ -1301,8 +1303,8 @@ func (b *Batcher) monitorSpendAndNotify(ctx context.Context, sweep *sweep,
// monitorConfAndNotify monitors the confirmation of a specific transaction and
// writes the response back to the response channel. It is called if the batch
// is fully confirmed and we just need to deliver the data back to the caller
// though SpendNotifier.
// is reorg-safely confirmed and we just need to deliver the data back to the
// caller though SpendNotifier.
func (b *Batcher) monitorConfAndNotify(ctx context.Context, sweep *sweep,
notifier *SpendNotifier, spendTx *wire.MsgTx,
onChainFeePortion btcutil.Amount) error {

View file

@ -1360,7 +1360,12 @@ func testPresigned_purging(t *testing.T, numSwaps, numConfirmedSwaps int,
require.LessOrEqual(t, numConfirmedSwaps, numSwaps)
const sweepsPerSwap = 2
const (
sweepsPerSwap = 2
feeRate = chainfee.SatPerKWeight(10_000)
swapAmount = 3_000_001
)
sweepAmounts := []btcutil.Amount{1_000_001, 2_000_000}
lnd := test.NewMockLnd()
@ -1370,7 +1375,7 @@ func testPresigned_purging(t *testing.T, numSwaps, numConfirmedSwaps int,
customFeeRate := func(_ context.Context,
_ lntypes.Hash) (chainfee.SatPerKWeight, error) {
return chainfee.SatPerKWeight(10_000), nil
return feeRate, nil
}
presignedHelper := newMockPresignedHelper()
@ -1388,12 +1393,17 @@ func testPresigned_purging(t *testing.T, numSwaps, numConfirmedSwaps int,
checkBatcherError(t, err)
}()
swapHashes := make([]lntypes.Hash, numSwaps)
groups := make([][]Input, numSwaps)
txs := make([]*wire.MsgTx, numSwaps)
allOps := make([]wire.OutPoint, 0, numSwaps*sweepsPerSwap)
spendChans := make([]<-chan *SpendDetail, numSwaps)
confChans := make([]<-chan *ConfDetail, numSwaps)
for i := range numSwaps {
// Create a swap of sweepsPerSwap sweeps.
swapHash := lntypes.Hash{byte(i + 1)}
swapHashes[i] = swapHash
ops := make([]wire.OutPoint, sweepsPerSwap)
group := make([]Input, sweepsPerSwap)
for j := range sweepsPerSwap {
@ -1405,15 +1415,16 @@ func testPresigned_purging(t *testing.T, numSwaps, numConfirmedSwaps int,
group[j] = Input{
Outpoint: ops[j],
Value: btcutil.Amount(1_000_000 * (j + 1)),
Value: sweepAmounts[j],
}
}
groups[i] = group
// Create a swap in DB.
swap := &loopdb.LoopOutContract{
SwapContract: loopdb.SwapContract{
CltvExpiry: 111,
AmountRequested: 3_000_000,
AmountRequested: swapAmount,
ProtocolVersion: loopdb.ProtocolVersionMuSig2,
HtlcKeys: htlcKeys,
@ -1440,11 +1451,24 @@ func testPresigned_purging(t *testing.T, numSwaps, numConfirmedSwaps int,
)
require.NoError(t, err)
// Create a spending notification channel.
spendChan := make(chan *SpendDetail, 1)
spendChans[i] = spendChan
confChan := make(chan *ConfDetail, 1)
confChans[i] = confChan
notifier := &SpendNotifier{
SpendChan: spendChan,
SpendErrChan: make(chan error, 1),
ConfChan: confChan,
ConfErrChan: make(chan error, 1),
QuitChan: make(chan bool, 1),
}
// Add the sweep, triggering the publish attempt.
require.NoError(t, batcher.AddSweep(ctx, &SweepRequest{
SwapHash: swapHash,
Inputs: group,
Notifier: &dummyNotifier,
Notifier: notifier,
}))
// For the first group it should register for the sweep's spend
@ -1543,6 +1567,13 @@ func testPresigned_purging(t *testing.T, numSwaps, numConfirmedSwaps int,
SpendingHeight: int32(601 + numSwaps + 1),
}
lnd.SpendChannel <- spendDetail
// Make sure that notifiers of confirmed sweeps received notifications.
for i := range numConfirmedSwaps {
spend := <-spendChans[i]
require.Equal(t, txHash, spend.Tx.TxHash())
}
<-lnd.RegisterConfChannel
require.NoError(t, lnd.NotifyHeight(
int32(601+numSwaps+1+batchConfHeight),
@ -1554,12 +1585,18 @@ func testPresigned_purging(t *testing.T, numSwaps, numConfirmedSwaps int,
// CleanupTransactions is called here.
<-presignedHelper.cleanupCalled
// If all the swaps were confirmed, stop.
if numConfirmedSwaps == numSwaps {
return
// Increasing block height caused the second batch to re-publish.
if online && numConfirmedSwaps < numSwaps {
<-lnd.TxPublishChannel
}
if !online {
// Make sure that notifiers of confirmed sweeps received notifications.
for i := range numConfirmedSwaps {
conf := <-confChans[i]
require.Equal(t, txHash, conf.Tx.TxHash())
}
if !online && numConfirmedSwaps != numSwaps {
// If the sweeps are offline, the missing sweeps in the
// confirmed transaction should be re-added to the batcher as
// new batch. The groups are added incrementally, so we need
@ -1568,6 +1605,47 @@ func testPresigned_purging(t *testing.T, numSwaps, numConfirmedSwaps int,
<-lnd.TxPublishChannel
}
// Now make sure that a correct spend and conf contification is sent if
// AddSweep is called after confirming the sweeps.
for i := range numConfirmedSwaps {
// Create a spending notification channel.
spendChan := make(chan *SpendDetail, 1)
confChan := make(chan *ConfDetail)
notifier := &SpendNotifier{
SpendChan: spendChan,
SpendErrChan: make(chan error, 1),
ConfChan: confChan,
ConfErrChan: make(chan error, 1),
QuitChan: make(chan bool, 1),
}
// Add the sweep, triggering the publish attempt.
require.NoError(t, batcher.AddSweep(ctx, &SweepRequest{
SwapHash: swapHashes[i],
Inputs: groups[i],
Notifier: notifier,
}))
spendReg := <-lnd.RegisterSpendChannel
spendReg.SpendChannel <- spendDetail
spend := <-spendChan
require.Equal(t, txHash, spend.Tx.TxHash())
<-lnd.RegisterConfChannel
lnd.ConfChannel <- &chainntnfs.TxConfirmation{
Tx: tx,
}
conf := <-confChan
require.Equal(t, tx.TxHash(), conf.Tx.TxHash())
}
// If all the swaps were confirmed, stop.
if numConfirmedSwaps == numSwaps {
return
}
// Wait to new batch to appear and to have the expected size.
wantSize := (numSwaps - numConfirmedSwaps) * sweepsPerSwap
if online {
@ -1675,11 +1753,13 @@ func TestPresigned(t *testing.T) {
testPurging(3, 1, false)
testPurging(3, 2, false)
testPurging(5, 2, false)
testPurging(5, 3, false)
// Test cases in which the sweeps are online.
testPurging(2, 1, true)
testPurging(3, 1, true)
testPurging(3, 2, true)
testPurging(5, 2, true)
testPurging(5, 3, true)
})
}

View file

@ -2457,22 +2457,15 @@ func testSweepBatcherSweepReentry(t *testing.T, store testStore,
return b.state == Closed
}, test.Timeout, eventuallyCheckFrequency)
// Since second batch was created we check that it registered for its
// primary sweep's spend.
<-lnd.RegisterSpendChannel
// While handling the spend notification the batch should detect that
// some sweeps did not appear in the spending tx, therefore it redirects
// them back to the batcher and the batcher inserts them in a new batch.
require.Eventually(t, func() bool {
return batcher.numBatches(ctx) == 2
}, test.Timeout, eventuallyCheckFrequency)
// We mock the confirmation notification.
lnd.ConfChannel <- &chainntnfs.TxConfirmation{
Tx: spendingTx,
}
// Since second batch was created we check that it registered for its
// primary sweep's spend.
<-lnd.RegisterSpendChannel
// Wait for tx to be published.
// Here is a race condition, which is unlikely to cause a crash: if we
// wait for publish tx before sending a conf notification (previous