diff --git a/api/v1alpha1/seinodetask_types.go b/api/v1alpha1/seinodetask_types.go index f377d5b0..0bdd03b6 100644 --- a/api/v1alpha1/seinodetask_types.go +++ b/api/v1alpha1/seinodetask_types.go @@ -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. diff --git a/sidecar/tasks/gov_result.go b/sidecar/tasks/gov_result.go index f1cf83a0..afa43367 100644 --- a/sidecar/tasks/gov_result.go +++ b/sidecar/tasks/gov_result.go @@ -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)) diff --git a/sidecar/tasks/gov_vote.go b/sidecar/tasks/gov_vote.go index 8f7686ae..6195b18b 100644 --- a/sidecar/tasks/gov_vote.go +++ b/sidecar/tasks/gov_vote.go @@ -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" @@ -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. @@ -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, @@ -75,6 +95,10 @@ 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, @@ -82,7 +106,8 @@ func (g *GovVoter) Handler() engine.TaskHandler { "option", params.Option, "txHash", out.TxHash, "height", out.Height, - "inclusionStatus", out.InclusionStatus) + "inclusionStatus", out.InclusionStatus, + "votedByState", votedByState) return out, cerr }) } @@ -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) { + 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 +} diff --git a/sidecar/tasks/gov_vote_test.go b/sidecar/tasks/gov_vote_test.go index 7e592d65..89479556 100644 --- a/sidecar/tasks/gov_vote_test.go +++ b/sidecar/tasks/gov_vote_test.go @@ -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) { @@ -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) + } +} diff --git a/sidecarapi/wire/wire.go b/sidecarapi/wire/wire.go index 63afad48..f62d9d41 100644 --- a/sidecarapi/wire/wire.go +++ b/sidecarapi/wire/wire.go @@ -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" )