From c9c10d91dfdbfd89221fb30aff44b4da65ac7aff Mon Sep 17 00:00:00 2001 From: iurii Date: Thu, 2 Apr 2026 18:47:31 +0300 Subject: [PATCH 01/14] runner: optimize away BaseRunner.mtx --- protocol/v2/ssv/runner/runner.go | 34 +++++++------------ protocol/v2/ssv/validator/committee.go | 7 +++- protocol/v2/ssv/validator/timer.go | 45 +++++++++----------------- protocol/v2/ssv/validator/validator.go | 5 +++ 4 files changed, 38 insertions(+), 53 deletions(-) diff --git a/protocol/v2/ssv/runner/runner.go b/protocol/v2/ssv/runner/runner.go index 710bd23e05..6fb608a93d 100644 --- a/protocol/v2/ssv/runner/runner.go +++ b/protocol/v2/ssv/runner/runner.go @@ -6,7 +6,6 @@ import ( "fmt" "maps" "slices" - "sync" "github.com/attestantio/go-eth2-client/spec/phase0" ssz "github.com/ferranbt/fastssz" @@ -82,8 +81,17 @@ type DoppelgangerProvider interface { var _ Runner = new(CommitteeRunner) type BaseRunner struct { - mtx sync.RWMutex - State *State + // State stores the current runner state, this state corresponds to 1 particular duty the runner is + // currently busy with at the moment. The BaseRunner is not responsible for synchronizing any updates + // State might need to record - the caller is responsible to ensure the updates/reads (these can happen + // whenever runner's method is called to process a p2p message, or an event) are applied sequentially, + // plus the caller is also responsible for ensuring there is no race with moving on to the next duty + // (the baseSetupForNewDuty call). + // Note, the current implementation achieves concurrent safety by making sure every State read/update + // is done by the same go-routing, handling all the messages in queue.SSVMessage (p2p messages and events) + // sequentially. + State *State + Share map[phase0.ValidatorIndex]*spectypes.Share QBFTController *controller.Controller NetworkConfig *networkconfig.Network @@ -193,16 +201,7 @@ func (b *BaseRunner) SetHighestDecidedSlot(slot phase0.Slot) { // baseSetupForNewDuty is sets the runner for a new duty func (b *BaseRunner) baseSetupForNewDuty(duty spectypes.Duty, quorum uint64) { - // start new state - // start new state - // TODO nicer way to get quorum - state := NewRunnerState(quorum, duty) - - // TODO: potentially incomplete locking of b.State. runner.Execute(duty) has access to - // b.State but currently does not write to it - b.mtx.Lock() // writes to b.State - b.State = state - b.mtx.Unlock() + b.State = NewRunnerState(quorum, duty) } // baseStartNewDuty is a base func that all runner implementation can call to start a duty @@ -465,9 +464,6 @@ func (b *BaseRunner) decide( // hasRunningDuty returns true if a new duty didn't start or an existing duty marked as finished func (b *BaseRunner) hasRunningDuty() bool { - b.mtx.RLock() // reads b.State - defer b.mtx.RUnlock() - if b.State == nil { return false } @@ -476,16 +472,10 @@ func (b *BaseRunner) hasRunningDuty() bool { } func (b *BaseRunner) hasDutyAssigned() bool { - b.mtx.RLock() // reads b.State - defer b.mtx.RUnlock() - return b.State != nil } func (b *BaseRunner) hasDutyFinished() bool { - b.mtx.RLock() // reads b.State - defer b.mtx.RUnlock() - if b.State == nil { return false } diff --git a/protocol/v2/ssv/validator/committee.go b/protocol/v2/ssv/validator/committee.go index d0977279ca..8f1cfccb83 100644 --- a/protocol/v2/ssv/validator/committee.go +++ b/protocol/v2/ssv/validator/committee.go @@ -325,7 +325,12 @@ func (c *Committee) ProcessMessage(ctx context.Context, logger *zap.Logger, msg dutyRunner, found := c.Runners[slot] c.mtx.RUnlock() if !found { - return fmt.Errorf("event message: no committee runner found for slot %d", slot) + // Old runners are pruned, timeout-event issuer is unaware of that - that's why we can end up here + return nil + } + if !dutyRunner.HasRunningDuty() { + // Duties terminate eventually, timeout-event issuer is unaware of that - that's why we can end up here + return nil } timeoutData, err := eventMsg.GetTimeoutData() diff --git a/protocol/v2/ssv/validator/timer.go b/protocol/v2/ssv/validator/timer.go index 218a82adf7..57c6910d19 100644 --- a/protocol/v2/ssv/validator/timer.go +++ b/protocol/v2/ssv/validator/timer.go @@ -22,38 +22,26 @@ func (v *Validator) onTimeout(ctx context.Context, logger *zap.Logger, identifie v.mtx.RLock() // read-lock for v.Queues defer v.mtx.RUnlock() - // The relevant queue might not have been initialized yet, hence we need to check for nil here + // If the relevant queue hasn't been initialized yet, there probably isn't a running duty we can issue a + // timeout for, in practice this should never happen - but we need to handle this just in case. q := v.Queues[identifier.GetRoleType()] if q == nil { return } - dr := v.DutyRunners[identifier.GetRoleType()] - if dr == nil { - // runner can be nil: expired committee runners are removed, but timeout event can still be. in this case we should just skip it - logger.Warn("❗no duty runner found for role", fields.RunnerRole(identifier.GetRoleType())) - return - } - hasDuty := dr.HasRunningDuty() - if !hasDuty { - return - } - msg, err := v.createTimerMessage(identifier, height, round) if err != nil { - logger.Debug("❗ failed to create timer msg", zap.Error(err)) + logger.Error("❗ failed to create timer msg", zap.Error(err)) return } dec, err := queue.DecodeSSVMessage(msg) if err != nil { - logger.Debug("❌ failed to decode timer msg", zap.Error(err)) + logger.Error("❌ failed to decode timer msg", zap.Error(err)) return } if pushed := q.TryPush(dec); !pushed { - logger.Warn("❗️ dropping timeout message because the queue is full", - fields.RunnerRole(identifier.GetRoleType()), - ) + logger.Error("❗️ dropping timeout message because the queue is full", fields.RunnerRole(identifier.GetRoleType())) return } } @@ -86,33 +74,30 @@ func (v *Validator) createTimerMessage(identifier spectypes.MessageID, height sp func (c *Committee) onTimeout(ctx context.Context, logger *zap.Logger, identifier spectypes.MessageID, height specqbft.Height) roundtimer.OnRoundTimeoutF { return func(round specqbft.Round) { - c.mtx.RLock() // read-lock for c.Queues, c.Runners + c.mtx.RLock() // read-lock for c.Queues defer c.mtx.RUnlock() - dr := c.Runners[phase0.Slot(height)] - if dr == nil { // only happens when we prune expired runners - logger.Debug("❗no committee runner found for slot") - return - } - - hasDuty := dr.HasRunningDuty() - if !hasDuty { + // If the relevant queue hasn't been initialized yet, there probably isn't a running duty we can issue a + // timeout for, in practice this should never happen - but we need to handle this just in case. + // This is also possible if the queue got pruned already (due to becoming old and irrelevant). + q := c.Queues[phase0.Slot(height)] + if q.Q == nil { return } msg, err := c.createTimerMessage(identifier, height, round) if err != nil { - logger.Debug("❗ failed to create timer msg", zap.Error(err)) + logger.Error("❗ failed to create timer msg", zap.Error(err)) return } dec, err := queue.DecodeSSVMessage(msg) if err != nil { - logger.Debug("❌ failed to decode timer msg", zap.Error(err)) + logger.Error("❌ failed to decode timer msg", zap.Error(err)) return } - if pushed := c.Queues[phase0.Slot(height)].Q.TryPush(dec); !pushed { - logger.Warn("❗️ dropping timeout message because the queue is full", fields.RunnerRole(identifier.GetRoleType())) + if pushed := q.Q.TryPush(dec); !pushed { + logger.Error("❗️ dropping timeout message because the queue is full", fields.RunnerRole(identifier.GetRoleType())) } } } diff --git a/protocol/v2/ssv/validator/validator.go b/protocol/v2/ssv/validator/validator.go index 4be9981aa1..2afafce153 100644 --- a/protocol/v2/ssv/validator/validator.go +++ b/protocol/v2/ssv/validator/validator.go @@ -210,6 +210,11 @@ func (v *Validator) ProcessMessage(ctx context.Context, logger *zap.Logger, msg case ssvtypes.Timeout: span.AddEvent("process validator message = event(timeout)") + if !dutyRunner.HasRunningDuty() { + // Duties terminate eventually, timeout-event issuer is unaware of that - that's why we can end up here + return nil + } + timeoutData, err := eventMsg.GetTimeoutData() if err != nil { return fmt.Errorf("get event message timeout data: %w", err) From 2ad6ce7bf33464c77a4a7e7472911fe8b22f7494 Mon Sep 17 00:00:00 2001 From: iurii Date: Thu, 2 Apr 2026 18:59:09 +0300 Subject: [PATCH 02/14] fix typos --- protocol/v2/ssv/runner/runner.go | 2 +- protocol/v2/ssv/validator/timer.go | 4 ++-- 2 files changed, 3 insertions(+), 3 deletions(-) diff --git a/protocol/v2/ssv/runner/runner.go b/protocol/v2/ssv/runner/runner.go index 6fb608a93d..5d9e9fa3e8 100644 --- a/protocol/v2/ssv/runner/runner.go +++ b/protocol/v2/ssv/runner/runner.go @@ -88,7 +88,7 @@ type BaseRunner struct { // plus the caller is also responsible for ensuring there is no race with moving on to the next duty // (the baseSetupForNewDuty call). // Note, the current implementation achieves concurrent safety by making sure every State read/update - // is done by the same go-routing, handling all the messages in queue.SSVMessage (p2p messages and events) + // is done by the same go-routine, handling all the messages in queue.SSVMessage (p2p messages and events) // sequentially. State *State diff --git a/protocol/v2/ssv/validator/timer.go b/protocol/v2/ssv/validator/timer.go index 57c6910d19..c6a258febb 100644 --- a/protocol/v2/ssv/validator/timer.go +++ b/protocol/v2/ssv/validator/timer.go @@ -22,7 +22,7 @@ func (v *Validator) onTimeout(ctx context.Context, logger *zap.Logger, identifie v.mtx.RLock() // read-lock for v.Queues defer v.mtx.RUnlock() - // If the relevant queue hasn't been initialized yet, there probably isn't a running duty we can issue a + // If the relevant queue hasn't been initialized yet, there isn't a running duty we can issue a // timeout for, in practice this should never happen - but we need to handle this just in case. q := v.Queues[identifier.GetRoleType()] if q == nil { @@ -77,7 +77,7 @@ func (c *Committee) onTimeout(ctx context.Context, logger *zap.Logger, identifie c.mtx.RLock() // read-lock for c.Queues defer c.mtx.RUnlock() - // If the relevant queue hasn't been initialized yet, there probably isn't a running duty we can issue a + // If the relevant queue hasn't been initialized yet, there isn't a running duty we can issue a // timeout for, in practice this should never happen - but we need to handle this just in case. // This is also possible if the queue got pruned already (due to becoming old and irrelevant). q := c.Queues[phase0.Slot(height)] From 7460ccfe2b48be48aa45c2a7868b4f0b1786bedf Mon Sep 17 00:00:00 2001 From: iurii Date: Thu, 2 Apr 2026 21:08:08 +0300 Subject: [PATCH 03/14] add logging for better visibility --- protocol/v2/ssv/validator/committee.go | 2 ++ protocol/v2/ssv/validator/timer.go | 6 ++++-- protocol/v2/ssv/validator/validator.go | 1 + 3 files changed, 7 insertions(+), 2 deletions(-) diff --git a/protocol/v2/ssv/validator/committee.go b/protocol/v2/ssv/validator/committee.go index 8f1cfccb83..f076e262dd 100644 --- a/protocol/v2/ssv/validator/committee.go +++ b/protocol/v2/ssv/validator/committee.go @@ -326,10 +326,12 @@ func (c *Committee) ProcessMessage(ctx context.Context, logger *zap.Logger, msg c.mtx.RUnlock() if !found { // Old runners are pruned, timeout-event issuer is unaware of that - that's why we can end up here + logger.Debug("event message: timeout event arrived, but targeted runner not found (likely was pruned)") return nil } if !dutyRunner.HasRunningDuty() { // Duties terminate eventually, timeout-event issuer is unaware of that - that's why we can end up here + logger.Debug("event message: timeout event arrived after duty has finished") return nil } diff --git a/protocol/v2/ssv/validator/timer.go b/protocol/v2/ssv/validator/timer.go index c6a258febb..f49d4ff380 100644 --- a/protocol/v2/ssv/validator/timer.go +++ b/protocol/v2/ssv/validator/timer.go @@ -26,12 +26,13 @@ func (v *Validator) onTimeout(ctx context.Context, logger *zap.Logger, identifie // timeout for, in practice this should never happen - but we need to handle this just in case. q := v.Queues[identifier.GetRoleType()] if q == nil { + logger.Error("❗ couldn't schedule timeout event due to missing queue") return } msg, err := v.createTimerMessage(identifier, height, round) if err != nil { - logger.Error("❗ failed to create timer msg", zap.Error(err)) + logger.Error("❌ failed to create timer msg", zap.Error(err)) return } dec, err := queue.DecodeSSVMessage(msg) @@ -82,12 +83,13 @@ func (c *Committee) onTimeout(ctx context.Context, logger *zap.Logger, identifie // This is also possible if the queue got pruned already (due to becoming old and irrelevant). q := c.Queues[phase0.Slot(height)] if q.Q == nil { + logger.Debug("couldn't schedule timeout event due to missing queue (likely was pruned)") return } msg, err := c.createTimerMessage(identifier, height, round) if err != nil { - logger.Error("❗ failed to create timer msg", zap.Error(err)) + logger.Error("❌ failed to create timer msg", zap.Error(err)) return } dec, err := queue.DecodeSSVMessage(msg) diff --git a/protocol/v2/ssv/validator/validator.go b/protocol/v2/ssv/validator/validator.go index 2afafce153..02f8424fee 100644 --- a/protocol/v2/ssv/validator/validator.go +++ b/protocol/v2/ssv/validator/validator.go @@ -212,6 +212,7 @@ func (v *Validator) ProcessMessage(ctx context.Context, logger *zap.Logger, msg if !dutyRunner.HasRunningDuty() { // Duties terminate eventually, timeout-event issuer is unaware of that - that's why we can end up here + logger.Debug("event message: timeout event arrived after duty has finished") return nil } From fcff15c607b7b80bbaaaf2765688c797f0355176 Mon Sep 17 00:00:00 2001 From: iurii Date: Fri, 3 Apr 2026 11:35:08 +0300 Subject: [PATCH 04/14] simplify & clarify --- protocol/v2/ssv/runner/committee_test.go | 2 +- protocol/v2/ssv/runner/proposer_test.go | 2 +- protocol/v2/ssv/runner/runner.go | 29 +++++++----------------- 3 files changed, 10 insertions(+), 23 deletions(-) diff --git a/protocol/v2/ssv/runner/committee_test.go b/protocol/v2/ssv/runner/committee_test.go index fd9a1806d5..ab829d3e1e 100644 --- a/protocol/v2/ssv/runner/committee_test.go +++ b/protocol/v2/ssv/runner/committee_test.go @@ -259,7 +259,7 @@ func TestCommitteeRunnerExecuteDuty_FetchesAttestationDataAndStartsConsensus(t * env := newCommitteeRunnerEnv(t, []int{1}, &committeeDutyGuardStub{}, &doppelgangerStub{}) duty := spectestingutils.TestingAttesterDuty(spec.DataVersionElectra) - env.runner.BaseRunner.baseSetupForNewDuty(duty, env.sampleKey.Threshold) + env.runner.BaseRunner.State = NewRunnerState(env.sampleKey.Threshold, duty) require.NoError(t, env.runner.executeDuty(context.Background(), env.logger, duty)) require.NotNil(t, env.runner.ValCheck) diff --git a/protocol/v2/ssv/runner/proposer_test.go b/protocol/v2/ssv/runner/proposer_test.go index b671ddbe34..cf5f874a0c 100644 --- a/protocol/v2/ssv/runner/proposer_test.go +++ b/protocol/v2/ssv/runner/proposer_test.go @@ -431,7 +431,7 @@ func setupRunnerForPostConsensus( t.Helper() duty := spectestingutils.TestingProposerDutyV(consensusData.Version) - runner.BaseRunner.baseSetupForNewDuty(duty, keySet.Threshold) + runner.BaseRunner.State = NewRunnerState(keySet.Threshold, duty) runner.measurements.StartDutyFlow() runner.measurements.StartConsensus() runner.measurements.EndConsensus() diff --git a/protocol/v2/ssv/runner/runner.go b/protocol/v2/ssv/runner/runner.go index 5d9e9fa3e8..5eae14c39e 100644 --- a/protocol/v2/ssv/runner/runner.go +++ b/protocol/v2/ssv/runner/runner.go @@ -199,18 +199,13 @@ func (b *BaseRunner) SetHighestDecidedSlot(slot phase0.Slot) { b.highestDecidedSlot = slot } -// baseSetupForNewDuty is sets the runner for a new duty -func (b *BaseRunner) baseSetupForNewDuty(duty spectypes.Duty, quorum uint64) { - b.State = NewRunnerState(quorum, duty) -} - // baseStartNewDuty is a base func that all runner implementation can call to start a duty func (b *BaseRunner) baseStartNewDuty(ctx context.Context, logger *zap.Logger, runner Runner, duty spectypes.Duty, quorum uint64) error { if err := b.ShouldProcessDuty(duty); err != nil { return fmt.Errorf("can't start duty: %w", err) } - b.baseSetupForNewDuty(duty, quorum) + b.State = NewRunnerState(quorum, duty) if err := runner.executeDuty(ctx, logger, duty); err != nil { return fmt.Errorf("failed to execute duty: %w", err) @@ -223,7 +218,7 @@ func (b *BaseRunner) baseStartNewNonBeaconDuty(ctx context.Context, logger *zap. if err := b.ShouldProcessNonBeaconDuty(duty); err != nil { return fmt.Errorf("can't start non-beacon duty: %w", err) } - b.baseSetupForNewDuty(duty, quorum) + b.State = NewRunnerState(quorum, duty) return runner.executeDuty(ctx, logger, duty) } @@ -462,25 +457,17 @@ func (b *BaseRunner) decide( return nil } -// hasRunningDuty returns true if a new duty didn't start or an existing duty marked as finished -func (b *BaseRunner) hasRunningDuty() bool { - if b.State == nil { - return false - } - - return !b.State.Finished -} - func (b *BaseRunner) hasDutyAssigned() bool { return b.State != nil } -func (b *BaseRunner) hasDutyFinished() bool { - if b.State == nil { - return false - } +// hasRunningDuty returns true if a new duty didn't start or an existing duty marked as finished +func (b *BaseRunner) hasRunningDuty() bool { + return b.hasDutyAssigned() && !b.State.Finished +} - return b.State.Finished +func (b *BaseRunner) hasDutyFinished() bool { + return b.hasDutyAssigned() && b.State.Finished } func (b *BaseRunner) ShouldProcessDuty(duty spectypes.Duty) error { From ecca5075a016bd86ba6cacb682cd5e428aebda9f Mon Sep 17 00:00:00 2001 From: iurii Date: Fri, 3 Apr 2026 13:47:15 +0300 Subject: [PATCH 05/14] remove noisy logs --- protocol/v2/ssv/validator/committee.go | 1 - protocol/v2/ssv/validator/validator.go | 1 - 2 files changed, 2 deletions(-) diff --git a/protocol/v2/ssv/validator/committee.go b/protocol/v2/ssv/validator/committee.go index f076e262dd..be3629c9b1 100644 --- a/protocol/v2/ssv/validator/committee.go +++ b/protocol/v2/ssv/validator/committee.go @@ -331,7 +331,6 @@ func (c *Committee) ProcessMessage(ctx context.Context, logger *zap.Logger, msg } if !dutyRunner.HasRunningDuty() { // Duties terminate eventually, timeout-event issuer is unaware of that - that's why we can end up here - logger.Debug("event message: timeout event arrived after duty has finished") return nil } diff --git a/protocol/v2/ssv/validator/validator.go b/protocol/v2/ssv/validator/validator.go index 02f8424fee..2afafce153 100644 --- a/protocol/v2/ssv/validator/validator.go +++ b/protocol/v2/ssv/validator/validator.go @@ -212,7 +212,6 @@ func (v *Validator) ProcessMessage(ctx context.Context, logger *zap.Logger, msg if !dutyRunner.HasRunningDuty() { // Duties terminate eventually, timeout-event issuer is unaware of that - that's why we can end up here - logger.Debug("event message: timeout event arrived after duty has finished") return nil } From 001e2c3d1a9a2e3334e50767df9f04b0fd4730fb Mon Sep 17 00:00:00 2001 From: iurii Date: Fri, 3 Apr 2026 18:58:40 +0300 Subject: [PATCH 06/14] skip timeout events targeting past duties --- protocol/v2/ssv/runner/aggregator.go | 4 ++++ protocol/v2/ssv/runner/committee.go | 4 ++++ protocol/v2/ssv/runner/proposer.go | 4 ++++ protocol/v2/ssv/runner/runner.go | 13 +++++++++++-- protocol/v2/ssv/runner/runner_state.go | 17 +++++++---------- .../ssv/runner/sync_committee_contribution.go | 4 ++++ .../v2/ssv/runner/validator_registration.go | 6 +++++- protocol/v2/ssv/runner/voluntary_exit.go | 4 ++++ protocol/v2/ssv/validator/validator.go | 11 ++++++++++- 9 files changed, 53 insertions(+), 14 deletions(-) diff --git a/protocol/v2/ssv/runner/aggregator.go b/protocol/v2/ssv/runner/aggregator.go index 7c164b2be3..4d97b4c723 100644 --- a/protocol/v2/ssv/runner/aggregator.go +++ b/protocol/v2/ssv/runner/aggregator.go @@ -492,6 +492,10 @@ func (r *AggregatorRunner) GetRole() spectypes.RunnerRole { return r.BaseRunner.GetRole() } +func (r *AggregatorRunner) GetCurrentDutySlot() (phase0.Slot, bool) { + return r.BaseRunner.GetCurrentDutySlot() +} + func (r *AggregatorRunner) GetLastHeight() specqbft.Height { return r.BaseRunner.GetLastHeight() } diff --git a/protocol/v2/ssv/runner/committee.go b/protocol/v2/ssv/runner/committee.go index 5961427f99..14a1da0282 100644 --- a/protocol/v2/ssv/runner/committee.go +++ b/protocol/v2/ssv/runner/committee.go @@ -212,6 +212,10 @@ func (r *CommitteeRunner) GetRole() spectypes.RunnerRole { return r.BaseRunner.GetRole() } +func (r *CommitteeRunner) GetCurrentDutySlot() (phase0.Slot, bool) { + return r.BaseRunner.GetCurrentDutySlot() +} + func (r *CommitteeRunner) GetLastHeight() specqbft.Height { return r.BaseRunner.GetLastHeight() } diff --git a/protocol/v2/ssv/runner/proposer.go b/protocol/v2/ssv/runner/proposer.go index afb006fdc5..bdf4224a4a 100644 --- a/protocol/v2/ssv/runner/proposer.go +++ b/protocol/v2/ssv/runner/proposer.go @@ -547,6 +547,10 @@ func (r *ProposerRunner) GetRole() spectypes.RunnerRole { return r.BaseRunner.GetRole() } +func (r *ProposerRunner) GetCurrentDutySlot() (phase0.Slot, bool) { + return r.BaseRunner.GetCurrentDutySlot() +} + func (r *ProposerRunner) GetLastHeight() specqbft.Height { return r.BaseRunner.GetLastHeight() } diff --git a/protocol/v2/ssv/runner/runner.go b/protocol/v2/ssv/runner/runner.go index 5eae14c39e..9e320d08f1 100644 --- a/protocol/v2/ssv/runner/runner.go +++ b/protocol/v2/ssv/runner/runner.go @@ -31,6 +31,7 @@ type Getters interface { HasAcceptedProposalForCurrentRound() bool GetShares() map[phase0.ValidatorIndex]*spectypes.Share GetRole() spectypes.RunnerRole + GetCurrentDutySlot() (phase0.Slot, bool) GetLastHeight() specqbft.Height GetLastRound() specqbft.Round GetStateRoot() ([32]byte, error) @@ -136,6 +137,14 @@ func (b *BaseRunner) GetRole() spectypes.RunnerRole { return b.RunnerRoleType } +func (b *BaseRunner) GetCurrentDutySlot() (phase0.Slot, bool) { + if !b.hasDutyAssigned() { + return 0, false + } + + return b.State.CurrentDuty.DutySlot(), true +} + func (b *BaseRunner) GetLastHeight() specqbft.Height { if ctrl := b.QBFTController; ctrl != nil { return ctrl.Height @@ -260,7 +269,7 @@ func (b *BaseRunner) baseConsensusMsgProcessing(ctx context.Context, logger *zap span := trace.SpanFromContext(ctx) prevDecided := false - if b.hasRunningDuty() && b.State != nil && b.State.RunningInstance != nil { + if b.hasRunningDuty() && b.hasDutyAssigned() && b.State.RunningInstance != nil { prevDecided, _ = b.State.RunningInstance.IsDecided() } if prevDecided { @@ -482,7 +491,7 @@ func (b *BaseRunner) ShouldProcessDuty(duty spectypes.Duty) error { func (b *BaseRunner) ShouldProcessNonBeaconDuty(duty spectypes.Duty) error { // assume CurrentDuty is not nil if state is not nil - if b.State != nil && b.State.CurrentDuty.DutySlot() >= duty.DutySlot() { + if b.hasDutyAssigned() && b.State.CurrentDuty.DutySlot() >= duty.DutySlot() { return spectypes.NewError( spectypes.DutyAlreadyPassedErrorCode, fmt.Sprintf("duty for slot %d already passed. Current slot is %d", duty.DutySlot(), b.State.CurrentDuty.DutySlot()), diff --git a/protocol/v2/ssv/runner/runner_state.go b/protocol/v2/ssv/runner/runner_state.go index dea49b898f..8235583e71 100644 --- a/protocol/v2/ssv/runner/runner_state.go +++ b/protocol/v2/ssv/runner/runner_state.go @@ -84,18 +84,15 @@ func (pcs *State) MarshalJSON() ([]byte, error) { Finished: pcs.Finished, } - if pcs.CurrentDuty != nil { - if ValidatorDuty, ok := pcs.CurrentDuty.(*spectypes.ValidatorDuty); ok { - alias.ValidatorDuty = ValidatorDuty - } else if committeeDuty, ok := pcs.CurrentDuty.(*spectypes.CommitteeDuty); ok { - alias.CommitteeDuty = committeeDuty - } else { - return nil, errors.New("can't marshal because BaseRunner.State.CurrentDuty isn't ValidatorDuty or CommitteeDuty") - } + if ValidatorDuty, ok := pcs.CurrentDuty.(*spectypes.ValidatorDuty); ok { + alias.ValidatorDuty = ValidatorDuty + } else if committeeDuty, ok := pcs.CurrentDuty.(*spectypes.CommitteeDuty); ok { + alias.CommitteeDuty = committeeDuty + } else { + return nil, errors.New("can't marshal because BaseRunner.State.CurrentDuty isn't ValidatorDuty or CommitteeDuty") } - byts, err := json.Marshal(alias) - return byts, err + return json.Marshal(alias) } func (pcs *State) UnmarshalJSON(data []byte) error { diff --git a/protocol/v2/ssv/runner/sync_committee_contribution.go b/protocol/v2/ssv/runner/sync_committee_contribution.go index fc364a1a44..a26dbf8cf7 100644 --- a/protocol/v2/ssv/runner/sync_committee_contribution.go +++ b/protocol/v2/ssv/runner/sync_committee_contribution.go @@ -564,6 +564,10 @@ func (r *SyncCommitteeAggregatorRunner) GetRole() spectypes.RunnerRole { return r.BaseRunner.GetRole() } +func (r *SyncCommitteeAggregatorRunner) GetCurrentDutySlot() (phase0.Slot, bool) { + return r.BaseRunner.GetCurrentDutySlot() +} + func (r *SyncCommitteeAggregatorRunner) GetLastHeight() specqbft.Height { return r.BaseRunner.GetLastHeight() } diff --git a/protocol/v2/ssv/runner/validator_registration.go b/protocol/v2/ssv/runner/validator_registration.go index c30cd74fb5..2fcf4134a8 100644 --- a/protocol/v2/ssv/runner/validator_registration.go +++ b/protocol/v2/ssv/runner/validator_registration.go @@ -168,7 +168,7 @@ func (r *ValidatorRegistrationRunner) ProcessPostConsensus(ctx context.Context, } func (r *ValidatorRegistrationRunner) expectedPreConsensusRootsAndDomain() ([]ssz.HashRoot, phase0.DomainType, error) { - if r.BaseRunner.State == nil || r.BaseRunner.State.CurrentDuty == nil { + if !r.BaseRunner.hasDutyAssigned() { return nil, spectypes.DomainError, fmt.Errorf("no running duty to compute preconsensus roots and domain") } vr, err := r.buildValidatorRegistration(r.BaseRunner.State.CurrentDuty.DutySlot()) @@ -285,6 +285,10 @@ func (r *ValidatorRegistrationRunner) GetRole() spectypes.RunnerRole { return r.BaseRunner.GetRole() } +func (r *ValidatorRegistrationRunner) GetCurrentDutySlot() (phase0.Slot, bool) { + return r.BaseRunner.GetCurrentDutySlot() +} + func (r *ValidatorRegistrationRunner) GetLastHeight() specqbft.Height { return r.BaseRunner.GetLastHeight() } diff --git a/protocol/v2/ssv/runner/voluntary_exit.go b/protocol/v2/ssv/runner/voluntary_exit.go index 5505eef98e..cfd22f7d43 100644 --- a/protocol/v2/ssv/runner/voluntary_exit.go +++ b/protocol/v2/ssv/runner/voluntary_exit.go @@ -254,6 +254,10 @@ func (r *VoluntaryExitRunner) GetRole() spectypes.RunnerRole { return r.BaseRunner.GetRole() } +func (r *VoluntaryExitRunner) GetCurrentDutySlot() (phase0.Slot, bool) { + return r.BaseRunner.GetCurrentDutySlot() +} + func (r *VoluntaryExitRunner) GetLastHeight() specqbft.Height { return r.BaseRunner.GetLastHeight() } diff --git a/protocol/v2/ssv/validator/validator.go b/protocol/v2/ssv/validator/validator.go index 2afafce153..313d7a3ca4 100644 --- a/protocol/v2/ssv/validator/validator.go +++ b/protocol/v2/ssv/validator/validator.go @@ -217,7 +217,16 @@ func (v *Validator) ProcessMessage(ctx context.Context, logger *zap.Logger, msg timeoutData, err := eventMsg.GetTimeoutData() if err != nil { - return fmt.Errorf("get event message timeout data: %w", err) + return fmt.Errorf("event message: get timeout data: %w", err) + } + currentDutySlot, ok := dutyRunner.GetCurrentDutySlot() + if !ok { + return fmt.Errorf("event message: get current duty slot to compare vs timeout event slot: %w", err) + } + if timeoutData.Height != specqbft.Height(currentDutySlot) { + // Timeout events can be delayed in the queue until the runner already moved on to a new duty, we can + // safely skip these + return nil } if err := dutyRunner.OnTimeoutQBFT(ctx, logger, timeoutData); err != nil { From 8baf1ce47391276dc73035d330ea73347c7d214d Mon Sep 17 00:00:00 2001 From: iurii Date: Fri, 3 Apr 2026 19:48:41 +0300 Subject: [PATCH 07/14] standardize state/qbft-instance checks --- protocol/v2/ssv/runner/aggregator.go | 2 +- protocol/v2/ssv/runner/committee.go | 2 +- protocol/v2/ssv/runner/proposer.go | 2 +- protocol/v2/ssv/runner/runner.go | 28 ++++++++----------- protocol/v2/ssv/runner/runner_validations.go | 2 +- .../ssv/runner/sync_committee_contribution.go | 2 +- .../v2/ssv/runner/validator_registration.go | 2 +- protocol/v2/ssv/runner/voluntary_exit.go | 2 +- .../multi_start_new_runner_duty_type.go | 8 +++--- protocol/v2/ssv/spectest/ssv_mapping_test.go | 8 ++---- protocol/v2/ssv/spectest/util.go | 10 +++---- 11 files changed, 31 insertions(+), 37 deletions(-) diff --git a/protocol/v2/ssv/runner/aggregator.go b/protocol/v2/ssv/runner/aggregator.go index 4d97b4c723..46efa71a53 100644 --- a/protocol/v2/ssv/runner/aggregator.go +++ b/protocol/v2/ssv/runner/aggregator.go @@ -101,7 +101,7 @@ func (r *AggregatorRunner) StartNewDuty(ctx context.Context, logger *zap.Logger, // HasRunningDuty returns true if a duty is already running (StartNewDuty called and returned nil) func (r *AggregatorRunner) HasRunningDuty() bool { - return r.BaseRunner.hasRunningDuty() + return r.BaseRunner.hasDutyRunning() } func (r *AggregatorRunner) ProcessPreConsensus(ctx context.Context, logger *zap.Logger, signedMsg *spectypes.PartialSignatureMessages) error { diff --git a/protocol/v2/ssv/runner/committee.go b/protocol/v2/ssv/runner/committee.go index 14a1da0282..714bd40d4f 100644 --- a/protocol/v2/ssv/runner/committee.go +++ b/protocol/v2/ssv/runner/committee.go @@ -253,7 +253,7 @@ func (r *CommitteeRunner) state() *State { } func (r *CommitteeRunner) HasRunningDuty() bool { - return r.BaseRunner.hasRunningDuty() + return r.BaseRunner.hasDutyRunning() } func (r *CommitteeRunner) ProcessPreConsensus(ctx context.Context, logger *zap.Logger, signedMsg *spectypes.PartialSignatureMessages) error { diff --git a/protocol/v2/ssv/runner/proposer.go b/protocol/v2/ssv/runner/proposer.go index bdf4224a4a..ceb864ae35 100644 --- a/protocol/v2/ssv/runner/proposer.go +++ b/protocol/v2/ssv/runner/proposer.go @@ -108,7 +108,7 @@ func (r *ProposerRunner) StartNewDuty(ctx context.Context, logger *zap.Logger, d // HasRunningDuty returns true if a duty is already running (StartNewDuty called and returned nil) func (r *ProposerRunner) HasRunningDuty() bool { - return r.BaseRunner.hasRunningDuty() + return r.BaseRunner.hasDutyRunning() } func (r *ProposerRunner) ProcessPreConsensus(ctx context.Context, logger *zap.Logger, signedMsg *spectypes.PartialSignatureMessages) error { diff --git a/protocol/v2/ssv/runner/runner.go b/protocol/v2/ssv/runner/runner.go index 9e320d08f1..98872909c5 100644 --- a/protocol/v2/ssv/runner/runner.go +++ b/protocol/v2/ssv/runner/runner.go @@ -106,21 +106,18 @@ type BaseRunner struct { highestDecidedSlot phase0.Slot } +func (b *BaseRunner) HasStartedQBFTInstance() bool { + return b.hasDutyAssigned() && b.State.RunningInstance != nil +} + func (b *BaseRunner) HasRunningQBFTInstance() bool { - var runningInstance *instance.Instance - if b.hasRunningDuty() { - runningInstance = b.State.RunningInstance - if runningInstance != nil { - decided, _ := runningInstance.IsDecided() - return !decided - } - } - return false + // Note: RunningInstance.State cannot be nil for existing RunningInstance by construction. + return b.HasStartedQBFTInstance() && !b.State.RunningInstance.State.Decided } func (b *BaseRunner) HasAcceptedProposalForCurrentRound() bool { var runningInstance *instance.Instance - if b.hasRunningDuty() { + if b.hasDutyRunning() { runningInstance = b.State.RunningInstance if runningInstance != nil { return runningInstance.State.ProposalAcceptedForCurrentRound != nil @@ -153,7 +150,7 @@ func (b *BaseRunner) GetLastHeight() specqbft.Height { } func (b *BaseRunner) GetLastRound() specqbft.Round { - if b.hasRunningDuty() { + if b.hasDutyRunning() { inst := b.State.RunningInstance if inst != nil { return inst.State.Round @@ -269,7 +266,7 @@ func (b *BaseRunner) baseConsensusMsgProcessing(ctx context.Context, logger *zap span := trace.SpanFromContext(ctx) prevDecided := false - if b.hasRunningDuty() && b.hasDutyAssigned() && b.State.RunningInstance != nil { + if b.hasDutyRunning() && b.hasDutyAssigned() && b.HasStartedQBFTInstance() { prevDecided, _ = b.State.RunningInstance.IsDecided() } if prevDecided { @@ -284,7 +281,7 @@ func (b *BaseRunner) baseConsensusMsgProcessing(ctx context.Context, logger *zap return false, nil, err } - if !b.hasRunningDuty() { + if !b.hasDutyRunning() { logger.Debug("no running duty, applied consensus message but cannot progress further") return false, nil, nil } @@ -405,7 +402,7 @@ func (b *BaseRunner) didDecideCorrectly(prevDecided bool, signedMessage *spectyp return false, nil } - if b.State.RunningInstance == nil { + if !b.HasStartedQBFTInstance() { return false, spectypes.NewError(spectypes.DecidedWrongInstanceErrorCode, "decided wrong instance (running instance is nil)") } @@ -470,8 +467,7 @@ func (b *BaseRunner) hasDutyAssigned() bool { return b.State != nil } -// hasRunningDuty returns true if a new duty didn't start or an existing duty marked as finished -func (b *BaseRunner) hasRunningDuty() bool { +func (b *BaseRunner) hasDutyRunning() bool { return b.hasDutyAssigned() && !b.State.Finished } diff --git a/protocol/v2/ssv/runner/runner_validations.go b/protocol/v2/ssv/runner/runner_validations.go index 9968531b75..c54b54267a 100644 --- a/protocol/v2/ssv/runner/runner_validations.go +++ b/protocol/v2/ssv/runner/runner_validations.go @@ -90,7 +90,7 @@ func (b *BaseRunner) ValidatePostConsensusMsg(ctx context.Context, runner Runner return err } - if b.State.RunningInstance == nil { + if !b.HasStartedQBFTInstance() { return NewRetryableError(spectypes.WrapError(spectypes.NoRunningConsensusInstanceErrorCode, ErrInstanceNotFound)) } diff --git a/protocol/v2/ssv/runner/sync_committee_contribution.go b/protocol/v2/ssv/runner/sync_committee_contribution.go index a26dbf8cf7..409e750608 100644 --- a/protocol/v2/ssv/runner/sync_committee_contribution.go +++ b/protocol/v2/ssv/runner/sync_committee_contribution.go @@ -83,7 +83,7 @@ func (r *SyncCommitteeAggregatorRunner) StartNewDuty(ctx context.Context, logger // HasRunningDuty returns true if a duty is already running (StartNewDuty called and returned nil) func (r *SyncCommitteeAggregatorRunner) HasRunningDuty() bool { - return r.BaseRunner.hasRunningDuty() + return r.BaseRunner.hasDutyRunning() } func (r *SyncCommitteeAggregatorRunner) ProcessPreConsensus(ctx context.Context, logger *zap.Logger, signedMsg *spectypes.PartialSignatureMessages) error { diff --git a/protocol/v2/ssv/runner/validator_registration.go b/protocol/v2/ssv/runner/validator_registration.go index 2fcf4134a8..aa1abb8017 100644 --- a/protocol/v2/ssv/runner/validator_registration.go +++ b/protocol/v2/ssv/runner/validator_registration.go @@ -87,7 +87,7 @@ func (r *ValidatorRegistrationRunner) StartNewDuty(ctx context.Context, logger * // HasRunningDuty returns true if a duty is already running (StartNewDuty called and returned nil) func (r *ValidatorRegistrationRunner) HasRunningDuty() bool { - return r.BaseRunner.hasRunningDuty() + return r.BaseRunner.hasDutyRunning() } func (r *ValidatorRegistrationRunner) ProcessPreConsensus(ctx context.Context, logger *zap.Logger, signedMsg *spectypes.PartialSignatureMessages) error { diff --git a/protocol/v2/ssv/runner/voluntary_exit.go b/protocol/v2/ssv/runner/voluntary_exit.go index cfd22f7d43..b951d05f67 100644 --- a/protocol/v2/ssv/runner/voluntary_exit.go +++ b/protocol/v2/ssv/runner/voluntary_exit.go @@ -70,7 +70,7 @@ func (r *VoluntaryExitRunner) StartNewDuty(ctx context.Context, logger *zap.Logg // HasRunningDuty returns true if a duty is already running (StartNewDuty called and returned nil) func (r *VoluntaryExitRunner) HasRunningDuty() bool { - return r.BaseRunner.hasRunningDuty() + return r.BaseRunner.hasDutyRunning() } // ProcessPreConsensus Check for quorum of partial signatures over VoluntaryExit and, diff --git a/protocol/v2/ssv/spectest/multi_start_new_runner_duty_type.go b/protocol/v2/ssv/spectest/multi_start_new_runner_duty_type.go index 5de5257cbe..88108b7341 100644 --- a/protocol/v2/ssv/spectest/multi_start_new_runner_duty_type.go +++ b/protocol/v2/ssv/spectest/multi_start_new_runner_duty_type.go @@ -96,28 +96,28 @@ func (test *StartNewRunnerDutySpecTest) RunAsPartOfMultiTest(t *testing.T, logge for _, inst := range r.BaseRunner.QBFTController.StoredInstances { inst.ValueChecker = protocoltesting.TestingValueChecker{} } - if r.BaseRunner.State.RunningInstance != nil { + if r.BaseRunner.HasStartedQBFTInstance() { r.BaseRunner.State.RunningInstance.ValueChecker = protocoltesting.TestingValueChecker{} } case *runner.AggregatorRunner: for _, inst := range r.BaseRunner.QBFTController.StoredInstances { inst.ValueChecker = protocoltesting.TestingValueChecker{} } - if r.BaseRunner.State.RunningInstance != nil { + if r.BaseRunner.HasStartedQBFTInstance() { r.BaseRunner.State.RunningInstance.ValueChecker = protocoltesting.TestingValueChecker{} } case *runner.ProposerRunner: for _, inst := range r.BaseRunner.QBFTController.StoredInstances { inst.ValueChecker = protocoltesting.TestingValueChecker{} } - if r.BaseRunner.State.RunningInstance != nil { + if r.BaseRunner.HasStartedQBFTInstance() { r.BaseRunner.State.RunningInstance.ValueChecker = protocoltesting.TestingValueChecker{} } case *runner.SyncCommitteeAggregatorRunner: for _, inst := range r.BaseRunner.QBFTController.StoredInstances { inst.ValueChecker = protocoltesting.TestingValueChecker{} } - if r.BaseRunner.State.RunningInstance != nil { + if r.BaseRunner.HasStartedQBFTInstance() { r.BaseRunner.State.RunningInstance.ValueChecker = protocoltesting.TestingValueChecker{} } } diff --git a/protocol/v2/ssv/spectest/ssv_mapping_test.go b/protocol/v2/ssv/spectest/ssv_mapping_test.go index 2b51755d2e..de145697ab 100644 --- a/protocol/v2/ssv/spectest/ssv_mapping_test.go +++ b/protocol/v2/ssv/spectest/ssv_mapping_test.go @@ -389,11 +389,9 @@ func fixRunnerForRun(t *testing.T, runnerMap map[string]any, ks *spectestingutil if baseRunner.QBFTController != nil { baseRunner.QBFTController = fixControllerForRun(logger, baseRunner.QBFTController, ks) - if baseRunner.State != nil { - if baseRunner.State.RunningInstance != nil { - operator := spectestingutils.TestingCommitteeMember(ks) - baseRunner.State.RunningInstance = fixInstanceForRun(logger, ks, baseRunner.State.RunningInstance, baseRunner.QBFTController, operator) - } + if baseRunner.HasStartedQBFTInstance() { + operator := spectestingutils.TestingCommitteeMember(ks) + baseRunner.State.RunningInstance = fixInstanceForRun(logger, ks, baseRunner.State.RunningInstance, baseRunner.QBFTController, operator) } } diff --git a/protocol/v2/ssv/spectest/util.go b/protocol/v2/ssv/spectest/util.go index a7927f13f7..60778a0ba3 100644 --- a/protocol/v2/ssv/spectest/util.go +++ b/protocol/v2/ssv/spectest/util.go @@ -49,7 +49,7 @@ func runnerForTest(t *testing.T, runnerType runner.Runner, name string, testType for _, inst := range cr.BaseRunner.QBFTController.StoredInstances { inst.ValueChecker = valCheck } - if cr.BaseRunner.State != nil && cr.BaseRunner.State.RunningInstance != nil { + if cr.BaseRunner.HasStartedQBFTInstance() { cr.BaseRunner.State.RunningInstance.ValueChecker = valCheck } case *runner.AggregatorRunner: @@ -60,7 +60,7 @@ func runnerForTest(t *testing.T, runnerType runner.Runner, name string, testType for _, inst := range ar.BaseRunner.QBFTController.StoredInstances { inst.ValueChecker = valCheck } - if ar.BaseRunner.State != nil && ar.BaseRunner.State.RunningInstance != nil { + if ar.BaseRunner.HasStartedQBFTInstance() { ar.BaseRunner.State.RunningInstance.ValueChecker = valCheck } case *runner.ProposerRunner: @@ -71,7 +71,7 @@ func runnerForTest(t *testing.T, runnerType runner.Runner, name string, testType for _, inst := range pr.BaseRunner.QBFTController.StoredInstances { inst.ValueChecker = valCheck } - if pr.BaseRunner.State != nil && pr.BaseRunner.State.RunningInstance != nil { + if pr.BaseRunner.HasStartedQBFTInstance() { pr.BaseRunner.State.RunningInstance.ValueChecker = valCheck } case *runner.SyncCommitteeAggregatorRunner: @@ -82,7 +82,7 @@ func runnerForTest(t *testing.T, runnerType runner.Runner, name string, testType for _, inst := range scr.BaseRunner.QBFTController.StoredInstances { inst.ValueChecker = valCheck } - if scr.BaseRunner.State != nil && scr.BaseRunner.State.RunningInstance != nil { + if scr.BaseRunner.HasStartedQBFTInstance() { scr.BaseRunner.State.RunningInstance.ValueChecker = valCheck } case *runner.ValidatorRegistrationRunner: @@ -102,7 +102,7 @@ func normalizeExpectedProposerStartValues(pr *runner.ProposerRunner) { } if state := pr.BaseRunner.State; state != nil { state.DecidedValue = normalizeProposerConsensusValue(state.DecidedValue) - if state.RunningInstance != nil { + if pr.BaseRunner.HasStartedQBFTInstance() { state.RunningInstance.StartValue = normalizeProposerConsensusValue(state.RunningInstance.StartValue) if state.RunningInstance.State != nil { state.RunningInstance.State.LastPreparedValue = normalizeProposerConsensusValue(state.RunningInstance.State.LastPreparedValue) From 108219159abab9dd3919fe6f5b8f93d86d49b074 Mon Sep 17 00:00:00 2001 From: iurii Date: Fri, 3 Apr 2026 20:00:10 +0300 Subject: [PATCH 08/14] remove redundant condition in baseConsensusMsgProcessing --- protocol/v2/ssv/runner/runner.go | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/protocol/v2/ssv/runner/runner.go b/protocol/v2/ssv/runner/runner.go index 98872909c5..9be63542af 100644 --- a/protocol/v2/ssv/runner/runner.go +++ b/protocol/v2/ssv/runner/runner.go @@ -266,7 +266,7 @@ func (b *BaseRunner) baseConsensusMsgProcessing(ctx context.Context, logger *zap span := trace.SpanFromContext(ctx) prevDecided := false - if b.hasDutyRunning() && b.hasDutyAssigned() && b.HasStartedQBFTInstance() { + if b.hasDutyRunning() && b.HasStartedQBFTInstance() { prevDecided, _ = b.State.RunningInstance.IsDecided() } if prevDecided { From 48a5d79912e83ee57a7a155876a9f102ed4e5e81 Mon Sep 17 00:00:00 2001 From: iurii Date: Fri, 3 Apr 2026 20:23:05 +0300 Subject: [PATCH 09/14] adjust unit-tests to accomodate HasRunningQBFTInstance behavior change --- .../v2/ssv/validator/committee_queue_test.go | 97 +++++++++++-------- 1 file changed, 57 insertions(+), 40 deletions(-) diff --git a/protocol/v2/ssv/validator/committee_queue_test.go b/protocol/v2/ssv/validator/committee_queue_test.go index 4668012ff4..9982f39bcd 100644 --- a/protocol/v2/ssv/validator/committee_queue_test.go +++ b/protocol/v2/ssv/validator/committee_queue_test.go @@ -691,44 +691,58 @@ func TestChangingFilterState(t *testing.T) { // - The count of processed messages matches the expected number for that scenario. func TestCommitteeQueueFilteringScenarios(t *testing.T) { testCases := []struct { - name string - hasRunningDuty bool - decided bool - proposalAccepted bool - messagesTypes []specqbft.MessageType - expectedProcessed []bool + name string + hasRunningDuty bool + hasRunningInstance bool + decided bool + proposalAccepted bool + messagesTypes []specqbft.MessageType + expectedProcessed []bool }{ { - name: "no active duty", - hasRunningDuty: false, - decided: false, - proposalAccepted: false, - messagesTypes: []specqbft.MessageType{specqbft.ProposalMsgType, specqbft.PrepareMsgType, specqbft.CommitMsgType}, - expectedProcessed: []bool{true, true, true}, // All messages processed with queue.FilterAny when no active duty + name: "no active duty, no running instance", + hasRunningDuty: false, + hasRunningInstance: false, + decided: false, + proposalAccepted: false, + messagesTypes: []specqbft.MessageType{specqbft.ProposalMsgType, specqbft.PrepareMsgType, specqbft.CommitMsgType}, + expectedProcessed: []bool{true, true, true}, // All messages processed with queue.FilterAny when no active duty }, { - name: "no proposal accepted", - hasRunningDuty: true, - decided: false, - proposalAccepted: false, - messagesTypes: []specqbft.MessageType{specqbft.ProposalMsgType, specqbft.PrepareMsgType, specqbft.CommitMsgType, specqbft.RoundChangeMsgType}, - expectedProcessed: []bool{true, false, false, true}, // Proposal and RoundChange should be processed + name: "no active duty, lingering undecided instance", + hasRunningDuty: false, + hasRunningInstance: true, + decided: false, + proposalAccepted: false, + messagesTypes: []specqbft.MessageType{specqbft.ProposalMsgType, specqbft.PrepareMsgType, specqbft.CommitMsgType}, + expectedProcessed: []bool{true, false, false}, // Lingering undecided instance still triggers the consensus filter }, { - name: "proposal accepted but not decided", - hasRunningDuty: true, - decided: false, - proposalAccepted: true, - messagesTypes: []specqbft.MessageType{specqbft.ProposalMsgType, specqbft.PrepareMsgType, specqbft.CommitMsgType}, - expectedProcessed: []bool{true, true, true}, // All consensus messages should be processed + name: "no proposal accepted", + hasRunningDuty: true, + hasRunningInstance: true, + decided: false, + proposalAccepted: false, + messagesTypes: []specqbft.MessageType{specqbft.ProposalMsgType, specqbft.PrepareMsgType, specqbft.CommitMsgType, specqbft.RoundChangeMsgType}, + expectedProcessed: []bool{true, false, false, true}, // Proposal and RoundChange should be processed }, { - name: "decided", - hasRunningDuty: true, - decided: true, - proposalAccepted: true, - messagesTypes: []specqbft.MessageType{specqbft.ProposalMsgType, specqbft.PrepareMsgType, specqbft.CommitMsgType}, - expectedProcessed: []bool{true, true, true}, // All should be processed + name: "proposal accepted but not decided", + hasRunningDuty: true, + hasRunningInstance: true, + decided: false, + proposalAccepted: true, + messagesTypes: []specqbft.MessageType{specqbft.ProposalMsgType, specqbft.PrepareMsgType, specqbft.CommitMsgType}, + expectedProcessed: []bool{true, true, true}, // All consensus messages should be processed + }, + { + name: "decided", + hasRunningDuty: true, + hasRunningInstance: true, + decided: true, + proposalAccepted: true, + messagesTypes: []specqbft.MessageType{specqbft.ProposalMsgType, specqbft.PrepareMsgType, specqbft.CommitMsgType}, + expectedProcessed: []bool{true, true, true}, // All should be processed }, } @@ -750,7 +764,7 @@ func TestCommitteeQueueFilteringScenarios(t *testing.T) { q := queueContainer{ Q: queue.New(logger, 10), queueState: &queue.State{ - HasRunningInstance: tc.hasRunningDuty, + HasRunningInstance: tc.hasRunningInstance, Height: specqbft.Height(slot), Slot: slot, Round: 1, @@ -764,17 +778,20 @@ func TestCommitteeQueueFilteringScenarios(t *testing.T) { } } + state := &runner.State{} + if tc.hasRunningInstance { + state.RunningInstance = &instance.Instance{ + State: &specqbft.State{ + Decided: tc.decided, + ProposalAcceptedForCurrentRound: proposalMsg, + Round: 1, + }, + } + } + committeeRunner := &runner.CommitteeRunner{ BaseRunner: &runner.BaseRunner{ - State: &runner.State{ - RunningInstance: &instance.Instance{ - State: &specqbft.State{ - Decided: tc.decided, - ProposalAcceptedForCurrentRound: proposalMsg, - Round: 1, - }, - }, - }, + State: state, }, } From 6954d336602b9a370fe9240f6c9c9be5d3417fe2 Mon Sep 17 00:00:00 2001 From: iurii Date: Fri, 3 Apr 2026 20:25:06 +0300 Subject: [PATCH 10/14] preserve old behavior for HasRunningQBFTInstance func --- protocol/v2/ssv/runner/runner.go | 2 +- .../v2/ssv/validator/committee_queue_test.go | 97 ++++++++----------- 2 files changed, 41 insertions(+), 58 deletions(-) diff --git a/protocol/v2/ssv/runner/runner.go b/protocol/v2/ssv/runner/runner.go index 9be63542af..c958d2fe6b 100644 --- a/protocol/v2/ssv/runner/runner.go +++ b/protocol/v2/ssv/runner/runner.go @@ -112,7 +112,7 @@ func (b *BaseRunner) HasStartedQBFTInstance() bool { func (b *BaseRunner) HasRunningQBFTInstance() bool { // Note: RunningInstance.State cannot be nil for existing RunningInstance by construction. - return b.HasStartedQBFTInstance() && !b.State.RunningInstance.State.Decided + return b.hasDutyRunning() && b.HasStartedQBFTInstance() && !b.State.RunningInstance.State.Decided } func (b *BaseRunner) HasAcceptedProposalForCurrentRound() bool { diff --git a/protocol/v2/ssv/validator/committee_queue_test.go b/protocol/v2/ssv/validator/committee_queue_test.go index 9982f39bcd..4668012ff4 100644 --- a/protocol/v2/ssv/validator/committee_queue_test.go +++ b/protocol/v2/ssv/validator/committee_queue_test.go @@ -691,58 +691,44 @@ func TestChangingFilterState(t *testing.T) { // - The count of processed messages matches the expected number for that scenario. func TestCommitteeQueueFilteringScenarios(t *testing.T) { testCases := []struct { - name string - hasRunningDuty bool - hasRunningInstance bool - decided bool - proposalAccepted bool - messagesTypes []specqbft.MessageType - expectedProcessed []bool + name string + hasRunningDuty bool + decided bool + proposalAccepted bool + messagesTypes []specqbft.MessageType + expectedProcessed []bool }{ { - name: "no active duty, no running instance", - hasRunningDuty: false, - hasRunningInstance: false, - decided: false, - proposalAccepted: false, - messagesTypes: []specqbft.MessageType{specqbft.ProposalMsgType, specqbft.PrepareMsgType, specqbft.CommitMsgType}, - expectedProcessed: []bool{true, true, true}, // All messages processed with queue.FilterAny when no active duty + name: "no active duty", + hasRunningDuty: false, + decided: false, + proposalAccepted: false, + messagesTypes: []specqbft.MessageType{specqbft.ProposalMsgType, specqbft.PrepareMsgType, specqbft.CommitMsgType}, + expectedProcessed: []bool{true, true, true}, // All messages processed with queue.FilterAny when no active duty }, { - name: "no active duty, lingering undecided instance", - hasRunningDuty: false, - hasRunningInstance: true, - decided: false, - proposalAccepted: false, - messagesTypes: []specqbft.MessageType{specqbft.ProposalMsgType, specqbft.PrepareMsgType, specqbft.CommitMsgType}, - expectedProcessed: []bool{true, false, false}, // Lingering undecided instance still triggers the consensus filter + name: "no proposal accepted", + hasRunningDuty: true, + decided: false, + proposalAccepted: false, + messagesTypes: []specqbft.MessageType{specqbft.ProposalMsgType, specqbft.PrepareMsgType, specqbft.CommitMsgType, specqbft.RoundChangeMsgType}, + expectedProcessed: []bool{true, false, false, true}, // Proposal and RoundChange should be processed }, { - name: "no proposal accepted", - hasRunningDuty: true, - hasRunningInstance: true, - decided: false, - proposalAccepted: false, - messagesTypes: []specqbft.MessageType{specqbft.ProposalMsgType, specqbft.PrepareMsgType, specqbft.CommitMsgType, specqbft.RoundChangeMsgType}, - expectedProcessed: []bool{true, false, false, true}, // Proposal and RoundChange should be processed + name: "proposal accepted but not decided", + hasRunningDuty: true, + decided: false, + proposalAccepted: true, + messagesTypes: []specqbft.MessageType{specqbft.ProposalMsgType, specqbft.PrepareMsgType, specqbft.CommitMsgType}, + expectedProcessed: []bool{true, true, true}, // All consensus messages should be processed }, { - name: "proposal accepted but not decided", - hasRunningDuty: true, - hasRunningInstance: true, - decided: false, - proposalAccepted: true, - messagesTypes: []specqbft.MessageType{specqbft.ProposalMsgType, specqbft.PrepareMsgType, specqbft.CommitMsgType}, - expectedProcessed: []bool{true, true, true}, // All consensus messages should be processed - }, - { - name: "decided", - hasRunningDuty: true, - hasRunningInstance: true, - decided: true, - proposalAccepted: true, - messagesTypes: []specqbft.MessageType{specqbft.ProposalMsgType, specqbft.PrepareMsgType, specqbft.CommitMsgType}, - expectedProcessed: []bool{true, true, true}, // All should be processed + name: "decided", + hasRunningDuty: true, + decided: true, + proposalAccepted: true, + messagesTypes: []specqbft.MessageType{specqbft.ProposalMsgType, specqbft.PrepareMsgType, specqbft.CommitMsgType}, + expectedProcessed: []bool{true, true, true}, // All should be processed }, } @@ -764,7 +750,7 @@ func TestCommitteeQueueFilteringScenarios(t *testing.T) { q := queueContainer{ Q: queue.New(logger, 10), queueState: &queue.State{ - HasRunningInstance: tc.hasRunningInstance, + HasRunningInstance: tc.hasRunningDuty, Height: specqbft.Height(slot), Slot: slot, Round: 1, @@ -778,20 +764,17 @@ func TestCommitteeQueueFilteringScenarios(t *testing.T) { } } - state := &runner.State{} - if tc.hasRunningInstance { - state.RunningInstance = &instance.Instance{ - State: &specqbft.State{ - Decided: tc.decided, - ProposalAcceptedForCurrentRound: proposalMsg, - Round: 1, - }, - } - } - committeeRunner := &runner.CommitteeRunner{ BaseRunner: &runner.BaseRunner{ - State: state, + State: &runner.State{ + RunningInstance: &instance.Instance{ + State: &specqbft.State{ + Decided: tc.decided, + ProposalAcceptedForCurrentRound: proposalMsg, + Round: 1, + }, + }, + }, }, } From 35394a62863ab624e5b8f7f30143b1596538a329 Mon Sep 17 00:00:00 2001 From: iurii Date: Fri, 3 Apr 2026 20:40:55 +0300 Subject: [PATCH 11/14] cleanup --- protocol/v2/ssv/runner/aggregator.go | 2 +- protocol/v2/ssv/runner/committee.go | 2 +- protocol/v2/ssv/runner/proposer.go | 2 +- protocol/v2/ssv/runner/runner.go | 8 ++++---- protocol/v2/ssv/runner/sync_committee_contribution.go | 2 +- protocol/v2/ssv/runner/validator_registration.go | 2 +- protocol/v2/ssv/runner/voluntary_exit.go | 2 +- protocol/v2/ssv/validator/validator.go | 6 +----- 8 files changed, 11 insertions(+), 15 deletions(-) diff --git a/protocol/v2/ssv/runner/aggregator.go b/protocol/v2/ssv/runner/aggregator.go index 46efa71a53..80451149fd 100644 --- a/protocol/v2/ssv/runner/aggregator.go +++ b/protocol/v2/ssv/runner/aggregator.go @@ -492,7 +492,7 @@ func (r *AggregatorRunner) GetRole() spectypes.RunnerRole { return r.BaseRunner.GetRole() } -func (r *AggregatorRunner) GetCurrentDutySlot() (phase0.Slot, bool) { +func (r *AggregatorRunner) GetCurrentDutySlot() phase0.Slot { return r.BaseRunner.GetCurrentDutySlot() } diff --git a/protocol/v2/ssv/runner/committee.go b/protocol/v2/ssv/runner/committee.go index 714bd40d4f..78bd0cc27e 100644 --- a/protocol/v2/ssv/runner/committee.go +++ b/protocol/v2/ssv/runner/committee.go @@ -212,7 +212,7 @@ func (r *CommitteeRunner) GetRole() spectypes.RunnerRole { return r.BaseRunner.GetRole() } -func (r *CommitteeRunner) GetCurrentDutySlot() (phase0.Slot, bool) { +func (r *CommitteeRunner) GetCurrentDutySlot() phase0.Slot { return r.BaseRunner.GetCurrentDutySlot() } diff --git a/protocol/v2/ssv/runner/proposer.go b/protocol/v2/ssv/runner/proposer.go index ceb864ae35..05258db0e7 100644 --- a/protocol/v2/ssv/runner/proposer.go +++ b/protocol/v2/ssv/runner/proposer.go @@ -547,7 +547,7 @@ func (r *ProposerRunner) GetRole() spectypes.RunnerRole { return r.BaseRunner.GetRole() } -func (r *ProposerRunner) GetCurrentDutySlot() (phase0.Slot, bool) { +func (r *ProposerRunner) GetCurrentDutySlot() phase0.Slot { return r.BaseRunner.GetCurrentDutySlot() } diff --git a/protocol/v2/ssv/runner/runner.go b/protocol/v2/ssv/runner/runner.go index c958d2fe6b..8f7a50a34a 100644 --- a/protocol/v2/ssv/runner/runner.go +++ b/protocol/v2/ssv/runner/runner.go @@ -31,7 +31,7 @@ type Getters interface { HasAcceptedProposalForCurrentRound() bool GetShares() map[phase0.ValidatorIndex]*spectypes.Share GetRole() spectypes.RunnerRole - GetCurrentDutySlot() (phase0.Slot, bool) + GetCurrentDutySlot() phase0.Slot GetLastHeight() specqbft.Height GetLastRound() specqbft.Round GetStateRoot() ([32]byte, error) @@ -134,12 +134,12 @@ func (b *BaseRunner) GetRole() spectypes.RunnerRole { return b.RunnerRoleType } -func (b *BaseRunner) GetCurrentDutySlot() (phase0.Slot, bool) { +func (b *BaseRunner) GetCurrentDutySlot() phase0.Slot { if !b.hasDutyAssigned() { - return 0, false + return 0 } - return b.State.CurrentDuty.DutySlot(), true + return b.State.CurrentDuty.DutySlot() } func (b *BaseRunner) GetLastHeight() specqbft.Height { diff --git a/protocol/v2/ssv/runner/sync_committee_contribution.go b/protocol/v2/ssv/runner/sync_committee_contribution.go index 409e750608..c90e06d997 100644 --- a/protocol/v2/ssv/runner/sync_committee_contribution.go +++ b/protocol/v2/ssv/runner/sync_committee_contribution.go @@ -564,7 +564,7 @@ func (r *SyncCommitteeAggregatorRunner) GetRole() spectypes.RunnerRole { return r.BaseRunner.GetRole() } -func (r *SyncCommitteeAggregatorRunner) GetCurrentDutySlot() (phase0.Slot, bool) { +func (r *SyncCommitteeAggregatorRunner) GetCurrentDutySlot() phase0.Slot { return r.BaseRunner.GetCurrentDutySlot() } diff --git a/protocol/v2/ssv/runner/validator_registration.go b/protocol/v2/ssv/runner/validator_registration.go index aa1abb8017..40c9eb496a 100644 --- a/protocol/v2/ssv/runner/validator_registration.go +++ b/protocol/v2/ssv/runner/validator_registration.go @@ -285,7 +285,7 @@ func (r *ValidatorRegistrationRunner) GetRole() spectypes.RunnerRole { return r.BaseRunner.GetRole() } -func (r *ValidatorRegistrationRunner) GetCurrentDutySlot() (phase0.Slot, bool) { +func (r *ValidatorRegistrationRunner) GetCurrentDutySlot() phase0.Slot { return r.BaseRunner.GetCurrentDutySlot() } diff --git a/protocol/v2/ssv/runner/voluntary_exit.go b/protocol/v2/ssv/runner/voluntary_exit.go index b951d05f67..d61c4cc2a1 100644 --- a/protocol/v2/ssv/runner/voluntary_exit.go +++ b/protocol/v2/ssv/runner/voluntary_exit.go @@ -254,7 +254,7 @@ func (r *VoluntaryExitRunner) GetRole() spectypes.RunnerRole { return r.BaseRunner.GetRole() } -func (r *VoluntaryExitRunner) GetCurrentDutySlot() (phase0.Slot, bool) { +func (r *VoluntaryExitRunner) GetCurrentDutySlot() phase0.Slot { return r.BaseRunner.GetCurrentDutySlot() } diff --git a/protocol/v2/ssv/validator/validator.go b/protocol/v2/ssv/validator/validator.go index 313d7a3ca4..a603fd30f2 100644 --- a/protocol/v2/ssv/validator/validator.go +++ b/protocol/v2/ssv/validator/validator.go @@ -219,11 +219,7 @@ func (v *Validator) ProcessMessage(ctx context.Context, logger *zap.Logger, msg if err != nil { return fmt.Errorf("event message: get timeout data: %w", err) } - currentDutySlot, ok := dutyRunner.GetCurrentDutySlot() - if !ok { - return fmt.Errorf("event message: get current duty slot to compare vs timeout event slot: %w", err) - } - if timeoutData.Height != specqbft.Height(currentDutySlot) { + if timeoutData.Height != specqbft.Height(dutyRunner.GetCurrentDutySlot()) { // Timeout events can be delayed in the queue until the runner already moved on to a new duty, we can // safely skip these return nil From e232cdab3946dbba7bfdb60efe77149160633c89 Mon Sep 17 00:00:00 2001 From: iurii Date: Mon, 6 Apr 2026 17:05:05 +0300 Subject: [PATCH 12/14] revert the removal of queue filter --- protocol/v2/ssv/runner/runner.go | 5 +++++ protocol/v2/ssv/validator/queue_validator.go | 11 ++++++++++- 2 files changed, 15 insertions(+), 1 deletion(-) diff --git a/protocol/v2/ssv/runner/runner.go b/protocol/v2/ssv/runner/runner.go index 7990bd4ffe..ef5081d7a9 100644 --- a/protocol/v2/ssv/runner/runner.go +++ b/protocol/v2/ssv/runner/runner.go @@ -27,6 +27,7 @@ import ( ) type Getters interface { + HasRunningDuty() bool HasRunningQBFTInstance() bool HasAcceptedProposalForCurrentRound() bool GetShares() map[phase0.ValidatorIndex]*spectypes.Share @@ -102,6 +103,10 @@ type BaseRunner struct { highestDecidedSlot phase0.Slot } +func (b *BaseRunner) HasRunningDuty() bool { + return b.hasDutyRunning() +} + func (b *BaseRunner) HasStartedQBFTInstance() bool { return b.hasDutyAssigned() && b.State.RunningInstance != nil } diff --git a/protocol/v2/ssv/validator/queue_validator.go b/protocol/v2/ssv/validator/queue_validator.go index d6ba11374d..337efe74cb 100644 --- a/protocol/v2/ssv/validator/queue_validator.go +++ b/protocol/v2/ssv/validator/queue_validator.go @@ -123,7 +123,16 @@ func (v *Validator) StartQueueConsumer( state.Quorum = v.Operator.GetQuorum() filter := queue.FilterAny - if state.HasRunningInstance && !r.HasAcceptedProposalForCurrentRound() { + if !r.HasRunningDuty() { + // If no duty is running, pop only ExecuteDuty messages. + filter = func(m *queue.SSVMessage) bool { + e, ok := m.Body.(*types.EventMsg) + if !ok || e == nil { + return false + } + return e.Type == types.ExecuteDuty + } + } else if state.HasRunningInstance && !r.HasAcceptedProposalForCurrentRound() { // If no proposal was accepted for the current round, skip prepare & commit messages // for the current height and round. filter = func(m *queue.SSVMessage) bool { From 0fc95c4dd187081122840f5660450db73aaa2438 Mon Sep 17 00:00:00 2001 From: iurii Date: Mon, 6 Apr 2026 19:17:15 +0300 Subject: [PATCH 13/14] runner: clarify queue state management --- protocol/v2/ssv/queue/message_prioritizer.go | 4 +- protocol/v2/ssv/runner/runner.go | 18 +- protocol/v2/ssv/validator/committee.go | 44 +- protocol/v2/ssv/validator/committee_queue.go | 38 +- .../v2/ssv/validator/committee_queue_test.go | 447 +++++------------- protocol/v2/ssv/validator/queue_validator.go | 24 +- protocol/v2/ssv/validator/timer.go | 4 +- 7 files changed, 189 insertions(+), 390 deletions(-) diff --git a/protocol/v2/ssv/queue/message_prioritizer.go b/protocol/v2/ssv/queue/message_prioritizer.go index d6eaa747b1..e22a198b8e 100644 --- a/protocol/v2/ssv/queue/message_prioritizer.go +++ b/protocol/v2/ssv/queue/message_prioritizer.go @@ -5,8 +5,8 @@ import ( specqbft "github.com/ssvlabs/ssv-spec/qbft" ) -// State represents a portion of the current state -// that is relevant to the prioritization of messages. +// State represents Runner state that is useful for comparing the priority of various messages (message priority +// depends on what the current runner state is). type State struct { HasRunningInstance bool Height specqbft.Height diff --git a/protocol/v2/ssv/runner/runner.go b/protocol/v2/ssv/runner/runner.go index ef5081d7a9..ff398cec45 100644 --- a/protocol/v2/ssv/runner/runner.go +++ b/protocol/v2/ssv/runner/runner.go @@ -135,14 +135,6 @@ func (b *BaseRunner) GetRole() spectypes.RunnerRole { return b.RunnerRoleType } -func (b *BaseRunner) GetCurrentDutySlot() phase0.Slot { - if !b.hasDutyAssigned() { - return 0 - } - - return b.State.CurrentDuty.DutySlot() -} - func (b *BaseRunner) GetLastHeight() specqbft.Height { if ctrl := b.QBFTController; ctrl != nil { return ctrl.Height @@ -480,6 +472,14 @@ func (b *BaseRunner) hasDutyFinished() bool { return b.hasDutyAssigned() && b.State.Finished } +func (b *BaseRunner) currentDutySlot() phase0.Slot { + if !b.hasDutyAssigned() { + return 0 + } + // State.CurrentDuty cannot be nil for non-nil State by construction. + return b.State.CurrentDuty.DutySlot() +} + func (b *BaseRunner) ShouldProcessDuty(duty spectypes.Duty) error { if b.QBFTController.Height >= specqbft.Height(duty.DutySlot()) && b.QBFTController.Height != 0 { return spectypes.NewError( @@ -507,7 +507,7 @@ func (b *BaseRunner) OnTimeoutQBFT(ctx context.Context, logger *zap.Logger, time return nil } - if timeoutData.Height != specqbft.Height(b.GetCurrentDutySlot()) { + if timeoutData.Height != specqbft.Height(b.currentDutySlot()) { // Validator-Runners are re-used to process duties targeting different slots (unlike Committee-Runners that // are working with exactly one slot), thus for Validator-Runners timeout events can be delayed in the queue // until the runner has already moved on to a new duty/slot - this is why timeout-event height(== slot) diff --git a/protocol/v2/ssv/validator/committee.go b/protocol/v2/ssv/validator/committee.go index ecac50eb22..a673ca423a 100644 --- a/protocol/v2/ssv/validator/committee.go +++ b/protocol/v2/ssv/validator/committee.go @@ -35,7 +35,7 @@ type Committee struct { // mtx syncs access to Queues, Runners, Shares. mtx sync.RWMutex - Queues map[phase0.Slot]queueContainer + Queues map[phase0.Slot]queue.Queue Runners map[phase0.Slot]*runner.CommitteeRunner Shares map[phase0.ValidatorIndex]*spectypes.Share @@ -65,7 +65,7 @@ func NewCommittee( return &Committee{ logger: logger, networkConfig: networkConfig, - Queues: make(map[phase0.Slot]queueContainer), + Queues: make(map[phase0.Slot]queue.Queue), Runners: make(map[phase0.Slot]*runner.CommitteeRunner), Shares: shares, CommitteeMember: operator, @@ -92,7 +92,7 @@ func (c *Committee) RemoveShare(validatorIndex phase0.ValidatorIndex) { // StartDuty starts a new duty for the given slot. func (c *Committee) StartDuty(ctx context.Context, logger *zap.Logger, duty *spectypes.CommitteeDuty) ( *runner.CommitteeRunner, - queueContainer, + queue.Queue, error, ) { ctx, span := tracer.Start(ctx, @@ -106,13 +106,13 @@ func (c *Committee) StartDuty(ctx context.Context, logger *zap.Logger, duty *spe span.AddEvent("prepare duty and runner") r, q, runnableDuty, err := c.prepareDutyAndRunner(ctx, logger, duty) if err != nil { - return nil, queueContainer{}, traces.Errorf(span, "prepare duty and runner: %w", err) + return nil, nil, traces.Errorf(span, "prepare duty and runner: %w", err) } logger.Info("ℹ️ starting duty processing") err = r.StartNewDuty(ctx, logger, runnableDuty, c.CommitteeMember.GetQuorum()) if err != nil { - return nil, queueContainer{}, traces.Errorf(span, "runner failed to start duty: %w", err) + return nil, nil, traces.Errorf(span, "runner failed to start duty: %w", err) } span.SetStatus(codes.Ok, "") @@ -121,7 +121,7 @@ func (c *Committee) StartDuty(ctx context.Context, logger *zap.Logger, duty *spe func (c *Committee) prepareDutyAndRunner(ctx context.Context, logger *zap.Logger, duty *spectypes.CommitteeDuty) ( r *runner.CommitteeRunner, - q queueContainer, + q queue.Queue, runnableDuty *spectypes.CommitteeDuty, err error, ) { @@ -137,18 +137,18 @@ func (c *Committee) prepareDutyAndRunner(ctx context.Context, logger *zap.Logger defer c.mtx.Unlock() if _, exists := c.Runners[duty.Slot]; exists { - return nil, queueContainer{}, nil, traces.Errorf(span, "CommitteeRunner for slot %d already exists", duty.Slot) + return nil, nil, nil, traces.Errorf(span, "CommitteeRunner for slot %d already exists", duty.Slot) } shares, attesters, runnableDuty, err := c.prepareDuty(logger, duty) if err != nil { - return nil, queueContainer{}, nil, traces.Error(span, err) + return nil, nil, nil, traces.Error(span, err) } // Create the corresponding runner. r, err = c.CreateRunnerFn(duty.Slot, shares, attesters, c.dutyGuard) if err != nil { - return nil, queueContainer{}, nil, traces.Errorf(span, "could not create CommitteeRunner: %w", err) + return nil, nil, nil, traces.Errorf(span, "could not create CommitteeRunner: %w", err) } r.SetTimeoutFunc(c.onTimeout) c.Runners[duty.Slot] = r @@ -165,26 +165,18 @@ func (c *Committee) prepareDutyAndRunner(ctx context.Context, logger *zap.Logger // getQueue returns queue for the provided slot, lazily initializing it if it didn't exist previously. // MUST be called with c.mtx locked! -func (c *Committee) getQueue(logger *zap.Logger, slot phase0.Slot) queueContainer { +func (c *Committee) getQueue(logger *zap.Logger, slot phase0.Slot) queue.Queue { q, exists := c.Queues[slot] if !exists { - q = queueContainer{ - Q: queue.New( - logger, - 1000, - queue.WithInboxSizeMetric( - queue.InboxSizeMetric, - queue.CommitteeQueueMetricType, - queue.CommitteeMetricID(slot), - ), + q = queue.New( + logger, + 1000, + queue.WithInboxSizeMetric( + queue.InboxSizeMetric, + queue.CommitteeQueueMetricType, + queue.CommitteeMetricID(slot), ), - queueState: &queue.State{ - HasRunningInstance: false, - Height: specqbft.Height(slot), - Slot: slot, - Quorum: c.CommitteeMember.GetQuorum(), - }, - } + ) c.Queues[slot] = q } diff --git a/protocol/v2/ssv/validator/committee_queue.go b/protocol/v2/ssv/validator/committee_queue.go index 757c8754ed..880cd31d80 100644 --- a/protocol/v2/ssv/validator/committee_queue.go +++ b/protocol/v2/ssv/validator/committee_queue.go @@ -23,12 +23,6 @@ import ( "github.com/ssvlabs/ssv/protocol/v2/types" ) -// queueContainer wraps a queue with its corresponding state -type queueContainer struct { - Q queue.Queue - queueState *queue.State -} - // EnqueueMessage enqueues a spectypes.SSVMessage for processing. // TODO: accept DecodedSSVMessage once p2p is upgraded to decode messages during validation. func (c *Committee) EnqueueMessage(ctx context.Context, msg *queue.SSVMessage) { @@ -67,7 +61,7 @@ func (c *Committee) EnqueueMessage(ctx context.Context, msg *queue.SSVMessage) { c.mtx.Unlock() span.AddEvent("pushing message to the queue") - if pushed := q.Q.TryPush(msg); !pushed { + if pushed := q.TryPush(msg); !pushed { const errMsg = "❗ dropping message because the queue is full" logger.Warn(errMsg) span.SetStatus(codes.Error, errMsg) @@ -82,16 +76,13 @@ func (c *Committee) EnqueueMessage(ctx context.Context, msg *queue.SSVMessage) { func (c *Committee) ConsumeQueue( ctx context.Context, logger *zap.Logger, - q queueContainer, + q queue.Queue, handler MessageHandler, // should be c.ProcessMessage, it is a param so can be mocked out for testing - rnr *runner.CommitteeRunner, + r *runner.CommitteeRunner, ) { logger.Debug("📬 queue consumer is running") defer logger.Debug("📪 queue consumer is closed") - // Construct a representation of the current state. - state := *q.queueState - // msgStates keeps track of in-flight processing state (retry count + span context) per message. // Since this map grows over time, we need to clean it up automatically. There is no specific TTL value // to use for its entries - it just needs to be large enough to prevent unnecessary (but non-harmful) @@ -102,11 +93,21 @@ func (c *Committee) ConsumeQueue( go msgStates.Start() defer msgStates.Stop() + // rState defines current runner state that will be used for deciding which messages we want to process + // sooner (vs which ones can wait till later). + rState := queue.State{ + Quorum: c.CommitteeMember.GetQuorum(), // never changes for duty runner + // Slot: slot, // Slot is not used to prioritize messages + } + for ctx.Err() == nil { - state.HasRunningInstance = rnr.HasRunningQBFTInstance() + // Update rState to incorporate the effects previously handled message might have had on the runner state. + rState.HasRunningInstance = r.HasRunningQBFTInstance() + rState.Height = r.GetLastHeight() + rState.Round = r.GetLastRound() filter := queue.FilterAny - if state.HasRunningInstance && !rnr.HasAcceptedProposalForCurrentRound() { + if rState.HasRunningInstance && !r.HasAcceptedProposalForCurrentRound() { // If no proposal was accepted for the current round, skip prepare & commit messages // for the current round. filter = func(m *queue.SSVMessage) bool { @@ -115,13 +116,13 @@ func (c *Committee) ConsumeQueue( return m.MsgType != spectypes.SSVPartialSignatureMsgType } - if sm.Round != state.Round { // allow next round or change round messages. + if sm.Round != rState.Round { // allow next round or change round messages. return true } return sm.MsgType != specqbft.PrepareMsgType && sm.MsgType != specqbft.CommitMsgType } - } else if state.HasRunningInstance { + } else if rState.HasRunningInstance { filter = func(ssvMessage *queue.SSVMessage) bool { // don't read post consensus until decided return ssvMessage.MsgType != spectypes.SSVPartialSignatureMsgType @@ -129,8 +130,7 @@ func (c *Committee) ConsumeQueue( } // Pop the highest priority message for the current state. - // TODO: (Alan) bring back filter - msg := q.Q.Pop(ctx, queue.NewCommitteeQueuePrioritizer(&state), filter) + msg := q.Pop(ctx, queue.NewCommitteeQueuePrioritizer(&rState), filter) if ctx.Err() != nil { // Optimization: terminate fast if we can. return @@ -243,7 +243,7 @@ func (c *Committee) ConsumeQueue( case <-msgState.ctx.Done(): return } - if pushed := q.Q.TryPush(msg); !pushed { + if pushed := q.TryPush(msg); !pushed { const droppingMsgDueToQueueIsFullEvent = "❗ not gonna replay message because the queue is full" msgLogger.Error(droppingMsgDueToQueueIsFullEvent) msgState.span.AddEvent(droppingMsgDueToQueueIsFullEvent, trace.WithAttributes( diff --git a/protocol/v2/ssv/validator/committee_queue_test.go b/protocol/v2/ssv/validator/committee_queue_test.go index 4c31962cda..0e195dda7c 100644 --- a/protocol/v2/ssv/validator/committee_queue_test.go +++ b/protocol/v2/ssv/validator/committee_queue_test.go @@ -22,6 +22,7 @@ import ( "github.com/ssvlabs/ssv/networkconfig" "github.com/ssvlabs/ssv/observability/log" "github.com/ssvlabs/ssv/protocol/v2/message" + "github.com/ssvlabs/ssv/protocol/v2/qbft/controller" "github.com/ssvlabs/ssv/protocol/v2/qbft/instance" "github.com/ssvlabs/ssv/protocol/v2/ssv/queue" "github.com/ssvlabs/ssv/protocol/v2/ssv/runner" @@ -80,7 +81,7 @@ func runConsumeQueueAsync( t *testing.T, ctx context.Context, committee *Committee, - q queueContainer, + q queue.Queue, logger *zap.Logger, handler MessageHandler, committeeRunner *runner.CommitteeRunner, @@ -146,6 +147,40 @@ func setupMessageCollection(capacity int) (chan *queue.SSVMessage, MessageHandle return msgChannel, handler } +func newCommitteeQueueStateForTest(slot phase0.Slot, round specqbft.Round, hasRunningInstance bool, quorum uint64) *queue.State { + return &queue.State{ + HasRunningInstance: hasRunningInstance, + Height: specqbft.Height(slot), + Slot: slot, + Round: round, + Quorum: quorum, + } +} + +func newCommitteeRunnerForTest( + slot phase0.Slot, + round specqbft.Round, + decided bool, + proposal *specqbft.ProcessingMessage, +) *runner.CommitteeRunner { + return &runner.CommitteeRunner{ + BaseRunner: &runner.BaseRunner{ + QBFTController: &controller.Controller{ + Height: specqbft.Height(slot), + }, + State: &runner.State{ + RunningInstance: &instance.Instance{ + State: &specqbft.State{ + Decided: decided, + ProposalAcceptedForCurrentRound: proposal, + Round: round, + }, + }, + }, + }, + } +} + // TestHandleMessageCreatesQueue verifies that the HandleMessage method correctly // initializes a new queue when receiving a message for a slot that doesn't have // an associated queue yet. @@ -170,7 +205,7 @@ func TestHandleMessageCreatesQueue(t *testing.T) { committee := &Committee{ logger: logger, networkConfig: networkconfig.TestNetwork, - Queues: make(map[phase0.Slot]queueContainer), + Queues: make(map[phase0.Slot]queue.Queue), Runners: make(map[phase0.Slot]*runner.CommitteeRunner), CommitteeMember: &spectypes.CommitteeMember{}, } @@ -189,13 +224,18 @@ func TestHandleMessageCreatesQueue(t *testing.T) { require.True(t, ok) - assert.NotNil(t, q.Q) - assert.Equal(t, slot, q.queueState.Slot) - assert.False(t, q.queueState.HasRunningInstance) - assert.Equal(t, specqbft.Height(slot), q.queueState.Height) + assert.NotNil(t, q) + assert.Equal(t, 1, q.Len()) - // default, the queueState.Round is not explicitly initialized from the incoming message - assert.Equal(t, specqbft.Round(0), q.queueState.Round) + queuedMsg := q.TryPop( + queue.NewCommitteeQueuePrioritizer( + newCommitteeQueueStateForTest(slot, 0, false, committee.CommitteeMember.GetQuorum()), + ), + queue.FilterAny, + ) + require.NotNil(t, queuedMsg) + assert.Equal(t, testMsg.MsgID, queuedMsg.MsgID) + assert.Equal(t, testMsg.MsgType, queuedMsg.MsgType) } // TestConsumeQueueBasic tests the fundamental queue consumption functionality @@ -222,7 +262,7 @@ func TestConsumeQueueBasic(t *testing.T) { committee := &Committee{ logger: logger, networkConfig: networkconfig.TestNetwork, - Queues: make(map[phase0.Slot]queueContainer), + Queues: make(map[phase0.Slot]queue.Queue), Runners: make(map[phase0.Slot]*runner.CommitteeRunner), CommitteeMember: &spectypes.CommitteeMember{}, } @@ -245,35 +285,15 @@ func TestConsumeQueueBasic(t *testing.T) { } testMsg2 := makeTestSSVMessage(t, spectypes.SSVConsensusMsgType, msgID2, qbftMsg2) - q := queueContainer{ - Q: queue.New(logger, 1000), - queueState: &queue.State{ - HasRunningInstance: true, - Height: specqbft.Height(slot), - Slot: slot, - Round: 1, - }, - } - q.Q.TryPush(testMsg1) - q.Q.TryPush(testMsg2) + q := queue.New(logger, 1000) + q.TryPush(testMsg1) + q.TryPush(testMsg2) proposalMsg := &specqbft.ProcessingMessage{ QBFTMessage: qbftMsg1, } - committeeRunner := &runner.CommitteeRunner{ - BaseRunner: &runner.BaseRunner{ - State: &runner.State{ - RunningInstance: &instance.Instance{ - State: &specqbft.State{ - Decided: false, - ProposalAcceptedForCurrentRound: proposalMsg, - Round: 1, - }, - }, - }, - }, - } + committeeRunner := newCommitteeRunnerForTest(slot, 1, false, proposalMsg) msgChannel, handler := setupMessageCollection(2) runConsumeQueueAsync(t, ctx, committee, q, logger, handler, committeeRunner) @@ -307,7 +327,7 @@ func TestFilterNoProposalAccepted(t *testing.T) { committee := &Committee{ networkConfig: networkconfig.TestNetwork, - Queues: make(map[phase0.Slot]queueContainer), + Queues: make(map[phase0.Slot]queue.Queue), Runners: make(map[phase0.Slot]*runner.CommitteeRunner), CommitteeMember: &spectypes.CommitteeMember{}, } @@ -351,34 +371,14 @@ func TestFilterNoProposalAccepted(t *testing.T) { combinedMessages[i], combinedMessages[j] = combinedMessages[j], combinedMessages[i] }) - q := queueContainer{ - Q: queue.New(logger, 1000), - queueState: &queue.State{ - HasRunningInstance: true, - Height: specqbft.Height(slot), - Slot: slot, - Round: currentRound, - }, - } + q := queue.New(logger, 1000) for _, combined := range combinedMessages { testMsg := makeTestSSVMessage(t, spectypes.SSVConsensusMsgType, combined.ID, combined.Message) - q.Q.TryPush(testMsg) + q.TryPush(testMsg) } - committeeRunner := &runner.CommitteeRunner{ - BaseRunner: &runner.BaseRunner{ - State: &runner.State{ - RunningInstance: &instance.Instance{ - State: &specqbft.State{ - Decided: false, - ProposalAcceptedForCurrentRound: nil, - Round: currentRound, - }, - }, - }, - }, - } + committeeRunner := newCommitteeRunnerForTest(slot, currentRound, false, nil) msgChannel, handler := setupMessageCollection(4) runConsumeQueueAsync(t, ctx, committee, q, logger, handler, committeeRunner) @@ -427,7 +427,7 @@ func TestFilterNotDecidedSkipsPartialSignatures(t *testing.T) { committee := &Committee{ networkConfig: networkconfig.TestNetwork, - Queues: make(map[phase0.Slot]queueContainer), + Queues: make(map[phase0.Slot]queue.Queue), Runners: make(map[phase0.Slot]*runner.CommitteeRunner), CommitteeMember: &spectypes.CommitteeMember{}, } @@ -458,36 +458,16 @@ func TestFilterNotDecidedSkipsPartialSignatures(t *testing.T) { testMsg1 := makeTestSSVMessage(t, spectypes.SSVConsensusMsgType, msgID1, qbftMsg) testMsg2 := makeTestSSVMessage(t, spectypes.SSVPartialSignatureMsgType, msgID2, partialSigMsg) - q := queueContainer{ - Q: queue.New(logger, 1000), - queueState: &queue.State{ - HasRunningInstance: true, - Height: specqbft.Height(slot), - Slot: slot, - Round: 1, - }, - } + q := queue.New(logger, 1000) - q.Q.TryPush(testMsg1) - q.Q.TryPush(testMsg2) + q.TryPush(testMsg1) + q.TryPush(testMsg2) proposalMsg := &specqbft.ProcessingMessage{ QBFTMessage: qbftMsg, } - committeeRunner := &runner.CommitteeRunner{ - BaseRunner: &runner.BaseRunner{ - State: &runner.State{ - RunningInstance: &instance.Instance{ - State: &specqbft.State{ - Decided: false, - ProposalAcceptedForCurrentRound: proposalMsg, - Round: 1, - }, - }, - }, - }, - } + committeeRunner := newCommitteeRunnerForTest(slot, 1, false, proposalMsg) msgChannel, handler := setupMessageCollection(2) runConsumeQueueAsync(t, ctx, committee, q, logger, handler, committeeRunner) @@ -507,7 +487,7 @@ func TestFilterDecidedAllowsAll(t *testing.T) { committee := &Committee{ networkConfig: networkconfig.TestNetwork, - Queues: make(map[phase0.Slot]queueContainer), + Queues: make(map[phase0.Slot]queue.Queue), Runners: make(map[phase0.Slot]*runner.CommitteeRunner), CommitteeMember: &spectypes.CommitteeMember{}, } @@ -538,36 +518,16 @@ func TestFilterDecidedAllowsAll(t *testing.T) { testMsg1 := makeTestSSVMessage(t, spectypes.SSVConsensusMsgType, msgID1, qbftMsg) testMsg2 := makeTestSSVMessage(t, spectypes.SSVPartialSignatureMsgType, msgID2, partialSigMsg) - q := queueContainer{ - Q: queue.New(logger, 1000), - queueState: &queue.State{ - HasRunningInstance: true, - Height: specqbft.Height(slot), - Slot: slot, - Round: 1, - }, - } + q := queue.New(logger, 1000) - q.Q.TryPush(testMsg1) - q.Q.TryPush(testMsg2) + q.TryPush(testMsg1) + q.TryPush(testMsg2) proposalMsg := &specqbft.ProcessingMessage{ QBFTMessage: qbftMsg, } - committeeRunner := &runner.CommitteeRunner{ - BaseRunner: &runner.BaseRunner{ - State: &runner.State{ - RunningInstance: &instance.Instance{ - State: &specqbft.State{ - Decided: true, - ProposalAcceptedForCurrentRound: proposalMsg, - Round: 1, - }, - }, - }, - }, - } + committeeRunner := newCommitteeRunnerForTest(slot, 1, true, proposalMsg) msgChannel, handler := setupMessageCollection(2) runConsumeQueueAsync(t, ctx, committee, q, logger, handler, committeeRunner) @@ -619,16 +579,8 @@ func TestChangingFilterState(t *testing.T) { return fmt.Errorf("intentionally stopping ConsumeQueue after first message") } - q := queueContainer{ - Q: queue.New(logger, 1), - queueState: &queue.State{ - HasRunningInstance: true, - Height: specqbft.Height(slot), - Slot: slot, - Round: round, - }, - } - q.Q.TryPush(prepareMsg) + q := queue.New(logger, 1) + q.TryPush(prepareMsg) c := &Committee{ networkConfig: networkconfig.TestNetwork, @@ -639,36 +591,12 @@ func TestChangingFilterState(t *testing.T) { } // 1) No proposal accepted => Prepare should be filtered out - r1 := &runner.CommitteeRunner{ - BaseRunner: &runner.BaseRunner{ - State: &runner.State{ - RunningInstance: &instance.Instance{ - State: &specqbft.State{ - Decided: false, - ProposalAcceptedForCurrentRound: nil, - Round: round, - }, - }, - }, - }, - } + r1 := newCommitteeRunnerForTest(slot, round, false, nil) seen1 := runOnce(r1) assert.Nil(t, seen1) // 2) Proposal accepted => now we should see exactly one Prepare - r2 := &runner.CommitteeRunner{ - BaseRunner: &runner.BaseRunner{ - State: &runner.State{ - RunningInstance: &instance.Instance{ - State: &specqbft.State{ - Decided: false, - ProposalAcceptedForCurrentRound: &specqbft.ProcessingMessage{QBFTMessage: prepareBody}, - Round: round, - }, - }, - }, - }, - } + r2 := newCommitteeRunnerForTest(slot, round, false, &specqbft.ProcessingMessage{QBFTMessage: prepareBody}) seen2 := runOnce(r2) require.NotNil(t, seen2) @@ -740,22 +668,14 @@ func TestCommitteeQueueFilteringScenarios(t *testing.T) { committee := &Committee{ networkConfig: networkconfig.TestNetwork, - Queues: make(map[phase0.Slot]queueContainer), + Queues: make(map[phase0.Slot]queue.Queue), Runners: make(map[phase0.Slot]*runner.CommitteeRunner), CommitteeMember: &spectypes.CommitteeMember{}, } slot := phase0.Slot(123) - q := queueContainer{ - Q: queue.New(logger, 10), - queueState: &queue.State{ - HasRunningInstance: tc.hasRunningDuty, - Height: specqbft.Height(slot), - Slot: slot, - Round: 1, - }, - } + q := queue.New(logger, 10) var proposalMsg *specqbft.ProcessingMessage if tc.proposalAccepted { @@ -764,19 +684,7 @@ func TestCommitteeQueueFilteringScenarios(t *testing.T) { } } - committeeRunner := &runner.CommitteeRunner{ - BaseRunner: &runner.BaseRunner{ - State: &runner.State{ - RunningInstance: &instance.Instance{ - State: &specqbft.State{ - Decided: tc.decided, - ProposalAcceptedForCurrentRound: proposalMsg, - Round: 1, - }, - }, - }, - }, - } + committeeRunner := newCommitteeRunnerForTest(slot, 1, tc.decided, proposalMsg) // Set runner state based on hasRunningDuty parameter if !tc.hasRunningDuty { @@ -799,7 +707,7 @@ func TestCommitteeQueueFilteringScenarios(t *testing.T) { MsgType: msgType, } testMsg := makeTestSSVMessage(t, spectypes.SSVConsensusMsgType, msgID, qbftMsg) - pushed := q.Q.TryPush(testMsg) + pushed := q.TryPush(testMsg) require.True(t, pushed) } @@ -902,36 +810,16 @@ func TestFilterPartialSignatureMessages(t *testing.T) { committee := &Committee{ networkConfig: networkconfig.TestNetwork, - Queues: make(map[phase0.Slot]queueContainer), + Queues: make(map[phase0.Slot]queue.Queue), Runners: make(map[phase0.Slot]*runner.CommitteeRunner), CommitteeMember: &spectypes.CommitteeMember{}, } slot := phase0.Slot(123) - q := queueContainer{ - Q: queue.New(logger, 10), - queueState: &queue.State{ - HasRunningInstance: true, - Height: specqbft.Height(slot), - Slot: slot, - Round: 1, - }, - } + q := queue.New(logger, 10) - committeeRunner := &runner.CommitteeRunner{ - BaseRunner: &runner.BaseRunner{ - State: &runner.State{ - RunningInstance: &instance.Instance{ - State: &specqbft.State{ - Decided: tc.decided, - ProposalAcceptedForCurrentRound: &specqbft.ProcessingMessage{}, - Round: 1, - }, - }, - }, - }, - } + committeeRunner := newCommitteeRunnerForTest(slot, 1, tc.decided, &specqbft.ProcessingMessage{}) msgID := spectypes.MessageID{0x10} partialSigMsg := &spectypes.PartialSignatureMessages{ @@ -947,7 +835,7 @@ func TestFilterPartialSignatureMessages(t *testing.T) { } testMsg := makeTestSSVMessage(t, spectypes.SSVPartialSignatureMsgType, msgID, partialSigMsg) - pushed := q.Q.TryPush(testMsg) + pushed := q.TryPush(testMsg) require.True(t, pushed) if tc.shouldBeFiltered { @@ -989,7 +877,7 @@ func TestConsumeQueuePrioritization(t *testing.T) { committee := &Committee{ networkConfig: networkconfig.TestNetwork, - Queues: make(map[phase0.Slot]queueContainer), + Queues: make(map[phase0.Slot]queue.Queue), Runners: make(map[phase0.Slot]*runner.CommitteeRunner), CommitteeMember: &spectypes.CommitteeMember{}, } @@ -1018,30 +906,14 @@ func TestConsumeQueuePrioritization(t *testing.T) { makeTestSSVMessage(t, message.SSVEventMsgType, spectypes.MessageID{5}, eventMsgBody), } - q := queueContainer{ - Q: queue.New(logger, 10), - queueState: &queue.State{ - HasRunningInstance: true, - Height: specqbft.Height(slot), - Slot: slot, - Round: currentRound, - }, - } + q := queue.New(logger, 10) for _, msg := range testMessages { - q.Q.TryPush(msg) + q.TryPush(msg) } // Runner with a proposal already accepted, not yet decided acceptedProposal := &specqbft.ProcessingMessage{QBFTMessage: proposalMsgBody} - committeeRunner := &runner.CommitteeRunner{ - BaseRunner: &runner.BaseRunner{ - State: &runner.State{RunningInstance: &instance.Instance{State: &specqbft.State{ - Decided: false, - ProposalAcceptedForCurrentRound: acceptedProposal, - Round: currentRound, - }}}, - }, - } + committeeRunner := newCommitteeRunnerForTest(slot, currentRound, false, acceptedProposal) msgChannel := make(chan *queue.SSVMessage, len(testMessages)) @@ -1116,19 +988,13 @@ func TestHandleMessageQueueFullAndDropping(t *testing.T) { committee := &Committee{ logger: logger, networkConfig: networkconfig.TestNetwork, - Queues: make(map[phase0.Slot]queueContainer), + Queues: make(map[phase0.Slot]queue.Queue), CommitteeMember: &spectypes.CommitteeMember{}, } // Step 0: Create the queue container with the desired small capacity and add it to the committee - qContainer := queueContainer{ - Q: queue.New(logger, queueCapacity), - queueState: &queue.State{ - HasRunningInstance: false, - Height: specqbft.Height(slot), - Slot: slot, - }, - } + qContainer := queue.New(logger, queueCapacity) + qState := newCommitteeQueueStateForTest(slot, 0, false, committee.CommitteeMember.GetQuorum()) committee.Queues[slot] = qContainer // Step 1: Fill the pre-made queue to its capacity by calling HandleMessage @@ -1142,7 +1008,7 @@ func TestHandleMessageQueueFullAndDropping(t *testing.T) { committee.EnqueueMessage(ctx, testMsg) } - require.Equal(t, queueCapacity, qContainer.Q.Len()) + require.Equal(t, queueCapacity, qContainer.Len()) // Step 2: Clear log buffer and attempt to push one more message (this one should be dropped) droppedMsgID := msgIDBase @@ -1152,7 +1018,7 @@ func TestHandleMessageQueueFullAndDropping(t *testing.T) { committee.EnqueueMessage(ctx, testMsgDrop) - assert.Equal(t, queueCapacity, qContainer.Q.Len()) + assert.Equal(t, queueCapacity, qContainer.Len()) // Step 3: Verify that the dropped message is not in the queue and original messages are intact. // Pop messages one by one and check their MsgID and Type. @@ -1164,7 +1030,7 @@ func TestHandleMessageQueueFullAndDropping(t *testing.T) { popCtx, popCancel := context.WithTimeout(t.Context(), 200*time.Millisecond) // Use FilterAny since we are just checking the contents, not a live consumption scenario. // The prioritizer does not matter here as we drain the queue completely. - msg := qContainer.Q.Pop(popCtx, queue.NewCommitteeQueuePrioritizer(qContainer.queueState), queue.FilterAny) + msg := qContainer.Pop(popCtx, queue.NewCommitteeQueuePrioritizer(qState), queue.FilterAny) popCancel() require.NotNil(t, msg) @@ -1194,7 +1060,7 @@ func TestHandleMessageQueueFullAndDropping(t *testing.T) { finalPopCtx, finalPopCancel := context.WithTimeout(t.Context(), 200*time.Millisecond) defer finalPopCancel() - assert.Nil(t, qContainer.Q.Pop(finalPopCtx, queue.NewCommitteeQueuePrioritizer(qContainer.queueState), queue.FilterAny)) + assert.Nil(t, qContainer.Pop(finalPopCtx, queue.NewCommitteeQueuePrioritizer(qState), queue.FilterAny)) } // TestConsumeQueueStopsOnErrNoValidDuties verifies that ConsumeQueue stops @@ -1220,27 +1086,17 @@ func TestConsumeQueueStopsOnErrNoValidDuties(t *testing.T) { } slot := phase0.Slot(123) - q := queueContainer{ - Q: queue.New(logger, 10), - queueState: &queue.State{ - HasRunningInstance: true, - Height: specqbft.Height(slot), - Slot: slot, - Round: 1, - }, - } + q := queue.New(logger, 10) // Add multiple messages - msg1 := makeTestSSVMessage(t, spectypes.SSVConsensusMsgType, spectypes.MessageID{1}, &specqbft.Message{Height: specqbft.Height(slot), MsgType: specqbft.ProposalMsgType}) - msg2 := makeTestSSVMessage(t, spectypes.SSVConsensusMsgType, spectypes.MessageID{2}, &specqbft.Message{Height: specqbft.Height(slot), MsgType: specqbft.PrepareMsgType}) - msg3 := makeTestSSVMessage(t, spectypes.SSVConsensusMsgType, spectypes.MessageID{3}, &specqbft.Message{Height: specqbft.Height(slot), MsgType: specqbft.CommitMsgType}) - q.Q.TryPush(msg1) - q.Q.TryPush(msg2) - q.Q.TryPush(msg3) - - committeeRunner := &runner.CommitteeRunner{ - BaseRunner: &runner.BaseRunner{State: &runner.State{RunningInstance: &instance.Instance{State: &specqbft.State{}}}}, - } + msg1 := makeTestSSVMessage(t, spectypes.SSVConsensusMsgType, spectypes.MessageID{1}, &specqbft.Message{Height: specqbft.Height(slot), Round: 1, MsgType: specqbft.ProposalMsgType}) + msg2 := makeTestSSVMessage(t, spectypes.SSVConsensusMsgType, spectypes.MessageID{2}, &specqbft.Message{Height: specqbft.Height(slot), Round: 1, MsgType: specqbft.PrepareMsgType}) + msg3 := makeTestSSVMessage(t, spectypes.SSVConsensusMsgType, spectypes.MessageID{3}, &specqbft.Message{Height: specqbft.Height(slot), Round: 1, MsgType: specqbft.CommitMsgType}) + q.TryPush(msg1) + q.TryPush(msg2) + q.TryPush(msg3) + + committeeRunner := newCommitteeRunnerForTest(slot, 1, false, nil) var processedMessagesCount int32 handler := func(ctx context.Context, _ *zap.Logger, msg *queue.SSVMessage) error { @@ -1259,7 +1115,7 @@ func TestConsumeQueueStopsOnErrNoValidDuties(t *testing.T) { committee.ConsumeQueue(ctx, logger, q, handler, committeeRunner) assert.Equal(t, int32(1), atomic.LoadInt32(&processedMessagesCount)) - assert.Equal(t, 2, q.Q.Len()) + assert.Equal(t, 2, q.Len()) } // TestConsumeQueueBurstTraffic verifies that under a burst of interleaved messages, @@ -1283,19 +1139,11 @@ func TestConsumeQueueBurstTraffic(t *testing.T) { slot := phase0.Slot(42) committee := &Committee{ networkConfig: networkconfig.TestNetwork, - Queues: make(map[phase0.Slot]queueContainer), + Queues: make(map[phase0.Slot]queue.Queue), Runners: make(map[phase0.Slot]*runner.CommitteeRunner), CommitteeMember: &spectypes.CommitteeMember{}, } - qc := queueContainer{ - Q: queue.New(logger, 1000), - queueState: &queue.State{ - HasRunningInstance: true, - Height: specqbft.Height(slot), - Slot: slot, - Round: 1, - }, - } + qc := queue.New(logger, 1000) committee.Queues[slot] = qc // Mark that consensus is already decided & proposal accepted → partial-sigs allowed @@ -1306,19 +1154,7 @@ func TestConsumeQueueBurstTraffic(t *testing.T) { MsgType: specqbft.ProposalMsgType, }, } - committee.Runners[slot] = &runner.CommitteeRunner{ - BaseRunner: &runner.BaseRunner{ - State: &runner.State{ - RunningInstance: &instance.Instance{ - State: &specqbft.State{ - Decided: true, - ProposalAcceptedForCurrentRound: acceptedProposal, - Round: 1, - }, - }, - }, - }, - } + committee.Runners[slot] = newCommitteeRunnerForTest(slot, 1, true, acceptedProposal) // --- Build 200 randomized messages and count expected per priority bucket --- var ( @@ -1404,7 +1240,7 @@ func TestConsumeQueueBurstTraffic(t *testing.T) { allMsgs[i], allMsgs[j] = allMsgs[j], allMsgs[i] }) for _, m := range allMsgs { - require.True(t, qc.Q.TryPush(m)) + require.True(t, qc.TryPush(m)) } // --- Drain the queue, capturing the priority bucket of each popped message --- @@ -1497,7 +1333,7 @@ func TestQueueLoadAndSaturationScenarios(t *testing.T) { committee := &Committee{ logger: logger, networkConfig: networkconfig.TestNetwork, - Queues: make(map[phase0.Slot]queueContainer), + Queues: make(map[phase0.Slot]queue.Queue), Runners: make(map[phase0.Slot]*runner.CommitteeRunner), CommitteeMember: &spectypes.CommitteeMember{}, } @@ -1506,15 +1342,8 @@ func TestQueueLoadAndSaturationScenarios(t *testing.T) { nextRound := specqbft.Round(2) queueCapacity := 3 - qContainer := queueContainer{ - Q: queue.New(logger, queueCapacity), - queueState: &queue.State{ - HasRunningInstance: true, - Height: specqbft.Height(slot), - Slot: slot, - Round: currentRound, - }, - } + qContainer := queue.New(logger, queueCapacity) + qState := newCommitteeQueueStateForTest(slot, currentRound, true, committee.CommitteeMember.GetQuorum()) committee.Queues[slot] = qContainer // 1. Fill the queue's inbox channel to capacity using HandleMessage. @@ -1524,7 +1353,7 @@ func TestQueueLoadAndSaturationScenarios(t *testing.T) { testMsg := makeTestSSVMessage(t, spectypes.SSVConsensusMsgType, msgID, prepareMsgBody) committee.EnqueueMessage(ctx, testMsg) } - require.Equal(t, queueCapacity, qContainer.Q.Len()) + require.Equal(t, queueCapacity, qContainer.Len()) // 2. Attempt to HandleMessage a new Prepare message for the *nextRound*. poppableMsgBody := &specqbft.Message{Height: specqbft.Height(slot), Round: nextRound, MsgType: specqbft.PrepareMsgType} @@ -1532,13 +1361,13 @@ func TestQueueLoadAndSaturationScenarios(t *testing.T) { committee.EnqueueMessage(ctx, poppableTestMsg) // 3. Verify the poppable message was dropped. - assert.Equal(t, queueCapacity, qContainer.Q.Len()) + assert.Equal(t, queueCapacity, qContainer.Len()) // 4. Verify the content of the queue. drainedMessages := make([]*queue.SSVMessage, 0, queueCapacity) for i := 0; i < queueCapacity; i++ { popCtx, popCancel := context.WithTimeout(t.Context(), 200*time.Millisecond) - msg := qContainer.Q.Pop(popCtx, queue.NewCommitteeQueuePrioritizer(qContainer.queueState), queue.FilterAny) + msg := qContainer.Pop(popCtx, queue.NewCommitteeQueuePrioritizer(qState), queue.FilterAny) popCancel() // Ensure cancellation happens after Pop or timeout require.NotNil(t, msg) drainedMessages = append(drainedMessages, msg) @@ -1547,7 +1376,7 @@ func TestQueueLoadAndSaturationScenarios(t *testing.T) { finalPopCtx, finalPopCancel := context.WithTimeout(t.Context(), 200*time.Millisecond) defer finalPopCancel() - assert.Nil(t, qContainer.Q.Pop(finalPopCtx, queue.NewCommitteeQueuePrioritizer(qContainer.queueState), queue.FilterAny), "Queue should be empty after draining initial messages") + assert.Nil(t, qContainer.Pop(finalPopCtx, queue.NewCommitteeQueuePrioritizer(qState), queue.FilterAny), "Queue should be empty after draining initial messages") foundNextRoundMessage := false for _, msg := range drainedMessages { @@ -1573,7 +1402,7 @@ func TestQueueLoadAndSaturationScenarios(t *testing.T) { committee := &Committee{ logger: logger, networkConfig: networkconfig.TestNetwork, - Queues: make(map[phase0.Slot]queueContainer), + Queues: make(map[phase0.Slot]queue.Queue), Runners: make(map[phase0.Slot]*runner.CommitteeRunner), CommitteeMember: &spectypes.CommitteeMember{}, } @@ -1581,15 +1410,8 @@ func TestQueueLoadAndSaturationScenarios(t *testing.T) { currentRound := specqbft.Round(1) queueCapacity := 3 - qContainer := queueContainer{ - Q: queue.New(logger, queueCapacity), - queueState: &queue.State{ - HasRunningInstance: true, - Height: specqbft.Height(slot), - Slot: slot, - Round: currentRound, - }, - } + qContainer := queue.New(logger, queueCapacity) + qState := newCommitteeQueueStateForTest(slot, currentRound, true, committee.CommitteeMember.GetQuorum()) committee.Queues[slot] = qContainer // 1. Fill the queue with low-priority consensus messages @@ -1603,7 +1425,7 @@ func TestQueueLoadAndSaturationScenarios(t *testing.T) { testMsg := makeTestSSVMessage(t, spectypes.SSVConsensusMsgType, msgID, commitMsgBody) committee.EnqueueMessage(ctx, testMsg) } - require.Equal(t, queueCapacity, qContainer.Q.Len(), "Queue should be at capacity") + require.Equal(t, queueCapacity, qContainer.Len(), "Queue should be at capacity") // 2. Try to add a high-priority proposal message (proposals are higher priority than commits) highPriorityMsgBody := &specqbft.Message{ @@ -1623,13 +1445,13 @@ func TestQueueLoadAndSaturationScenarios(t *testing.T) { committee.EnqueueMessage(ctx, highPriorityMsg) // 3. Verify queue length still at capacity - assert.Equal(t, queueCapacity, qContainer.Q.Len()) + assert.Equal(t, queueCapacity, qContainer.Len()) // 4. Verify only the original messages are in the queue drainedMessages := make([]*queue.SSVMessage, 0, queueCapacity) for i := 0; i < queueCapacity; i++ { popCtx, popCancel := context.WithTimeout(t.Context(), 200*time.Millisecond) - msg := qContainer.Q.Pop(popCtx, queue.NewCommitteeQueuePrioritizer(qContainer.queueState), queue.FilterAny) + msg := qContainer.Pop(popCtx, queue.NewCommitteeQueuePrioritizer(qState), queue.FilterAny) popCancel() require.NotNil(t, msg) drainedMessages = append(drainedMessages, msg) @@ -1667,39 +1489,18 @@ func TestQueueLoadAndSaturationScenarios(t *testing.T) { committee := &Committee{ networkConfig: networkconfig.TestNetwork, - Queues: make(map[phase0.Slot]queueContainer), + Queues: make(map[phase0.Slot]queue.Queue), Runners: make(map[phase0.Slot]*runner.CommitteeRunner), CommitteeMember: &spectypes.CommitteeMember{}, } queueCapacity := 5 currentRound := specqbft.Round(1) - - committeeRunner := &runner.CommitteeRunner{ - BaseRunner: &runner.BaseRunner{ - State: &runner.State{ - RunningInstance: &instance.Instance{ - State: &specqbft.State{ - Decided: false, - ProposalAcceptedForCurrentRound: nil, - Round: currentRound, - }, - }, - }, - }, - } - slot := phase0.Slot(456) - q := queueContainer{ - Q: queue.New(logger, queueCapacity), - queueState: &queue.State{ - HasRunningInstance: true, - Height: specqbft.Height(slot), - Slot: slot, - Round: currentRound, - }, - } + committeeRunner := newCommitteeRunnerForTest(slot, currentRound, false, nil) + + q := queue.New(logger, queueCapacity) var ( processedMsgs []*queue.SSVMessage @@ -1728,7 +1529,7 @@ func TestQueueLoadAndSaturationScenarios(t *testing.T) { for i := 0; i < queueCapacity; i++ { msgID := spectypes.MessageID{byte(i + 1)} prepare := &specqbft.Message{Height: specqbft.Height(slot), Round: currentRound, MsgType: specqbft.PrepareMsgType} - require.True(t, q.Q.TryPush(makeTestSSVMessage(t, spectypes.SSVConsensusMsgType, msgID, prepare))) + require.True(t, q.TryPush(makeTestSSVMessage(t, spectypes.SSVConsensusMsgType, msgID, prepare))) } time.Sleep(400 * time.Millisecond) // Give time for the consumer to process (and filter) messages @@ -1742,7 +1543,7 @@ func TestQueueLoadAndSaturationScenarios(t *testing.T) { // Push ExecuteDuty execData, _ := json.Marshal(&types.ExecuteCommitteeDutyData{Duty: &spectypes.CommitteeDuty{Slot: slot}}) execMsg := makeTestSSVMessage(t, message.SSVEventMsgType, spectypes.MessageID{0xEE}, &types.EventMsg{Type: types.ExecuteDuty, Data: execData}) - require.True(t, q.Q.TryPush(execMsg)) + require.True(t, q.TryPush(execMsg)) select { case <-handlerCalled: // Good case <-time.After(1 * time.Second): @@ -1752,7 +1553,7 @@ func TestQueueLoadAndSaturationScenarios(t *testing.T) { // Push Proposal proposal := &specqbft.Message{Height: specqbft.Height(slot), Round: currentRound, MsgType: specqbft.ProposalMsgType} propMsg := makeTestSSVMessage(t, spectypes.SSVConsensusMsgType, spectypes.MessageID{0xFF}, proposal) - require.True(t, q.Q.TryPush(propMsg)) + require.True(t, q.TryPush(propMsg)) select { case <-handlerCalled: // Good case <-time.After(1 * time.Second): diff --git a/protocol/v2/ssv/validator/queue_validator.go b/protocol/v2/ssv/validator/queue_validator.go index 337efe74cb..937f088483 100644 --- a/protocol/v2/ssv/validator/queue_validator.go +++ b/protocol/v2/ssv/validator/queue_validator.go @@ -110,17 +110,23 @@ func (v *Validator) StartQueueConsumer( go msgStates.Start() defer msgStates.Stop() + // rState defines current runner state that will be used for deciding which messages we want to process + // sooner (vs which ones can wait till later). + rState := queue.State{ + Quorum: v.Operator.GetQuorum(), // never changes for duty runner + } + for ctx.Err() == nil { - // Construct a representation of the current state. - state := queue.State{} r := v.DutyRunners.DutyRunnerForMsgID(msgID) if r == nil { return fmt.Errorf("could not get duty runner for msg ID %v", msgID) } - state.HasRunningInstance = r.HasRunningQBFTInstance() - state.Height = r.GetLastHeight() - state.Round = r.GetLastRound() - state.Quorum = v.Operator.GetQuorum() + + // Update rState to incorporate the effects previously handled message might have had on the runner state. + rState.HasRunningInstance = r.HasRunningQBFTInstance() + rState.Height = r.GetLastHeight() + rState.Round = r.GetLastRound() + // rState.Slot = slot // Slot is not used to prioritize messages filter := queue.FilterAny if !r.HasRunningDuty() { @@ -132,7 +138,7 @@ func (v *Validator) StartQueueConsumer( } return e.Type == types.ExecuteDuty } - } else if state.HasRunningInstance && !r.HasAcceptedProposalForCurrentRound() { + } else if rState.HasRunningInstance && !r.HasAcceptedProposalForCurrentRound() { // If no proposal was accepted for the current round, skip prepare & commit messages // for the current height and round. filter = func(m *queue.SSVMessage) bool { @@ -141,7 +147,7 @@ func (v *Validator) StartQueueConsumer( return true } - if qbftMsg.Height != state.Height || qbftMsg.Round != state.Round { + if qbftMsg.Height != rState.Height || qbftMsg.Round != rState.Round { return true } return qbftMsg.MsgType != specqbft.PrepareMsgType && qbftMsg.MsgType != specqbft.CommitMsgType @@ -149,7 +155,7 @@ func (v *Validator) StartQueueConsumer( } // Pop the highest priority message for the current state. - msg := q.Pop(ctx, queue.NewMessagePrioritizer(&state), filter) + msg := q.Pop(ctx, queue.NewMessagePrioritizer(&rState), filter) if ctx.Err() != nil { // Optimization: terminate fast if we can. return nil diff --git a/protocol/v2/ssv/validator/timer.go b/protocol/v2/ssv/validator/timer.go index f49d4ff380..1b6186aad2 100644 --- a/protocol/v2/ssv/validator/timer.go +++ b/protocol/v2/ssv/validator/timer.go @@ -82,7 +82,7 @@ func (c *Committee) onTimeout(ctx context.Context, logger *zap.Logger, identifie // timeout for, in practice this should never happen - but we need to handle this just in case. // This is also possible if the queue got pruned already (due to becoming old and irrelevant). q := c.Queues[phase0.Slot(height)] - if q.Q == nil { + if q == nil { logger.Debug("couldn't schedule timeout event due to missing queue (likely was pruned)") return } @@ -98,7 +98,7 @@ func (c *Committee) onTimeout(ctx context.Context, logger *zap.Logger, identifie return } - if pushed := q.Q.TryPush(dec); !pushed { + if pushed := q.TryPush(dec); !pushed { logger.Error("❗️ dropping timeout message because the queue is full", fields.RunnerRole(identifier.GetRoleType())) } } From 624971a659b40e1e9b60bb38b879e224e49c91df Mon Sep 17 00:00:00 2001 From: iurii Date: Mon, 6 Apr 2026 19:39:21 +0300 Subject: [PATCH 14/14] runner: clarify queue state management (with respect to slot) --- protocol/v2/ssv/queue/messages.go | 2 +- protocol/v2/ssv/runner/runner.go | 19 ++++++++++--------- protocol/v2/ssv/validator/committee_queue.go | 2 +- .../v2/ssv/validator/committee_queue_test.go | 3 +++ protocol/v2/ssv/validator/queue_validator.go | 2 +- 5 files changed, 16 insertions(+), 12 deletions(-) diff --git a/protocol/v2/ssv/queue/messages.go b/protocol/v2/ssv/queue/messages.go index 1dd734b6ec..d47c5d81f8 100644 --- a/protocol/v2/ssv/queue/messages.go +++ b/protocol/v2/ssv/queue/messages.go @@ -134,7 +134,7 @@ func compareHeightOrSlot(state *State, m *SSVMessage) int { if qbftMsg.Height > state.Height { return 1 } - } else if pms, ok := m.Body.(*spectypes.PartialSignatureMessages); ok && pms != nil { // everyone likes pms + } else if pms, ok := m.Body.(*spectypes.PartialSignatureMessages); ok && pms != nil { if pms.Slot == state.Slot { return 0 } diff --git a/protocol/v2/ssv/runner/runner.go b/protocol/v2/ssv/runner/runner.go index ff398cec45..54fca1e5d3 100644 --- a/protocol/v2/ssv/runner/runner.go +++ b/protocol/v2/ssv/runner/runner.go @@ -32,6 +32,7 @@ type Getters interface { HasAcceptedProposalForCurrentRound() bool GetShares() map[phase0.ValidatorIndex]*spectypes.Share GetRole() spectypes.RunnerRole + GetCurrentDutySlot() phase0.Slot GetLastHeight() specqbft.Height GetLastRound() specqbft.Round GetStateRoot() ([32]byte, error) @@ -135,6 +136,14 @@ func (b *BaseRunner) GetRole() spectypes.RunnerRole { return b.RunnerRoleType } +func (b *BaseRunner) GetCurrentDutySlot() phase0.Slot { + if !b.hasDutyAssigned() { + return 0 + } + // State.CurrentDuty cannot be nil for non-nil State by construction. + return b.State.CurrentDuty.DutySlot() +} + func (b *BaseRunner) GetLastHeight() specqbft.Height { if ctrl := b.QBFTController; ctrl != nil { return ctrl.Height @@ -472,14 +481,6 @@ func (b *BaseRunner) hasDutyFinished() bool { return b.hasDutyAssigned() && b.State.Finished } -func (b *BaseRunner) currentDutySlot() phase0.Slot { - if !b.hasDutyAssigned() { - return 0 - } - // State.CurrentDuty cannot be nil for non-nil State by construction. - return b.State.CurrentDuty.DutySlot() -} - func (b *BaseRunner) ShouldProcessDuty(duty spectypes.Duty) error { if b.QBFTController.Height >= specqbft.Height(duty.DutySlot()) && b.QBFTController.Height != 0 { return spectypes.NewError( @@ -507,7 +508,7 @@ func (b *BaseRunner) OnTimeoutQBFT(ctx context.Context, logger *zap.Logger, time return nil } - if timeoutData.Height != specqbft.Height(b.currentDutySlot()) { + if timeoutData.Height != specqbft.Height(b.GetCurrentDutySlot()) { // Validator-Runners are re-used to process duties targeting different slots (unlike Committee-Runners that // are working with exactly one slot), thus for Validator-Runners timeout events can be delayed in the queue // until the runner has already moved on to a new duty/slot - this is why timeout-event height(== slot) diff --git a/protocol/v2/ssv/validator/committee_queue.go b/protocol/v2/ssv/validator/committee_queue.go index 880cd31d80..30e1ed9ffb 100644 --- a/protocol/v2/ssv/validator/committee_queue.go +++ b/protocol/v2/ssv/validator/committee_queue.go @@ -97,7 +97,7 @@ func (c *Committee) ConsumeQueue( // sooner (vs which ones can wait till later). rState := queue.State{ Quorum: c.CommitteeMember.GetQuorum(), // never changes for duty runner - // Slot: slot, // Slot is not used to prioritize messages + Slot: r.GetCurrentDutySlot(), } for ctx.Err() == nil { diff --git a/protocol/v2/ssv/validator/committee_queue_test.go b/protocol/v2/ssv/validator/committee_queue_test.go index 0e195dda7c..95be2e033a 100644 --- a/protocol/v2/ssv/validator/committee_queue_test.go +++ b/protocol/v2/ssv/validator/committee_queue_test.go @@ -176,6 +176,9 @@ func newCommitteeRunnerForTest( Round: round, }, }, + CurrentDuty: &spectypes.CommitteeDuty{ + Slot: slot, + }, }, }, } diff --git a/protocol/v2/ssv/validator/queue_validator.go b/protocol/v2/ssv/validator/queue_validator.go index 937f088483..11e870bfcf 100644 --- a/protocol/v2/ssv/validator/queue_validator.go +++ b/protocol/v2/ssv/validator/queue_validator.go @@ -126,7 +126,7 @@ func (v *Validator) StartQueueConsumer( rState.HasRunningInstance = r.HasRunningQBFTInstance() rState.Height = r.GetLastHeight() rState.Round = r.GetLastRound() - // rState.Slot = slot // Slot is not used to prioritize messages + rState.Slot = r.GetCurrentDutySlot() filter := queue.FilterAny if !r.HasRunningDuty() {