Skip to content
78 changes: 72 additions & 6 deletions autoclaim/bridgedetector/l2_to_lx.go
Original file line number Diff line number Diff line change
Expand Up @@ -26,6 +26,13 @@ var (
// ErrCandidatesNotSynced signals that the source bridge service has not yet synced the requested
// LER. The detector treats it as "retry later": it skips the source without advancing its LER cursor.
ErrCandidatesNotSynced = errors.New("autoclaim l2-to-lx bridge detector: source bridge service not synced yet")
// ErrCannotResolveInitialLER signals that a newly discovered source network's initial LER cursor
// could not be derived from the configured (non-zero) StartL1Block: it either predates the first
// L1 info tree leaf or has no LER yet at that block. The detector treats it as "retry later" --
// it skips the source without advancing its LER cursor -- rather than silently falling back to
// full history, which for an established chain can mean scanning years of stale deposits.
ErrCannotResolveInitialLER = errors.New(
"autoclaim l2-to-lx bridge detector: cannot resolve initial LER cursor from configured StartL1Block")
// ErrSourceDisabled signals that a source network is permanently excluded from bridge service
// resolution, so its URL can never resolve, however long the detector waits. It is the permanent
// counterpart of ErrURLNotFound: the detector treats it as "nothing to do" rather than "retry
Expand Down Expand Up @@ -252,8 +259,15 @@ type L2ToLx struct {
// report is made once per source instead of once per poll (the detector polls indefinitely).
// Only ever touched from the single poll goroutine.
disabledSourceLogged map[uint32]struct{}
// retryLogged records when each retry-later condition was last reported, keyed by
// "<condition>/<sourceID>", so a persistent condition logs at most once per retryLogInterval
// rather than on every poll. Only ever touched from the single poll goroutine.
retryLogged map[string]time.Time
}

// retryLogInterval bounds how often a persistent retry-later condition is logged per source.
const retryLogInterval = 10 * time.Minute

// NewL2ToLx creates an L2-to-Lx Auto Claim bridge detector.
func NewL2ToLx(
source VerifiedBatchSource,
Expand Down Expand Up @@ -301,6 +315,7 @@ func NewL2ToLx(
return time.Now().UTC()
},
disabledSourceLogged: make(map[uint32]struct{}),
retryLogged: make(map[string]time.Time),
}
for _, option := range options {
option(detector)
Expand Down Expand Up @@ -509,6 +524,30 @@ func (w *L2ToLx) processSource(

pending, err := w.resolveFromLER(ctx, source, destinationIDs)
if err != nil {
if errors.Is(err, l1infotreesync.ErrBlockNotProcessed) {
// l1infotreesync has not reached StartL1Block yet (expected on a fresh datadir, where the
// auto-resolved block is far ahead of the sync). Transient: retry this source next poll
// instead of failing the whole poll.
if w.shouldLogRetry("not-processed", source.sourceID) {
w.logInfof("autoclaim l2-to-lx bridge detector: skip source %d (l1infotreesync has not "+
"reached StartL1Block yet, will retry every poll): %v", source.sourceID, err)
}
return sourceRetryLater, nil
}
if errors.Is(err, ErrCannotResolveInitialLER) {
// Loud and retried, never silent: unlike a stale bridge-service URL this is a config
// problem an operator must fix (lower StartL1Block, or set it to 0 to accept scanning
// full history), so it is worth a WARN rather than the Info level used for transient
// finder/sync misses below. It only holds back this one source -- other sources in the
// same poll are unaffected -- and is retried every poll until the config is corrected.
if w.shouldLogRetry("initial-ler", source.sourceID) {
w.logWarnf("autoclaim l2-to-lx bridge detector: skip source %d (cannot resolve initial "+
"LER cursor, will retry every poll until fixed; logged at most every %s): %v -- bridges "+
"from this source will not be autoclaimed until this is resolved",
source.sourceID, retryLogInterval, err)
}
return sourceRetryLater, nil
}
return sourceUpToDate, err
}
if len(pending) == 0 {
Expand Down Expand Up @@ -692,8 +731,17 @@ func sortedNetworks(networks []uint32) []uint32 {
}

// initialFromLER derives the exclusive lower-bound LER the first time a source network is seen. When
// StartL1Block is 0, the full history is requested (nil). Otherwise the source's LER at StartL1Block
// is used; a zero LER (the network had no LER yet at that block) also requests the full history.
// StartL1Block is 0, the full history is requested (nil) -- StartL1Block can only be literally 0
// when an operator sets it explicitly, since automatic resolution (autoclaim/bridgedetector.
// ResolveStartBlock) is always clamped to a strictly positive minimum block, so a 0 here is always
// deliberate operator intent, never a resolution artifact.
//
// For any other configured StartL1Block, the source's LER at that block is used. If it cannot be
// derived -- StartL1Block predates the first L1 info tree leaf, or the network had no LER yet at
// that block -- this used to silently fall back to full history too. That is dangerous for an
// established chain (it can mean scanning years of stale deposits the moment DryRun is lifted), so
// it now returns ErrCannotResolveInitialLER instead; the caller treats that as "retry later" for
// just this source; see processSource.
func (w *L2ToLx) initialFromLER(ctx context.Context, sourceID uint32) (*common.Hash, error) {
if w.startL1Block == 0 {
return nil, nil
Expand All @@ -702,9 +750,8 @@ func (w *L2ToLx) initialFromLER(ctx context.Context, sourceID uint32) (*common.H
leaf, err := w.source.GetLatestL1InfoLeafUntilBlock(ctx, w.startL1Block)
if err != nil {
if errors.Is(err, l1infotreesync.ErrNotFound) {
// StartL1Block predates the first L1 info tree leaf, so there is no baseline to derive a
// lower-bound LER from. Same situation as a zero LER at that block: fetch the full history.
return nil, nil
return nil, fmt.Errorf("%w: source %d: StartL1Block %d predates the first L1 info tree leaf",
ErrCannotResolveInitialLER, sourceID, w.startL1Block)
}
return nil, fmt.Errorf("get latest l1 info leaf until block %d for source %d: %w",
w.startL1Block, sourceID, err)
Expand All @@ -716,7 +763,8 @@ func (w *L2ToLx) initialFromLER(ctx context.Context, sourceID uint32) (*common.H
sourceID, leaf.RollupExitRoot, err)
}
if ler == (common.Hash{}) {
return nil, nil
return nil, fmt.Errorf("%w: source %d has no LER yet at StartL1Block %d",
ErrCannotResolveInitialLER, sourceID, w.startL1Block)
}
return &ler, nil
}
Expand Down Expand Up @@ -891,6 +939,24 @@ func (w *L2ToLx) logErrorf(format string, args ...interface{}) {
}
}

func (w *L2ToLx) logWarnf(format string, args ...interface{}) {
if w.log != nil {
w.log.Warnf(format, args...)
}
}

// shouldLogRetry reports whether the given retry-later condition for sourceID is due to be logged,
// and records the time if so.
func (w *L2ToLx) shouldLogRetry(condition string, sourceID uint32) bool {
key := fmt.Sprintf("%s/%d", condition, sourceID)
now := w.now()
if last, ok := w.retryLogged[key]; ok && now.Sub(last) < retryLogInterval {
return false
}
w.retryLogged[key] = now
return true
}

func (w *L2ToLx) logInfof(format string, args ...interface{}) {
if w.log != nil {
w.log.Infof(format, args...)
Expand Down
78 changes: 67 additions & 11 deletions autoclaim/bridgedetector/l2_to_lx_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -141,13 +141,15 @@ func TestL2ToLxInitialCursorFromStartL1Block(t *testing.T) {
require.Equal(t, newLER, fetcher.queries[0].ToLER)
}

func TestL2ToLxInitialCursorZeroLEROmitsFromLER(t *testing.T) {
func TestL2ToLxInitialCursorZeroLERExplicitGenesisOmitsFromLER(t *testing.T) {
// StartL1Block=0 is explicit operator intent to scan full history; it must never reach
// initialFromLER's LER lookup at all.
ctx := context.Background()
newLER := lerHash(9)
source := &fakeVerifiedBatchSource{
lastProcessedBlock: 100,
rowsByRange: map[blockRange][]*l1infotreesync.VerifyBatches{
{from: 40, to: 100}: {makeVerifyRow(1, newLER, 60)},
{from: 0, to: 99}: {makeVerifyRow(1, newLER, 60)},
},
latestLeaf: &l1infotreesync.L1InfoTreeLeaf{BlockNumber: 40, RollupExitRoot: lerHash(555)},
localExitRoots: map[uint32]common.Hash{1: {}}, // network had no LER yet at StartL1Block
Expand All @@ -158,13 +160,41 @@ func TestL2ToLxInitialCursorZeroLEROmitsFromLER(t *testing.T) {
claimer0 := &fakeClaimer{target: autoclaimtypes.ClaimerTarget{ID: fakeClaimer0ID, DestinationNetwork: 0}}
detector := newTestL2ToLxDetector(
t, source, fetcher, newFakeRegistry(claimer0), newMemoryCursorStore(), newFakePerPairLERStore(), newFakeEnqueuer(),
WithL2ToLxStartL1Block(40), WithL2ToLxBlockWindow(100),
WithL2ToLxStartL1Block(0), WithL2ToLxBlockWindow(100),
)

_, err := detector.PollOnce(ctx)
require.NoError(t, err)
require.Len(t, fetcher.queries, 1)
require.Nil(t, fetcher.queries[0].FromLER, "a zero initial LER must omit from_ler (full history)")
require.Nil(t, fetcher.queries[0].FromLER, "explicit StartL1Block=0 must omit from_ler (full history)")
}

func TestL2ToLxInitialCursorZeroLERWithNonZeroStartBlockRetriesInsteadOfFallingBack(t *testing.T) {
// A non-zero, non-explicit-genesis StartL1Block that resolves to a zero LER used to silently
// fall back to full history. It must now retry the source instead of fetching candidates.
ctx := context.Background()
newLER := lerHash(9)
source := &fakeVerifiedBatchSource{
lastProcessedBlock: 100,
rowsByRange: map[blockRange][]*l1infotreesync.VerifyBatches{
{from: 40, to: 100}: {makeVerifyRow(1, newLER, 60)},
},
latestLeaf: &l1infotreesync.L1InfoTreeLeaf{BlockNumber: 40, RollupExitRoot: lerHash(555)},
localExitRoots: map[uint32]common.Hash{1: {}}, // network had no LER yet at StartL1Block
}
fetcher := newFakeFetcher()
fetcher.urls[1] = fakeSrcURL1
claimer0 := &fakeClaimer{target: autoclaimtypes.ClaimerTarget{ID: fakeClaimer0ID, DestinationNetwork: 0}}
detector := newTestL2ToLxDetector(
t, source, fetcher, newFakeRegistry(claimer0), newMemoryCursorStore(), newFakePerPairLERStore(), newFakeEnqueuer(),
WithL2ToLxStartL1Block(40), WithL2ToLxBlockWindow(100),
)

result, err := detector.PollOnce(ctx)
require.NoError(t, err, "an unresolvable initial LER must not fail the whole poll")
require.Empty(t, fetcher.queries, "must not silently fetch full history")
require.Equal(t, 1, result.SkippedSourceCount)
require.Equal(t, 0, result.ProcessedSourceCount)
}

func TestL2ToLxFinderMissSkipsSourceWithoutAdvancingCursor(t *testing.T) {
Expand Down Expand Up @@ -593,15 +623,15 @@ func TestL2ToLxFetchAllCandidatesMergesAcrossBatches(t *testing.T) {
"fetchAllCandidates must merge every batch's candidates, not keep only the last batch's")
}

func TestL2ToLxInitialCursorLeafNotFoundOmitsFromLER(t *testing.T) {
// A StartL1Block that predates the first L1 info tree leaf has no baseline to derive a
// lower-bound LER from; it must behave like the zero-LER case and request the full history.
func TestL2ToLxInitialCursorLeafNotFoundExplicitGenesisOmitsFromLER(t *testing.T) {
// StartL1Block=0 is explicit operator intent to scan full history; it must never reach
// initialFromLER's leaf lookup at all.
ctx := context.Background()
newLER := lerHash(9)
source := &fakeVerifiedBatchSource{
lastProcessedBlock: 100,
rowsByRange: map[blockRange][]*l1infotreesync.VerifyBatches{
{from: 40, to: 100}: {makeVerifyRow(1, newLER, 60)},
{from: 0, to: 99}: {makeVerifyRow(1, newLER, 60)},
},
latestLeafErr: l1infotreesync.ErrNotFound,
}
Expand All @@ -611,14 +641,40 @@ func TestL2ToLxInitialCursorLeafNotFoundOmitsFromLER(t *testing.T) {
claimer0 := &fakeClaimer{target: autoclaimtypes.ClaimerTarget{ID: fakeClaimer0ID, DestinationNetwork: 0}}
detector := newTestL2ToLxDetector(
t, source, fetcher, newFakeRegistry(claimer0), newMemoryCursorStore(), newFakePerPairLERStore(), newFakeEnqueuer(),
WithL2ToLxStartL1Block(40), WithL2ToLxBlockWindow(100),
WithL2ToLxStartL1Block(0), WithL2ToLxBlockWindow(100),
)

_, err := detector.PollOnce(ctx)
require.NoError(t, err)
require.Len(t, fetcher.queries, 1)
require.Nil(t, fetcher.queries[0].FromLER,
"a StartL1Block older than the first L1 info tree leaf must omit from_ler (full history)")
require.Nil(t, fetcher.queries[0].FromLER, "explicit StartL1Block=0 must omit from_ler (full history)")
}

func TestL2ToLxInitialCursorLeafNotFoundWithNonZeroStartBlockRetriesInsteadOfFallingBack(t *testing.T) {
// A non-zero, non-explicit-genesis StartL1Block that predates the first L1 info tree leaf used
// to silently fall back to full history. It must now retry the source instead.
ctx := context.Background()
newLER := lerHash(9)
source := &fakeVerifiedBatchSource{
lastProcessedBlock: 100,
rowsByRange: map[blockRange][]*l1infotreesync.VerifyBatches{
{from: 40, to: 100}: {makeVerifyRow(1, newLER, 60)},
},
latestLeafErr: l1infotreesync.ErrNotFound,
}
fetcher := newFakeFetcher()
fetcher.urls[1] = fakeSrcURL1
claimer0 := &fakeClaimer{target: autoclaimtypes.ClaimerTarget{ID: fakeClaimer0ID, DestinationNetwork: 0}}
detector := newTestL2ToLxDetector(
t, source, fetcher, newFakeRegistry(claimer0), newMemoryCursorStore(), newFakePerPairLERStore(), newFakeEnqueuer(),
WithL2ToLxStartL1Block(40), WithL2ToLxBlockWindow(100),
)

result, err := detector.PollOnce(ctx)
require.NoError(t, err, "an unresolvable initial LER must not fail the whole poll")
require.Empty(t, fetcher.queries, "must not silently fetch full history")
require.Equal(t, 1, result.SkippedSourceCount)
require.Equal(t, 0, result.ProcessedSourceCount)
}

func TestL2ToLxPaginationAndDedup(t *testing.T) {
Expand Down
132 changes: 132 additions & 0 deletions autoclaim/bridgedetector/startblock.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,132 @@
package bridgedetector

import (
"context"
"fmt"
"time"

aggkitcommon "github.com/agglayer/aggkit/common"
aggkittypes "github.com/agglayer/aggkit/types"
)

// StartBlockResolution is the outcome of resolving a bridge detector's configured start-block field
// (L1ToL2BridgeDetector.StartBlock or L2ToLxBridgeDetector.StartL1Block).
type StartBlockResolution struct {
// Block is the concrete start block to hand to the detector.
Block uint64
}

// ResolveStartBlock resolves a bridge detector's start-block config field into a concrete block
// number:
//
// - An explicit value (configured != nil), including 0, is used verbatim: no RPC call, no
// clamping. This is what makes the change backwards compatible -- a deployment that already
// pins a start block keeps behaving exactly as it did before this field became a pointer.
// - An unset value (configured == nil) is resolved by binary-searching L1 block headers for the
// highest block at or before "now - lookback" (about log2(head-minBlock) header lookups), then
// clamped to [minBlock, head] so it can never predate the contract this detector cares about
// (minBlock) nor exceed the current chain head.
//
// It deliberately does not estimate the block number from an average block-time constant: a fixed
// seconds-per-block assumption has been observed, in an orchestration tool in this same codebase, to
// be wrong enough (12.0s assumed vs. 12.512s actual on Sepolia) to skew a multi-hour estimate by
// over ten hours. Reading real timestamps back from the chain avoids that class of error entirely.
func ResolveStartBlock(
ctx context.Context,
l1Client aggkittypes.EthClienter,
configured *uint64,
lookback time.Duration,
minBlock uint64,
now time.Time,
logger aggkitcommon.Logger,
detectorName string,
) (StartBlockResolution, error) {
if configured != nil {
if *configured == 0 && logger != nil {
logger.Warnf(
"autoclaim %s: start block explicitly set to 0 (genesis); this scans the entire "+
"chain history from block 0 and may take a very long time", detectorName)
}
return StartBlockResolution{Block: *configured}, nil
}

block, blockTime, err := resolveBlockAtLookback(ctx, l1Client, lookback, minBlock, now)
if err != nil {
return StartBlockResolution{}, fmt.Errorf(
"resolve %s start block from %s lookback: %w", detectorName, lookback, err)
}

if logger != nil {
logger.Infof(
"autoclaim %s: no start block configured; resolved block %d (unix timestamp %d) from a "+
"%s lookback ending %s UTC",
detectorName, block, blockTime,
lookback, now.UTC().Format(time.RFC3339))
logger.Warnf(
"autoclaim %s: bridges originating before block %d will never be autoclaimed and "+
"require manual claiming", detectorName, block)
}

return StartBlockResolution{Block: block}, nil
}

// resolveBlockAtLookback binary-searches L1 block headers for the highest block whose timestamp is
// at or before now-lookback, then clamps the result to [minBlock, head]. If the head itself is at or below minBlock (a
// misconfigured floor), the head is returned as-is: there is nothing newer to resolve to.
func resolveBlockAtLookback(
ctx context.Context,
l1Client aggkittypes.EthClienter,
lookback time.Duration,
minBlock uint64,
now time.Time,
) (block, blockTime uint64, err error) {
head, err := l1Client.CustomHeaderByNumber(ctx, &aggkittypes.LatestBlock)
if err != nil {
return 0, 0, fmt.Errorf("get L1 head header: %w", err)
}
if head.Number <= minBlock {
return head.Number, head.Time, nil
}

var targetTime uint64
if nowUnix := now.Unix(); nowUnix > int64(lookback.Seconds()) {
targetTime = uint64(nowUnix) - uint64(lookback.Seconds())
}

minHeader, err := headerAt(ctx, l1Client, minBlock)
if err != nil {
return 0, 0, err
}
if minHeader.Time >= targetTime {
// The lookback window predates (or starts exactly at) minBlock itself: clamp to the floor.
return minHeader.Number, minHeader.Time, nil
}

// Binary search (minBlock, head.Number] for the highest block whose timestamp is <= targetTime.
// Using an upper mid keeps the search converging without risking a uint64 underflow at 0.
lo, hi := minBlock, head.Number
best := minHeader
for lo < hi {
mid := lo + (hi-lo+1)/2 //nolint:mnd
h, err := headerAt(ctx, l1Client, mid)
if err != nil {
return 0, 0, err
}
if h.Time <= targetTime {
lo = mid
best = h
} else {
hi = mid - 1
}
}

return best.Number, best.Time, nil
}

func headerAt(ctx context.Context, l1Client aggkittypes.EthClienter, block uint64) (*aggkittypes.BlockHeader, error) {
h, err := l1Client.CustomHeaderByNumber(ctx, aggkittypes.NewBlockNumber(block))
if err != nil {
return nil, fmt.Errorf("get header at block %d: %w", block, err)
}
return h, nil
}
Loading
Loading