From 840fd9e3e25b787e0be308152ee18648356d99d9 Mon Sep 17 00:00:00 2001 From: faisal-link Date: Wed, 8 Jul 2026 17:42:06 +0400 Subject: [PATCH 01/11] update core configs plugin yaml --- scripts/core.config.toml | 12 ++++++++++++ 1 file changed, 12 insertions(+) diff --git a/scripts/core.config.toml b/scripts/core.config.toml index 190b4ae38..4112d2ce4 100644 --- a/scripts/core.config.toml +++ b/scripts/core.config.toml @@ -32,6 +32,8 @@ NetworkNameFull = "sui-localnet" # optional, CLL naming convention Name = 'sui-node-1' URL = 'http://localhost:9000' # For local development SolidityURL = 'http://localhost:9000' +GrpcTarget = 'localhost:9000' # gRPC endpoint (host:port); required for checkpoint-based indexing +GrpcToken = 'test' # gRPC auth token; gRPC is only enabled when both target and token are set [Sui.TransactionManager] BroadcastChanSize = 100 @@ -50,6 +52,16 @@ SyncTimeoutSecs = 3 PollingIntervalSecs = 3 SyncTimeoutSecs = 3 +# ChainPoller drives checkpoint-based indexing over gRPC and fans checkpoint +# contents out to the Events and Transactions indexers. Requires a node with +# GrpcTarget/GrpcToken configured above. +[Sui.ChainPoller] +PollingIntervalSecs = 2 # how often live polling checks for new checkpoints +SyncTimeoutSecs = 60 # timeout for processing a single checkpoint +ChannelBufferSize = 16 # buffer size of the events/transactions channels +BackfillCheckpointCount = 100 # start at (latest - N); also the rewind span on rescan +# StartCheckpointSequence = 12345 # optional; explicit start checkpoint (overrides backfill) + # [Tracing] # # Enabled turns trace collection on or off. On requires an OTEL Tracing Collector. From bf0c4494da3606e01a6171e2316c17bd25c88198 Mon Sep 17 00:00:00 2001 From: faisal-link Date: Thu, 9 Jul 2026 01:13:44 +0400 Subject: [PATCH 02/11] guard empty response when caching CR results --- relayer/chainreader/reader/chainreader.go | 6 ++++++ relayer/client/grpc_client.go | 16 +++++++++++----- 2 files changed, 17 insertions(+), 5 deletions(-) diff --git a/relayer/chainreader/reader/chainreader.go b/relayer/chainreader/reader/chainreader.go index b846314b7..ed42770bd 100644 --- a/relayer/chainreader/reader/chainreader.go +++ b/relayer/chainreader/reader/chainreader.go @@ -1039,6 +1039,12 @@ func (s *suiChainReader) executeFunction(ctx context.Context, parsed *readIdenti if err != nil { return []any{}, err } + // Guard against an empty package ID: simulating against it would resolve to the zero package + // (0x0) and fail with a confusing "Dependent package not found on-chain" error. Fail loudly + // so the read is retried instead of silently targeting 0x0. + if latestPackageId == "" { + return []any{}, fmt.Errorf("resolved empty latest package id for contract %s (module %s)", parsed.contractName, common.GetModuleForContract(parsed.contractName)) + } // this is the upgraded pkgID parsed.address = latestPackageId diff --git a/relayer/client/grpc_client.go b/relayer/client/grpc_client.go index 3644b97d4..bcffa82f7 100644 --- a/relayer/client/grpc_client.go +++ b/relayer/client/grpc_client.go @@ -1531,10 +1531,11 @@ func (c *PTBClient) loadModulePackageIdsInternal(ctx context.Context, packageId } func (c *PTBClient) GetLatestPackageId(ctx context.Context, packageId string, module string) (string, error) { - // attempt reading the value from cache first + // attempt reading the value from cache first. Only non-empty entries are ever cached, so a + // cache hit is always a valid package ID. cacheKey := "latest_pkg_id_fetch:" + packageId + ":" + module if cachedID, found := c.cache.Get(cacheKey); found { - if id, ok := cachedID.(string); ok { + if id, ok := cachedID.(string); ok && id != "" { return id, nil } } @@ -1545,9 +1546,14 @@ func (c *PTBClient) GetLatestPackageId(ctx context.Context, packageId string, mo var err error result, err = c.getLatestPackageIdInternal(ctx, packageId, module) - // Package IDs are stable for the duration of a CCIP deployment; a longer TTL avoids - // repeated GetFunction/GetPackage/ListOwnedObjects storms during config polling. - c.cache.Set(cacheKey, result, DefaultPackageIdCacheTTL) + // Only cache successful, non-empty resolutions. Caching an empty result on error would + // poison the cache: subsequent calls would return ("", nil) for the TTL window, causing + // reads to target the zero package address (0x0). Package IDs are stable for the duration + // of a CCIP deployment, so a longer TTL on success avoids repeated + // GetFunction/GetPackage/ListOwnedObjects storms during config polling. + if err == nil && result != "" { + c.cache.Set(cacheKey, result, DefaultPackageIdCacheTTL) + } return err }) From 805c5ba7af53d6861a6bce83c9dd11ba7e9f7ed4 Mon Sep 17 00:00:00 2001 From: faisal-link Date: Thu, 9 Jul 2026 18:25:21 +0400 Subject: [PATCH 03/11] apply caching to GetCCIPPackageId in gRPC client --- relayer/client/grpc_client.go | 19 ++++++++++++++++++- 1 file changed, 18 insertions(+), 1 deletion(-) diff --git a/relayer/client/grpc_client.go b/relayer/client/grpc_client.go index bcffa82f7..2afa29446 100644 --- a/relayer/client/grpc_client.go +++ b/relayer/client/grpc_client.go @@ -1609,6 +1609,13 @@ func (c *PTBClient) SetCachedValues(keyValues map[string]any) { // GetCCIPPackageId gets the CCIP package ID from the offramp package ID. // IMPORTANT: This function expects to call the original (un-upgraded / first version) offramp package ID. func (c *PTBClient) GetCCIPPackageID(ctx context.Context, offRampPackageID string) (string, error) { + cacheKey := "ccip_package_id_fetch:" + offRampPackageID + if cachedID, found := c.GetCachedValue(cacheKey); found { + if id, ok := cachedID.(string); ok && id != "" { + return id, nil + } + } + response, err := c.ReadFunction( ctx, offRampPackageID, @@ -1622,7 +1629,17 @@ func (c *PTBClient) GetCCIPPackageID(ctx context.Context, offRampPackageID strin return "", err } - return response[0].(string), nil + ccipPackageID := response[0].(string) + if ccipPackageID == "" { + return "", fmt.Errorf("no CCIP package ID found for offramp package %s", offRampPackageID) + } + + // Since the original CCIP package ID is the value we require here, we can cache it + // without a specific TTL. Upgrades that change the latest package ID will be resolved + // using `GetLatestPackageId` which will re-read the value from the chain and applies + // a TTL if it uses a cache. + c.SetCachedValue(cacheKey, ccipPackageID) + return ccipPackageID, nil } // GetValueFromPackageOwnedObjectField gets the value of a field from a package owned object. From 882a494b8ad64525c9a1b79de24163150157af92 Mon Sep 17 00:00:00 2001 From: faisal-link Date: Thu, 9 Jul 2026 19:24:37 +0400 Subject: [PATCH 04/11] enable read cache by default --- relayer/chainreader/reader/reader_cache.go | 108 +++++++++++++++--- .../chainreader/reader/reader_cache_test.go | 77 +++++++++++++ 2 files changed, 166 insertions(+), 19 deletions(-) diff --git a/relayer/chainreader/reader/reader_cache.go b/relayer/chainreader/reader/reader_cache.go index 104812ea4..4044b2417 100644 --- a/relayer/chainreader/reader/reader_cache.go +++ b/relayer/chainreader/reader/reader_cache.go @@ -23,23 +23,35 @@ type CacheConfig struct { // refs are version-stable), so this can be minutes without serving a stale ref. ObjectTTL time.Duration // ReadCacheEnabled caches decoded read-call (devInspect) results keyed by read identifier + params. - // It is OFF by default: config reads change rarely, but caching them trades a little staleness for - // fewer node round-trips, so it is opt-in and should be enabled only with a short ReadTTL. + // Config reads (OffRamp OCR config, source-chain configs, static/dynamic config, price seq number) + // change rarely, so caching them for a short ReadTTL cuts the redundant devInspect fan-out that + // otherwise floods the node every config-poll cycle. Combined with StaleReadTTL below, it also makes + // reads resilient to transient RPC cancellation (the EVM->Sui commit blocker). ReadCacheEnabled bool - // ReadTTL bounds how long a cached read result may be served before it is re-fetched. + // ReadTTL bounds how long a cached read result is served as fresh before it is re-fetched. ReadTTL time.Duration - // CleanupInterval is how often expired entries are purged from both underlying caches. + // StaleReadTTL bounds how long the last successfully-read value may be served as a fallback when a + // fresh read fails (see GetReadResults). This turns a transient read failure — e.g. a context + // cancellation while a slow config poll is in flight — into "serve the last known-good config" + // instead of "return a zero/empty config", which is what makes the Sui commit plugin reject every + // EVM->Sui report. Only consulted when ReadCacheEnabled is true; a zero value falls back to the + // default via withDefaults. + StaleReadTTL time.Duration + // CleanupInterval is how often expired entries are purged from the underlying caches. CleanupInterval time.Duration } // DefaultCacheConfig returns safe defaults: object caching ON (high value, no staleness risk for -// version-stable objects) and read-call caching OFF (opt-in, since it can briefly mask config changes). +// version-stable objects) and read-call caching ON with a short freshness TTL plus a bounded +// serve-stale fallback (config reads change rarely, and serving a slightly-stale config is far safer +// than serving a zero config on a transient read failure). func DefaultCacheConfig() CacheConfig { return CacheConfig{ ObjectCacheEnabled: true, ObjectTTL: 5 * time.Minute, - ReadCacheEnabled: false, - ReadTTL: 2 * time.Second, + ReadCacheEnabled: true, + ReadTTL: 10 * time.Second, + StaleReadTTL: 5 * time.Minute, CleanupInterval: 1 * time.Minute, } } @@ -52,6 +64,9 @@ func (c CacheConfig) withDefaults() CacheConfig { if c.ReadTTL <= 0 { c.ReadTTL = d.ReadTTL } + if c.StaleReadTTL <= 0 { + c.StaleReadTTL = d.StaleReadTTL + } if c.CleanupInterval <= 0 { c.CleanupInterval = d.CleanupInterval } @@ -79,16 +94,20 @@ type Cache struct { readCache *cache.Cache readGroup singleflight.Group + // staleReadCache retains the last successfully-read value per key for StaleReadTTL, so it can be + // served as a fallback when a fresh read fails. It intentionally outlives readCache entries. + staleReadCache *cache.Cache } // NewCache builds a Cache from cfg (missing durations are defaulted). func NewCache(lggr logger.Logger, cfg CacheConfig) *Cache { cfg = cfg.withDefaults() return &Cache{ - lggr: logger.Named(lggr, "ReaderCache"), - cfg: cfg, - objectCache: cache.New(cfg.ObjectTTL, cfg.CleanupInterval), - readCache: cache.New(cfg.ReadTTL, cfg.CleanupInterval), + lggr: logger.Named(lggr, "ReaderCache"), + cfg: cfg, + objectCache: cache.New(cfg.ObjectTTL, cfg.CleanupInterval), + readCache: cache.New(cfg.ReadTTL, cfg.CleanupInterval), + staleReadCache: cache.New(cfg.StaleReadTTL, cfg.CleanupInterval), } } @@ -133,9 +152,15 @@ func (rc *Cache) GetObjectMetadata( } // GetReadResults returns the decoded results of a read call, serving them from cache when enabled and -// otherwise invoking loader exactly once across concurrent callers for the same key. The cached value is -// the raw []any decoded result; callers decode it into their own return value on each call, so no shared -// mutable state escapes. Disabled by default — see CacheConfig.ReadCacheEnabled. +// otherwise invoking loader exactly once across concurrent callers for the same key. Disabled by +// default — see CacheConfig.ReadCacheEnabled. +// +// Every cache hit returns a deep copy of the cached value. This is required for correctness: callers +// (GetLatestValue -> prepareFunctionReadResult -> applyResultFieldRenames/MaybeRenameFields, and the +// NormalizeReturnValuesToHex path) mutate the decoded result in place — e.g. renaming the OffRamp OCR +// config's `big_f` field to `f`. Handing out the cached reference would let the first read mutate the +// shared value, so a later cache hit would see an already-transformed object (e.g. `big_f` missing) and +// fail. Copying on the way out keeps the cached value pristine and each caller's result independent. func (rc *Cache) GetReadResults( ctx context.Context, key string, @@ -146,18 +171,27 @@ func (rc *Cache) GetReadResults( } if res, ok := getTyped[[]any](rc.readCache, key); ok { - return res, nil + return deepCopyAnySlice(res), nil } v, err, _ := rc.readGroup.Do(key, func() (any, error) { if res, ok := getTyped[[]any](rc.readCache, key); ok { return res, nil } - res, err := loader(ctx) - if err != nil { - return nil, err + res, lErr := loader(ctx) + if lErr != nil { + // Serve-stale: a fresh read failed (commonly a context cancellation while a slow config + // poll is in flight). Rather than surfacing a zero/empty result — which upstream turns into + // a zero OCR config that blocks EVM->Sui commits — fall back to the last successfully-read + // value if one is still within StaleReadTTL. + if stale, ok := getTyped[[]any](rc.staleReadCache, key); ok { + rc.lggr.Warnw("read failed; serving last-known-good cached result", "key", key, "err", lErr) + return stale, nil + } + return nil, lErr } rc.readCache.Set(key, res, cache.DefaultExpiration) + rc.staleReadCache.Set(key, res, cache.DefaultExpiration) return res, nil }) if err != nil { @@ -165,7 +199,43 @@ func (rc *Cache) GetReadResults( } res, _ := v.([]any) - return res, nil + // Copy on the way out so concurrent singleflight sharers and later cache hits never mutate the + // cached value (see doc comment). + return deepCopyAnySlice(res), nil +} + +// deepCopyAnySlice returns a structural deep copy of a decoded read result. The values originate from +// protobuf Struct.AsInterface, so they are only ever composed of map[string]any, []any and immutable +// scalars (string, float64, bool, nil); those containers are the ones downstream transforms mutate. +func deepCopyAnySlice(in []any) []any { + if in == nil { + return nil + } + out := make([]any, len(in)) + for i, v := range in { + out[i] = deepCopyAny(v) + } + return out +} + +func deepCopyAny(v any) any { + switch t := v.(type) { + case map[string]any: + m := make(map[string]any, len(t)) + for k, val := range t { + m[k] = deepCopyAny(val) + } + return m + case []any: + return deepCopyAnySlice(t) + case []byte: + cp := make([]byte, len(t)) + copy(cp, t) + return cp + default: + // Scalars (string, float64, bool, nil, ...) are immutable; safe to share. + return v + } } // getTyped fetches key from c and asserts it to T, returning ok=false on miss or type mismatch. diff --git a/relayer/chainreader/reader/reader_cache_test.go b/relayer/chainreader/reader/reader_cache_test.go index e91fc7b39..f114965ff 100644 --- a/relayer/chainreader/reader/reader_cache_test.go +++ b/relayer/chainreader/reader/reader_cache_test.go @@ -135,3 +135,80 @@ func TestReaderCache_ReadResults(t *testing.T) { } require.Equal(t, int32(3), atomic.LoadInt32(&dCalls), "disabled read cache should pass through") } + +// A cache hit must return an independent deep copy: callers mutate the decoded result in place (e.g. +// GetLatestValue renames the OffRamp OCR config's `big_f` field), so handing out the shared cached +// reference would corrupt the cached value and break every subsequent read of the same key. This is the +// regression that made the EVM->Sui config read return a broken/zero OCR config once caching was on. +func TestReaderCache_ReadResults_HitReturnsIndependentCopy(t *testing.T) { + t.Parallel() + + rc := NewCache(logger.Test(t), CacheConfig{ReadCacheEnabled: true, ReadTTL: time.Minute}) + var calls int32 + loader := func(context.Context) ([]any, error) { + atomic.AddInt32(&calls, 1) + // Mirrors the shape of a Sui OCR config read (nested map with the big_f field). + return []any{map[string]any{ + "config_info": map[string]any{"big_f": float64(1), "n": float64(2)}, + "signers": []any{"a", "b"}, + }}, nil + } + + // First read populates the cache, then mutates its own copy the way GetLatestValue does (rename + // big_f -> f in place). + first, err := rc.GetReadResults(context.Background(), "k", loader) + require.NoError(t, err) + ci := first[0].(map[string]any)["config_info"].(map[string]any) + ci["f"] = ci["big_f"] + delete(ci, "big_f") + first[0].(map[string]any)["signers"].([]any)[0] = "MUTATED" + + // Second read is a cache hit (loader not called again) and must see the pristine original value, + // unaffected by the first caller's in-place mutation. + second, err := rc.GetReadResults(context.Background(), "k", loader) + require.NoError(t, err) + require.Equal(t, int32(1), atomic.LoadInt32(&calls), "second read must be served from cache") + + ci2 := second[0].(map[string]any)["config_info"].(map[string]any) + require.Equal(t, float64(1), ci2["big_f"], "cached big_f must survive a prior caller's rename") + _, renamed := ci2["f"] + require.False(t, renamed, "prior caller's rename must not leak into the cached value") + require.Equal(t, "a", second[0].(map[string]any)["signers"].([]any)[0], + "prior caller's slice mutation must not leak into the cached value") +} + +// A transient read failure after the fresh TTL has expired must be served from the last known-good +// value (serve-stale), not surfaced as an error/zero result. This is the direct fix for the EVM->Sui +// commit blocker where a cancelled config read produced a zero OCR config. +func TestReaderCache_ReadResults_ServesStaleOnTransientFailure(t *testing.T) { + t.Parallel() + + rc := NewCache(logger.Test(t), CacheConfig{ + ReadCacheEnabled: true, + ReadTTL: 20 * time.Millisecond, // short fresh window so we can force a re-fetch + StaleReadTTL: time.Minute, + }) + + good := []any{map[string]any{"config_digest": "digest-1"}} + _, err := rc.GetReadResults(context.Background(), "k", func(context.Context) ([]any, error) { + return good, nil + }) + require.NoError(t, err) + + // Let the fresh entry expire so the next read must call the loader. + time.Sleep(40 * time.Millisecond) + + // Loader now fails (simulating a context cancellation during a slow poll). Serve-stale must return + // the last good value instead of the error. + res, err := rc.GetReadResults(context.Background(), "k", func(context.Context) ([]any, error) { + return nil, context.Canceled + }) + require.NoError(t, err, "transient failure after TTL expiry must be served from the stale cache") + require.Equal(t, "digest-1", res[0].(map[string]any)["config_digest"]) + + // With no prior success there is nothing to serve stale, so the error propagates. + _, err = rc.GetReadResults(context.Background(), "never-loaded", func(context.Context) ([]any, error) { + return nil, context.Canceled + }) + require.ErrorIs(t, err, context.Canceled, "a cold key with no stale entry must surface the read error") +} From cf9f44d18326b0ce65fcc0a1aae0d4f76e1d1a3b Mon Sep 17 00:00:00 2001 From: faisal-link Date: Thu, 9 Jul 2026 19:39:13 +0400 Subject: [PATCH 05/11] disable simulateTransaction checks for ReadFunction --- relayer/client/grpc_client.go | 11 ++++++++--- 1 file changed, 8 insertions(+), 3 deletions(-) diff --git a/relayer/client/grpc_client.go b/relayer/client/grpc_client.go index 2afa29446..b69f46c66 100644 --- a/relayer/client/grpc_client.go +++ b/relayer/client/grpc_client.go @@ -690,7 +690,7 @@ func (c *PTBClient) readFunctionInternal(ctx context.Context, packageId string, buildDur = time.Since(buildStart) simStart := time.Now() - results, err = c.simulatePTBInternal(readCtx, txExecService, bcsBytes) + results, err = c.simulatePTBInternal(readCtx, txExecService, bcsBytes, false) simDur = time.Since(simStart) totalDur := time.Since(rfStart) if totalDur > 3*time.Second { @@ -715,14 +715,18 @@ func (c *PTBClient) SimulatePTB(ctx context.Context, bcsBytes []byte) ([]any, er } var simErr error - results, simErr = c.simulatePTBInternal(ctx, txExecService, bcsBytes) + results, simErr = c.simulatePTBInternal(ctx, txExecService, bcsBytes, true) return simErr }) return results, err } -func (c *PTBClient) simulatePTBInternal(ctx context.Context, txExecService suirpcv2.TransactionExecutionServiceClient, bcsBytes []byte) ([]any, error) { +func (c *PTBClient) simulatePTBInternal(ctx context.Context, txExecService suirpcv2.TransactionExecutionServiceClient, bcsBytes []byte, checks bool) ([]any, error) { doGasSelection := false + checksEnum := suirpcv2.SimulateTransactionRequest_DISABLED.Enum() + if checks { + checksEnum = suirpcv2.SimulateTransactionRequest_ENABLED.Enum() + } // measure the raw SimulateTransaction RPC latency and the number of simulate calls // concurrently hitting the single gRPC connection / local Sui node. @@ -730,6 +734,7 @@ func (c *PTBClient) simulatePTBInternal(ctx context.Context, txExecService suirp response, err := txExecService.SimulateTransaction(ctx, &suirpcv2.SimulateTransactionRequest{ Transaction: &suirpcv2.Transaction{Bcs: &suirpcv2.Bcs{Value: bcsBytes}}, DoGasSelection: &doGasSelection, + Checks: checksEnum, }) simDur := time.Since(simStart) if simDur > time.Second { From 311e022cdec011d48b03659720ab229cb72343c7 Mon Sep 17 00:00:00 2001 From: faisal-link Date: Thu, 9 Jul 2026 21:58:39 +0400 Subject: [PATCH 06/11] update ExtendedPTBClient and ChainPoller --- relayer/chainreader/indexer/chain_poller.go | 112 ++++++++++++++--- .../chainreader/indexer/chain_poller_test.go | 114 ++++++++++++++++++ relayer/chainreader/reader/reader_cache.go | 2 +- relayer/client/grpc_client.go | 24 ++++ 4 files changed, 233 insertions(+), 19 deletions(-) create mode 100644 relayer/chainreader/indexer/chain_poller_test.go diff --git a/relayer/chainreader/indexer/chain_poller.go b/relayer/chainreader/indexer/chain_poller.go index 0fec02bff..bee57e7ea 100644 --- a/relayer/chainreader/indexer/chain_poller.go +++ b/relayer/chainreader/indexer/chain_poller.go @@ -40,6 +40,7 @@ type SelectorProvider func() []*client.EventSelector // to the indexers. type ChainPoller struct { client client.SuiPTBClient + extendedClient client.ExtendedPTBClient logger logger.Logger config config.ChainPollerConfig eventsCh chan CheckpointEventsBatch @@ -67,6 +68,7 @@ func NewChainPoller( return &ChainPoller{ client: client, + extendedClient: asExtendedPTBClient(client), logger: logger.Named(log, "ChainPoller"), config: cfg, eventsCh: make(chan CheckpointEventsBatch, bufferSize), @@ -178,31 +180,65 @@ func (cp *ChainPoller) run(ctx context.Context) { // computeStartSequence calculates the starting checkpoint sequence number. func (cp *ChainPoller) computeStartSequence(ctx context.Context) (uint64, error) { + var startSeq uint64 + // If StartCheckpointSequence is configured, use it directly if cp.config.StartCheckpointSequence != nil { - return *cp.config.StartCheckpointSequence, nil + startSeq = *cp.config.StartCheckpointSequence + } else { + // Get the latest checkpoint + latestSeq, err := cp.getLatestCheckpointSequence(ctx) + if err != nil { + return 0, fmt.Errorf("failed to get latest checkpoint: %w", err) + } + + cp.logger.Infow("Latest checkpoint fetched in chain poller", "sequence", latestSeq) + + // If BackfillCheckpointCount is configured, start from latest - N + if cp.config.BackfillCheckpointCount != nil { + count := *cp.config.BackfillCheckpointCount + if latestSeq > count { + startSeq = latestSeq - count + } + } else { + // Default: start from latest (no backfill) + startSeq = latestSeq + } } - // Get the latest checkpoint - latestSeq, err := cp.getLatestCheckpointSequence(ctx) - if err != nil { - return 0, fmt.Errorf("failed to get latest checkpoint: %w", err) + return cp.clampToProviderFloor(ctx, startSeq), nil +} + +// clampToProviderFloor raises startSeq to the provider's lowest available checkpoint when history +// has been pruned below the configured backfill/start point. +func (cp *ChainPoller) clampToProviderFloor(ctx context.Context, startSeq uint64) uint64 { + if cp.extendedClient == nil { + cp.logger.Warnw("Failed to get provider checkpoint availability, using requested start sequence", + "startSequence", startSeq, + "error", errors.New("client does not implement ExtendedPTBClient"), + ) + return startSeq } - cp.logger.Infow("Latest checkpoint fetched in chain poller", "sequence", latestSeq) + info, err := cp.extendedClient.GetCheckpointAvailability(ctx) + if err != nil { + cp.logger.Warnw("Failed to get provider checkpoint availability, using requested start sequence", + "startSequence", startSeq, + "error", err, + ) + return startSeq + } - // If BackfillCheckpointCount is configured, start from latest - N - if cp.config.BackfillCheckpointCount != nil { - count := *cp.config.BackfillCheckpointCount - if latestSeq > count { - return latestSeq - count, nil - } - // If latest < count, start from 0 - return 0, nil + lowest := info.GetLowestAvailableCheckpoint() + if lowest > 0 && startSeq < lowest { + cp.logger.Warnw("Start sequence below provider history floor, clamping", + "requested", startSeq, + "lowestAvailable", lowest, + ) + return lowest } - // Default: start from latest (no backfill) - return latestSeq, nil + return startSeq } // getLatestCheckpointSequence fetches the latest checkpoint and returns its sequence number. @@ -249,6 +285,11 @@ func (cp *ChainPoller) RescanRecent() { // catchUp processes checkpoints from startSeq to endSeq (inclusive). func (cp *ChainPoller) catchUp(ctx context.Context, startSeq, endSeq uint64) { + startSeq = cp.clampToProviderFloor(ctx, startSeq) + if startSeq > endSeq { + return + } + for seq := startSeq; seq <= endSeq; seq++ { select { case <-ctx.Done(): @@ -264,11 +305,22 @@ func (cp *ChainPoller) catchUp(ctx context.Context, startSeq, endSeq uint64) { ) return } - cp.logger.Warnw("Checkpoint not found during catch-up, skipping", + + if lowest := cp.providerLowestAvailable(ctx); lowest > 0 && seq < lowest { + cp.logger.Errorw("Checkpoint below provider history floor during catch-up, jumping to lowest available", + "sequence", seq, + "lowestAvailable", lowest, + "error", err, + ) + seq = lowest - 1 + continue + } + + cp.logger.Warnw("Checkpoint not found during catch-up, will retry on next poll", "sequence", seq, "error", err, ) - continue + return } cp.logger.Errorw("Failed to process checkpoint, will retry on next poll", "sequence", seq, @@ -281,6 +333,30 @@ func (cp *ChainPoller) catchUp(ctx context.Context, startSeq, endSeq uint64) { } } +func (cp *ChainPoller) providerLowestAvailable(ctx context.Context) uint64 { + if cp.extendedClient == nil { + cp.logger.Warnw("Failed to get provider checkpoint availability", + "error", errors.New("client does not implement ExtendedPTBClient"), + ) + return 0 + } + + info, err := cp.extendedClient.GetCheckpointAvailability(ctx) + if err != nil { + cp.logger.Warnw("Failed to get provider checkpoint availability", "error", err) + return 0 + } + return info.GetLowestAvailableCheckpoint() +} + +func asExtendedPTBClient(suiClient client.SuiPTBClient) client.ExtendedPTBClient { + ext, ok := suiClient.(client.ExtendedPTBClient) + if !ok { + return nil + } + return ext +} + func isCheckpointNotFound(err error) bool { for err != nil { if st, ok := status.FromError(err); ok && st.Code() == codes.NotFound { diff --git a/relayer/chainreader/indexer/chain_poller_test.go b/relayer/chainreader/indexer/chain_poller_test.go new file mode 100644 index 000000000..e78f0796a --- /dev/null +++ b/relayer/chainreader/indexer/chain_poller_test.go @@ -0,0 +1,114 @@ +package indexer + +import ( + "context" + "errors" + "testing" + "time" + + suirpcv2 "github.com/block-vision/sui-go-sdk/pb/sui/rpc/v2" + "github.com/stretchr/testify/require" + "google.golang.org/grpc/codes" + "google.golang.org/grpc/status" + + "github.com/smartcontractkit/chainlink-common/pkg/logger" + "github.com/smartcontractkit/chainlink-sui/relayer/chainreader/config" + "github.com/smartcontractkit/chainlink-sui/relayer/client" + "github.com/smartcontractkit/chainlink-sui/relayer/testutils" +) + +type checkpointTestClient struct { + testutils.FakeSuiPTBClient + lowestAvailable uint64 + notFoundBelow uint64 + processed []uint64 +} + +func (c *checkpointTestClient) GetCheckpointAvailability(ctx context.Context) (*suirpcv2.GetServiceInfoResponse, error) { + lowest := c.lowestAvailable + return &suirpcv2.GetServiceInfoResponse{ + LowestAvailableCheckpoint: &lowest, + }, nil +} + +func (c *checkpointTestClient) GetTransaction(ctx context.Context, digest string) (client.TransactionDetails, error) { + return client.TransactionDetails{}, nil +} + +func (c *checkpointTestClient) GetCheckpointData(ctx context.Context, checkpointSequenceNumber uint64) (*client.CheckpointData, error) { + if checkpointSequenceNumber < c.notFoundBelow { + return nil, status.Error(codes.NotFound, "checkpoint not found") + } + + c.processed = append(c.processed, checkpointSequenceNumber) + seq := checkpointSequenceNumber + return &client.CheckpointData{ + Checkpoint: &suirpcv2.Checkpoint{ + SequenceNumber: &seq, + }, + }, nil +} + +func TestClampToProviderFloor(t *testing.T) { + t.Parallel() + + mockClient := &checkpointTestClient{ + lowestAvailable: 500, + } + cp := &ChainPoller{ + client: mockClient, + extendedClient: asExtendedPTBClient(mockClient), + logger: logger.Test(t), + } + + require.Equal(t, uint64(500), cp.clampToProviderFloor(context.Background(), 100)) + require.Equal(t, uint64(600), cp.clampToProviderFloor(context.Background(), 600)) +} + +func TestCatchUpJumpsOverPrunedCheckpoints(t *testing.T) { + t.Parallel() + + mockClient := &checkpointTestClient{ + lowestAvailable: 500, + notFoundBelow: 500, + } + + cp := NewChainPoller(mockClient, logger.Test(t), config.ChainPollerConfig{ + SyncTimeout: time.Minute, + }, func() []*client.EventSelector { + return nil + }) + + cp.catchUp(context.Background(), 400, 502) + + require.Equal(t, []uint64{500, 501, 502}, mockClient.processed) + require.Equal(t, uint64(502), cp.lastProcessed) +} + +func TestCatchUpRetriesInRangeNotFound(t *testing.T) { + t.Parallel() + + mockClient := &checkpointTestClient{ + lowestAvailable: 100, + notFoundBelow: 501, + } + + cp := NewChainPoller(mockClient, logger.Test(t), config.ChainPollerConfig{ + SyncTimeout: time.Minute, + }, func() []*client.EventSelector { + return nil + }) + + cp.catchUp(context.Background(), 500, 502) + + require.Empty(t, mockClient.processed) + require.Equal(t, uint64(0), cp.lastProcessed) +} + +func TestIsCheckpointNotFound(t *testing.T) { + t.Parallel() + + require.True(t, isCheckpointNotFound(status.Error(codes.NotFound, "missing"))) + require.True(t, isCheckpointNotFound(errors.New("wrapped: not found"))) + require.False(t, isCheckpointNotFound(errors.New("timeout"))) +} diff --git a/relayer/chainreader/reader/reader_cache.go b/relayer/chainreader/reader/reader_cache.go index 4044b2417..e6a3276c0 100644 --- a/relayer/chainreader/reader/reader_cache.go +++ b/relayer/chainreader/reader/reader_cache.go @@ -49,7 +49,7 @@ func DefaultCacheConfig() CacheConfig { return CacheConfig{ ObjectCacheEnabled: true, ObjectTTL: 5 * time.Minute, - ReadCacheEnabled: true, + ReadCacheEnabled: false, ReadTTL: 10 * time.Second, StaleReadTTL: 5 * time.Minute, CleanupInterval: 1 * time.Minute, diff --git a/relayer/client/grpc_client.go b/relayer/client/grpc_client.go index b69f46c66..0bb7ec201 100644 --- a/relayer/client/grpc_client.go +++ b/relayer/client/grpc_client.go @@ -74,6 +74,7 @@ var RateLimitWeights = map[string]int64{ "GetTransactionsByCheckpoint": 0, "GetLatestCheckpoint": 0, "GetCheckpointData": 0, + "GetCheckpointAvailability": 0, "SimulatePTB": 0, "GetCoinMetadata": 0, } @@ -174,6 +175,7 @@ var _ SuiPTBClient = (*PTBClient)(nil) type ExtendedPTBClient interface { SuiPTBClient GetTransaction(ctx context.Context, digest string) (TransactionDetails, error) + GetCheckpointAvailability(ctx context.Context) (*suirpcv2.GetServiceInfoResponse, error) } var _ ExtendedPTBClient = (*PTBClient)(nil) @@ -1735,3 +1737,25 @@ func (c *PTBClient) getParentObjectIDInternal(ctx context.Context, packageID str return parentObjectID, nil } + +// GetCheckpointAvailability returns the provider's checkpoint history bounds from GetServiceInfo. +func (c *PTBClient) GetCheckpointAvailability(ctx context.Context) (*suirpcv2.GetServiceInfoResponse, error) { + var result *suirpcv2.GetServiceInfoResponse + + err := c.WithRateLimit(ctx, "GetCheckpointAvailability", func(ctx context.Context) error { + service, err := c.getLedgerService(ctx) + if err != nil { + return fmt.Errorf("failed to get ledger service: %w", err) + } + + resp, err := service.GetServiceInfo(ctx, &suirpcv2.GetServiceInfoRequest{}) + if err != nil { + return fmt.Errorf("GetServiceInfo failed: %w", err) + } + + result = resp + return nil + }) + + return result, err +} From aa70668f55ee63568334ccfbc0ac7a55315f53bb Mon Sep 17 00:00:00 2001 From: faisal-link Date: Fri, 10 Jul 2026 00:09:34 +0400 Subject: [PATCH 07/11] use EventIndex for offset instead of count-based offset --- relayer/chainreader/indexer/events_indexer.go | 12 ++---------- 1 file changed, 2 insertions(+), 10 deletions(-) diff --git a/relayer/chainreader/indexer/events_indexer.go b/relayer/chainreader/indexer/events_indexer.go index 71e54f026..22ac37504 100644 --- a/relayer/chainreader/indexer/events_indexer.go +++ b/relayer/chainreader/indexer/events_indexer.go @@ -205,14 +205,9 @@ func (eIndexer *EventsIndexer) processEventsForHandle( } packageID := parts[0] - totalCount, err := eIndexer.db.GetTotalCount(ctx, packageID, handle) - if err != nil { - return fmt.Errorf("failed to get total count: %w", err) - } - // Build event records records := make([]database.EventRecord, 0, len(items)) - for i, item := range items { + for _, item := range items { // Convert event JSON data using SDK's AsInterface() data := make(map[string]any) if jsonVal := item.Event.GetJson(); jsonVal != nil { @@ -238,9 +233,6 @@ func (eIndexer *EventsIndexer) processEventsForHandle( } } - // Calculate offset (totalCount + i + 1) - offset := totalCount + uint64(i) + 1 - // Convert block hash to bytes blockHashBytes := []byte(meta.Digest) if decoded, err := base58.Decode(meta.Digest); err == nil { @@ -250,7 +242,7 @@ func (eIndexer *EventsIndexer) processEventsForHandle( record := database.EventRecord{ EventAccountAddress: packageID, EventHandle: handle, - EventOffset: offset, + EventOffset: uint64(item.EventIndex), TxDigest: txDigestHex, BlockVersion: 0, BlockHeight: strconv.FormatUint(meta.SequenceNumber, 10), From 1dbce5f0bcf8211a166298d08ebdee564dd7e710 Mon Sep 17 00:00:00 2001 From: faisal-link Date: Fri, 10 Jul 2026 01:35:32 +0400 Subject: [PATCH 08/11] update gRPC keepalive values --- relayer/client/suigrpcconn/connection.go | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/relayer/client/suigrpcconn/connection.go b/relayer/client/suigrpcconn/connection.go index c3c4fa899..04f6de69b 100644 --- a/relayer/client/suigrpcconn/connection.go +++ b/relayer/client/suigrpcconn/connection.go @@ -15,8 +15,8 @@ import ( ) const ( - defaultKeepAlive = 30 * time.Second - defaultKeepaliveTimeout = 5 * time.Second + defaultKeepAlive = 60 * time.Second + defaultKeepaliveTimeout = 20 * time.Second defaultMaxRecvMsgSize = 20 * 1024 * 1024 ) From c0e402bd45edf4a5f37ee4dc000795607b90053e Mon Sep 17 00:00:00 2001 From: faisal-link Date: Fri, 10 Jul 2026 01:54:02 +0400 Subject: [PATCH 09/11] lint --- codec/decoder.go | 1 - codec/encoder.go | 1 - relayer/chainreader/reader/reader_cache_test.go | 2 +- relayer/codec/decoder.go | 1 - 4 files changed, 1 insertion(+), 4 deletions(-) diff --git a/codec/decoder.go b/codec/decoder.go index d3c191677..468ee89ec 100644 --- a/codec/decoder.go +++ b/codec/decoder.go @@ -13,7 +13,6 @@ import ( aptosBCS "github.com/aptos-labs/aptos-go-sdk/bcs" "github.com/block-vision/sui-go-sdk/models" - ) const ( diff --git a/codec/encoder.go b/codec/encoder.go index ac520d469..9e61e227d 100644 --- a/codec/encoder.go +++ b/codec/encoder.go @@ -9,7 +9,6 @@ import ( "reflect" "strconv" "strings" - ) const ( diff --git a/relayer/chainreader/reader/reader_cache_test.go b/relayer/chainreader/reader/reader_cache_test.go index f114965ff..e081e772f 100644 --- a/relayer/chainreader/reader/reader_cache_test.go +++ b/relayer/chainreader/reader/reader_cache_test.go @@ -170,7 +170,7 @@ func TestReaderCache_ReadResults_HitReturnsIndependentCopy(t *testing.T) { require.Equal(t, int32(1), atomic.LoadInt32(&calls), "second read must be served from cache") ci2 := second[0].(map[string]any)["config_info"].(map[string]any) - require.Equal(t, float64(1), ci2["big_f"], "cached big_f must survive a prior caller's rename") + require.InEpsilon(t, float64(1), ci2["big_f"], 0.0000000000000001, "cached big_f must survive a prior caller's rename") _, renamed := ci2["f"] require.False(t, renamed, "prior caller's rename must not leak into the cached value") require.Equal(t, "a", second[0].(map[string]any)["signers"].([]any)[0], diff --git a/relayer/codec/decoder.go b/relayer/codec/decoder.go index 686533602..834f01dfc 100644 --- a/relayer/codec/decoder.go +++ b/relayer/codec/decoder.go @@ -6,7 +6,6 @@ import ( "github.com/smartcontractkit/chainlink-sui/codec" ) - // Deprecated: use codec.DecodeSuiJsonValue // //go:fix inline From 64dbd5f100c081fb29d81b29e77993915543d197 Mon Sep 17 00:00:00 2001 From: faisal-link Date: Fri, 10 Jul 2026 03:43:16 +0400 Subject: [PATCH 10/11] re-enable read cache and add logs --- relayer/chainreader/indexer/chain_poller.go | 2 ++ relayer/chainreader/indexer/events_indexer.go | 2 ++ relayer/chainreader/reader/reader_cache.go | 4 ++-- 3 files changed, 6 insertions(+), 2 deletions(-) diff --git a/relayer/chainreader/indexer/chain_poller.go b/relayer/chainreader/indexer/chain_poller.go index bee57e7ea..169d2d990 100644 --- a/relayer/chainreader/indexer/chain_poller.go +++ b/relayer/chainreader/indexer/chain_poller.go @@ -494,6 +494,8 @@ func (cp *ChainPoller) filterEvents(meta CheckpointMeta, transactions []*suirpcv // Check if event matches any selector for _, sel := range selectors { + cp.logger.Debugw("Checking if event matches selector", "event", event.GetEventType()) + if eventMatchesSelector(event, sel) { item := CheckpointEventItem{ Event: event, diff --git a/relayer/chainreader/indexer/events_indexer.go b/relayer/chainreader/indexer/events_indexer.go index 22ac37504..9cf3bfedf 100644 --- a/relayer/chainreader/indexer/events_indexer.go +++ b/relayer/chainreader/indexer/events_indexer.go @@ -252,6 +252,8 @@ func (eIndexer *EventsIndexer) processEventsForHandle( IsSynthetic: false, } records = append(records, record) + + eIndexer.logger.Debugw("Prepared event record to insert", "event", record) } // Batch insert with fallback diff --git a/relayer/chainreader/reader/reader_cache.go b/relayer/chainreader/reader/reader_cache.go index e6a3276c0..6c41e3ec7 100644 --- a/relayer/chainreader/reader/reader_cache.go +++ b/relayer/chainreader/reader/reader_cache.go @@ -49,8 +49,8 @@ func DefaultCacheConfig() CacheConfig { return CacheConfig{ ObjectCacheEnabled: true, ObjectTTL: 5 * time.Minute, - ReadCacheEnabled: false, - ReadTTL: 10 * time.Second, + ReadCacheEnabled: true, + ReadTTL: 1 * time.Minute, StaleReadTTL: 5 * time.Minute, CleanupInterval: 1 * time.Minute, } From 9f553f39c466b142c84da7317182e32837cef330 Mon Sep 17 00:00:00 2001 From: faisal-link Date: Fri, 10 Jul 2026 03:48:17 +0400 Subject: [PATCH 11/11] add cache to parent object ID fetch --- relayer/client/grpc_client.go | 12 ++++++++++++ 1 file changed, 12 insertions(+) diff --git a/relayer/client/grpc_client.go b/relayer/client/grpc_client.go index 0bb7ec201..41c183f18 100644 --- a/relayer/client/grpc_client.go +++ b/relayer/client/grpc_client.go @@ -1688,10 +1688,22 @@ func (c *PTBClient) getValuesFromPackageOwnedObjectFieldInternal(ctx context.Con // With derived objects, pointers now store a reference to the parent "Object" struct (e.g., OffRampObject, CCIPObject). // e.g. OffRampStatePointer contains "off_ramp_object_id" field pointing to OffRampObject. func (c *PTBClient) GetParentObjectID(ctx context.Context, packageID string, moduleID string, pointerObjectName string) (string, error) { + cacheKey := "parent_object_id_fetch:" + packageID + ":" + moduleID + ":" + pointerObjectName + if cachedID, found := c.GetCachedValue(cacheKey); found { + if id, ok := cachedID.(string); ok && id != "" { + return id, nil + } + } + var result string err := c.WithRateLimit(ctx, "GetParentObjectID", func(ctx context.Context) error { var err error result, err = c.getParentObjectIDInternal(ctx, packageID, moduleID, pointerObjectName) + + if err == nil && result != "" { + c.SetCachedValue(cacheKey, result) + } + return err }) return result, err