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
87 changes: 57 additions & 30 deletions sei-db/db_engine/litt/disktable/forward_iterator.go
Original file line number Diff line number Diff line change
Expand Up @@ -20,8 +20,11 @@ var _ litt.Iterator = (*forwardIterator)(nil)
// segments, so their files remain on disk until Close releases them — even if garbage collection collects those
// segments meanwhile. Close is therefore mandatory: a leaked iterator pins its segments' files indefinitely.
type forwardIterator struct {
// table is the owning disk table, used to issue the close request.
table *DiskTable
// onClose is called once, by Close, after the buffered reader is closed. It performs whatever
// cleanup this iterator's owner requires: for a live table, releasing the reservation on each
// snapshot segment and notifying the control loop; for an offline iterator, releasing a directory
// lock instead.
onClose func() error

// segs is the ordered (lowest-to-highest index) snapshot of sealed segments in scope.
segs []*segment.Segment
Expand Down Expand Up @@ -63,12 +66,13 @@ type forwardIterator struct {
groupValue []byte
}

// newForwardIterator creates a forward iterator over the given snapshot of sealed segments.
// newForwardIterator creates a forward iterator over the given snapshot of sealed segments, owned by a
// live table.
func newForwardIterator(table *DiskTable, segs []*segment.Segment) *forwardIterator {
return &forwardIterator{
table: table,
segs: segs,
segPos: 0,
onClose: closeLiveIterator(table, segs),
segs: segs,
segPos: 0,
}
}

Expand All @@ -84,11 +88,25 @@ func newForwardIteratorAt(
keyPos int,
) *forwardIterator {
return &forwardIterator{
table: table,
onClose: closeLiveIterator(table, segs),
segs: segs,
segPos: segPos,
keys: keys,
keyPos: keyPos,
}
}

// NewOfflineForwardIterator creates a forward iterator over the given snapshot of segments, gathered
// directly from disk rather than from a live table. release is called once, by Close, in place of the
// live path's segment-reservation release and control-loop notification.
func NewOfflineForwardIterator(segs []*segment.Segment, release func()) litt.Iterator {
return &forwardIterator{
onClose: func() error {
release()
return nil
},
segs: segs,
segPos: segPos,
keys: keys,
keyPos: keyPos,
segPos: 0,
}
}

Expand Down Expand Up @@ -236,8 +254,7 @@ func (it *forwardIterator) secondaryWithinGroup(addr types.Address) bool {
uint64(addr.Offset())+uint64(addr.ValueSize()) <= end
}

// Close releases the resources held by the iterator, including the reservations on its snapshot segments
// (allowing any segment GC collected while it was open to finally be deleted from disk).
// Close releases the resources held by the iterator, via onClose.
func (it *forwardIterator) Close() error {
if it.closed {
return nil
Expand All @@ -251,27 +268,37 @@ func (it *forwardIterator) Close() error {
it.reader = nil
}

// Release the reservation on each snapshot segment. This must happen even on the error paths below: a missed
// release pins those segments' files on disk indefinitely.
for _, seg := range it.segs {
seg.Release()
}
closeErr := it.onClose()
it.segs = nil

// Notify the control loop so the open-iterator metric is updated.
request := &controlLoopCloseIteratorRequest{
completionChan: make(chan struct{}, 1),
}
err := it.table.controlLoop.enqueue(request)
if err != nil {
return fmt.Errorf("failed to send close iterator request: %w", err)
}
_, err = util.Await(it.table.errorMonitor, request.completionChan)
if err != nil {
return fmt.Errorf("failed to await iterator close: %w", err)
}
if readerErr != nil {
return fmt.Errorf("failed to close segment reader: %w", readerErr)
}
return nil
return closeErr
}

// closeLiveIterator returns the onClose function for an iterator owned by a live table: it releases the
// reservation on each snapshot segment (allowing any segment GC collected while the iterator was open to
// finally be deleted from disk), then notifies the control loop so the open-iterator metric is updated.
func closeLiveIterator(table *DiskTable, segs []*segment.Segment) func() error {
return func() error {
// This must happen even if the notification below fails: a missed release pins those segments'
// files on disk indefinitely.
for _, seg := range segs {
seg.Release()
}

request := &controlLoopCloseIteratorRequest{
completionChan: make(chan struct{}, 1),
}
err := table.controlLoop.enqueue(request)
if err != nil {
return fmt.Errorf("failed to send close iterator request: %w", err)
}
_, err = util.Await(table.errorMonitor, request.completionChan)
if err != nil {
return fmt.Errorf("failed to await iterator close: %w", err)
}
return nil
}
}
62 changes: 30 additions & 32 deletions sei-db/db_engine/litt/disktable/reverse_iterator.go
Original file line number Diff line number Diff line change
Expand Up @@ -6,7 +6,6 @@ import (
"github.com/sei-protocol/sei-chain/sei-db/db_engine/litt"
"github.com/sei-protocol/sei-chain/sei-db/db_engine/litt/disktable/segment"
"github.com/sei-protocol/sei-chain/sei-db/db_engine/litt/types"
"github.com/sei-protocol/sei-chain/sei-db/db_engine/litt/util"
)

var _ litt.Iterator = (*reverseIterator)(nil)
Expand All @@ -20,8 +19,10 @@ var _ litt.Iterator = (*reverseIterator)(nil)
// segments, so their files remain on disk until Close releases them — even if garbage collection collects those
// segments meanwhile. Close is therefore mandatory: a leaked iterator pins its segments' files indefinitely.
type reverseIterator struct {
// table is the owning disk table, used to issue the close request.
table *DiskTable
// onClose is called once, by Close. It performs whatever cleanup this iterator's owner requires:
// for a live table, releasing the reservation on each snapshot segment and notifying the control
// loop; for an offline iterator, releasing a directory lock instead.
onClose func() error

// segs is the ordered (lowest-to-highest index) snapshot of sealed segments in scope.
segs []*segment.Segment
Expand All @@ -45,12 +46,13 @@ type reverseIterator struct {
closed bool
}

// newReverseIterator creates a reverse iterator over the given snapshot of sealed segments.
// newReverseIterator creates a reverse iterator over the given snapshot of sealed segments, owned by a
// live table.
func newReverseIterator(table *DiskTable, segs []*segment.Segment) *reverseIterator {
return &reverseIterator{
table: table,
segs: segs,
segPos: len(segs) - 1,
onClose: closeLiveIterator(table, segs),
segs: segs,
segPos: len(segs) - 1,
}
}

Expand All @@ -66,11 +68,25 @@ func newReverseIteratorAt(
keyPos int,
) *reverseIterator {
return &reverseIterator{
table: table,
onClose: closeLiveIterator(table, segs),
segs: segs,
segPos: segPos,
keys: keys,
keyPos: keyPos,
}
}

// NewOfflineReverseIterator creates a reverse iterator over the given snapshot of segments, gathered
// directly from disk rather than from a live table. release is called once, by Close, in place of the
// live path's segment-reservation release and control-loop notification.
func NewOfflineReverseIterator(segs []*segment.Segment, release func()) litt.Iterator {
return &reverseIterator{
onClose: func() error {
release()
return nil
},
segs: segs,
segPos: segPos,
keys: keys,
keyPos: keyPos,
segPos: len(segs) - 1,
}
}

Expand Down Expand Up @@ -143,32 +159,14 @@ func (it *reverseIterator) GetValue() (value []byte, err error) {
return value, nil
}

// Close releases the resources held by the iterator, including the reservations on its snapshot segments
// (allowing any segment GC collected while it was open to finally be deleted from disk).
// Close releases the resources held by the iterator, via onClose.
func (it *reverseIterator) Close() error {
if it.closed {
return nil
}
it.closed = true

// Release the reservation on each snapshot segment. This must happen even on the error paths below: a missed
// release pins those segments' files on disk indefinitely.
for _, seg := range it.segs {
seg.Release()
}
closeErr := it.onClose()
it.segs = nil

// Notify the control loop so the open-iterator metric is updated.
request := &controlLoopCloseIteratorRequest{
completionChan: make(chan struct{}, 1),
}
err := it.table.controlLoop.enqueue(request)
if err != nil {
return fmt.Errorf("failed to send close iterator request: %w", err)
}
_, err = util.Await(it.table.errorMonitor, request.completionChan)
if err != nil {
return fmt.Errorf("failed to await iterator close: %w", err)
}
return nil
return closeErr
}
126 changes: 126 additions & 0 deletions sei-db/db_engine/litt/offline/iterator.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,126 @@
package offline

import (
"context"
"fmt"
"log/slog"
"path/filepath"
"time"

"github.com/sei-protocol/sei-chain/sei-db/db_engine/litt"
"github.com/sei-protocol/sei-chain/sei-db/db_engine/litt/disktable"
"github.com/sei-protocol/sei-chain/sei-db/db_engine/litt/disktable/segment"
"github.com/sei-protocol/sei-chain/sei-db/db_engine/litt/util"
)

// NewIterator opens an offline iterator over the table at tableName under config.Paths, without starting
// a live database. It iterates oldest-to-newest if reverse is false, newest-to-oldest if reverse is true.
//
// The database must NOT be running while this is called: NewIterator takes the same directory lock the
// database uses, held for the iterator's lifetime and released by Close, so it will fail rather than read
// alongside a live database.
func NewIterator(config *litt.Config, tableName string, reverse bool) (litt.Iterator, error) {
logger := slog.Default()

if config == nil || len(config.Paths) == 0 {
return nil, fmt.Errorf("at least one path must be provided")
}
if err := config.SanitizePaths(); err != nil {
return nil, fmt.Errorf("failed to sanitize data directories: %w", err)
}
roots := config.Paths

for _, root := range roots {
if err := util.EnsureDirectoryExists(root, config.Fsync); err != nil {
return nil, fmt.Errorf("failed to ensure data directory %q: %w", root, err)
}
}

releaseLocks, err := util.LockDirectories(logger, roots, util.LockfileName, config.Fsync)
if err != nil {
return nil, fmt.Errorf("failed to lock data directories %v: %w", roots, err)
}

segs, err := gatherOrderedSegments(logger, roots, tableName, config.Fsync)
if err != nil {
releaseLocks()
return nil, err
}

if reverse {
return disktable.NewOfflineReverseIterator(segs, releaseLocks), nil
}
return disktable.NewOfflineForwardIterator(segs, releaseLocks), nil
}

// gatherOrderedSegments enumerates a table's segments across roots, offline, and returns them ordered from
// lowest to highest index. Segments already logically garbage collected (below the table's durable
// gc-watermark) are excluded, matching what a live table's own reads would see. Returns an empty slice, not
// an error, if the table has no segments.
func gatherOrderedSegments(
logger *slog.Logger,
roots []string,
tableName string,
fsync bool,
) ([]*segment.Segment, error) {
exists, err := tableExists(roots, tableName)
if err != nil {
return nil, err
}
if !exists {
return nil, nil
}

errorMonitor := util.NewErrorMonitor(context.Background(), logger, nil)

segmentPaths, err := segment.BuildSegmentPaths(roots, "", tableName)
if err != nil {
return nil, fmt.Errorf("failed to build segment paths: %w", err)
}

lowestSegmentIndex, highestSegmentIndex, segments, err := segment.GatherSegmentFiles(
logger, errorMonitor, segmentPaths, false /* snapshottingEnabled */, time.Now(),
true /* cleanOrphans */, fsync)
Comment thread
cody-littley marked this conversation as resolved.

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[suggestion] Re-raising with a consequence that is new in this revision rather than the earlier read-only argument (which you declined). PruneAfter now calls GetRange from inside requireTargetWithinStoredRange, i.e. as part of its pre-mutation guard, and its godoc promises "Both are checked before any receipt is touched, so a refusal leaves the store's data unchanged."

That is not quite true: reaching the guard runs GatherSegmentFiles with cleanOrphans=true, which deletes orphaned segment files and garbage files, and NewIterator creates the data directories at line 34 first. So a PruneAfter that refuses — the below-the-floor and below-the-oldest-receipt paths both — has already mutated the directory, destroying exactly the leftovers an operator would want after the refusal prompts a post-mortem.

Passing false from gatherOrderedSegments loads the same segments without the cleanup; the rollback path is the one that legitimately wants true. Alternatively, weaken the guarantee in PruneAfter's doc so it does not claim more than the code delivers.

if err != nil {
return nil, fmt.Errorf("failed to gather segment files: %w", err)
}
if len(segments) == 0 {
return nil, nil
}

watermark, defined, err := highestGCWatermark(roots, tableName)
if err != nil {
return nil, err
}
if defined && watermark > highestSegmentIndex {
// Everything present is below the watermark; no readable segments remain.
return nil, nil
}
floor := lowestSegmentIndex
if defined && watermark > floor {
floor = watermark
}

ordered := make([]*segment.Segment, 0, highestSegmentIndex-floor+1)
for index := floor; index <= highestSegmentIndex; index++ {
ordered = append(ordered, segments[index])
}
return ordered, nil
}

// tableExists reports whether tableName has ever been created under any of roots. A table that has never
// been written has no segments directory, and scanning for its segment files would otherwise fail with a
// "no such file or directory" error rather than reporting it as merely empty.
func tableExists(roots []string, tableName string) (bool, error) {
for _, root := range roots {
segmentsDir := filepath.Join(root, tableName, segment.SegmentDirectory)
isDir, err := util.IsDirectory(segmentsDir)
if err != nil {
return false, fmt.Errorf("failed to check directory %s: %w", segmentsDir, err)
}
if isDir {
return true, nil
}
}
return false, nil
}
Loading
Loading