Skip to content
Open
Show file tree
Hide file tree
Changes from 3 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
46 changes: 46 additions & 0 deletions docs/architecture/service-lifecycle.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,46 @@
# Service lifecycle core

`BaseService.Start` gives `OnStart` an attempt-scoped child context. Parent
cancellation and manual `Stop` both cancel it. `Go(ctx, fn)` registers workers
using that context or a derived context; registration is synchronized with
shutdown and rejected for a different attempt or after shutdown begins.

Shutdown closes admission, marks `IsRunning` false, cancels the context and
calls `OnStop`. It then joins registered workers, calls optional `OnDrain`, and
releases `Wait`. `Stopping()` signals cancellation separately from completion.
`Stop` invokes `OnStop` synchronously but does not join registered workers.
Concurrent or recursive calls to `Stop` return immediately.

Hooks run outside the lifecycle mutex. `OnStop` must unblock workers without
waiting for its own registered work. Workers may call `Stop`; neither hooks nor
workers may call their own service's `Wait`. Release resources shared with
workers in `OnDrain`, or after `Stop` followed by `Wait`.

`IsRunning` is false during startup and shutdown. `Wait` joins the current
attempt, including startup and cleanup; it returns immediately before startup.
Failed `OnStart` cancels and joins its workers before retry becomes possible.
The implementation must roll back acquired resources: failed attempts do not
invoke `OnStop` or `OnDrain`. Successfully stopped services cannot restart.
Composite retries also require recreating children already successfully stopped.

## First migration stage

Consensus State registers its receive loop and joins its queue fan-in and WAL
before returning. The timeout ticker, WAL and autofile group use shared worker
tracking. AutoFile explicitly joins its timer/signal worker. WAL ownership lasts
until the receive loop's final write. Consensus handoff admission is also joined,
so a late blocksync handoff cannot start State after reactor shutdown.

Node waits for reactors before closing event sinks and stores. The blocked
receive-loop regression and original node teardown tests verify this path
without using leaktest as an extra production-shutdown barrier.

This is a partial migration. Plain goroutines are not tracked automatically.
RPC/WebSocket handlers, mempool rechecks, reactor gossip and other service
workers retain their existing ownership mechanisms. This change does not claim
that Node.Wait joins every helper or introduce the broader durable-finalization
and dependency-lifetime changes from #1515.

Two compatibility bridges preserve existing parent-context lifetimes: blocksync
application during handoff and connection error callbacks after I/O shutdown.
Their local cancellation and join mechanisms remain in place in this stage.
10 changes: 9 additions & 1 deletion internal/blocksync/synchronizer.go
Original file line number Diff line number Diff line change
Expand Up @@ -236,6 +236,13 @@ func (s *Synchronizer) setStartHeight(height int64) {
s.jobGen.mtx.Unlock()
}

type applicationContextKey struct{}

// Start preserves the caller's application lifetime across a block-sync handover.
func (s *Synchronizer) Start(ctx context.Context) error {
return s.BaseService.Start(context.WithValue(ctx, applicationContextKey{}, ctx))
}

