staticaddr/deposit: reconcile active deposits with wallet

Treat lnd's wallet view as the source of spendable static-address
outpoints while keeping historical deposit records in the DB. Reconcile
active FSMs against the current wallet view, reactivate known deposits
that reappear, and hide stale Deposited records from the visible
deposit set.
This commit is contained in:
Slyghtning 2026-07-08 13:51:12 +02:00
parent ac12d251f5
commit e4bcc94a36
No known key found for this signature in database
GPG key ID: F82D456EA023C9BF
2 changed files with 437 additions and 3 deletions

View file

@ -284,6 +284,15 @@ func (m *Manager) pollDeposits(ctx context.Context) {
}()
}
// EnsureDepositsFresh reconciles the cached active deposit set with lnd's
// current wallet view. Spending paths call this before selecting deposits so
// stale persisted records are not treated as live funds. This can happen when
// an unconfirmed funding transaction is replaced, a confirmed deposit is
// reorged out, or the output was spent outside the active manager path.
func (m *Manager) EnsureDepositsFresh(ctx context.Context) error {
return m.reconcileDeposits(ctx)
}
// reconcileDeposits fetches all spends to our static addresses from our lnd
// wallet and matches it against the deposits in our memory that we've seen so
// far. It picks the newly identified deposits and starts a state machine per
@ -308,6 +317,11 @@ func (m *Manager) reconcileDeposits(ctx context.Context) error {
"confirmations: %w", err)
}
err = m.syncActiveDeposits(ctx, utxos)
if err != nil {
return fmt.Errorf("unable to sync active deposits: %w", err)
}
newDeposits := m.filterNewDeposits(utxos)
if len(newDeposits) == 0 {
log.Tracef("No new deposits...")
@ -456,6 +470,96 @@ func (m *Manager) updateDepositConfirmations(ctx context.Context,
return nil
}
// syncActiveDeposits reconciles the live active set with lnd's current wallet
// view. Known Deposited records that are visible but inactive become active
// again, and active Deposited records that are no longer wallet-visible are
// removed from the live set. The DB record is left untouched as historical
// evidence that the outpoint was once detected.
func (m *Manager) syncActiveDeposits(ctx context.Context,
utxos []*lnwallet.Utxo) error {
currentUtxos := make(map[wire.OutPoint]struct{}, len(utxos))
for _, utxo := range utxos {
currentUtxos[utxo.OutPoint] = struct{}{}
}
type deactivatedDeposit struct {
outpoint wire.OutPoint
fsm *FSM
}
toActivate := make([]*Deposit, 0, len(utxos))
var toDeactivate []deactivatedDeposit
func() {
m.mu.Lock()
defer m.mu.Unlock()
toDeactivate = make(
[]deactivatedDeposit, 0, len(m.activeDeposits),
)
for _, utxo := range utxos {
deposit, ok := m.deposits[utxo.OutPoint]
if !ok {
continue
}
if _, active := m.activeDeposits[utxo.OutPoint]; active {
continue
}
if !deposit.IsInState(Deposited) {
continue
}
toActivate = append(toActivate, deposit)
}
for outpoint, fsm := range m.activeDeposits {
if _, ok := currentUtxos[outpoint]; ok {
continue
}
if fsm == nil || fsm.deposit == nil {
continue
}
if !fsm.deposit.IsInState(Deposited) {
continue
}
delete(m.activeDeposits, outpoint)
toDeactivate = append(toDeactivate, deactivatedDeposit{
outpoint: outpoint,
fsm: fsm,
})
}
}()
for _, deactivated := range toDeactivate {
deactivated.fsm.Stop()
log.Infof("Removed vanished deposit %v from active set",
deactivated.outpoint)
}
for _, deposit := range toActivate {
if !deposit.IsInState(Deposited) {
continue
}
err := m.startDepositFsm(ctx, deposit)
if err != nil {
m.removeActiveDeposit(deposit.OutPoint)
return err
}
log.Infof("Reactivated visible deposit %v", deposit.OutPoint)
}
return nil
}
// filterNewDeposits filters the given utxos for new deposits that we haven't
// seen before.
func (m *Manager) filterNewDeposits(utxos []*lnwallet.Utxo) []*lnwallet.Utxo {
@ -687,11 +791,41 @@ func (m *Manager) removeActiveDeposit(outpoint wire.OutPoint) {
}
}
// GetAllDeposits returns all active deposits.
// GetAllDeposits returns all known deposits from the database.
func (m *Manager) GetAllDeposits(ctx context.Context) ([]*Deposit, error) {
return m.cfg.Store.AllDeposits(ctx)
}
// GetVisibleDeposits returns deposits that should be exposed through normal
// user-facing views. The database can contain historical Deposited rows whose
// outpoints are no longer present in lnd's current wallet view, for example
// after replacement or reorg. Once the manager has recovered its live cache,
// plain Deposited records are only visible while their outpoint is in the
// active set.
func (m *Manager) GetVisibleDeposits(ctx context.Context) ([]*Deposit, error) {
deposits, err := m.cfg.Store.AllDeposits(ctx)
if err != nil {
return nil, err
}
m.mu.Lock()
defer m.mu.Unlock()
liveCacheReady := len(m.deposits) > 0
filtered := make([]*Deposit, 0, len(deposits))
for _, d := range deposits {
if liveCacheReady && d.IsInState(Deposited) {
if _, ok := m.activeDeposits[d.OutPoint]; !ok {
continue
}
}
filtered = append(filtered, d)
}
return filtered, nil
}
// UpdateDeposit overrides all fields of the deposit with given ID in the store.
func (m *Manager) UpdateDeposit(ctx context.Context, d *Deposit) error {
d.Lock()

View file

@ -3,6 +3,7 @@ package deposit
import (
"context"
"errors"
"strings"
"sync"
"sync/atomic"
"testing"
@ -11,6 +12,7 @@ import (
"github.com/btcsuite/btcd/btcutil"
"github.com/btcsuite/btcd/chaincfg/chainhash"
"github.com/btcsuite/btcd/wire"
"github.com/lightninglabs/loop/fsm"
"github.com/lightninglabs/loop/staticaddr/script"
"github.com/lightninglabs/loop/staticaddr/version"
"github.com/lightninglabs/loop/test"
@ -101,9 +103,18 @@ func TestReconcileDepositsSerialized(t *testing.T) {
}
errCount++
require.ErrorContains(t, err, "unable to start new deposit FSM")
errMsg := err.Error()
require.True(
t,
strings.Contains(
errMsg, "unable to start new deposit FSM",
) || strings.Contains(
errMsg, "unable to sync active deposits",
),
"unexpected error: %v", err,
)
}
require.Equal(t, 1, errCount)
require.Equal(t, 2, errCount)
}
// TestReconcileConfirmedDepositUsesCurrentHeight verifies confirmation heights
@ -232,6 +243,102 @@ func TestUpdateDepositConfirmationsRecomputesPositiveHeight(t *testing.T) {
mockStore.AssertExpectations(t)
}
// TestReconcileDepositsDeactivatesVanishedUnconfirmedDeposit verifies that a
// missing wallet outpoint is removed from the live active set without mutating
// its historical DB state.
func TestReconcileDepositsDeactivatesVanishedUnconfirmedDeposit(t *testing.T) {
ctx := t.Context()
outpoint := wire.OutPoint{
Hash: chainhash.Hash{2},
Index: 7,
}
deposit := &Deposit{
OutPoint: outpoint,
}
deposit.SetState(Deposited)
mockAddressManager := new(mockAddressManager)
mockAddressManager.On(
"ListUnspent", mock.Anything, int32(0), int32(MaxConfs),
).Return([]*lnwallet.Utxo{}, nil)
manager := NewManager(&ManagerConfig{
AddressManager: mockAddressManager,
Store: new(mockStore),
})
manager.deposits[outpoint] = deposit
fsm := &FSM{
deposit: deposit,
stopChan: make(chan struct{}),
quitChan: make(chan struct{}),
}
go func() {
<-fsm.stopChan
close(fsm.quitChan)
}()
manager.activeDeposits[outpoint] = fsm
require.NoError(t, manager.reconcileDeposits(ctx))
require.Equal(t, Deposited, deposit.GetState())
require.Empty(t, manager.activeDeposits)
select {
case <-fsm.quitChan:
case <-time.After(time.Second):
t.Fatal("fsm did not stop after deposit vanished")
}
}
// TestReconcileDepositsDeactivatesVanishedConfirmedDeposit verifies that a
// previously confirmed deposit is also removed from the live active set if it
// vanishes from the wallet view.
func TestReconcileDepositsDeactivatesVanishedConfirmedDeposit(t *testing.T) {
ctx := context.Background()
outpoint := wire.OutPoint{
Hash: chainhash.Hash{9},
Index: 4,
}
deposit := &Deposit{
OutPoint: outpoint,
ConfirmationHeight: 123,
}
deposit.SetState(Deposited)
mockAddressManager := new(mockAddressManager)
mockAddressManager.On(
"ListUnspent", mock.Anything, int32(0), int32(MaxConfs),
).Return([]*lnwallet.Utxo{}, nil)
manager := NewManager(&ManagerConfig{
AddressManager: mockAddressManager,
Store: new(mockStore),
})
manager.deposits[outpoint] = deposit
fsm := &FSM{
deposit: deposit,
stopChan: make(chan struct{}),
quitChan: make(chan struct{}),
}
go func() {
<-fsm.stopChan
close(fsm.quitChan)
}()
manager.activeDeposits[outpoint] = fsm
require.NoError(t, manager.reconcileDeposits(ctx))
require.Equal(t, Deposited, deposit.GetState())
require.EqualValues(t, 123, deposit.ConfirmationHeight)
require.Empty(t, manager.activeDeposits)
select {
case <-fsm.quitChan:
case <-time.After(time.Second):
t.Fatal("fsm did not stop after confirmed deposit vanished")
}
}
// TestAllOutpointsActiveDepositsRejectsDuplicateOutpoints verifies that a
// duplicated selection is rejected before the manager tries to lock the same
// deposit twice.
@ -350,6 +457,199 @@ func TestLockDepositsAllowsReversedConcurrentRequests(t *testing.T) {
}
}
// TestReconcileDepositsReactivatesReappearedDeposit verifies that the same
// outpoint can become active again if lnd reports it after a prior wallet-view
// miss.
func TestReconcileDepositsReactivatesReappearedDeposit(t *testing.T) {
ctx := context.Background()
outpoint := wire.OutPoint{
Hash: chainhash.Hash{3},
Index: 5,
}
deposit := &Deposit{
OutPoint: outpoint,
Value: btcutil.Amount(100_000),
ConfirmationHeight: 77,
}
deposit.SetState(Deposited)
utxo := &lnwallet.Utxo{
OutPoint: outpoint,
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 updateStates []fsm.StateType
mockStore.On(
"UpdateDeposit", mock.Anything, mock.Anything,
).Return(nil).Run(func(args mock.Arguments) {
updatedDeposit := args.Get(1).(*Deposit)
updateStates = append(updateStates, updatedDeposit.state)
if updatedDeposit.isInStateNoLock(Deposited) {
require.Zero(t, updatedDeposit.ConfirmationHeight)
}
})
manager := NewManager(&ManagerConfig{
AddressManager: mockAddressManager,
Store: mockStore,
})
manager.deposits[outpoint] = deposit
// Reconciliation should reactivate the existing record instead of
// creating a second deposit entry for the same outpoint.
require.NoError(t, manager.reconcileDeposits(ctx))
require.Equal(t, Deposited, deposit.GetState())
require.Zero(t, deposit.ConfirmationHeight)
require.Len(t, manager.activeDeposits, 1)
require.Equal(t, []fsm.StateType{Deposited}, updateStates)
}
// TestReconcileDepositsKeepsInactiveOnFSMStartFailure verifies that a failed
// reactivation does not leave memory saying a deposit is active without an FSM.
func TestReconcileDepositsKeepsInactiveOnFSMStartFailure(t *testing.T) {
ctx := context.Background()
outpoint := wire.OutPoint{
Hash: chainhash.Hash{11},
Index: 5,
}
deposit := &Deposit{
OutPoint: outpoint,
Value: btcutil.Amount(100_000),
ConfirmationHeight: 77,
}
deposit.SetState(Deposited)
utxo := &lnwallet.Utxo{
OutPoint: outpoint,
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)(nil), errors.New("fsm init failed"))
var (
updateStates []fsm.StateType
updateHeights []int64
)
mockStore := new(mockStore)
mockStore.On(
"UpdateDeposit", mock.Anything, mock.Anything,
).Return(nil).Run(func(args mock.Arguments) {
updatedDeposit := args.Get(1).(*Deposit)
updateStates = append(updateStates, updatedDeposit.state)
updateHeights = append(
updateHeights, updatedDeposit.ConfirmationHeight,
)
})
manager := NewManager(&ManagerConfig{
AddressManager: mockAddressManager,
Store: mockStore,
})
manager.deposits[outpoint] = deposit
err := manager.reconcileDeposits(ctx)
require.ErrorContains(t, err, "unable to sync active deposits")
require.Equal(t, Deposited, deposit.GetState())
require.Zero(t, deposit.ConfirmationHeight)
require.Empty(t, manager.activeDeposits)
require.Equal(t, []fsm.StateType{Deposited}, updateStates)
require.EqualValues(t, []int64{0}, updateHeights)
}
// TestReconcileDepositsDeactivatesBeforeActivationFailure verifies that a
// failed reactivation of one visible deposit does not leave another vanished
// deposit in the live active set.
func TestReconcileDepositsDeactivatesBeforeActivationFailure(t *testing.T) {
ctx := context.Background()
visibleOutpoint := wire.OutPoint{
Hash: chainhash.Hash{21},
Index: 5,
}
vanishedOutpoint := wire.OutPoint{
Hash: chainhash.Hash{22},
Index: 6,
}
visibleDeposit := &Deposit{
OutPoint: visibleOutpoint,
Value: btcutil.Amount(100_000),
}
visibleDeposit.SetState(Deposited)
vanishedDeposit := &Deposit{
OutPoint: vanishedOutpoint,
Value: btcutil.Amount(100_000),
}
vanishedDeposit.SetState(Deposited)
utxo := &lnwallet.Utxo{
OutPoint: visibleOutpoint,
Value: visibleDeposit.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)(nil), errors.New("fsm init failed"))
manager := NewManager(&ManagerConfig{
AddressManager: mockAddressManager,
Store: new(mockStore),
})
manager.deposits[visibleOutpoint] = visibleDeposit
manager.deposits[vanishedOutpoint] = vanishedDeposit
vanishedFsm := &FSM{
deposit: vanishedDeposit,
stopChan: make(chan struct{}),
quitChan: make(chan struct{}),
}
go func() {
<-vanishedFsm.stopChan
close(vanishedFsm.quitChan)
}()
manager.activeDeposits[vanishedOutpoint] = vanishedFsm
err := manager.reconcileDeposits(ctx)
require.ErrorContains(t, err, "unable to sync active deposits")
require.Empty(t, manager.activeDeposits)
select {
case <-vanishedFsm.quitChan:
case <-time.After(time.Second):
t.Fatal("vanished deposit fsm did not stop")
}
}
// 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.