diff --git a/docs/development/sql-query-dsl.md b/docs/development/sql-query-dsl.md index a9f23d9f73..5fe7302b9a 100644 --- a/docs/development/sql-query-dsl.md +++ b/docs/development/sql-query-dsl.md @@ -15,7 +15,7 @@ adding a store method or a new condition, not at application authors. | :--- | :--- | | `query` | Entry points: `Select()`, `Insert()`, `Update()`, `Delete()`, `Table()`. | | `query/common` | The `Builder` that accumulates SQL text plus bound parameters, and the `Serializable` / `Condition` / `CondInterpreter` contracts. | -| `query/cond` | Condition constructors: `Eq`, `Cmp`, `In`, `InTuple`, `And`, `Or`, `Exists`, `BetweenTimestamps`, … | +| `query/cond` | Condition constructors: `Eq`, `Cmp`, `In`, `InTuple`, `And`, `Or`, `Exists`, `NotExists`, `BetweenTimestamps`, … | | `query/select`, `query/insert`, `query/update`, `query/delete` | Per-statement builders. | | `query/pagination` | Pagination strategies and the interpreter that turns them into `LIMIT`/`OFFSET`/`WHERE` clauses. | diff --git a/docs/services/selector.md b/docs/services/selector.md index 39c18d5ade..1c47834550 100644 --- a/docs/services/selector.md +++ b/docs/services/selector.md @@ -7,7 +7,7 @@ The **Selector Service** (`token/services/selector`) picks the unspent tokens (U The Selector Service is responsible for: * **UTXO Selection**: Finding a set of spendable tokens that cover the total quantity required for a transfer operation. * **Double-Spending Mitigation**: Temporarily locking selected tokens during the transaction assembly phase to prevent multiple concurrent transactions from attempting to spend the same tokens. -* **Candidate Enumeration**: Walking the wallet's candidate tokens in randomized order, locking each one as it is encountered, and stopping as soon as the accumulated amount covers the request. Token amounts do not order or rank the candidates. +* **Candidate Enumeration**: Walking the wallet's candidate tokens, locking each one as it is encountered, and stopping as soon as the accumulated amount covers the request. Under `sherdlock`, candidates already locked by another process are excluded from the query itself, and the remaining candidates are ordered ascending by amount with only same-amount candidates shuffled against each other (see [Token Selection Algorithm](#token-selection-algorithm)). Under `simple`, candidates are walked in database order with no amount ranking. ## Interaction with TTX and Storage @@ -24,8 +24,8 @@ graph LR end subgraph "Selection Logic" - Query[Query Spendable Tokens] - Pick[Take Next Candidate - randomized order] + Query[Query Spendable Tokens - excludes locked, ordered by amount] + Pick[Take Next Candidate - ascending, shuffled within same-amount bucket] Lock[Acquire Temporary Lock] Done[Return Locked Tokens] end @@ -42,7 +42,7 @@ graph LR - **Selector Service**: Creates a selector instance per transaction and orchestrates the Selection Logic steps - **Query Spendable Tokens**: Selector calls the Fetcher to retrieve available tokens - **Fetcher Logic**: Checks cache first (fast path), queries Token Store - TokenDB on cache miss (slow path) -- **Take Next Candidate**: Selector takes the next token from the randomized candidate set; the token's amount plays no part in the choice +- **Take Next Candidate**: Under `sherdlock`, the selector takes the next token from a candidate set ordered ascending by amount, with only same-amount candidates shuffled against each other, and already-locked tokens excluded from the set entirely. Under `simple`, candidates are taken in database order and amount plays no part in the choice. - **Acquire Temporary Lock**: Selector locks each candidate as it is encountered, before it knows whether the request can be covered at all; a candidate already locked by another process is skipped and the loop moves on ## Key Components @@ -52,38 +52,61 @@ The `SelectorManager` is the entry point for obtaining a `Selector` instance anc ### Token Selection Algorithm -Selection is a **randomized greedy first-fit**. It is not configurable, and it is not -amount-aware. `Selector.selectInternal` (`token/services/selector/sherdlock/selector.go`) -does the following: +Selection is a **greedy first-fit**, not configurable. It is amount-aware only to the extent +described below (see [#2395](https://github.com/LFDT-Panurus/panurus/issues/2395) mechanisms +2–3); it is not a smallest-fit or largest-fit strategy. `Selector.selectInternal` +(`token/services/selector/sherdlock/selector.go`) does the following: -1. the candidate tokens of the wallet and token type are enumerated in randomized order, +1. the candidate tokens of the wallet and token type are enumerated — under `sherdlock`, + already-locked candidates are excluded from the query (the anti-join, below) and the + remainder is ordered ascending by amount with same-amount runs shuffled against each other + (the bucketed shuffle, below); under `simple`, candidates are walked in database order, 2. each candidate is locked as it is encountered — a candidate already locked by another - process is skipped; a lock failure wrapping `token.SelectorRateLimited` is a hard abort - (not a skip), + process is skipped, and (`sherdlock` only) blacklisted for the remainder of this `Select` + call so a refetch does not immediately re-attempt and re-lose the same race; a lock failure + wrapping `token.SelectorRateLimited` is a hard abort (not a skip), 3. the amounts of the successfully locked tokens are added up, and 4. the selector returns as soon as the running sum reaches the requested quantity. -A token's amount therefore only decides *when* the loop stops, never *which* candidate is -picked. Two consequences worth planning for: - -* **The number and size of the inputs is not minimized.** A request that a single large - token could have covered may well be funded by several small ones. -* **The result is not deterministic.** The same request against the same wallet can select - a different set of tokens, and a different number of inputs, on each run. - -**The randomization is deliberate.** It is what spreads concurrent selectors of the same -wallet across different candidates: walking a fixed order would make every selector contend -for the same first tokens, driving up lock failures and, with them, the immediate-retry path -that gives up with `token.SelectorSufficientButLockedFunds`, and beyond it the backoff path -that ends in `token.SelectorInsufficientFunds`. - -The shuffle lives in the sherdlock fetcher, not in the selection loop -(`token/services/selector/sherdlock/fetcher.go`): the lazy fetcher wraps the database -iterator in `collections.NewPermutatedIterator`, and the cached fetcher hands out a fresh -permutation of the cached slice on every query. The `simple` driver does **not** shuffle — it -walks the database iterator in the order the token store returns it -(`token/services/selector/simple/selector.go`) — so concurrent selectors under `simple` are -more exposed to colliding on the same leading candidates. +Two consequences worth planning for still hold: + +* **The number and size of the inputs is not minimized.** Ordering ascending by amount + means a request is preferentially funded by several small tokens before a large one is + even considered, which can *increase* the number of inputs relative to a single large + token that could have covered the request alone. +* **The result is not fully deterministic.** Candidates of the same amount are shuffled + against each other, so the same request against the same wallet can still select a + different set of same-amount tokens on each run; the amount ordering across different + amounts, however, is deterministic. + +**`sherdlock`-only: anti-join against locked tokens.** The candidate query excludes any token +currently held by a lock in the `TokenLocks` table (`NOT EXISTS` against `TokenLocks`, added +to `buildSpendableTokensIteratorByQuery` in `token/services/storage/db/sql/common/tokens.go`). +This stops a selector from *starting* a race it is bound to lose; the `INSERT`-based lock +acquisition (below) remains the race-safe backstop, since the anti-join is read-then-act and +therefore not itself race-free. Because the anti-join can hide every remaining token from a +wallet that is not actually out of funds — everything left is simply locked by someone else — +`Selector.selectInternal` disambiguates an empty scan with +`TokenFetcher.HasAnySpendableTokens`, a lock-ignoring existence check, before returning +`token.SelectorInsufficientFunds`. + +**`sherdlock`-only: size-ordered, bucket-shuffled candidates.** Candidates are ordered +ascending by amount (an `ORDER BY` added to the same query), then shuffled only *within* runs +of equal amount — `bucketedIterator.NewPermutation()` in +`token/services/selector/sherdlock/fetcher.go`. A strictly deterministic smallest-fit rule was +deliberately avoided: it would just relocate all contention onto the single smallest token +instead of spreading it. **The shuffle is still deliberate** for the reason it always was: it +spreads concurrent selectors of the same wallet across different same-amount candidates — +walking a fixed order within a bucket would make every selector contend for the same leading +candidate, driving up lock failures and, with them, the immediate-retry path that gives up +with `token.SelectorSufficientButLockedFunds`, and beyond it the backoff path that ends in +`token.SelectorInsufficientFunds`. The bucketing lives in the sherdlock fetcher, not in the +selection loop: the lazy fetcher and the cached fetcher both hand out a fresh +`bucketedIterator` permutation on every query. The `simple` driver does **neither** the +anti-join nor the size ordering — it walks the database iterator in the order the token store +returns it (`token/services/selector/simple/selector.go`), unordered and un-shuffled — so +concurrent selectors under `simple` remain fully exposed to colliding on the same leading +candidates and to starting races against already-locked tokens. **How it works in the flow (see "Selection Logic" subgraph in diagram):** 1. **TTX Request**: TTX Service requests token selection for a transfer operation @@ -105,10 +128,12 @@ more exposed to colliding on the same leading candidates. #### Strategies that are not implemented -Amount-aware strategies — smallest-first, largest-first, First-In-First-Out, or minimizing -the number of inputs — are **not** implemented and cannot be configured. There is no -strategy abstraction in the code and no configuration key that selects one. Making selection -amount-aware is tracked in +Deterministic amount-aware strategies — strict smallest-first, largest-first, +First-In-First-Out, or minimizing the number of inputs — are **not** implemented and cannot +be configured, even though `sherdlock` now orders candidates ascending by amount (see above): +that ordering is bucket-shuffled specifically to avoid becoming a deterministic smallest-fit +rule. There is no strategy abstraction in the code and no configuration key that selects one. +Making selection fully amount-aware (e.g. minimizing input count) is tracked in [issue #2017](https://github.com/LFDT-Panurus/panurus/issues/2017). ### Locking Mechanism @@ -239,7 +264,9 @@ token: It does **not** select a selection algorithm: both drivers walk candidates greedily and stop on first cover, but they diverge in several ways beyond the shuffle: -- `sherdlock` randomizes the candidate order; `simple` walks tokens in database order. +- `sherdlock` orders candidates ascending by amount (shuffling only within same-amount runs) + and excludes already-locked tokens from the query via an anti-join; `simple` walks tokens + in unordered database order and does not exclude locked tokens from its query. - `sherdlock` holds already-acquired locks across immediate retries; `simple` releases all locks between every retry attempt. - `simple` runs a `GetTokens` concurrency check after a successful cover and can return a diff --git a/token/services/selector/sherdlock/fetcher.go b/token/services/selector/sherdlock/fetcher.go index f386666dc1..c94c255962 100644 --- a/token/services/selector/sherdlock/fetcher.go +++ b/token/services/selector/sherdlock/fetcher.go @@ -9,6 +9,7 @@ package sherdlock import ( "context" "io" + "math/rand/v2" "sync" "sync/atomic" "time" @@ -122,6 +123,13 @@ func (f *mixedFetcher) UnspentTokensIteratorBy(ctx context.Context, walletID str return f.lazyFetcher.UnspentTokensIteratorBy(ctx, walletID, currency) } +// HasAnySpendableTokens delegates to the lazy fetcher's underlying DB, since +// this check must never be answered from a cache that may itself be behind +// the anti-join (see TokenDB.HasAnySpendableTokens). +func (f *mixedFetcher) HasAnySpendableTokens(ctx context.Context, walletID string, currency token2.Type) (bool, error) { + return f.lazyFetcher.HasAnySpendableTokens(ctx, walletID, currency) +} + // peekedIterator replays an already-consumed first item before delegating // subsequent Next calls to the wrapped iterator. type peekedIterator[T any] struct { @@ -183,8 +191,20 @@ func (f *lazyFetcher) UnspentTokensIteratorBy(ctx context.Context, walletID stri if err != nil { return nil, err } + defer it.Close() + + items, err := iterators.ReadAllPointers[token2.UnspentTokenInWallet](it) + if err != nil { + return nil, err + } - return collections.NewPermutatedIterator[token2.UnspentTokenInWallet](it) + return newBucketedIterator(items).NewPermutation(), nil +} + +// HasAnySpendableTokens delegates straight to the DB (see +// TokenDB.HasAnySpendableTokens): the lazy fetcher has no cache to consult. +func (f *lazyFetcher) HasAnySpendableTokens(ctx context.Context, walletID string, currency token2.Type) (bool, error) { + return f.tokenDB.HasAnySpendableTokens(ctx, walletID, currency) } type permutatableIterator[T any] interface { @@ -192,6 +212,63 @@ type permutatableIterator[T any] interface { NewPermutation() iterators.Iterator[T] } +// bucketedIterator wraps a slice of tokens already ordered ascending by +// amount (see buildSpendableTokensIteratorByQuery's ORDER BY, #2395 phase +// 4b) and permutes it by shuffling only within contiguous runs of tokens +// with equal Quantity, preserving the size ordering across runs. This +// answers the incident's "why did a 1 CHF request grab a 200 CHF token +// instead of a same-size one" question without introducing a new hot spot: +// a strictly deterministic smallest-fit rule would just relocate all +// contention onto the single smallest token. +type bucketedIterator struct { + items []*token2.UnspentTokenInWallet + pos int +} + +// newBucketedIterator wraps items, which must already be ordered ascending +// by amount, for later shuffling via NewPermutation. +func newBucketedIterator(items []*token2.UnspentTokenInWallet) *bucketedIterator { + return &bucketedIterator{items: items} +} + +// Next returns items in the order they were stored, and (nil, nil) once +// exhausted: sherdlock.Iterator's contract (see selectInternal's t == nil +// refetch branch) signals exhaustion with a nil element and nil error, not +// io.EOF — returning io.EOF here made every lazy-fetch refetch cycle look +// like a hard failure instead of "cache exhausted, fetch more" (#2395). +func (b *bucketedIterator) Next() (*token2.UnspentTokenInWallet, error) { + if b.pos >= len(b.items) { + return nil, nil + } + item := b.items[b.pos] + b.pos++ + + return item, nil +} + +func (b *bucketedIterator) Close() {} + +// NewPermutation returns a fresh iterator over the same items: still +// ascending by amount overall, but with each run of equal-Quantity tokens +// independently shuffled, so equally-good candidates are still randomized +// against each other while the size ordering across runs survives. +func (b *bucketedIterator) NewPermutation() iterators.Iterator[*token2.UnspentTokenInWallet] { + shuffled := make([]*token2.UnspentTokenInWallet, len(b.items)) + copy(shuffled, b.items) + + for start := 0; start < len(shuffled); { + end := start + 1 + for end < len(shuffled) && shuffled[end].Quantity == shuffled[start].Quantity { + end++ + } + bucket := shuffled[start:end] + rand.Shuffle(len(bucket), func(i, j int) { bucket[i], bucket[j] = bucket[j], bucket[i] }) + start = end + } + + return newBucketedIterator(shuffled) +} + type tokenCache interface { Get(key string) (permutatableIterator[*token2.UnspentTokenInWallet], bool) Add(key string, value permutatableIterator[*token2.UnspentTokenInWallet]) @@ -357,7 +434,11 @@ func (f *cachedFetcher) updateCache(ctx context.Context, tokensByKey map[string] // Step 1: Add/update new entries first newKeys := make(map[string]struct{}, len(tokensByKey)) for key, toks := range tokensByKey { - f.cache.Add(key, iterators.Slice(toks)) + // toks arrived from SpendableTokensIteratorBy already ascending by + // amount (#2395 phase 4b) and groupTokensByKey preserves that order + // per key, so bucketedIterator's within-bucket shuffle on + // NewPermutation still shuffles only among equally-good candidates. + f.cache.Add(key, newBucketedIterator(toks)) newKeys[key] = struct{}{} } @@ -399,6 +480,13 @@ func (f *cachedFetcher) UnspentTokensIteratorBy(ctx context.Context, walletID st return collections.NewEmptyIterator[*token2.UnspentTokenInWallet](), nil } +// HasAnySpendableTokens bypasses the cache and asks the DB directly (see +// TokenDB.HasAnySpendableTokens): the cache is itself populated from the +// anti-joined query, so it cannot answer this question. +func (f *cachedFetcher) HasAnySpendableTokens(ctx context.Context, walletID string, currency token2.Type) (bool, error) { + return f.tokenDB.HasAnySpendableTokens(ctx, walletID, currency) +} + // isCacheOverused checks if the cache has been queried too many times since the last refresh. func (f *cachedFetcher) isCacheOverused() bool { return f.queriesResponded.Load() >= f.maxQueriesBeforeRefresh diff --git a/token/services/selector/sherdlock/fetcher_test.go b/token/services/selector/sherdlock/fetcher_test.go index 15df742dde..96b9854547 100644 --- a/token/services/selector/sherdlock/fetcher_test.go +++ b/token/services/selector/sherdlock/fetcher_test.go @@ -44,6 +44,12 @@ func (m *mockTokenDB) SpendableTokensIteratorBy(ctx context.Context, walletID st return args.Get(0).(driver.SpendableTokensIterator), args.Error(1) } +func (m *mockTokenDB) HasAnySpendableTokens(ctx context.Context, walletID string, typ token2.Type) (bool, error) { + args := m.Called(ctx, walletID, typ) + + return args.Bool(0), args.Error(1) +} + func TestNewCachedFetcher_WithDefaults(t *testing.T) { mockDB := new(mockTokenDB) diff --git a/token/services/selector/sherdlock/interfaces.go b/token/services/selector/sherdlock/interfaces.go index 4d3bc2e162..01848b37a8 100644 --- a/token/services/selector/sherdlock/interfaces.go +++ b/token/services/selector/sherdlock/interfaces.go @@ -41,6 +41,10 @@ type TokenLocker interface { //go:generate counterfeiter -o mocks/token_fetcher.go -fake-name FakeTokenFetcher . TokenFetcher type TokenFetcher interface { UnspentTokensIteratorBy(ctx context.Context, walletID string, currency token2.Type) (Iterator[*token2.UnspentTokenInWallet], error) + // HasAnySpendableTokens reports whether the wallet has at least one + // spendable token of the given type, ignoring locks. See TokenDB's + // method of the same name for why the selector needs this. + HasAnySpendableTokens(ctx context.Context, walletID string, currency token2.Type) (bool, error) } // FetcherProvider interface for providing fetcher instances. @@ -55,6 +59,15 @@ type FetcherProvider interface { //go:generate counterfeiter -o mocks/tokendb.go -fake-name FakeTokenDB . TokenDB type TokenDB interface { SpendableTokensIteratorBy(ctx context.Context, walletID string, typ token2.Type) (driver.SpendableTokensIterator, error) + // HasAnySpendableTokens reports whether the wallet has at least one + // spendable token of the given type, ignoring locks. SpendableTokensIteratorBy + // excludes already-locked tokens (#2395, mechanism 3), so an empty result from + // it does not prove the wallet has no funds at all — it may just mean every + // token is momentarily locked by another process. This method answers that + // question directly, without the lock exclusion, so the selector can tell + // "genuinely insufficient funds" apart from "funds exist but are all locked + // right now" (see selector.go's use of it). + HasAnySpendableTokens(ctx context.Context, walletID string, typ token2.Type) (bool, error) } // ConfigProvider interface for configuration provider. diff --git a/token/services/selector/sherdlock/manager_test.go b/token/services/selector/sherdlock/manager_test.go index 6f37e40d9a..b50f6714f6 100644 --- a/token/services/selector/sherdlock/manager_test.go +++ b/token/services/selector/sherdlock/manager_test.go @@ -805,6 +805,7 @@ func TestManager_NewSelector_WithDifferentPrecisions(t *testing.T) { type mockTokenFetcher struct { unspentTokensIteratorByFunc func(ctx context.Context, walletID string, currency token2.Type) (Iterator[*token2.UnspentTokenInWallet], error) + hasAnySpendableTokensFunc func(ctx context.Context, walletID string, currency token2.Type) (bool, error) } func (m *mockTokenFetcher) UnspentTokensIteratorBy(ctx context.Context, walletID string, currency token2.Type) (Iterator[*token2.UnspentTokenInWallet], error) { @@ -815,6 +816,14 @@ func (m *mockTokenFetcher) UnspentTokensIteratorBy(ctx context.Context, walletID return &mockIterator{}, nil } +func (m *mockTokenFetcher) HasAnySpendableTokens(ctx context.Context, walletID string, currency token2.Type) (bool, error) { + if m.hasAnySpendableTokensFunc != nil { + return m.hasAnySpendableTokensFunc(ctx, walletID, currency) + } + + return false, nil +} + type mockLocker struct { lockFunc func(ctx context.Context, tokenID *token2.ID, consumerTxID transaction.ID) error unlockByTxIDFunc func(ctx context.Context, consumerTxID transaction.ID) error diff --git a/token/services/selector/sherdlock/mocks/token_fetcher.go b/token/services/selector/sherdlock/mocks/token_fetcher.go index e3d6e8f6d3..e8921caad5 100644 --- a/token/services/selector/sherdlock/mocks/token_fetcher.go +++ b/token/services/selector/sherdlock/mocks/token_fetcher.go @@ -10,6 +10,21 @@ import ( ) type FakeTokenFetcher struct { + HasAnySpendableTokensStub func(context.Context, string, token.Type) (bool, error) + hasAnySpendableTokensMutex sync.RWMutex + hasAnySpendableTokensArgsForCall []struct { + arg1 context.Context + arg2 string + arg3 token.Type + } + hasAnySpendableTokensReturns struct { + result1 bool + result2 error + } + hasAnySpendableTokensReturnsOnCall map[int]struct { + result1 bool + result2 error + } UnspentTokensIteratorByStub func(context.Context, string, token.Type) (sherdlock.Iterator[*token.UnspentTokenInWallet], error) unspentTokensIteratorByMutex sync.RWMutex unspentTokensIteratorByArgsForCall []struct { @@ -29,6 +44,72 @@ type FakeTokenFetcher struct { invocationsMutex sync.RWMutex } +func (fake *FakeTokenFetcher) HasAnySpendableTokens(arg1 context.Context, arg2 string, arg3 token.Type) (bool, error) { + fake.hasAnySpendableTokensMutex.Lock() + ret, specificReturn := fake.hasAnySpendableTokensReturnsOnCall[len(fake.hasAnySpendableTokensArgsForCall)] + fake.hasAnySpendableTokensArgsForCall = append(fake.hasAnySpendableTokensArgsForCall, struct { + arg1 context.Context + arg2 string + arg3 token.Type + }{arg1, arg2, arg3}) + stub := fake.HasAnySpendableTokensStub + fakeReturns := fake.hasAnySpendableTokensReturns + fake.recordInvocation("HasAnySpendableTokens", []interface{}{arg1, arg2, arg3}) + fake.hasAnySpendableTokensMutex.Unlock() + if stub != nil { + return stub(arg1, arg2, arg3) + } + if specificReturn { + return ret.result1, ret.result2 + } + return fakeReturns.result1, fakeReturns.result2 +} + +func (fake *FakeTokenFetcher) HasAnySpendableTokensCallCount() int { + fake.hasAnySpendableTokensMutex.RLock() + defer fake.hasAnySpendableTokensMutex.RUnlock() + return len(fake.hasAnySpendableTokensArgsForCall) +} + +func (fake *FakeTokenFetcher) HasAnySpendableTokensCalls(stub func(context.Context, string, token.Type) (bool, error)) { + fake.hasAnySpendableTokensMutex.Lock() + defer fake.hasAnySpendableTokensMutex.Unlock() + fake.HasAnySpendableTokensStub = stub +} + +func (fake *FakeTokenFetcher) HasAnySpendableTokensArgsForCall(i int) (context.Context, string, token.Type) { + fake.hasAnySpendableTokensMutex.RLock() + defer fake.hasAnySpendableTokensMutex.RUnlock() + argsForCall := fake.hasAnySpendableTokensArgsForCall[i] + return argsForCall.arg1, argsForCall.arg2, argsForCall.arg3 +} + +func (fake *FakeTokenFetcher) HasAnySpendableTokensReturns(result1 bool, result2 error) { + fake.hasAnySpendableTokensMutex.Lock() + defer fake.hasAnySpendableTokensMutex.Unlock() + fake.HasAnySpendableTokensStub = nil + fake.hasAnySpendableTokensReturns = struct { + result1 bool + result2 error + }{result1, result2} +} + +func (fake *FakeTokenFetcher) HasAnySpendableTokensReturnsOnCall(i int, result1 bool, result2 error) { + fake.hasAnySpendableTokensMutex.Lock() + defer fake.hasAnySpendableTokensMutex.Unlock() + fake.HasAnySpendableTokensStub = nil + if fake.hasAnySpendableTokensReturnsOnCall == nil { + fake.hasAnySpendableTokensReturnsOnCall = make(map[int]struct { + result1 bool + result2 error + }) + } + fake.hasAnySpendableTokensReturnsOnCall[i] = struct { + result1 bool + result2 error + }{result1, result2} +} + func (fake *FakeTokenFetcher) UnspentTokensIteratorBy(arg1 context.Context, arg2 string, arg3 token.Type) (sherdlock.Iterator[*token.UnspentTokenInWallet], error) { fake.unspentTokensIteratorByMutex.Lock() ret, specificReturn := fake.unspentTokensIteratorByReturnsOnCall[len(fake.unspentTokensIteratorByArgsForCall)] diff --git a/token/services/selector/sherdlock/mocks/tokendb.go b/token/services/selector/sherdlock/mocks/tokendb.go index e892d5f231..5d0140a56c 100644 --- a/token/services/selector/sherdlock/mocks/tokendb.go +++ b/token/services/selector/sherdlock/mocks/tokendb.go @@ -11,6 +11,21 @@ import ( ) type FakeTokenDB struct { + HasAnySpendableTokensStub func(context.Context, string, token.Type) (bool, error) + hasAnySpendableTokensMutex sync.RWMutex + hasAnySpendableTokensArgsForCall []struct { + arg1 context.Context + arg2 string + arg3 token.Type + } + hasAnySpendableTokensReturns struct { + result1 bool + result2 error + } + hasAnySpendableTokensReturnsOnCall map[int]struct { + result1 bool + result2 error + } SpendableTokensIteratorByStub func(context.Context, string, token.Type) (driver.SpendableTokensIterator, error) spendableTokensIteratorByMutex sync.RWMutex spendableTokensIteratorByArgsForCall []struct { @@ -30,6 +45,72 @@ type FakeTokenDB struct { invocationsMutex sync.RWMutex } +func (fake *FakeTokenDB) HasAnySpendableTokens(arg1 context.Context, arg2 string, arg3 token.Type) (bool, error) { + fake.hasAnySpendableTokensMutex.Lock() + ret, specificReturn := fake.hasAnySpendableTokensReturnsOnCall[len(fake.hasAnySpendableTokensArgsForCall)] + fake.hasAnySpendableTokensArgsForCall = append(fake.hasAnySpendableTokensArgsForCall, struct { + arg1 context.Context + arg2 string + arg3 token.Type + }{arg1, arg2, arg3}) + stub := fake.HasAnySpendableTokensStub + fakeReturns := fake.hasAnySpendableTokensReturns + fake.recordInvocation("HasAnySpendableTokens", []interface{}{arg1, arg2, arg3}) + fake.hasAnySpendableTokensMutex.Unlock() + if stub != nil { + return stub(arg1, arg2, arg3) + } + if specificReturn { + return ret.result1, ret.result2 + } + return fakeReturns.result1, fakeReturns.result2 +} + +func (fake *FakeTokenDB) HasAnySpendableTokensCallCount() int { + fake.hasAnySpendableTokensMutex.RLock() + defer fake.hasAnySpendableTokensMutex.RUnlock() + return len(fake.hasAnySpendableTokensArgsForCall) +} + +func (fake *FakeTokenDB) HasAnySpendableTokensCalls(stub func(context.Context, string, token.Type) (bool, error)) { + fake.hasAnySpendableTokensMutex.Lock() + defer fake.hasAnySpendableTokensMutex.Unlock() + fake.HasAnySpendableTokensStub = stub +} + +func (fake *FakeTokenDB) HasAnySpendableTokensArgsForCall(i int) (context.Context, string, token.Type) { + fake.hasAnySpendableTokensMutex.RLock() + defer fake.hasAnySpendableTokensMutex.RUnlock() + argsForCall := fake.hasAnySpendableTokensArgsForCall[i] + return argsForCall.arg1, argsForCall.arg2, argsForCall.arg3 +} + +func (fake *FakeTokenDB) HasAnySpendableTokensReturns(result1 bool, result2 error) { + fake.hasAnySpendableTokensMutex.Lock() + defer fake.hasAnySpendableTokensMutex.Unlock() + fake.HasAnySpendableTokensStub = nil + fake.hasAnySpendableTokensReturns = struct { + result1 bool + result2 error + }{result1, result2} +} + +func (fake *FakeTokenDB) HasAnySpendableTokensReturnsOnCall(i int, result1 bool, result2 error) { + fake.hasAnySpendableTokensMutex.Lock() + defer fake.hasAnySpendableTokensMutex.Unlock() + fake.HasAnySpendableTokensStub = nil + if fake.hasAnySpendableTokensReturnsOnCall == nil { + fake.hasAnySpendableTokensReturnsOnCall = make(map[int]struct { + result1 bool + result2 error + }) + } + fake.hasAnySpendableTokensReturnsOnCall[i] = struct { + result1 bool + result2 error + }{result1, result2} +} + func (fake *FakeTokenDB) SpendableTokensIteratorBy(arg1 context.Context, arg2 string, arg3 token.Type) (driver.SpendableTokensIterator, error) { fake.spendableTokensIteratorByMutex.Lock() ret, specificReturn := fake.spendableTokensIteratorByReturnsOnCall[len(fake.spendableTokensIteratorByArgsForCall)] diff --git a/token/services/selector/sherdlock/selector.go b/token/services/selector/sherdlock/selector.go index 13586ab434..093c4dbc8d 100644 --- a/token/services/selector/sherdlock/selector.go +++ b/token/services/selector/sherdlock/selector.go @@ -192,13 +192,25 @@ func (s *Selector) selectInternal(ctx context.Context, owner token.OwnerFilter, return nil, nil, immediateRetries, errors.Wrapf(err, "failed to get tokens for [%s:%s]", owner.ID(), tokenType) } else if t == nil { if !tokensLockedByOthersExist { - return nil, nil, immediateRetries, errors.Wrapf( - token.SelectorInsufficientFunds, - "insufficient funds, only [%s] tokens of type [%s] are available, but [%s] were requested and no other process has any tokens locked", - sum.Decimal(), - tokenType, - quantity.Decimal(), - ) + // The candidate query excludes already-locked tokens (#2395, + // mechanism 3), so an empty scan that never saw a lock conflict is + // ambiguous: it may mean this wallet truly has no more funds, or + // that every remaining token is currently locked by someone else + // and was hidden from us entirely. Disambiguate with a direct, + // lock-ignoring existence check before giving up. + hasAny, hasAnyErr := s.fetcher.HasAnySpendableTokens(ctx, owner.ID(), tokenType) + if hasAnyErr != nil { + return nil, nil, immediateRetries, errors.Wrapf(hasAnyErr, "failed to check for locked tokens for [%s:%s]", owner.ID(), tokenType) + } + if !hasAny { + return nil, nil, immediateRetries, errors.Wrapf( + token.SelectorInsufficientFunds, + "insufficient funds, only [%s] tokens of type [%s] are available, but [%s] were requested and no other process has any tokens locked", + sum.Decimal(), + tokenType, + quantity.Decimal(), + ) + } } if !sawNonBlacklistedCandidate && !blacklisted.Empty() { diff --git a/token/services/selector/testutils/testutils.go b/token/services/selector/testutils/testutils.go index 8d8f96d495..e2c5730860 100644 --- a/token/services/selector/testutils/testutils.go +++ b/token/services/selector/testutils/testutils.go @@ -149,6 +149,13 @@ func (q *MockQueryService) UnspentTokensIteratorBy(_ context.Context, walletID s return &token.UnspentTokensIterator{UnspentTokensIterator: &MockIterator{q, q.cache[walletID], 0}}, nil } +// HasAnySpendableTokens reports whether walletID has at least one cached +// token, mirroring SpendableTokensIteratorBy's ignore-locks semantics (this +// mock has no lock concept at all). +func (q *MockQueryService) HasAnySpendableTokens(_ context.Context, walletID string, _ token2.Type) (bool, error) { + return len(q.cache[walletID]) > 0, nil +} + func (q *MockQueryService) GetTokens(ctx context.Context, inputs ...*token2.ID) ([]*token2.Token, error) { ts := make([]*token2.Token, len(inputs)) for i, input := range inputs { diff --git a/token/services/storage/db/driver/token.go b/token/services/storage/db/driver/token.go index 53a5ec78a7..f8a1fc8d21 100644 --- a/token/services/storage/db/driver/token.go +++ b/token/services/storage/db/driver/token.go @@ -202,6 +202,12 @@ type TokenStore interface { UnspentTokensIteratorBy(ctx context.Context, walletID string, tokenType token.Type) (driver.UnspentTokensIterator, error) // SpendableTokensIteratorBy returns an iterator over all tokens owned solely by the passed wallet identifier and of a given type SpendableTokensIteratorBy(ctx context.Context, walletID string, typ token.Type) (driver.SpendableTokensIterator, error) + // HasAnySpendableTokens reports whether the wallet has at least one spendable + // token of the given type, ignoring any lock currently held on it. Used to + // disambiguate "no funds at all" from "funds exist but are all currently + // locked" when SpendableTokensIteratorBy's anti-join against locked tokens + // (#2395) hides every candidate from the caller. + HasAnySpendableTokens(ctx context.Context, walletID string, typ token.Type) (bool, error) // UnsupportedTokensIteratorBy returns the minimum information for upgrade about the tokens that are not supported UnsupportedTokensIteratorBy(ctx context.Context, walletID string, tokenType token.Type) (driver.UnsupportedTokensIterator, error) // ListUnspentTokensBy returns the list of all tokens owned by the passed identifier of a given type diff --git a/token/services/storage/db/sql/common/tokens.go b/token/services/storage/db/sql/common/tokens.go index 720a7a7373..d014ee106d 100644 --- a/token/services/storage/db/sql/common/tokens.go +++ b/token/services/storage/db/sql/common/tokens.go @@ -40,6 +40,10 @@ type tokenTables struct { Certifications string Requests string TokenSKICleanups string + // TokenLocks is only used by buildSpendableTokensIteratorByQuery's + // anti-join against already-locked tokens (#2395); no other TokenStore + // query touches it. + TokenLocks string } type TokenStore struct { @@ -93,6 +97,7 @@ func NewTokenStoreWithNotifier(readDB, writeDB *sql.DB, tables TableNames, ci co Certifications: tables.Certifications, Requests: tables.Requests, TokenSKICleanups: tables.TokenSKICleanups, + TokenLocks: tables.TokenLocks, }, ci, notifier, nil), nil } @@ -110,6 +115,7 @@ func NewTokenStoreWithNotifierAndCleanup( Certifications: tables.Certifications, Requests: tables.Requests, TokenSKICleanups: tables.TokenSKICleanups, + TokenLocks: tables.TokenLocks, }, ci, notifier, cleanupLeaderFactory), nil } @@ -391,18 +397,47 @@ func (it *dedupedTokenRowsIterator) Next() (*token.UnspentToken, error) { // can compare the dynamic path against a prepared-once path using identical // SQL (see #1919). func buildSpendableTokensIteratorByQuery(db *TokenStore, walletID string, typ token.Type) (string, []any) { + tokenTable := q.Table(db.table.Tokens) + tokenLocksTable := q.Table(db.table.TokenLocks) + return q.Select(). FieldsByName("tx_id", "idx", "token_type", "quantity", "owner_wallet_id"). - From(q.Table(db.table.Tokens)). - Where(HasTokenDetails(driver.QueryTokenDetailsParams{ - WalletID: walletID, - TokenType: typ, - Spendable: driver.SpendableOnly, - LedgerTokenFormats: db.getSupportedTokenFormats(), - }, nil)). + From(tokenTable). + Where(cond.And( + HasTokenDetails(driver.QueryTokenDetailsParams{ + WalletID: walletID, + TokenType: typ, + Spendable: driver.SpendableOnly, + LedgerTokenFormats: db.getSupportedTokenFormats(), + }, nil), + notLocked(tokenTable, tokenLocksTable), + )). + OrderBy(q.Asc(common3.FieldName("amount"))). Format(db.ci) } +// notLocked excludes tokens that currently have a row in TokenLocks, so +// concurrent selectors stop racing to lock a token they can already see is +// held by someone else (#2395, mechanism 3). The row-level INSERT into +// TokenLocks remains the race-safe backstop for the window between this read +// and that INSERT; this anti-join only stops selectors from starting a race +// they are very likely to lose. +// +// This is deliberately not folded into HasTokenDetails: balance and audit +// queries need to see locked tokens too, only the spendable-tokens query +// used by the selector should exclude them. +func notLocked(tokenTable, tokenLocksTable common3.Table) cond.Condition { + return cond.NotExists( + q.Select(). + Fields(common3.FieldName("1")). + From(tokenLocksTable). + Where(cond.And( + cond.Cmp(tokenLocksTable.Field("tx_id"), "=", tokenTable.Field("tx_id")), + cond.Cmp(tokenLocksTable.Field("idx"), "=", tokenTable.Field("idx")), + )), + ) +} + func (db *TokenStore) SpendableTokensIteratorBy(ctx context.Context, walletID string, typ token.Type) (tdriver.SpendableTokensIterator, error) { key := unspentTokensStmtKey(walletID, typ) rows, err := db.spendableTokensStmts.Execute(ctx, db.readDB, key, func() (string, []any, error) { @@ -419,6 +454,44 @@ func (db *TokenStore) SpendableTokensIteratorBy(ctx context.Context, walletID st }), nil } +// buildHasAnySpendableTokensQuery builds the SQL query and args for +// HasAnySpendableTokens without executing it: the same candidate filter as +// buildSpendableTokensIteratorByQuery, minus the notLocked anti-join, capped +// at one row. +func buildHasAnySpendableTokensQuery(db *TokenStore, walletID string, typ token.Type) (string, []any) { + return q.Select(). + Fields(common3.FieldName("1")). + From(q.Table(db.table.Tokens)). + Where(HasTokenDetails(driver.QueryTokenDetailsParams{ + WalletID: walletID, + TokenType: typ, + Spendable: driver.SpendableOnly, + LedgerTokenFormats: db.getSupportedTokenFormats(), + }, nil)). + Limit(1). + Format(db.ci) +} + +// HasAnySpendableTokens reports whether the wallet has at least one +// spendable token of the given type, ignoring locks (see the Godoc on the +// driver.TokenStore method for why the selector needs this). +func (db *TokenStore) HasAnySpendableTokens(ctx context.Context, walletID string, typ token.Type) (bool, error) { + query, args := buildHasAnySpendableTokensQuery(db, walletID, typ) + + rows, err := db.readDB.QueryContext(ctx, query, args...) + if err != nil { + return false, errors.Wrapf(err, "error querying db") + } + defer Close(rows) + + has := rows.Next() + if err := rows.Err(); err != nil { + return false, err + } + + return has, nil +} + // UnspentLedgerTokensIteratorBy returns an iterator over all unspent ledger tokens func (db *TokenStore) UnspentLedgerTokensIteratorBy(ctx context.Context) (tdriver.LedgerTokensIterator, error) { return db.queryLedgerTokens(ctx, driver.QueryTokenDetailsParams{Spendable: driver.Any}) @@ -1518,6 +1591,20 @@ func (db *TokenStore) GetSchema() string { FOREIGN KEY (tx_id, idx) REFERENCES %s ); CREATE INDEX IF NOT EXISTS idx_cleaned_at_%s ON %s ( cleaned_at ); + + -- TokenLocks: created here too (idempotently, alongside + -- TokenLockStore.GetSchema) because buildSpendableTokensIteratorByQuery's + -- notLocked anti-join (#2395, mechanism 3) makes this table a hard + -- dependency of the token store itself, not just of the locker. + CREATE TABLE IF NOT EXISTS %s ( + tx_id TEXT NOT NULL, + idx INT NOT NULL, + consumer_tx_id TEXT NOT NULL, + created_at TIMESTAMPTZ NOT NULL, + PRIMARY KEY(tx_id, idx), + FOREIGN KEY (tx_id, idx) REFERENCES %s + ); + CREATE INDEX IF NOT EXISTS idx_consumer_tx_id_%s ON %s ( consumer_tx_id ); `, db.table.Requests, db.table.Requests, db.table.Requests, db.table.Requests, db.table.Requests, db.table.Requests, db.table.Requests, db.table.Tokens, @@ -1530,6 +1617,8 @@ func (db *TokenStore) GetSchema() string { db.table.PublicParams, db.table.PublicParams, db.table.PublicParams, db.table.Certifications, db.table.Tokens, db.table.TokenSKICleanups, db.table.Tokens, db.table.TokenSKICleanups, db.table.TokenSKICleanups, + db.table.TokenLocks, db.table.Tokens, + db.table.TokenLocks, db.table.TokenLocks, ) } diff --git a/token/services/storage/db/sql/common/tokens_prepared_reuse_test.go b/token/services/storage/db/sql/common/tokens_prepared_reuse_test.go index e2f4fb776f..dbbaa2f38d 100644 --- a/token/services/storage/db/sql/common/tokens_prepared_reuse_test.go +++ b/token/services/storage/db/sql/common/tokens_prepared_reuse_test.go @@ -103,7 +103,8 @@ func TestSpendableTokensIteratorByPreparedReuse(t *testing.T) { store := &TokenStore{ readDB: db, table: tokenTables{ - Tokens: "tokens", + Tokens: "tokens", + TokenLocks: "token_locks", }, ci: stubCondInterpreter{}, spendableTokensStmts: newPreparedStmtHolder[string](), diff --git a/token/services/tokens/mock/token_store.go b/token/services/tokens/mock/token_store.go index d7263275e1..ed24784566 100644 --- a/token/services/tokens/mock/token_store.go +++ b/token/services/tokens/mock/token_store.go @@ -179,6 +179,21 @@ type FakeTokenStore struct { result1 []*token.Token result2 error } + HasAnySpendableTokensStub func(context.Context, string, token.Type) (bool, error) + hasAnySpendableTokensMutex sync.RWMutex + hasAnySpendableTokensArgsForCall []struct { + arg1 context.Context + arg2 string + arg3 token.Type + } + hasAnySpendableTokensReturns struct { + result1 bool + result2 error + } + hasAnySpendableTokensReturnsOnCall map[int]struct { + result1 bool + result2 error + } IsMineStub func(context.Context, string, uint64) (bool, error) isMineMutex sync.RWMutex isMineArgsForCall []struct { @@ -1314,6 +1329,72 @@ func (fake *FakeTokenStore) GetTokensReturnsOnCall(i int, result1 []*token.Token }{result1, result2} } +func (fake *FakeTokenStore) HasAnySpendableTokens(arg1 context.Context, arg2 string, arg3 token.Type) (bool, error) { + fake.hasAnySpendableTokensMutex.Lock() + ret, specificReturn := fake.hasAnySpendableTokensReturnsOnCall[len(fake.hasAnySpendableTokensArgsForCall)] + fake.hasAnySpendableTokensArgsForCall = append(fake.hasAnySpendableTokensArgsForCall, struct { + arg1 context.Context + arg2 string + arg3 token.Type + }{arg1, arg2, arg3}) + stub := fake.HasAnySpendableTokensStub + fakeReturns := fake.hasAnySpendableTokensReturns + fake.recordInvocation("HasAnySpendableTokens", []interface{}{arg1, arg2, arg3}) + fake.hasAnySpendableTokensMutex.Unlock() + if stub != nil { + return stub(arg1, arg2, arg3) + } + if specificReturn { + return ret.result1, ret.result2 + } + return fakeReturns.result1, fakeReturns.result2 +} + +func (fake *FakeTokenStore) HasAnySpendableTokensCallCount() int { + fake.hasAnySpendableTokensMutex.RLock() + defer fake.hasAnySpendableTokensMutex.RUnlock() + return len(fake.hasAnySpendableTokensArgsForCall) +} + +func (fake *FakeTokenStore) HasAnySpendableTokensCalls(stub func(context.Context, string, token.Type) (bool, error)) { + fake.hasAnySpendableTokensMutex.Lock() + defer fake.hasAnySpendableTokensMutex.Unlock() + fake.HasAnySpendableTokensStub = stub +} + +func (fake *FakeTokenStore) HasAnySpendableTokensArgsForCall(i int) (context.Context, string, token.Type) { + fake.hasAnySpendableTokensMutex.RLock() + defer fake.hasAnySpendableTokensMutex.RUnlock() + argsForCall := fake.hasAnySpendableTokensArgsForCall[i] + return argsForCall.arg1, argsForCall.arg2, argsForCall.arg3 +} + +func (fake *FakeTokenStore) HasAnySpendableTokensReturns(result1 bool, result2 error) { + fake.hasAnySpendableTokensMutex.Lock() + defer fake.hasAnySpendableTokensMutex.Unlock() + fake.HasAnySpendableTokensStub = nil + fake.hasAnySpendableTokensReturns = struct { + result1 bool + result2 error + }{result1, result2} +} + +func (fake *FakeTokenStore) HasAnySpendableTokensReturnsOnCall(i int, result1 bool, result2 error) { + fake.hasAnySpendableTokensMutex.Lock() + defer fake.hasAnySpendableTokensMutex.Unlock() + fake.HasAnySpendableTokensStub = nil + if fake.hasAnySpendableTokensReturnsOnCall == nil { + fake.hasAnySpendableTokensReturnsOnCall = make(map[int]struct { + result1 bool + result2 error + }) + } + fake.hasAnySpendableTokensReturnsOnCall[i] = struct { + result1 bool + result2 error + }{result1, result2} +} + func (fake *FakeTokenStore) IsMine(arg1 context.Context, arg2 string, arg3 uint64) (bool, error) { fake.isMineMutex.Lock() ret, specificReturn := fake.isMineReturnsOnCall[len(fake.isMineArgsForCall)]