Run executes the complete extract logic.
()
| 154 | |
| 155 | // Run executes the complete extract logic. |
| 156 | func (e *ExtractorOracle) Run() { |
| 157 | var err error |
| 158 | |
| 159 | // PutConfig before WatchNats |
| 160 | err = e.storeManager.PutConfig(e.subject, e.mysqlContext) |
| 161 | if err != nil { |
| 162 | e.onError(common.TaskStateDead, errors.Wrap(err, "PutConfig")) |
| 163 | return |
| 164 | } |
| 165 | |
| 166 | e.logger.Info("src watch Nats") |
| 167 | e.natsAddr, err = e.storeManager.SrcWatchNats(e.subject, e.shutdownCh, func(err error) { |
| 168 | e.onError(common.TaskStateDead, err) |
| 169 | }) |
| 170 | if err != nil { |
| 171 | e.onError(common.TaskStateDead, errors.Wrap(err, "SrcWatchNats")) |
| 172 | return |
| 173 | } |
| 174 | // init nats |
| 175 | e.logger.Info("initNatsPubClient") |
| 176 | e.logger.Debug("begin Connect nats server", "NatAddr", e.natsAddr) |
| 177 | sc, err := gonats.Connect(e.natsAddr) |
| 178 | if err != nil { |
| 179 | e.logger.Error("cannot connect nats server", "natsAddr", e.natsAddr, "err", err) |
| 180 | e.onError(common.TaskStateDead, errors.Wrap(err, "Connect Nats")) |
| 181 | return |
| 182 | } |
| 183 | e.logger.Info("Connect nats server", "natsAddr", e.natsAddr) |
| 184 | e.natsConn = sc |
| 185 | |
| 186 | _, err = e.natsConn.Subscribe(fmt.Sprintf("%s_control2", e.subject), func(m *gonats.Msg) { |
| 187 | if m.Data == nil { |
| 188 | e.onError(common.TaskStateDead, fmt.Errorf("zero-byte control msg")) |
| 189 | return |
| 190 | } |
| 191 | |
| 192 | ctrlMsg := &common.ControlMsg{} |
| 193 | _, err := ctrlMsg.Unmarshal(m.Data) |
| 194 | if err != nil { |
| 195 | e.onError(common.TaskStateDead, fmt.Errorf("failed to unmarshal a control msg")) |
| 196 | return |
| 197 | } |
| 198 | |
| 199 | switch ctrlMsg.Type { |
| 200 | case common.ControlMsgError: |
| 201 | e.onError(common.TaskStateDead, fmt.Errorf("applier error: %v", ctrlMsg.Msg)) |
| 202 | return |
| 203 | } |
| 204 | }) |
| 205 | if err != nil { |
| 206 | e.onError(common.TaskStateDead, errors.Wrap(err, "Subscribe control2")) |
| 207 | return |
| 208 | } |
| 209 | |
| 210 | e.initDBConnections() |
| 211 | |
| 212 | startSCN, committedSCN, err := e.calculateSCNPos() |
| 213 | if err != nil { |
nothing calls this directly
no test coverage detected