diff --git a/sei-db/state_db/sc/memiavl/db.go b/sei-db/state_db/sc/memiavl/db.go index ddd18f675c..feb1fee485 100644 --- a/sei-db/state_db/sc/memiavl/db.go +++ b/sei-db/state_db/sc/memiavl/db.go @@ -27,7 +27,21 @@ import ( const LockFileName = "LOCK" -var errReadOnly = errors.New("db is read-only") +var ( + errReadOnly = errors.New("db is read-only") + + // ErrReadOnlyWALCorrupt means a read-only open observed an incomplete, + // corrupt, or concurrently recovering changelog. The source WAL is left + // untouched; callers can retry after the writer finishes its current WAL + // operation. + ErrReadOnlyWALCorrupt = errors.New("read-only changelog is incomplete or corrupt") + + // ErrReadOnlyWALUnavailable means the immutable WAL view cannot replay + // every version from the selected snapshot through the requested target. + // The live writer may have pruned or advanced the changelog while the + // reader opened it; callers can retry against a new point-in-time view. + ErrReadOnlyWALUnavailable = errors.New("read-only changelog cannot reach the requested version") +) // DB implements DB-like functionalities on top of MultiTree: // - async snapshot rewriting @@ -158,9 +172,26 @@ func OpenDB(targetVersion int64, opts Options) (database *DB, _err error) { ) }() var ( - err error - fileLock FileLock + err error + fileLock FileLock + mtree *MultiTree + streamHandler wal.ChangelogWAL ) + defer func() { + if _err == nil { + return + } + if streamHandler != nil { + _ = streamHandler.Close() + } + if mtree != nil { + _ = mtree.Close() + } + if fileLock != nil { + _ = fileLock.Unlock() + _ = fileLock.Destroy() + } + }() if err := opts.Validate(); err != nil { return nil, fmt.Errorf("invalid commit store options: %w", err) } @@ -194,19 +225,28 @@ func OpenDB(targetVersion int64, opts Options) (database *DB, _err error) { } path := filepath.Join(opts.Dir, snapshot) - mtree, err := LoadMultiTree(context.Background(), path, opts) + mtree, err = LoadMultiTree(context.Background(), path, opts) if err != nil { return nil, err } // Snapshot mmap files are loaded with MADV_RANDOM in OpenSnapshot(). - // MemIAVL owns changelog lifecycle: always open the WAL here. - // Even in read-only mode we may need WAL replay to reconstruct non-snapshot versions. - streamHandler, err := wal.NewChangelogWAL(utils.GetChangelogPath(opts.Dir), wal.Config{ - WriteBufferSize: opts.AsyncCommitBuffer, - }) + // MemIAVL owns changelog lifecycle: always open the WAL here. Read-only + // callers still need replay to reconstruct non-snapshot versions, but they + // must not use the writable opener: it repairs a torn tail by truncating it + // and completes interrupted WAL truncations by renaming or removing files. + if opts.ReadOnly { + streamHandler, err = wal.OpenReadOnlyChangelogWAL(utils.GetChangelogPath(opts.Dir)) + } else { + streamHandler, err = wal.NewChangelogWAL(utils.GetChangelogPath(opts.Dir), wal.Config{ + WriteBufferSize: opts.AsyncCommitBuffer, + }) + } if err != nil { + if opts.ReadOnly && errors.Is(err, wal.ErrCorrupt) { + return nil, fmt.Errorf("%w; source WAL was not modified: %w", ErrReadOnlyWALCorrupt, err) + } return nil, fmt.Errorf("failed to open changelog WAL: %w", err) } @@ -221,6 +261,23 @@ func OpenDB(targetVersion int64, opts Options) (database *DB, _err error) { if !walHasEntries { walIndexDelta = mtree.WorkingCommitInfo().Version - 1 } + if opts.ReadOnly && walHasEntries && (targetVersion == 0 || targetVersion > mtree.Version()) { + firstIndex, firstErr := streamHandler.FirstOffset() + if firstErr != nil { + return nil, fmt.Errorf("read changelog first offset: %w", firstErr) + } + if firstIndex > math.MaxInt64 { + return nil, fmt.Errorf("%w: first WAL offset %d overflows int64", ErrReadOnlyWALUnavailable, firstIndex) + } + firstVersion := int64(firstIndex) + walIndexDelta + firstNeeded := utils.NextVersion(mtree.Version(), mtree.initialVersion.Load()) + if firstVersion > firstNeeded { + snapshotVersion := mtree.Version() + return nil, fmt.Errorf("%w: selected snapshot version %d needs changelog version %d, "+ + "but the immutable WAL view starts at version %d", + ErrReadOnlyWALUnavailable, snapshotVersion, firstNeeded, firstVersion) + } + } // Replay WAL to catch up to target version (if WAL has entries) if walHasEntries && (targetVersion == 0 || targetVersion > mtree.Version()) { @@ -230,6 +287,11 @@ func OpenDB(targetVersion int64, opts Options) (database *DB, _err error) { } logger.Info("finished replay and caught up to target version", "version", targetVersion) } + if opts.ReadOnly && targetVersion > 0 && mtree.Version() != targetVersion { + reached := mtree.Version() + return nil, fmt.Errorf("%w: requested %d, reached %d", + ErrReadOnlyWALUnavailable, targetVersion, reached) + } if opts.LoadForOverwriting && targetVersion > 0 { currentSnapshot, err := os.Readlink(currentPath(opts.Dir)) @@ -296,6 +358,10 @@ func OpenDB(targetVersion int64, opts Options) (database *DB, _err error) { snapshotWriterPool: workerPool, opts: opts, } + // The DB owns these resources from this point forward. + mtree = nil + streamHandler = nil + fileLock = nil // Apply initial stores on a fresh DB (version 0) so they get persisted to WAL. // This creates the trees and populates pendingLogEntry, which will be written @@ -307,6 +373,7 @@ func OpenDB(targetVersion int64, opts Options) (database *DB, _err error) { upgrades = append(upgrades, &proto.TreeNameUpgrade{Name: name}) } if err := db.ApplyUpgrades(upgrades); err != nil { + _ = db.Close() return nil, fmt.Errorf("failed to apply initial stores: %w", err) } } diff --git a/sei-db/state_db/sc/memiavl/db_test.go b/sei-db/state_db/sc/memiavl/db_test.go index 3cddd00f77..d804189902 100644 --- a/sei-db/state_db/sc/memiavl/db_test.go +++ b/sei-db/state_db/sc/memiavl/db_test.go @@ -6,6 +6,7 @@ import ( "os" "path/filepath" "runtime/debug" + "sort" "strconv" "sync" "testing" @@ -1116,3 +1117,132 @@ func TestUpdateCurrentSymlinkClearsStaleTmp(t *testing.T) { require.NoError(t, err) require.Equal(t, "snapshot-1", target) } + +func TestReadOnlyOpenRejectsTornWALWithoutRepair(t *testing.T) { + dir := t.TempDir() + db, err := OpenDB(0, Options{ + Dir: dir, + CreateIfMissing: true, + InitialStores: []string{"test"}, + }) + require.NoError(t, err) + for i := 0; i < 3; i++ { + require.NoError(t, db.ApplyChangeSets([]*proto.NamedChangeSet{{ + Name: "test", + Changeset: ChangeSets[i], + }})) + _, err := db.Commit() + require.NoError(t, err) + } + require.NoError(t, db.Close()) + + segment := lastMemiAVLWALSegment(t, dir) + file, err := os.OpenFile(filepath.Clean(segment), os.O_WRONLY|os.O_APPEND, 0) + require.NoError(t, err) + _, err = file.Write([]byte{0x10}) + require.NoError(t, err) + require.NoError(t, file.Close()) + before, err := os.ReadFile(filepath.Clean(segment)) + require.NoError(t, err) + + _, err = OpenDB(0, Options{Dir: dir, ReadOnly: true}) + require.ErrorIs(t, err, ErrReadOnlyWALCorrupt) + after, readErr := os.ReadFile(filepath.Clean(segment)) + require.NoError(t, readErr) + require.Equal(t, before, after, "read-only open must leave a torn live tail untouched") + + repaired, err := OpenDB(0, Options{Dir: dir}) + require.NoError(t, err, "the writable owner must retain the existing tail-repair behavior") + require.Equal(t, int64(3), repaired.Version()) + require.NoError(t, repaired.Close()) +} + +func TestReadOnlyOpenRejectsWALGap(t *testing.T) { + dir := t.TempDir() + db, err := OpenDB(0, Options{ + Dir: dir, + CreateIfMissing: true, + InitialStores: []string{"test"}, + }) + require.NoError(t, err) + for i := 0; i < 3; i++ { + require.NoError(t, db.ApplyChangeSets([]*proto.NamedChangeSet{{ + Name: "test", + Changeset: ChangeSets[i], + }})) + _, err := db.Commit() + require.NoError(t, err) + } + require.NoError(t, db.GetWAL().TruncateBefore(2)) + require.NoError(t, db.Close()) + + _, err = OpenDB(3, Options{Dir: dir, ReadOnly: true}) + require.ErrorIs(t, err, ErrReadOnlyWALUnavailable) + require.Contains(t, err.Error(), "needs changelog version 1") + require.Contains(t, err.Error(), "starts at version 2") +} + +func TestReadOnlyOpenRejectsShortWAL(t *testing.T) { + dir := t.TempDir() + db, err := OpenDB(0, Options{ + Dir: dir, + CreateIfMissing: true, + InitialStores: []string{"test"}, + }) + require.NoError(t, err) + for i := 0; i < 3; i++ { + require.NoError(t, db.ApplyChangeSets([]*proto.NamedChangeSet{{ + Name: "test", + Changeset: ChangeSets[i], + }})) + _, err := db.Commit() + require.NoError(t, err) + } + require.NoError(t, db.GetWAL().TruncateAfter(2)) + require.NoError(t, db.Close()) + + _, err = OpenDB(3, Options{Dir: dir, ReadOnly: true}) + require.ErrorIs(t, err, ErrReadOnlyWALUnavailable) + require.Contains(t, err.Error(), "requested 3, reached 2") +} + +func TestOpenDBFailureReleasesFileLock(t *testing.T) { + dir := t.TempDir() + db, err := OpenDB(0, Options{ + Dir: dir, + CreateIfMissing: true, + InitialStores: []string{"test"}, + }) + require.NoError(t, err) + require.NoError(t, db.ApplyChangeSets([]*proto.NamedChangeSet{{ + Name: "test", + Changeset: ChangeSets[0], + }})) + _, err = db.Commit() + require.NoError(t, err) + require.NoError(t, db.Close()) + + require.NoError(t, os.Remove(currentPath(dir))) + _, err = OpenDB(1, Options{Dir: dir, LoadForOverwriting: true}) + require.ErrorContains(t, err, "fail to read current version") + + lock, err := LockFile(filepath.Join(dir, LockFileName)) + require.NoError(t, err, "failed OpenDB must release its exclusive lock") + require.NoError(t, lock.Unlock()) + require.NoError(t, lock.Destroy()) +} + +func lastMemiAVLWALSegment(t *testing.T, dir string) string { + t.Helper() + entries, err := os.ReadDir(utils.GetChangelogPath(dir)) + require.NoError(t, err) + var names []string + for _, entry := range entries { + if !entry.IsDir() && len(entry.Name()) == 20 { + names = append(names, entry.Name()) + } + } + require.NotEmpty(t, names) + sort.Strings(names) + return filepath.Join(utils.GetChangelogPath(dir), names[len(names)-1]) +} diff --git a/sei-db/tools/cmd/seidb/operations/evm_logical_digest.go b/sei-db/tools/cmd/seidb/operations/evm_logical_digest.go index 79ecd140c6..49a3693907 100644 --- a/sei-db/tools/cmd/seidb/operations/evm_logical_digest.go +++ b/sei-db/tools/cmd/seidb/operations/evm_logical_digest.go @@ -81,12 +81,14 @@ const ( // that exact height (or --height 0 for the current symlink). This is the // preferred mode whenever the target height lines up with an existing // snapshot boundary. -// - replay (SLOW): opens a read-only DB, replays the changelog up to -// --height, then walks the in-memory/mmap tree. Roughly an order of -// magnitude slower than snapshot (changelog replay + per-leaf tree walk -// instead of a sequential file read). Use it only when no snapshot exists -// at the target height — e.g. nodes whose snapshot rewrite lags the tip, so -// an arbitrary comparison height has no snapshot- on disk. +// - replay (SLOW): opens a non-mutating read-only WAL view, replays the +// changelog up to --height, then walks the in-memory/mmap tree. Roughly an +// order of magnitude slower than snapshot (changelog replay + per-leaf tree +// walk instead of a sequential file read). If a live writer leaves a torn +// tail in view, replay fails and asks the operator to rerun instead of +// repairing the source WAL. Use it only when no snapshot exists at the +// target height — e.g. nodes whose snapshot rewrite lags the tip, so an +// arbitrary comparison height has no snapshot- on disk. // // The flatkv side is always a pebble WAL-replay-to-height and is fast // regardless. So when comparing across nodes, pick a height that is an existing @@ -1111,8 +1113,26 @@ func openMemiAVLReplayReadOnly(dbDir string, height int64) (*memiavl.DB, error) ZeroCopy: true, }) if err != nil { + if errors.Is(err, memiavl.ErrReadOnlyWALCorrupt) { + return nil, fmt.Errorf("memiavl changelog tail is incomplete, corrupt, or changing; "+ + "live WAL was not modified; rerun the command, and if the error persists after stopping seid, "+ + "repair the WAL offline: %w", err) + } + if errors.Is(err, memiavl.ErrReadOnlyWALUnavailable) { + return nil, fmt.Errorf("the immutable memiavl changelog view could not reach height %d; "+ + "live WAL was not modified; rerun the command: %w", height, err) + } return nil, fmt.Errorf("open memiavl read-only replay: %w", err) } + if height > 0 && db.Version() != height { + versionErr := fmt.Errorf("memiavl replay version mismatch: requested %d, reached %d; "+ + "the live changelog did not provide a complete path to the target; rerun the command", + height, db.Version()) + if closeErr := db.Close(); closeErr != nil { + return nil, errors.Join(versionErr, fmt.Errorf("close memiavl read-only replay: %w", closeErr)) + } + return nil, versionErr + } return db, nil } diff --git a/sei-db/tools/cmd/seidb/operations/memiavl_open_test.go b/sei-db/tools/cmd/seidb/operations/memiavl_open_test.go new file mode 100644 index 0000000000..ad06579c35 --- /dev/null +++ b/sei-db/tools/cmd/seidb/operations/memiavl_open_test.go @@ -0,0 +1,61 @@ +package operations + +import ( + "os" + "path/filepath" + "sort" + "testing" + + "github.com/stretchr/testify/require" + + "github.com/sei-protocol/sei-chain/sei-db/common/keys" + "github.com/sei-protocol/sei-chain/sei-db/common/utils" + "github.com/sei-protocol/sei-chain/sei-db/proto" +) + +func TestOpenMemiAVLReplayReadOnlyReportsRetryWithoutRepair(t *testing.T) { + homeDir := t.TempDir() + store := newTestMemiavlStore(t, homeDir) + require.NoError(t, store.ApplyChangeSets([]*proto.NamedChangeSet{{ + Name: keys.EVMStoreKey, + Changeset: proto.ChangeSet{Pairs: []*proto.KVPair{noncePair(addrN(0xA1), 1)}}, + }})) + _, err := store.Commit() + require.NoError(t, err) + require.NoError(t, store.Close()) + + dbDir := utils.GetCosmosSCStorePath(homeDir) + segment := lastOperationsMemiAVLWALSegment(t, dbDir) + file, err := os.OpenFile(filepath.Clean(segment), os.O_WRONLY|os.O_APPEND, 0) + require.NoError(t, err) + _, err = file.Write([]byte{0x10}) + require.NoError(t, err) + require.NoError(t, file.Close()) + before, err := os.ReadFile(filepath.Clean(segment)) + require.NoError(t, err) + + _, err = openMemiAVLReplayReadOnly(dbDir, 0) + require.Error(t, err) + require.Contains(t, err.Error(), "live WAL was not modified") + require.Contains(t, err.Error(), "rerun the command") + + after, readErr := os.ReadFile(filepath.Clean(segment)) + require.NoError(t, readErr) + require.Equal(t, before, after) +} + +func lastOperationsMemiAVLWALSegment(t *testing.T, dbDir string) string { + t.Helper() + changelogDir := utils.GetChangelogPath(dbDir) + entries, err := os.ReadDir(changelogDir) + require.NoError(t, err) + var names []string + for _, entry := range entries { + if !entry.IsDir() && len(entry.Name()) == 20 { + names = append(names, entry.Name()) + } + } + require.NotEmpty(t, names) + sort.Strings(names) + return filepath.Join(changelogDir, names[len(names)-1]) +} diff --git a/sei-db/wal/changelog.go b/sei-db/wal/changelog.go index fadb054f31..c3071b9679 100644 --- a/sei-db/wal/changelog.go +++ b/sei-db/wal/changelog.go @@ -26,6 +26,25 @@ func NewChangelogWAL(dir string, config Config) (ChangelogWAL, error) { ) } +// OpenReadOnlyChangelogWAL opens an immutable point-in-time view of the +// changelog segment files. It never creates, truncates, removes, or renames WAL +// files. A torn tail or an in-progress recovery marker returns ErrCorrupt so +// callers can fail and retry after the writer moves on. +func OpenReadOnlyChangelogWAL(dir string) (ChangelogWAL, error) { + readOnly, err := openReadOnlyWAL( + dir, + func(data []byte) (proto.ChangelogEntry, error) { + var entry proto.ChangelogEntry + err := entry.Unmarshal(data) + return entry, err + }, + ) + if err != nil { + return nil, err + } + return readOnly, nil +} + // FindFirstOffsetAfterVersion returns the first WAL offset whose entry version is // strictly greater than targetVersion. If no such entry exists, it returns // lastOffset+1. Changelog versions are monotonic, but empty blocks can advance diff --git a/sei-db/wal/readonly.go b/sei-db/wal/readonly.go new file mode 100644 index 0000000000..884ca77394 --- /dev/null +++ b/sei-db/wal/readonly.go @@ -0,0 +1,248 @@ +package wal + +import ( + "encoding/binary" + "errors" + "fmt" + "io" + "os" + "path/filepath" + "sort" + "strconv" + "strings" + "sync" + "sync/atomic" + + tidwallwal "github.com/tidwall/wal" +) + +var ( + // ErrReadOnly is returned when a caller tries to mutate a read-only WAL. + ErrReadOnly = errors.New("WAL is read-only") + + // ErrCorrupt identifies a malformed or unstable WAL view. It aliases the + // underlying tidwall sentinel so callers do not need to import the storage + // implementation only to decide whether a read-only open can be retried. + ErrCorrupt = tidwallwal.ErrCorrupt +) + +type readOnlySegment struct { + name string + index uint64 +} + +type readOnlyEntry struct { + file *os.File + dataOffset int64 + size int +} + +// readOnlyWAL is an immutable view of the plain segment files present when it +// opens. It does not use tidwall/wal.Open because that function creates files, +// opens the tail for writing, and completes interrupted truncations by removing +// and renaming segment files. +type readOnlyWAL[T any] struct { + unmarshal UnmarshalFn[T] + files []*os.File + entries []readOnlyEntry + firstOffset uint64 + closed atomic.Bool + closeOnce sync.Once + closeErr error +} + +func openReadOnlyWAL[T any](dir string, unmarshal UnmarshalFn[T]) (*readOnlyWAL[T], error) { + segments, err := listReadOnlySegments(dir) + if err != nil { + return nil, err + } + + log := &readOnlyWAL[T]{unmarshal: unmarshal} + cleanup := func(err error) (*readOnlyWAL[T], error) { + _ = log.Close() + return nil, err + } + + var nextIndex uint64 + for i, segment := range segments { + if i > 0 && segment.index != nextIndex { + return cleanup(fmt.Errorf("%w: segment %s starts at index %d, expected %d", + ErrCorrupt, segment.name, segment.index, nextIndex)) + } + + path := filepath.Join(dir, segment.name) + file, err := os.Open(filepath.Clean(path)) + if err != nil { + return cleanup(fmt.Errorf("open WAL segment %s: %w", path, err)) + } + log.files = append(log.files, file) + + info, err := file.Stat() + if err != nil { + return cleanup(fmt.Errorf("stat WAL segment %s: %w", path, err)) + } + if !info.Mode().IsRegular() { + return cleanup(fmt.Errorf("%w: WAL segment %s is not a regular file", ErrCorrupt, path)) + } + + data, err := io.ReadAll(io.NewSectionReader(file, 0, info.Size())) + if err != nil { + return cleanup(fmt.Errorf("read WAL segment %s: %w", path, err)) + } + if int64(len(data)) != info.Size() { + return cleanup(fmt.Errorf("%w: WAL segment %s changed while it was read", + ErrCorrupt, path)) + } + + entries, err := indexReadOnlySegment(file, data) + if err != nil { + return cleanup(fmt.Errorf("index WAL segment %s: %w", path, err)) + } + if len(entries) == 0 && i != len(segments)-1 { + return cleanup(fmt.Errorf("%w: non-tail WAL segment %s is empty", ErrCorrupt, path)) + } + if len(log.entries) == 0 && len(entries) > 0 { + log.firstOffset = segment.index + } + log.entries = append(log.entries, entries...) + nextIndex = segment.index + uint64(len(entries)) + } + return log, nil +} + +func listReadOnlySegments(dir string) ([]readOnlySegment, error) { + entries, err := os.ReadDir(dir) + if err != nil { + return nil, fmt.Errorf("read WAL directory %s: %w", dir, err) + } + + segments := make([]readOnlySegment, 0, len(entries)) + for _, entry := range entries { + if entry.IsDir() { + continue + } + name := entry.Name() + if strings.HasSuffix(name, ".START") || strings.HasSuffix(name, ".END") { + return nil, fmt.Errorf("%w: WAL recovery marker %s is present; retry after the writer finishes", + ErrCorrupt, name) + } + if len(name) != 20 { + continue + } + index, err := strconv.ParseUint(name, 10, 64) + if err != nil || index == 0 { + continue + } + segments = append(segments, readOnlySegment{name: name, index: index}) + } + sort.Slice(segments, func(i, j int) bool { + return segments[i].index < segments[j].index + }) + return segments, nil +} + +func indexReadOnlySegment(file *os.File, data []byte) ([]readOnlyEntry, error) { + entries := make([]readOnlyEntry, 0) + for pos := 0; pos < len(data); { + recordLen, err := loadNextBinaryEntry(data[pos:]) + if err != nil { + return nil, err + } + size, prefixLen := binary.Uvarint(data[pos:]) + entrySize := int(size) //nolint:gosec // loadNextBinaryEntry rejects sizes above math.MaxInt32. + entries = append(entries, readOnlyEntry{ + file: file, + dataOffset: int64(pos + prefixLen), + size: entrySize, + }) + pos += recordLen + } + return entries, nil +} + +func (log *readOnlyWAL[T]) Write(T) error { + return ErrReadOnly +} + +func (log *readOnlyWAL[T]) TruncateBefore(uint64) error { + return ErrReadOnly +} + +func (log *readOnlyWAL[T]) TruncateAfter(uint64) error { + return ErrReadOnly +} + +func (log *readOnlyWAL[T]) TruncateAll() error { + return ErrReadOnly +} + +func (log *readOnlyWAL[T]) FirstOffset() (uint64, error) { + if log.closed.Load() { + return 0, os.ErrClosed + } + if len(log.entries) == 0 { + return 0, nil + } + return log.firstOffset, nil +} + +func (log *readOnlyWAL[T]) LastOffset() (uint64, error) { + if log.closed.Load() { + return 0, os.ErrClosed + } + if len(log.entries) == 0 { + return 0, nil + } + return log.firstOffset + uint64(len(log.entries)) - 1, nil +} + +func (log *readOnlyWAL[T]) ReadAt(index uint64) (T, error) { + var zero T + if log.closed.Load() { + return zero, os.ErrClosed + } + if index < log.firstOffset || index-log.firstOffset >= uint64(len(log.entries)) { + return zero, fmt.Errorf("read WAL offset %d: out of range", index) + } + + entry := log.entries[index-log.firstOffset] + data := make([]byte, entry.size) + if _, err := entry.file.ReadAt(data, entry.dataOffset); err != nil { + return zero, fmt.Errorf("read WAL offset %d: %w", index, err) + } + value, err := log.unmarshal(data) + if err != nil { + return zero, fmt.Errorf("unmarshal WAL offset %d: %w", index, err) + } + return value, nil +} + +func (log *readOnlyWAL[T]) Replay(start, end uint64, processFn func(index uint64, entry T) error) error { + if end < start { + return nil + } + for index := start; index <= end; index++ { + entry, err := log.ReadAt(index) + if err != nil { + return err + } + if err := processFn(index, entry); err != nil { + return fmt.Errorf("process WAL offset %d: %w", index, err) + } + } + return nil +} + +func (log *readOnlyWAL[T]) Close() error { + log.closeOnce.Do(func() { + log.closed.Store(true) + var errs []error + for _, file := range log.files { + if err := file.Close(); err != nil { + errs = append(errs, err) + } + } + log.closeErr = errors.Join(errs...) + }) + return log.closeErr +} diff --git a/sei-db/wal/readonly_test.go b/sei-db/wal/readonly_test.go new file mode 100644 index 0000000000..9fb0ebc3c0 --- /dev/null +++ b/sei-db/wal/readonly_test.go @@ -0,0 +1,240 @@ +package wal + +import ( + "os" + "path/filepath" + "sort" + "testing" + + "github.com/stretchr/testify/require" + + "github.com/sei-protocol/sei-chain/sei-db/proto" +) + +func TestOpenReadOnlyChangelogWALReplaysWithoutMutation(t *testing.T) { + dir := t.TempDir() + writable, err := NewChangelogWAL(dir, Config{}) + require.NoError(t, err) + writeReadOnlyTestData(t, writable) + require.NoError(t, writable.Close()) + + before := snapshotWALFiles(t, dir) + readOnly, err := OpenReadOnlyChangelogWAL(dir) + require.NoError(t, err) + + first, err := readOnly.FirstOffset() + require.NoError(t, err) + require.Equal(t, uint64(1), first) + last, err := readOnly.LastOffset() + require.NoError(t, err) + require.Equal(t, uint64(3), last) + + var names []string + require.NoError(t, readOnly.Replay(first, last, func(_ uint64, entry proto.ChangelogEntry) error { + names = append(names, entry.Changesets[0].Name) + return nil + })) + require.Equal(t, []string{"test", "test", "test"}, names) + require.Equal(t, before, snapshotWALFiles(t, dir)) + + require.ErrorIs(t, readOnly.Write(proto.ChangelogEntry{}), ErrReadOnly) + require.ErrorIs(t, readOnly.TruncateBefore(2), ErrReadOnly) + require.ErrorIs(t, readOnly.TruncateAfter(2), ErrReadOnly) + require.NoError(t, readOnly.Close()) + require.NoError(t, readOnly.Close()) +} + +func TestOpenReadOnlyChangelogWALRejectsTornTailWithoutRepair(t *testing.T) { + dir := t.TempDir() + writable, err := NewChangelogWAL(dir, Config{}) + require.NoError(t, err) + writeReadOnlyTestData(t, writable) + require.NoError(t, writable.Close()) + + segment := lastPlainWALSegment(t, dir) + file, err := os.OpenFile(filepath.Clean(segment), os.O_WRONLY|os.O_APPEND, 0) + require.NoError(t, err) + _, err = file.Write([]byte{0x10}) // declares a 16-byte record whose payload has not arrived + require.NoError(t, err) + require.NoError(t, file.Close()) + before := snapshotWALFiles(t, dir) + + _, err = OpenReadOnlyChangelogWAL(dir) + require.ErrorIs(t, err, ErrCorrupt) + require.Equal(t, before, snapshotWALFiles(t, dir), "read-only open must not repair the source tail") +} + +func TestOpenReadOnlyChangelogWALKeepsPointInTimeView(t *testing.T) { + dir := t.TempDir() + writable, err := NewChangelogWAL(dir, Config{}) + require.NoError(t, err) + require.NoError(t, writable.Write(proto.ChangelogEntry{Version: 1})) + + readOnly, err := OpenReadOnlyChangelogWAL(dir) + require.NoError(t, err) + t.Cleanup(func() { require.NoError(t, readOnly.Close()) }) + + require.NoError(t, writable.Write(proto.ChangelogEntry{Version: 2})) + require.NoError(t, writable.Close()) + + last, err := readOnly.LastOffset() + require.NoError(t, err) + require.Equal(t, uint64(1), last) + _, err = readOnly.ReadAt(2) + require.Error(t, err) +} + +func TestOpenReadOnlyChangelogWALRejectsRecoveryMarkers(t *testing.T) { + for _, suffix := range []string{".START", ".END"} { + t.Run(suffix, func(t *testing.T) { + dir := t.TempDir() + marker := filepath.Join(dir, "00000000000000000001"+suffix) + require.NoError(t, os.WriteFile(marker, nil, 0o600)) + before := snapshotWALFiles(t, dir) + + _, err := OpenReadOnlyChangelogWAL(dir) + require.ErrorIs(t, err, ErrCorrupt) + require.Equal(t, before, snapshotWALFiles(t, dir), + "read-only open must not complete writable WAL recovery") + }) + } +} + +func TestOpenReadOnlyChangelogWALEmptyDirectory(t *testing.T) { + dir := t.TempDir() + + readOnly, err := OpenReadOnlyChangelogWAL(dir) + require.NoError(t, err) + first, err := readOnly.FirstOffset() + require.NoError(t, err) + require.Zero(t, first) + last, err := readOnly.LastOffset() + require.NoError(t, err) + require.Zero(t, last) + require.NoError(t, readOnly.Close()) + + entries, err := os.ReadDir(dir) + require.NoError(t, err) + require.Empty(t, entries, "read-only open must not create an initial segment") +} + +func TestOpenReadOnlyChangelogWALDoesNotCreateMissingDirectory(t *testing.T) { + dir := filepath.Join(t.TempDir(), "missing") + + _, err := OpenReadOnlyChangelogWAL(dir) + require.Error(t, err) + require.NoDirExists(t, dir) +} + +func TestOpenReadOnlyChangelogWALConcurrentWriter(t *testing.T) { + dir := t.TempDir() + writable, err := NewChangelogWAL(dir, Config{}) + require.NoError(t, err) + + const versions = 500 + done := make(chan struct{}) + writerErr := make(chan error, 1) + go func() { + var runErr error + for version := 1; version <= versions && runErr == nil; version++ { + runErr = writable.Write(proto.ChangelogEntry{Version: int64(version)}) + if runErr != nil || version <= 20 || version%10 != 0 { + continue + } + first, err := writable.FirstOffset() + if err != nil { + runErr = err + break + } + keepFrom := uint64(version - 20) + if keepFrom > first { + runErr = writable.TruncateBefore(keepFrom) + } + } + if closeErr := writable.Close(); runErr == nil { + runErr = closeErr + } + writerErr <- runErr + close(done) + }() + +reading: + for { + select { + case <-done: + break reading + default: + } + + readOnly, err := OpenReadOnlyChangelogWAL(dir) + if err != nil { + continue // source changed while opening; fail-closed and retry is valid + } + first, err := readOnly.FirstOffset() + require.NoError(t, err) + last, err := readOnly.LastOffset() + require.NoError(t, err) + if first > 0 { + require.GreaterOrEqual(t, last, first) + entry, err := readOnly.ReadAt(last) + require.NoError(t, err) + require.Equal(t, int64(last), entry.Version) + } + require.NoError(t, readOnly.Close()) + } + require.NoError(t, <-writerErr) + + readOnly, err := OpenReadOnlyChangelogWAL(dir) + require.NoError(t, err) + defer func() { require.NoError(t, readOnly.Close()) }() + last, err := readOnly.LastOffset() + require.NoError(t, err) + require.Equal(t, uint64(versions), last) + entry, err := readOnly.ReadAt(last) + require.NoError(t, err) + require.Equal(t, int64(versions), entry.Version) +} + +func lastPlainWALSegment(t *testing.T, dir string) string { + t.Helper() + entries, err := os.ReadDir(dir) + require.NoError(t, err) + var names []string + for _, entry := range entries { + if !entry.IsDir() && len(entry.Name()) == 20 { + names = append(names, entry.Name()) + } + } + require.NotEmpty(t, names) + sort.Strings(names) + return filepath.Join(dir, names[len(names)-1]) +} + +func snapshotWALFiles(t *testing.T, dir string) map[string][]byte { + t.Helper() + files := make(map[string][]byte) + entries, err := os.ReadDir(dir) + require.NoError(t, err) + for _, entry := range entries { + if entry.IsDir() { + continue + } + data, err := os.ReadFile(filepath.Clean(filepath.Join(dir, entry.Name()))) + require.NoError(t, err) + files[entry.Name()] = data + } + return files +} + +func writeReadOnlyTestData(t *testing.T, changelog ChangelogWAL) { + t.Helper() + for i, changeset := range ChangeSets { + require.NoError(t, changelog.Write(proto.ChangelogEntry{ + Version: int64(i + 1), + Changesets: []*proto.NamedChangeSet{{ + Name: "test", + Changeset: changeset, + }}, + })) + } +}