mirror of
https://github.com/lightninglabs/loop.git
synced 2026-08-13 12:33:03 +02:00
staticaddr/deposit: track unconfirmed deposits
Retain static-address deposits as soon as lnd reports the UTXO, even when the output is still unconfirmed. Store the first confirmation height once the output confirms. Replay the startup block to recovered deposit FSMs so expiry handling can run immediately after restart. Derive confirmation heights from a stable wallet view because lnd reports confirmation counts.
This commit is contained in:
parent
a134cb1f22
commit
32b3b9650f
9 changed files with 913 additions and 62 deletions
|
|
@ -622,6 +622,7 @@ func (d *Daemon) initialize(withMacaroonService bool) error {
|
|||
depositStore := deposit.NewSqlStore(baseDb)
|
||||
depoCfg := &deposit.ManagerConfig{
|
||||
AddressManager: staticAddressManager,
|
||||
ChainKit: d.lnd.ChainKit,
|
||||
Store: depositStore,
|
||||
WalletKit: d.lnd.WalletKit,
|
||||
ChainNotifier: d.lnd.ChainNotifier,
|
||||
|
|
|
|||
|
|
@ -29,6 +29,10 @@ func (r *ID) FromByteSlice(b []byte) error {
|
|||
|
||||
// Deposit bundles an utxo at a static address together with manager-relevant
|
||||
// data.
|
||||
//
|
||||
// Lock order: if both Manager.mu and a Deposit lock are needed, acquire
|
||||
// Manager.mu before Deposit.Lock. Never acquire Manager.mu while holding a
|
||||
// Deposit lock.
|
||||
type Deposit struct {
|
||||
sync.Mutex
|
||||
|
||||
|
|
@ -45,7 +49,8 @@ type Deposit struct {
|
|||
Value btcutil.Amount
|
||||
|
||||
// ConfirmationHeight is the absolute height at which the deposit was
|
||||
// first confirmed.
|
||||
// first confirmed. A value of zero means the deposit is still
|
||||
// unconfirmed.
|
||||
ConfirmationHeight int64
|
||||
|
||||
// TimeOutSweepPkScript is the pk script that is used to sweep the
|
||||
|
|
@ -69,6 +74,10 @@ func (d *Deposit) IsInFinalState() bool {
|
|||
d.Lock()
|
||||
defer d.Unlock()
|
||||
|
||||
return d.isInFinalStateNoLock()
|
||||
}
|
||||
|
||||
func (d *Deposit) isInFinalStateNoLock() bool {
|
||||
return d.state == Expired || d.state == Withdrawn ||
|
||||
d.state == LoopedIn || d.state == HtlcTimeoutSwept ||
|
||||
d.state == ChannelPublished
|
||||
|
|
@ -78,6 +87,10 @@ func (d *Deposit) IsExpired(currentHeight, expiry uint32) bool {
|
|||
d.Lock()
|
||||
defer d.Unlock()
|
||||
|
||||
if d.ConfirmationHeight <= 0 {
|
||||
return false
|
||||
}
|
||||
|
||||
return currentHeight >= uint32(d.ConfirmationHeight)+expiry
|
||||
}
|
||||
|
||||
|
|
|
|||
17
staticaddr/deposit/deposit_test.go
Normal file
17
staticaddr/deposit/deposit_test.go
Normal file
|
|
@ -0,0 +1,17 @@
|
|||
package deposit
|
||||
|
||||
import (
|
||||
"testing"
|
||||
|
||||
"github.com/stretchr/testify/require"
|
||||
)
|
||||
|
||||
// TestDepositIsExpiredUnconfirmed verifies that unconfirmed deposits do not
|
||||
// expire because their CSV timeout has not started yet.
|
||||
func TestDepositIsExpiredUnconfirmed(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
d := &Deposit{}
|
||||
|
||||
require.False(t, d.IsExpired(1_000, 144))
|
||||
}
|
||||
|
|
@ -4,6 +4,7 @@ import (
|
|||
"context"
|
||||
"errors"
|
||||
"fmt"
|
||||
"sync"
|
||||
|
||||
"github.com/btcsuite/btcd/txscript"
|
||||
"github.com/btcsuite/btcd/wire"
|
||||
|
|
@ -41,8 +42,8 @@ var (
|
|||
|
||||
// States.
|
||||
var (
|
||||
// Deposited signals that funds at a static address have reached the
|
||||
// confirmation height.
|
||||
// Deposited signals that funds at a static address have been detected
|
||||
// and are available to the client.
|
||||
Deposited = fsm.StateType("Deposited")
|
||||
|
||||
// Withdrawing signals that the withdrawal transaction has been
|
||||
|
|
@ -92,8 +93,8 @@ var (
|
|||
// Events.
|
||||
var (
|
||||
// OnStart is sent to the fsm once the deposit outpoint has been
|
||||
// sufficiently confirmed. It transitions the fsm into the Deposited
|
||||
// state from where we can trigger a withdrawal, a loopin or an expiry.
|
||||
// detected. It transitions the fsm into the Deposited state from where
|
||||
// we can trigger a withdrawal, a loopin or an expiry.
|
||||
OnStart = fsm.EventType("OnStart")
|
||||
|
||||
// OnWithdrawInitiated is sent to the fsm when a withdrawal has been
|
||||
|
|
@ -160,12 +161,17 @@ type FSM struct {
|
|||
|
||||
blockNtfnChan chan uint32
|
||||
|
||||
// stopChan requests shutdown of the block notification loop.
|
||||
stopChan chan struct{}
|
||||
|
||||
// quitChan stops after the FSM stops consuming blockNtfnChan.
|
||||
quitChan chan struct{}
|
||||
|
||||
// finalizedDepositChan is used to signal that the deposit has been
|
||||
// finalized and the FSM can be removed from the manager's memory.
|
||||
finalizedDepositChan chan wire.OutPoint
|
||||
|
||||
stopOnce sync.Once
|
||||
}
|
||||
|
||||
// NewFSM creates a new state machine that can action on all static address
|
||||
|
|
@ -191,6 +197,7 @@ func NewFSM(ctx context.Context, deposit *Deposit, cfg *ManagerConfig,
|
|||
params: params,
|
||||
address: address,
|
||||
blockNtfnChan: make(chan uint32),
|
||||
stopChan: make(chan struct{}),
|
||||
quitChan: make(chan struct{}),
|
||||
finalizedDepositChan: finalizedDepositChan,
|
||||
}
|
||||
|
|
@ -226,6 +233,9 @@ func NewFSM(ctx context.Context, deposit *Deposit, cfg *ManagerConfig,
|
|||
ctx, currentHeight,
|
||||
)
|
||||
|
||||
case <-fsm.stopChan:
|
||||
return
|
||||
|
||||
case <-ctx.Done():
|
||||
return
|
||||
}
|
||||
|
|
@ -235,12 +245,27 @@ func NewFSM(ctx context.Context, deposit *Deposit, cfg *ManagerConfig,
|
|||
return depoFsm, nil
|
||||
}
|
||||
|
||||
// Stop requests shutdown of the FSM's block notification loop.
|
||||
func (f *FSM) Stop() {
|
||||
if f == nil || f.stopChan == nil {
|
||||
return
|
||||
}
|
||||
|
||||
f.stopOnce.Do(func() {
|
||||
close(f.stopChan)
|
||||
})
|
||||
}
|
||||
|
||||
// handleBlockNotification inspects the current block height and sends the
|
||||
// OnExpiry event to publish the expiry sweep transaction if the deposit timed
|
||||
// out, or it republishes the expiry sweep transaction if it was not yet swept.
|
||||
func (f *FSM) handleBlockNotification(ctx context.Context,
|
||||
currentHeight uint32) {
|
||||
|
||||
if f.deposit.IsInFinalState() {
|
||||
return
|
||||
}
|
||||
|
||||
// If the deposit is expired but not yet sufficiently confirmed, we
|
||||
// republish the expiry sweep transaction.
|
||||
if f.deposit.IsExpired(currentHeight, f.params.Expiry) {
|
||||
|
|
@ -353,6 +378,11 @@ func (f *FSM) DepositStatesV0() fsm.States {
|
|||
// still pending, we publish the expiry sweep.
|
||||
OnExpiry: PublishExpirySweep,
|
||||
|
||||
// If the server publishes the HTLC without
|
||||
// paying us, we need to keep the deposit locked
|
||||
// until the HTLC timeout path can be swept.
|
||||
OnSweepingHtlcTimeout: SweepHtlcTimeout,
|
||||
|
||||
OnLoopInInitiated: LoopingIn,
|
||||
|
||||
OnRecover: LoopingIn,
|
||||
|
|
|
|||
113
staticaddr/deposit/fsm_test.go
Normal file
113
staticaddr/deposit/fsm_test.go
Normal file
|
|
@ -0,0 +1,113 @@
|
|||
package deposit
|
||||
|
||||
import (
|
||||
"context"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/btcsuite/btcd/chaincfg/chainhash"
|
||||
"github.com/btcsuite/btcd/wire"
|
||||
"github.com/lightninglabs/loop/fsm"
|
||||
"github.com/lightninglabs/loop/staticaddr/script"
|
||||
"github.com/stretchr/testify/mock"
|
||||
"github.com/stretchr/testify/require"
|
||||
)
|
||||
|
||||
// TestHandleBlockNotificationIgnoresFinalStates verifies that a block-driven
|
||||
// expiry notification cannot mutate deposits that already reached a final
|
||||
// state but have not yet been removed from the manager's active set.
|
||||
func TestHandleBlockNotificationIgnoresFinalStates(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
finalStates := []fsm.StateType{
|
||||
Expired,
|
||||
Withdrawn,
|
||||
LoopedIn,
|
||||
HtlcTimeoutSwept,
|
||||
ChannelPublished,
|
||||
}
|
||||
|
||||
for i, state := range finalStates {
|
||||
t.Run(string(state), func(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
outpoint := wire.OutPoint{
|
||||
Hash: chainhash.Hash{byte(i + 1)},
|
||||
Index: uint32(i),
|
||||
}
|
||||
deposit := &Deposit{
|
||||
OutPoint: outpoint,
|
||||
ConfirmationHeight: 1,
|
||||
}
|
||||
deposit.SetState(state)
|
||||
|
||||
depositFSM := &FSM{
|
||||
cfg: &ManagerConfig{
|
||||
Store: new(mockStore),
|
||||
},
|
||||
deposit: deposit,
|
||||
params: &script.Parameters{Expiry: 1},
|
||||
quitChan: make(chan struct{}),
|
||||
finalizedDepositChan: make(chan wire.OutPoint, 1),
|
||||
}
|
||||
depositFSM.StateMachine = fsm.NewStateMachineWithState(
|
||||
depositFSM.DepositStatesV0(), state,
|
||||
DefaultObserverSize,
|
||||
)
|
||||
depositFSM.ActionEntryFunc = depositFSM.updateDeposit
|
||||
|
||||
depositFSM.handleBlockNotification(context.Background(), 3)
|
||||
|
||||
require.Never(t, func() bool {
|
||||
return deposit.GetState() != state
|
||||
}, 100*time.Millisecond, 10*time.Millisecond)
|
||||
|
||||
select {
|
||||
case finalized := <-depositFSM.finalizedDepositChan:
|
||||
t.Fatalf("unexpected finalization for %v", finalized)
|
||||
|
||||
default:
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
// TestLoopingInTransitionsToSweepHtlcTimeout verifies that a deposit selected
|
||||
// by a loop-in can be moved into the timeout sweep state if the server confirms
|
||||
// the HTLC without paying the invoice.
|
||||
func TestLoopingInTransitionsToSweepHtlcTimeout(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
outpoint := wire.OutPoint{
|
||||
Hash: chainhash.Hash{9},
|
||||
Index: 9,
|
||||
}
|
||||
deposit := &Deposit{
|
||||
OutPoint: outpoint,
|
||||
}
|
||||
deposit.SetState(LoopingIn)
|
||||
|
||||
store := new(mockStore)
|
||||
store.On(
|
||||
"UpdateDeposit", mock.Anything, mock.Anything,
|
||||
).Return(nil).Once()
|
||||
|
||||
depositFSM := &FSM{
|
||||
cfg: &ManagerConfig{
|
||||
Store: store,
|
||||
},
|
||||
deposit: deposit,
|
||||
params: &script.Parameters{Expiry: 1},
|
||||
}
|
||||
depositFSM.StateMachine = fsm.NewStateMachineWithState(
|
||||
depositFSM.DepositStatesV0(), LoopingIn, DefaultObserverSize,
|
||||
)
|
||||
depositFSM.ActionEntryFunc = depositFSM.updateDeposit
|
||||
|
||||
err := depositFSM.SendEvent(
|
||||
t.Context(), OnSweepingHtlcTimeout, nil,
|
||||
)
|
||||
require.NoError(t, err)
|
||||
require.Equal(t, SweepHtlcTimeout, deposit.GetState())
|
||||
store.AssertExpectations(t)
|
||||
}
|
||||
|
|
@ -6,6 +6,7 @@ import (
|
|||
"fmt"
|
||||
"sort"
|
||||
"sync"
|
||||
"sync/atomic"
|
||||
"time"
|
||||
|
||||
"github.com/btcsuite/btcd/txscript"
|
||||
|
|
@ -17,9 +18,8 @@ import (
|
|||
)
|
||||
|
||||
const (
|
||||
// MinConfs is the minimum number of confirmations we require for a
|
||||
// deposit to be considered available for loop-ins, coop-spends and
|
||||
// timeouts.
|
||||
// MinConfs is the legacy minimum confirmation target deposits had to
|
||||
// reach before they were considered ready to be used for swaps.
|
||||
MinConfs = 6
|
||||
|
||||
// MaxConfs is unset since we don't require a max number of
|
||||
|
|
@ -41,6 +41,10 @@ type ManagerConfig struct {
|
|||
// address parameters.
|
||||
AddressManager AddressManager
|
||||
|
||||
// ChainKit is used to query the best known chain tip when deriving
|
||||
// confirmation heights from wallet UTXOs.
|
||||
ChainKit lndclient.ChainKitClient
|
||||
|
||||
// Store is the database store that is used to store static address
|
||||
// related records.
|
||||
Store Store
|
||||
|
|
@ -58,12 +62,20 @@ type ManagerConfig struct {
|
|||
}
|
||||
|
||||
// Manager manages the address state machines.
|
||||
//
|
||||
// Lock order: if both Manager.mu and a Deposit lock are needed, acquire
|
||||
// Manager.mu before Deposit.Lock. Never acquire Manager.mu while holding a
|
||||
// Deposit lock.
|
||||
type Manager struct {
|
||||
cfg *ManagerConfig
|
||||
|
||||
// mu guards access to the activeDeposits map.
|
||||
mu sync.Mutex
|
||||
|
||||
// reconcileMu serializes deposit reconciliation so new deposits are
|
||||
// discovered and retained exactly once per outpoint.
|
||||
reconcileMu sync.Mutex
|
||||
|
||||
// activeDeposits contains all the active static address outputs.
|
||||
activeDeposits map[wire.OutPoint]*FSM
|
||||
|
||||
|
|
@ -77,6 +89,9 @@ type Manager struct {
|
|||
// been finalized. The manager will adjust its internal state and flush
|
||||
// finalized deposits from its memory.
|
||||
finalizedDepositChan chan wire.OutPoint
|
||||
|
||||
// currentHeight stores the currently best known block height.
|
||||
currentHeight atomic.Uint32
|
||||
}
|
||||
|
||||
// NewManager creates a new deposit manager.
|
||||
|
|
@ -98,6 +113,19 @@ func (m *Manager) Run(ctx context.Context, initChan chan struct{}) error {
|
|||
return err
|
||||
}
|
||||
|
||||
var startupHeight uint32
|
||||
select {
|
||||
case height := <-newBlockChan:
|
||||
startupHeight = uint32(height)
|
||||
m.currentHeight.Store(startupHeight)
|
||||
|
||||
case err = <-newBlockErrChan:
|
||||
return err
|
||||
|
||||
case <-ctx.Done():
|
||||
return ctx.Err()
|
||||
}
|
||||
|
||||
// Recover previous deposits and static address parameters from the DB.
|
||||
err = m.recoverDeposits(ctx)
|
||||
if err != nil {
|
||||
|
|
@ -113,6 +141,13 @@ func (m *Manager) Run(ctx context.Context, initChan chan struct{}) error {
|
|||
log.Errorf("unable to reconcile deposits: %v", err)
|
||||
}
|
||||
|
||||
// The startup height was consumed before recovered deposit FSMs existed.
|
||||
// Replay it so already-expired recovered deposits can act immediately.
|
||||
err = m.notifyActiveDeposits(ctx, startupHeight)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
// Start the deposit notifier.
|
||||
m.pollDeposits(ctx)
|
||||
|
||||
|
|
@ -123,32 +158,22 @@ func (m *Manager) Run(ctx context.Context, initChan chan struct{}) error {
|
|||
for {
|
||||
select {
|
||||
case height := <-newBlockChan:
|
||||
// Inform all active deposits about a new block arrival.
|
||||
m.mu.Lock()
|
||||
activeDeposits := make([]*FSM, 0, len(m.activeDeposits))
|
||||
for _, fsm := range m.activeDeposits {
|
||||
activeDeposits = append(activeDeposits, fsm)
|
||||
m.currentHeight.Store(uint32(height))
|
||||
|
||||
err := m.reconcileDeposits(ctx)
|
||||
if err != nil {
|
||||
log.Errorf("unable to reconcile deposits: %v", err)
|
||||
}
|
||||
m.mu.Unlock()
|
||||
|
||||
for _, fsm := range activeDeposits {
|
||||
select {
|
||||
case fsm.blockNtfnChan <- uint32(height):
|
||||
|
||||
case <-fsm.quitChan:
|
||||
continue
|
||||
|
||||
case <-ctx.Done():
|
||||
return ctx.Err()
|
||||
}
|
||||
err = m.notifyActiveDeposits(ctx, uint32(height))
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
case outpoint := <-m.finalizedDepositChan:
|
||||
// If deposits notify us about their finalization, flush
|
||||
// the finalized deposit from memory.
|
||||
m.mu.Lock()
|
||||
delete(m.activeDeposits, outpoint)
|
||||
m.mu.Unlock()
|
||||
m.removeActiveDeposit(outpoint)
|
||||
|
||||
case err = <-newBlockErrChan:
|
||||
return err
|
||||
|
|
@ -159,6 +184,33 @@ func (m *Manager) Run(ctx context.Context, initChan chan struct{}) error {
|
|||
}
|
||||
}
|
||||
|
||||
// notifyActiveDeposits informs all active deposit FSMs about a new block
|
||||
// height.
|
||||
func (m *Manager) notifyActiveDeposits(ctx context.Context,
|
||||
height uint32) error {
|
||||
|
||||
m.mu.Lock()
|
||||
activeDeposits := make([]*FSM, 0, len(m.activeDeposits))
|
||||
for _, fsm := range m.activeDeposits {
|
||||
activeDeposits = append(activeDeposits, fsm)
|
||||
}
|
||||
m.mu.Unlock()
|
||||
|
||||
for _, fsm := range activeDeposits {
|
||||
select {
|
||||
case fsm.blockNtfnChan <- height:
|
||||
|
||||
case <-fsm.quitChan:
|
||||
continue
|
||||
|
||||
case <-ctx.Done():
|
||||
return ctx.Err()
|
||||
}
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
// recoverDeposits recovers static address parameters, previous deposits and
|
||||
// state machines from the database and starts the deposit notifier.
|
||||
func (m *Manager) recoverDeposits(ctx context.Context) error {
|
||||
|
|
@ -207,8 +259,10 @@ func (m *Manager) recoverDeposits(ctx context.Context) error {
|
|||
return nil
|
||||
}
|
||||
|
||||
// pollDeposits polls new deposits to our static address and notifies the
|
||||
// manager's event loop about them.
|
||||
// pollDeposits periodically polls for new deposits to our static address. This
|
||||
// complements the block-driven reconciliation in the main event loop: while new
|
||||
// blocks trigger reconcileDeposits to promptly detect confirmations, the ticker
|
||||
// here catches deposits that appear in the mempool between blocks.
|
||||
func (m *Manager) pollDeposits(ctx context.Context) {
|
||||
log.Debugf("Waiting for new static address deposits...")
|
||||
|
||||
|
|
@ -236,13 +290,20 @@ func (m *Manager) pollDeposits(ctx context.Context) {
|
|||
// far. It picks the newly identified deposits and starts a state machine per
|
||||
// deposit to track its progress.
|
||||
func (m *Manager) reconcileDeposits(ctx context.Context) error {
|
||||
m.reconcileMu.Lock()
|
||||
defer m.reconcileMu.Unlock()
|
||||
|
||||
log.Tracef("Reconciling new deposits...")
|
||||
|
||||
utxos, err := m.cfg.AddressManager.ListUnspent(
|
||||
ctx, MinConfs, MaxConfs,
|
||||
)
|
||||
utxos, bestHeight, err := m.listUnspentWithBestHeight(ctx)
|
||||
if err != nil {
|
||||
return fmt.Errorf("unable to list new deposits: %w", err)
|
||||
return err
|
||||
}
|
||||
|
||||
err = m.updateDepositConfirmations(ctx, utxos, bestHeight)
|
||||
if err != nil {
|
||||
return fmt.Errorf("unable to update deposit "+
|
||||
"confirmations: %w", err)
|
||||
}
|
||||
|
||||
newDeposits := m.filterNewDeposits(utxos)
|
||||
|
|
@ -252,7 +313,7 @@ func (m *Manager) reconcileDeposits(ctx context.Context) error {
|
|||
}
|
||||
|
||||
for _, utxo := range newDeposits {
|
||||
deposit, err := m.createNewDeposit(ctx, utxo)
|
||||
deposit, err := m.createNewDeposit(ctx, utxo, bestHeight)
|
||||
if err != nil {
|
||||
return fmt.Errorf("unable to retain new deposit: %w",
|
||||
err)
|
||||
|
|
@ -269,12 +330,73 @@ func (m *Manager) reconcileDeposits(ctx context.Context) error {
|
|||
return nil
|
||||
}
|
||||
|
||||
// listUnspentWithBestHeight returns the wallet's current static-address UTXOs
|
||||
// together with a stable chain tip height for any confirmed outputs.
|
||||
func (m *Manager) listUnspentWithBestHeight(ctx context.Context) (
|
||||
[]*lnwallet.Utxo, int32, error) {
|
||||
|
||||
utxos, err := m.cfg.AddressManager.ListUnspent(
|
||||
ctx, 0, MaxConfs,
|
||||
)
|
||||
if err != nil {
|
||||
return nil, 0, fmt.Errorf("unable to list new deposits: %w", err)
|
||||
}
|
||||
|
||||
needsBestHeight := false
|
||||
for _, utxo := range utxos {
|
||||
if utxo.Confirmations > 0 {
|
||||
needsBestHeight = true
|
||||
break
|
||||
}
|
||||
}
|
||||
|
||||
if !needsBestHeight {
|
||||
return utxos, 0, nil
|
||||
}
|
||||
|
||||
if m.cfg.ChainKit == nil {
|
||||
return nil, 0, errors.New("chain kit client required for " +
|
||||
"confirmed deposits")
|
||||
}
|
||||
|
||||
const maxAttempts = 3
|
||||
for range maxAttempts {
|
||||
_, beforeHeight, err := m.cfg.ChainKit.GetBestBlock(ctx)
|
||||
if err != nil {
|
||||
return nil, 0, fmt.Errorf("unable to get best block "+
|
||||
"before listing deposits: %w", err)
|
||||
}
|
||||
|
||||
utxos, err = m.cfg.AddressManager.ListUnspent(ctx, 0, MaxConfs)
|
||||
if err != nil {
|
||||
return nil, 0, fmt.Errorf("unable to list new deposits: %w",
|
||||
err)
|
||||
}
|
||||
|
||||
_, afterHeight, err := m.cfg.ChainKit.GetBestBlock(ctx)
|
||||
if err != nil {
|
||||
return nil, 0, fmt.Errorf("unable to get best block "+
|
||||
"after listing deposits: %w", err)
|
||||
}
|
||||
|
||||
if beforeHeight == afterHeight {
|
||||
m.currentHeight.Store(uint32(afterHeight))
|
||||
return utxos, afterHeight, nil
|
||||
}
|
||||
}
|
||||
|
||||
return nil, 0, errors.New("unable to get stable best block while " +
|
||||
"listing deposits")
|
||||
}
|
||||
|
||||
// createNewDeposit transforms the wallet utxo into a deposit struct and stores
|
||||
// it in our database and manager memory.
|
||||
func (m *Manager) createNewDeposit(ctx context.Context,
|
||||
utxo *lnwallet.Utxo) (*Deposit, error) {
|
||||
utxo *lnwallet.Utxo, bestHeight int32) (*Deposit, error) {
|
||||
|
||||
blockHeight, err := m.getBlockHeight(ctx, utxo)
|
||||
confirmationHeight, err := confirmationHeightForUtxo(
|
||||
bestHeight, utxo,
|
||||
)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
|
@ -302,7 +424,7 @@ func (m *Manager) createNewDeposit(ctx context.Context,
|
|||
state: Deposited,
|
||||
OutPoint: utxo.OutPoint,
|
||||
Value: utxo.Value,
|
||||
ConfirmationHeight: int64(blockHeight),
|
||||
ConfirmationHeight: confirmationHeight,
|
||||
TimeOutSweepPkScript: timeoutSweepPkScript,
|
||||
}
|
||||
|
||||
|
|
@ -318,37 +440,70 @@ func (m *Manager) createNewDeposit(ctx context.Context,
|
|||
return deposit, nil
|
||||
}
|
||||
|
||||
// getBlockHeight retrieves the block height of a given utxo.
|
||||
func (m *Manager) getBlockHeight(ctx context.Context,
|
||||
utxo *lnwallet.Utxo) (uint32, error) {
|
||||
// confirmationHeightForUtxo derives the first confirmation height of a wallet
|
||||
// UTXO from a stable best-known chain tip. Unconfirmed UTXOs return 0.
|
||||
func confirmationHeightForUtxo(bestHeight int32,
|
||||
utxo *lnwallet.Utxo) (int64, error) {
|
||||
|
||||
addressParams, err := m.cfg.AddressManager.GetStaticAddressParameters(
|
||||
ctx,
|
||||
)
|
||||
if err != nil {
|
||||
return 0, fmt.Errorf("couldn't get confirmation height for "+
|
||||
"deposit, %w", err)
|
||||
if utxo.Confirmations <= 0 {
|
||||
return 0, nil
|
||||
}
|
||||
|
||||
notifChan, errChan, err :=
|
||||
m.cfg.ChainNotifier.RegisterConfirmationsNtfn(
|
||||
ctx, &utxo.OutPoint.Hash, addressParams.PkScript,
|
||||
MinConfs, addressParams.InitiationHeight,
|
||||
if bestHeight <= 0 {
|
||||
return 0, fmt.Errorf("invalid best height %d", bestHeight)
|
||||
}
|
||||
|
||||
firstConfirmationHeight := int64(bestHeight) - utxo.Confirmations + 1
|
||||
if firstConfirmationHeight <= 0 {
|
||||
return 0, fmt.Errorf("invalid confirmation height %d for %v "+
|
||||
"with best height %d and %d confirmations",
|
||||
firstConfirmationHeight, utxo.OutPoint, bestHeight,
|
||||
utxo.Confirmations)
|
||||
}
|
||||
|
||||
return firstConfirmationHeight, nil
|
||||
}
|
||||
|
||||
// updateDepositConfirmations syncs first confirmation heights for deposits that
|
||||
// are visible in lnd's wallet view.
|
||||
func (m *Manager) updateDepositConfirmations(ctx context.Context,
|
||||
utxos []*lnwallet.Utxo, bestHeight int32) error {
|
||||
|
||||
for _, utxo := range utxos {
|
||||
m.mu.Lock()
|
||||
deposit, ok := m.deposits[utxo.OutPoint]
|
||||
m.mu.Unlock()
|
||||
if !ok {
|
||||
continue
|
||||
}
|
||||
|
||||
confirmationHeight, err := confirmationHeightForUtxo(
|
||||
bestHeight, utxo,
|
||||
)
|
||||
if err != nil {
|
||||
return 0, err
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
deposit.Lock()
|
||||
if deposit.ConfirmationHeight == confirmationHeight {
|
||||
deposit.Unlock()
|
||||
continue
|
||||
}
|
||||
|
||||
previousConfirmationHeight := deposit.ConfirmationHeight
|
||||
deposit.ConfirmationHeight = confirmationHeight
|
||||
|
||||
err = m.cfg.Store.UpdateDeposit(ctx, deposit)
|
||||
if err != nil {
|
||||
deposit.ConfirmationHeight = previousConfirmationHeight
|
||||
deposit.Unlock()
|
||||
return err
|
||||
}
|
||||
|
||||
deposit.Unlock()
|
||||
}
|
||||
|
||||
select {
|
||||
case tx := <-notifChan:
|
||||
return tx.BlockHeight, nil
|
||||
|
||||
case err := <-errChan:
|
||||
return 0, err
|
||||
|
||||
case <-ctx.Done():
|
||||
return 0, ctx.Err()
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// filterNewDeposits filters the given utxos for new deposits that we haven't
|
||||
|
|
@ -439,6 +594,10 @@ func (m *Manager) GetActiveDepositsInState(stateFilter fsm.StateType) (
|
|||
func (m *Manager) AllOutpointsActiveDeposits(outpoints []wire.OutPoint,
|
||||
targetState fsm.StateType) ([]*Deposit, bool) {
|
||||
|
||||
if hasDuplicateOutpoints(outpoints) {
|
||||
return nil, false
|
||||
}
|
||||
|
||||
m.mu.Lock()
|
||||
defer m.mu.Unlock()
|
||||
|
||||
|
|
@ -495,8 +654,15 @@ func (m *Manager) TransitionDeposits(ctx context.Context, deposits []*Deposit,
|
|||
|
||||
outpoints := make([]wire.OutPoint, len(deposits))
|
||||
for i, d := range deposits {
|
||||
if d == nil {
|
||||
return fmt.Errorf("nil deposit at index %d", i)
|
||||
}
|
||||
|
||||
outpoints[i] = d.OutPoint
|
||||
}
|
||||
if hasDuplicateOutpoints(outpoints) {
|
||||
return fmt.Errorf("duplicate deposit outpoint")
|
||||
}
|
||||
|
||||
m.mu.Lock()
|
||||
stateMachines, _ := m.toActiveDeposits(&outpoints)
|
||||
|
|
@ -508,6 +674,13 @@ func (m *Manager) TransitionDeposits(ctx context.Context, deposits []*Deposit,
|
|||
|
||||
lockDeposits(deposits)
|
||||
defer unlockDeposits(deposits)
|
||||
for _, deposit := range deposits {
|
||||
if deposit.isInFinalStateNoLock() {
|
||||
return fmt.Errorf("deposit %v is no longer active in "+
|
||||
"state %v", deposit.OutPoint, deposit.state)
|
||||
}
|
||||
}
|
||||
|
||||
for _, sm := range stateMachines {
|
||||
err := sm.SendEvent(ctx, event, nil)
|
||||
if err != nil {
|
||||
|
|
@ -537,6 +710,37 @@ func unlockDeposits(deposits []*Deposit) {
|
|||
}
|
||||
}
|
||||
|
||||
// hasDuplicateOutpoints prevents callers from referencing the same deposit
|
||||
// more than once. Selection and transition paths later lock every returned
|
||||
// deposit, so allowing duplicate outpoints could put the same deposit pointer in
|
||||
// the list twice and self-deadlock.
|
||||
func hasDuplicateOutpoints(outpoints []wire.OutPoint) bool {
|
||||
seen := make(map[wire.OutPoint]struct{}, len(outpoints))
|
||||
for _, outpoint := range outpoints {
|
||||
if _, ok := seen[outpoint]; ok {
|
||||
return true
|
||||
}
|
||||
|
||||
seen[outpoint] = struct{}{}
|
||||
}
|
||||
|
||||
return false
|
||||
}
|
||||
|
||||
// removeActiveDeposit removes and stops the FSM for an active outpoint.
|
||||
func (m *Manager) removeActiveDeposit(outpoint wire.OutPoint) {
|
||||
m.mu.Lock()
|
||||
fsm, ok := m.activeDeposits[outpoint]
|
||||
if ok {
|
||||
delete(m.activeDeposits, outpoint)
|
||||
}
|
||||
m.mu.Unlock()
|
||||
|
||||
if ok {
|
||||
fsm.Stop()
|
||||
}
|
||||
}
|
||||
|
||||
// GetAllDeposits returns all active deposits.
|
||||
func (m *Manager) GetAllDeposits(ctx context.Context) ([]*Deposit, error) {
|
||||
return m.cfg.Store.AllDeposits(ctx)
|
||||
|
|
|
|||
37
staticaddr/deposit/manager_height_test.go
Normal file
37
staticaddr/deposit/manager_height_test.go
Normal file
|
|
@ -0,0 +1,37 @@
|
|||
package deposit
|
||||
|
||||
import (
|
||||
"testing"
|
||||
|
||||
"github.com/btcsuite/btcd/chaincfg/chainhash"
|
||||
"github.com/btcsuite/btcd/wire"
|
||||
"github.com/lightningnetwork/lnd/lnwallet"
|
||||
"github.com/stretchr/testify/require"
|
||||
)
|
||||
|
||||
func TestConfirmationHeightForUtxo(t *testing.T) {
|
||||
t.Run("unconfirmed", func(t *testing.T) {
|
||||
height, err := confirmationHeightForUtxo(0, &lnwallet.Utxo{})
|
||||
require.NoError(t, err)
|
||||
require.Zero(t, height)
|
||||
})
|
||||
|
||||
t.Run("confirmed", func(t *testing.T) {
|
||||
height, err := confirmationHeightForUtxo(101, &lnwallet.Utxo{
|
||||
OutPoint: wire.OutPoint{
|
||||
Hash: chainhash.Hash{1},
|
||||
Index: 2,
|
||||
},
|
||||
Confirmations: 6,
|
||||
})
|
||||
require.NoError(t, err)
|
||||
require.EqualValues(t, 96, height)
|
||||
})
|
||||
|
||||
t.Run("invalid best height", func(t *testing.T) {
|
||||
_, err := confirmationHeightForUtxo(2, &lnwallet.Utxo{
|
||||
Confirmations: 6,
|
||||
})
|
||||
require.ErrorContains(t, err, "invalid confirmation height")
|
||||
})
|
||||
}
|
||||
334
staticaddr/deposit/manager_reconcile_test.go
Normal file
334
staticaddr/deposit/manager_reconcile_test.go
Normal file
|
|
@ -0,0 +1,334 @@
|
|||
package deposit
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"sync"
|
||||
"sync/atomic"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/btcsuite/btcd/btcutil"
|
||||
"github.com/btcsuite/btcd/chaincfg/chainhash"
|
||||
"github.com/btcsuite/btcd/wire"
|
||||
"github.com/lightninglabs/loop/staticaddr/script"
|
||||
"github.com/lightninglabs/loop/staticaddr/version"
|
||||
"github.com/lightninglabs/loop/test"
|
||||
"github.com/lightningnetwork/lnd/lnwallet"
|
||||
"github.com/stretchr/testify/mock"
|
||||
"github.com/stretchr/testify/require"
|
||||
)
|
||||
|
||||
// expectStableBestBlock configures two stable best-block lookups.
|
||||
func expectStableBestBlock(mockChainKit *MockChainKit, height int32) {
|
||||
mockChainKit.On(
|
||||
"GetBestBlock", mock.Anything,
|
||||
).Return(chainhash.Hash{}, height, nil).Twice()
|
||||
}
|
||||
|
||||
// TestReconcileDepositsSerialized verifies reconciliation is serialized across
|
||||
// concurrent callers.
|
||||
func TestReconcileDepositsSerialized(t *testing.T) {
|
||||
ctx := context.Background()
|
||||
mockLnd := test.NewMockLnd()
|
||||
utxo := &lnwallet.Utxo{
|
||||
AddressType: lnwallet.TaprootPubkey,
|
||||
Value: btcutil.Amount(100_000),
|
||||
Confirmations: 0,
|
||||
OutPoint: wire.OutPoint{
|
||||
Hash: chainhash.Hash{1},
|
||||
Index: 1,
|
||||
},
|
||||
}
|
||||
|
||||
mockAddressManager := new(mockAddressManager)
|
||||
mockAddressManager.On(
|
||||
"ListUnspent", mock.Anything, int32(0), int32(MaxConfs),
|
||||
).Return([]*lnwallet.Utxo{utxo}, nil)
|
||||
mockAddressManager.On(
|
||||
"GetStaticAddressParameters", mock.Anything,
|
||||
).Return((*script.Parameters)(nil), errors.New("fsm init failed"))
|
||||
|
||||
mockStore := new(mockStore)
|
||||
var createCalls atomic.Int32
|
||||
createEntered := make(chan struct{})
|
||||
releaseCreate := make(chan struct{})
|
||||
mockStore.On(
|
||||
"CreateDeposit", mock.Anything, mock.Anything,
|
||||
).Return(nil).Run(func(mock.Arguments) {
|
||||
if createCalls.Add(1) == 1 {
|
||||
close(createEntered)
|
||||
}
|
||||
|
||||
<-releaseCreate
|
||||
})
|
||||
|
||||
manager := NewManager(&ManagerConfig{
|
||||
AddressManager: mockAddressManager,
|
||||
Store: mockStore,
|
||||
WalletKit: mockLnd.WalletKit,
|
||||
Signer: mockLnd.Signer,
|
||||
})
|
||||
|
||||
var wg sync.WaitGroup
|
||||
wg.Add(2)
|
||||
|
||||
errs := make(chan error, 2)
|
||||
go func() {
|
||||
defer wg.Done()
|
||||
errs <- manager.reconcileDeposits(ctx)
|
||||
}()
|
||||
|
||||
<-createEntered
|
||||
|
||||
go func() {
|
||||
defer wg.Done()
|
||||
errs <- manager.reconcileDeposits(ctx)
|
||||
}()
|
||||
|
||||
time.Sleep(100 * time.Millisecond)
|
||||
close(releaseCreate)
|
||||
wg.Wait()
|
||||
close(errs)
|
||||
|
||||
var gotErrs []error
|
||||
for err := range errs {
|
||||
gotErrs = append(gotErrs, err)
|
||||
}
|
||||
|
||||
require.EqualValues(t, 1, createCalls.Load())
|
||||
require.Len(t, manager.deposits, 1)
|
||||
require.Empty(t, manager.activeDeposits)
|
||||
require.Len(t, gotErrs, 2)
|
||||
|
||||
var errCount int
|
||||
for _, err := range gotErrs {
|
||||
if err == nil {
|
||||
continue
|
||||
}
|
||||
|
||||
errCount++
|
||||
require.ErrorContains(t, err, "unable to start new deposit FSM")
|
||||
}
|
||||
require.Equal(t, 1, errCount)
|
||||
}
|
||||
|
||||
// TestReconcileConfirmedDepositUsesBestBlockHeight verifies confirmation
|
||||
// heights are derived from a stable chain tip.
|
||||
func TestReconcileConfirmedDepositUsesBestBlockHeight(t *testing.T) {
|
||||
ctx := context.Background()
|
||||
mockLnd := test.NewMockLnd()
|
||||
utxo := &lnwallet.Utxo{
|
||||
AddressType: lnwallet.TaprootPubkey,
|
||||
Value: btcutil.Amount(100_000),
|
||||
Confirmations: 3,
|
||||
OutPoint: wire.OutPoint{
|
||||
Hash: chainhash.Hash{8},
|
||||
Index: 1,
|
||||
},
|
||||
}
|
||||
|
||||
mockAddressManager := new(mockAddressManager)
|
||||
mockAddressManager.On(
|
||||
"ListUnspent", mock.Anything, int32(0), int32(MaxConfs),
|
||||
).Return([]*lnwallet.Utxo{utxo}, nil)
|
||||
mockAddressManager.On(
|
||||
"GetStaticAddressParameters", mock.Anything,
|
||||
).Return((*script.Parameters)(nil), errors.New("fsm init failed"))
|
||||
|
||||
mockChainKit := new(MockChainKit)
|
||||
expectStableBestBlock(mockChainKit, 100)
|
||||
|
||||
mockStore := new(mockStore)
|
||||
mockStore.On(
|
||||
"CreateDeposit", mock.Anything, mock.Anything,
|
||||
).Return(nil).Run(func(args mock.Arguments) {
|
||||
createdDeposit := args.Get(1).(*Deposit)
|
||||
require.EqualValues(t, 98, createdDeposit.ConfirmationHeight)
|
||||
})
|
||||
|
||||
manager := NewManager(&ManagerConfig{
|
||||
AddressManager: mockAddressManager,
|
||||
ChainKit: mockChainKit,
|
||||
Store: mockStore,
|
||||
WalletKit: mockLnd.WalletKit,
|
||||
Signer: mockLnd.Signer,
|
||||
})
|
||||
|
||||
err := manager.reconcileDeposits(ctx)
|
||||
require.ErrorContains(t, err, "unable to start new deposit FSM")
|
||||
}
|
||||
|
||||
// TestUpdateDepositConfirmationsResetsReorgedDeposit verifies that a deposit
|
||||
// which remains wallet-visible but loses confirmations has its confirmation
|
||||
// height reset. This can happen if a confirmed transaction is reorged back into
|
||||
// the mempool.
|
||||
func TestUpdateDepositConfirmationsResetsReorgedDeposit(t *testing.T) {
|
||||
ctx := context.Background()
|
||||
outpoint := wire.OutPoint{
|
||||
Hash: chainhash.Hash{7},
|
||||
Index: 2,
|
||||
}
|
||||
|
||||
deposit := &Deposit{
|
||||
OutPoint: outpoint,
|
||||
ConfirmationHeight: 99,
|
||||
}
|
||||
deposit.SetState(Deposited)
|
||||
|
||||
utxo := &lnwallet.Utxo{
|
||||
OutPoint: outpoint,
|
||||
Confirmations: 0,
|
||||
}
|
||||
|
||||
mockStore := new(mockStore)
|
||||
mockStore.On(
|
||||
"UpdateDeposit", mock.Anything, mock.Anything,
|
||||
).Return(nil).Run(func(args mock.Arguments) {
|
||||
updatedDeposit := args.Get(1).(*Deposit)
|
||||
require.Zero(t, updatedDeposit.ConfirmationHeight)
|
||||
})
|
||||
|
||||
manager := NewManager(&ManagerConfig{
|
||||
Store: mockStore,
|
||||
})
|
||||
manager.deposits[outpoint] = deposit
|
||||
|
||||
err := manager.updateDepositConfirmations(ctx, []*lnwallet.Utxo{utxo}, 0)
|
||||
require.NoError(t, err)
|
||||
require.Zero(t, deposit.ConfirmationHeight)
|
||||
mockStore.AssertExpectations(t)
|
||||
}
|
||||
|
||||
// TestAllOutpointsActiveDepositsRejectsDuplicateOutpoints verifies that a
|
||||
// duplicated selection is rejected before the manager tries to lock the same
|
||||
// deposit twice.
|
||||
func TestAllOutpointsActiveDepositsRejectsDuplicateOutpoints(t *testing.T) {
|
||||
outpoint := wire.OutPoint{
|
||||
Hash: chainhash.Hash{12},
|
||||
Index: 6,
|
||||
}
|
||||
|
||||
deposit := &Deposit{
|
||||
OutPoint: outpoint,
|
||||
}
|
||||
deposit.SetState(Deposited)
|
||||
|
||||
manager := NewManager(&ManagerConfig{})
|
||||
manager.deposits[outpoint] = deposit
|
||||
manager.activeDeposits[outpoint] = &FSM{
|
||||
deposit: deposit,
|
||||
}
|
||||
|
||||
deposits, ok := manager.AllOutpointsActiveDeposits(
|
||||
[]wire.OutPoint{outpoint, outpoint}, Deposited,
|
||||
)
|
||||
require.False(t, ok)
|
||||
require.Nil(t, deposits)
|
||||
}
|
||||
|
||||
// TestTransitionDepositsRejectsDuplicateOutpoints verifies that transition
|
||||
// callers cannot deadlock the manager by passing the same deposit twice.
|
||||
func TestTransitionDepositsRejectsDuplicateOutpoints(t *testing.T) {
|
||||
outpoint := wire.OutPoint{
|
||||
Hash: chainhash.Hash{13},
|
||||
Index: 6,
|
||||
}
|
||||
|
||||
deposit := &Deposit{
|
||||
OutPoint: outpoint,
|
||||
}
|
||||
deposit.SetState(Deposited)
|
||||
|
||||
manager := NewManager(&ManagerConfig{})
|
||||
err := manager.TransitionDeposits(
|
||||
t.Context(), []*Deposit{deposit, deposit}, OnLoopInInitiated,
|
||||
LoopingIn,
|
||||
)
|
||||
require.ErrorContains(t, err, "duplicate deposit outpoint")
|
||||
require.Equal(t, Deposited, deposit.GetState())
|
||||
}
|
||||
|
||||
// TestReconcileReplacementDepositCreatesNewDeposit ensures that a replacement
|
||||
// UTXO is retained as a new deposit while an in-flight deposit remains tied to
|
||||
// the outpoint selected by a loop-in.
|
||||
func TestReconcileReplacementDepositCreatesNewDeposit(t *testing.T) {
|
||||
ctx := context.Background()
|
||||
mockLnd := test.NewMockLnd()
|
||||
oldOutpoint := wire.OutPoint{
|
||||
Hash: chainhash.Hash{4},
|
||||
Index: 8,
|
||||
}
|
||||
newOutpoint := wire.OutPoint{
|
||||
Hash: chainhash.Hash{5},
|
||||
Index: 9,
|
||||
}
|
||||
|
||||
depositID, err := GetRandomDepositID()
|
||||
require.NoError(t, err)
|
||||
|
||||
deposit := &Deposit{
|
||||
ID: depositID,
|
||||
OutPoint: oldOutpoint,
|
||||
Value: btcutil.Amount(100_000),
|
||||
}
|
||||
deposit.SetState(LoopingIn)
|
||||
|
||||
utxo := &lnwallet.Utxo{
|
||||
OutPoint: newOutpoint,
|
||||
Value: deposit.Value,
|
||||
Confirmations: 0,
|
||||
}
|
||||
|
||||
mockAddressManager := new(mockAddressManager)
|
||||
mockAddressManager.On(
|
||||
"ListUnspent", mock.Anything, int32(0), int32(MaxConfs),
|
||||
).Return([]*lnwallet.Utxo{utxo}, nil)
|
||||
mockAddressManager.On(
|
||||
"GetStaticAddressParameters", mock.Anything,
|
||||
).Return(&script.Parameters{
|
||||
ProtocolVersion: version.ProtocolVersion_V0,
|
||||
}, nil)
|
||||
mockAddressManager.On(
|
||||
"GetStaticAddress", mock.Anything,
|
||||
).Return((*script.StaticAddress)(nil), nil)
|
||||
|
||||
mockStore := new(mockStore)
|
||||
var createdDeposit *Deposit
|
||||
mockStore.On(
|
||||
"CreateDeposit", mock.Anything, mock.Anything,
|
||||
).Return(nil).Run(func(args mock.Arguments) {
|
||||
createdDeposit = args.Get(1).(*Deposit)
|
||||
})
|
||||
|
||||
manager := NewManager(&ManagerConfig{
|
||||
AddressManager: mockAddressManager,
|
||||
Store: mockStore,
|
||||
WalletKit: mockLnd.WalletKit,
|
||||
Signer: mockLnd.Signer,
|
||||
})
|
||||
manager.deposits[oldOutpoint] = deposit
|
||||
fsm := &FSM{}
|
||||
manager.activeDeposits[oldOutpoint] = fsm
|
||||
|
||||
require.NoError(t, manager.reconcileDeposits(ctx))
|
||||
|
||||
require.Same(t, deposit, manager.deposits[oldOutpoint])
|
||||
require.Equal(t, oldOutpoint, deposit.OutPoint)
|
||||
require.Equal(t, LoopingIn, deposit.GetState())
|
||||
|
||||
replacement, ok := manager.deposits[newOutpoint]
|
||||
require.True(t, ok)
|
||||
require.Same(t, createdDeposit, replacement)
|
||||
require.NotEqual(t, depositID, replacement.ID)
|
||||
require.Equal(t, newOutpoint, replacement.OutPoint)
|
||||
require.Equal(t, Deposited, replacement.GetState())
|
||||
require.Zero(t, replacement.ConfirmationHeight)
|
||||
|
||||
require.Same(t, fsm, manager.activeDeposits[oldOutpoint])
|
||||
require.NotSame(t, fsm, manager.activeDeposits[newOutpoint])
|
||||
|
||||
mockStore.AssertNotCalled(
|
||||
t, "UpdateDeposit", mock.Anything, mock.Anything,
|
||||
)
|
||||
}
|
||||
|
|
@ -219,6 +219,49 @@ func (m *MockChainNotifier) RegisterSpendNtfn(ctx context.Context,
|
|||
args.Get(1).(chan error), args.Error(2)
|
||||
}
|
||||
|
||||
type MockChainKit struct {
|
||||
mock.Mock
|
||||
}
|
||||
|
||||
// RawClientWithMacAuth implements lndclient.ChainKitClient for tests.
|
||||
func (m *MockChainKit) RawClientWithMacAuth(
|
||||
ctx context.Context) (context.Context, time.Duration,
|
||||
chainrpc.ChainKitClient) {
|
||||
|
||||
return ctx, 0, nil
|
||||
}
|
||||
|
||||
// GetBlock implements lndclient.ChainKitClient for tests.
|
||||
func (m *MockChainKit) GetBlock(context.Context, chainhash.Hash) (
|
||||
*wire.MsgBlock, error) {
|
||||
|
||||
panic("unexpected GetBlock call")
|
||||
}
|
||||
|
||||
// GetBlockHeader implements lndclient.ChainKitClient for tests.
|
||||
func (m *MockChainKit) GetBlockHeader(context.Context, chainhash.Hash) (
|
||||
*wire.BlockHeader, error) {
|
||||
|
||||
panic("unexpected GetBlockHeader call")
|
||||
}
|
||||
|
||||
// GetBestBlock returns the configured best-block mock response.
|
||||
func (m *MockChainKit) GetBestBlock(ctx context.Context) (
|
||||
chainhash.Hash, int32, error) {
|
||||
|
||||
args := m.Called(ctx)
|
||||
|
||||
return args.Get(0).(chainhash.Hash), args.Get(1).(int32),
|
||||
args.Error(2)
|
||||
}
|
||||
|
||||
// GetBlockHash implements lndclient.ChainKitClient for tests.
|
||||
func (m *MockChainKit) GetBlockHash(context.Context, int64) (
|
||||
chainhash.Hash, error) {
|
||||
|
||||
panic("unexpected GetBlockHash call")
|
||||
}
|
||||
|
||||
// TestManager checks that the manager processes the right channel notifications
|
||||
// while a deposit is expiring.
|
||||
func TestManager(t *testing.T) {
|
||||
|
|
@ -237,6 +280,10 @@ func TestManager(t *testing.T) {
|
|||
runErrChan <- testContext.manager.Run(ctx, initChan)
|
||||
}()
|
||||
|
||||
// Send an initial block so the manager can proceed past its startup
|
||||
// block wait.
|
||||
testContext.blockChan <- int32(defaultDepositConfirmations)
|
||||
|
||||
// Ensure that the manager has been initialized.
|
||||
select {
|
||||
case <-initChan:
|
||||
|
|
@ -307,6 +354,61 @@ func TestManager(t *testing.T) {
|
|||
}
|
||||
}
|
||||
|
||||
// TestManagerReplaysStartupBlockToRecoveredDeposits verifies that the initial
|
||||
// block epoch consumed during startup is delivered to recovered deposit FSMs.
|
||||
func TestManagerReplaysStartupBlockToRecoveredDeposits(t *testing.T) {
|
||||
ctx, cancel := context.WithCancel(t.Context())
|
||||
defer cancel()
|
||||
|
||||
const defaultTimeout = 30 * time.Second
|
||||
|
||||
testContext := newManagerTestContext(t)
|
||||
|
||||
initChan := make(chan struct{})
|
||||
runErrChan := make(chan error, 1)
|
||||
go func() {
|
||||
runErrChan <- testContext.manager.Run(ctx, initChan)
|
||||
}()
|
||||
|
||||
// Send only the startup block at the recovered deposit's expiry height.
|
||||
testContext.blockChan <- int32(
|
||||
defaultDepositConfirmations + defaultExpiry,
|
||||
)
|
||||
|
||||
select {
|
||||
case <-initChan:
|
||||
|
||||
case err := <-runErrChan:
|
||||
require.NoError(t, err, "manager failed to start")
|
||||
|
||||
case <-time.After(defaultTimeout):
|
||||
t.Fatal("manager timed out starting")
|
||||
}
|
||||
|
||||
select {
|
||||
case <-testContext.mockLnd.SignOutputRawChannel:
|
||||
|
||||
case <-time.After(defaultTimeout):
|
||||
t.Fatal("did not receive sign request")
|
||||
}
|
||||
|
||||
select {
|
||||
case <-testContext.mockLnd.TxPublishChannel:
|
||||
|
||||
case <-time.After(defaultTimeout):
|
||||
t.Fatal("did not receive published expiry tx")
|
||||
}
|
||||
|
||||
cancel()
|
||||
select {
|
||||
case err := <-runErrChan:
|
||||
require.ErrorIs(t, err, context.Canceled)
|
||||
|
||||
case <-time.After(defaultTimeout):
|
||||
t.Fatal("manager did not stop")
|
||||
}
|
||||
}
|
||||
|
||||
// ManagerTestContext is a helper struct that contains all the necessary
|
||||
// components to test the reservation manager.
|
||||
type ManagerTestContext struct {
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue