diff --git a/CLAUDE.md b/CLAUDE.md new file mode 100644 index 0000000..e84f834 --- /dev/null +++ b/CLAUDE.md @@ -0,0 +1,51 @@ +# CLAUDE.md + +Guidance for AI agents working in this repo. Keep changes small, tested, and consistent with existing patterns. + +## What this is + +Multi-chain blockchain transaction indexer (Go 1.25). Watches configured chains, fetches blocks, extracts transfers for monitored addresses, and emits events. Supported chains live in `internal/indexer/` (EVM, Solana, Bitcoin, Tron, Sui, Cosmos, Aptos, TON, XRP, Stellar). + +## Commands + +```bash +make build # -> ./indexer +./indexer index --chain= --catchup --debug # run one chain (name = config key, e.g. solana_mainnet) +./indexer index --catchup # run all enabled chains +make stop # pkill the running indexer +go test ./... # all tests (some hit live RPC / need network) +go test ./pkg/adaptive/ -race # a single package, with race detector +go fmt ./... # format before committing +``` + +- Config path defaults to `configs/config.yaml`. **`configs/config.yaml` is gitignored** (real keys/endpoints live there, local only). The tracked template is `configs/config.example.yaml` — update it with **placeholders** (`${HELIUS_KEY}`), never real keys. +- Running needs NATS + Redis reachable (`nats:` / `redis:` blocks in config). Without them the indexer fails at startup. + +## Architecture + +- `internal/indexer/` — per-chain `Indexer` implementations (`indexer.go` defines the interface: `GetBlock`, `GetBlocks`, `GetBlocksByNumbers`, `GetLatestBlockNumber`, `IsHealthy`). Each parses raw RPC blocks into `types.Transaction` transfers. +- `internal/worker/` — worker modes over an indexer, built once per chain in `factory.go` and shared across modes (`types.go`): `regular` (real-time head), `catchup` (backfill ranges), `rescanner` (retry failed blocks), `manual`, `mempool`. `base.go` holds shared block-handling/emit logic. +- `internal/rpc/` — `Failover[T]` provider pool: health tracking, blacklisting, latency-based rotation, error classification (`analyzeError`). `internal/rpc//` has the concrete RPC clients. +- `pkg/adaptive/` — AIMD concurrency limiter used to pace RPC calls to observed latency/errors. +- `pkg/store/`, `pkg/kvstore/`, `pkg/repository/` — persistence (latest block, failed blocks, catchup ranges). `pkg/events/` — emission. `pkg/ratelimiter/` — shared per-chain RPS limiter. `pkg/common/config/` — config types. + +## Conventions + +- **Write minimal comments.** Prefer self-explanatory code (clear names, small functions) over comments. Only comment non-obvious rationale — a gotcha, a "why", something that would bite the next reader. Never restate what the code plainly does. No doc-comment boilerplate on every function. +- **Adding/changing a chain:** implement `indexer.Indexer` in `internal/indexer/.go`, wire an RPC client in `internal/rpc//`, add a `buildIndexer` in `internal/worker/factory.go`, and a config block. +- **Failover tuning lives in code**, not yaml — `rpc.DefaultFailoverConfig()`. Every `NewFailover` call passes `nil` and gets those defaults. Don't reintroduce a yaml `failover:` block. +- **Solana specifics:** `getBlock` uses `encoding=json` (not `jsonParsed`) — cheaper. The parser resolves accounts by index + base58 data and must append `meta.loadedAddresses` (v0 ALT accounts: static + writable + readonly) via `solanaEffectiveAccountKeys`. Skipped slots are normal (`ErrorTypeBlockNotFound`); never treat them as failed blocks. +- **Free public RPCs cannot sustain Solana getBlock** at slot rate — expect 429s/lag without a keyed node. This is capacity, not code. +- **Errors:** classify recoverable RPC failures in `analyzeError` (rate_limit, timeout, chain_unavailable, ...) so blacklist/cooldown policy is consistent. + +## Testing + +- Unit tests are deterministic and offline; prefer them. Some indexer tests hit live mainnet RPC (network-dependent — a DNS/429 failure there is environmental, not your change). +- Use `-race` for anything with goroutines/locks (e.g. `pkg/adaptive`). +- After edits: `go build ./...` then run the affected package's tests. + +## Git + +- Branch off `main`; never commit directly to it. +- Conventional commits (`feat(scope):`, `fix(scope):`, `refactor:`, `chore:`). +- Stage only the files you changed — do not `git add -A` (unrelated WIP files may be dirty in the tree). diff --git a/configs/config.example.yaml b/configs/config.example.yaml index a334d9f..9f23e51 100644 --- a/configs/config.example.yaml +++ b/configs/config.example.yaml @@ -163,14 +163,19 @@ chains: start_block: 0 poll_interval: "2s" nodes: + # STRONGLY RECOMMENDED: put at least one keyed node first. Free public + # nodes cannot sustain getBlock at Solana's slot rate — expect heavy 429s + # and lag without a paid/free-tier-with-key endpoint. + # - url: "https://mainnet.helius-rpc.com/?api-key=${HELIUS_KEY}" # recommended + # - url: "https://solana-mainnet.g.alchemy.com/v2/${ALCHEMY_KEY}" # optional + # - url: "https://solana-mainnet.g.quicknode.pro/${QUICKNODE_KEY}/" # optional - url: "https://solana-rpc.publicnode.com" - url: "https://api.mainnet.solana.com" - - url: "https://solana.drpc.org" - url: "https://solana.leorpc.com/?api_key=FREE" + - url: "https://solana-mainnet.gateway.tatum.io" - url: "https://solana.api.pocket.network" - url: "https://public.rpc.solanavibestation.com" - # - url: "https://solana-mainnet.g.alchemy.com/v2/${ALCHEMY_KEY}" # optional - # - url: "https://solana-mainnet.g.quicknode.pro/${QUICKNODE_KEY}/" # optional + # NOTE: solana.drpc.org does NOT serve Solana on the free plan — removed. client: timeout: "15s" max_retries: 2 @@ -189,8 +194,13 @@ chains: start_block: 0 poll_interval: "2s" nodes: - - url: "https://api.devnet.solana.com" + # Keyed nodes (Ankr/Chainstack free tiers with a key) lead — verified to + # serve getBlock reliably, unlike the unkeyed public devnet endpoints. + # - url: "https://rpc.ankr.com/solana_devnet/${ANKR_KEY}" + # - url: "https://solana-devnet.core.chainstack.com/${CHAINSTACK_KEY}" - url: "https://solana-devnet.api.onfinality.io/public" + - url: "https://solana-devnet.gateway.tatum.io/" + - url: "https://api.devnet.solana.com" client: timeout: "12s" max_retries: 2 diff --git a/internal/indexer/solana.go b/internal/indexer/solana.go index 979d572..c14236c 100644 --- a/internal/indexer/solana.go +++ b/internal/indexer/solana.go @@ -10,6 +10,7 @@ import ( "github.com/fystack/multichain-indexer/internal/rpc" "github.com/fystack/multichain-indexer/internal/rpc/solana" + "github.com/fystack/multichain-indexer/pkg/adaptive" "github.com/fystack/multichain-indexer/pkg/common/config" "github.com/fystack/multichain-indexer/pkg/common/constant" "github.com/fystack/multichain-indexer/pkg/common/enum" @@ -25,6 +26,8 @@ type SolanaIndexer struct { config config.ChainConfig failover *rpc.Failover[solana.SolanaAPI] pubkeyStore PubkeyStore + // shared across worker modes: one congestion controller per chain. + limiter *adaptive.Limiter } func NewSolanaIndexer( @@ -33,7 +36,25 @@ func NewSolanaIndexer( failover *rpc.Failover[solana.SolanaAPI], pubkeyStore PubkeyStore, ) *SolanaIndexer { - return &SolanaIndexer{chainName: chainName, config: cfg, failover: failover, pubkeyStore: pubkeyStore} + maxConc := cfg.Throttle.Concurrency + if maxConc <= 0 { + maxConc = 1 + } + limiter := adaptive.New(adaptive.Config{ + Max: maxConc, + Min: 1, + LowLatency: 1200 * time.Millisecond, + HighLatency: 2500 * time.Millisecond, + AdjustInterval: time.Second, + GrowStreak: 10, + }) + return &SolanaIndexer{ + chainName: chainName, + config: cfg, + failover: failover, + pubkeyStore: pubkeyStore, + limiter: limiter, + } } func (s *SolanaIndexer) GetName() string { return strings.ToUpper(s.chainName) } @@ -150,36 +171,30 @@ func (s *SolanaIndexer) GetBlocksByNumbers(ctx context.Context, blockNumbers []u return []BlockResult{}, nil } - maxConc := s.config.Throttle.Concurrency - if maxConc <= 0 { - maxConc = 1 - } - results := make([]BlockResult, len(blockNumbers)) eg, egCtx := errgroup.WithContext(ctx) - sem := make(chan struct{}, maxConc) for i, slot := range blockNumbers { i := i slot := slot eg.Go(func() error { - select { - case sem <- struct{}{}: - defer func() { <-sem }() - case <-egCtx.Done(): - return egCtx.Err() + if err := s.limiter.Acquire(egCtx); err != nil { + return err } + defer s.limiter.Release() var ( b *solana.GetBlockResult berr error ) + fetchStart := time.Now() berr = s.failover.ExecuteWithRetry(egCtx, func(c solana.SolanaAPI) error { blk, err := c.GetBlock(egCtx, slot) b = blk return err }) + s.limiter.Observe(time.Since(fetchStart), berr == nil) if berr != nil { results[i] = BlockResult{Number: slot, Error: &Error{ErrorType: ErrorTypeUnknown, Message: berr.Error()}} @@ -464,6 +479,23 @@ func solanaParseTokenTransfer(ix solana.Instruction, accountKeys []solana.Accoun } } +// solanaEffectiveAccountKeys appends v0 ALT accounts (static + writable + readonly) +// so instruction/token-balance indices resolve under encoding=json. +func solanaEffectiveAccountKeys(static []solana.AccountKey, loaded *solana.LoadedAddresses) []solana.AccountKey { + if loaded == nil || (len(loaded.Writable) == 0 && len(loaded.Readonly) == 0) { + return static + } + out := make([]solana.AccountKey, 0, len(static)+len(loaded.Writable)+len(loaded.Readonly)) + out = append(out, static...) + for _, pk := range loaded.Writable { + out = append(out, solana.AccountKey{Pubkey: pk, Writable: true}) + } + for _, pk := range loaded.Readonly { + out = append(out, solana.AccountKey{Pubkey: pk}) + } + return out +} + func (s *SolanaIndexer) extractSolanaTransfers(networkID string, slot uint64, ts uint64, b *solana.GetBlockResult) []types.Transaction { out := make([]types.Transaction, 0) for txIdx, tx := range b.Transactions { @@ -478,7 +510,7 @@ func (s *SolanaIndexer) extractSolanaTransfers(networkID string, slot uint64, ts } txHash := tx.Transaction.Signatures[0] fee := decimal.NewFromInt(int64(tx.Meta.Fee)) - accountKeys := tx.Transaction.Message.AccountKeys + accountKeys := solanaEffectiveAccountKeys(tx.Transaction.Message.AccountKeys, tx.Meta.LoadedAddresses) // Build token-account -> (owner, mint) lookup from token balance metadata. // This isn't used to infer transfers; only to map SPL token accounts to owners/mints. diff --git a/internal/indexer/solana_test.go b/internal/indexer/solana_test.go index 91c455c..fe13de5 100644 --- a/internal/indexer/solana_test.go +++ b/internal/indexer/solana_test.go @@ -6,6 +6,7 @@ import ( "time" "github.com/fystack/multichain-indexer/internal/rpc/solana" + "github.com/fystack/multichain-indexer/pkg/adaptive" "github.com/fystack/multichain-indexer/pkg/common/config" "github.com/fystack/multichain-indexer/pkg/common/constant" "github.com/fystack/multichain-indexer/pkg/common/types" @@ -24,6 +25,7 @@ func newTestSolanaIndexer() *SolanaIndexer { chainName: "solana", config: config.ChainConfig{NetworkId: "solana-mainnet"}, pubkeyStore: nil, // no filtering + limiter: adaptive.New(adaptive.Config{Max: 4}), } } @@ -370,3 +372,29 @@ func TestParseSquadsMultisigTransfer(t *testing.T) { tokenTransfer.FromAddress, tokenTransfer.ToAddress, tokenTransfer.Amount, tokenTransfer.AssetAddress) } + +// TestSolanaEffectiveAccountKeys: v0 ALT accounts append as static+writable+readonly. +func TestSolanaEffectiveAccountKeys(t *testing.T) { + static := []solana.AccountKey{{Pubkey: "S0"}, {Pubkey: "S1"}} + + // No loaded addresses: returns the static slice unchanged (jsonParsed path). + assert.Equal(t, static, solanaEffectiveAccountKeys(static, nil)) + assert.Equal(t, static, solanaEffectiveAccountKeys(static, &solana.LoadedAddresses{})) + + loaded := &solana.LoadedAddresses{ + Writable: []string{"W0", "W1"}, + Readonly: []string{"R0"}, + } + got := solanaEffectiveAccountKeys(static, loaded) + require.Len(t, got, 5) + + pubkeys := make([]string, len(got)) + for i, k := range got { + pubkeys[i] = k.Pubkey + } + assert.Equal(t, []string{"S0", "S1", "W0", "W1", "R0"}, pubkeys, + "order must be static + loaded writable + loaded readonly") + assert.True(t, got[2].Writable, "loaded writable accounts must be marked writable") + assert.True(t, got[3].Writable) + assert.False(t, got[4].Writable, "loaded readonly accounts must not be writable") +} diff --git a/internal/rpc/failover.go b/internal/rpc/failover.go index 2904c72..b2fdf4a 100644 --- a/internal/rpc/failover.go +++ b/internal/rpc/failover.go @@ -2,6 +2,7 @@ package rpc import ( "context" + "errors" "fmt" "math/rand" "sort" @@ -21,16 +22,25 @@ type FailoverConfig struct { ErrorThreshold int ForceRotateThreshold int DefaultTimeout time.Duration + // SlowResponseThreshold rotates away a provider that is slow even on success. + SlowResponseThreshold time.Duration + SlowResponseCooldown time.Duration + // EmergencyRecoveryInterval spaces emergency recoveries when the whole pool + // is blacklisted, to avoid hot recover→fail→recover across callers. + EmergencyRecoveryInterval time.Duration } func DefaultFailoverConfig() FailoverConfig { return FailoverConfig{ - HealthCheckInterval: 30 * time.Second, - EnableBlacklisting: true, - MinActiveProviders: 2, - ErrorThreshold: 5, - ForceRotateThreshold: 3, - DefaultTimeout: 10 * time.Second, + HealthCheckInterval: 30 * time.Second, + EnableBlacklisting: true, + MinActiveProviders: 2, + ErrorThreshold: 5, + ForceRotateThreshold: 3, + DefaultTimeout: 10 * time.Second, + SlowResponseThreshold: 3 * time.Second, + SlowResponseCooldown: 2 * time.Minute, + EmergencyRecoveryInterval: 2 * time.Second, } } @@ -186,10 +196,13 @@ type Failover[T NetworkClient] struct { currentIndex int config FailoverConfig lastHealthCheck time.Time + lastEmergency time.Time metrics *FailoverMetrics logThrottler *LogThrottler } +var errAllProvidersBackoff = errors.New("all providers unavailable, backing off") + // NewFailover creates a new type-safe Failover[T] func NewFailover[T NetworkClient](config *FailoverConfig) *Failover[T] { if config == nil { @@ -202,6 +215,15 @@ func NewFailover[T NetworkClient](config *FailoverConfig) *Failover[T] { if config.ForceRotateThreshold <= 0 { config.ForceRotateThreshold = DefaultFailoverConfig().ForceRotateThreshold } + if config.SlowResponseThreshold <= 0 { + config.SlowResponseThreshold = DefaultFailoverConfig().SlowResponseThreshold + } + if config.SlowResponseCooldown <= 0 { + config.SlowResponseCooldown = DefaultFailoverConfig().SlowResponseCooldown + } + if config.EmergencyRecoveryInterval <= 0 { + config.EmergencyRecoveryInterval = DefaultFailoverConfig().EmergencyRecoveryInterval + } return &Failover[T]{ providers: make([]*Provider, 0), currentIndex: -1, @@ -354,6 +376,10 @@ func (f *Failover[T]) performEmergencyRecoveryLocked() (*Provider, error) { return nil, fmt.Errorf("no available providers") } + if !f.lastEmergency.IsZero() && time.Since(f.lastEmergency) < f.config.EmergencyRecoveryInterval { + return nil, errAllProvidersBackoff + } + var blacklisted []*Provider for _, p := range f.providers { if p.State == StateBlacklisted { @@ -373,6 +399,7 @@ func (f *Failover[T]) performEmergencyRecoveryLocked() (*Provider, error) { first := blacklisted[0] first.Recover() f.currentIndex = 0 + f.lastEmergency = time.Now() f.metrics.IncrementEmergencyRecovery() logger.Info("Emergency recovery", "name", first.Name) @@ -414,9 +441,48 @@ func (f *Failover[T]) executeCore(ctx context.Context, provider *Provider, fn fu f.metrics.IncrementSuccess() provider.Success(elapsed) + f.evaluateSlowSuccess(provider, elapsed) return nil } +// evaluateSlowSuccess blacklists a slow-but-successful provider, unless that +// would drop the available pool below MinActiveProviders. +func (f *Failover[T]) evaluateSlowSuccess(provider *Provider, elapsed time.Duration) { + if !f.config.EnableBlacklisting || f.config.SlowResponseThreshold <= 0 { + return + } + if elapsed <= f.config.SlowResponseThreshold { + return + } + if len(f.GetAvailableProviders()) <= f.config.MinActiveProviders { + if f.logThrottler.ShouldLog(fmt.Sprintf("slow_success_min_%s", provider.Name)) { + logger.Warn("Provider slow but kept to preserve minimum active providers", + "provider", provider.Name, + "latency_ms", elapsed.Milliseconds(), + "threshold_ms", f.config.SlowResponseThreshold.Milliseconds(), + "min_active", f.config.MinActiveProviders, + ) + } + return + } + + if f.logThrottler.ShouldLog(fmt.Sprintf("slow_success_%s", provider.Name)) { + provider.mu.RLock() + providerURL := provider.URL + provider.mu.RUnlock() + logger.Warn("Blacklisting slow provider on successful-but-slow response", + "provider", provider.Name, + "url", providerURL, + "latency_ms", elapsed.Milliseconds(), + "threshold_ms", f.config.SlowResponseThreshold.Milliseconds(), + "cooldown", f.config.SlowResponseCooldown, + ) + } + provider.Blacklist(f.config.SlowResponseCooldown) + f.metrics.IncrementBlacklist() + f.metrics.IncrementErrorType("slow_response") +} + // handleUnhealthyProvider marks provider as unhealthy and blacklists it func (f *Failover[T]) handleUnhealthyProvider(provider *Provider, issue ProviderIssue) { logKey := fmt.Sprintf("switch_%s", provider.Name) @@ -640,6 +706,13 @@ func (f *Failover[T]) analyzeError(err error, elapsed time.Duration) ProviderIss cooldown: 24 * time.Hour, markUnhealthy: true, }, + { + // Node does not serve this chain (e.g. drpc free plan) — permanent. + patterns: []string{"not available on free plan", "upgrade to paid plan", "\"code\":35", "\"code\": 35"}, + reason: "chain_unavailable", + cooldown: 24 * time.Hour, + markUnhealthy: true, + }, { patterns: []string{ "-32701", @@ -692,9 +765,9 @@ func (f *Failover[T]) analyzeError(err error, elapsed time.Duration) ProviderIss } // Check for slow response - if elapsed > 3*time.Second { + if f.config.SlowResponseThreshold > 0 && elapsed > f.config.SlowResponseThreshold { issue.Reason = "slow_response" - issue.Cooldown = 2 * time.Minute + issue.Cooldown = f.config.SlowResponseCooldown issue.MarkUnhealthy = true } diff --git a/internal/rpc/failover_test.go b/internal/rpc/failover_test.go index 0774374..d075def 100644 --- a/internal/rpc/failover_test.go +++ b/internal/rpc/failover_test.go @@ -118,6 +118,21 @@ func TestAnalyzeAndHandleError_RestrictedQuery(t *testing.T) { assert.Equal(t, int64(1), errorsByType["restricted_query"]) } +func TestAnalyzeAndHandleError_ChainUnavailable(t *testing.T) { + f, p := newTestFailover() + + err := fmt.Errorf(`getBlock failed: HTTP 400 Bad Request: {"error":{"message":"chain is not available on free plan, please upgrade to paid plan","code":35}}`) + f.AnalyzeAndHandleError(p, err, 100*time.Millisecond) + + assert.False(t, p.IsAvailable(), "a node that doesn't serve the chain should be blacklisted immediately") + assert.Equal(t, StateBlacklisted, p.State) + // Long cooldown (24h) — the condition is permanent for that node. + assert.True(t, time.Now().Add(23*time.Hour).Before(p.BlacklistedUntil)) + + errorsByType := f.GetMetrics()["errors_by_type"].(map[string]int64) + assert.Equal(t, int64(1), errorsByType["chain_unavailable"]) +} + func TestAnalyzeAndHandleError_ConnectionError(t *testing.T) { f, p := newTestFailover() @@ -228,6 +243,111 @@ func TestExecuteCore_GenericErrorsForceRotateToHealthySibling(t *testing.T) { assert.Equal(t, int64(1), metrics["provider_switches"]) } +func TestEmergencyRecovery_SpacedByInterval(t *testing.T) { + cfg := DefaultFailoverConfig() + cfg.EmergencyRecoveryInterval = 100 * time.Millisecond + f := NewFailover[NetworkClient](&cfg) + + a := newTestProvider("a") + b := newTestProvider("b") + require.NoError(t, f.AddProvider(a)) + require.NoError(t, f.AddProvider(b)) + + // Whole pool down (long blacklist so it doesn't expire mid-test). + a.Blacklist(time.Hour) + b.Blacklist(time.Hour) + + // First call recovers one provider. + p, err := f.GetBestProvider() + require.NoError(t, err) + require.NotNil(t, p) + + // Knock the recovered one back out so the pool is fully down again. + p.Blacklist(time.Hour) + + // Immediately: within the interval -> caller is told to back off, no churn. + _, err = f.GetBestProvider() + require.ErrorIs(t, err, errAllProvidersBackoff) + + // After the interval passes, emergency recovery is allowed again. + time.Sleep(120 * time.Millisecond) + p2, err := f.GetBestProvider() + require.NoError(t, err) + require.NotNil(t, p2) +} + +func TestExecuteCore_SlowSuccessBlacklistsProvider(t *testing.T) { + cfg := DefaultFailoverConfig() + cfg.SlowResponseThreshold = 10 * time.Millisecond + cfg.SlowResponseCooldown = time.Minute + cfg.MinActiveProviders = 2 + f := NewFailover[NetworkClient](&cfg) + + // Three providers so blacklisting one still leaves >= MinActiveProviders. + first := newTestProvider("first") + require.NoError(t, f.AddProvider(first)) + require.NoError(t, f.AddProvider(newTestProvider("second"))) + require.NoError(t, f.AddProvider(newTestProvider("third"))) + + err := f.executeCore(context.Background(), first, func(NetworkClient) error { + time.Sleep(30 * time.Millisecond) // exceeds SlowResponseThreshold + return nil + }) + require.NoError(t, err) + + assert.False(t, first.IsAvailable(), "slow-but-successful provider should be blacklisted") + + metrics := f.GetMetrics() + assert.Equal(t, int64(1), metrics["blacklist_events"]) + errorsByType := metrics["errors_by_type"].(map[string]int64) + assert.Equal(t, int64(1), errorsByType["slow_response"]) + + got, err := f.GetBestProvider() + require.NoError(t, err) + assert.NotEqual(t, first.Name, got.Name, "should rotate off the slow provider") +} + +func TestExecuteCore_SlowSuccessKeepsProviderWhenPoolAtMinimum(t *testing.T) { + cfg := DefaultFailoverConfig() + cfg.SlowResponseThreshold = 10 * time.Millisecond + cfg.SlowResponseCooldown = time.Minute + cfg.MinActiveProviders = 2 + f := NewFailover[NetworkClient](&cfg) + + // Only MinActiveProviders providers: blacklisting would starve the pool. + first := newTestProvider("first") + require.NoError(t, f.AddProvider(first)) + require.NoError(t, f.AddProvider(newTestProvider("second"))) + + err := f.executeCore(context.Background(), first, func(NetworkClient) error { + time.Sleep(30 * time.Millisecond) + return nil + }) + require.NoError(t, err) + + assert.True(t, first.IsAvailable(), "must keep slow provider to preserve minimum active pool") + assert.Equal(t, int64(0), f.GetMetrics()["blacklist_events"]) +} + +func TestExecuteCore_FastSuccessDoesNotBlacklist(t *testing.T) { + cfg := DefaultFailoverConfig() + cfg.SlowResponseThreshold = 500 * time.Millisecond + f := NewFailover[NetworkClient](&cfg) + + first := newTestProvider("first") + require.NoError(t, f.AddProvider(first)) + require.NoError(t, f.AddProvider(newTestProvider("second"))) + require.NoError(t, f.AddProvider(newTestProvider("third"))) + + err := f.executeCore(context.Background(), first, func(NetworkClient) error { + return nil // fast + }) + require.NoError(t, err) + + assert.True(t, first.IsAvailable()) + assert.Equal(t, int64(0), f.GetMetrics()["blacklist_events"]) +} + func TestExecuteCore_TransientGenericErrorsDoNotForceRotate(t *testing.T) { cfg := DefaultFailoverConfig() cfg.ForceRotateThreshold = 3 diff --git a/internal/rpc/solana/client.go b/internal/rpc/solana/client.go index e919d2c..f19aef8 100644 --- a/internal/rpc/solana/client.go +++ b/internal/rpc/solana/client.go @@ -74,8 +74,10 @@ func (c *Client) GetTransaction(ctx context.Context, signature string) (*GetTran } func (c *Client) GetBlock(ctx context.Context, slot uint64) (*GetBlockResult, error) { + // json is smaller/faster than jsonParsed; parser handles it via account + // indices + base58 data + meta.loadedAddresses (see extractSolanaTransfers). cfg := GetBlockConfig{ - Encoding: "jsonParsed", + Encoding: "json", TransactionDetails: "full", Rewards: false, MaxSupportedTransactionVersion: 0, diff --git a/internal/rpc/solana/types.go b/internal/rpc/solana/types.go index 8db5c51..43801fb 100644 --- a/internal/rpc/solana/types.go +++ b/internal/rpc/solana/types.go @@ -1,5 +1,7 @@ package solana +import "encoding/json" + // Minimal JSON-RPC types for Solana getBlock / getSlot. type jsonRPCRequest struct { @@ -52,6 +54,13 @@ type TxnMeta struct { PreTokenBalances []TokenBalance `json:"preTokenBalances"` PostTokenBalances []TokenBalance `json:"postTokenBalances"` InnerInstructions []InnerInstruction `json:"innerInstructions"` + // v0 ALT accounts; present only under encoding=json (jsonParsed pre-merges them). + LoadedAddresses *LoadedAddresses `json:"loadedAddresses"` +} + +type LoadedAddresses struct { + Writable []string `json:"writable"` + Readonly []string `json:"readonly"` } type InnerInstruction struct { @@ -90,6 +99,25 @@ type AccountKey struct { Writable bool `json:"writable"` } +// UnmarshalJSON accepts a bare pubkey string (json) or an object (jsonParsed). +func (a *AccountKey) UnmarshalJSON(data []byte) error { + if len(data) > 0 && data[0] == '"' { + var pubkey string + if err := json.Unmarshal(data, &pubkey); err != nil { + return err + } + a.Pubkey = pubkey + return nil + } + type alias AccountKey + var v alias + if err := json.Unmarshal(data, &v); err != nil { + return err + } + *a = AccountKey(v) + return nil +} + type Instruction struct { ProgramIdIndex uint64 `json:"programIdIndex"` Accounts any `json:"accounts"` diff --git a/internal/worker/base.go b/internal/worker/base.go index a36fb8a..cb01b5a 100644 --- a/internal/worker/base.go +++ b/internal/worker/base.go @@ -54,7 +54,7 @@ type BaseWorker struct { // Stop stops the worker and cleans up internal resources func (bw *BaseWorker) Stop() { bw.cancel() - bw.logger.Info("Worker stopped", "chain", bw.chain.GetName()) + bw.logger.Info("Worker stopped") } // newWorkerWithMode constructs a BaseWorker with the given mode and logger. @@ -161,7 +161,6 @@ func (bw *BaseWorker) handleBlockResult(result indexer.BlockResult) bool { } bw.logger.Error("Failed to process block", - "chain", bw.chain.GetName(), "block", result.Number, "err", result.Error.Message, ) @@ -176,7 +175,6 @@ func (bw *BaseWorker) handleBlockResult(result indexer.BlockResult) bool { if result.Block == nil { bw.logger.Error("Nil block result", - "chain", bw.chain.GetName(), "block", result.Number, ) bw.notifyObserver(result.Number, BlockStatusFailed) @@ -187,7 +185,6 @@ func (bw *BaseWorker) handleBlockResult(result indexer.BlockResult) bool { bw.emitBlock(result.Block) bw.logger.Info("Processed block successfully", - "chain", bw.chain.GetName(), "block", result.Block.Number, ) registry.ClearFailedBlocks(bw.chain.GetName(), []uint64{result.Number}) @@ -224,7 +221,6 @@ func (bw *BaseWorker) emitBlock(block *types.Block) { "direction", types.DirectionIn, "from", inTx.FromAddress, "to", inTx.ToAddress, - "chain", bw.chain.GetName(), "type", inTx.Type, "txhash", inTx.TxHash, "status", inTx.Status, @@ -241,7 +237,6 @@ func (bw *BaseWorker) emitBlock(block *types.Block) { "direction", types.DirectionOut, "from", outTx.FromAddress, "to", outTx.ToAddress, - "chain", bw.chain.GetName(), "type", outTx.Type, "txhash", outTx.TxHash, "status", outTx.Status, @@ -324,7 +319,6 @@ func (bw *BaseWorker) emitUTXOs(block *types.Block) { event.Spent = filteredSpent bw.logger.Info("Emitting UTXO event", - "chain", bw.chain.GetName(), "txhash", event.TxHash, "created", len(event.Created), "spent", len(event.Spent), diff --git a/internal/worker/catchup.go b/internal/worker/catchup.go index fe3f6ac..910ddd0 100644 --- a/internal/worker/catchup.go +++ b/internal/worker/catchup.go @@ -72,7 +72,6 @@ func (cw *CatchupWorker) Start() { } cw.logger.Info("Starting optimized catchup worker", - "chain", cw.chain.GetName(), "ranges", len(cw.blockRanges), "total_blocks", totalBlocks, "parallel_workers", CATCHUP_WORKERS, @@ -109,9 +108,7 @@ func (cw *CatchupWorker) runCatchup() { // No ranges left: stay alive and re-check the store for ranges queued // later, instead of exiting permanently. if len(cw.blockRanges) == 0 { - cw.logger.Debug("No catchup ranges, waiting for new work", - "chain", cw.chain.GetName(), - ) + cw.logger.Debug("No catchup ranges, waiting for new work") select { case <-cw.ctx.Done(): return @@ -130,14 +127,12 @@ func (cw *CatchupWorker) reloadStoredRanges() []blockstore.CatchupRange { progress, err := cw.blockStore.GetCatchupProgress(cw.chain.GetNetworkInternalCode()) if err != nil { cw.logger.Warn("Failed to reload catchup progress while idle", - "chain", cw.chain.GetName(), "error", err, ) return nil } if len(progress) > 0 { cw.logger.Info("Picked up newly queued catchup ranges", - "chain", cw.chain.GetName(), "ranges", len(progress), ) status.EnsureStatusRegistry(cw.statusRegistry).SetCatchupRanges(cw.chain.GetName(), progress) @@ -159,14 +154,12 @@ func (cw *CatchupWorker) loadCatchupProgress() []blockstore.CatchupRange { // Load existing catchup ranges from database (they're already split when saved) if progress, err := cw.blockStore.GetCatchupProgress(cw.chain.GetNetworkInternalCode()); err == nil { cw.logger.Info("Loading existing catchup progress", - "chain", cw.chain.GetName(), "progress_ranges", len(progress), ) ranges = progress registry.SetCatchupRanges(cw.chain.GetName(), progress) } else { cw.logger.Warn("Failed to load catchup progress, will create new range", - "chain", cw.chain.GetName(), "error", err, ) } @@ -181,7 +174,6 @@ func (cw *CatchupWorker) loadCatchupProgress() []blockstore.CatchupRange { } start, end := latest+1, head cw.logger.Info("Creating new catchup range", - "chain", cw.chain.GetName(), "latest_block", latest, "head_block", head, "catchup_start", start, "catchup_end", end, @@ -199,7 +191,6 @@ func (cw *CatchupWorker) loadCatchupProgress() []blockstore.CatchupRange { newRanges, ); err != nil { cw.logger.Error("Failed to batch save catchup ranges", - "chain", cw.chain.GetName(), "count", len(newRanges), "error", err, ) @@ -220,7 +211,6 @@ func (cw *CatchupWorker) splitLargeRange(r blockstore.CatchupRange) []blockstore if len(subRanges) > 1 { cw.logger.Info("Split large catchup range", - "chain", cw.chain.GetName(), "original_range", fmt.Sprintf("%d-%d", r.Start, r.End), "original_size", r.End-r.Start+1, "sub_ranges", len(subRanges), @@ -396,14 +386,12 @@ func (cw *CatchupWorker) saveProgress(r blockstore.CatchupRange, current uint64) defer cw.progressMu.Unlock() registry := status.EnsureStatusRegistry(cw.statusRegistry) cw.logger.Debug("Saving catchup progress", - "chain", cw.chain.GetName(), "range", fmt.Sprintf("%d-%d", r.Start, r.End), "current", current, ) current = min(current, r.End) if err := cw.blockStore.SaveCatchupProgress(cw.chain.GetNetworkInternalCode(), r.Start, r.End, current); err != nil { cw.logger.Warn("Failed to save catchup progress", - "chain", cw.chain.GetName(), "range", fmt.Sprintf("%d-%d", r.Start, r.End), "current", current, "error", err, @@ -429,13 +417,11 @@ func (cw *CatchupWorker) completeRange(r blockstore.CatchupRange) error { registry := status.EnsureStatusRegistry(cw.statusRegistry) cw.logger.Info("Completing catchup range", - "chain", cw.chain.GetName(), "range", fmt.Sprintf("%d-%d", r.Start, r.End), ) if err := cw.blockStore.DeleteCatchupRange(cw.chain.GetNetworkInternalCode(), r.Start, r.End); err != nil { cw.logger.Warn("Failed to delete catchup range", - "chain", cw.chain.GetName(), "range", fmt.Sprintf("%d-%d", r.Start, r.End), "error", err, ) @@ -456,7 +442,6 @@ func (cw *CatchupWorker) completeRange(r blockstore.CatchupRange) error { func (cw *CatchupWorker) Close() error { cw.logger.Info("Closing catchup worker, saving progress...", - "chain", cw.chain.GetName(), "ranges", len(cw.blockRanges), ) @@ -483,7 +468,6 @@ func (cw *CatchupWorker) Close() error { if err := cw.blockStore.SaveCatchupRanges(cw.chain.GetNetworkInternalCode(), rangesToSave); err != nil { cw.logger.Error("Failed to batch save progress on close", - "chain", cw.chain.GetName(), "ranges", len(rangesToSave), "error", err, ) diff --git a/internal/worker/default_db_loader.go b/internal/worker/default_db_loader.go index c93a387..1af060f 100644 --- a/internal/worker/default_db_loader.go +++ b/internal/worker/default_db_loader.go @@ -2,8 +2,10 @@ package worker import ( "context" + "errors" "github.com/fystack/multichain-indexer/pkg/model" + "github.com/jackc/pgx/v5/pgconn" "gorm.io/gorm" ) @@ -28,6 +30,11 @@ func (l *DefaultDBLoader) LoadAddresses(ctx context.Context, params AddressLoade Limit(params.Limit). Find(&rows).Error if err != nil { + // Filtering by an enum value the DB type does not define (22P02): no rows. + var pgErr *pgconn.PgError + if errors.As(err, &pgErr) && pgErr.Code == "22P02" { + return nil, nil + } return nil, err } diff --git a/internal/worker/manual.go b/internal/worker/manual.go index 64ebe62..d3a0450 100644 --- a/internal/worker/manual.go +++ b/internal/worker/manual.go @@ -63,7 +63,7 @@ func NewManualWorker( } func (mw *ManualWorker) Start() { - mw.logger.Info("Starting manual worker", "chain", mw.chain.GetName()) + mw.logger.Info("Starting manual worker") // Periodic metrics mw.executeWithRecovery("manual metrics", func() { @@ -89,14 +89,14 @@ func (mw *ManualWorker) loop() { for { select { case <-ctx.Done(): - mw.logger.Info("Manual worker stopped", "chain", mw.chain.GetName()) + mw.logger.Info("Manual worker stopped") return default: } start, end, err := mw.mbs.GetNextRange(ctx, mw.chain.GetNetworkInternalCode()) if err != nil { - mw.logger.Error("GetNextRange failed", "err", err, "chain", mw.chain.GetName()) + mw.logger.Error("GetNextRange failed", "err", err) time.Sleep(time.Second) continue } @@ -105,7 +105,6 @@ func (mw *ManualWorker) loop() { count, _ := mw.mbs.CountRanges(ctx, mw.chain.GetNetworkInternalCode()) if emptyAttempts >= mw.config.MaxEmptyAttempts { mw.logger.Info("No ranges to process, sleeping", - "chain", mw.chain.GetName(), "sleep", mw.config.EmptySleep, "queued_ranges", count, ) @@ -127,14 +126,13 @@ func (mw *ManualWorker) loop() { func (mw *ManualWorker) handleRange(ctx context.Context, start, end uint64) { mw.logger.Info("Processing range", - "chain", mw.chain.GetName(), "start", start, "end", end, ) results, err := mw.chain.GetBlocks(ctx, start, end, false) if err != nil { - mw.logger.Error("GetBlocks failed", "err", err, "chain", mw.chain.GetName()) + mw.logger.Error("GetBlocks failed", "err", err) time.Sleep(time.Second) return } @@ -147,7 +145,6 @@ func (mw *ManualWorker) handleRange(ctx context.Context, start, end uint64) { } mw.logger.Info("Finished processing", - "chain", mw.chain.GetName(), "start", start, "end", end, "lastSuccess", lastSuccess, @@ -158,7 +155,7 @@ func (mw *ManualWorker) handleRange(ctx context.Context, start, end uint64) { } if lastSuccess >= end { if err := mw.mbs.RemoveRange(ctx, mw.chain.GetNetworkInternalCode(), start, end); err != nil { - mw.logger.Error("RemoveRange failed", "err", err, "chain", mw.chain.GetName()) + mw.logger.Error("RemoveRange failed", "err", err) } } } @@ -166,7 +163,7 @@ func (mw *ManualWorker) handleRange(ctx context.Context, start, end uint64) { func (mw *ManualWorker) logMissingRangesMetric() { ranges, err := mw.mbs.ListRanges(mw.ctx, mw.chain.GetNetworkInternalCode()) if err != nil { - mw.logger.Warn("ListRanges failed", "chain", mw.chain.GetName(), "err", err) + mw.logger.Warn("ListRanges failed", "err", err) return } @@ -180,7 +177,6 @@ func (mw *ManualWorker) logMissingRangesMetric() { } mw.logger.Info("Missing block ranges status", - "chain", mw.chain.GetName(), "status", status, "missing_count", rangeCount, ) diff --git a/internal/worker/mempool.go b/internal/worker/mempool.go index 9bf9f75..c393e26 100644 --- a/internal/worker/mempool.go +++ b/internal/worker/mempool.go @@ -72,7 +72,6 @@ func NewMempoolWorker( // Start begins the mempool polling loop func (mw *MempoolWorker) Start() { mw.logger.Info("Starting mempool worker", - "chain", mw.chain.GetName(), "poll_interval", mw.pollInterval, ) go mw.run(mw.processMempool) @@ -80,13 +79,13 @@ func (mw *MempoolWorker) Start() { // Stop stops the mempool worker func (mw *MempoolWorker) Stop() { - mw.logger.Info("Stopping mempool worker", "chain", mw.chain.GetName()) + mw.logger.Info("Stopping mempool worker") mw.BaseWorker.Stop() } // processMempool polls the mempool for new transactions func (mw *MempoolWorker) processMempool() error { - mw.logger.Debug("Polling mempool", "chain", mw.chain.GetName()) + mw.logger.Debug("Polling mempool") transactions, utxoEvents, err := mw.btcIndexer.GetMempoolTransactions(mw.ctx) if err != nil { diff --git a/internal/worker/regular.go b/internal/worker/regular.go index ec93d59..afaf2e8 100644 --- a/internal/worker/regular.go +++ b/internal/worker/regular.go @@ -71,7 +71,6 @@ func NewRegularWorker( func (rw *RegularWorker) Start() { rw.logger.Info("Starting regular worker", - "chain", rw.chain.GetName(), "start_block", rw.currentBlock, ) rw.persistTicker = time.NewTicker(blockHashPersistInterval) @@ -122,7 +121,6 @@ func (rw *RegularWorker) processRegularBlocks() error { end := min(start+uint64(rw.config.Throttle.BatchSize)-1, latest) startTime := time.Now() rw.logger.Info("Processing range", - "chain", rw.chain.GetName(), "start", start, "end", end, "size", end-start+1, ) @@ -140,7 +138,6 @@ func (rw *RegularWorker) processRegularBlocks() error { rw.updateHeadStatus(latest, indexedAt) rw.logger.Info("Processed latest blocks", - "chain", rw.chain.GetName(), "start", start, "end", end, "elapsed", time.Since(startTime), "last_success", lastSuccess, @@ -165,6 +162,14 @@ func (rw *RegularWorker) processBatch( } for _, res := range results { + // Skipped slots are normal on Solana: advance past them, don't fail them. + if rw.isSolanaSkippedSlot(res) { + rw.notifyObserver(res.Number, BlockStatusNotFound) + if res.Number > lastSuccess { + lastSuccess = res.Number + } + continue + } if rw.handleBlockResult(res) { lastSuccess = res.Number lastSuccessHash = res.Block.Hash @@ -173,6 +178,12 @@ func (rw *RegularWorker) processBatch( return lastSuccess, lastSuccessHash, false, nil } +func (rw *RegularWorker) isSolanaSkippedSlot(res indexer.BlockResult) bool { + return res.Error != nil && + res.Error.ErrorType == indexer.ErrorTypeBlockNotFound && + rw.chain.GetNetworkType() == enum.NetworkTypeSol +} + // commitProgress advances currentBlock past the last indexed block, persisting // the checkpoint and its hash. Returns the indexing timestamp, or zero if no new // block was indexed. @@ -200,13 +211,12 @@ func (rw *RegularWorker) determineStartingBlock() uint64 { chainLatest, chainErr := rw.getLatestBlockWithRetry() if chainErr != nil { rw.logger.Warn("Chain RPC failed, resuming from KV latest", - "chain", rw.chain.GetName(), "kvLatest", kvLatest) + "kvLatest", kvLatest) return kvLatest } if chainLatest > kvLatest { ranges := rw.queueCatchupRanges(kvLatest+1, chainLatest) rw.logger.Info("Queued catchup ranges", - "chain", rw.chain.GetName(), "gap", fmt.Sprintf("%d-%d", kvLatest+1, chainLatest), "ranges_created", len(ranges), ) @@ -216,7 +226,7 @@ func (rw *RegularWorker) determineStartingBlock() uint64 { if kvErr != nil { rw.logger.Error("Block store unavailable, starting from chain head", - "chain", rw.chain.GetName(), "error", kvErr) + "error", kvErr) } return rw.waitForChainHead() } @@ -230,7 +240,6 @@ func (rw *RegularWorker) queueCatchupRanges(start, end uint64) []blockstore.Catc if err := rw.blockStore.SaveCatchupRanges(rw.chain.GetNetworkInternalCode(), ranges); err != nil { rw.logger.Error("Failed to save catchup ranges", - "chain", rw.chain.GetName(), "count", len(ranges), "error", err, ) @@ -263,7 +272,7 @@ func (rw *RegularWorker) waitForChainHead() uint64 { return latest } rw.logger.Warn("Waiting for chain head before starting", - "chain", rw.chain.GetName(), "error", err) + "error", err) select { case <-rw.ctx.Done(): return 0 @@ -292,7 +301,6 @@ func (rw *RegularWorker) detectAndHandleReorg(res *indexer.BlockResult) (bool, e reorgStart = prevNum - rollbackWindow } rw.logger.Warn("Reorg detected; rolling back", - "chain", rw.chain.GetName(), "at_block", prevNum, "expected_parent", storedHash, "actual_parent", res.Block.ParentHash, @@ -359,7 +367,6 @@ func (rw *RegularWorker) loadBlockHashes() { } rw.blockHashes = hashes rw.logger.Info("Loaded persisted block hashes", - "chain", rw.chain.GetName(), "count", len(hashes), ) } @@ -383,7 +390,6 @@ func (rw *RegularWorker) flushBlockHashes() { } if err := rw.blockStore.SaveBlockHashes(rw.chain.GetNetworkInternalCode(), rw.blockHashes); err != nil { rw.logger.Error("Failed to persist block hashes", - "chain", rw.chain.GetName(), "error", err, ) return @@ -407,7 +413,6 @@ func (rw *RegularWorker) skipAheadIfLagging(latest uint64) bool { skipEnd := latest - 1 rw.logger.Warn("Lag threshold exceeded, skipping ahead to chain head", - "chain", rw.chain.GetName(), "current_block", rw.currentBlock, "chain_head", latest, "lag", latest-rw.currentBlock, @@ -422,7 +427,6 @@ func (rw *RegularWorker) skipAheadIfLagging(latest uint64) bool { rw.clearBlockHashes() rw.logger.Info("Skip-ahead complete, queued catchup ranges", - "chain", rw.chain.GetName(), "new_current", rw.currentBlock, "catchup_ranges", len(ranges), ) diff --git a/internal/worker/regular_test.go b/internal/worker/regular_test.go index 8852084..222c587 100644 --- a/internal/worker/regular_test.go +++ b/internal/worker/regular_test.go @@ -119,6 +119,46 @@ func TestRegularWorkerProcessRegularBlocksMarksUnresolvedGapFailed(t *testing.T) require.Equal(t, []uint64{100, 100}, chain.getBlockCalls) } +func TestRegularWorkerProcessRegularBlocksSkipsSolanaSkippedSlot(t *testing.T) { + t.Parallel() + + chain := &stubIndexer{ + name: "solana", + internalCode: "sol", + networkType: enum.NetworkTypeSol, + latest: 102, + getBlocksFunc: func(context.Context, uint64, uint64, bool) ([]indexer.BlockResult, error) { + return []indexer.BlockResult{ + { + Number: 100, + Block: &types.Block{Number: 100, Hash: "h100", ParentHash: "h099"}, + }, + { + // Skipped slot: normal on Solana, must not be marked failed. + Number: 101, + Error: &indexer.Error{ErrorType: indexer.ErrorTypeBlockNotFound, Message: "block not found (skipped slot?)"}, + }, + { + Number: 102, + Block: &types.Block{Number: 102, Hash: "h102", ParentHash: "h101"}, + }, + }, nil + }, + } + store := &stubBlockStore{} + rw := newTestRegularWorker(chain, store, 100, 3) + + err := rw.processRegularBlocks() + require.NoError(t, err) + // currentBlock advances past the skipped slot to the next unindexed slot. + require.Equal(t, uint64(103), rw.currentBlock) + require.Equal(t, []uint64{102}, store.savedLatest) + // The skipped slot must NOT be persisted as a failed block. + require.Empty(t, store.failedBlocks) + // No single-block recovery is attempted for a skipped slot. + require.Empty(t, chain.getBlockCalls) +} + func TestBaseWorkerExecuteRecoverableConvertsPanicToError(t *testing.T) { t.Parallel() diff --git a/internal/worker/rescanner.go b/internal/worker/rescanner.go index 72ff67e..612a104 100644 --- a/internal/worker/rescanner.go +++ b/internal/worker/rescanner.go @@ -71,7 +71,6 @@ func NewRescannerWorker( func (rw *RescannerWorker) Start() { rw.logger.Info("Starting rescanner worker", - "chain", rw.chain.GetName(), "interval", rw.interval, "maxRetries", rw.maxRetries, ) @@ -157,7 +156,7 @@ func (rw *RescannerWorker) incrementRetry(block uint64) { delete(rw.failedBlocks, block) rw.addRemove(block) rw.logger.Error("Max retries reached; giving up", - "chain", rw.chain.GetName(), "block", block) + "block", block) } else { rw.failedBlocks[block] = count + 1 rw.addSave(block) @@ -184,7 +183,7 @@ func (rw *RescannerWorker) processRescan() error { time.Sleep(rw.interval) return nil } - rw.logger.Info("Got blocks for rescan", "chain", rw.chain.GetName(), "blocks", len(blocks)) + rw.logger.Info("Got blocks for rescan", "blocks", len(blocks)) return rw.processBatch(blocks) } @@ -220,7 +219,6 @@ func (rw *RescannerWorker) processBatch(blocks []uint64) error { } rw.logger.Info("Rescanner pass", - "chain", rw.chain.GetName(), "retried", len(blocks), "success", success, "remaining", len(rw.failedBlocks), @@ -257,14 +255,12 @@ func (rw *RescannerWorker) flushUnsafe() { if len(rw.pendingSaves) > 0 { _ = rw.blockStore.SaveFailedBlocks(rw.chain.GetNetworkInternalCode(), rw.pendingSaves) rw.logger.Debug("Batch saved failed blocks", - "chain", rw.chain.GetName(), "count", len(rw.pendingSaves)) rw.pendingSaves = rw.pendingSaves[:0] } if len(rw.pendingRemoves) > 0 { _ = rw.blockStore.RemoveFailedBlocks(rw.chain.GetNetworkInternalCode(), rw.pendingRemoves) rw.logger.Debug("Batch removed failed blocks", - "chain", rw.chain.GetName(), "count", len(rw.pendingRemoves)) rw.pendingRemoves = rw.pendingRemoves[:0] } diff --git a/pkg/adaptive/limiter.go b/pkg/adaptive/limiter.go new file mode 100644 index 0000000..0d211bf --- /dev/null +++ b/pkg/adaptive/limiter.go @@ -0,0 +1,141 @@ +// Package adaptive provides an AIMD concurrency limiter: the configured +// concurrency is a ceiling the limiter backs off from under latency/errors and +// recovers toward when calls are fast. +package adaptive + +import ( + "context" + "sync" + "time" +) + +type Config struct { + Min int + Max int + Start int + HighLatency time.Duration + LowLatency time.Duration + AdjustInterval time.Duration + GrowStreak int +} + +func (c *Config) withDefaults() { + if c.Max < 1 { + c.Max = 1 + } + if c.Min < 1 { + c.Min = 1 + } + if c.Min > c.Max { + c.Min = c.Max + } + if c.Start <= 0 || c.Start > c.Max { + c.Start = c.Max + } + if c.Start < c.Min { + c.Start = c.Min + } + if c.HighLatency <= 0 { + c.HighLatency = 2500 * time.Millisecond + } + if c.LowLatency <= 0 || c.LowLatency >= c.HighLatency { + c.LowLatency = c.HighLatency / 2 + } + if c.AdjustInterval <= 0 { + c.AdjustInterval = time.Second + } + if c.GrowStreak <= 0 { + c.GrowStreak = 10 + } +} + +type Limiter struct { + cfg Config + mu sync.Mutex + cond *sync.Cond + + limit int + inflight int + goodStreak int + lastAdjust time.Time +} + +func New(cfg Config) *Limiter { + cfg.withDefaults() + l := &Limiter{cfg: cfg, limit: cfg.Start} + l.cond = sync.NewCond(&l.mu) + return l +} + +func (l *Limiter) Acquire(ctx context.Context) error { + l.mu.Lock() + defer l.mu.Unlock() + + for { + if err := ctx.Err(); err != nil { + return err + } + if l.inflight < l.limit { + l.inflight++ + return nil + } + // sync.Cond.Wait ignores ctx; broadcast on cancel to re-evaluate. + stop := context.AfterFunc(ctx, func() { + l.mu.Lock() + l.cond.Broadcast() + l.mu.Unlock() + }) + l.cond.Wait() + stop() + } +} + +func (l *Limiter) Release() { + l.mu.Lock() + if l.inflight > 0 { + l.inflight-- + } + l.mu.Unlock() + l.cond.Signal() +} + +func (l *Limiter) Observe(latency time.Duration, ok bool) { + l.mu.Lock() + defer l.mu.Unlock() + + if !ok || latency >= l.cfg.HighLatency { + l.goodStreak = 0 + if l.limit > l.cfg.Min && time.Since(l.lastAdjust) >= l.cfg.AdjustInterval { + l.limit = maxInt(l.cfg.Min, l.limit/2) + l.lastAdjust = time.Now() + } + return + } + + if latency > l.cfg.LowLatency { + return + } + + l.goodStreak++ + if l.goodStreak >= l.cfg.GrowStreak && + l.limit < l.cfg.Max && + time.Since(l.lastAdjust) >= l.cfg.AdjustInterval { + l.limit++ + l.goodStreak = 0 + l.lastAdjust = time.Now() + l.cond.Signal() + } +} + +func (l *Limiter) Limit() int { + l.mu.Lock() + defer l.mu.Unlock() + return l.limit +} + +func maxInt(a, b int) int { + if a > b { + return a + } + return b +} diff --git a/pkg/adaptive/limiter_test.go b/pkg/adaptive/limiter_test.go new file mode 100644 index 0000000..1e64576 --- /dev/null +++ b/pkg/adaptive/limiter_test.go @@ -0,0 +1,161 @@ +package adaptive + +import ( + "context" + "sync" + "testing" + "time" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +func TestConfigDefaults(t *testing.T) { + l := New(Config{Max: 16}) + assert.Equal(t, 16, l.Limit(), "start defaults to max") + assert.Equal(t, 1, l.cfg.Min) + assert.Equal(t, 2500*time.Millisecond, l.cfg.HighLatency) + assert.Equal(t, 1250*time.Millisecond, l.cfg.LowLatency) +} + +func TestObserveMultiplicativeDecreaseOnError(t *testing.T) { + // A short AdjustInterval lets each spaced-out failure halve the limit. + l := New(Config{Max: 16, Min: 1, AdjustInterval: time.Millisecond}) + for _, want := range []int{8, 4, 2, 1, 1} { + l.Observe(time.Second, false) + assert.Equal(t, want, l.Limit()) + time.Sleep(2 * time.Millisecond) // pass the adjust interval + } +} + +func TestObserveDecreaseOnHighLatency(t *testing.T) { + l := New(Config{Max: 10, Min: 2, HighLatency: 2 * time.Second, AdjustInterval: 0}) + l.Observe(3*time.Second, true) // success but slow -> congestion + assert.Equal(t, 5, l.Limit()) +} + +func TestObserveAdditiveIncreaseOnFastSuccess(t *testing.T) { + l := New(Config{ + Max: 16, Min: 1, Start: 4, + LowLatency: time.Second, HighLatency: 2 * time.Second, + AdjustInterval: 0, GrowStreak: 3, + }) + // Fewer than GrowStreak good samples: no growth yet. + l.Observe(100*time.Millisecond, true) + l.Observe(100*time.Millisecond, true) + assert.Equal(t, 4, l.Limit()) + // Third good sample crosses the streak threshold -> +1. + l.Observe(100*time.Millisecond, true) + assert.Equal(t, 5, l.Limit()) +} + +func TestObserveNeverExceedsMax(t *testing.T) { + l := New(Config{Max: 3, Min: 1, Start: 3, LowLatency: time.Second, AdjustInterval: 0, GrowStreak: 1}) + for i := 0; i < 20; i++ { + l.Observe(10*time.Millisecond, true) + } + assert.Equal(t, 3, l.Limit(), "cannot grow past Max") +} + +func TestAdjustIntervalDampsBurst(t *testing.T) { + // A burst of failures within one interval collapses the limit only once. + l := New(Config{Max: 16, Min: 1, AdjustInterval: time.Hour}) + for i := 0; i < 10; i++ { + l.Observe(time.Second, false) + } + assert.Equal(t, 8, l.Limit(), "only one decrease per AdjustInterval") +} + +func TestAcquireRespectsLimitAndReleases(t *testing.T) { + l := New(Config{Max: 2, Min: 1, Start: 2}) + ctx := context.Background() + require.NoError(t, l.Acquire(ctx)) + require.NoError(t, l.Acquire(ctx)) + + // Third acquire must block until a Release happens. + acquired := make(chan struct{}) + go func() { + _ = l.Acquire(ctx) + close(acquired) + }() + + select { + case <-acquired: + t.Fatal("acquire should block when at limit") + case <-time.After(50 * time.Millisecond): + } + + l.Release() + select { + case <-acquired: + case <-time.After(time.Second): + t.Fatal("acquire should proceed after release") + } +} + +func TestAcquireCancelledContext(t *testing.T) { + l := New(Config{Max: 1, Min: 1, Start: 1}) + require.NoError(t, l.Acquire(context.Background())) // fill the only slot + + ctx, cancel := context.WithCancel(context.Background()) + errc := make(chan error, 1) + go func() { errc <- l.Acquire(ctx) }() + + time.Sleep(20 * time.Millisecond) + cancel() // cancellation alone must wake the waiter — no Release needed + + select { + case err := <-errc: + assert.ErrorIs(t, err, context.Canceled) + case <-time.After(time.Second): + t.Fatal("cancelled acquire should return") + } +} + +// TestAcquireWakesAllWaitersOnCancel: blocked waiters all return on ctx cancel. +func TestAcquireWakesAllWaitersOnCancel(t *testing.T) { + l := New(Config{Max: 2, Min: 1, Start: 2}) + require.NoError(t, l.Acquire(context.Background())) + require.NoError(t, l.Acquire(context.Background())) // both slots held, never released + + ctx, cancel := context.WithCancel(context.Background()) + const waiters = 20 + done := make(chan error, waiters) + for i := 0; i < waiters; i++ { + go func() { done <- l.Acquire(ctx) }() + } + + time.Sleep(30 * time.Millisecond) // let them all park in Wait + cancel() + + timeout := time.After(2 * time.Second) + for i := 0; i < waiters; i++ { + select { + case err := <-done: + assert.ErrorIs(t, err, context.Canceled) + case <-timeout: + t.Fatalf("waiter %d did not wake on cancel", i) + } + } +} + +func TestConcurrentAcquireReleaseNoLeak(t *testing.T) { + l := New(Config{Max: 4, Min: 1, Start: 4}) + ctx := context.Background() + var wg sync.WaitGroup + for i := 0; i < 100; i++ { + wg.Add(1) + go func() { + defer wg.Done() + if l.Acquire(ctx) == nil { + time.Sleep(time.Millisecond) + l.Release() + } + }() + } + wg.Wait() + l.mu.Lock() + inflight := l.inflight + l.mu.Unlock() + assert.Equal(t, 0, inflight, "all slots released") +} diff --git a/pkg/common/config/types.go b/pkg/common/config/types.go index fd41a7c..1d6ff0f 100644 --- a/pkg/common/config/types.go +++ b/pkg/common/config/types.go @@ -3,7 +3,6 @@ package config import ( "time" - "github.com/fystack/multichain-indexer/internal/rpc" "github.com/fystack/multichain-indexer/pkg/common/enum" ) @@ -36,7 +35,6 @@ type Defaults struct { Status StatusConfig `yaml:"status"` Client ClientConfig `yaml:"client"` Throttle Throttle `yaml:"throttle"` - Failover rpc.FailoverConfig `yaml:"failover"` } type Chains map[string]ChainConfig diff --git a/pkg/repository/repository.go b/pkg/repository/repository.go index 308aad8..56b869e 100644 --- a/pkg/repository/repository.go +++ b/pkg/repository/repository.go @@ -27,6 +27,9 @@ var ( // https://github.com/jackc/pgerrcode/blob/master/errcode.go UniqueViolation = "23505" ForeignKeyViolation = "23503" + // InvalidTextRepresentation (22P02) e.g. filtering by an enum value the DB + // type does not define — treated as no matching rows. + InvalidTextRepresentation = "22P02" ) type Repository[T any] interface { @@ -82,6 +85,13 @@ func (r *repository[T]) WrapError(ctx context.Context, err error) error { return nil } + // Filtering by an enum value the DB type does not define (22P02): no such + // rows exist, so treat it as an empty result rather than a hard error. + var pgErr *pgconn.PgError + if errors.As(err, &pgErr) && pgErr.Code == InvalidTextRepresentation { + return nil + } + // otherwise, return original error // this is usually an unidentified internal error if err != nil {