openchannel: recover deposits at runtime after PSBT finalize failure

If the PSBT finalize step succeeds but the stream fails before
ChanPending, deposits would remain stuck in OpeningChannel until the
next daemon restart. Run the recovery logic immediately so deposits are
resolved without requiring a restart.

Also, add tests for the following edge cases requested in review:

- Reorg: channel tx reorged, UTXOs reappear as unspent, deposits
  return to Deposited state.
- Daemon restart during channel opening: deposits in OpeningChannel
  recovered based on UTXO status (spent → ChannelPublished, unspent →
  Deposited).
- Mempool eviction: tx evicted, UTXOs unspent, deposits return to
  Deposited.
- Mempool rejection: tx never accepted, same recovery as eviction.
- Stream errors: lnd stream fails before PSBT finalize, error returned
  without errPsbtFinalized so deposits can be safely rolled back.
- PSBT finalize then stream abort: finalize succeeds but stream dies
  before ChanPending, error wrapped with errPsbtFinalized so caller
  triggers recovery instead of blind rollback.
- Duplicate outpoints: already covered by TestOpenChannelDuplicateOutpoints.
This commit is contained in:
Slyghtning 2026-02-19 09:56:10 +01:00
parent e710ea5e8f
commit cd377f35f8
No known key found for this signature in database
GPG key ID: F82D456EA023C9BF
2 changed files with 485 additions and 3 deletions

View file

