diff --git a/docs/release-notes/release-notes-next.md b/docs/release-notes/release-notes-next.md index b2c458152..b604b1d4a 100644 --- a/docs/release-notes/release-notes-next.md +++ b/docs/release-notes/release-notes-next.md @@ -53,6 +53,13 @@ `loopd` failed with `exec format error` on ARM hosts. [Issue #1211](https://github.com/lightninglabs/loop/issues/1211) +* A Loop In whose initiation fails now cancels its probe invoice. A probe + payment that reached the node after the failure was previously held for + hours, until shortly before its expiry. The probe watcher and failure + cleanup share successful cancellation results to avoid duplicate requests. + Cancellation failures are logged, and a failed attempt does not prevent + the other caller from trying. + #### Maintenance * Align the standalone `looprpc` module's OpenTelemetry SDK and OTLP trace diff --git a/loopin.go b/loopin.go index c047aefc1..7d2549b51 100644 --- a/loopin.go +++ b/loopin.go @@ -7,6 +7,7 @@ import ( "errors" "fmt" "sync" + "time" "github.com/btcsuite/btcd/btcec/v2" "github.com/btcsuite/btcd/btcutil" @@ -227,6 +228,21 @@ func newLoopInSwap(globalCtx context.Context, cfg *swapConfig, return nil, err } + // The server may still be routing the probe when its call fails or is + // canceled. It may also answer before the probe reaches us. Once we stop + // waiting for the probe, a later HTLC would stay held until lnd releases + // it shortly before expiry. Cancel the probe invoice on every subsequent + // initiation failure, sharing the watcher's cancellation result. + cancelProbe := newProbeInvoiceCanceler(cfg.lnd.Invoices, probeHash) + initiationSucceeded := false + defer func() { + if !initiationSucceeded { + cancelProbeInvoice( + globalCtx, probeHash, cancelProbe, + ) + } + }() + // Default the HTLC internal key to our sender key. senderInternalPubKey := senderKey @@ -248,7 +264,9 @@ func newLoopInSwap(globalCtx context.Context, cfg *swapConfig, probeWaitCtx, probeWaitCancel := context.WithCancel(globalCtx) // Launch a goroutine to monitor the probe. - probeResult, err := awaitProbe(probeWaitCtx, *cfg.lnd, probeHash) + probeResult, err := awaitProbe( + probeWaitCtx, *cfg.lnd, probeHash, cancelProbe, + ) if err != nil { probeWaitCancel() return nil, fmt.Errorf("probe failed: %v", err) @@ -346,6 +364,8 @@ func newLoopInSwap(globalCtx context.Context, cfg *swapConfig, swap.abandonChan = make(chan struct{}, 1) + initiationSucceeded = true + return &loopInInitResult{ swap: swap, serverMessage: swapResp.serverMessage, @@ -355,7 +375,8 @@ func newLoopInSwap(globalCtx context.Context, cfg *swapConfig, // awaitProbe waits for a probe payment to arrive and cancels it. This is a // workaround for the current lack of multi-path probing. func awaitProbe(ctx context.Context, lnd lndclient.LndServices, - probeHash lntypes.Hash) (chan error, error) { + probeHash lntypes.Hash, + cancelProbe func(context.Context) error) (chan error, error) { // Subscribe to the probe invoice. updateChan, errChan, err := lnd.Invoices.SubscribeSingleInvoice( @@ -381,9 +402,8 @@ func awaitProbe(ctx context.Context, lnd lndclient.LndServices, // Cancel probe invoice so that the // server will know that its probe was // successful. - err := lnd.Invoices.CancelInvoice( + err := cancelProbe( context.WithoutCancel(ctx), - probeHash, ) if err != nil { log.Errorf("Cancel probe "+ @@ -420,6 +440,76 @@ func awaitProbe(ctx context.Context, lnd lndclient.LndServices, return probeResult, nil } +// invoiceCanceler is the invoice RPC needed for probe cancellation. +type invoiceCanceler interface { + CancelInvoice(context.Context, lntypes.Hash) error +} + +// newProbeInvoiceCanceler returns a cancellation function shared by the probe +// watcher and initiation cleanup. It serializes attempts and remembers only +// successful cancellations, so a failed attempt can be retried by the other +// caller. Its state lives only as long as this probe's callers. +func newProbeInvoiceCanceler(invoices invoiceCanceler, + probeHash lntypes.Hash) func(context.Context) error { + + // The gate acts as a mutex protecting canceled and the cancellation RPC, + // while allowing a waiting caller to stop when its context is canceled + // or its deadline expires. + gate := make(chan struct{}, 1) + canceled := false + + return func(ctx context.Context) error { + select { + case gate <- struct{}{}: + case <-ctx.Done(): + return ctx.Err() + } + defer func() { <-gate }() + + if canceled { + return nil + } + + // The context may have expired while waiting, even if the select + // chose the gate. Do not start another RPC with an expired budget. + if err := ctx.Err(); err != nil { + return err + } + + if err := invoices.CancelInvoice(ctx, probeHash); err != nil { + return err + } + + canceled = true + + return nil + } +} + +// probeInvoiceCleanupTimeout bounds both waiting for an in-flight cancellation +// and canceling a probe invoice after a failed swap initiation. +const probeInvoiceCleanupTimeout = 10 * time.Second + +// cancelProbeInvoice cancels the probe invoice of a swap whose initiation +// failed. The cancellation uses its own context, since the initiation's +// context may already be canceled. The shared cancelProbe function skips the +// RPC if either caller has already successfully canceled this probe. All +// errors, including an invoice deleted without a confirmed cancellation, are +// logged and leave a subsequent caller free to try again. +func cancelProbeInvoice(ctx context.Context, probeHash lntypes.Hash, + cancelProbe func(context.Context) error) { + + cleanupCtx, cancel := context.WithTimeout( + context.WithoutCancel(ctx), probeInvoiceCleanupTimeout, + ) + defer cancel() + + if err := cancelProbe(cleanupCtx); err != nil { + log.Warnf("Unable to cancel probe invoice %v: %v", probeHash, + err) + } +} + // resumeLoopInSwap returns a swap object representing a pending swap that has // been restored from the database. func resumeLoopInSwap(_ context.Context, cfg *swapConfig, diff --git a/loopin_test.go b/loopin_test.go index 93d019da8..193406995 100644 --- a/loopin_test.go +++ b/loopin_test.go @@ -3,6 +3,7 @@ package loop import ( "context" "fmt" + "sync/atomic" "testing" "time" @@ -46,6 +47,18 @@ type probeInvoicesMock struct { cancelCtxErr chan error } +// probeInvoiceCancelMock allows tests to control individual cancellation RPCs. +type probeInvoiceCancelMock struct { + cancelInvoice func(context.Context, lntypes.Hash) error +} + +// CancelInvoice invokes the cancellation configured by the test. +func (p *probeInvoiceCancelMock) CancelInvoice(ctx context.Context, + hash lntypes.Hash) error { + + return p.cancelInvoice(ctx, hash) +} + // cancelErrorInvoicesMock is an InvoicesClient that returns a configured // cancellation error. type cancelErrorInvoicesMock struct { @@ -182,6 +195,224 @@ func TestProcessHtlcSpendIgnoresGRPCAlreadySettled(t *testing.T) { require.Equal(t, loopdb.StateFailTimeout, initResult.swap.state) } +// TestLoopInCancelsProbeInvoiceOnInitiationFailure checks that a failed swap +// initiation cancels the probe invoice. The probe watcher stops with the +// initiation call, so nothing else would fail back a probe HTLC that the server +// routes afterwards. This covers a failed call, and a call that succeeds +// without the probe having reached the client, which fails the initiation +// with a probe error. +func TestLoopInCancelsProbeInvoiceOnInitiationFailure(t *testing.T) { + testCases := []struct { + name string + serverErr bool + expectErr string + }{ + { + name: "server error", + serverErr: true, + expectErr: "cannot initiate swap", + }, + { + name: "missing probe", + expectErr: "probe error", + }, + } + + for _, testCase := range testCases { + t.Run(testCase.name, func(t *testing.T) { + defer test.Guard(t)() + + testCtx := newLoopInTestContext(t) + if testCase.serverErr { + testCtx.server.expectedSwapAmt = + testLoopInRequest.Amount + 1 + } else { + testCtx.server.skipProbe = true + } + cfg := newSwapConfig( + &testCtx.lnd.LndServices, testCtx.store, + testCtx.server, nil, + clock.NewTestClock(time.Unix(123, 0)), + ) + + _, err := newLoopInSwap( + context.Background(), cfg, 600, + &testLoopInRequest, + ) + require.ErrorContains(t, err, testCase.expectErr) + + probeSubscription := + <-testCtx.lnd.SingleInvoiceSubcribeChannel + select { + case canceled := <-testCtx.lnd.FailInvoiceChannel: + require.Equal( + t, probeSubscription.Hash, canceled, + ) + + case <-time.After(test.Timeout): + t.Fatal("probe invoice was not canceled") + } + }) + } +} + +// TestLoopInSkipsCanceledProbeOnInitiationFailure checks that a failure after +// the probe has arrived does not repeat the watcher's successful cancellation. +func TestLoopInSkipsCanceledProbeOnInitiationFailure(t *testing.T) { + defer test.Guard(t)() + + testCtx := newLoopInTestContext(t) + cfg := newSwapConfig( + &testCtx.lnd.LndServices, testCtx.store, testCtx.server, nil, + clock.NewTestClock(time.Unix(123, 0)), + ) + + // The server completes the probe, but returns an expiry too far in the + // future for the client to accept the swap. + testCtx.server.height = 600 + MaxLoopInAcceptDelta + _, err := newLoopInSwap( + context.Background(), cfg, 600, &testLoopInRequest, + ) + require.ErrorIs(t, err, ErrExpiryTooFar) + + // The server mock consumed the watcher's cancellation. Deferred cleanup + // must observe its success without sending another cancellation. + select { + case <-testCtx.lnd.FailInvoiceChannel: + t.Fatal("probe invoice cancellation was repeated") + default: + } +} + +// TestProbeInvoiceCancelerRemembersSuccess verifies that only a successful +// cancellation suppresses future RPCs. All failures, including invoice-not-found +// errors from lnd, must remain visible and allow another attempt. +func TestProbeInvoiceCancelerRemembersSuccess(t *testing.T) { + t.Parallel() + + testCases := []struct { + name string + firstErr error + wantCalls int + }{ + { + name: "success", + wantCalls: 1, + }, + { + name: "RPC failure", + firstErr: fmt.Errorf("cancel failed"), + wantCalls: 2, + }, + { + name: "invoice not found", + firstErr: invpkg.ErrInvoiceNotFound, + wantCalls: 2, + }, + { + name: "invoice not found over gRPC", + firstErr: status.Error( + codes.Unknown, invpkg.ErrInvoiceNotFound.Error(), + ), + wantCalls: 2, + }, + } + + for _, testCase := range testCases { + t.Run(testCase.name, func(t *testing.T) { + probeHash := lntypes.Hash{1} + calls := 0 + invoices := &probeInvoiceCancelMock{ + cancelInvoice: func(_ context.Context, + hash lntypes.Hash) error { + + require.Equal(t, probeHash, hash) + calls++ + if calls == 1 { + return testCase.firstErr + } + return nil + }, + } + cancelProbe := newProbeInvoiceCanceler(invoices, probeHash) + ctx := context.Background() + + require.ErrorIs(t, cancelProbe(ctx), testCase.firstErr) + require.NoError(t, cancelProbe(ctx)) + require.NoError(t, cancelProbe(ctx)) + require.Equal(t, testCase.wantCalls, calls) + }) + } +} + +// TestProbeInvoiceCancelerSerializesCalls verifies that concurrent callers +// share a successful cancellation and that a waiting caller's deadline does +// not depend on the in-flight RPC returning. +func TestProbeInvoiceCancelerSerializesCalls(t *testing.T) { + t.Parallel() + + const waiters = 8 + ctx, cancel := context.WithTimeout(context.Background(), test.Timeout) + defer cancel() + + var calls atomic.Int32 + started := make(chan struct{}, waiters+2) + release := make(chan struct{}) + invoices := &probeInvoiceCancelMock{ + cancelInvoice: func(ctx context.Context, _ lntypes.Hash) error { + calls.Add(1) + started <- struct{}{} + select { + case <-release: + return nil + case <-ctx.Done(): + return ctx.Err() + } + }, + } + cancelProbe := newProbeInvoiceCanceler(invoices, lntypes.Hash{1}) + results := make(chan error, waiters+1) + go func() { + results <- cancelProbe(ctx) + }() + + select { + case <-started: + case <-ctx.Done(): + t.Fatal("first cancellation did not start") + } + + // The first RPC stays blocked while the second caller exhausts its + // budget. No second RPC should be started, nor should the first be + // canceled by the waiting caller's timeout. + waitCtx, waitCancel := context.WithTimeout(ctx, 50*time.Millisecond) + defer waitCancel() + require.ErrorIs(t, cancelProbe(waitCtx), context.DeadlineExceeded) + require.EqualValues(t, 1, calls.Load()) + select { + case err := <-results: + t.Fatalf("in-flight cancellation stopped early: %v", err) + default: + } + + for range waiters { + go func() { + results <- cancelProbe(ctx) + }() + } + close(release) + + for range waiters + 1 { + select { + case err := <-results: + require.NoError(t, err) + case <-ctx.Done(): + t.Fatal("cancellation callers did not finish") + } + } + require.EqualValues(t, 1, calls.Load()) +} + // SubscribeSingleInvoice returns the mock's preconfigured channels. func (p *probeInvoicesMock) SubscribeSingleInvoice(_ context.Context, _ lntypes.Hash) (<-chan lndclient.InvoiceUpdate, <-chan error, error) { @@ -220,7 +451,10 @@ func TestAwaitProbeCancelInvoiceUsesLiveContext(t *testing.T) { } ctx, cancel := context.WithCancel(context.Background()) - probeResult, err := awaitProbe(ctx, lnd, lntypes.Hash{1}) + defer cancel() + probeHash := lntypes.Hash{1} + cancelProbe := newProbeInvoiceCanceler(invoices, probeHash) + probeResult, err := awaitProbe(ctx, lnd, probeHash, cancelProbe) require.NoError(t, err) invoices.updateChan <- lndclient.InvoiceUpdate{ diff --git a/server_mock_test.go b/server_mock_test.go index e8bc47277..ca28e6d1b 100644 --- a/server_mock_test.go +++ b/server_mock_test.go @@ -54,6 +54,10 @@ type serverMock struct { swapHash lntypes.Hash prepayHash lntypes.Hash + // skipProbe makes NewLoopInSwap answer without paying the probe + // invoice. + skipProbe bool + // preimagePush is a channel that preimage pushes are sent into. preimagePush chan lntypes.Preimage @@ -190,13 +194,15 @@ func (s *serverMock) NewLoopInSwap(_ context.Context, swapHash lntypes.Hash, // Simulate the server paying the probe invoice and expect the client to // cancel the probe payment. - probeSub := <-s.lnd.SingleInvoiceSubcribeChannel - probeSub.Update <- lndclient.InvoiceUpdate{ - Invoice: lndclient.Invoice{ - State: invpkg.ContractAccepted, - }, + if !s.skipProbe { + probeSub := <-s.lnd.SingleInvoiceSubcribeChannel + probeSub.Update <- lndclient.InvoiceUpdate{ + Invoice: lndclient.Invoice{ + State: invpkg.ContractAccepted, + }, + } + <-s.lnd.FailInvoiceChannel } - <-s.lnd.FailInvoiceChannel resp := &newLoopInResponse{ expiry: s.height + testChargeOnChainCltvDelta,