From a3097946a2dcdeb30eed8b550da5860efe7822ab Mon Sep 17 00:00:00 2001 From: Boris Nagaev Date: Sat, 26 Sep 2026 09:34:04 +0000 Subject: [PATCH] loopin: cancel the probe invoice when initiation fails Before a Loop In, the server probes the route to the client by paying a probe hold invoice whose preimage nobody knows. The client's watcher cancels the invoice once the probe HTLC is accepted, failing the probe back to the server. When the server call returns, the client stops waiting for the probe. If initiation fails, the server may still be routing a probe, and an HTLC that arrives afterwards can remain held until shortly before its CLTV expiry. The same applies when the server answers successfully before the probe reaches the client, which fails initiation with a probe error. Defer probe cancellation after successful invoice creation and skip this cleanup only when initiation succeeds. This covers all subsequent failure paths, including future ones. Cleanup uses a context detached from the initiation context, with a ten-second budget covering both waiting for another cancellation and the RPC itself. Give the watcher and deferred cleanup one shared cancellation function per probe. It serializes attempts with a context-aware gate and remembers only cancellations that returned success. Later calls then skip the RPC, including when lnd has already garbage-collected the canceled invoice. The watcher retains its existing cancellation context and RPC timeout. Failed attempts are logged and leave the other caller free to try within its remaining budget; there is no automatic retry loop. Invoice-notfound errors are not suppressed. They can still occur if lnd canceled and deleted the invoice but its response was lost, or if another caller canceled it. The remembered success is local to this probe's lifetime. Add regression coverage for successful cancellation reuse, failure after the watcher has canceled, retrying failed attempts without suppressing invoice-not-found, concurrent callers, and deadlines while waiting. --- docs/release-notes/release-notes-next.md | 7 + loopin.go | 98 +++++++++- loopin_test.go | 236 ++++++++++++++++++++++- server_mock_test.go | 18 +- 4 files changed, 348 insertions(+), 11 deletions(-) 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,