Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 3 additions & 1 deletion api/v1alpha1/seinodetask_types.go
Original file line number Diff line number Diff line change
Expand Up @@ -19,7 +19,9 @@ const (

// SeiNodeTaskKindGovVote backs the sidecar `gov-vote` task. Submits
// MsgVote on an existing proposal. Chain-idempotent (last-write-wins on
// proposalId/voter).
// proposalId/voter). On a node with the tx index off, as on most
// validators, the sidecar confirms the vote from the gov module's recorded
// choice instead of the tx.
SeiNodeTaskKindGovVote SeiNodeTaskKind = "GovVote"

// SeiNodeTaskKindGovParamChange backs the sidecar `gov-param-change` task.
Expand Down
2 changes: 1 addition & 1 deletion sidecar/tasks/gov_result.go
Original file line number Diff line number Diff line change
Expand Up @@ -55,7 +55,7 @@ func classifyGovResult(taskType engine.TaskType, r *SignAndBroadcastResult) (*wi
// Broadcast accepted but the node's tx index is off, so the outcome is
// unobservable. Terminal (retrying this node is futile) but NOT
// committed_failed — the operator must verify via an indexed RPC. The
// Unjail handler overrides this when the jail state shows the release.
// Unjail and GovVote handlers override this when state shows the effect.
out.InclusionStatus = wire.InclusionUnverifiable
txBroadcastTotal.WithLabelValues(string(taskType), wire.InclusionUnverifiable).Inc()
return out, Terminal(fmt.Errorf("tx %s inclusion unverifiable: %w", r.TxHash, errTxIndexingDisabled))
Expand Down
123 changes: 120 additions & 3 deletions sidecar/tasks/gov_vote.go
Original file line number Diff line number Diff line change
Expand Up @@ -10,8 +10,11 @@ import (
"context"
"errors"
"fmt"
"time"

sdk "github.com/sei-protocol/sei-chain/sei-cosmos/types"
govtypes "github.com/sei-protocol/sei-chain/sei-cosmos/x/gov/types"
"github.com/sei-protocol/sei-chain/sei-tendermint/rpc/coretypes"

"github.com/sei-protocol/seilog"

Expand All @@ -38,10 +41,27 @@ type GovVoteRequest struct {
// documented read-only after startup, so the copy is safe.
type GovVoter struct {
cfg engine.ExecutionConfig

// broadcast and readVote are test seams. They default to SignAndBroadcast
// and the chain's gov vote query.
broadcast func(ctx context.Context, cfg engine.ExecutionConfig, in SignAndBroadcastInput) (*SignAndBroadcastResult, error)
readVote func(ctx context.Context, cfg engine.ExecutionConfig, chainID string, proposalID uint64, voter sdk.AccAddress) (*govtypes.Vote, error)

// confirmWait bounds how long a vote whose tx the node cannot look up
// waits for the vote to show in state, reads included; confirmEvery is
// the poll interval.
confirmWait time.Duration
confirmEvery time.Duration
}

func NewGovVoter(cfg engine.ExecutionConfig) *GovVoter {
return &GovVoter{cfg: cfg}
return &GovVoter{
cfg: cfg,
broadcast: SignAndBroadcast,
readVote: chainVote,
confirmWait: 30 * time.Second,
confirmEvery: time.Second,
}
}

// Handler delegates to SignAndBroadcast after MsgVote construction.
Expand All @@ -62,7 +82,7 @@ func (g *GovVoter) Handler() engine.TaskHandler {
if err != nil {
return nil, err
}
result, err := SignAndBroadcast(ctx, g.cfg, SignAndBroadcastInput{
result, err := g.broadcast(ctx, g.cfg, SignAndBroadcastInput{
ChainID: params.ChainID,
KeyName: params.KeyName,
Msg: msg,
Expand All @@ -75,14 +95,19 @@ func (g *GovVoter) Handler() engine.TaskHandler {
return nil, err
}
out, cerr := classifyGovResult(engine.TaskGovVote, result)
votedByState := result.Unverifiable && g.votedAfter(ctx, params.ChainID, msg)
if votedByState {
cerr = nil
}
govVoteLog.Info("vote broadcast",
"taskId", engine.TaskIDFromContext(ctx),
"chainId", params.ChainID,
"proposalId", params.ProposalID,
"option", params.Option,
"txHash", out.TxHash,
"height", out.Height,
"inclusionStatus", out.InclusionStatus)
"inclusionStatus", out.InclusionStatus,
"votedByState", votedByState)
return out, cerr
})
}
Expand All @@ -107,3 +132,95 @@ func buildVoteMsg(cfg engine.ExecutionConfig, params GovVoteRequest) (*govtypes.
}
return govtypes.NewMsgVote(info.GetAddress(), params.ProposalID, govtypes.VoteOption(option)), nil
}

// votedAfter settles a vote whose tx the node cannot look up because its tx
// index is off, as on most validators. The vote's effect shows in state: the
// gov module records the voter's choice on the proposal. A vote that reads the
// requested option landed; an earlier vote with the same option reads the
// same, and the recorded choice is right either way. It polls until
// confirmWait passes; the deadline also bounds each read.
func (g *GovVoter) votedAfter(ctx context.Context, chainID string, msg *govtypes.MsgVote) bool {
voter, err := sdk.AccAddressFromBech32(msg.Voter)
if err != nil {
return false
}
ctx, cancel := context.WithTimeout(ctx, g.confirmWait)
defer cancel()
var lastErr error
for {
v, err := g.readVote(ctx, g.cfg, chainID, msg.ProposalId, voter)
if err == nil && voteHasOption(v, msg.Option) {
Comment thread
seidroid[bot] marked this conversation as resolved.
return true
}
if err != nil {
lastErr = err
}
select {
case <-ctx.Done():
govVoteLog.Warn("vote not confirmed from the gov vote query",
"proposalId", msg.ProposalId, "voter", msg.Voter, "wait", g.confirmWait, "lastReadErr", lastErr)
return false
case <-time.After(g.confirmEvery):
}
}
}

// voteHasOption reports whether v records option as its one full-weight
// choice, which is how the gov keeper stores a MsgVote.
func voteHasOption(v *govtypes.Vote, option govtypes.VoteOption) bool {
if v == nil || len(v.Options) != 1 {
return false
}
o := v.Options[0]
return o.Option == option && o.Weight.Equal(sdk.OneDec())
}

// chainVote reads the voter's recorded vote on proposalID from the local seid
// over gRPC-over-ABCI. The SDK query path does not honor ctx, so the read runs
// off-goroutine and ctx cancellation returns at once, as in chainJailState.
func chainVote(ctx context.Context, cfg engine.ExecutionConfig, chainID string, proposalID uint64, voter sdk.AccAddress) (*govtypes.Vote, error) {
clientCtx, err := newSignTxClientContext(cfg, SignAndBroadcastInput{ChainID: chainID}, voter)
if err != nil {
return nil, err
}
type res struct {
v *govtypes.Vote
err error
}
ch := make(chan res, 1)
go func() {
v, err := readVoteState(ctx, clientCtx.Client.Status, govtypes.NewQueryClient(clientCtx), proposalID, voter)
ch <- res{v, err}
}()
select {
case <-ctx.Done():
return nil, ctx.Err()
case r := <-ch:
return r.v, r.err
}
}

// readVoteState reads the voter's recorded vote through narrow seams, so a test
// can fake each read. A node that is catching up answers from an old height,
// where an earlier vote may still show; it reports an error instead, and the
// confirmation keeps polling.
func readVoteState(
ctx context.Context,
statusOf func(context.Context) (*coretypes.ResultStatus, error),
gov govtypes.QueryClient,
proposalID uint64,
voter sdk.AccAddress,
) (*govtypes.Vote, error) {
s, err := statusOf(ctx)
if err != nil {
return nil, fmt.Errorf("query local seid /status: %w", err)
}
if s.SyncInfo.CatchingUp {
return nil, errors.New("local seid is catching up, so its vote record may be stale")
}
r, err := gov.Vote(ctx, &govtypes.QueryVoteRequest{ProposalId: proposalID, Voter: voter.String()})
if err != nil {
return nil, err
}
return &r.Vote, nil
}
134 changes: 134 additions & 0 deletions sidecar/tasks/gov_vote_test.go
Original file line number Diff line number Diff line change
@@ -1,12 +1,20 @@
package tasks

import (
"context"
"encoding/json"
"errors"
"strings"
"testing"
"time"

sdk "github.com/sei-protocol/sei-chain/sei-cosmos/types"
govtypes "github.com/sei-protocol/sei-chain/sei-cosmos/x/gov/types"

"google.golang.org/grpc"

"github.com/sei-protocol/sei-k8s-controller/sidecar/engine"
"github.com/sei-protocol/sei-k8s-controller/sidecarapi/wire"
)

func TestBuildVoteMsg(t *testing.T) {
Expand Down Expand Up @@ -98,3 +106,129 @@ func TestBuildVoteMsg(t *testing.T) {
}
})
}

// govVoteHarness wires a GovVoter with a fake broadcast and a fake vote read
// that returns votes in order, then repeats the last one.
type govVoteHarness struct {
g *GovVoter
reads int
}

func newGovVoteHarness(t *testing.T, result *SignAndBroadcastResult, votes ...*govtypes.Vote) *govVoteHarness {
t.Helper()
kr, _ := testKeyring(t)
h := &govVoteHarness{}
h.g = &GovVoter{
cfg: engine.ExecutionConfig{Keyring: kr},
broadcast: func(context.Context, engine.ExecutionConfig, SignAndBroadcastInput) (*SignAndBroadcastResult, error) {
return result, nil
},
readVote: func(context.Context, engine.ExecutionConfig, string, uint64, sdk.AccAddress) (*govtypes.Vote, error) {
h.reads++
if len(votes) == 0 {
return nil, errors.New("voter not found for proposal")
}
return votes[min(h.reads-1, len(votes)-1)], nil
},
confirmWait: 50 * time.Millisecond,
confirmEvery: time.Millisecond,
}
return h
}

func runGovVote(t *testing.T, g *GovVoter, option string) (*wire.GovTxResult, error) {
t.Helper()
ctx := engine.WithTaskID(context.Background(), "gov-vote-test")
raw, err := g.Handler()(ctx, map[string]any{
"chainId": "sei-test", "keyName": "node_admin", "proposalId": 7, "option": option, "fees": "4000usei", "gas": 200000,
})
if len(raw) == 0 || string(raw) == "null" {
return nil, err
}
var out wire.GovTxResult
if uerr := json.Unmarshal(raw, &out); uerr != nil {
t.Fatalf("decode result: %v", uerr)
}
return &out, err
}

func recordedVote(option govtypes.VoteOption) *govtypes.Vote {
return &govtypes.Vote{ProposalId: 7, Options: govtypes.NewNonSplitVoteOption(option)}
}

// PLT-1401: validators often run with the tx index off, so the node cannot
// look up the vote tx. The task then confirms the vote from the gov module's
// recorded choice: Complete once the vote reads the requested option.
func TestGovVoteUnverifiableTxConfirmedByVoteQuery(t *testing.T) {
h := newGovVoteHarness(t, &SignAndBroadcastResult{TxHash: "ABCD", Unverifiable: true},
nil, recordedVote(govtypes.OptionYes))
out, err := runGovVote(t, h.g, "yes")
if err != nil {
t.Fatalf("unexpected err: %v", err)
}
if out == nil || out.TxHash != "ABCD" || out.InclusionStatus != wire.InclusionUnverifiable {
t.Errorf("result = %+v; want the tx hash, still marked unverifiable", out)
}
if h.reads < 2 {
t.Errorf("read the vote %d times, want at least 2 (not yet recorded, then recorded)", h.reads)
}
}

// A recorded vote with a different option is not this vote: the task keeps
// the unverifiable failure, and the operator checks the tx.
func TestGovVoteUnverifiableTxDifferentOptionStaysUnverifiable(t *testing.T) {
h := newGovVoteHarness(t, &SignAndBroadcastResult{TxHash: "ABCD", Unverifiable: true}, recordedVote(govtypes.OptionNo))
out, err := runGovVote(t, h.g, "yes")
if !IsTerminal(err) || !strings.Contains(err.Error(), "inclusion unverifiable") {
t.Fatalf("want terminal inclusion-unverifiable error, got %v", err)
}
if out == nil || out.TxHash != "ABCD" {
t.Errorf("result = %+v", out)
}
}

// A vote the node can look up never reaches the vote query.
func TestGovVoteCommittedTxSkipsVoteQuery(t *testing.T) {
now := time.Now()
h := newGovVoteHarness(t, &SignAndBroadcastResult{TxHash: "ABCD", Height: 9, IncludedAt: &now})
if _, err := runGovVote(t, h.g, "yes"); err != nil {
t.Fatalf("unexpected err: %v", err)
}
if h.reads != 0 {
t.Errorf("read the vote %d times, want 0", h.reads)
}
}

type fakeGovQuery struct {
govtypes.QueryClient
vote *govtypes.Vote
calls int
}

func (f *fakeGovQuery) Vote(context.Context, *govtypes.QueryVoteRequest, ...grpc.CallOption) (*govtypes.QueryVoteResponse, error) {
f.calls++
return &govtypes.QueryVoteResponse{Vote: *f.vote}, nil
}

// seidroid on #604: a node that is catching up can still show an older vote,
// so the read refuses to trust it, as the Unjail jail-state read does.
func TestReadVoteState(t *testing.T) {
_, voter := testKeyring(t)
gov := &fakeGovQuery{vote: recordedVote(govtypes.OptionYes)}

if _, err := readVoteState(context.Background(), statusAt(testBlockTime, true), gov, 7, voter); err == nil ||
!strings.Contains(err.Error(), "catching up") {
t.Fatalf("catching up: want an error, got %v", err)
}
if gov.calls != 0 {
t.Errorf("queried the vote %d times while catching up, want 0", gov.calls)
}

v, err := readVoteState(context.Background(), statusAt(testBlockTime, false), gov, 7, voter)
if err != nil {
t.Fatalf("caught up: unexpected err: %v", err)
}
if !voteHasOption(v, govtypes.OptionYes) {
t.Errorf("vote = %+v, want the recorded yes", v)
}
}
9 changes: 5 additions & 4 deletions sidecarapi/wire/wire.go
Original file line number Diff line number Diff line change
Expand Up @@ -103,10 +103,11 @@ const (
// its on-chain outcome cannot be observed from the target node (its tx
// index is disabled), so the sidecar can confirm neither success nor
// failure. Terminal — retrying the same node is futile — but distinct from
// committed_failed: the operator must verify via an indexed RPC. One
// exception: an unjail confirms its effect from the validator's jail state,
// so an Unjail task can complete with this status once the validator reads
// released. Key completion on the task phase, not on this value alone.
// committed_failed: the operator must verify via an indexed RPC. Two
// exceptions confirm their effect from state instead, so the task can
// complete with this status: an Unjail once the validator reads released,
// and a GovVote once the gov module records the requested option. Key
// completion on the task phase, not on this value alone.
InclusionUnverifiable = "unverifiable"
)

Expand Down
Loading