| 1007 | } |
| 1008 | |
| 1009 | func (l *validator) stateMachine() error { |
| 1010 | defer l.cancel() |
| 1011 | |
| 1012 | var err error |
| 1013 | |
| 1014 | var sub pubsub.Subscriber |
| 1015 | |
| 1016 | sub, err = l.pubsub.Subscribe() |
| 1017 | if err != nil { |
| 1018 | return err |
| 1019 | } |
| 1020 | |
| 1021 | blocksCount := 0 |
| 1022 | replayDone := false |
| 1023 | stage := nodeTestStagePreUpgrade |
| 1024 | |
| 1025 | wdCtrl := func(ctx context.Context, ctrl watchdogCtrl) { |
| 1026 | resp := make(chan struct{}, 1) |
| 1027 | _ = l.pubsub.Publish(wdReq{ |
| 1028 | event: ctrl, |
| 1029 | resp: resp, |
| 1030 | }) |
| 1031 | |
| 1032 | select { |
| 1033 | case <-ctx.Done(): |
| 1034 | case <-resp: |
| 1035 | } |
| 1036 | } |
| 1037 | |
| 1038 | loop: |
| 1039 | for { |
| 1040 | select { |
| 1041 | case <-l.ctx.Done(): |
| 1042 | err = l.ctx.Err() |
| 1043 | break loop |
| 1044 | case ev := <-sub.Events(): |
| 1045 | switch evt := ev.(type) { |
| 1046 | case event: |
| 1047 | switch evt.id { |
| 1048 | case nodeEventStart: |
| 1049 | l.t.Logf("[%s][%s]: node started", l.params.name, nodeTestStageMapStr[stage]) |
| 1050 | if stage == nodeTestStageUpgrade { |
| 1051 | stage = nodeTestStagePostUpgrade1 |
| 1052 | blocksCount = 0 |
| 1053 | replayDone = false |
| 1054 | } |
| 1055 | case nodeEventReplayBlocksStart: |
| 1056 | l.t.Logf("[%s][%s]: node started replaying blocks", l.params.name, nodeTestStageMapStr[stage]) |
| 1057 | case nodeEventReplayBlocksDone: |
| 1058 | l.t.Logf("[%s][%s]: node done replaying blocks", l.params.name, nodeTestStageMapStr[stage]) |
| 1059 | wdCtrl(l.ctx, watchdogCtrlStart) |
| 1060 | replayDone = true |
| 1061 | case nodeEventBlockCommited: |
| 1062 | // ignore index events until replay done |
| 1063 | if !replayDone { |
| 1064 | break |
| 1065 | } |
| 1066 | |