diff --git a/cmd/mithril/node/node.go b/cmd/mithril/node/node.go index c9c964064..2bf22f378 100644 --- a/cmd/mithril/node/node.go +++ b/cmd/mithril/node/node.go @@ -2566,6 +2566,10 @@ postBootstrap: klog.Fatalf("invalid port: %d", rpcPort) } else if rpcPort != 0 { rpcServer = rpcserver.NewRpcServer(accountsDb, uint16(rpcPort), epochScheduleFromState(mithrilState), solana.MustHashFromBase58(networkGenesisHash)) + retentionSlots := epochScheduleFromState(mithrilState).SlotsPerEpoch * 2 + if err := rpcServer.EnableBlockHistory(filepath.Join(accountsPath, "rpc-block-history"), retentionSlots, mithrilState.LastRootedSlot); err != nil { + klog.Fatalf("enable RPC block history: %v", err) + } rpcServer.Start() mlog.Log.Infof("Started RPC server on port %d", rpcPort) } @@ -4602,6 +4606,13 @@ func runReplayWithRecovery( break } } + if history, ok := rpcServer.(interface{ RewindBlockHistory(uint64) error }); ok { + if err := history.RewindBlockHistory(mithrilState.LastRootedSlot); err != nil { + result.Error = fmt.Errorf("fork switch: rewind RPC block history to slot %d: %w", mithrilState.LastRootedSlot, err) + mlog.Log.Errorf("%v; halting", result.Error) + break + } + } if mithrilState.LastRootedContext == nil { mlog.Log.Errorf("fork switch: no rooted checkpoint context to re-replay from; halting") break diff --git a/pkg/blockhistory/store.go b/pkg/blockhistory/store.go new file mode 100644 index 000000000..2d34fb56f --- /dev/null +++ b/pkg/blockhistory/store.go @@ -0,0 +1,353 @@ +package blockhistory + +import ( + "encoding/json" + "errors" + "fmt" + "os" + "path/filepath" + "sort" + "strconv" + "strings" + "sync" + "sync/atomic" + + b "github.com/Overclock-Validator/mithril/pkg/block" + "github.com/gagliardetto/solana-go" + "github.com/gagliardetto/solana-go/rpc" +) + +const recordVersion = 1 + +var ( + ErrNotAvailable = errors.New("block history is not available") + ErrSlotSkipped = errors.New("slot was skipped") +) + +type Record struct { + Slot uint64 `json:"slot"` + ParentSlot uint64 `json:"parentSlot,omitempty"` + BlockHeight uint64 `json:"blockHeight,omitempty"` + Blockhash string `json:"blockhash,omitempty"` + PreviousBlockhash string `json:"previousBlockhash,omitempty"` + BlockTime *int64 `json:"blockTime"` + Rewards []rpc.BlockReward `json:"rewards"` + Skipped bool `json:"skipped,omitempty"` +} + +type batch struct { + Version uint32 `json:"version"` + Through uint64 `json:"through"` + Records []Record `json:"records"` +} + +// Store retains speculative summaries in memory and writes one immutable batch +// immediately before the matching AccountsDB fold commits. Readers are also +// bounded by rooted, so an abandoned fork or interrupted fold stays invisible. +type Store struct { + dir string + retentionSlots uint64 + rooted atomic.Uint64 + mu sync.RWMutex + pending map[uint64]Record + persisted map[uint64]Record + sourceBatch map[uint64]uint64 + batches map[uint64][]uint64 +} + +func Open(dir string, retentionSlots uint64) (*Store, error) { + if retentionSlots == 0 { + return nil, errors.New("block history retention must be greater than zero") + } + if err := os.MkdirAll(dir, 0o755); err != nil { + return nil, fmt.Errorf("create block history directory: %w", err) + } + store := &Store{ + dir: dir, + retentionSlots: retentionSlots, + pending: make(map[uint64]Record), + persisted: make(map[uint64]Record), + sourceBatch: make(map[uint64]uint64), + batches: make(map[uint64][]uint64), + } + entries, err := os.ReadDir(dir) + if err != nil { + return nil, fmt.Errorf("read block history directory: %w", err) + } + for _, entry := range entries { + if entry.IsDir() || !strings.HasPrefix(entry.Name(), "batch-") || !strings.HasSuffix(entry.Name(), ".json") { + continue + } + through, err := strconv.ParseUint(strings.TrimSuffix(strings.TrimPrefix(entry.Name(), "batch-"), ".json"), 10, 64) + if err != nil { + continue + } + loaded, err := readBatch(filepath.Join(dir, entry.Name())) + if err != nil { + return nil, fmt.Errorf("load block history batch %d: %w", through, err) + } + if loaded.Through != through { + return nil, fmt.Errorf("load block history batch %d: record contains through %d", through, loaded.Through) + } + store.install(loaded) + } + return store, nil +} + +func (s *Store) RecordBlock(block *b.Block) error { + if block == nil { + return errors.New("record block history: nil block") + } + rewards := make([]rpc.BlockReward, len(block.Rewards)) + copy(rewards, block.Rewards) + record := Record{ + Slot: block.Slot, + ParentSlot: block.ParentSlot, + BlockHeight: block.BlockHeight, + Blockhash: solana.Hash(block.Blockhash).String(), + PreviousBlockhash: solana.Hash(block.LastBlockhash).String(), + Rewards: rewards, + } + timestamp := block.UnixTimestamp + if timestamp == 0 && block.FooterProducerTimeNanos != 0 { + timestamp = int64(block.FooterProducerTimeNanos / 1_000_000_000) + } + if timestamp != 0 { + record.BlockTime = ×tamp + } + s.mu.Lock() + s.pending[record.Slot] = record + s.mu.Unlock() + return nil +} + +func (s *Store) RecordSkipped(slot uint64) error { + s.mu.Lock() + s.pending[slot] = Record{Slot: slot, Skipped: true} + s.mu.Unlock() + return nil +} + +// Prepare persists every pending outcome through the fold boundary. It must run +// before the corresponding AccountsDB commit; a leftover unselected batch is +// harmless because Get also checks the recovered rooted watermark. +func (s *Store) Prepare(through uint64) error { + s.mu.RLock() + records := make([]Record, 0) + for slot, record := range s.pending { + if slot <= through { + records = append(records, record) + } + } + s.mu.RUnlock() + if len(records) == 0 { + return nil + } + sort.Slice(records, func(i, j int) bool { return records[i].Slot < records[j].Slot }) + payload := batch{Version: recordVersion, Through: through, Records: records} + tmp, err := os.CreateTemp(s.dir, ".block-history-*") + if err != nil { + return fmt.Errorf("create block history temp file: %w", err) + } + tmpPath := tmp.Name() + defer os.Remove(tmpPath) + writeErr := json.NewEncoder(tmp).Encode(payload) + if writeErr == nil { + writeErr = tmp.Sync() + } + closeErr := tmp.Close() + if writeErr != nil { + return fmt.Errorf("write block history through slot %d: %w", through, writeErr) + } + if closeErr != nil { + return fmt.Errorf("close block history through slot %d: %w", through, closeErr) + } + if err := os.Rename(tmpPath, s.batchPath(through)); err != nil { + return fmt.Errorf("publish block history through slot %d: %w", through, err) + } + if err := syncDir(s.dir); err != nil { + return err + } + s.mu.Lock() + s.install(payload) + s.mu.Unlock() + return nil +} + +func (s *Store) install(value batch) { + slots := make([]uint64, 0, len(value.Records)) + for _, record := range value.Records { + s.persisted[record.Slot] = record + s.sourceBatch[record.Slot] = value.Through + slots = append(slots, record.Slot) + } + s.batches[value.Through] = slots +} + +func (s *Store) SetRooted(slot uint64) error { + return s.setRooted(slot, false) +} + +// Rewind adopts an older durable root and drops every speculative record from +// the abandoned replay attempt before that attempt is run again. +func (s *Store) Rewind(slot uint64) error { + return s.setRooted(slot, true) +} + +// DiscardUnrootedFrom drops a discarded fork suffix without changing the +// durable root. Pending records below from remain available for the next fold. +func (s *Store) DiscardUnrootedFrom(from uint64) error { + s.mu.Lock() + rooted := s.rooted.Load() + if from <= rooted { + s.mu.Unlock() + return fmt.Errorf("cannot discard block history at or below rooted slot %d", rooted) + } + for slot := range s.pending { + if slot >= from { + delete(s.pending, slot) + } + } + removed := false + for through, slots := range s.batches { + if through < from || through <= rooted { + continue + } + if err := os.Remove(s.batchPath(through)); err != nil && !os.IsNotExist(err) { + s.mu.Unlock() + return fmt.Errorf("discard unrooted block history batch %d: %w", through, err) + } + removed = true + for _, slot := range slots { + if s.sourceBatch[slot] == through { + delete(s.persisted, slot) + delete(s.sourceBatch, slot) + } + } + delete(s.batches, through) + } + s.mu.Unlock() + if removed { + return syncDir(s.dir) + } + return nil +} + +func (s *Store) setRooted(slot uint64, rewind bool) error { + s.rooted.Store(slot) + s.mu.Lock() + for pendingSlot := range s.pending { + if rewind || pendingSlot <= slot { + delete(s.pending, pendingSlot) + } + } + pruneBefore := uint64(0) + if slot > s.retentionSlots { + pruneBefore = slot - s.retentionSlots + } + removed := false + for through, slots := range s.batches { + orphaned := through > slot + expired := pruneBefore > 0 && through <= pruneBefore + if !orphaned && !expired { + continue + } + if err := os.Remove(s.batchPath(through)); err != nil && !os.IsNotExist(err) { + s.mu.Unlock() + return fmt.Errorf("prune block history batch %d: %w", through, err) + } + removed = true + for _, recordSlot := range slots { + if s.sourceBatch[recordSlot] == through { + delete(s.persisted, recordSlot) + delete(s.sourceBatch, recordSlot) + } + } + delete(s.batches, through) + } + s.mu.Unlock() + if removed { + if err := syncDir(s.dir); err != nil { + return err + } + } + return nil +} + +func (s *Store) Rooted() uint64 { return s.rooted.Load() } + +func (s *Store) Get(slot uint64) (Record, error) { + if slot > s.rooted.Load() { + return Record{}, ErrNotAvailable + } + s.mu.RLock() + record, ok := s.persisted[slot] + sourceBatch := s.sourceBatch[slot] + s.mu.RUnlock() + if !ok || sourceBatch > s.rooted.Load() { + return Record{}, ErrNotAvailable + } + if record.Skipped { + return Record{}, ErrSlotSkipped + } + return record, nil +} + +func (s *Store) MinimumLedgerSlot() (uint64, error) { return s.first(false) } + +func (s *Store) FirstAvailableBlock() (uint64, error) { return s.first(true) } + +func (s *Store) first(requireBlock bool) (uint64, error) { + rooted := s.rooted.Load() + var first uint64 + found := false + s.mu.RLock() + defer s.mu.RUnlock() + for slot, record := range s.persisted { + if slot > rooted || s.sourceBatch[slot] > rooted || (requireBlock && record.Skipped) { + continue + } + if !found || slot < first { + first = slot + found = true + } + } + if !found { + return 0, ErrNotAvailable + } + return first, nil +} + +func (s *Store) batchPath(through uint64) string { + return filepath.Join(s.dir, fmt.Sprintf("batch-%020d.json", through)) +} + +func readBatch(path string) (batch, error) { + data, err := os.ReadFile(path) + if err != nil { + return batch{}, err + } + var value batch + if err := json.Unmarshal(data, &value); err != nil { + return batch{}, err + } + if value.Version != recordVersion { + return batch{}, fmt.Errorf("unsupported block history version %d", value.Version) + } + return value, nil +} + +func syncDir(path string) error { + dir, err := os.Open(path) + if err != nil { + return fmt.Errorf("open block history directory: %w", err) + } + if err := dir.Sync(); err != nil { + _ = dir.Close() + return fmt.Errorf("sync block history directory: %w", err) + } + if err := dir.Close(); err != nil { + return fmt.Errorf("close block history directory: %w", err) + } + return nil +} diff --git a/pkg/blockhistory/store_test.go b/pkg/blockhistory/store_test.go new file mode 100644 index 000000000..a74616674 --- /dev/null +++ b/pkg/blockhistory/store_test.go @@ -0,0 +1,200 @@ +package blockhistory + +import ( + "errors" + "sync" + "testing" + + b "github.com/Overclock-Validator/mithril/pkg/block" + "github.com/gagliardetto/solana-go" + "github.com/stretchr/testify/require" +) + +func TestStoreConcurrentReadAndPrepare(t *testing.T) { + store, err := Open(t.TempDir(), 100) + require.NoError(t, err) + require.NoError(t, store.RecordBlock(&b.Block{Slot: 1})) + require.NoError(t, store.Prepare(10)) + require.NoError(t, store.SetRooted(10)) + + var readers sync.WaitGroup + readers.Add(1) + go func() { + defer readers.Done() + for i := 0; i < 100000; i++ { + _, _ = store.Get(1) + } + }() + for i := 0; i < 20; i++ { + require.NoError(t, store.RecordBlock(&b.Block{Slot: 1})) + require.NoError(t, store.Prepare(10)) + } + readers.Wait() +} + +func TestStoreOnlyPublishesRootedHistoryAndReplacesFork(t *testing.T) { + store, err := Open(t.TempDir(), 100) + require.NoError(t, err) + + firstHash := solana.Hash{1} + secondHash := solana.Hash{2} + require.NoError(t, store.RecordBlock(&b.Block{Slot: 10, Blockhash: firstHash})) + _, err = store.Get(10) + require.ErrorIs(t, err, ErrNotAvailable) + + require.NoError(t, store.RecordBlock(&b.Block{Slot: 10, Blockhash: secondHash})) + require.NoError(t, store.Prepare(10)) + require.NoError(t, store.SetRooted(10)) + record, err := store.Get(10) + require.NoError(t, err) + require.Equal(t, secondHash.String(), record.Blockhash) +} + +func TestStoreHidesPreparedBatchUntilItsFoldCommits(t *testing.T) { + store, err := Open(t.TempDir(), 100) + require.NoError(t, err) + require.NoError(t, store.RecordBlock(&b.Block{Slot: 11})) + require.NoError(t, store.Prepare(12)) + store.rooted.Store(11) // The account fold through slot 12 has not committed. + _, err = store.Get(11) + require.ErrorIs(t, err, ErrNotAvailable) + _, err = store.FirstAvailableBlock() + require.ErrorIs(t, err, ErrNotAvailable) + require.NoError(t, store.SetRooted(12)) + _, err = store.Get(11) + require.NoError(t, err) +} + +func TestStoreUsesAlpenglowFooterBlockTime(t *testing.T) { + store, err := Open(t.TempDir(), 100) + require.NoError(t, err) + require.NoError(t, store.RecordBlock(&b.Block{Slot: 1, FooterProducerTimeNanos: 1_790_110_767_123_456_789})) + require.NoError(t, store.Prepare(1)) + require.NoError(t, store.SetRooted(1)) + record, err := store.Get(1) + require.NoError(t, err) + require.NotNil(t, record.BlockTime) + require.Equal(t, int64(1_790_110_767), *record.BlockTime) +} + +func TestStorePrepareWithoutPendingHistoryIsNoOp(t *testing.T) { + store, err := Open(t.TempDir(), 100) + require.NoError(t, err) + require.NoError(t, store.Prepare(10)) + require.NoError(t, store.SetRooted(10)) +} + +func TestStoreKeepsLedgerAndBlockFloorsDistinct(t *testing.T) { + store, err := Open(t.TempDir(), 100) + require.NoError(t, err) + require.NoError(t, store.RecordSkipped(8)) + require.NoError(t, store.RecordBlock(&b.Block{Slot: 9})) + require.NoError(t, store.Prepare(9)) + require.NoError(t, store.SetRooted(9)) + + minimum, err := store.MinimumLedgerSlot() + require.NoError(t, err) + require.Equal(t, uint64(8), minimum) + firstBlock, err := store.FirstAvailableBlock() + require.NoError(t, err) + require.Equal(t, uint64(9), firstBlock) + _, err = store.Get(8) + require.ErrorIs(t, err, ErrSlotSkipped) +} + +func TestStoreRetentionPrunesOldRecords(t *testing.T) { + dir := t.TempDir() + store, err := Open(dir, 2) + require.NoError(t, err) + for slot := uint64(1); slot <= 4; slot++ { + require.NoError(t, store.RecordBlock(&b.Block{Slot: slot})) + require.NoError(t, store.Prepare(slot)) + require.NoError(t, store.SetRooted(slot)) + } + _, err = store.Get(1) + require.True(t, errors.Is(err, ErrNotAvailable)) + minimum, err := store.MinimumLedgerSlot() + require.NoError(t, err) + require.Equal(t, uint64(3), minimum) + + reopened, err := Open(dir, 2) + require.NoError(t, err) + require.NoError(t, reopened.SetRooted(4)) + first, err := reopened.FirstAvailableBlock() + require.NoError(t, err) + require.Equal(t, uint64(3), first) +} + +func TestStoreDropsUnselectedBatchAfterRestart(t *testing.T) { + dir := t.TempDir() + store, err := Open(dir, 100) + require.NoError(t, err) + require.NoError(t, store.RecordBlock(&b.Block{Slot: 10, Blockhash: solana.Hash{1}})) + require.NoError(t, store.Prepare(10)) // Simulate crash before AccountsDB commits slot 10. + + reopened, err := Open(dir, 100) + require.NoError(t, err) + require.NoError(t, reopened.SetRooted(9)) + require.NoError(t, reopened.RecordBlock(&b.Block{Slot: 10, Blockhash: solana.Hash{2}})) + require.NoError(t, reopened.Prepare(10)) + require.NoError(t, reopened.SetRooted(10)) + record, err := reopened.Get(10) + require.NoError(t, err) + require.Equal(t, solana.Hash{2}.String(), record.Blockhash) +} + +func TestStoreRewindDropsPersistedAndSpeculativeForkHistory(t *testing.T) { + dir := t.TempDir() + store, err := Open(dir, 100) + require.NoError(t, err) + for slot := uint64(10); slot <= 12; slot++ { + require.NoError(t, store.RecordBlock(&b.Block{Slot: slot, Blockhash: solana.Hash{byte(slot)}})) + require.NoError(t, store.Prepare(slot)) + require.NoError(t, store.SetRooted(slot)) + } + require.NoError(t, store.RecordBlock(&b.Block{Slot: 13, Blockhash: solana.Hash{13}})) + + require.NoError(t, store.Rewind(10)) + _, err = store.Get(11) + require.ErrorIs(t, err, ErrNotAvailable) + + // A replacement fold must not inherit the abandoned slot 13 record. + require.NoError(t, store.RecordBlock(&b.Block{Slot: 11, Blockhash: solana.Hash{21}})) + require.NoError(t, store.Prepare(13)) + require.NoError(t, store.SetRooted(13)) + _, err = store.Get(13) + require.ErrorIs(t, err, ErrNotAvailable) + record, err := store.Get(11) + require.NoError(t, err) + require.Equal(t, solana.Hash{21}.String(), record.Blockhash) +} + +func TestStoreDiscardsUnrootedForkSuffixWithoutAdvancingRoot(t *testing.T) { + dir := t.TempDir() + store, err := Open(dir, 100) + require.NoError(t, err) + require.NoError(t, store.RecordBlock(&b.Block{Slot: 10, Blockhash: solana.Hash{10}})) + require.NoError(t, store.Prepare(10)) + require.NoError(t, store.SetRooted(10)) + require.NoError(t, store.RecordBlock(&b.Block{Slot: 11, Blockhash: solana.Hash{11}})) + require.NoError(t, store.RecordBlock(&b.Block{Slot: 12, Blockhash: solana.Hash{12}})) + require.NoError(t, store.Prepare(12)) + require.NoError(t, store.RecordBlock(&b.Block{Slot: 13, Blockhash: solana.Hash{13}})) + + require.NoError(t, store.DiscardUnrootedFrom(12)) + require.Equal(t, uint64(10), store.Rooted()) + require.NoError(t, store.RecordBlock(&b.Block{Slot: 12, Blockhash: solana.Hash{22}})) + require.NoError(t, store.Prepare(13)) + require.NoError(t, store.SetRooted(13)) + + reopened, err := Open(dir, 100) + require.NoError(t, err) + require.NoError(t, reopened.SetRooted(13)) + for slot, hash := range map[uint64]solana.Hash{10: {10}, 11: {11}, 12: {22}} { + record, err := reopened.Get(slot) + require.NoError(t, err) + require.Equal(t, hash.String(), record.Blockhash) + } + _, err = reopened.Get(13) + require.ErrorIs(t, err, ErrNotAvailable) +} diff --git a/pkg/replay/block.go b/pkg/replay/block.go index e571c98d0..1aec50e02 100644 --- a/pkg/replay/block.go +++ b/pkg/replay/block.go @@ -52,6 +52,14 @@ type SlotCtxSetter interface { SetSlotCtx(slotCtx *sealevel.SlotCtx) } +type blockHistoryPublisher interface { + RecordBlockHistory(block *b.Block) error + RecordSkippedBlockHistory(slot uint64) error + PrepareBlockHistory(through uint64) error + SetRootedBlockHistorySlot(slot uint64) error + DiscardUnrootedBlockHistory(from uint64) error +} + // BlockFetchOpts contains options for parallel block fetching type BlockFetchOpts struct { MaxRPS int // Rate limit (requests per second), 0 = use default @@ -1688,6 +1696,7 @@ func ReplayBlocks( onCancelWriteState OnCancelWriteState, // callback to write state immediately on cancellation (can be nil) ) *ReplayResult { result := &ReplayResult{} + historyPublisher, _ := rpcServer.(blockHistoryPublisher) alpenglowMode := consensusOpts != nil && consensusOpts.Alpenglow replayFrontier := uint64(0) if startSlot > 0 { @@ -2006,7 +2015,16 @@ func ReplayBlocks( // on the existing fold worker, after releasing the live cache lock. Capture: transactionStatuses.CaptureSnapshotThrough, Install: func(through uint64, payload []byte) (*state.TransactionStatusCheckpointRef, error) { - return PrepareTransactionStatusCheckpoint(acctsDbPath, through, payload) + ref, err := PrepareTransactionStatusCheckpoint(acctsDbPath, through, payload) + if err != nil { + return nil, err + } + if historyPublisher != nil { + if err := historyPublisher.PrepareBlockHistory(through); err != nil { + return nil, fmt.Errorf("prepare block history: %w", err) + } + } + return ref, nil }, AfterCommit: checkpointAfterCommit, }); hookErr != nil { @@ -2076,6 +2094,11 @@ func ReplayBlocks( mithrilState.LastRootedSlot = promotedThrough mithrilState.LastRootedBankhash = rootedCtx.Bankhash mithrilState.LastRootedContext = rootedCtx + if historyPublisher != nil { + if err := historyPublisher.SetRootedBlockHistorySlot(promotedThrough); err != nil { + mlog.Log.Errorf("block history retention after rooted slot %d: %v", promotedThrough, err) + } + } if rewardsCompletion.retire(&partitionedRewardsInfo, promotedThrough) { rewardsHoldBelowSlot = 0 mlog.Log.Infof("epoch rewards bookkeeping retired through durable slot %d; later fork switches may unwind in memory", promotedThrough) @@ -2472,6 +2495,13 @@ func ReplayBlocks( mlog.Log.Warnf("%v — block source rejected the fork rewind", sw) return false } + if historyPublisher != nil { + if err := historyPublisher.DiscardUnrootedBlockHistory(sw.Slot); err != nil { + result.Error = fmt.Errorf("%w: discard unrooted block history: %v", sw, err) + mlog.Log.Warnf("%v", result.Error) + return false + } + } if !parentSwitchNeedsStateUnwind(sw.Slot, currentExecutedAnchorSlot()) { // Trailing skips advance replay's consumed frontier without creating // account-state layers. Finalized ancestry can later require a block @@ -2897,6 +2927,13 @@ func ReplayBlocks( // Handle skipped slots - log and continue without execution if block.IsSkipped { streamer.discardSlot(block.Slot, "skipped") + if historyPublisher != nil { + if err := historyPublisher.RecordSkippedBlockHistory(block.Slot); err != nil { + result.Error = fmt.Errorf("record skipped block history at slot %d: %w", block.Slot, err) + mlog.Log.Errorf("%v", result.Error) + break + } + } // Zero is the explicit locally consumed outcome for a skip. Parent-ID // gap inference is provisional; recording it lets a later certificate // or discovered ancestry require a source rewind and, when necessary, @@ -3155,6 +3192,13 @@ func ReplayBlocks( global.ClearPendingStakePubkeys() break } + if historyPublisher != nil { + if err := historyPublisher.RecordBlockHistory(block); err != nil { + result.Error = fmt.Errorf("record block history at slot %d: %w", block.Slot, err) + mlog.Log.Errorf("%v", result.Error) + break + } + } // The successful child now owns its derived snapshot. Any later bank uses // lastSlotCtx; the one-shot retained unwind bridge is no longer needed. rewardsCompletion.observeBank(partitionedRewardsInfo, lastSlotCtx.BankSysvars()) diff --git a/pkg/replay/block_execution.go b/pkg/replay/block_execution.go index 5e5e5518e..f40aebee2 100644 --- a/pkg/replay/block_execution.go +++ b/pkg/replay/block_execution.go @@ -561,6 +561,9 @@ func (exec *blockExecution) finalize() (*sealevel.SlotCtx, error) { slotCtx.LamportsBurnt = fees.DistributeTxFeesToSlotLeader(acctsDb, slotCtx, block.Leader, &txFeeAccumulator) slotCtx.RecordModifiedAcct(block.Leader) } + if err := recordBlockFeeReward(block, slotCtx, &txFeeAccumulator); err != nil { + return nil, err + } metrics.GlobalBlockReplay.Reward.AddTimingSince(start) start = time.Now() diff --git a/pkg/replay/block_fee_reward.go b/pkg/replay/block_fee_reward.go new file mode 100644 index 000000000..bf11a9f8d --- /dev/null +++ b/pkg/replay/block_fee_reward.go @@ -0,0 +1,52 @@ +package replay + +import ( + "fmt" + "math" + + b "github.com/Overclock-Validator/mithril/pkg/block" + "github.com/Overclock-Validator/mithril/pkg/fees" + "github.com/Overclock-Validator/mithril/pkg/global" + "github.com/Overclock-Validator/mithril/pkg/sealevel" + "github.com/gagliardetto/solana-go" + "github.com/gagliardetto/solana-go/rpc" +) + +func recordBlockFeeReward(block *b.Block, slotCtx *sealevel.SlotCtx, accumulated *fees.TxFeeInfoAccumulator) error { + if block == nil || slotCtx == nil || accumulated == nil || len(block.Transactions) == 0 || accumulated.TotalFees == 0 { + return nil + } + leader := block.Leader + if !global.ManageLeaderSchedule() && block.BlockReward != nil { + leader = block.BlockReward.Leader + } + if leader == (solana.PublicKey{}) { + return nil + } + if slotCtx.LamportsBurnt > accumulated.TotalFees { + return fmt.Errorf("record fee reward at slot %d: burnt fees %d exceed total fees %d", block.Slot, slotCtx.LamportsBurnt, accumulated.TotalFees) + } + lamports := accumulated.TotalFees - slotCtx.LamportsBurnt + if lamports > math.MaxInt64 { + return fmt.Errorf("record fee reward at slot %d: reward %d overflows int64", block.Slot, lamports) + } + leaderAccount, err := slotCtx.GetAccount(leader) + if err != nil { + return fmt.Errorf("record fee reward at slot %d: read leader %s: %w", block.Slot, leader, err) + } + reward := rpc.BlockReward{ + Pubkey: leader, + Lamports: int64(lamports), + PostBalance: leaderAccount.Lamports, + RewardType: rpc.RewardTypeFee, + } + block.BlockReward = &b.BlockRewardsInfo{Leader: leader, Lamports: lamports, PostBalance: leaderAccount.Lamports} + for i := range block.Rewards { + if block.Rewards[i].RewardType == rpc.RewardTypeFee { + block.Rewards[i] = reward + return nil + } + } + block.Rewards = append(block.Rewards, reward) + return nil +} diff --git a/pkg/replay/block_fee_reward_test.go b/pkg/replay/block_fee_reward_test.go new file mode 100644 index 000000000..59b4eefa7 --- /dev/null +++ b/pkg/replay/block_fee_reward_test.go @@ -0,0 +1,39 @@ +package replay + +import ( + "testing" + + "github.com/Overclock-Validator/mithril/pkg/accounts" + b "github.com/Overclock-Validator/mithril/pkg/block" + "github.com/Overclock-Validator/mithril/pkg/fees" + "github.com/Overclock-Validator/mithril/pkg/global" + "github.com/Overclock-Validator/mithril/pkg/sealevel" + "github.com/gagliardetto/solana-go" + "github.com/gagliardetto/solana-go/rpc" + "github.com/stretchr/testify/require" +) + +func TestRecordBlockFeeRewardUsesExecutedBankResult(t *testing.T) { + global.SetManageLeaderSchedule(true) + defer global.SetManageLeaderSchedule(false) + leader := solana.PublicKey{1} + mem := accounts.NewMemAccounts() + require.NoError(t, mem.SetAccountWithoutLock(leader, &accounts.Account{Key: leader, Lamports: 1_025})) + slotCtx := &sealevel.SlotCtx{Slot: 7, Accounts: mem, LamportsBurnt: 25} + block := &b.Block{Slot: 7, Leader: leader, Transactions: []*solana.Transaction{{}}} + + require.NoError(t, recordBlockFeeReward(block, slotCtx, &fees.TxFeeInfoAccumulator{TotalFees: 50})) + require.Len(t, block.Rewards, 1) + require.Equal(t, rpc.RewardTypeFee, block.Rewards[0].RewardType) + require.Equal(t, int64(25), block.Rewards[0].Lamports) + require.Equal(t, uint64(1_025), block.Rewards[0].PostBalance) + require.Equal(t, uint64(25), block.BlockReward.Lamports) +} + +func TestRecordBlockFeeRewardIgnoresBlockWithoutLeader(t *testing.T) { + require.NoError(t, recordBlockFeeReward( + &b.Block{Slot: 8, Transactions: []*solana.Transaction{{}}}, + &sealevel.SlotCtx{Slot: 8}, + &fees.TxFeeInfoAccumulator{TotalFees: 50}, + )) +} diff --git a/pkg/replay/leader_finalize.go b/pkg/replay/leader_finalize.go index 78c0b90d9..3a62d0b74 100644 --- a/pkg/replay/leader_finalize.go +++ b/pkg/replay/leader_finalize.go @@ -143,6 +143,9 @@ func CommitLeaderSlot(in CommitLeaderInput) (*sealevel.SlotCtx, error) { slotCtx.LamportsBurnt = fees.DistributeTxFeesToSlotLeader(in.AcctsDb, slotCtx, block.Leader, &in.TxFeeAccumulator) slotCtx.RecordModifiedAcct(block.Leader) } + if err := recordBlockFeeReward(block, slotCtx, &in.TxFeeAccumulator); err != nil { + return nil, err + } var rentSysvar *sealevel.SysvarRent if bankSysvars := slotCtx.BankSysvars(); bankSysvars != nil { if bankRent, ok := bankSysvars.Rent(); ok { diff --git a/pkg/replay/streaming_lifecycle_test.go b/pkg/replay/streaming_lifecycle_test.go index 6ee885ef8..d5237dcf1 100644 --- a/pkg/replay/streaming_lifecycle_test.go +++ b/pkg/replay/streaming_lifecycle_test.go @@ -22,6 +22,7 @@ import ( "github.com/Overclock-Validator/mithril/pkg/txverify" bin "github.com/gagliardetto/binary" "github.com/gagliardetto/solana-go" + "github.com/gagliardetto/solana-go/rpc" "github.com/stretchr/testify/require" ) @@ -71,6 +72,7 @@ func (t *lifecycleTail) Add(slot uint64, delta []*accounts.Account, bankhash []b func (t *lifecycleTail) OverCap() bool { return false } type lifecycleEnv struct { + leader solana.PublicKey feats *features.Features durable accounts.MemAccounts acctsDb *accountsdb.AccountsDb @@ -165,6 +167,7 @@ func newLifecycleEnv(t *testing.T) *lifecycleEnv { // executed parent; txs nil is the streaming shell. func (env *lifecycleEnv) block(txs []*solana.Transaction) *b.Block { return &b.Block{ + Leader: env.leader, Slot: lifecycleSlot, VoteTimestamps: make(map[solana.PublicKey]sealevel.BlockTimestamp), Epoch: 0, @@ -183,6 +186,7 @@ func (env *lifecycleEnv) block(txs []*solana.Transaction) *b.Block { } type lifecycleOutcome struct { + rewards []rpc.BlockReward bankhash []byte numSignatures uint64 computeUnits uint64 @@ -245,7 +249,9 @@ func lifecycleWholeBlock(t *testing.T, env *lifecycleEnv, txs []*solana.Transact block.MarkTransactionSignaturesVerified() slotCtx, err := ProcessBlock(env.acctsDb, block, env.epochSchedule, txParallelism, nil, &persistedTracker{}, tail, statuses, false, env.parent) require.NoError(t, err) - return lifecycleOutcomeOf(t, slotCtx, tail) + outcome := lifecycleOutcomeOf(t, slotCtx, tail) + outcome.rewards = block.Rewards + return outcome } // lifecycleStream opens the bank on a transaction-less shell, feeds txs in @@ -312,7 +318,35 @@ func lifecycleStream(t *testing.T, env *lifecycleEnv, txs []*solana.Transaction, require.True(t, exec.closed) require.Equal(t, uint64(1), metrics.GlobalBlockReplay.StreamingExecution.Opened) require.Equal(t, uint64(len(txs)), metrics.GlobalBlockReplay.StreamingExecution.Transactions) - return lifecycleOutcomeOf(t, slotCtx, tail), s + outcome := lifecycleOutcomeOf(t, slotCtx, tail) + outcome.rewards = block.Rewards + return outcome, s +} + +func TestStreamingLifecycleRecordsExecutedFeeRewards(t *testing.T) { + previousStreaming, previousSchedule := StreamingExecutionCfg, global.ManageLeaderSchedule() + StreamingExecutionCfg = StreamingExecutionConfig{Enabled: true} + global.SetManageLeaderSchedule(true) + t.Cleanup(func() { + StreamingExecutionCfg = previousStreaming + global.SetManageLeaderSchedule(previousSchedule) + }) + txs := transferTransactions(t, 4, 1_200) + leader := txfixture.DestPubkey() + wholeEnv := newLifecycleEnv(t) + wholeEnv.leader = leader + whole := lifecycleWholeBlock(t, wholeEnv, txs, 2) + require.Len(t, whole.rewards, 1) + require.Equal(t, rpc.RewardTypeFee, whole.rewards[0].RewardType) + require.Equal(t, leader, whole.rewards[0].Pubkey) + require.Equal(t, int64(10_000), whole.rewards[0].Lamports) + require.Equal(t, whole.delta[leader].Lamports, whole.rewards[0].PostBalance) + + streamedEnv := newLifecycleEnv(t) + streamedEnv.leader = leader + streamed, _ := lifecycleStream(t, streamedEnv, txs, []int{2}, 1, 2) + requireSameLifecycleOutcome(t, whole, streamed) + require.Equal(t, whole.rewards, streamed.rewards) } func TestStreamingLifecycleMatchesWholeBlock(t *testing.T) { diff --git a/pkg/rpcserver/errors.go b/pkg/rpcserver/errors.go index ee2a1654f..2f328d325 100644 --- a/pkg/rpcserver/errors.go +++ b/pkg/rpcserver/errors.go @@ -13,8 +13,12 @@ const ( rpcCodeInvalidParams jsonrpc.ErrorCode = -32602 // -32002 matches Agave's SendTransactionPreflightFailure. rpcCodeSendTransactionPreflightFailure jsonrpc.ErrorCode = -32002 + // -32004 matches Agave's BlockNotAvailable error. + rpcCodeBlockNotAvailable jsonrpc.ErrorCode = -32004 // -32016 is Agave's reserved code for MinContextSlotNotReached. rpcCodeMinContextSlotNotReached jsonrpc.ErrorCode = -32016 + // -32007 matches Agave's SlotSkipped error. + rpcCodeSlotSkipped jsonrpc.ErrorCode = -32007 ) type MinContextSlotNotReachedError struct { @@ -59,6 +63,44 @@ type InvalidParamsError struct { Message string } +type SlotSkippedError struct { + Slot uint64 +} + +type BlockNotAvailableError struct { + Slot uint64 +} + +func (e *BlockNotAvailableError) Error() string { + return fmt.Sprintf("Block not available for slot %d", e.Slot) +} + +func (e *BlockNotAvailableError) ToJSONRPCError() (jsonrpc.JSONRPCError, error) { + return jsonrpc.JSONRPCError{Code: rpcCodeBlockNotAvailable, Message: e.Error()}, nil +} + +func (e *BlockNotAvailableError) FromJSONRPCError(rpcErr jsonrpc.JSONRPCError) error { + if rpcErr.Code != rpcCodeBlockNotAvailable { + return fmt.Errorf("unexpected code %d for BlockNotAvailableError", rpcErr.Code) + } + return nil +} + +func (e *SlotSkippedError) Error() string { + return fmt.Sprintf("Slot %d was skipped", e.Slot) +} + +func (e *SlotSkippedError) ToJSONRPCError() (jsonrpc.JSONRPCError, error) { + return jsonrpc.JSONRPCError{Code: rpcCodeSlotSkipped, Message: e.Error()}, nil +} + +func (e *SlotSkippedError) FromJSONRPCError(rpcErr jsonrpc.JSONRPCError) error { + if rpcErr.Code != rpcCodeSlotSkipped { + return fmt.Errorf("unexpected code %d for SlotSkippedError", rpcErr.Code) + } + return nil +} + func (e *InvalidParamsError) Error() string { return e.Message } func (e *InvalidParamsError) ToJSONRPCError() (jsonrpc.JSONRPCError, error) { @@ -114,6 +156,8 @@ func rpcErrorRegistry() jsonrpc.Errors { errs := jsonrpc.NewErrors() errs.Register(rpcCodeInvalidParams, new(*InvalidParamsError)) errs.Register(rpcCodeSendTransactionPreflightFailure, new(*SendTransactionPreflightFailureError)) + errs.Register(rpcCodeBlockNotAvailable, new(*BlockNotAvailableError)) errs.Register(rpcCodeMinContextSlotNotReached, new(*MinContextSlotNotReachedError)) + errs.Register(rpcCodeSlotSkipped, new(*SlotSkippedError)) return errs } diff --git a/pkg/rpcserver/get_block_history.go b/pkg/rpcserver/get_block_history.go new file mode 100644 index 000000000..a53264fae --- /dev/null +++ b/pkg/rpcserver/get_block_history.go @@ -0,0 +1,150 @@ +package rpcserver + +import ( + "context" + "errors" + "fmt" + "math" + + "github.com/Overclock-Validator/mithril/pkg/blockhistory" + "github.com/filecoin-project/go-jsonrpc" + "github.com/gagliardetto/solana-go/rpc" +) + +type GetBlockResp struct { + BlockHeight uint64 `json:"blockHeight"` + BlockTime *int64 `json:"blockTime"` + Blockhash string `json:"blockhash"` + PreviousBlockhash string `json:"previousBlockhash"` + ParentSlot uint64 `json:"parentSlot"` + Rewards []blockRewardResponse `json:"rewards"` +} + +type blockRewardResponse struct { + rpc.BlockReward + Commission *uint8 `json:"commission"` +} + +func (rpcServer *RpcServer) MinimumLedgerSlot(ctx context.Context, p jsonrpc.RawParams) (uint64, error) { + if err := requireNoParams(p, "minimumLedgerSlot"); err != nil { + return 0, err + } + if rpcServer.blockHistory == nil { + return 0, errors.New("node has no retained block history") + } + return rpcServer.blockHistory.MinimumLedgerSlot() +} + +func (rpcServer *RpcServer) GetFirstAvailableBlock(ctx context.Context, p jsonrpc.RawParams) (uint64, error) { + if err := requireNoParams(p, "getFirstAvailableBlock"); err != nil { + return 0, err + } + if rpcServer.blockHistory == nil { + return 0, errors.New("node has no retained block history") + } + return rpcServer.blockHistory.FirstAvailableBlock() +} + +func (rpcServer *RpcServer) GetBlock(ctx context.Context, p jsonrpc.RawParams) (GetBlockResp, error) { + params, err := jsonrpc.DecodeParams[[]interface{}](p) + if err != nil { + return GetBlockResp{}, &InvalidParamsError{Message: fmt.Sprintf("decoding params: %v", err)} + } + if len(params) < 1 || len(params) > 2 { + return GetBlockResp{}, &InvalidParamsError{Message: "getBlock requires a slot and optional config"} + } + slot, err := blockRPCUint64(params[0], "slot") + if err != nil { + return GetBlockResp{}, err + } + details, err := validateBlockConfig(params) + if err != nil { + return GetBlockResp{}, err + } + if rpcServer.blockHistory == nil { + return GetBlockResp{}, errors.New("node has no retained block history") + } + if details == "full" { + return GetBlockResp{}, &InvalidParamsError{Message: "getBlock transactionDetails: full is not supported because transactions are not retained"} + } + record, err := rpcServer.blockHistory.Get(slot) + if errors.Is(err, blockhistory.ErrSlotSkipped) { + return GetBlockResp{}, &SlotSkippedError{Slot: slot} + } + if errors.Is(err, blockhistory.ErrNotAvailable) { + return GetBlockResp{}, &BlockNotAvailableError{Slot: slot} + } + if err != nil { + return GetBlockResp{}, err + } + response := GetBlockResp{ + BlockHeight: record.BlockHeight, + BlockTime: record.BlockTime, + Blockhash: record.Blockhash, + PreviousBlockhash: record.PreviousBlockhash, + ParentSlot: record.ParentSlot, + Rewards: make([]blockRewardResponse, len(record.Rewards)), + } + for i, reward := range record.Rewards { + response.Rewards[i] = blockRewardResponse{BlockReward: reward, Commission: reward.Commission} + } + return response, nil +} + +func blockRPCUint64(raw interface{}, name string) (uint64, error) { + value, ok := raw.(float64) + if !ok || value < 0 || value >= math.Exp2(64) || math.Trunc(value) != value { + return 0, &InvalidParamsError{Message: fmt.Sprintf("invalid %s", name)} + } + return uint64(value), nil +} + +func requireNoParams(p jsonrpc.RawParams, method string) error { + if len(p) == 0 { + return nil + } + params, err := jsonrpc.DecodeParams[[]interface{}](p) + if err != nil { + return &InvalidParamsError{Message: fmt.Sprintf("decoding params: %v", err)} + } + if len(params) != 0 { + return &InvalidParamsError{Message: method + " does not accept parameters"} + } + return nil +} + +func validateBlockConfig(params []interface{}) (string, error) { + if len(params) == 1 || params[1] == nil { + return "", &InvalidParamsError{Message: "getBlock block history requires transactionDetails: none or full"} + } + config, ok := params[1].(map[string]interface{}) + if !ok { + return "", &InvalidParamsError{Message: "invalid getBlock config"} + } + if raw := config["maxSupportedTransactionVersion"]; raw != nil { + version, err := blockRPCUint64(raw, "maxSupportedTransactionVersion") + if err != nil || version > math.MaxUint8 { + return "", &InvalidParamsError{Message: "invalid maxSupportedTransactionVersion"} + } + } + if commitment, exists := config["commitment"]; exists { + value, ok := commitment.(string) + if !ok || (value != "confirmed" && value != "finalized") { + return "", &InvalidParamsError{Message: "getBlock requires confirmed or finalized commitment"} + } + } + if encoding, exists := config["encoding"]; exists && encoding != "json" { + return "", &InvalidParamsError{Message: "getBlock block history supports json encoding"} + } + details, exists := config["transactionDetails"] + if !exists || (details != "none" && details != "full") { + return "", &InvalidParamsError{Message: "getBlock block history supports transactionDetails: none or full"} + } + if rewards, exists := config["rewards"]; exists { + value, ok := rewards.(bool) + if !ok || !value { + return "", &InvalidParamsError{Message: "getBlock block history requires rewards: true"} + } + } + return details.(string), nil +} diff --git a/pkg/rpcserver/get_block_history_test.go b/pkg/rpcserver/get_block_history_test.go new file mode 100644 index 000000000..ba245363b --- /dev/null +++ b/pkg/rpcserver/get_block_history_test.go @@ -0,0 +1,82 @@ +package rpcserver + +import ( + "encoding/json" + "testing" + + b "github.com/Overclock-Validator/mithril/pkg/block" + "github.com/Overclock-Validator/mithril/pkg/blockhistory" + "github.com/filecoin-project/go-jsonrpc" + "github.com/gagliardetto/solana-go" + "github.com/gagliardetto/solana-go/rpc" + "github.com/stretchr/testify/require" +) + +func TestBlockHistoryRPCServesExporterContract(t *testing.T) { + store, err := blockhistory.Open(t.TempDir(), 100) + require.NoError(t, err) + leader := solana.PublicKey{1} + require.NoError(t, store.RecordSkipped(40)) + require.NoError(t, store.RecordBlock(&b.Block{ + Slot: 41, + ParentSlot: 39, + BlockHeight: 30, + Blockhash: solana.Hash{2}, + LastBlockhash: solana.Hash{3}, + Rewards: []rpc.BlockReward{{ + Pubkey: leader, Lamports: 25, PostBalance: 1_025, RewardType: rpc.RewardTypeFee, + }}, + })) + require.NoError(t, store.Prepare(41)) + require.NoError(t, store.SetRooted(41)) + server := &RpcServer{blockHistory: store} + + minimum, err := server.MinimumLedgerSlot(t.Context(), jsonrpc.RawParams(`[]`)) + require.NoError(t, err) + require.Equal(t, uint64(40), minimum) + first, err := server.GetFirstAvailableBlock(t.Context(), jsonrpc.RawParams(`[]`)) + require.NoError(t, err) + require.Equal(t, uint64(41), first) + for _, params := range []jsonrpc.RawParams{nil, jsonrpc.RawParams(`null`)} { + minimum, err = server.MinimumLedgerSlot(t.Context(), params) + require.NoError(t, err) + require.Equal(t, uint64(40), minimum) + first, err = server.GetFirstAvailableBlock(t.Context(), params) + require.NoError(t, err) + require.Equal(t, uint64(41), first) + } + + response, err := server.GetBlock(t.Context(), jsonrpc.RawParams(`[41,{"commitment":"confirmed","encoding":"json","transactionDetails":"none","rewards":true,"maxSupportedTransactionVersion":0}]`)) + require.NoError(t, err) + require.Equal(t, uint64(30), response.BlockHeight) + require.Equal(t, int64(25), response.Rewards[0].Lamports) + require.Equal(t, leader, response.Rewards[0].Pubkey) + encoded, err := json.Marshal(response) + require.NoError(t, err) + require.Contains(t, string(encoded), `"commission":null`) + + _, err = server.GetBlock(t.Context(), jsonrpc.RawParams(`[40,{"transactionDetails":"none"}]`)) + var skipped *SlotSkippedError + require.ErrorAs(t, err, &skipped) + require.Equal(t, uint64(40), skipped.Slot) + + _, err = server.GetBlock(t.Context(), jsonrpc.RawParams(`[42,{"transactionDetails":"none"}]`)) + var unavailable *BlockNotAvailableError + require.ErrorAs(t, err, &unavailable) + require.Equal(t, uint64(42), unavailable.Slot) +} + +func TestGetBlockRejectsUnretainedTransactionDetails(t *testing.T) { + store, err := blockhistory.Open(t.TempDir(), 100) + require.NoError(t, err) + require.NoError(t, store.RecordBlock(&b.Block{Slot: 1})) + require.NoError(t, store.Prepare(1)) + require.NoError(t, store.SetRooted(1)) + server := &RpcServer{blockHistory: store} + _, err = server.GetBlock(t.Context(), jsonrpc.RawParams(`[1,{"transactionDetails":"full"}]`)) + var invalid *InvalidParamsError + require.ErrorAs(t, err, &invalid) + require.Contains(t, invalid.Error(), "transactions are not retained") + _, err = server.GetBlock(t.Context(), jsonrpc.RawParams(`[2,{"transactionDetails":"full"}]`)) + require.ErrorAs(t, err, &invalid) +} diff --git a/pkg/rpcserver/rpcserver.go b/pkg/rpcserver/rpcserver.go index 8a3dd7afe..ebc1a9267 100644 --- a/pkg/rpcserver/rpcserver.go +++ b/pkg/rpcserver/rpcserver.go @@ -14,6 +14,8 @@ import ( "time" "github.com/Overclock-Validator/mithril/pkg/accountsdb" + b "github.com/Overclock-Validator/mithril/pkg/block" + "github.com/Overclock-Validator/mithril/pkg/blockhistory" "github.com/Overclock-Validator/mithril/pkg/mlog" "github.com/Overclock-Validator/mithril/pkg/sealevel" "github.com/filecoin-project/go-jsonrpc" @@ -31,6 +33,7 @@ type RpcServer struct { slotCtx *sealevel.SlotCtx slotCtxMu sync.RWMutex genesisHash string + blockHistory *blockhistory.Store leaderTPUCacheMu sync.RWMutex leaderTPUByIdentity map[solana.PublicKey]tpuEndpoint @@ -48,14 +51,71 @@ type RpcServer struct { const maxQuietMethodProbeBody = 64 << 10 var supportedRPCMethods = map[string]struct{}{ - "getAccountInfo": {}, - "getBankHash": {}, - "getBlockHeight": {}, - "getEpochInfo": {}, - "getGenesisHash": {}, - "getLatestBlockhash": {}, - "sendTransaction": {}, - "simulateTransaction": {}, + "getAccountInfo": {}, + "getBankHash": {}, + "getBlock": {}, + "getBlockHeight": {}, + "getEpochInfo": {}, + "getFirstAvailableBlock": {}, + "getGenesisHash": {}, + "getLatestBlockhash": {}, + "minimumLedgerSlot": {}, + "sendTransaction": {}, + "simulateTransaction": {}, +} + +func (rpcServer *RpcServer) EnableBlockHistory(dir string, retentionSlots, rootedSlot uint64) error { + store, err := blockhistory.Open(dir, retentionSlots) + if err != nil { + return err + } + if err := store.SetRooted(rootedSlot); err != nil { + return err + } + rpcServer.blockHistory = store + return nil +} + +func (rpcServer *RpcServer) RecordBlockHistory(block *b.Block) error { + if rpcServer == nil || rpcServer.blockHistory == nil { + return nil + } + return rpcServer.blockHistory.RecordBlock(block) +} + +func (rpcServer *RpcServer) RecordSkippedBlockHistory(slot uint64) error { + if rpcServer == nil || rpcServer.blockHistory == nil { + return nil + } + return rpcServer.blockHistory.RecordSkipped(slot) +} + +func (rpcServer *RpcServer) PrepareBlockHistory(through uint64) error { + if rpcServer == nil || rpcServer.blockHistory == nil { + return nil + } + return rpcServer.blockHistory.Prepare(through) +} + +func (rpcServer *RpcServer) SetRootedBlockHistorySlot(slot uint64) error { + if rpcServer == nil || rpcServer.blockHistory == nil { + return nil + } + return rpcServer.blockHistory.SetRooted(slot) +} + +func (rpcServer *RpcServer) RewindBlockHistory(slot uint64) error { + if rpcServer == nil || rpcServer.blockHistory == nil { + return nil + } + return rpcServer.blockHistory.Rewind(slot) +} + +func (rpcServer *RpcServer) DiscardUnrootedBlockHistory(from uint64) error { + if rpcServer == nil || rpcServer.blockHistory == nil { + return nil + } + return rpcServer.blockHistory.DiscardUnrootedFrom(from) } func NewRpcServer(acctsDb *accountsdb.AccountsDb, port uint16, epochSchedule *sealevel.SysvarEpochSchedule, genesisHash solana.Hash) *RpcServer { diff --git a/pkg/snapshot/build_db.go b/pkg/snapshot/build_db.go index a7e2ed934..b1f285ffe 100644 --- a/pkg/snapshot/build_db.go +++ b/pkg/snapshot/build_db.go @@ -51,6 +51,7 @@ func CleanAccountsDbDir(accountsDbDir string) { // so retaining them would waste space and make stale diagnostics look // actionable. "transaction-status-checkpoints", + "rpc-block-history", "mithril_db", "mithril_db_log_shards", "bankhash_db",