From fcea4dc4fb400c19a4d32f418be1f0a1fdf9349d Mon Sep 17 00:00:00 2001 From: Jason-Gregoire Date: Mon, 28 Sep 2026 05:31:09 -0400 Subject: [PATCH] Queue no block that is still in flight, so a discovery that falls behind stops searching the same blocks again while new ones wait, and treat a panicking search as a failed search instead of ending the worker --- pkg/intent/discovery.go | 72 +++++++++++++++-- pkg/intent/discovery_inflight_test.go | 110 ++++++++++++++++++++++++++ 2 files changed, 175 insertions(+), 7 deletions(-) create mode 100644 pkg/intent/discovery_inflight_test.go diff --git a/pkg/intent/discovery.go b/pkg/intent/discovery.go index 637e2849..f44db830 100644 --- a/pkg/intent/discovery.go +++ b/pkg/intent/discovery.go @@ -18,6 +18,7 @@ import ( "github.com/certen/independant-validator/pkg/entitlement" "github.com/certen/independant-validator/pkg/envvar" "log" + "runtime/debug" "sort" "strings" "sync" @@ -197,7 +198,7 @@ type IntentDiscovery struct { // Block monitoring state lastProcessedBlock uint64 - lastQueuedBlock uint64 // highest block sent to workers (prevents re-queuing) + lastQueuedBlock uint64 // highest block height a tick has seen as the head finalizeCeiling uint64 // watermark is not finalized past this (= latest - confirmLag); the last few heights stay re-scannable // chainHead is the latest height the last successful poll observed, and lastAdvanceAt is // when the watermark last moved. Together they answer "is discovery alive and keeping up", @@ -209,7 +210,7 @@ type IntentDiscovery struct { stopCh chan struct{} blockProcessCh chan *BlockProcessJob processedBlocks map[uint64]bool // tracks out-of-order block completions for watermark - watermarkMu sync.Mutex // protects lastProcessedBlock, lastQueuedBlock, processedBlocks + watermarkMu sync.Mutex // protects lastProcessedBlock, lastQueuedBlock, processedBlocks, inFlight mu sync.RWMutex // Intent tracking - E.4 remediation: Two-phase status tracking @@ -235,6 +236,13 @@ type IntentDiscovery struct { // was present at the last look (guarded by watermarkMu). pauseFile string paused bool + + // inFlight are the heights sent to the workers whose job has not finished (guarded by watermarkMu). + // A tick queues no height that is in flight (RB3-F151): it used to queue the whole range from the + // watermark to the head on every tick, so while the workers were behind - after a restart's rewind - + // each tick added the same few hundred blocks again, the workers searched blocks the watermark had + // long passed, and new blocks waited behind them until discovery reported itself stalled. + inFlight map[uint64]bool } // LedgerStoreInterface defines the interface for ledger operations needed by intent discovery @@ -356,6 +364,7 @@ func (id *IntentDiscovery) StartMonitoring() { id.blockProcessCh = make(chan *BlockProcessJob, id.config.MaxConcurrentBlocks) id.retryCh = make(chan *intentRetryJob, 256) id.processedBlocks = make(map[uint64]bool) + id.inFlight = make(map[uint64]bool) id.lastQueuedBlock = id.lastProcessedBlock // reset queue tracker to current watermark // Keep intent status across restarts to avoid reprocessing // E.4 remediation: Two-phase status tracking @@ -622,6 +631,7 @@ func (id *IntentDiscovery) checkForNewBlocks(ctx context.Context) error { id.lastQueuedBlock = latest id.finalizeCeiling = latest id.processedBlocks = make(map[uint64]bool) + id.inFlight = make(map[uint64]bool) id.watermarkMu.Unlock() if id.ledgerStore != nil { if err := id.ledgerStore.SaveIntentLastBlock(latest); err != nil { @@ -654,24 +664,42 @@ func (id *IntentDiscovery) checkForNewBlocks(ctx context.Context) error { hi = from + maxPerTick - 1 } id.lastQueuedBlock = latest + // The heights this tick hands to the workers: every one from the watermark up that is neither in flight + // nor already searched (RB3-F151). A block searched while it was above the finalize ceiling is not + // marked searched (advanceWatermark), so it is searched again - the re-scan of the unconfirmed tip. + if id.inFlight == nil { + id.inFlight = make(map[uint64]bool) + } + var queue []uint64 + if from <= latest { + for h := from; h <= hi; h++ { + if id.inFlight[h] || id.processedBlocks[h] { + continue + } + id.inFlight[h] = true + queue = append(queue, h) + } + } id.watermarkMu.Unlock() - if from > latest { + if len(queue) == 0 { return nil } // Only log genuine forward progress (more than just the re-scan window), to avoid // per-tick noise while idle. - if hi-from+1 > confirmLag+1 { - id.logger.Printf("🔎 Scanning blocks [%d -> %d] (latest %d, finalize<=%d)", from, hi, latest, ceiling) + if uint64(len(queue)) > confirmLag+1 { + id.logger.Printf("🔎 Scanning %d blocks [%d -> %d] (latest %d, finalize<=%d)", len(queue), queue[0], queue[len(queue)-1], latest, ceiling) } - for h := from; h <= hi; h++ { + for i, h := range queue { select { case id.blockProcessCh <- &BlockProcessJob{ PartitionURL: "acc://dn.acme", BlockHeight: h, }: case <-id.stopCh: + // Nothing will search the rest; they are no longer in flight. + id.releaseInFlight(queue[i:]) return nil } } @@ -679,6 +707,15 @@ func (id *IntentDiscovery) checkForNewBlocks(ctx context.Context) error { return nil } +// releaseInFlight marks heights as no longer in flight, so a later tick queues them again. +func (id *IntentDiscovery) releaseInFlight(heights []uint64) { + id.watermarkMu.Lock() + defer id.watermarkMu.Unlock() + for _, h := range heights { + delete(id.inFlight, h) + } +} + // blockProcessor processes blocks to find Certen intents func (id *IntentDiscovery) blockProcessor(workerID string) { defer func() { @@ -701,7 +738,7 @@ func (id *IntentDiscovery) blockProcessor(workerID string) { return } id.logger.Printf("📦 Worker %s received job for block %d", workerID, job.BlockHeight) - if err := id.processBlock(job, workerID); err != nil { + if err := id.searchBlock(job, workerID); err != nil { id.logger.Printf("❌ Worker %s failed to search block %d: %v", workerID, job.BlockHeight, err) // The watermark passes a block only once it is searched or kept to be searched // (RB3-F125). It used to pass it regardless, on the theory that it would "appear again @@ -716,6 +753,19 @@ func (id *IntentDiscovery) blockProcessor(workerID string) { } } +// searchBlock searches one block. A panic in the search is a failed search (RB3-F151): the block is kept or +// handed back like any other failure, and the worker goes on. It used to end the worker - one fewer, for good +// - and the block was searched again only because every tick queued it again. +func (id *IntentDiscovery) searchBlock(job *BlockProcessJob, workerID string) (err error) { + defer func() { + if r := recover(); r != nil { + id.logger.Printf("🚨 PANIC searching block %d in %s: %v\n%s", job.BlockHeight, workerID, r, debug.Stack()) + err = fmt.Errorf("search of block %d panicked: %v", job.BlockHeight, r) + } + }() + return id.processBlock(job, workerID) +} + // keepUnsearched records a block whose search failed; false when it could not be kept (the caller then // must not let the watermark pass it). func (id *IntentDiscovery) keepUnsearched(height uint64, cause error) bool { @@ -906,11 +956,19 @@ func (id *IntentDiscovery) advanceWatermark(height uint64) { id.watermarkMu.Lock() defer id.watermarkMu.Unlock() + // Its job has finished (RB3-F151). + delete(id.inFlight, height) + // Ignore stale blocks (already processed or from before a network switch reset) if height <= id.lastProcessedBlock || height > id.lastQueuedBlock { return } + // A block searched while above the finalize ceiling is left unmarked, so the next tick searches it + // again: an intent whose block became queryable a tick after the first look is still found. + if height > id.finalizeCeiling { + return + } id.processedBlocks[height] = true // Advance lastProcessedBlock through contiguous completed blocks, but never past diff --git a/pkg/intent/discovery_inflight_test.go b/pkg/intent/discovery_inflight_test.go new file mode 100644 index 00000000..c038ebda --- /dev/null +++ b/pkg/intent/discovery_inflight_test.go @@ -0,0 +1,110 @@ +package intent + +import ( + "context" + "io" + "log" + "sort" + "sync" + "testing" + "time" +) + +func drainHeights(ch chan *BlockProcessJob) []uint64 { + var out []uint64 + for { + select { + case j := <-ch: + out = append(out, j.BlockHeight) + default: + sort.Slice(out, func(i, k int) bool { return out[i] < out[k] }) + return out + } + } +} + +// RB3-F151: while the workers are behind, every tick used to queue the whole range from the watermark to the +// head again, so each block was queued once per tick and new blocks waited behind the duplicates. A height +// in flight is queued once. +func TestATickQueuesNoBlockThatIsInFlight(t *testing.T) { + id := &IntentDiscovery{client: tipClient{tip: 120}, logger: log.New(io.Discard, "", 0), lastProcessedBlock: 100, + blockProcessCh: make(chan *BlockProcessJob, 1000), stopCh: make(chan struct{}), processedBlocks: map[uint64]bool{}} + for i := 0; i < 3; i++ { + if err := id.checkForNewBlocks(context.Background()); err != nil { + t.Fatal(err) + } + } + got := drainHeights(id.blockProcessCh) + if len(got) != 20 || got[0] != 101 || got[19] != 120 { + t.Fatalf("three ticks queued %d jobs (%v...), want each of 101-120 once", len(got), got[:min(5, len(got))]) + } +} + +// The tip is still searched again: a block searched while above the finalize ceiling is re-queued on the next +// tick, so an intent whose block became queryable a tick later is found. Blocks at or below the ceiling are +// finalized and not searched again. +func TestTheUnconfirmedTipIsSearchedAgain(t *testing.T) { + id := &IntentDiscovery{client: tipClient{tip: 120}, logger: log.New(io.Discard, "", 0), lastProcessedBlock: 100, + blockProcessCh: make(chan *BlockProcessJob, 1000), stopCh: make(chan struct{}), processedBlocks: map[uint64]bool{}} + if err := id.checkForNewBlocks(context.Background()); err != nil { + t.Fatal(err) + } + for _, h := range drainHeights(id.blockProcessCh) { + id.advanceWatermark(h) // every block searched + } + if id.lastProcessedBlock != 118 { + t.Fatalf("watermark %d, want the ceiling 118", id.lastProcessedBlock) + } + if err := id.checkForNewBlocks(context.Background()); err != nil { + t.Fatal(err) + } + if got := drainHeights(id.blockProcessCh); len(got) != 2 || got[0] != 119 || got[1] != 120 { + t.Fatalf("re-queued %v, want the unconfirmed tip 119-120", got) + } +} + +type recordingUnsearched struct { + mu sync.Mutex + kept []uint64 +} + +func (r *recordingUnsearched) Put(h uint64, _ string) error { + r.mu.Lock() + defer r.mu.Unlock() + r.kept = append(r.kept, h) + return nil +} +func (r *recordingUnsearched) List() ([]UnsearchedBlock, error) { return nil, nil } +func (r *recordingUnsearched) Remove(uint64) error { return nil } +func (r *recordingUnsearched) count() int { + r.mu.Lock() + defer r.mu.Unlock() + return len(r.kept) +} + +// A search that panics is a failed search: the block is kept for another search and leaves flight, and the +// worker goes on to the next job. It used to end the worker. +func TestAPanickingSearchIsAFailedSearchAndTheWorkerGoesOn(t *testing.T) { + store := &recordingUnsearched{} + id := &IntentDiscovery{client: tipClient{tip: 120}, // SearchCertenTransactions is not implemented: it panics + logger: log.New(io.Discard, "", 0), lastProcessedBlock: 100, lastQueuedBlock: 120, finalizeCeiling: 118, + blockProcessCh: make(chan *BlockProcessJob, 4), stopCh: make(chan struct{}), processedBlocks: map[uint64]bool{}, + inFlight: map[uint64]bool{101: true, 102: true}, unsearched: store} + defer close(id.stopCh) + go id.blockProcessor("w") + id.blockProcessCh <- &BlockProcessJob{BlockHeight: 101} + id.blockProcessCh <- &BlockProcessJob{BlockHeight: 102} + deadline := time.Now().Add(5 * time.Second) + for store.count() < 2 && time.Now().Before(deadline) { + time.Sleep(10 * time.Millisecond) + } + if store.count() != 2 { + t.Fatalf("%d blocks kept after panicking searches, want 2 (the worker must survive the first)", store.count()) + } + id.watermarkMu.Lock() + inFlight := len(id.inFlight) + id.watermarkMu.Unlock() + if inFlight != 0 { + t.Fatalf("%d heights still in flight after their jobs finished", inFlight) + } +}