// OnStart implements service.Service by spawning requesters routine and recording
// synchronizer's start time.
func (s *Synchronizer) OnStart(ctx context.Context) error {
Expand All @@ -245,14 +252,15 @@ func (s *Synchronizer) OnStart(ctx context.Context) error {
s.lastAdvance = s.clock.Now()
s.lastMonitorUpdate = s.lastAdvance
s.ctx, s.cancel = context.WithCancel(ctx)
applicationCtx := ctx.Value(applicationContextKey{}).(context.Context)
Comment thread
lklimek marked this conversation as resolved.
Outdated
s.consumerDone = make(chan struct{})
s.workerPool.Run(s.ctx)
go s.runHandler(s.ctx, s.produceJob)
go func() {
defer close(s.consumerDone)
s.runHandler(s.ctx, func(handlerCtx context.Context) error {
// Handover cancels consumer I/O; only node shutdown may cancel application.
return s.consumeJobResult(handlerCtx, ctx)
return s.consumeJobResult(handlerCtx, applicationCtx)
})
}()
return nil
Expand Down
1 change: 1 addition & 0 deletions internal/consensus/common_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -1192,6 +1192,7 @@ type mockTicker struct {
}

func (m *mockTicker) Start(context.Context) error { return nil }
func (m *mockTicker) Wait() {}
func (m *mockTicker) Stop() {}
func (m *mockTicker) IsRunning() bool { return false }

Expand Down
40 changes: 24 additions & 16 deletions internal/consensus/reactor.go
Original file line number Diff line number Diff line change
Expand Up @@ -413,20 +413,20 @@ func (r *Reactor) OnStart(ctx context.Context) error {
return nil
}

// OnStop stops the reactor by signaling to all spawned goroutines to exit and
// blocking until they all exit, as well as unsubscribing from events and stopping
// state.
// OnStop requests reactor and consensus shutdown.
func (r *Reactor) OnStop() {
// Cancel the reactor context to signal all goroutines to stop
if r.cancel != nil {
r.cancel()
}

r.state.Stop()
}

if !r.WaitSync() {
r.state.Wait()
}
// OnDrain joins consensus after any admitted handoff has finished.
func (r *Reactor) OnDrain() {
r.state.Stop()
r.state.Wait()
}

// WaitSync returns whether the consensus reactor is waiting for state/block sync.
Expand All @@ -447,6 +447,20 @@ func (r *Reactor) WaitSync() bool {
// targetHeight is the highest committed block height reported during block sync.
// skipWAL says the node needs no WAL catchup.
func (r *Reactor) SwitchToConsensus(ctx context.Context, state sm.State, skipWAL bool, targetHeight int64) {
Comment thread
lklimek marked this conversation as resolved.
if r.ctx == nil || ctx.Err() != nil {
return
}
done := make(chan struct{})
if !r.Go(r.ctx, func(ctx context.Context) {
defer close(done)
r.switchToConsensus(ctx, state, skipWAL, targetHeight)
}) {
return
}
<-done
}

func (r *Reactor) switchToConsensus(ctx context.Context, state sm.State, skipWAL bool, targetHeight int64) {
r.logger.Info("switching to consensus", "target_height", targetHeight)

if targetHeight > state.LastBlockHeight {
Expand All @@ -473,6 +487,9 @@ func (r *Reactor) SwitchToConsensus(ctx context.Context, state sm.State, skipWAL
r.state.eventPublisher.PublishNewRoundStepEvent(stateData.RoundState)

if err := r.state.Start(ctx); err != nil {
if ctx.Err() != nil {
return
}
panic(fmt.Sprintf(`failed to start consensus state: %v

conS:
Expand Down Expand Up @@ -520,8 +537,6 @@ func (r *Reactor) GetPeerState(peerID types.NodeID) (*PeerState, bool) {
// internal pubsub defined in the consensus state to broadcast them to peers
// upon receiving.
func (r *Reactor) subscribeToBroadcastEvents(ctx context.Context, stateCh p2p.Channel) {
onStopCh := r.state.getOnStopCh()

r.state.emitter.AddListener(
types.EventNewRoundStepValue,
func(data eventemitter.EventData) error {
Expand All @@ -531,14 +546,7 @@ func (r *Reactor) subscribeToBroadcastEvents(ctx context.Context, stateCh p2p.Ch
return err
}
r.logResult(err, r.logger, "broadcasting round step message", "height", rs.Height, "round", rs.Round)
select {
case onStopCh <- data.(*cstypes.RoundState):
return nil
case <-ctx.Done():
return ctx.Err()
default:
return nil
}
return nil
},
)

Expand Down
14 changes: 14 additions & 0 deletions internal/consensus/reactor_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -834,3 +834,17 @@ func TestReactorPeerDownDeletesPeerStateSynchronously(t *testing.T) {
"peerDown must delete the peer state before returning, so a fast reconnect cannot reuse it")
require.False(t, ps.IsRunning(), "peerDown must stop the peer state before returning")
}

func TestReactorRejectsConsensusHandoffAfterShutdown(t *testing.T) {
ctx, cancel := context.WithCancel(context.Background())
defer cancel()
cs, _ := makeState(ctx, t, makeStateArgs{})
rts := setup(ctx, t, 1, []*State{cs}, 32)
for _, r := range rts.reactors {
r.Stop()
r.Wait()
r.SwitchToConsensus(ctx, cs.GetStateData().state, false, 0)
require.False(t, cs.IsRunning(), "handoff started consensus after reactor shutdown")
require.True(t, r.WaitSync(), "rejected handoff changed sync state")
}
}
4 changes: 0 additions & 4 deletions internal/consensus/replay_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -1521,10 +1521,6 @@ func testWALRoundsSkipper(t *testing.T, slowProposer bool) {
<-replayStopped
cs.Stop()
cs.Wait()
// The stop predicate bypasses receiveRoutine's WAL and queue shutdown.
cs.wal.Stop()
cs.wal.Wait()
cs.msgInfoQueue.stop()
}()

newBlockSub, err := cs.eventBus.SubscribeWithArgs(ctx, pubsub.SubscribeArgs{
Expand Down
71 changes: 31 additions & 40 deletions internal/consensus/state.go
Original file line number Diff line number Diff line change
Expand Up @@ -191,9 +191,6 @@ type State struct {
// proposer's latest available app protocol version that goes to block header
proposedAppVersion uint64

// wait the channel event happening for shutting down the state gracefully
onStopCh chan *cstypes.RoundState

msgInfoQueue *msgInfoQueue
msgDispatcher *msgInfoDispatcher
blockExecutor *blockExecutor
Expand Down Expand Up @@ -285,7 +282,6 @@ func NewState(
evpool: evpool,
emitter: eventemitter.New(eventemitter.WithLogger(logger)),
metrics: NopMetrics(),
onStopCh: make(chan *cstypes.RoundState),
verificationBudget: newVerificationBudget(
cfg.VerificationRateLimit,
),
Expand Down Expand Up @@ -471,7 +467,20 @@ func (cs *State) SetPrivValidator(ctx context.Context, priv types.PrivValidator)

// OnStart loads the latest state via the WAL, and starts the timeout and
// receive routines.
func (cs *State) OnStart(ctx context.Context) error {
func (cs *State) OnStart(ctx context.Context) (err error) {
defer func() {
if err != nil {
cs.wal.Stop()
cs.wal.Wait()
cs.wal = nilWAL{}
cs.timeoutTicker.Stop()
cs.timeoutTicker.Wait()
if _, ok := cs.timeoutTicker.(*timeoutTicker); ok {
cs.timeoutTicker = NewTimeoutTicker(cs.logger)
cs.roundScheduler.timeoutTicker = cs.timeoutTicker
}
}
}()
if err := cs.updateStateFromStore(); err != nil {
return err
}
Expand Down Expand Up @@ -518,6 +527,7 @@ func (cs *State) OnStart(ctx context.Context) error {

// 1) prep work
cs.wal.Stop()
cs.wal.Wait()

repairAttempted = true

Expand Down Expand Up @@ -545,7 +555,9 @@ func (cs *State) OnStart(ctx context.Context) error {
}

// now start the receiveRoutine
go cs.receiveRoutine(ctx, cs.stopFn)
if !cs.Go(ctx, func(ctx context.Context) { cs.receiveRoutine(ctx, cs.stopFn) }) {
return context.Canceled
}

// schedule the first round!
// use GetRoundState so we don't race the receiveRoutine for access
Expand All @@ -556,7 +568,8 @@ func (cs *State) OnStart(ctx context.Context) error {

// loadWalFile loads WAL data from file. It overwrites cs.wal.
func (cs *State) loadWalFile(ctx context.Context) error {
wal, err := cs.OpenWAL(ctx, cs.config.WalFile())
// receiveRoutine owns the WAL until its final write and explicit Stop/Wait.
wal, err := cs.OpenWAL(context.WithoutCancel(ctx), cs.config.WalFile())
if err != nil {
cs.logger.Error("failed to load state WAL", "err", err)
return err
Expand All @@ -566,31 +579,11 @@ func (cs *State) loadWalFile(ctx context.Context) error {
return nil
}

func (cs *State) getOnStopCh() chan *cstypes.RoundState {
cs.mtx.RLock()
defer cs.mtx.RUnlock()
// OnStop requests timeout worker shutdown. Wait also joins receiveRoutine.
func (cs *State) OnStop() { cs.timeoutTicker.Stop() }
Comment thread
lklimek marked this conversation as resolved.

return cs.onStopCh
}

// OnStop implements service.Service.
func (cs *State) OnStop() {
stateData := cs.stateDataStore.Get()
// If the node is committing a new block, wait until it is finished!
if cs.GetRoundState().Step == cstypes.RoundStepApplyCommit {
select {
case <-cs.getOnStopCh():
case <-time.After(stateData.state.ConsensusParams.Timeout.Vote):
// we wait vote timeout, just in case
cs.logger.Error("OnStop: timeout waiting for commit to finish", "time", stateData.state.ConsensusParams.Timeout.Vote)
}
}

if cs.timeoutTicker.IsRunning() {
cs.timeoutTicker.Stop()
}
// WAL is stopped in receiveRoutine.
}
// OnDrain joins the timeout worker before callers release consensus stores.
func (cs *State) OnDrain() { cs.timeoutTicker.Wait() }

// OpenWAL opens a file to log all consensus messages and timeouts for
// deterministic accountability.
Expand All @@ -603,6 +596,7 @@ func (cs *State) OpenWAL(ctx context.Context, walFile string) (WAL, error) {

if err := wal.Start(ctx); err != nil {
cs.logger.Error("failed to start WAL", "err", err)
wal.Group().Close()
return nil, err
}

Expand Down Expand Up @@ -700,12 +694,6 @@ func (cs *State) receiveRoutine(ctx context.Context, stopFn func(*State) bool) {
if r := recover(); r != nil {
cs.logger.Error("CONSENSUS FAILURE!!!", "err", r, "stack", string(debug.Stack()))

// Make a best-effort attempt to close the WAL, but otherwise do not
// attempt to gracefully terminate. Once consensus has irrecoverably
// failed, any additional progress we permit the node to make may
// complicate diagnosing and recovering from the failure.
onExit(cs)

// There are a couple of cases where the we
// panic with an error from deeper within the
// state machine and in these cases, typically
Expand Down Expand Up @@ -739,7 +727,12 @@ func (cs *State) receiveRoutine(ctx context.Context, stopFn func(*State) bool) {
}
}()

go cs.msgInfoQueue.fanIn(ctx)
queueDone := make(chan struct{})
go func() {
defer close(queueDone)
cs.msgInfoQueue.fanIn(ctx)
}()
defer func() { onExit(cs); <-queueDone }()

for {
if stopFn != nil && stopFn(cs) {
Expand Down Expand Up @@ -769,7 +762,6 @@ func (cs *State) receiveRoutine(ctx context.Context, stopFn func(*State) bool) {
if err != nil {
// Leaving without this would strand the queue's reader
// goroutines: nothing else ever tells them to stop.
onExit(cs)
return
}
err = stateData.Save()
Expand All @@ -790,7 +782,6 @@ func (cs *State) receiveRoutine(ctx context.Context, stopFn func(*State) bool) {
cs.logger.Error("failed update state-data", "err", err)
}
case <-ctx.Done():
onExit(cs)
return

}
Expand Down
Loading
Loading