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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
72 changes: 65 additions & 7 deletions pkg/intent/discovery.go
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down Expand Up @@ -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",
Expand All @@ -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
Expand All @@ -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
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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 {
Expand Down Expand Up @@ -654,31 +664,58 @@ 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
}
}

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() {
Expand All @@ -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
Expand All @@ -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 {
Expand Down Expand Up @@ -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
Expand Down
110 changes: 110 additions & 0 deletions pkg/intent/discovery_inflight_test.go
Original file line number Diff line number Diff line change
@@ -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)
}
}
Loading