Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion docs/development/sql-query-dsl.md
Original file line number Diff line number Diff line change
Expand Up @@ -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. |

Expand Down
99 changes: 63 additions & 36 deletions docs/services/selector.md
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand All @@ -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
Expand All @@ -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
Expand All @@ -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
Expand All @@ -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
Expand Down Expand Up @@ -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
Expand Down
92 changes: 90 additions & 2 deletions token/services/selector/sherdlock/fetcher.go
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,7 @@ package sherdlock
import (
"context"
"io"
"math/rand/v2"
"sync"
"sync/atomic"
"time"
Expand Down Expand Up @@ -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 {
Expand Down Expand Up @@ -183,15 +191,84 @@ 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 {
iterators.Iterator[T]
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) {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Nit: this Godoc justifies itself with a bug that never existed on this path.

The comment says returning io.EOF "made every lazy-fetch refetch cycle look like a hard failure instead of 'cache exhausted, fetch more' (#2395)". The code being replaced was collections.NewPermutatedIterator → FSC iterators.Permutate → iterators.Slice, whose Next returns (zero, nil) — i.e. (nil, nil) for a pointer element — at exhaustion (fabric-smart-client@v0.16.0/.../iterators/slice.go:43-46).

The (nil, nil) contract is right; only the stated rationale is wrong, and it will mislead the next reader into thinking a regression was fixed here.

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 {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Buckets are detected on the Quantity string while the ordering comes from the amount column.

This compares the quantity TEXT column, but the run structure it assumes was produced by ORDER BY amount on a separate NUMERIC(78,0) column. The two are written independently — token/services/tokens/storage.go:208-210: Quantity: tta.Tok.Quantity verbatim from the action, Amount: q.ToBigInt().Uint64() derived.

Any non-canonical spelling difference between two equal-valued tokens (e.g. 0x0a vs 0xa, or a decimal vs hex encoding from a different driver/peer version) splits one amount bucket into singleton buckets that are never shuffled against each other, restoring exactly the deterministic hot-spot the bucketing exists to prevent — and silently, since the ordering still looks correct. Comparing the numeric amount (or carrying it on UnspentTokenInWallet) would make the bucket boundary agree with the ORDER BY by construction.

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])
Expand Down Expand Up @@ -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{}{}
}

Expand Down Expand Up @@ -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
Expand Down
6 changes: 6 additions & 0 deletions token/services/selector/sherdlock/fetcher_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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) {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

None of the new behaviour has a test.

The only test-file changes in this PR are mock/stub methods to satisfy the widened interfaces. Nothing asserts that:

  • the emitted SQL contains the NOT EXISTS clause or the ORDER BY amount — TestSpendableTokensIteratorByPreparedReuse only compares the dynamic and prepared paths to each other, so it would pass identically if either clause were dropped;
  • a locked token is actually excluded end-to-end;
  • bucketedIterator.NewPermutation preserves cross-bucket ordering while permuting within a bucket;
  • HasAnySpendableTokens ignores locks.

Per AGENTS.md's testing conventions each of these is a cheap table-driven or SQL-text assertion, and the HasAnySpendableTokens / MaxInputs / Quantity-bucketing issues flagged elsewhere in this review are all the kind of regression such a test would have caught.

args := m.Called(ctx, walletID, typ)

return args.Bool(0), args.Error(1)
}

func TestNewCachedFetcher_WithDefaults(t *testing.T) {
mockDB := new(mockTokenDB)

Expand Down
13 changes: 13 additions & 0 deletions token/services/selector/sherdlock/interfaces.go
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand All @@ -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.
Expand Down
Loading
Loading