From 318ddb22f2d34cd551c3dbb2e99e2b615b3bbb7d Mon Sep 17 00:00:00 2001 From: 7layermagik <7layermagik@users.noreply.github.com> Date: Thu, 24 Sep 2026 00:16:49 -0500 Subject: [PATCH] alpenglow: separate opt-in signing reservations and ordered history persistence Preserve reviewed persistence changes from fde37516, 75f4b7e7, c2cbc4a8 and e10b9537; transport lives in the certificate/delivery prerequisite. --- .github/workflows/go_build.yml | 28 ++ cmd/mithril/node/node.go | 103 +++--- cmd/mithril/node/vote_startup.go | 40 +++ cmd/mithril/node/vote_startup_test.go | 73 ++++ config.example.toml | 10 + docs/reserved-vote-history.md | 162 +++++++++ docs/vote-delivery-persistence-evidence.md | 14 + pkg/alpenglow/vote_history.go | 113 ++++-- pkg/alpenglow/vote_history_snapshot_test.go | 49 +++ pkg/alpenglow/vote_reservation.go | 104 ++++++ pkg/alpenglow/vote_reservation_test.go | 71 ++++ pkg/blockprod/leader.go | 6 + .../leader_signing_reservation_test.go | 17 + pkg/config/config.go | 2 + pkg/consensus/engine.go | 40 ++- pkg/consensus/vote_history_writer.go | 125 +++++++ pkg/consensus/vote_history_writer_test.go | 201 +++++++++++ pkg/consensus/vote_reservation.go | 296 ++++++++++++++++ pkg/consensus/vote_reservation_test.go | 330 ++++++++++++++++++ pkg/consensus/voter.go | 279 +++++++++++---- pkg/consensus/voter_finality_ordering_test.go | 192 ++++++++++ pkg/consensus/voter_wait_slot_test.go | 91 +++++ 22 files changed, 2209 insertions(+), 137 deletions(-) create mode 100644 cmd/mithril/node/vote_startup.go create mode 100644 cmd/mithril/node/vote_startup_test.go create mode 100644 docs/reserved-vote-history.md create mode 100644 docs/vote-delivery-persistence-evidence.md create mode 100644 pkg/alpenglow/vote_history_snapshot_test.go create mode 100644 pkg/alpenglow/vote_reservation.go create mode 100644 pkg/alpenglow/vote_reservation_test.go create mode 100644 pkg/blockprod/leader_signing_reservation_test.go create mode 100644 pkg/consensus/vote_history_writer.go create mode 100644 pkg/consensus/vote_history_writer_test.go create mode 100644 pkg/consensus/vote_reservation.go create mode 100644 pkg/consensus/vote_reservation_test.go create mode 100644 pkg/consensus/voter_finality_ordering_test.go create mode 100644 pkg/consensus/voter_wait_slot_test.go diff --git a/.github/workflows/go_build.yml b/.github/workflows/go_build.yml index 50383467f..4fa7306c9 100644 --- a/.github/workflows/go_build.yml +++ b/.github/workflows/go_build.yml @@ -17,3 +17,31 @@ jobs: - name: Build run: go build -v ./cmd/mithril + + regression-tests: + runs-on: ubuntu-latest + timeout-minutes: 20 + permissions: + contents: read + env: + GOMAXPROCS: "2" + steps: + - uses: actions/checkout@v3 + + - name: Setup Go + uses: actions/setup-go@v4 + with: + go-version: 1.26.4 + + - name: Voting, checkpoint, streaming and scheduler race regressions + # Run the complete affected package suites, including subprocess crash + # recovery and cancellation tests. The independent sealevel suite has + # known base-branch failures documented in the validation report. + run: >- + go test -race -p 2 -count=1 + ./pkg/alpenglow ./pkg/consensus ./pkg/replay + ./pkg/turbine ./pkg/sigverify ./pkg/blockprod/... + ./cmd/mithril/node ./cmd/mithril/configcmd + + - name: Vote-program deque ownership race regression + run: go test -race -count=1 ./pkg/sealevel -run '^TestProcessNewVoteStateOwnsRetainedDeque$' diff --git a/cmd/mithril/node/node.go b/cmd/mithril/node/node.go index c452d494e..1f6a21d8e 100644 --- a/cmd/mithril/node/node.go +++ b/cmd/mithril/node/node.go @@ -103,6 +103,9 @@ var ( validatorTPUQUICBind string validatorAdvertisedIP string validatorSigverifyWorkers int + validatorWaitToVoteSlot uint64 + validatorReservedHistory bool + validatorInitializeReservation bool // Mode thresholds blockNearTipThreshold int // Enter near-tip when gap <= this @@ -547,6 +550,9 @@ func init() { Run.Flags().StringVar(&validatorTPUQUICBind, "tpu-quic-bind-addr", "", "Validator TPU QUIC listen address (default 0.0.0.0:8004)") Run.Flags().StringVar(&validatorAdvertisedIP, "validator-advertised-ip", "", "Public IP advertised for validator TPU QUIC") Run.Flags().IntVar(&validatorSigverifyWorkers, "tpu-sigverify-workers", 0, "TPU signature verification workers (0 = GOMAXPROCS)") + Run.Flags().BoolVar(&validatorReservedHistory, "reserved-vote-history", false, "Use durable signing reservations with unsynchronized per-vote history writes") + Run.Flags().BoolVar(&validatorInitializeReservation, "initialize-vote-reservation", false, "Enroll complete synchronous vote history in reserved mode (one-time migration)") + Run.Flags().Uint64Var(&validatorWaitToVoteSlot, "wait-to-vote-slot", 0, "Do not cast new votes below this slot; the automatic startup cutoff still applies (0 = automatic only)") // [tuning] section flags Run.Flags().Uint64Var(¶mArenaSizeMB, "param-arena-size-mb", 512, "Size in MB for serialized parameter arena (0 to disable)") @@ -639,6 +645,11 @@ func initConfigAndBindFlags(cmd *cobra.Command) error { if err := config.InitConfig(); err != nil { return err } + if slot, err := configuredWaitToVoteSlot(cmd); err != nil { + return err + } else { + validatorWaitToVoteSlot = slot + } // Check if a CLI flag was explicitly set by the user flagChanged := func(name string) bool { @@ -2642,52 +2653,9 @@ postBootstrap: } global.SeedWallClockSlot(wallClockSeed) startupWallSlot := global.WallClockSlot() - waitToVoteSlot := startupWallSlot - startupWallSlot%alpenglow.LeaderWindowSlots - if waitToVoteSlot <= math.MaxUint64-2*alpenglow.LeaderWindowSlots { - waitToVoteSlot += 2 * alpenglow.LeaderWindowSlots - } else { - waitToVoteSlot = math.MaxUint64 - } - mlog.Log.Infof("ALPENGLOW voting startup watermark: wall_clock=%d wait_to_vote=%d", startupWallSlot, waitToVoteSlot) + waitToVoteSlot := effectiveWaitToVoteSlot(startupWallSlot, validatorWaitToVoteSlot) + mlog.Log.Infof("ALPENGLOW voting startup watermark: wall_clock=%d configured_wait_to_vote=%d wait_to_vote=%d", startupWallSlot, validatorWaitToVoteSlot, waitToVoteSlot) - identityPubkey := solana.PrivateKey(validatorIdentity).PublicKey() - if err := consensusEngine.EnableVoting(consensusengine.VotingConfig{ - Identity: validatorIdentity, - AuthorizedVoter: validatorAuthorizedVoter, - VoteAccount: validatorVoteAccount, - HistoryDir: blockstorePath, - EpochForSlot: epochSchedule.GetEpoch, - SlotDuration: blockprod.AlpenglowSlotDuration, - WaitToVoteSlot: waitToVoteSlot, - ReadyToVote: func(slot uint64) bool { - wallSlot := global.WallClockSlot() - if liveSlot, ok := consensusEngine.AlpenglowLiveSlot(); ok { - wallSlot = liveSlot - } - return slot >= wallSlot || wallSlot-slot <= alpenglow.LeaderWindowSlots - }, - Peers: func(validators []alpenglow.ValidatorStake) []alpenglow.VotorPeer { - peers := make([]alpenglow.VotorPeer, 0, len(validators)) - seen := make(map[solana.PublicKey]struct{}, len(validators)) - for _, validator := range validators { - if validator.Stake == 0 || validator.NodePubkey == identityPubkey { - continue - } - addr, ok := sharedGossip.LookupAlpenglow(validator.NodePubkey) - if !ok { - continue - } - if _, duplicate := seen[validator.NodePubkey]; duplicate { - continue - } - seen[validator.NodePubkey] = struct{}{} - peers = append(peers, alpenglow.VotorPeer{Identity: validator.NodePubkey, Addr: addr}) - } - return peers - }, - }); err != nil { - klog.Fatalf("enable Alpenglow voting: %v", err) - } broadcaster, err := turbine.NewTurbineBroadcaster(turbine.TurbineBroadcasterConfig{ Self: solana.PrivateKey(validatorIdentity).PublicKey(), Peers: sharedGossip, @@ -2744,6 +2712,50 @@ postBootstrap: mlog.Log.Warnf("validator gossip TPU advertisement: %v", err) } + // Bind and validate local transports before consuming the durable clean + // voting marker. Startup configuration failures must not force recovery. + identityPubkey := solana.PrivateKey(validatorIdentity).PublicKey() + if err := consensusEngine.EnableVoting(consensusengine.VotingConfig{ + Identity: validatorIdentity, + AuthorizedVoter: validatorAuthorizedVoter, + VoteAccount: validatorVoteAccount, + HistoryDir: blockstorePath, + ReservedHistory: validatorReservedHistory, + InitializeVoteReservation: validatorInitializeReservation, + Genesis: solana.MustHashFromBase58(networkGenesisHash), + EpochForSlot: epochSchedule.GetEpoch, + SlotDuration: blockprod.AlpenglowSlotDuration, + WaitToVoteSlot: waitToVoteSlot, + ReadyToVote: func(slot uint64) bool { + wallSlot := global.WallClockSlot() + if liveSlot, ok := consensusEngine.AlpenglowLiveSlot(); ok { + wallSlot = liveSlot + } + return slot >= wallSlot || wallSlot-slot <= alpenglow.LeaderWindowSlots + }, + Peers: func(validators []alpenglow.ValidatorStake) []alpenglow.VotorPeer { + peers := make([]alpenglow.VotorPeer, 0, len(validators)) + seen := make(map[solana.PublicKey]struct{}, len(validators)) + for _, validator := range validators { + if validator.Stake == 0 || validator.NodePubkey == identityPubkey { + continue + } + addr, ok := sharedGossip.LookupAlpenglow(validator.NodePubkey) + if !ok { + continue + } + if _, duplicate := seen[validator.NodePubkey]; duplicate { + continue + } + seen[validator.NodePubkey] = struct{}{} + peers = append(peers, alpenglow.VotorPeer{Identity: validator.NodePubkey, Addr: addr}) + } + return peers + }, + }); err != nil { + klog.Fatalf("enable Alpenglow voting: %v", err) + } + rewardBuilder := rewardcerts.NewBuilder(rewardcerts.BuilderConfig{ RootSlot: global.Slot, BeforeBuild: consensusEngine.FlushAlpenglowRewardVotes, @@ -2796,6 +2808,7 @@ postBootstrap: } }, ProductionParent: consensusEngine.AlpenglowBlockProductionParent, + CanSignSlot: consensusEngine.AlpenglowCanSignLeaderSlot, CurrentSlot: func() uint64 { if slot, ok := consensusEngine.AlpenglowLiveSlot(); ok { return slot diff --git a/cmd/mithril/node/vote_startup.go b/cmd/mithril/node/vote_startup.go new file mode 100644 index 000000000..547e55194 --- /dev/null +++ b/cmd/mithril/node/vote_startup.go @@ -0,0 +1,40 @@ +package node + +import ( + "fmt" + "math" + "strconv" + + "github.com/Overclock-Validator/mithril/pkg/alpenglow" + "github.com/Overclock-Validator/mithril/pkg/config" + "github.com/spf13/cobra" +) + +func configuredWaitToVoteSlot(cmd *cobra.Command) (uint64, error) { + if flag := cmd.Flags().Lookup("wait-to-vote-slot"); flag != nil && flag.Changed { + return cmd.Flags().GetUint64("wait-to-vote-slot") + } + const key = "validator.wait_to_vote_slot" + if !config.IsSet(key) { + return 0, nil + } + // Unlike GetUint64, parsing explicitly must not turn an invalid operator + // cutoff into zero and silently remove the requested voting restriction. + slot, err := strconv.ParseUint(config.GetString(key), 10, 64) + if err != nil { + return 0, fmt.Errorf("%s must be an unsigned 64-bit slot: %w", key, err) + } + return slot, nil +} + +// The operator cutoff can postpone voting but cannot weaken the existing +// startup guard. Equality permits voting, subject to all other Votor checks. +func effectiveWaitToVoteSlot(startupWallSlot, configured uint64) uint64 { + automatic := startupWallSlot - startupWallSlot%alpenglow.LeaderWindowSlots + if automatic <= math.MaxUint64-2*alpenglow.LeaderWindowSlots { + automatic += 2 * alpenglow.LeaderWindowSlots + } else { + automatic = math.MaxUint64 + } + return max(automatic, configured) +} diff --git a/cmd/mithril/node/vote_startup_test.go b/cmd/mithril/node/vote_startup_test.go new file mode 100644 index 000000000..d86a02528 --- /dev/null +++ b/cmd/mithril/node/vote_startup_test.go @@ -0,0 +1,73 @@ +package node + +import ( + "math" + "strings" + "testing" + + "github.com/Overclock-Validator/mithril/pkg/config" + "github.com/spf13/cobra" + "github.com/spf13/viper" + "github.com/stretchr/testify/require" +) + +func TestConfiguredWaitToVoteSlot(t *testing.T) { + for _, tc := range []struct { + name, toml, cli string + want uint64 + invalid bool + }{ + {name: "default"}, + {name: "toml", toml: "wait_to_vote_slot = 1234", want: 1234}, + {name: "cli wins", toml: "wait_to_vote_slot = 1234", cli: "5678", want: 5678}, + {name: "explicit zero wins", toml: "wait_to_vote_slot = 1234", cli: "0"}, + {name: "maximum CLI", cli: "18446744073709551615", want: math.MaxUint64}, + {name: "negative TOML", toml: "wait_to_vote_slot = -1", invalid: true}, + {name: "fractional TOML", toml: "wait_to_vote_slot = 1.5", invalid: true}, + {name: "malformed TOML value", toml: `wait_to_vote_slot = "oops"`, invalid: true}, + {name: "empty TOML value", toml: `wait_to_vote_slot = ""`, invalid: true}, + {name: "overflow TOML value", toml: `wait_to_vote_slot = "18446744073709551616"`, invalid: true}, + {name: "negative CLI", cli: "-1", invalid: true}, + {name: "overflow CLI", cli: "18446744073709551616", invalid: true}, + } { + t.Run(tc.name, func(t *testing.T) { + viper.Reset() + t.Cleanup(viper.Reset) + config.ApplyDefaults(viper.GetViper()) + viper.SetConfigType("toml") + require.NoError(t, viper.ReadConfig(strings.NewReader("[validator]\n"+tc.toml))) + cmd := &cobra.Command{} + cmd.Flags().Uint64("wait-to-vote-slot", 0, "") + var err error + if tc.cli != "" { + err = cmd.Flags().Set("wait-to-vote-slot", tc.cli) + } + var got uint64 + if err == nil { + got, err = configuredWaitToVoteSlot(cmd) + } + if tc.invalid { + require.Error(t, err) + return + } + require.NoError(t, err) + require.Equal(t, tc.want, got) + }) + } + require.NotNil(t, Run.Flags().Lookup("wait-to-vote-slot")) +} + +func TestEffectiveWaitToVoteSlot(t *testing.T) { + for _, tc := range []struct{ startup, configured, want uint64 }{ + {100, 0, 108}, + {103, 0, 108}, + {103, 104, 108}, + {103, 108, 108}, + {103, 123, 123}, // Operator cutoff need not align with a leader window. + {103, math.MaxUint64, math.MaxUint64}, + {math.MaxUint64 - 7, 0, math.MaxUint64}, + {math.MaxUint64, 0, math.MaxUint64}, + } { + require.Equal(t, tc.want, effectiveWaitToVoteSlot(tc.startup, tc.configured), "%+v", tc) + } +} diff --git a/config.example.toml b/config.example.toml index 7e3833087..f46e78549 100644 --- a/config.example.toml +++ b/config.example.toml @@ -328,6 +328,16 @@ name = "mithril" # Signature-verification workers (0 = GOMAXPROCS). tpu_sigverify_workers = 0 + # Optional minimum slot for NEW votes (inclusive), also --wait-to-vote-slot. + # Useful when rejoining after recovery. CLI overrides this setting. + # Zero adds no operator cutoff; the automatic startup cutoff and normal + # consensus checks still apply. A lower value cannot bypass those checks. + # Replay/repair continue while waiting. Previously recorded, authenticated + # votes can still be restored/rebroadcast under the existing recovery rules. + # This does not coordinate a cluster restart or wait for supermajority, + # and does not allow resetting a corrupt vote-history file. + wait_to_vote_slot = 0 + # ============================================================================ # [consensus] - Alpenglow Consensus # ============================================================================ diff --git a/docs/reserved-vote-history.md b/docs/reserved-vote-history.md new file mode 100644 index 000000000..bd94e2e57 --- /dev/null +++ b/docs/reserved-vote-history.md @@ -0,0 +1,162 @@ +# Voting persistence and crash recovery + +## Intended guarantee + +A crash must not lead to conflicting externally published votes or reserved-mode +leader actions because the validator forgot its earlier local decisions. The design permits +loss of recent detailed history and sacrifices voting availability when its +completeness is uncertain. It does not promise immediate restart voting, that +every signed vote reaches disk, or recovery from rolled-back safety files. +Normal anti-equivocation, execution, parent and validator-binding checks remain +necessary; this is a persistence contract, not a proof of the entire protocol. + +The failure model includes process termination and host/power failure, provided +successful file and directory syncs survive, the current safety files are +preserved, and only one fenced owner uses the signing identity. Software tests +exercise the recovery decisions; they do not qualify actual storage against +power loss. Valid signatures on saved files establish integrity, not freshness. + +## Two different publication guarantees + +| Mode | Required before a vote can escape to the pool/network | What a restart may trust | +| --- | --- | --- | +| Default synchronous history | Exact validated history is written, file-synced, renamed and directory-synced before local pool admission or network enqueue. | Retained exact decisions and their rooted boundary, subject to normal restoration checks. | +| Opt-in reserved history | The vote's slot is covered by an acknowledged durable reservation before signing. Its exact history snapshot is prepared and queued before publication, without waiting for per-vote I/O. | The startup reservation bound, unless a separately validated clean-history seal proves exact retained history. | + +In synchronous mode, BLS bytes may be computed privately in RAM before history +is synced. The guarantee is **persist before publication**, including local +pool admission because it can publish a certificate. In reserved mode, a +successful queue submission, background rename, or `written` counter is **not** +a durable acknowledgement of that vote. Replacing a whole file atomically is +not the same as making its bytes/directory entry survive power loss. + +Both modes retain complete decisions in memory while running and prune them +only through the normal verified-root rules. The synchronous history guarantee +covers votes; it is not a complete journal of produced leader blocks. The +additional leader reservation barrier applies only in reserved mode. + +## Reserved-mode restart rule + +Let **H** be the reservation's `Through` value loaded at startup, **F** the +verified finality/checkpoint floor, and **S** a proposed signing slot. H bounds +what the previous process *might* have signed; it is not its last actual vote. + +Without a validated clean seal, every vote type and historical local-vote +restoration must obey all of these conditions: + +- **S > H**: never re-sign in the uncertain range during this run, even when + an older history file contains that particular vote. +- **F >= H**: do not sign above the range until verified finality/checkpoint + state has reached its end. F equal to H is sufficient; S equal to H is not. +- **S <= the current acknowledged reservation**, plus all ordinary protocol, + live-joining and configured minimum-slot checks. + +The startup recovery bound stays fixed during the run. Background renewal may +raise the current signing allowance; it does not move the recovery target. +Newly received blocks, elapsed wall time, an RPC tip, replay progress alone, and +`--wait-to-vote-slot` cannot substitute for verified finality/checkpoint state. +The recovery wait can be indefinite if the cluster halts below H. + +For example, suppose votes through 1,015 escaped, detailed history survives only +through 1,012, and the durable reservation is 1,032. After an unclean restart, +slots <= 1,032 remain forbidden. Slot 1,033 is also forbidden while F < 1,032. +Once F >= 1,032, it can pass the recovery gate only after an acknowledged grant +covers 1,033 and all normal voting checks pass. Seeing a new block after restart +does not by itself meet these conditions. + +## Clean shutdown is a durable protocol + +1. Stop/join the voter and leader producer. Halt/join the reservation worker; + drain/join the history writer so an older rename cannot overwrite the seal. +2. Require no latched safety fault or history-write failure, no unresolved + reservation-write uncertainty, and verified recovery through the startup H. +3. Sync exact retained history, including the verified rooted boundary. +4. Sync a reservation record containing the digest of that exact history. + +A successful process exit, a signal handler running, or an attempted final save +is not sufficient. Startup must load and validate the history and reservation, +match the clean digest, and **durably consume the clean marker with a dirty +successor before new vote/leader signing or new history decisions**. A later crash uses +the reservation again, even if the old history still looks valid. + +| Restart state | Vote recovery | Leader recovery | +| --- | --- | --- | +| Valid history and reservation, no matching clean seal | Enforce S > H and F >= H, including restoration. | Enforce S > H and F >= H. | +| Valid matching clean seal, successfully consumed | Resume using exact retained voting decisions and ordinary checks; no additional vote quarantine through H. | Still enforce S > H and F >= H: vote history does not enumerate every block that may have been signed. | +| Missing, corrupt, unreadable or incompatible enrolled state | Refuse automatic reset/startup; do not infer safety from an RPC tip. | Same refusal. | + +Errors during sealing do not authorize treating the session as clean; startup +must validate whichever durable record survived. A failed reservation write +may have reached storage despite its error. Running signers retain only the +previous acknowledged allowance; a later monotonic successful write can resolve +that uncertainty. An unresolved error prevents deliberately sealing clean. + +## Lifetime, storage and enrollment + +The reservation grants up to 32 slots beyond the requested slot and renews when +16 or fewer remain. Only successful file-and-directory sync acknowledgement +publishes new permission. Exhaustion pauses signing while replay/verification +continue. Repeated restarts without new grants do not advance the bound. + +Detailed history uses one ordered writer with one in-flight and at most one +newer pending complete snapshot. New complete snapshots may supersede unwritten +ones. Validation/encoding/signing of the snapshot remain on the voter goroutine; +only file I/O is asynchronous. Writer errors latch a safety fault and stop +voting. An already in-flight vote is still covered by its durable reservation. + +Preserve both `vote_history-.mithril.json` and +`vote_reservation-.mithril.json` independently of AccountsDB snapshots. +The checkpoint encoding cache is only an encoding optimization; it is not the +vote journal or a replacement for the reservation. Rolling back AccountsDB or +an application binary must not roll back either signing-safety file. + +The directory lock only excludes concurrent owners of the same history path. +It cannot fence copies of the identity on other hosts/paths. Signed records and +generation numbers cannot detect an operator restoring an old valid pair of +safety files. Media loss, stale safety-file restoration, dishonest sync behavior +and compromised/copied signing keys are outside this automatic recovery contract. + +Use `--reserved-vote-history` to opt in; default persistence remains synchronous. +First enrollment additionally requires `--initialize-vote-reservation`, exclusive +identity ownership and complete synchronous history from the stopped previous +writer. An empty directory is appropriate only for a previously unused identity. +The software cannot distinguish that case from deleting both files for an old +identity; initialization is not a safe disaster-recovery reset. Remove the +initialization flag afterward. These are CLI flags, not TOML settings. + +Enrolled history uses version 2, rejected by older binaries that do not enforce +the reservation. Missing/corrupt enrolled state or a changed identity, vote +account, authorized voter, genesis or shred version must not be automatically +re-enrolled. Disabling the flag or deleting files is not a supported downgrade. +Transport binding/validation runs before the clean marker is consumed. + +## Live admission is separate from restart recovery + +The live vote admission floor uses retained consensus-pool state and the history +root, not the newest finalization certificate. Replay can still contribute a +notarization within the bounded 16-slot retained tail before later certificates +are collected; finalized retained slots can also receive skips. Durable-root +pruning is ordered behind completed replay events on the voter. These admission +changes also apply to synchronous mode. They do not weaken reserved-mode F >= H. +`--wait-to-vote-slot N` adds an inclusive minimum; it cannot bypass recovery, +execution, retained-root or parent checks. + +## Evidence and limits + +Existing tests map the recovery contract to these cases: + +| Contract | Tests | +| --- | --- | +| Lost valid history suffix; fixed bound across repeated crashes | `TestReservedVotingLostHistorySuffixAndRepeatedCrash` | +| All five vote types, restoration and leader gates at H/H+1 | `TestReservedVotingEverySignatureTypeAtBound` | +| Clean-marker consumption and stricter leader restart | `TestReservedVotingCleanMarkerConsumedBeforeSigning` | +| Digest mismatch, missing/corrupt/domain-mismatched state | `TestReservedCleanDigestMismatchUsesCrashRecovery`, `TestReservedVotingRejectsMissingCorruptOrWrongDomain` | +| Pending/uncertain sync cannot authorize signing | `TestSigningReservationUnacknowledgedSyncCannotAuthorize`, `TestSigningReservationUncertainWriteSurvivesRestart` | +| Writer drain, failure and lost pending snapshots | `TestAsyncHistoryBlockedWriteDoesNotDelayVotesAndCleanCloseDrains`, `TestAsyncHistoryFailureStopsVoterAndPreventsCleanMarker`, `TestAsyncHistoryProcessCrashLosesPendingSnapshots` | + +The subprocess-kill test exercises lost in-flight/pending application snapshots; +it does not power-cycle a host or prove filesystem durability. Fault-injection +and older-history fixtures check the algorithm's decisions under the stated +storage contract. This is not a formal consensus proof or mainnet storage +qualification. Persistence benchmarks likewise do not establish replay or FAST +performance. diff --git a/docs/vote-delivery-persistence-evidence.md b/docs/vote-delivery-persistence-evidence.md new file mode 100644 index 000000000..1ae1f5cfa --- /dev/null +++ b/docs/vote-delivery-persistence-evidence.md @@ -0,0 +1,14 @@ +# Vote Delivery Persistence: benchmark evidence + +The maintained subsystem documentation and reusable Go benchmarks describe the +implementation and reproduction method. Historical raw results and session +notes are retained at [the tested source snapshot](https://github.com/Overclock-Validator/mithril/tree/54b233ff0e27fb929644f7d53bf4a699cb590cd8) +(tag `review-evidence-20260916-vote-delivery-persistence`). They are omitted from this proposed merge. + +[Historical result files](https://github.com/Overclock-Validator/mithril/tree/54b233ff0e27fb929644f7d53bf4a699cb590cd8/docs/results) + +Measurements retain their original baselines. Rebasing onto PR #278 does not +turn an intermediate-version benchmark into a comparison with the new base. +Component timings and short live observations do not establish sustained FAST +inclusion gains. The final review description records validation of the rebased +source separately from historical benchmark results. diff --git a/pkg/alpenglow/vote_history.go b/pkg/alpenglow/vote_history.go index d227bbe12..448977fa3 100644 --- a/pkg/alpenglow/vote_history.go +++ b/pkg/alpenglow/vote_history.go @@ -15,6 +15,7 @@ import ( ) const voteHistoryVersion = 1 +const reservedVoteHistoryVersion = 2 var ErrVoteHistoryNotFound = errors.New("alpenglow vote history not found") @@ -23,6 +24,7 @@ var ErrVoteHistoryNotFound = errors.New("alpenglow vote history not found") // Agave when resuming (notarized blocks and ParentReady edges), in addition to // the anti-equivocation vote sets. type VoteHistory struct { + ReservationRequired bool `json:"reservation_required,omitempty"` Version uint32 `json:"version"` NodePubkey solana.PublicKey `json:"node_pubkey"` Root uint64 `json:"root"` @@ -420,50 +422,111 @@ func VoteHistoryFilename(dir string, node solana.PublicKey) string { return filepath.Join(dir, fmt.Sprintf("vote_history-%s.mithril.json", node)) } -// SaveVoteHistory signs the exact serialized history with the validator -// identity and atomically replaces the previous file before a vote can be -// admitted to consensus or sent to the network. +// SaveVoteHistory authenticates and durably replaces the exact history: write, +// file sync, rename, then directory sync. Synchronous-mode callers require +// success before pool admission (which can publish certificates) or network +// enqueue. The BLS signature may already have been computed privately in RAM; +// this is persist-before-publication, not persist-before-BLS-computation. func SaveVoteHistory(dir string, h *VoteHistory, identity ed25519.PrivateKey) error { + return saveVoteHistory(dir, h, identity, true) +} + +// SaveReservedVoteHistory writes and renames without per-vote sync. Success +// does not prove that this history survived a host/power failure. The caller +// must enforce an independently durable signing reservation and its restart +// quarantine; a valid-looking older history is not evidence of completeness. +func SaveReservedVoteHistory(dir string, h *VoteHistory, identity ed25519.PrivateKey) error { + if h == nil || !h.ReservationRequired { + return fmt.Errorf("reserved history requires durable reservation enrollment") + } + return saveVoteHistory(dir, h, identity, false) +} + +// VoteHistorySnapshot owns signed, immutable bytes. It retains no reference to +// the voter's mutable maps or signing key and can be saved by a worker. +type VoteHistorySnapshot struct { + node solana.PublicKey + encoded []byte +} + +// PrepareReservedVoteHistory validates and signs the complete current history +// without filesystem access. Only the owner of h may call this while mutating h. +func PrepareReservedVoteHistory(h *VoteHistory, identity ed25519.PrivateKey) (*VoteHistorySnapshot, error) { + if h == nil || !h.ReservationRequired { + return nil, fmt.Errorf("reserved history requires durable reservation enrollment") + } + encoded, err := encodeVoteHistory(h, identity) + if err != nil { + return nil, err + } + return &VoteHistorySnapshot{node: h.NodePubkey, encoded: encoded}, nil +} + +// SaveReservedVoteHistorySnapshot replaces history without an explicit sync. +// Success means replacement completed, not durable vote acknowledgement. Callers +// must serialize writes and enforce the independent durable reservation. On an +// unclean restart even an intact snapshot cannot bypass the startup bound. +func SaveReservedVoteHistorySnapshot(dir string, snapshot *VoteHistorySnapshot) error { + if snapshot == nil || len(snapshot.encoded) == 0 { + return fmt.Errorf("save reserved history: empty snapshot") + } + if err := ensureDurableVoteHistoryDirectory(dir); err != nil { + return err + } + return replaceVoteHistoryFile(dir, VoteHistoryFilename(dir, snapshot.node), snapshot.encoded, false) +} + +func saveVoteHistory(dir string, h *VoteHistory, identity ed25519.PrivateKey, durable bool) error { + encoded, err := encodeVoteHistory(h, identity) + if err != nil { + return err + } + if err := ensureDurableVoteHistoryDirectory(dir); err != nil { + return err + } + return replaceVoteHistoryFile(dir, VoteHistoryFilename(dir, h.NodePubkey), encoded, durable) +} + +func encodeVoteHistory(h *VoteHistory, identity ed25519.PrivateKey) ([]byte, error) { if h == nil { - return fmt.Errorf("save vote history: nil history") + return nil, fmt.Errorf("save vote history: nil history") } if len(identity) != ed25519.PrivateKeySize { - return fmt.Errorf("save vote history: invalid identity key size %d", len(identity)) + return nil, fmt.Errorf("save vote history: invalid identity key size %d", len(identity)) } node := solana.PublicKey(identity.Public().(ed25519.PublicKey)) if node != h.NodePubkey { - return fmt.Errorf("save vote history: identity %s does not match history %s", node, h.NodePubkey) + return nil, fmt.Errorf("save vote history: identity %s does not match history %s", node, h.NodePubkey) } h.Version = voteHistoryVersion + if h.ReservationRequired { + h.Version = reservedVoteHistoryVersion + } if err := h.preparePersistedViews(); err != nil { - return fmt.Errorf("save vote history: %w", err) + return nil, fmt.Errorf("save vote history: %w", err) } defer func() { h.PersistedNotarized = nil h.PersistedParentReady = nil }() if err := h.validatePersistedState(); err != nil { - return fmt.Errorf("save vote history: %w", err) + return nil, fmt.Errorf("save vote history: %w", err) } data, err := json.Marshal(h) if err != nil { - return fmt.Errorf("serialize vote history: %w", err) + return nil, fmt.Errorf("serialize vote history: %w", err) } envelope := savedVoteHistory{ - Version: voteHistoryVersion, + Version: h.Version, Node: node, Data: data, Signature: ed25519.Sign(identity, data), } encoded, err := json.Marshal(envelope) if err != nil { - return fmt.Errorf("serialize saved vote history: %w", err) + return nil, fmt.Errorf("serialize saved vote history: %w", err) } - if err := ensureDurableVoteHistoryDirectory(dir); err != nil { - return err - } - filename := VoteHistoryFilename(dir, node) - return persistVoteHistoryFile(dir, filename, encoded) + return encoded, nil } // ensureDurableVoteHistoryDirectory creates each missing path component and @@ -515,6 +578,10 @@ func ensureDurableVoteHistoryDirectory(dir string) error { // alone does not guarantee that either the bytes or the new directory entry // survives a crash. func persistVoteHistoryFile(dir, filename string, encoded []byte) error { + return replaceVoteHistoryFile(dir, filename, encoded, true) +} + +func replaceVoteHistoryFile(dir, filename string, encoded []byte, durable bool) error { temporary, err := os.CreateTemp(dir, "."+filepath.Base(filename)+".tmp-") if err != nil { return fmt.Errorf("create temporary vote history: %w", err) @@ -538,8 +605,10 @@ func persistVoteHistoryFile(dir, filename string, encoded []byte) error { if n != len(encoded) { return fmt.Errorf("write temporary vote history: %w", io.ErrShortWrite) } - if err := temporary.Sync(); err != nil { - return fmt.Errorf("sync temporary vote history: %w", err) + if durable { + if err := temporary.Sync(); err != nil { + return fmt.Errorf("sync temporary vote history: %w", err) + } } closeErr := temporary.Close() closed = true @@ -551,8 +620,8 @@ func persistVoteHistoryFile(dir, filename string, encoded []byte) error { } renamed = true - if err := syncVoteHistoryDirectory(dir); err != nil { - return err + if durable { + return syncVoteHistoryDirectory(dir) } return nil } @@ -586,7 +655,7 @@ func LoadVoteHistory(dir string, node solana.PublicKey) (*VoteHistory, error) { if err := json.Unmarshal(encoded, &envelope); err != nil { return nil, fmt.Errorf("decode saved vote history: %w", err) } - if envelope.Version != voteHistoryVersion || envelope.Node != node { + if (envelope.Version != voteHistoryVersion && envelope.Version != reservedVoteHistoryVersion) || envelope.Node != node { return nil, fmt.Errorf("saved vote history identity/version mismatch") } if !ed25519.Verify(ed25519.PublicKey(node[:]), envelope.Data, envelope.Signature) { @@ -596,7 +665,7 @@ func LoadVoteHistory(dir string, node solana.PublicKey) (*VoteHistory, error) { if err := json.Unmarshal(envelope.Data, &h); err != nil { return nil, fmt.Errorf("decode vote history: %w", err) } - if h.Version != voteHistoryVersion || h.NodePubkey != node { + if h.Version != envelope.Version || h.NodePubkey != node || h.ReservationRequired != (h.Version == reservedVoteHistoryVersion) { return nil, fmt.Errorf("vote history identity/version mismatch") } if err := h.validatePersistedState(); err != nil { diff --git a/pkg/alpenglow/vote_history_snapshot_test.go b/pkg/alpenglow/vote_history_snapshot_test.go new file mode 100644 index 000000000..540e6fa14 --- /dev/null +++ b/pkg/alpenglow/vote_history_snapshot_test.go @@ -0,0 +1,49 @@ +package alpenglow + +import ( + "crypto/ed25519" + "testing" + + "github.com/gagliardetto/solana-go" + "github.com/stretchr/testify/require" +) + +func TestReservedHistorySnapshotIsImmutable(t *testing.T) { + identity := ed25519.NewKeyFromSeed(make([]byte, ed25519.SeedSize)) + node := solana.PublicKey(identity.Public().(ed25519.PublicKey)) + h := NewVoteHistory(node, 10) + h.ReservationRequired = true + block := BlockID{Slot: 11, Hash: solana.Hash{11}} + require.NoError(t, h.AddVote(NewNotarizationVote(11, block.Hash))) + h.NotarizedBlocks[block] = true + h.AddParentReady(12, block) + snapshot, err := PrepareReservedVoteHistory(h, identity) + require.NoError(t, err) + // Mutate/prune every transport-backed collection after taking the snapshot. + h.SetRoot(20) + require.NoError(t, h.AddVote(NewSkipVote(21))) + for i := range identity { + identity[i] = 0 + } + dir := t.TempDir() + require.NoError(t, SaveReservedVoteHistorySnapshot(dir, snapshot)) + loaded, err := LoadVoteHistory(dir, node) + require.NoError(t, err) + require.Equal(t, uint64(10), loaded.Root) + require.True(t, loaded.VotedAt(11)) + require.True(t, loaded.IsBlockNotarized(block)) + require.True(t, loaded.IsParentReady(12, block)) + require.False(t, loaded.HasSkipped(21)) +} + +func TestReservedHistorySnapshotRequiresEnrollmentAndValidHistory(t *testing.T) { + identity := ed25519.NewKeyFromSeed(make([]byte, ed25519.SeedSize)) + h := NewVoteHistory(solana.PublicKey(identity.Public().(ed25519.PublicKey)), 10) + _, err := PrepareReservedVoteHistory(h, identity) + require.Error(t, err) + h.ReservationRequired = true + h.Voted[11] = true // Inconsistent with canonical VotesCast. + _, err = PrepareReservedVoteHistory(h, identity) + require.Error(t, err) + require.Error(t, SaveReservedVoteHistorySnapshot(t.TempDir(), &VoteHistorySnapshot{})) +} diff --git a/pkg/alpenglow/vote_reservation.go b/pkg/alpenglow/vote_reservation.go new file mode 100644 index 000000000..b577da189 --- /dev/null +++ b/pkg/alpenglow/vote_reservation.go @@ -0,0 +1,104 @@ +package alpenglow + +import ( + "crypto/ed25519" + "crypto/sha256" + "encoding/json" + "errors" + "fmt" + "os" + "path/filepath" + + "github.com/gagliardetto/solana-go" + "golang.org/x/sys/unix" +) + +// VoteReservation is the durable upper bound on slots this identity may sign, +// not a record of slots actually signed. It must survive independently of +// AccountsDB checkpoints and must never be restored from an older snapshot. +// A signature authenticates this file; it does not prove freshness against +// rollback, nor does Generation provide an external monotonic counter. +// CleanHistoryDigest permits exact-history vote recovery only after validation +// and durable consumption by a dirty successor before signing. It does not +// certify complete leader-block history. See docs/reserved-vote-history.md. +type VoteReservation struct { + Version uint32 `json:"version"` + Node solana.PublicKey `json:"node"` + VoteAccount solana.PublicKey `json:"vote_account"` + AuthorizedVoter solana.PublicKey `json:"authorized_voter"` + Genesis solana.Hash `json:"genesis"` + ShredVersion uint16 `json:"shred_version"` + Generation uint64 `json:"generation"` + Through uint64 `json:"through"` + CleanHistoryDigest []byte `json:"clean_history_digest,omitempty"` +} + +func VoteReservationFilename(dir string, node solana.PublicKey) string { + return filepath.Join(dir, fmt.Sprintf("vote_reservation-%s.mithril.json", node)) +} + +// LockVoteHistory excludes concurrent owners of this directory. Operators must +// still fence copies of the same identity on other hosts or in other paths. +func LockVoteHistory(dir string, node solana.PublicKey) (*os.File, error) { + if err := ensureDurableVoteHistoryDirectory(dir); err != nil { + return nil, err + } + f, err := os.OpenFile(filepath.Join(dir, ".vote_history-"+node.String()+".lock"), os.O_CREATE|os.O_RDWR, 0600) + if err != nil { + return nil, err + } + if err := unix.Flock(int(f.Fd()), unix.LOCK_EX|unix.LOCK_NB); err != nil { + f.Close() + return nil, fmt.Errorf("vote history already owned: %w", err) + } + return f, nil +} + +func LoadVoteReservation(dir string, node solana.PublicKey) (VoteReservation, error) { + var r VoteReservation + encoded, err := os.ReadFile(VoteReservationFilename(dir, node)) + if err != nil { + return r, err + } + var envelope savedVoteHistory + if err := json.Unmarshal(encoded, &envelope); err != nil { + return r, err + } + if envelope.Version != 1 || envelope.Node != node || !ed25519.Verify(ed25519.PublicKey(node[:]), envelope.Data, envelope.Signature) { + return r, errors.New("invalid vote reservation signature/version/identity") + } + if err := json.Unmarshal(envelope.Data, &r); err != nil { + return r, err + } + if r.Version != 1 || r.Node != node || r.Generation == 0 || (len(r.CleanHistoryDigest) != 0 && len(r.CleanHistoryDigest) != sha256.Size) { + return r, errors.New("invalid vote reservation record") + } + return r, nil +} + +func SaveVoteReservation(dir string, r VoteReservation, identity ed25519.PrivateKey) error { + if len(identity) != ed25519.PrivateKeySize || solana.PublicKey(identity.Public().(ed25519.PublicKey)) != r.Node || r.Version != 1 || r.Generation == 0 { + return errors.New("invalid vote reservation signer/record") + } + data, err := json.Marshal(r) + if err != nil { + return err + } + encoded, err := json.Marshal(savedVoteHistory{Version: 1, Node: r.Node, Data: data, Signature: ed25519.Sign(identity, data)}) + if err != nil { + return err + } + if err := ensureDurableVoteHistoryDirectory(dir); err != nil { + return err + } + return persistVoteHistoryFile(dir, VoteReservationFilename(dir, r.Node), encoded) +} + +func VoteHistoryDigest(dir string, node solana.PublicKey) ([]byte, error) { + data, err := os.ReadFile(VoteHistoryFilename(dir, node)) + if err != nil { + return nil, err + } + digest := sha256.Sum256(data) + return digest[:], nil +} diff --git a/pkg/alpenglow/vote_reservation_test.go b/pkg/alpenglow/vote_reservation_test.go new file mode 100644 index 000000000..9e1840ade --- /dev/null +++ b/pkg/alpenglow/vote_reservation_test.go @@ -0,0 +1,71 @@ +package alpenglow + +import ( + "crypto/ed25519" + "encoding/json" + "os" + "testing" + + "github.com/gagliardetto/solana-go" + "github.com/stretchr/testify/require" +) + +func TestReservedHistoryFormatAndIntegrity(t *testing.T) { + key := ed25519.NewKeyFromSeed(make([]byte, 32)) + node := solana.PublicKey(key.Public().(ed25519.PublicKey)) + dir := t.TempDir() + h := NewVoteHistory(node, 39) + require.Error(t, SaveReservedVoteHistory(dir, h, key)) + h.ReservationRequired = true + require.NoError(t, h.AddVote(NewSkipVote(44))) + require.NoError(t, SaveReservedVoteHistory(dir, h, key)) + raw, err := os.ReadFile(VoteHistoryFilename(dir, node)) + require.NoError(t, err) + var envelope savedVoteHistory + require.NoError(t, json.Unmarshal(raw, &envelope)) + require.Equal(t, uint32(2), envelope.Version, "legacy reader must reject hybrid history") + loaded, err := LoadVoteHistory(dir, node) + require.NoError(t, err) + require.True(t, loaded.HasSkipped(44)) + require.True(t, loaded.ReservationRequired) + envelope.Data[10] ^= 1 + raw, err = json.Marshal(envelope) + require.NoError(t, err) + require.NoError(t, os.WriteFile(VoteHistoryFilename(dir, node), raw, 0600)) + _, err = LoadVoteHistory(dir, node) + require.Error(t, err) +} + +func BenchmarkVoteHistoryPersistence(b *testing.B) { + key := ed25519.NewKeyFromSeed(make([]byte, 32)) + node := solana.PublicKey(key.Public().(ed25519.PublicKey)) + for _, reserved := range []bool{false, true} { + name := "synchronous" + if reserved { + name = "reserved-write-rename" + } + b.Run(name, func(b *testing.B) { + dir := b.TempDir() + h := NewVoteHistory(node, 39) + h.ReservationRequired = reserved + for slot := uint64(40); slot < 72; slot++ { + if err := h.AddVote(NewNotarizationVote(slot, solana.Hash{byte(slot)})); err != nil { + b.Fatal(err) + } + } + save := SaveVoteHistory + if reserved { + save = SaveReservedVoteHistory + } + if err := SaveVoteHistory(dir, h, key); err != nil { + b.Fatal(err) + } + b.ResetTimer() + for i := 0; i < b.N; i++ { + if err := save(dir, h, key); err != nil { + b.Fatal(err) + } + } + }) + } +} diff --git a/pkg/blockprod/leader.go b/pkg/blockprod/leader.go index e304ffb8d..65c062173 100644 --- a/pkg/blockprod/leader.go +++ b/pkg/blockprod/leader.go @@ -97,6 +97,7 @@ type LeaderLoop struct { alpenglowClock bool parentContext func(uint64) ParentContext productionParent func(uint64) alpenglow.BlockProductionParent + canSignSlot func(uint64) bool onBlock func(*b.Block) commitLeaderSlot func(replay.CommitLeaderInput) (*sealevel.SlotCtx, error) @@ -145,6 +146,7 @@ type LeaderLoopConfig struct { AlpenglowClock bool ParentContext func(uint64) ParentContext ProductionParent func(slot uint64) alpenglow.BlockProductionParent + CanSignSlot func(slot uint64) bool OnBlock func(*b.Block) CurrentSlot func() uint64 LeaderForSlot func(uint64) (solana.PublicKey, bool) @@ -181,6 +183,7 @@ func NewLeaderLoop(cfg LeaderLoopConfig) *LeaderLoop { alpenglowClock: cfg.AlpenglowClock, parentContext: cfg.ParentContext, productionParent: cfg.ProductionParent, + canSignSlot: cfg.CanSignSlot, onBlock: cfg.OnBlock, currentSlot: cfg.CurrentSlot, leaderForSlot: cfg.LeaderForSlot, @@ -1101,6 +1104,9 @@ func (l *LeaderLoop) revalidateProductionParentForStartLocked(slot uint64, selec } func (l *LeaderLoop) startSlotLocked(slot uint64) error { + if l.canSignSlot != nil && !l.canSignSlot(slot) { + return fmt.Errorf("%w: waiting for durable signing reservation", errParentNotReady) + } selectedParent, parentReadyRequired, err := l.resolveProductionParent(slot) if err != nil { return err diff --git a/pkg/blockprod/leader_signing_reservation_test.go b/pkg/blockprod/leader_signing_reservation_test.go new file mode 100644 index 000000000..4b7e09c10 --- /dev/null +++ b/pkg/blockprod/leader_signing_reservation_test.go @@ -0,0 +1,17 @@ +package blockprod + +import ( + "github.com/stretchr/testify/require" + "testing" +) + +func TestSigningReservationGatesEveryLeaderSlotBeforeBuild(t *testing.T) { + for slot := uint64(40); slot < 44; slot++ { + var checked uint64 + l := &LeaderLoop{canSignSlot: func(s uint64) bool { checked = s; return false }} + require.ErrorIs(t, l.startSlotLocked(slot), errParentNotReady) + require.Equal(t, slot, checked) + // All other builder dependencies are deliberately nil: rejection must happen + // before accessing a working bank, executing or signing any shreds. + } +} diff --git a/pkg/config/config.go b/pkg/config/config.go index 1e7416c90..c31bf7448 100644 --- a/pkg/config/config.go +++ b/pkg/config/config.go @@ -25,6 +25,7 @@ func ApplyDefaults(v *viper.Viper) { v.SetDefault("validator.tpu_quic_bind_addr", "") v.SetDefault("validator.advertised_ip", "") v.SetDefault("validator.tpu_sigverify_workers", 0) + v.SetDefault("validator.wait_to_vote_slot", uint64(0)) } // LedgerConfig holds ledger-related configuration (matches Firedancer [ledger] section) @@ -235,6 +236,7 @@ type ValidatorConfig struct { TPUQUICBindAddr string `toml:"tpu_quic_bind_addr" mapstructure:"tpu_quic_bind_addr"` AdvertisedIP string `toml:"advertised_ip" mapstructure:"advertised_ip"` TPUSigverifyWorkers int `toml:"tpu_sigverify_workers" mapstructure:"tpu_sigverify_workers"` + WaitToVoteSlot uint64 `toml:"wait_to_vote_slot" mapstructure:"wait_to_vote_slot"` // Minimum slot for new votes; automatic startup cutoff still applies } // Config holds all configuration options for Mithril (Firedancer-style hierarchy) diff --git a/pkg/consensus/engine.go b/pkg/consensus/engine.go index f57cb2ab5..5cf2cb203 100644 --- a/pkg/consensus/engine.go +++ b/pkg/consensus/engine.go @@ -389,7 +389,7 @@ func (e *AlpenglowObserverEngine) EnableVoting(cfg VotingConfig) error { // events, matching Agave's initial_parent_ready selection. if slot, parent, ok := voter.history.HighestParentReadyMatching(func(parent alpenglow.BlockID) bool { return !e.ensureChain().IsObjectivelyInvalidBlock(parent) - }); ok && slot > root.Slot { + }); ok && slot > root.Slot && (voter.reservation == nil || voter.reservation.recoverThrough == 0) { if !e.ensurePool().RestoreParentReady(slot, parent) { mlog.Log.FileOnlyf("ALPENGLOW voting: ignored persisted ParentReady slot=%d parent=%s because newer root/live tracker state is authoritative", slot, parent) } @@ -975,11 +975,18 @@ func (e *AlpenglowObserverEngine) injectLocalVote(message alpenglow.VoteMessage, } } -// alpenglowVoteActionFloor is the highest slot on which this validator must -// not initiate a new vote. The pool root is its strict admission boundary; -// direct finality is included because the pool deliberately retains a short -// reward-accounting tail behind finality where network votes remain useful. +// alpenglowVoteActionFloor is the retained pool's strict admission boundary. +// Network finality is not a voting root: a replayed block may still contribute +// a notarization to a later fast certificate and its slot+8 reward certificate. +// The voter also checks its own persisted history root before signing. func (e *AlpenglowObserverEngine) alpenglowVoteActionFloor() uint64 { + return e.ensurePool().Snapshot().RootSlot +} + +// alpenglowVerifiedFinalityFloor releases crash-recovery reservations. Keep this +// independent of live vote admission: a retained reward window must not weaken +// the requirement to pass every slot that may have been signed before a crash. +func (e *AlpenglowObserverEngine) alpenglowVerifiedFinalityFloor() uint64 { floor := e.ensurePool().Snapshot().RootSlot if finalized := e.ensureChain().Snapshot().LatestDirectFinalizedBlock.Slot; finalized > floor { floor = finalized @@ -1542,6 +1549,25 @@ func (e *AlpenglowObserverEngine) PruneAlpenglowBefore(slot uint64) { if slot == 0 { return } + // Replay enqueues its completed-block event before publishing a durable + // promotion. Retire the pool, execution proof and history on that same + // ordered voter stream, so a fast checkpoint cannot overtake the vote. + e.voterMu.RLock() + voter := e.voter + e.voterMu.RUnlock() + if voter != nil { + if err := voter.enqueue(voterEvent{kind: voterEventDurableRoot, slot: slot}); err != nil { + e.latchSafetyError(err) + } + return + } + e.applyAlpenglowDurableRoot(slot) +} + +// applyAlpenglowDurableRoot is called by the voter after earlier replay events, +// or synchronously by an observer without a voting loop. Startup root restore +// remains a separate, immediate barrier in SetAlpenglowRoot. +func (e *AlpenglowObserverEngine) applyAlpenglowDurableRoot(slot uint64) alpenglow.BlockID { e.poolOutputMu.Lock() defer e.poolOutputMu.Unlock() @@ -1557,13 +1583,11 @@ func (e *AlpenglowObserverEngine) PruneAlpenglowBefore(slot uint64) { if e.certPool != nil { e.certPool.ObserveFloor(slot) } - if err := e.enqueueVoter(voterEvent{kind: voterEventRoot, root: root}); err != nil { - e.latchSafetyError(err) - } // Replay calls this only after the fold through slot is durably committed. // Keep the transport peer window tied to that local root, not to speculative // certificate finality or the pool's reward-retention floor. e.advanceVotorPeerRoot(slot) + return root } func (e *AlpenglowObserverEngine) pruneInvalidBlockIDsBefore(slot uint64) { diff --git a/pkg/consensus/vote_history_writer.go b/pkg/consensus/vote_history_writer.go new file mode 100644 index 000000000..9d175d6d4 --- /dev/null +++ b/pkg/consensus/vote_history_writer.go @@ -0,0 +1,125 @@ +package consensus + +import ( + "errors" + "fmt" + "sync" + + "github.com/Overclock-Validator/mithril/pkg/alpenglow" +) + +// One serial writer, one in-flight snapshot and at most one newer pending +// snapshot. A newer complete history supersedes an unwritten snapshot; every +// retained, unrooted voting decision is still present in that newer history. +// The independently durable reservation, not this queue, authorizes signing. +// In-flight/pending snapshots may be lost on process death; even a completed +// unsynced replacement may be lost on host/power failure. Neither submitted nor +// written is a durable vote acknowledgement. Recovery must use the startup +// reservation unless the separate clean-history seal validates. +type voteHistoryWriter struct { + mu sync.Mutex + pending *alpenglow.VoteHistorySnapshot + closing bool + failure error + submitted uint64 + written uint64 + coalesced uint64 + wake chan struct{} + done chan struct{} + persist func(*alpenglow.VoteHistorySnapshot) error + onError func(error) +} + +func newVoteHistoryWriter(persist func(*alpenglow.VoteHistorySnapshot) error, onError func(error)) *voteHistoryWriter { + w := &voteHistoryWriter{wake: make(chan struct{}, 1), done: make(chan struct{}), persist: persist, onError: onError} + go w.run() + return w +} + +// submit does no I/O and never waits for the writer. The mutex only protects +// pointer/counter changes; neither persistence nor error callbacks hold it. +// A nil return means queued only. It must never replace the reservation check. +func (w *voteHistoryWriter) submit(snapshot *alpenglow.VoteHistorySnapshot) error { + if snapshot == nil { + return errors.New("nil vote-history snapshot") + } + w.mu.Lock() + if w.failure != nil { + err := w.failure + w.mu.Unlock() + return err + } + if w.closing { + w.mu.Unlock() + return errors.New("vote-history writer is closed") + } + if w.pending != nil { + w.coalesced++ + } + w.pending = snapshot + w.submitted++ + w.mu.Unlock() + w.notify() + return nil +} + +func (w *voteHistoryWriter) notify() { + select { + case w.wake <- struct{}{}: + default: + } +} + +func (w *voteHistoryWriter) run() { + defer close(w.done) + for range w.wake { + for { + w.mu.Lock() + snapshot := w.pending + w.pending = nil + closing := w.closing + w.mu.Unlock() + if snapshot == nil { + if closing { + return + } + break + } + if err := w.persist(snapshot); err != nil { + err = fmt.Errorf("background vote-history write: %w", err) + w.mu.Lock() + w.failure = err + w.pending = nil + w.closing = true + w.mu.Unlock() + if w.onError != nil { + w.onError(err) + } + return + } + w.mu.Lock() + w.written++ + w.mu.Unlock() + } + } +} + +// close rejects new submissions and drains every retained snapshot. The voter +// must join this worker before writing and syncing its final clean history, +// otherwise an older in-flight rename could overwrite the sealed history. +func (w *voteHistoryWriter) close() error { + w.mu.Lock() + w.closing = true + w.mu.Unlock() + w.notify() + <-w.done + w.mu.Lock() + defer w.mu.Unlock() + return w.failure +} + +func (w *voteHistoryWriter) counters() (submitted, written, coalesced uint64) { + w.mu.Lock() + defer w.mu.Unlock() + return w.submitted, w.written, w.coalesced +} diff --git a/pkg/consensus/vote_history_writer_test.go b/pkg/consensus/vote_history_writer_test.go new file mode 100644 index 000000000..74c7bbfcd --- /dev/null +++ b/pkg/consensus/vote_history_writer_test.go @@ -0,0 +1,201 @@ +package consensus + +import ( + "errors" + "os" + "os/exec" + "path/filepath" + "sync" + "testing" + "time" + + "github.com/Overclock-Validator/mithril/pkg/alpenglow" + "github.com/gagliardetto/solana-go" + "github.com/stretchr/testify/require" +) + +func TestAsyncHistoryBlockedWriteDoesNotDelayVotesAndCleanCloseDrains(t *testing.T) { + cfg := reservedTestConfig(t.TempDir()) + v, err := openReservedTestVoter(t, cfg, 39) + require.NoError(t, err) + reserveThrough(t, v.reservation, 44) + require.NoError(t, v.historyWriter.close()) + entered, release := make(chan struct{}), make(chan struct{}) + var once sync.Once + unblock := func() { once.Do(func() { close(release) }) } + t.Cleanup(unblock) + first := true // Owned only by the serial writer. + v.historyWriter = newVoteHistoryWriter(func(s *alpenglow.VoteHistorySnapshot) error { + if first { + first = false + close(entered) + <-release + } + return alpenglow.SaveReservedVoteHistorySnapshot(cfg.HistoryDir, s) + }, nil) + voted, err := v.cast(alpenglow.NewSkipVote(44), false) + require.NoError(t, err) + require.True(t, voted) + <-entered + castDone := make(chan error, 1) + go func() { + for _, slot := range []uint64{45, 46} { + ok, err := v.cast(alpenglow.NewSkipVote(slot), false) + if err != nil || !ok { + castDone <- errors.New("vote failed while history writer was blocked") + return + } + } + castDone <- nil + }() + select { + case err := <-castDone: + require.NoError(t, err) + case <-time.After(time.Second): + t.Fatal("disk writer blocked voting") + } + submitted, written, coalesced := v.historyWriter.counters() + require.Equal(t, uint64(3), submitted) + require.Zero(t, written) + require.Equal(t, uint64(1), coalesced) + onDisk, err := alpenglow.LoadVoteHistory(cfg.HistoryDir, v.node) + require.NoError(t, err) + require.False(t, onDisk.HasSkipped(44), "I/O should still be blocked") + // Even when disk history lags, complete in-memory decisions forbid conflict. + voted, err = v.cast(alpenglow.NewNotarizationVote(44, solana.Hash{7}), false) + require.NoError(t, err) + require.False(t, voted) + closed := make(chan error, 1) + go func() { closed <- v.close() }() + require.Eventually(t, func() bool { + v.historyWriter.mu.Lock() + defer v.historyWriter.mu.Unlock() + return v.historyWriter.closing + }, time.Second, time.Millisecond) + select { + case err := <-closed: + t.Fatalf("close returned before draining its writer: %v", err) + default: + } + r, err := alpenglow.LoadVoteReservation(cfg.HistoryDir, v.node) + require.NoError(t, err) + require.Empty(t, r.CleanHistoryDigest) + unblock() + require.NoError(t, <-closed) + onDisk, err = alpenglow.LoadVoteHistory(cfg.HistoryDir, v.node) + require.NoError(t, err) + for _, slot := range []uint64{44, 45, 46} { + require.True(t, onDisk.HasSkipped(slot)) + } + r, err = alpenglow.LoadVoteReservation(cfg.HistoryDir, v.node) + require.NoError(t, err) + digest, err := alpenglow.VoteHistoryDigest(cfg.HistoryDir, v.node) + require.NoError(t, err) + require.Equal(t, digest, r.CleanHistoryDigest) + _, written, _ = v.historyWriter.counters() + require.Equal(t, uint64(2), written, "old in-flight snapshot must finish before newest complete snapshot") +} + +func TestAsyncHistoryFailureIsStickyAndReportedWithoutMoreSubmissions(t *testing.T) { + cfg := reservedTestConfig(t.TempDir()) + h := alpenglow.NewVoteHistory(voterTestValidatorSet(t, cfg.Identity, cfg.AuthorizedVoter, cfg.VoteAccount).Validators[0].NodePubkey, 39) + h.ReservationRequired = true + snapshot, err := alpenglow.PrepareReservedVoteHistory(h, cfg.Identity) + require.NoError(t, err) + diskErr := errors.New("injected disk failure") + reported := make(chan error, 1) + w := newVoteHistoryWriter(func(*alpenglow.VoteHistorySnapshot) error { return diskErr }, func(err error) { reported <- err }) + require.NoError(t, w.submit(snapshot)) + select { + case err := <-reported: + require.ErrorIs(t, err, diskErr) + case <-time.After(time.Second): + t.Fatal("background failure was not reported") + } + require.ErrorIs(t, w.submit(snapshot), diskErr) + require.ErrorIs(t, w.close(), diskErr) + require.ErrorIs(t, w.close(), diskErr) +} + +func TestAsyncHistoryFailureStopsVoterAndPreventsCleanMarker(t *testing.T) { + cfg := reservedTestConfig(t.TempDir()) + e, err := NewEngine(Config{AlpenglowIdentity: cfg.Identity, AlpenglowShredVersion: 0x1234}) + require.NoError(t, err) + t.Cleanup(func() { _ = e.Close() }) + set := voterTestValidatorSet(t, cfg.Identity, cfg.AuthorizedVoter, cfg.VoteAccount) + e.SetAlpenglowEpochLookup(cfg.EpochForSlot) + require.NoError(t, e.SetAlpenglowValidatorSet(set)) + root := alpenglow.BlockID{Slot: 39, Hash: solana.Hash{39}} + e.SetAlpenglowRoot(root) + v, err := newAlpenglowVoterUnstarted(e, cfg, root, []alpenglow.ValidatorSet{set}) + require.NoError(t, err) + t.Cleanup(func() { _ = v.close() }) + reserveThrough(t, v.reservation, 44) + require.NoError(t, v.historyWriter.close()) + diskErr := errors.New("injected background disk failure") + v.historyWriter = newVoteHistoryWriter(func(*alpenglow.VoteHistorySnapshot) error { return diskErr }, v.failHistoryWrite) + require.NoError(t, v.history.AddVote(alpenglow.NewSkipVote(44))) + require.NoError(t, v.saveHistory()) + select { + case <-v.done: + case <-time.After(time.Second): + t.Fatal("disk failure did not stop the voter") + } + require.ErrorIs(t, e.safetyError(), diskErr) + _, _, err = v.sign(alpenglow.NewSkipVote(45), false) + require.ErrorIs(t, err, diskErr) + require.Error(t, v.enqueue(voterEvent{kind: voterEventBlockTimeout, slot: 45})) + require.ErrorIs(t, v.close(), diskErr) + r, err := alpenglow.LoadVoteReservation(cfg.HistoryDir, v.node) + require.NoError(t, err) + require.Empty(t, r.CleanHistoryDigest) +} + +// Kill a subprocess while its first history write is blocked and a newer +// snapshot is pending. Successful earlier reservation syncs survive this +// process crash; this is deliberately not a host power-loss test. +func TestAsyncHistoryProcessCrashLosesPendingSnapshots(t *testing.T) { + const childEnv = "MITHRIL_ASYNC_HISTORY_TEST_DIR" + if dir := os.Getenv(childEnv); dir != "" { + cfg := reservedTestConfig(dir) + v, err := openReservedTestVoter(t, cfg, 39) + require.NoError(t, err) + reserveThrough(t, v.reservation, 44) + require.NoError(t, v.historyWriter.close()) + entered := make(chan struct{}) + v.historyWriter = newVoteHistoryWriter(func(*alpenglow.VoteHistorySnapshot) error { + close(entered) + select {} + }, nil) + for _, slot := range []uint64{44, 45} { + voted, err := v.cast(alpenglow.NewSkipVote(slot), false) + require.NoError(t, err) + require.True(t, voted) + if slot == 44 { + <-entered + } + } + require.NoError(t, os.WriteFile(filepath.Join(dir, "ready"), []byte("ready"), 0600)) + select {} + } + dir := t.TempDir() + cmd := exec.Command(os.Args[0], "-test.run=^TestAsyncHistoryProcessCrashLosesPendingSnapshots$", "-test.count=1") + cmd.Env = append(os.Environ(), childEnv+"="+dir) + require.NoError(t, cmd.Start()) + t.Cleanup(func() { _ = cmd.Process.Kill() }) + require.Eventually(t, func() bool { _, err := os.Stat(filepath.Join(dir, "ready")); return err == nil }, 10*time.Second, time.Millisecond) + require.NoError(t, cmd.Process.Kill()) + require.Error(t, cmd.Wait()) + cfg := reservedTestConfig(dir) + cfg.InitializeVoteReservation = false + v, err := openReservedTestVoter(t, cfg, 39) + require.NoError(t, err) + require.False(t, v.history.HasSkipped(44)) + require.False(t, v.history.HasSkipped(45)) + h := v.reservation.recoverThrough + require.GreaterOrEqual(t, h, uint64(45)) + _, _, err = v.sign(alpenglow.NewSkipVote(44), false) + require.ErrorIs(t, err, errVoterNotReady) + _, _, err = v.sign(alpenglow.NewSkipVote(h+1), false) + require.ErrorIs(t, err, errVoterNotReady, "verified finality must reach the lost history's bound") +} diff --git a/pkg/consensus/vote_reservation.go b/pkg/consensus/vote_reservation.go new file mode 100644 index 000000000..564cfc29d --- /dev/null +++ b/pkg/consensus/vote_reservation.go @@ -0,0 +1,296 @@ +package consensus + +import ( + "bytes" + "crypto/ed25519" + "errors" + "fmt" + "math" + "os" + "sync" + "sync/atomic" + "time" + + "github.com/Overclock-Validator/mithril/pkg/alpenglow" + "github.com/Overclock-Validator/mithril/pkg/mlog" + "github.com/gagliardetto/solana-go" +) + +const signingReserveSlots = uint64(32) +const signingRenewRemaining = uint64(16) + +// signingReservation bounds what a crash may erase from detailed vote history. +// Intended guarantee: losing recent history must not authorize conflicting +// voting/leader actions after restart. It does NOT guarantee that every vote +// survives on disk, immediate restart voting, or recovery from safety-file rollback. +// +// On an enrolled restart, H is the startup reservation. Without a clean-history +// seal, all vote types (including restored votes) are forbidden at slots <= H; +// slots > H also wait until verified finality/checkpoint state reaches H. +// Leaders always obey that startup barrier, even after a clean vote-history seal. +// The current acknowledged Through separately caps every new signing permission. +// Renewal can raise Through, but never moves this run's fixed recovery barrier. +// +// The worker owns record after startup. Only a successful file+directory sync +// publishes Through to signers; a request, queued write or uncertain sync cannot. +// This assumes storage honors sync and a single fenced identity owner preserves +// the current reservation independently of AccountsDB. See docs/reserved-vote-history.md. +type signingReservation struct { + uncertain bool // Worker only, read after halt. Failed sync may have reached storage. + record alpenglow.VoteReservation + through atomic.Uint64 + desired atomic.Uint64 + stopped atomic.Bool + recoverThrough uint64 // Startup H, or zero after first enrollment / a validated clean-history seal. + leaderThrough uint64 // Exact leader production history is not saved: always skip the old range. + wake chan struct{} + changed chan struct{} + stop chan struct{} + done chan struct{} + stopOnce sync.Once + persist func(alpenglow.VoteReservation) error +} + +func openSigningReservation(cfg VotingConfig, node solana.PublicKey, shredVersion uint16, history *alpenglow.VoteHistory) (*signingReservation, error) { + if cfg.Genesis == (solana.Hash{}) { + return nil, errors.New("reserved voting requires the bound genesis hash") + } + expected := alpenglow.VoteReservation{Version: 1, Node: node, VoteAccount: cfg.VoteAccount, AuthorizedVoter: solana.PublicKey(cfg.AuthorizedVoter.Public().(ed25519.PublicKey)), Genesis: cfg.Genesis, ShredVersion: shredVersion} + record, err := alpenglow.LoadVoteReservation(cfg.HistoryDir, node) + initializing := errors.Is(err, os.ErrNotExist) + if initializing { + if !cfg.InitializeVoteReservation || history.ReservationRequired { + return nil, errors.New("missing vote reservation; explicit first enrollment with complete synchronous history is required") + } + record = expected + record.Generation = 1 + record.Through = history.Root + for slot := range history.VotesCast { + record.Through = max(record.Through, slot) + } + // The baseline is made durable before the first reservation is created. + if err := alpenglow.SaveVoteHistory(cfg.HistoryDir, history, cfg.Identity); err != nil { + return nil, err + } + } else if err != nil { + return nil, fmt.Errorf("refuse unsafe reservation reset: %w", err) + } + if record.Node != expected.Node || record.VoteAccount != expected.VoteAccount || record.AuthorizedVoter != expected.AuthorizedVoter || record.Genesis != expected.Genesis || record.ShredVersion != expected.ShredVersion { + return nil, errors.New("vote reservation cluster or signing identity mismatch; explicit domain migration is required") + } + if record.Through == math.MaxUint64 || record.Generation == math.MaxUint64 { + return nil, errors.New("vote reservation exhausted") + } + for slot := range history.VotesCast { + if slot > record.Through { + return nil, fmt.Errorf("history slot %d exceeds durable reservation %d", slot, record.Through) + } + } + r := &signingReservation{record: record, recoverThrough: record.Through, leaderThrough: record.Through, wake: make(chan struct{}, 1), changed: make(chan struct{}, 1), stop: make(chan struct{}), done: make(chan struct{})} + r.persist = func(next alpenglow.VoteReservation) error { + return alpenglow.SaveVoteReservation(cfg.HistoryDir, next, cfg.Identity) + } + if initializing { + r.recoverThrough = 0 + } else if len(record.CleanHistoryDigest) != 0 { + digest, err := alpenglow.VoteHistoryDigest(cfg.HistoryDir, node) + if err != nil { + return nil, err + } + if bytes.Equal(digest, record.CleanHistoryDigest) { + r.recoverThrough = 0 + } + } + // A matching digest proves exact history only for the sealed session. Consume + // that exception with a durably acknowledged dirty successor before allowing + // new vote/leader signing or new detailed-history decisions. A crash after + // this write must use H, + // even if the detailed history file still looks valid or matches the old seal. + r.record.CleanHistoryDigest = nil + r.record.Generation++ + if err := r.persist(r.record); err != nil { + return nil, fmt.Errorf("consume vote reservation session: %w", err) + } + history.ReservationRequired = true + // Version 2 is deliberately rejected by older binaries that do not enforce H. + if err := alpenglow.SaveVoteHistory(cfg.HistoryDir, history, cfg.Identity); err != nil { + return nil, err + } + r.through.Store(r.record.Through) + mlog.Log.Infof("ALPENGLOW signing reservation: through=%d recovery_through=%d leader_recovery_through=%d", r.record.Through, r.recoverThrough, r.leaderThrough) + go r.run() + return r, nil +} + +// allow is nonblocking and protects vote signing/restoration and leader slots. +// For a nonzero startup barrier H, require slot > H AND finalized >= H. Merely +// observing a new block, waiting elapsed time, or replaying past H is not proof +// that prior decisions can be forgotten. Finality comes from verified consensus/ +// checkpoint state, never an RPC tip; --wait-to-vote-slot cannot override it. +// Passing this gate is necessary, not sufficient: normal protocol checks apply. +func (r *signingReservation) allow(slot, finalized uint64, leader bool) bool { + if r.stopped.Load() { + return false + } + floor := r.recoverThrough + if leader { + floor = r.leaderThrough + } + if floor != 0 && (slot <= floor || finalized < floor) { + return false + } + r.request(slot) + return slot <= r.through.Load() +} + +func (r *signingReservation) request(slot uint64) { + for { + old := r.desired.Load() + if slot <= old || r.desired.CompareAndSwap(old, slot) { + break + } + } + through := r.through.Load() + if slot > through || through-slot <= signingRenewRemaining { + select { + case r.wake <- struct{}{}: + default: + } + } +} + +func (r *signingReservation) run() { + defer close(r.done) + retry := time.NewTicker(250 * time.Millisecond) + defer retry.Stop() + var warned bool + for { + select { + case <-r.stop: + return + case <-r.wake: + case <-retry.C: + } + if r.stopped.Load() { + return + } + slot := r.desired.Load() + if slot == 0 || (slot <= r.record.Through && r.record.Through-slot > signingRenewRemaining) { + continue + } + if slot > math.MaxUint64-signingReserveSlots || r.record.Generation == math.MaxUint64 { + continue + } // Never wrap or grant permission. + next := r.record + next.Through = max(next.Through, slot+signingReserveSlots) + next.Generation++ + next.CleanHistoryDigest = nil + if err := r.persist(next); err != nil { + r.uncertain = true + if !warned { + mlog.Log.Errorf("ALPENGLOW signing reservation renewal failed; permission remains through %d: %v", r.through.Load(), err) + warned = true + } + continue // Retry the same or a greater bound; an uncertain sync never grants permission. + } + warned = false + r.uncertain = false + r.record = next + r.through.Store(next.Through) + select { + case r.changed <- struct{}{}: + default: + } + } +} + +func (r *signingReservation) halt() { + r.stopOnce.Do(func() { r.stopped.Store(true); close(r.stop) }) + <-r.done +} + +// seal may be called only after the voter loop and leader producer have stopped, +// and after the ordered history writer has been drained/joined. The caller must +// also establish verified finality >= recoverThrough and no latched safety fault. +// Sync exact history first, then sync its digest in the reservation. A normal +// process exit or successful unsynced rename alone is not a clean seal. On an +// error the caller must not assume cleanliness; restart validates whichever +// durable record survived. The seal never relaxes the next run's leader barrier. +func (r *signingReservation) seal(dir string, history *alpenglow.VoteHistory, identity ed25519.PrivateKey) error { + r.halt() + if r.uncertain { + return errors.New("uncertain reservation write; retaining unclean recovery") + } + if err := alpenglow.SaveVoteHistory(dir, history, identity); err != nil { + return err + } + digest, err := alpenglow.VoteHistoryDigest(dir, history.NodePubkey) + if err != nil { + return err + } + if r.record.Generation == math.MaxUint64 { + return errors.New("vote reservation generation exhausted") + } + next := r.record + next.Generation++ + next.CleanHistoryDigest = digest + return r.persist(next) +} + +// Retain only events blocked on renewal, not historical catch-up traffic. +// Replay the original event after acknowledgement so normal finality, parent, +// execution and invalidation checks still decide whether to vote. +func (v *alpenglowVoter) retainReservationEvent(event voterEvent) { + r := v.reservation + if r == nil { + return + } + var slot uint64 + switch event.kind { + case voterEventBlock: + slot = event.block.Block.Slot + case voterEventBlockTimeout, voterEventCrashedLeaderTimeout: + slot = event.slot | (alpenglow.LeaderWindowSlots - 1) + case voterEventConsensus: + switch event.consensus.Kind { + case alpenglow.ConsensusEventBlockNotarized, alpenglow.ConsensusEventParentReady, alpenglow.ConsensusEventSafeToNotar, alpenglow.ConsensusEventSafeToSkip: + slot = event.consensus.Slot + if event.consensus.Kind == alpenglow.ConsensusEventSafeToNotar || event.consensus.Kind == alpenglow.ConsensusEventSafeToSkip { + slot |= alpenglow.LeaderWindowSlots - 1 + } + default: + return + } + default: + return + } + if slot <= r.through.Load() || slot <= v.admissionFloor() || slot < v.waitToVoteSlot || v.engine.alpenglowVerifiedFinalityFloor() < r.recoverThrough || slot <= r.recoverThrough { + return + } + if !v.votingStarted && v.readyToVote != nil && !v.readyToVote(slot) { + return + } + r.request(slot) + if len(v.reservationEvents) < votorEventQueueSize { + v.reservationEvents = append(v.reservationEvents, event) + } +} + +// AlpenglowCanSignLeaderSlot protects every produced slot, including the +// trailing slots of a leader window. A clean vote-history marker is not a +// complete leader-block history, so leaders always skip the old reservation. +func (e *AlpenglowObserverEngine) AlpenglowCanSignLeaderSlot(slot uint64) bool { + if e.safetyError() != nil { + return false + } + e.voterMu.RLock() + defer e.voterMu.RUnlock() + v := e.voter + if v == nil { + return false + } + if v.reservation == nil { + return true + } + return v.reservation.allow(slot, e.alpenglowVerifiedFinalityFloor(), true) +} diff --git a/pkg/consensus/vote_reservation_test.go b/pkg/consensus/vote_reservation_test.go new file mode 100644 index 000000000..f4e5c418d --- /dev/null +++ b/pkg/consensus/vote_reservation_test.go @@ -0,0 +1,330 @@ +package consensus + +import ( + "crypto/ed25519" + "errors" + "math" + "os" + "sync/atomic" + "testing" + "time" + + "github.com/Overclock-Validator/mithril/pkg/alpenglow" + "github.com/gagliardetto/solana-go" + "github.com/stretchr/testify/require" +) + +func reservedTestConfig(dir string) VotingConfig { + return VotingConfig{Identity: voterTestKey(11), AuthorizedVoter: voterTestKey(12), VoteAccount: solana.PublicKey(voterTestKey(13).Public().(ed25519.PublicKey)), HistoryDir: dir, Genesis: solana.Hash{1}, ReservedHistory: true, InitializeVoteReservation: true, WaitToVoteSlot: 40, ReadyToVote: func(uint64) bool { return true }, EpochForSlot: func(uint64) uint64 { return 7 }, Peers: func([]alpenglow.ValidatorStake) []alpenglow.VotorPeer { return nil }} +} + +func openReservedTestVoter(t *testing.T, cfg VotingConfig, root uint64) (*alpenglowVoter, error) { + t.Helper() + e, err := NewEngine(Config{AlpenglowIdentity: cfg.Identity, AlpenglowShredVersion: 0x1234}) + require.NoError(t, err) + t.Cleanup(func() { require.NoError(t, e.Close()) }) + set := voterTestValidatorSet(t, cfg.Identity, cfg.AuthorizedVoter, cfg.VoteAccount) + e.SetAlpenglowEpochLookup(cfg.EpochForSlot) + require.NoError(t, e.SetAlpenglowValidatorSet(set)) + block := alpenglow.BlockID{Slot: root, Hash: solana.Hash{byte(root)}} + e.SetAlpenglowRoot(block) + v, err := newAlpenglowVoterUnstarted(e, cfg, block, []alpenglow.ValidatorSet{set}) + if err == nil { + t.Cleanup(func() { require.NoError(t, v.close()) }) + } + return v, err +} + +func reserveThrough(t *testing.T, r *signingReservation, slot uint64) { + t.Helper() + r.request(slot) + require.Eventually(t, func() bool { return r.through.Load() >= slot }, time.Second, time.Millisecond) +} + +// Simulate loss of this process without executing the clean shutdown protocol. +func crashReservedTestVoter(t *testing.T, v *alpenglowVoter) { + t.Helper() + v.shutdownOnce.Do(func() { + v.closeOnce.Do(func() { close(v.done) }) + v.wg.Wait() + v.reservation.halt() + if v.historyWriter != nil { + require.NoError(t, v.historyWriter.close()) + } + require.NoError(t, v.broadcaster.Close()) + require.NoError(t, v.historyLock.Close()) + }) +} + +func TestReservedVotingLostHistorySuffixAndRepeatedCrash(t *testing.T) { + cfg := reservedTestConfig(t.TempDir()) + v, err := openReservedTestVoter(t, cfg, 39) + require.NoError(t, err) + reserveThrough(t, v.reservation, 60) + baseline, err := os.ReadFile(alpenglow.VoteHistoryFilename(cfg.HistoryDir, v.node)) + require.NoError(t, err) + voted, err := v.cast(alpenglow.NewNotarizationVote(60, solana.Hash{1}), false) + require.NoError(t, err) + require.True(t, voted) + oldH := v.reservation.through.Load() + crashReservedTestVoter(t, v) + // Reproduce a host crash retaining a valid older version of detailed history. + require.NoError(t, os.WriteFile(alpenglow.VoteHistoryFilename(cfg.HistoryDir, v.node), baseline, 0600)) + cfg.InitializeVoteReservation = false + resumed, err := openReservedTestVoter(t, cfg, 39) + require.NoError(t, err) + require.Equal(t, oldH, resumed.reservation.recoverThrough) + for _, vote := range []alpenglow.Vote{alpenglow.NewNotarizationVote(60, solana.Hash{2}), alpenglow.NewSkipVote(60), alpenglow.NewFinalizationVote(60), alpenglow.NewNotarizationFallbackVote(60, solana.Hash{2}), alpenglow.NewSkipFallbackVote(60), alpenglow.NewSkipVote(oldH + 1)} { + // Check at the signing boundary, including the restoration bypass. + for _, normal := range []bool{false, true} { + _, _, err := resumed.sign(vote, normal) + require.ErrorIs(t, err, errVoterNotReady) + } + } + require.Equal(t, oldH, resumed.reservation.through.Load(), "recovery must not keep moving its target") + crashReservedTestVoter(t, resumed) + resumed, err = openReservedTestVoter(t, cfg, oldH) + require.NoError(t, err) + require.Equal(t, oldH, resumed.reservation.recoverThrough) + reserveThrough(t, resumed.reservation, oldH+1) + voted, err = resumed.cast(alpenglow.NewSkipVote(oldH+1), false) + require.NoError(t, err) + require.True(t, voted) + voted, err = resumed.cast(alpenglow.NewSkipVote(oldH), false) + require.NoError(t, err) + require.False(t, voted) +} + +func TestReservedVotingCleanMarkerConsumedBeforeSigning(t *testing.T) { + cfg := reservedTestConfig(t.TempDir()) + v, err := openReservedTestVoter(t, cfg, 39) + require.NoError(t, err) + reserveThrough(t, v.reservation, 44) + voted, err := v.cast(alpenglow.NewSkipVote(44), false) + require.NoError(t, err) + require.True(t, voted) + require.NoError(t, v.close()) + r, err := alpenglow.LoadVoteReservation(cfg.HistoryDir, v.node) + require.NoError(t, err) + require.NotEmpty(t, r.CleanHistoryDigest) + cfg.InitializeVoteReservation = false + resumed, err := openReservedTestVoter(t, cfg, 39) + require.NoError(t, err) + require.Zero(t, resumed.reservation.recoverThrough) + require.True(t, resumed.history.HasSkipped(44)) + r, err = alpenglow.LoadVoteReservation(cfg.HistoryDir, v.node) + require.NoError(t, err) + require.Empty(t, r.CleanHistoryDigest) + voted, err = resumed.cast(alpenglow.NewSkipVote(45), false) + require.NoError(t, err) + require.True(t, voted) + require.False(t, resumed.reservation.allow(45, 39, true), "clean vote history does not authorize repeating leader blocks") + crashReservedTestVoter(t, resumed) + again, err := openReservedTestVoter(t, cfg, 39) + require.NoError(t, err) + require.Equal(t, r.Through, again.reservation.recoverThrough) + // A shutdown before recovering the uncertain range must not mark it clean. + require.NoError(t, again.close()) + r, err = alpenglow.LoadVoteReservation(cfg.HistoryDir, v.node) + require.NoError(t, err) + require.Empty(t, r.CleanHistoryDigest) +} + +func TestReservedVotingRejectsMissingCorruptOrWrongDomain(t *testing.T) { + for _, which := range []string{"missing_history", "missing_bound", "corrupt_bound", "genesis", "authorized", "vote_account", "synchronous"} { + t.Run(which, func(t *testing.T) { + cfg := reservedTestConfig(t.TempDir()) + v, err := openReservedTestVoter(t, cfg, 39) + require.NoError(t, err) + reserveThrough(t, v.reservation, 44) + crashReservedTestVoter(t, v) + cfg.InitializeVoteReservation = false + switch which { + case "missing_history": + require.NoError(t, os.Remove(alpenglow.VoteHistoryFilename(cfg.HistoryDir, v.node))) + case "missing_bound": + require.NoError(t, os.Remove(alpenglow.VoteReservationFilename(cfg.HistoryDir, v.node))) + cfg.InitializeVoteReservation = true + case "corrupt_bound": + require.NoError(t, os.WriteFile(alpenglow.VoteReservationFilename(cfg.HistoryDir, v.node), []byte("{"), 0600)) + case "genesis": + cfg.Genesis = solana.Hash{2} + case "authorized": + cfg.AuthorizedVoter = voterTestKey(19) + case "vote_account": + cfg.VoteAccount = solana.PublicKey{9} + case "synchronous": + cfg.ReservedHistory = false + } + _, err = openReservedTestVoter(t, cfg, 39) + require.Error(t, err) + }) + } +} + +func TestReservedVotingRequiresEnrollmentAndExclusiveOwner(t *testing.T) { + cfg := reservedTestConfig(t.TempDir()) + cfg.InitializeVoteReservation = false + _, err := openReservedTestVoter(t, cfg, 39) + require.Error(t, err) + cfg.InitializeVoteReservation = true + v, err := openReservedTestVoter(t, cfg, 39) + require.NoError(t, err) + _, err = openReservedTestVoter(t, cfg, 39) + require.ErrorContains(t, err, "already owned") + require.NoError(t, v.close()) +} + +func TestSigningReservationUnacknowledgedSyncCannotAuthorize(t *testing.T) { + entered, release := make(chan struct{}), make(chan struct{}) + r := &signingReservation{record: alpenglow.VoteReservation{Through: 64, Generation: 1}, wake: make(chan struct{}, 1), changed: make(chan struct{}, 1), stop: make(chan struct{}), done: make(chan struct{})} + r.through.Store(64) + var calls atomic.Uint64 + r.persist = func(next alpenglow.VoteReservation) error { + calls.Add(1) + close(entered) + <-release + return nil + } + go r.run() + require.True(t, r.allow(64, 64, false)) + <-entered + require.Equal(t, uint64(64), r.through.Load()) + require.False(t, r.allow(65, 64, false)) + close(release) + require.Eventually(t, func() bool { return r.through.Load() > 64 }, time.Second, time.Millisecond) + require.True(t, r.allow(65, 64, false)) + r.halt() + require.Equal(t, uint64(1), calls.Load()) +} + +func TestSigningReservationUncertainWriteSurvivesRestart(t *testing.T) { + cfg := reservedTestConfig(t.TempDir()) + v, err := openReservedTestVoter(t, cfg, 39) + require.NoError(t, err) + // Stop the original worker, then exercise a fresh worker against the real record. + v.reservation.halt() + record, err := alpenglow.LoadVoteReservation(cfg.HistoryDir, v.node) + require.NoError(t, err) + r := &signingReservation{record: record, wake: make(chan struct{}, 1), changed: make(chan struct{}, 1), stop: make(chan struct{}), done: make(chan struct{})} + r.through.Store(record.Through) + wrote := make(chan struct{}, 1) + r.persist = func(next alpenglow.VoteReservation) error { + err := alpenglow.SaveVoteReservation(cfg.HistoryDir, next, cfg.Identity) + select { + case wrote <- struct{}{}: + default: + } + if err != nil { + return err + } + return errors.New("injected lost sync acknowledgement") + } + v.reservation = r + go r.run() + require.False(t, r.allow(60, 39, false)) + <-wrote + r.halt() + require.Equal(t, record.Through, r.through.Load()) + durable, err := alpenglow.LoadVoteReservation(cfg.HistoryDir, v.node) + require.NoError(t, err) + require.Greater(t, durable.Through, record.Through) + require.ErrorContains(t, r.seal(cfg.HistoryDir, v.history, cfg.Identity), "uncertain") + crashReservedTestVoter(t, v) + cfg.InitializeVoteReservation = false + resumed, err := openReservedTestVoter(t, cfg, 39) + require.NoError(t, err) + require.Equal(t, durable.Through, resumed.reservation.recoverThrough) +} + +func TestSigningReservationNeverWraps(t *testing.T) { + r := &signingReservation{record: alpenglow.VoteReservation{Through: 64, Generation: 1}, wake: make(chan struct{}, 1), changed: make(chan struct{}, 1), stop: make(chan struct{}), done: make(chan struct{}), persist: func(alpenglow.VoteReservation) error { t.Error("overflow attempted persistence"); return nil }} + r.through.Store(64) + go r.run() + require.False(t, r.allow(math.MaxUint64, 64, false)) + r.halt() + require.Equal(t, uint64(64), r.through.Load()) +} + +func TestReservationRetryRechecksFinality(t *testing.T) { + cfg := reservedTestConfig(t.TempDir()) + v, err := openReservedTestVoter(t, cfg, 39) + require.NoError(t, err) + event := voterEvent{kind: voterEventBlockTimeout, slot: 44} + require.NoError(t, v.handle(event)) + require.NotEmpty(t, v.reservationEvents) + reserveThrough(t, v.reservation, 44) + v.engine.SetAlpenglowRoot(alpenglow.BlockID{Slot: 47, Hash: solana.Hash{47}}) + pending := v.reservationEvents + v.reservationEvents = nil + for _, e := range pending { + require.NoError(t, v.handle(e)) + } + require.False(t, v.history.HasSkipped(44), "finalized work must not be signed after a delayed ack") +} + +func TestReservedVotingEverySignatureTypeAtBound(t *testing.T) { + cfg := reservedTestConfig(t.TempDir()) + v, err := openReservedTestVoter(t, cfg, 39) + require.NoError(t, err) + reserveThrough(t, v.reservation, 44) + h := v.reservation.through.Load() + // Freeze acknowledgement while leaving the signing guard active. + v.reservation.halt() + v.reservation.stopped.Store(false) + defer v.reservation.stopped.Store(true) + for _, slot := range []uint64{h, h + 1} { + votes := []alpenglow.Vote{alpenglow.NewNotarizationVote(slot, solana.Hash{2}), alpenglow.NewSkipVote(slot), alpenglow.NewFinalizationVote(slot), alpenglow.NewNotarizationFallbackVote(slot, solana.Hash{2}), alpenglow.NewSkipFallbackVote(slot)} + for _, vote := range votes { + for _, normal := range []bool{false, true} { + _, _, err := v.sign(vote, normal) + if slot == h { + require.NoError(t, err) + } else { + require.ErrorIs(t, err, errVoterNotReady) + } + } + } + } +} + +func TestReservedCleanDigestMismatchUsesCrashRecovery(t *testing.T) { + cfg := reservedTestConfig(t.TempDir()) + v, err := openReservedTestVoter(t, cfg, 39) + require.NoError(t, err) + baseline, err := os.ReadFile(alpenglow.VoteHistoryFilename(cfg.HistoryDir, v.node)) + require.NoError(t, err) + reserveThrough(t, v.reservation, 44) + voted, err := v.cast(alpenglow.NewSkipVote(44), false) + require.NoError(t, err) + require.True(t, voted) + require.NoError(t, v.close()) + require.NoError(t, os.WriteFile(alpenglow.VoteHistoryFilename(cfg.HistoryDir, v.node), baseline, 0600)) + resumed, err := openReservedTestVoter(t, cfg, 39) + require.NoError(t, err) + require.NotZero(t, resumed.reservation.recoverThrough) +} + +func TestReservationLoopRetriesAfterAcknowledgement(t *testing.T) { + cfg := reservedTestConfig(t.TempDir()) + v, err := openReservedTestVoter(t, cfg, 39) + require.NoError(t, err) + v.start() + require.NoError(t, v.enqueue(voterEvent{kind: voterEventBlockTimeout, slot: 44})) + // Read via stats, not mutable voter history, while its loop is running. + require.Eventually(t, func() bool { v.landingMu.RLock(); defer v.landingMu.RUnlock(); return v.stats.VotesCastThisRun == 4 }, time.Second, time.Millisecond) + require.NoError(t, v.close()) + for slot := uint64(44); slot <= 47; slot++ { + require.True(t, v.history.HasSkipped(slot)) + } +} + +func TestReservationRetainsWindowCrossingBound(t *testing.T) { + cfg := reservedTestConfig(t.TempDir()) + v, err := openReservedTestVoter(t, cfg, 39) + require.NoError(t, err) + reserveThrough(t, v.reservation, 44) + h := v.reservation.through.Load() // 76: window ends at 79. + v.retainReservationEvent(voterEvent{kind: voterEventBlockTimeout, slot: h}) + require.Len(t, v.reservationEvents, 1, "trailing skip slots require retry even if the first slot fits") +} diff --git a/pkg/consensus/voter.go b/pkg/consensus/voter.go index e5a220d96..4a9cbbc15 100644 --- a/pkg/consensus/voter.go +++ b/pkg/consensus/voter.go @@ -5,6 +5,7 @@ import ( "errors" "fmt" "math" + "os" "sort" "sync" "time" @@ -38,21 +39,30 @@ type VotingPeerSource func(validators []alpenglow.ValidatorStake) []alpenglow.Vo // account address; AuthorizedVoter is the Ed25519 signer from which its BLS key // was registered. type VotingConfig struct { - Identity ed25519.PrivateKey - AuthorizedVoter ed25519.PrivateKey - VoteAccount solana.PublicKey - HistoryDir string - EpochForSlot func(slot uint64) uint64 - Peers VotingPeerSource - SlotDuration time.Duration - WaitToVoteSlot uint64 - ReadyToVote func(slot uint64) bool + Identity ed25519.PrivateKey + AuthorizedVoter ed25519.PrivateKey + VoteAccount solana.PublicKey + HistoryDir string + ReservedHistory bool + InitializeVoteReservation bool + Genesis solana.Hash + EpochForSlot func(slot uint64) uint64 + Peers VotingPeerSource + SlotDuration time.Duration + WaitToVoteSlot uint64 // Inclusive minimum for new votes; authenticated history restoration is separate + ReadyToVote func(slot uint64) bool } // VotingStats exposes positive network evidence separately from local casting. // NetworkLandedVotes counts unique persisted votes whose rank appeared in the // exact BLS-verified certificate proof received over Votor QUIC. type VotingStats struct { + HistorySnapshotsSubmitted uint64 `json:"history_snapshots_submitted,omitempty"` + HistorySnapshotsWritten uint64 `json:"history_snapshots_written,omitempty"` + HistorySnapshotsCoalesced uint64 `json:"history_snapshots_coalesced,omitempty"` + ReservedHistory bool `json:"reserved_history"` + SigningReservedThrough uint64 `json:"signing_reserved_through,omitempty"` + RecoveryThrough uint64 `json:"recovery_through,omitempty"` Enabled bool `json:"enabled"` VotesCastThisRun uint64 `json:"votes_cast_this_run"` NetworkLandedVotes uint64 `json:"network_landed_votes"` @@ -68,7 +78,7 @@ type VotingStats struct { BroadcastPeerQueueDrops uint64 `json:"broadcast_peer_queue_drops"` BroadcastPeerQueueDiscarded uint64 `json:"broadcast_peer_queue_discarded"` BroadcastPeerSendTimeouts uint64 `json:"broadcast_peer_send_timeouts"` - BroadcastPeerQueueMaxDelay time.Duration `json:"broadcast_peer_queue_max_delay"` + BroadcastPeerQueueMaxDelay time.Duration `json:"broadcast_peer_queue_max_delay_ns"` BroadcastPeerQueues []alpenglow.VotorPeerQueueStats `json:"broadcast_peer_queues,omitempty"` BroadcastDesiredPeers int `json:"broadcast_desired_peers"` BroadcastActiveConnections int `json:"broadcast_active_connections"` @@ -93,6 +103,7 @@ const ( voterEventValidatorSet voterEventRoot voterEventNetworkCertificate + voterEventDurableRoot ) type voterEvent struct { @@ -114,43 +125,49 @@ type pendingVotorBlock struct { // history decisions run on loop; validator-set snapshots are protected only so // the outbound peer callback can read them from broadcast workers. type alpenglowVoter struct { - engine *AlpenglowObserverEngine - identity ed25519.PrivateKey - node solana.PublicKey - voteAccount solana.PublicKey - signer *alpenglow.BLSSigner - historyDir string - history *alpenglow.VoteHistory - epochForSlot func(uint64) uint64 - peerSource VotingPeerSource - slotDuration time.Duration - waitToVoteSlot uint64 - readyToVote func(slot uint64) bool - broadcaster *alpenglow.VotorBroadcaster - events chan voterEvent - done chan struct{} - startOnce sync.Once - closeOnce sync.Once - wg sync.WaitGroup - setsMu sync.RWMutex - sets map[uint64]alpenglow.ValidatorSet - restored map[alpenglow.VoteMessageKey]bool - pending map[uint64][]pendingVotorBlock - receivedShred map[uint64]bool - timeoutsSet map[uint64]bool - executedBlocks map[alpenglow.BlockID]bool - highestFinal uint64 - lastFinalizedAt time.Time - votingStarted bool - latestLiveSlot uint64 - standstillSlot *uint64 - refreshQueue []alpenglow.Message - refreshCursor int - lastWarn map[uint64]time.Time - landingMu sync.RWMutex - landed map[alpenglow.VoteMessageKey]struct{} - stats VotingStats - lastStatsLog time.Time + engine *AlpenglowObserverEngine + identity ed25519.PrivateKey + node solana.PublicKey + voteAccount solana.PublicKey + signer *alpenglow.BLSSigner + historyDir string + historyLock *os.File + reservation *signingReservation + historyWriter *voteHistoryWriter + reservationEvents []voterEvent + shutdownOnce sync.Once + shutdownErr error + history *alpenglow.VoteHistory + epochForSlot func(uint64) uint64 + peerSource VotingPeerSource + slotDuration time.Duration + waitToVoteSlot uint64 + readyToVote func(slot uint64) bool + broadcaster *alpenglow.VotorBroadcaster + events chan voterEvent + done chan struct{} + startOnce sync.Once + closeOnce sync.Once + wg sync.WaitGroup + setsMu sync.RWMutex + sets map[uint64]alpenglow.ValidatorSet + restored map[alpenglow.VoteMessageKey]bool + pending map[uint64][]pendingVotorBlock + receivedShred map[uint64]bool + timeoutsSet map[uint64]bool + executedBlocks map[alpenglow.BlockID]bool + highestFinal uint64 + lastFinalizedAt time.Time + votingStarted bool + latestLiveSlot uint64 + standstillSlot *uint64 + refreshQueue []alpenglow.Message + refreshCursor int + lastWarn map[uint64]time.Time + landingMu sync.RWMutex + landed map[alpenglow.VoteMessageKey]struct{} + stats VotingStats + lastStatsLog time.Time // beforeVoteGuard is a deterministic test seam for invalidation races. It // is nil in production. beforeVoteGuard func(alpenglow.BlockID) @@ -168,6 +185,9 @@ func newAlpenglowVoterUnstarted(engine *AlpenglowObserverEngine, cfg VotingConfi } func newAlpenglowVoterWithStart(engine *AlpenglowObserverEngine, cfg VotingConfig, root alpenglow.BlockID, start bool, sets []alpenglow.ValidatorSet) (*alpenglowVoter, error) { + if cfg.InitializeVoteReservation && !cfg.ReservedHistory { + return nil, errors.New("initialize-vote-reservation requires reserved-vote-history") + } if engine == nil { return nil, fmt.Errorf("enable Alpenglow voting: nil consensus engine") } @@ -197,17 +217,52 @@ func newAlpenglowVoterWithStart(engine *AlpenglowObserverEngine, cfg VotingConfi if node != engineNode { return nil, fmt.Errorf("enable Alpenglow voting: identity %s does not match consensus transport identity %s", node, engineNode) } + historyLock, err := alpenglow.LockVoteHistory(cfg.HistoryDir, node) + if err != nil { + return nil, err + } + keepLock := false + defer func() { + if !keepLock { + historyLock.Close() + } + }() + if !cfg.ReservedHistory { + if _, err := os.Stat(alpenglow.VoteReservationFilename(cfg.HistoryDir, node)); !errors.Is(err, os.ErrNotExist) { + return nil, fmt.Errorf("existing or unreadable vote reservation requires reserved history mode") + } + } history, err := alpenglow.LoadVoteHistory(cfg.HistoryDir, node) if err != nil { if !errors.Is(err, alpenglow.ErrVoteHistoryNotFound) { return nil, fmt.Errorf("enable Alpenglow voting: refuse unsafe vote-history reset: %w", err) } + if cfg.ReservedHistory { + if _, err := os.Stat(alpenglow.VoteReservationFilename(cfg.HistoryDir, node)); !errors.Is(err, os.ErrNotExist) || !cfg.InitializeVoteReservation { + return nil, fmt.Errorf("missing history for reserved voter; refusing automatic reset") + } + } history = alpenglow.NewVoteHistory(node, root.Slot) if err := alpenglow.SaveVoteHistory(cfg.HistoryDir, history, cfg.Identity); err != nil { return nil, fmt.Errorf("initialize Alpenglow vote history: %w", err) } mlog.Log.FileOnlyf("ALPENGLOW voting: created new vote history for %s at root %d; do not reuse this vote account on another validator", node, root.Slot) } + if history.ReservationRequired && !cfg.ReservedHistory { + return nil, errors.New("reserved vote history cannot be opened in synchronous mode") + } + var reservation *signingReservation + if cfg.ReservedHistory { + reservation, err = openSigningReservation(cfg, node, engine.shredVersion, history) + if err != nil { + return nil, err + } + defer func() { + if !keepLock { + reservation.halt() + } + }() + } if history.Root < root.Slot { history.SetRoot(root.Slot) if err := alpenglow.SaveVoteHistory(cfg.HistoryDir, history, cfg.Identity); err != nil { @@ -225,6 +280,8 @@ func newAlpenglowVoterWithStart(engine *AlpenglowObserverEngine, cfg VotingConfi voteAccount: cfg.VoteAccount, signer: signer, historyDir: cfg.HistoryDir, + historyLock: historyLock, + reservation: reservation, history: history, epochForSlot: cfg.EpochForSlot, peerSource: cfg.Peers, @@ -266,6 +323,12 @@ func newAlpenglowVoterWithStart(engine *AlpenglowObserverEngine, cfg VotingConfi return nil, err } v.broadcaster = broadcaster + if reservation != nil { + v.historyWriter = newVoteHistoryWriter(func(snapshot *alpenglow.VoteHistorySnapshot) error { + return alpenglow.SaveReservedVoteHistorySnapshot(v.historyDir, snapshot) + }, v.failHistoryWrite) + } + keepLock = true if start { v.start() } @@ -305,12 +368,25 @@ func (v *alpenglowVoter) loop() { v.closeOnce.Do(func() { close(v.done) }) v.wg.Done() }() + var reservationChanged <-chan struct{} + if v.reservation != nil { + reservationChanged = v.reservation.changed + } ticker := time.NewTicker(time.Second) defer ticker.Stop() for { select { case <-v.done: return + case <-reservationChanged: + pending := v.reservationEvents + v.reservationEvents = nil + for _, event := range pending { + if err := v.handle(event); err != nil { + v.engine.latchSafetyError(fmt.Errorf("reservation retry: %w", err)) + return + } + } case event := <-v.events: if err := v.handle(event); err != nil { v.engine.latchSafetyError(fmt.Errorf("voting engine: %w", err)) @@ -329,6 +405,11 @@ func (v *alpenglowVoter) loop() { } func (v *alpenglowVoter) handle(event voterEvent) error { + v.retainReservationEvent(event) + return v.handleEvent(event) +} + +func (v *alpenglowVoter) handleEvent(event voterEvent) error { floor := v.admissionFloor() if event.kind == voterEventBlock && v.engine.ensureChain().IsObjectivelyInvalidBlock(event.block.Block) { return nil @@ -368,6 +449,9 @@ func (v *alpenglowVoter) handle(event voterEvent) error { return nil } return v.saveHistory() + case voterEventDurableRoot: + root := v.engine.applyAlpenglowDurableRoot(event.slot) + return v.handleEvent(voterEvent{kind: voterEventRoot, root: root}) case voterEventNetworkCertificate: v.recordNetworkCertificate(event.certificate) return nil @@ -482,9 +566,9 @@ func (v *alpenglowVoter) handleConsensus(event alpenglow.ConsensusEvent) error { func (v *alpenglowVoter) admissionFloor() uint64 { floor := v.history.Root - if v.highestFinal > floor { - floor = v.highestFinal - } + // highestFinal tracks network progress and standstill, not retirement of + // our own decisions. Keep ParentReady, pending replay and exact vote history + // available until the retained pool or an ordered durable root retires them. if engineFloor := v.engine.alpenglowVoteActionFloor(); engineFloor > floor { floor = engineFloor } @@ -666,7 +750,7 @@ func (v *alpenglowVoter) castTarget(vote alpenglow.Vote, restoring bool, guarded if !restoring && v.beforeVoteGuard != nil { v.beforeVoteGuard(guardedBlock) } - // Hold through signing, durable history, pool admission, and broadcast. + // Hold through signing, history recording, pool admission, and broadcast. // Objective invalidation takes the write side before changing the chain, // so a new or restored vote is wholly before it or sees the tombstone. v.engine.invalidActionMu.RLock() @@ -687,18 +771,22 @@ func (v *alpenglowVoter) castTarget(vote alpenglow.Vote, restoring bool, guarded return false, nil } if !restoring { - // Finality can advance while the BLS signature is computed. Avoid a - // durable stale record when that race is already visible here; if it - // advances later, atomic admission below classifies it benignly. + // Retention/root pruning can advance while the BLS signature is computed. + // Avoid an expired record if that race is already visible here; atomic + // admission below classifies a later pruning race benignly. if vote.Slot <= v.admissionFloor() { return false, nil } if err := v.history.AddVote(vote); err != nil { return false, fmt.Errorf("record %s vote at slot %d: %w", vote.Type, vote.Slot, err) } - // Pool admission may synchronously assemble and publish a certificate. - // Persist the anti-equivocation record first so no externally visible - // proof can survive a crash without its signed local history. + // Pool admission can publish a certificate. In reserved mode the durable + // upper bound covers loss of this unsynchronized history replacement; + // synchronous mode still persists the exact history before admission. + // sign computed BLS bytes in RAM, but nothing may expose them before + // this boundary succeeds. In reserved mode saveHistory only queues the + // snapshot: restart safety comes from the durable reservation checked + // before sign, not from assuming this snapshot reached durable storage. if err := v.saveHistory(); err != nil { return false, err } @@ -732,7 +820,16 @@ func (v *alpenglowVoter) castTarget(vote alpenglow.Vote, restoring bool, guarded return true, nil } +// sign checks reservation recovery even when restoration bypasses the live +// joining gate. Re-signing a saved vote is still signing; the presence of an +// older valid history file cannot prove that its lost suffix was conflict-free. func (v *alpenglowVoter) sign(vote alpenglow.Vote, respectVotingGate bool) (alpenglow.VoteMessage, alpenglow.VoteVerifyResult, error) { + if err := v.engine.safetyError(); err != nil { + return alpenglow.VoteMessage{}, alpenglow.VoteVerifyResult{}, err + } + if v.reservation != nil && !v.reservation.allow(vote.Slot, v.engine.alpenglowVerifiedFinalityFloor(), false) { + return alpenglow.VoteMessage{}, alpenglow.VoteVerifyResult{}, fmt.Errorf("%w: waiting for verified recovery or durable signing reservation", errVoterNotReady) + } if respectVotingGate { if err := v.votingGateError(vote.Slot); err != nil { return alpenglow.VoteMessage{}, alpenglow.VoteVerifyResult{}, err @@ -773,7 +870,7 @@ func (v *alpenglowVoter) votingGateError(slot uint64) error { return fmt.Errorf("%w: slot %d is at or below consensus action floor %d", errVoterNotReady, slot, floor) } if slot < v.waitToVoteSlot { - return fmt.Errorf("%w: waiting for startup watermark slot %d", errVoterNotReady, v.waitToVoteSlot) + return fmt.Errorf("%w: waiting for voting cutoff slot %d", errVoterNotReady, v.waitToVoteSlot) } // ReadyToVote is a startup join guard, not a perpetual clock check. Once an // accepted live block or vote joins Votor, verified ParentReady and timeout @@ -782,6 +879,9 @@ func (v *alpenglowVoter) votingGateError(slot uint64) error { if !v.votingStarted && v.readyToVote != nil && !v.readyToVote(slot) { return fmt.Errorf("%w: slot is still behind the startup live voting window", errVoterNotReady) } + if v.reservation != nil && !v.reservation.allow(slot, v.engine.alpenglowVerifiedFinalityFloor(), false) { + return fmt.Errorf("%w: waiting for verified recovery or durable signing reservation", errVoterNotReady) + } return nil } @@ -842,6 +942,9 @@ func (v *alpenglowVoter) votorTransportValidators() []alpenglow.ValidatorStake { } func (v *alpenglowVoter) restoreVotesForEpoch(epoch uint64) error { + if v.reservation != nil && v.reservation.recoverThrough != 0 { + return nil + } floor := v.admissionFloor() for _, vote := range v.history.VotesAfter(v.history.Root - minU64(v.history.Root, 1)) { // Keep every signed history entry for anti-equivocation, but do not @@ -889,6 +992,9 @@ func (v *alpenglowVoter) restoreVotesForEpoch(epoch uint64) error { // persist-before-admission crash window without trusting process-local invalid // block state across a restart. func (v *alpenglowVoter) restoreVotesForBlock(block alpenglow.BlockID) (bool, error) { + if v.reservation != nil && v.reservation.recoverThrough != 0 { + return false, nil + } if block.Slot <= v.admissionFloor() { return false, nil } @@ -1091,13 +1197,29 @@ func (v *alpenglowVoter) isIdentityStaked(slot uint64) bool { return false } +// saveHistory is a publication boundary with different persistence semantics: +// synchronous mode acknowledges durable exact history; reserved mode acknowledges +// an immutable queued snapshot only. The latter relies on the signing reservation. func (v *alpenglowVoter) saveHistory() error { + if v.historyWriter != nil { + snapshot, err := alpenglow.PrepareReservedVoteHistory(v.history, v.identity) + if err != nil { + return err + } + return v.historyWriter.submit(snapshot) + } if err := alpenglow.SaveVoteHistory(v.historyDir, v.history, v.identity); err != nil { return fmt.Errorf("persist vote history before consensus publication: %w", err) } return nil } +func (v *alpenglowVoter) failHistoryWrite(err error) { + v.engine.latchSafetyError(err) + mlog.Log.Errorf("ALPENGLOW VOTING SAFETY: %v", err) + v.closeOnce.Do(func() { close(v.done) }) +} + func (v *alpenglowVoter) recordNetworkCertificate(cert alpenglow.Certificate) { set, ok := v.validatorSet(cert.Slot) if !ok { @@ -1206,6 +1328,14 @@ func (v *alpenglowVoter) snapshot() VotingStats { stats := v.stats v.landingMu.RUnlock() stats.Enabled = true + if v.historyWriter != nil { + stats.HistorySnapshotsSubmitted, stats.HistorySnapshotsWritten, stats.HistorySnapshotsCoalesced = v.historyWriter.counters() + } + if v.reservation != nil { + stats.ReservedHistory = true + stats.SigningReservedThrough = v.reservation.through.Load() + stats.RecoveryThrough = v.reservation.recoverThrough + } if v.broadcaster != nil { broadcast := v.broadcaster.Stats() stats.BroadcastMessagesQueued = broadcast.MessagesQueued @@ -1238,7 +1368,7 @@ func (v *alpenglowVoter) maybeLogStats() { } v.lastStatsLog = time.Now() stats := v.snapshot() - mlog.Log.FileOnlyf("alpenglow voting stats: votes_cast_this_run=%d network_landed=%d last_landed_slot=%d broadcast_queued=%d broadcast_dropped=%d peer_sends=%d peer_sends_skipped=%d peer_send_errors=%d desired_peers=%d active_connections=%d pending_connections=%d connection_attempts=%d connection_errors=%d connection_jobs_dropped=%d", + mlog.Log.FileOnlyf("alpenglow voting stats: votes_cast_this_run=%d network_landed=%d last_landed_slot=%d broadcast_queued=%d broadcast_dropped=%d peer_sends=%d peer_sends_skipped=%d peer_send_errors=%d peer_queue_drops=%d peer_queue_discarded=%d peer_send_timeouts=%d peer_queue_max_delay=%s desired_peers=%d active_connections=%d pending_connections=%d connection_attempts=%d connection_errors=%d connection_jobs_dropped=%d reserved_history=%t signing_through=%d recovery_through=%d history_submitted=%d history_written=%d history_coalesced=%d", stats.VotesCastThisRun, stats.NetworkLandedVotes, stats.LastNetworkLandedSlot, @@ -1247,12 +1377,18 @@ func (v *alpenglowVoter) maybeLogStats() { stats.BroadcastPeerSends, stats.BroadcastPeerSendsSkipped, stats.BroadcastPeerSendErrors, + stats.BroadcastPeerQueueDrops, + stats.BroadcastPeerQueueDiscarded, + stats.BroadcastPeerSendTimeouts, + stats.BroadcastPeerQueueMaxDelay, stats.BroadcastDesiredPeers, stats.BroadcastActiveConnections, stats.BroadcastPendingConnections, stats.BroadcastConnectionAttempts, stats.BroadcastConnectionErrors, stats.BroadcastConnectionJobsDropped, + stats.ReservedHistory, stats.SigningReservedThrough, stats.RecoveryThrough, + stats.HistorySnapshotsSubmitted, stats.HistorySnapshotsWritten, stats.HistorySnapshotsCoalesced, ) } @@ -1284,9 +1420,28 @@ func (v *alpenglowVoter) close() error { if v == nil { return nil } - v.closeOnce.Do(func() { close(v.done) }) - v.wg.Wait() - return v.broadcaster.Close() + v.shutdownOnce.Do(func() { + v.closeOnce.Do(func() { close(v.done) }) + v.wg.Wait() + if v.reservation != nil { + v.reservation.halt() + if v.historyWriter != nil { + v.shutdownErr = v.historyWriter.close() + } + // Exiting normally is insufficient: an unresolved recovery barrier, + // writer failure or safety fault must leave the session unsealed. + floor := v.engine.alpenglowVerifiedFinalityFloor() + if v.shutdownErr == nil && v.engine.safetyError() == nil && floor >= v.reservation.recoverThrough { + v.history.SetRoot(floor) + v.shutdownErr = v.reservation.seal(v.historyDir, v.history, v.identity) + if v.shutdownErr == nil { + mlog.Log.Infof("ALPENGLOW signing reservation: clean history sealed at root=%d through=%d", v.history.Root, v.reservation.through.Load()) + } + } + } + v.shutdownErr = errors.Join(v.shutdownErr, v.broadcaster.Close(), v.historyLock.Close()) + }) + return v.shutdownErr } func pendingContains(blocks []pendingVotorBlock, candidate pendingVotorBlock) bool { diff --git a/pkg/consensus/voter_finality_ordering_test.go b/pkg/consensus/voter_finality_ordering_test.go new file mode 100644 index 000000000..a53abb640 --- /dev/null +++ b/pkg/consensus/voter_finality_ordering_test.go @@ -0,0 +1,192 @@ +package consensus + +import ( + "context" + "crypto/ed25519" + "testing" + + "github.com/Overclock-Validator/mithril/pkg/alpenglow" + "github.com/Overclock-Validator/mithril/pkg/block" + "github.com/gagliardetto/solana-go" + "github.com/stretchr/testify/require" +) + +// Keep the actor unstarted so tests can place finality, replay and durable +// promotion in a deterministic order, using the real engine event queue. +func newOrderingTestVoter(t *testing.T, reserved bool) (*AlpenglowObserverEngine, *alpenglowVoter) { + t.Helper() + cfg := reservedTestConfig(t.TempDir()) + cfg.ReservedHistory = reserved + cfg.InitializeVoteReservation = reserved + root := alpenglow.BlockID{Slot: 39, Hash: solana.Hash{39}} + e, err := NewEngine(Config{AlpenglowIdentity: cfg.Identity, AlpenglowShredVersion: 0x1234}) + require.NoError(t, err) + t.Cleanup(func() { require.NoError(t, e.Close()) }) + set := voterTestValidatorSet(t, cfg.Identity, cfg.AuthorizedVoter, cfg.VoteAccount) + e.SetAlpenglowEpochLookup(cfg.EpochForSlot) + require.NoError(t, e.SetAlpenglowValidatorSet(set)) + e.SetAlpenglowRoot(root) + v, err := newAlpenglowVoterUnstarted(e, cfg, root, []alpenglow.ValidatorSet{set}) + require.NoError(t, err) + e.voter = v + if reserved { + reserveThrough(t, v.reservation, 44) + } + // Equivalent to startup's trusted-root ParentReady seed. + require.True(t, v.history.AddParentReady(40, root)) + return e, v +} + +func drainOrderingEvents(t *testing.T, v *alpenglowVoter) { + t.Helper() + for i := 0; i < 1000; i++ { + select { + case event := <-v.events: + require.NoError(t, v.handle(event)) + default: + return + } + } + t.Fatal("voter event queue did not drain") +} + +func observeOrderingBlock(t *testing.T, e *AlpenglowObserverEngine, slot uint64) alpenglow.BlockID { + t.Helper() + id := alpenglow.BlockID{Slot: slot, Hash: solana.Hash{byte(slot)}} + require.NoError(t, e.ObserveBlock(context.Background(), BlockObservation{Source: "ordering-test", Block: &block.Block{ + Slot: slot, ParentSlot: slot - 1, + AlpenglowBlockID: [32]byte(id.Hash), HasAlpenglowBlockID: true, + AlpenglowParentBlockID: [32]byte{byte(slot - 1)}, HasAlpenglowParentBlockID: true, + }})) + return id +} + +func finalizeOrderingBlock(t *testing.T, e *AlpenglowObserverEngine, v *alpenglowVoter, id alpenglow.BlockID) { + t.Helper() + // The two peers supply the 60% slow-finality quorum without our vote; + // adding our 30% notarization later can produce a fast certificate. + for _, vote := range []alpenglow.Vote{alpenglow.NewNotarizationVote(id.Slot, id.Hash), alpenglow.NewFinalizationVote(id.Slot)} { + for rank, key := range []ed25519.PrivateKey{voterTestKey(21), voterTestKey(22)} { + peer := signedVerifiedVoterPeerVote(t, e, v.sets[7], key, uint16(rank+1), vote) + _, err := e.acceptVerifiedVoteResult(peer) + require.NoError(t, err) + } + } + require.Equal(t, id.Slot, e.ensureChain().Snapshot().LatestDirectFinalizedBlock.Slot) +} + +func TestAlpenglowVoterReplaysFourBlocksAfterNetworkFinality(t *testing.T) { + for _, reserved := range []bool{false, true} { + name := "synchronous" + if reserved { + name = "reserved" + } + t.Run(name, func(t *testing.T) { + e, v := newOrderingTestVoter(t, reserved) + for slot := uint64(40); slot <= 43; slot++ { + id := observeOrderingBlock(t, e, slot) + finalizeOrderingBlock(t, e, v, id) + drainOrderingEvents(t, v) + require.Equal(t, slot, v.highestFinal) + require.Less(t, v.admissionFloor(), slot) + require.False(t, v.history.VotedAt(slot), "network finality does not prove local execution") + before := v.snapshot() + require.NoError(t, e.OnReplayResult(context.Background(), SlotReplayResult{Slot: slot, Source: "ordering-test"})) + drainOrderingEvents(t, v) + hash, ok := v.history.NotarizedVote(slot) + require.True(t, ok, "replay must still notarize after slow finality") + require.Equal(t, id.Hash, hash) + require.Greater(t, v.snapshot().BroadcastMessagesQueued, before.BroadcastMessagesQueued) + require.NoError(t, e.AlpenglowSafetyError()) + } + }) + } +} + +func TestAlpenglowVoterNetworkFinalityDuringLocalAdmission(t *testing.T) { + e, v := newOrderingTestVoter(t, false) + id := observeOrderingBlock(t, e, 40) + called := false + v.beforeLocalVoteInject = func(vote alpenglow.Vote) { + if called { + return + } + called = true + require.Equal(t, alpenglow.NewNotarizationVote(id.Slot, id.Hash), vote) + finalizeOrderingBlock(t, e, v, id) + } + require.NoError(t, e.OnReplayResult(context.Background(), SlotReplayResult{Slot: id.Slot})) + drainOrderingEvents(t, v) + require.True(t, called) + require.Positive(t, v.snapshot().VotesCastThisRun) + message, _, err := v.sign(alpenglow.NewNotarizationVote(id.Slot, id.Hash), false) + require.NoError(t, err) + require.True(t, e.ensurePool().HasVerifiedVote(message)) + require.NoError(t, e.AlpenglowSafetyError()) +} + +func TestAlpenglowVoterDurableRootCannotOvertakeQueuedReplay(t *testing.T) { + e, v := newOrderingTestVoter(t, true) + id := observeOrderingBlock(t, e, 40) + finalizeOrderingBlock(t, e, v, id) + drainOrderingEvents(t, v) + require.NoError(t, e.OnReplayResult(context.Background(), SlotReplayResult{Slot: id.Slot})) + e.PruneAlpenglowBefore(id.Slot) + require.Equal(t, uint64(39), e.ensurePool().Snapshot().RootSlot) + require.Contains(t, e.executedReplayBlocks, id, "queued replay must retain its execution proof") + queuedBefore := v.snapshot().BroadcastMessagesQueued + drainOrderingEvents(t, v) + require.Greater(t, v.snapshot().BroadcastMessagesQueued, queuedBefore) + require.Equal(t, id.Slot, v.history.Root) + require.Equal(t, id.Slot, e.ensurePool().Snapshot().RootSlot) + require.NotContains(t, e.executedReplayBlocks, id) + _, ok := v.history.NotarizedVote(id.Slot) + require.True(t, ok, "root vote must remain available to its intra-window child") + child := observeOrderingBlock(t, e, 41) + finalizeOrderingBlock(t, e, v, child) + require.NoError(t, e.OnReplayResult(context.Background(), SlotReplayResult{Slot: child.Slot})) + drainOrderingEvents(t, v) + hash, ok := v.history.NotarizedVote(child.Slot) + require.True(t, ok) + require.Equal(t, child.Hash, hash) + require.NoError(t, e.AlpenglowSafetyError()) +} + +func TestAlpenglowVoterFinalityPreservesEarlierSkipDecision(t *testing.T) { + e, v := newOrderingTestVoter(t, true) + require.NoError(t, v.history.AddVote(alpenglow.NewSkipVote(40))) + id := observeOrderingBlock(t, e, 40) + finalizeOrderingBlock(t, e, v, id) + drainOrderingEvents(t, v) + require.NoError(t, e.OnReplayResult(context.Background(), SlotReplayResult{Slot: 40})) + drainOrderingEvents(t, v) + require.True(t, v.history.HasSkipped(40)) + _, ok := v.history.NotarizedVote(40) + require.False(t, ok, "late replay must never replace an earlier round-one decision") + require.NoError(t, e.AlpenglowSafetyError()) +} + +func TestReservedRecoveryUsesVerifiedFinalityNotLiveAdmissionFloor(t *testing.T) { + cfg := reservedTestConfig(t.TempDir()) + v, err := openReservedTestVoter(t, cfg, 39) + require.NoError(t, err) + reserveThrough(t, v.reservation, 40) + h := v.reservation.through.Load() + crashReservedTestVoter(t, v) + cfg.InitializeVoteReservation = false + v, err = openReservedTestVoter(t, cfg, 39) + require.NoError(t, err) + _, _, err = v.sign(alpenglow.NewSkipVote(h+1), false) + require.ErrorIs(t, err, errVoterNotReady) + id := observeOrderingBlock(t, v.engine, h) + finalizeOrderingBlock(t, v.engine, v, id) + require.Less(t, v.admissionFloor(), h) + require.Equal(t, h, v.engine.alpenglowVerifiedFinalityFloor()) + reserveThrough(t, v.reservation, h+1) + for _, normal := range []bool{false, true} { + _, _, err = v.sign(alpenglow.NewSkipVote(h), normal) + require.ErrorIs(t, err, errVoterNotReady) + _, _, err = v.sign(alpenglow.NewSkipVote(h+1), normal) + require.NoError(t, err) + } +} diff --git a/pkg/consensus/voter_wait_slot_test.go b/pkg/consensus/voter_wait_slot_test.go new file mode 100644 index 000000000..b06a35819 --- /dev/null +++ b/pkg/consensus/voter_wait_slot_test.go @@ -0,0 +1,91 @@ +package consensus + +import ( + "crypto/ed25519" + "testing" + + "github.com/Overclock-Validator/mithril/pkg/alpenglow" + "github.com/gagliardetto/solana-go" + "github.com/stretchr/testify/require" +) + +func waitSlotTestVoter(t *testing.T, cutoff uint64) *alpenglowVoter { + t.Helper() + identity, authorized := voterTestKey(11), voterTestKey(12) + voteAccount := solana.PublicKey(voterTestKey(13).Public().(ed25519.PublicKey)) + set := voterTestValidatorSet(t, identity, authorized, voteAccount) + root := alpenglow.BlockID{Slot: 39, Hash: solana.Hash{0x39}} + engine, err := NewEngine(Config{AlpenglowShredVersion: 0x1234, AlpenglowIdentity: identity}) + require.NoError(t, err) + t.Cleanup(func() { require.NoError(t, engine.Close()) }) + engine.SetAlpenglowEpochLookup(func(uint64) uint64 { return set.Epoch }) + require.NoError(t, engine.SetAlpenglowValidatorSet(set)) + engine.SetAlpenglowRoot(root) + voter, err := newAlpenglowVoterUnstarted(engine, VotingConfig{ + Identity: identity, AuthorizedVoter: authorized, VoteAccount: voteAccount, + HistoryDir: t.TempDir(), EpochForSlot: func(uint64) uint64 { return set.Epoch }, + Peers: func([]alpenglow.ValidatorStake) []alpenglow.VotorPeer { return nil }, + WaitToVoteSlot: cutoff, ReadyToVote: func(uint64) bool { return true }, + }, root, []alpenglow.ValidatorSet{set}) + require.NoError(t, err) + t.Cleanup(func() { require.NoError(t, voter.close()) }) + return voter +} + +func TestWaitToVoteSlotGatesEveryNewVoteType(t *testing.T) { + voter := waitSlotTestVoter(t, 44) + constructors := []func(uint64) alpenglow.Vote{ + func(slot uint64) alpenglow.Vote { return alpenglow.NewNotarizationVote(slot, solana.Hash{1}) }, + alpenglow.NewFinalizationVote, + alpenglow.NewSkipVote, + func(slot uint64) alpenglow.Vote { return alpenglow.NewNotarizationFallbackVote(slot, solana.Hash{1}) }, + alpenglow.NewSkipFallbackVote, + } + for _, started := range []bool{false, true} { + voter.votingStarted = started + for _, voteAt := range constructors { + vote := voteAt(43) + _, _, err := voter.sign(vote, true) + require.ErrorIs(t, err, errVoterNotReady, "%s started=%t", vote.Type, started) + for _, slot := range []uint64{44, 45} { + message, _, err := voter.sign(voteAt(slot), true) + require.NoError(t, err, "%s slot=%d started=%t", vote.Type, slot, started) + require.Equal(t, voteAt(slot), message.Vote) + } + } + } +} + +func TestWaitToVoteSlotSplitsSkipWindowAndPersistsOnlyAllowedVotes(t *testing.T) { + voter := waitSlotTestVoter(t, 42) + require.NoError(t, voter.trySkipWindow(40)) + for _, slot := range []uint64{40, 41} { + require.False(t, voter.history.VotedAt(slot)) + } + for _, slot := range []uint64{42, 43} { + require.True(t, voter.history.HasSkipped(slot)) + } + restored, err := alpenglow.LoadVoteHistory(voter.historyDir, voter.node) + require.NoError(t, err) + require.Equal(t, voter.history.VotesCast, restored.VotesCast) + require.EqualValues(t, 2, voter.engine.ensurePool().Snapshot().VerifiedVotes) + // Joining live voting must not make older slots eligible afterward. + require.True(t, voter.votingStarted) + voted, err := voter.cast(alpenglow.NewSkipVote(41), false) + require.NoError(t, err) + require.False(t, voted) + require.False(t, voter.history.VotedAt(41)) +} + +func TestWaitToVoteSlotPreservesAuthenticatedHistoryRestoration(t *testing.T) { + voter := waitSlotTestVoter(t, 44) + require.NoError(t, voter.history.AddVote(alpenglow.NewSkipVote(40))) + require.NoError(t, voter.saveHistory()) + restored, err := alpenglow.LoadVoteHistory(voter.historyDir, voter.node) + require.NoError(t, err) + voter.history = restored + require.NoError(t, voter.restoreVotesForEpoch(voter.epochForSlot(40))) + require.EqualValues(t, 1, voter.engine.ensurePool().Snapshot().VerifiedVotes) + require.False(t, voter.votingStarted, "restoring a recorded vote must not bypass startup readiness") + require.ErrorIs(t, voter.votingGateError(41), errVoterNotReady) +}