@ -398,9 +398,17 @@ func (m *Manager) OpenChannel(ctx context.Context,
// If the PSBT was already finalized and sent to lnd, the
// funding transaction may have been broadcast. In that case
// we must not roll back the deposits to Deposited as they
// may already be spent on-chain.
if !errors.Is(err, errPsbtFinalized) {
// we must not blindly roll back. Instead, try to recover
// the deposits now so they don't remain stuck in
// OpeningChannel until the next restart.
if errors.Is(err, errPsbtFinalized) {
recoverErr := m.recoverOpeningChannelDeposits(ctx)
if recoverErr != nil {
log.Errorf("failed recovering deposits "+
"after PSBT finalize: %v",
recoverErr)
}
} else {
err2 := m.cfg.DepositManager.TransitionDeposits(
ctx, deposits, fsm.OnError,
deposit.Deposited,

View file

@ -3,8 +3,12 @@ package openchannel
import (
"context"
"errors"
"sync"
"testing"
"time"
"github.com/btcsuite/btcd/btcutil"
"github.com/btcsuite/btcd/chaincfg"
"github.com/btcsuite/btcd/chaincfg/chainhash"
"github.com/btcsuite/btcd/wire"
"github.com/lightninglabs/lndclient"
@ -12,7 +16,10 @@ import (
"github.com/lightninglabs/loop/staticaddr/deposit"
"github.com/lightningnetwork/lnd/lnrpc"
"github.com/lightningnetwork/lnd/lnwallet"
"github.com/lightningnetwork/lnd/lnwallet/chainfee"
"github.com/stretchr/testify/require"
"google.golang.org/grpc"
"google.golang.org/grpc/metadata"
)
type transitionCall struct {
@ -89,6 +96,9 @@ func (m *mockWalletKit) ListUnspent(_ context.Context, _, _ int32,
return m.utxos, nil
}
// TestRecoverOpeningChannelDepositsMixed verifies that recovery correctly
// classifies deposits based on UTXO status: unspent deposits are moved back to
// Deposited and spent deposits are transitioned to ChannelPublished.
func TestRecoverOpeningChannelDepositsMixed(t *testing.T) {
t.Parallel()
@ -135,6 +145,8 @@ func TestRecoverOpeningChannelDepositsMixed(t *testing.T) {
)
}
// TestRecoverOpeningChannelDepositsNoDeposits verifies that recovery is a
// no-op when there are no deposits in the OpeningChannel state.
func TestRecoverOpeningChannelDepositsNoDeposits(t *testing.T) {
t.Parallel()
@ -153,6 +165,8 @@ func TestRecoverOpeningChannelDepositsNoDeposits(t *testing.T) {
require.Empty(t, depositManager.calls)
}
// TestRecoverOpeningChannelDepositsListUnspentError verifies that a
// ListUnspent failure during recovery is propagated to the caller.
func TestRecoverOpeningChannelDepositsListUnspentError(t *testing.T) {
t.Parallel()
@ -176,6 +190,8 @@ func TestRecoverOpeningChannelDepositsListUnspentError(t *testing.T) {
require.Empty(t, depositManager.calls)
}
// TestRecoverOpeningChannelDepositsTransitionError verifies that a transition
// failure when moving unspent deposits back to Deposited is propagated.
func TestRecoverOpeningChannelDepositsTransitionError(t *testing.T) {
t.Parallel()
@ -203,6 +219,213 @@ func TestRecoverOpeningChannelDepositsTransitionError(t *testing.T) {
require.Len(t, depositManager.calls, 1)
}
// TestRecoverAfterReorg simulates a reorg where a channel funding transaction
// was confirmed but then reorged out. After the reorg the deposit UTXOs
// reappear as unspent, so recovery should move all deposits back to Deposited.
func TestRecoverAfterReorg(t *testing.T) {
t.Parallel()
d1 := &deposit.Deposit{OutPoint: testOutPoint(1)}
d2 := &deposit.Deposit{OutPoint: testOutPoint(2)}
depositManager := &mockDepositManager{
openingDeposits: []*deposit.Deposit{d1, d2},
}
walletKit := &mockWalletKit{
utxos: []*lnwallet.Utxo{
// After reorg both UTXOs are unspent again.
{OutPoint: d1.OutPoint},
{OutPoint: d2.OutPoint},
},
}
manager := &Manager{
cfg: &Config{
DepositManager: depositManager,
WalletKit: walletKit,
},
}
err := manager.recoverOpeningChannelDeposits(context.Background())
require.NoError(t, err)
require.Len(t, depositManager.calls, 1)
// Both deposits should transition back to Deposited.
require.Equal(t, fsm.OnError, depositManager.calls[0].event)
require.Equal(
t, deposit.Deposited, depositManager.calls[0].expectedState,
)
require.ElementsMatch(
t,
[]wire.OutPoint{d1.OutPoint, d2.OutPoint},
depositManager.calls[0].outpoints,
)
}
// TestRecoverAfterMempoolEviction simulates the case where the channel funding
// transaction was evicted from the mempool. The deposit UTXOs reappear as
// unspent, so recovery should move them back to Deposited.
func TestRecoverAfterMempoolEviction(t *testing.T) {
t.Parallel()
d := &deposit.Deposit{OutPoint: testOutPoint(1)}
depositManager := &mockDepositManager{
openingDeposits: []*deposit.Deposit{d},
}
walletKit := &mockWalletKit{
utxos: []*lnwallet.Utxo{
// UTXO reappears after mempool eviction.
{OutPoint: d.OutPoint},
},
}
manager := &Manager{
cfg: &Config{
DepositManager: depositManager,
WalletKit: walletKit,
},
}
err := manager.recoverOpeningChannelDeposits(context.Background())
require.NoError(t, err)
require.Len(t, depositManager.calls, 1)
require.Equal(t, fsm.OnError, depositManager.calls[0].event)
require.Equal(
t, deposit.Deposited, depositManager.calls[0].expectedState,
)
}
// TestRecoverAfterMempoolRejection simulates the case where the channel
// funding transaction was rejected from the mempool (e.g. fee too low). The
// deposit UTXOs were never spent, so recovery should move them back to
// Deposited.
func TestRecoverAfterMempoolRejection(t *testing.T) {
t.Parallel()
d1 := &deposit.Deposit{OutPoint: testOutPoint(1)}
d2 := &deposit.Deposit{OutPoint: testOutPoint(2)}
d3 := &deposit.Deposit{OutPoint: testOutPoint(3)}
depositManager := &mockDepositManager{
openingDeposits: []*deposit.Deposit{d1, d2, d3},
}
walletKit := &mockWalletKit{
utxos: []*lnwallet.Utxo{
// All UTXOs still unspent since tx was never accepted.
{OutPoint: d1.OutPoint},
{OutPoint: d2.OutPoint},
{OutPoint: d3.OutPoint},
},
}
manager := &Manager{
cfg: &Config{
DepositManager: depositManager,
WalletKit: walletKit,
},
}
err := manager.recoverOpeningChannelDeposits(context.Background())
require.NoError(t, err)
require.Len(t, depositManager.calls, 1)
require.Equal(t, fsm.OnError, depositManager.calls[0].event)
require.Equal(
t, deposit.Deposited, depositManager.calls[0].expectedState,
)
require.Len(t, depositManager.calls[0].outpoints, 3)
}
// TestRecoverDaemonRestartChannelPublished simulates a daemon restart where
// the channel funding tx was successfully broadcast and the deposit UTXOs are
// all spent. Recovery should move them to ChannelPublished.
func TestRecoverDaemonRestartChannelPublished(t *testing.T) {
t.Parallel()
d1 := &deposit.Deposit{OutPoint: testOutPoint(1)}
d2 := &deposit.Deposit{OutPoint: testOutPoint(2)}
depositManager := &mockDepositManager{
openingDeposits: []*deposit.Deposit{d1, d2},
}
walletKit := &mockWalletKit{
// No UTXOs returned - all deposit outpoints have been spent.
utxos: []*lnwallet.Utxo{},
}
manager := &Manager{
cfg: &Config{
DepositManager: depositManager,
WalletKit: walletKit,
},
}
err := manager.recoverOpeningChannelDeposits(context.Background())
require.NoError(t, err)
require.Len(t, depositManager.calls, 1)
require.Equal(
t, deposit.OnChannelPublished,
depositManager.calls[0].event,
)
require.Equal(
t, deposit.ChannelPublished,
depositManager.calls[0].expectedState,
)
require.ElementsMatch(
t,
[]wire.OutPoint{d1.OutPoint, d2.OutPoint},
depositManager.calls[0].outpoints,
)
}
// TestRecoverChannelPublishedTransitionError verifies that an error
// transitioning deposits to ChannelPublished during recovery is returned.
func TestRecoverChannelPublishedTransitionError(t *testing.T) {
t.Parallel()
d := &deposit.Deposit{OutPoint: testOutPoint(1)}
depositManager := &mockDepositManager{
openingDeposits: []*deposit.Deposit{d},
transitionErrs: map[fsm.EventType]error{
deposit.OnChannelPublished: errors.New(
"transition failed",
),
},
}
walletKit := &mockWalletKit{
// UTXO is spent, so recovery tries ChannelPublished transition.
utxos: []*lnwallet.Utxo{},
}
manager := &Manager{
cfg: &Config{
DepositManager: depositManager,
WalletKit: walletKit,
},
}
err := manager.recoverOpeningChannelDeposits(context.Background())
require.ErrorContains(t, err, "unable to recover spent opening deposits")
}
// TestRecoverGetActiveDepositsError verifies that a failure to fetch opening
// channel deposits is surfaced.
func TestRecoverGetActiveDepositsError(t *testing.T) {
t.Parallel()
depositManager := &mockDepositManager{
getErr: errors.New("db connection lost"),
}
manager := &Manager{
cfg: &Config{
DepositManager: depositManager,
},
}
err := manager.recoverOpeningChannelDeposits(context.Background())
require.ErrorContains(t, err, "unable to fetch opening channel deposits")
}
func testOutPoint(b byte) wire.OutPoint {
return wire.OutPoint{
Hash: chainhash.Hash{b},
@ -210,6 +433,9 @@ func testOutPoint(b byte) wire.OutPoint {
}
}
// TestOpenChannelDuplicateOutpoints verifies that OpenChannel rejects requests
// containing duplicate outpoints, which would cause fee miscalculation and an
// invalid PSBT with the same input listed twice.
func TestOpenChannelDuplicateOutpoints(t *testing.T) {
t.Parallel()
@ -238,6 +464,8 @@ func TestOpenChannelDuplicateOutpoints(t *testing.T) {
require.ErrorContains(t, err, "duplicate outpoint")
}
// TestValidateInitialPsbtFlags verifies that request fields incompatible with
// PSBT funding are rejected early, before any deposits are locked.
func TestValidateInitialPsbtFlags(t *testing.T) {
t.Parallel()
@ -289,6 +517,8 @@ func TestValidateInitialPsbtFlags(t *testing.T) {
}
}
// TestResolveCommitmentType verifies that supported commitment types are
// resolved correctly and unsupported types are rejected.
func TestResolveCommitmentType(t *testing.T) {
t.Parallel()
@ -340,3 +570,247 @@ func TestResolveCommitmentType(t *testing.T) {
})
}
}
// ---------------------------------------------------------------------------
// Mock types for PSBT channel open flow tests.
// ---------------------------------------------------------------------------
// mockLndClient implements lndclient.LightningClient for testing. Embedding
// the interface means unimplemented methods panic if called, which is
// desirable in tests to surface unexpected interactions.
type mockLndClient struct {
lndclient.LightningClient
rawClient lnrpc.LightningClient
mu sync.Mutex
fundingStepIdx int
fundingStepErr error
}
func (m *mockLndClient) RawClientWithMacAuth(
ctx context.Context) (context.Context, time.Duration,
lnrpc.LightningClient) {
return ctx, 0, m.rawClient
}
func (m *mockLndClient) FundingStateStep(_ context.Context,
_ *lnrpc.FundingTransitionMsg) (*lnrpc.FundingStateStepResp, error) {
m.mu.Lock()
defer m.mu.Unlock()
m.fundingStepIdx++
return &lnrpc.FundingStateStepResp{}, m.fundingStepErr
}
// mockRawLnrpcClient implements the raw gRPC lnrpc.LightningClient.
type mockRawLnrpcClient struct {
lnrpc.LightningClient
stream lnrpc.Lightning_OpenChannelClient
openErr error
}
func (m *mockRawLnrpcClient) OpenChannel(_ context.Context,
_ *lnrpc.OpenChannelRequest,
_ ...grpc.CallOption) (lnrpc.Lightning_OpenChannelClient, error) {
return m.stream, m.openErr
}
// mockClientStream implements grpc.ClientStream for embedding in
// mockOpenChanStream.
type mockClientStream struct{}
func (m *mockClientStream) Header() (metadata.MD, error) {
return nil, nil
}
func (m *mockClientStream) Trailer() metadata.MD { return nil }
func (m *mockClientStream) CloseSend() error { return nil }
func (m *mockClientStream) Context() context.Context {
return context.Background()
}
func (m *mockClientStream) SendMsg(_ interface{}) error { return nil }
func (m *mockClientStream) RecvMsg(_ interface{}) error { return nil }
// mockOpenChanStream implements lnrpc.Lightning_OpenChannelClient. It returns
// queued messages from Recv(), then returns finalErr once the queue is
// exhausted.
type mockOpenChanStream struct {
*mockClientStream
mu sync.Mutex
msgs []*lnrpc.OpenStatusUpdate
finalErr error
idx int
}
func (m *mockOpenChanStream) Recv() (*lnrpc.OpenStatusUpdate, error) {
m.mu.Lock()
defer m.mu.Unlock()
if m.idx >= len(m.msgs) {
return nil, m.finalErr
}
msg := m.msgs[m.idx]
m.idx++
return msg, nil
}
// mockWithdrawManager implements the WithdrawalManager interface.
type mockWithdrawManager struct {
tx *wire.MsgTx
psbt []byte
err error
}
func (m *mockWithdrawManager) CreateFinalizedWithdrawalTx(
_ context.Context, _ []*deposit.Deposit,
_ btcutil.Address, _ chainfee.SatPerKWeight, _ int64,
_ lnrpc.CommitmentType) (*wire.MsgTx, []byte, error) {
return m.tx, m.psbt, m.err
}
// testFundingAddress returns a valid regtest P2WPKH address for use in tests.
func testFundingAddress() string {
addr, _ := btcutil.NewAddressWitnessPubKeyHash(
make([]byte, 20), &chaincfg.RegressionNetParams,
)
return addr.EncodeAddress()
}
// ---------------------------------------------------------------------------
// Stream-level tests for the PSBT channel open flow.
// ---------------------------------------------------------------------------
// TestStreamOpenError verifies that when the lnd OpenChannel stream fails to
// open, the error is returned and the shim is cleaned up.
func TestStreamOpenError(t *testing.T) {
t.Parallel()
mockRaw := &mockRawLnrpcClient{
openErr: errors.New("connection refused"),
}
lnClient := &mockLndClient{rawClient: mockRaw}
manager := &Manager{
cfg: &Config{
LightningClient: lnClient,
ChainParams: &chaincfg.RegressionNetParams,
},
}
req := &lnrpc.OpenChannelRequest{
LocalFundingAmount: 100000,
MinConfs: defaultUtxoMinConf,
}
_, err := manager.openChannelPsbt(
context.Background(), req, nil, 0,
)
require.ErrorContains(t, err, "opening stream to server failed")
require.False(t, errors.Is(err, errPsbtFinalized))
// Verify that the shim was canceled via FundingStateStep.
lnClient.mu.Lock()
require.Equal(t, 1, lnClient.fundingStepIdx)
lnClient.mu.Unlock()
}
// TestStreamErrorBeforePsbtFinalize verifies that when the lnd stream returns
// an error before the PSBT is finalized, deposits are NOT wrapped in
// errPsbtFinalized so the caller can safely roll them back.
func TestStreamErrorBeforePsbtFinalize(t *testing.T) {
t.Parallel()
stream := &mockOpenChanStream{
mockClientStream: &mockClientStream{},
finalErr: errors.New("peer disconnected"),
}
mockRaw := &mockRawLnrpcClient{stream: stream}
lnClient := &mockLndClient{rawClient: mockRaw}
manager := &Manager{
cfg: &Config{
LightningClient: lnClient,
ChainParams: &chaincfg.RegressionNetParams,
},
}
req := &lnrpc.OpenChannelRequest{
LocalFundingAmount: 100000,
MinConfs: defaultUtxoMinConf,
}
_, err := manager.openChannelPsbt(
context.Background(), req, nil, 0,
)
require.Error(t, err)
require.False(t, errors.Is(err, errPsbtFinalized))
}
// TestPsbtFinalizeThenStreamAbort verifies that when the PSBT finalize step
// succeeds but the stream dies before ChanPending, the error is wrapped with
// errPsbtFinalized so that the caller knows deposits must not be blindly
// rolled back.
func TestPsbtFinalizeThenStreamAbort(t *testing.T) {
t.Parallel()
fundingAmt := int64(100000)
fundingAddr := testFundingAddress()
stream := &mockOpenChanStream{
mockClientStream: &mockClientStream{},
msgs: []*lnrpc.OpenStatusUpdate{
{
Update: &lnrpc.OpenStatusUpdate_PsbtFund{
PsbtFund: &lnrpc.ReadyForPsbtFunding{
FundingAmount: fundingAmt,
FundingAddress: fundingAddr,
},
},
},
},
finalErr: errors.New("stream died after finalize"),
}
mockRaw := &mockRawLnrpcClient{stream: stream}
lnClient := &mockLndClient{rawClient: mockRaw}
// Provide a minimal transaction that can be serialized.
withdrawMgr := &mockWithdrawManager{
tx: wire.NewMsgTx(2),
psbt: []byte("unsigned-psbt"),
}
manager := &Manager{
cfg: &Config{
LightningClient: lnClient,
WithdrawalManager: withdrawMgr,
ChainParams: &chaincfg.RegressionNetParams,
},
}
req := &lnrpc.OpenChannelRequest{
LocalFundingAmount: fundingAmt,
MinConfs: defaultUtxoMinConf,
}
_, err := manager.openChannelPsbt(
context.Background(), req, nil, 0,
)
require.Error(t, err)
require.True(t, errors.Is(err, errPsbtFinalized))
// FundingStateStep should have been called 3 times: verify, finalize,
// and shim cancel (from defer).
lnClient.mu.Lock()
require.Equal(t, 3, lnClient.fundingStepIdx)
lnClient.mu.Unlock()
}