(d *Driver)
| 178 | } |
| 179 | |
| 180 | func (h *taskHandle) NewRunner(d *Driver) (runner DriverHandle, err error) { |
| 181 | ctx := &common.ExecContext{ |
| 182 | Subject: h.taskConfig.JobName, |
| 183 | StateDir: d.config.DataDir, |
| 184 | } |
| 185 | |
| 186 | logger := h.logger |
| 187 | // create a new log file for debug job |
| 188 | if d.config.DebugJob == h.taskConfig.JobName && hclog.LevelFromString(d.config.LogLevel) != hclog.Debug { |
| 189 | logger = setupLogger(d.config.LogFile, "dtle_debug_job.log") |
| 190 | logger.SetLevel(hclog.Debug) |
| 191 | } |
| 192 | |
| 193 | switch common.TaskTypeFromString(h.taskConfig.Name) { |
| 194 | case common.TaskTypeSrc: |
| 195 | if h.driverConfig.SrcOracleConfig != nil { |
| 196 | h.logger.Debug("found oracle src", "SrcOracleConfig", h.driverConfig.SrcOracleConfig) |
| 197 | runner, err = extractor.NewExtractorOracle(ctx, h.driverConfig, logger, d.storeManager, h.waitCh, h.ctx) |
| 198 | if err != nil { |
| 199 | return nil, errors.Wrap(err, "NewExtractor") |
| 200 | } |
| 201 | } else { |
| 202 | e, err := mysql.NewExtractor(ctx, h.driverConfig, logger, d.storeManager, h.waitCh, h.ctx) |
| 203 | if err != nil { |
| 204 | return nil, errors.Wrap(err, "NewOracleExtractor") |
| 205 | } |
| 206 | runner = e |
| 207 | if h.driverConfig.TwoWaySync { |
| 208 | ctx2 := &common.ExecContext{ |
| 209 | Subject: ctx.Subject + "_dtrev", |
| 210 | StateDir: d.config.DataDir, |
| 211 | } |
| 212 | cfg2 := &common.MySQLDriverConfig{ |
| 213 | DtleTaskConfig: common.DtleTaskConfig{ |
| 214 | DestType: "mysql", |
| 215 | }, |
| 216 | } |
| 217 | e.RevApplier, err = mysql.NewApplier(ctx2, cfg2, logger, d.storeManager, d.config.NatsAdvertise, |
| 218 | h.waitCh, d.eventer, h.taskConfig, h.ctx) |
| 219 | } |
| 220 | |
| 221 | } |
| 222 | case common.TaskTypeDest: |
| 223 | h.logger.Debug("found dest", "allConfig", h.driverConfig) |
| 224 | switch strings.ToLower(h.driverConfig.DestType) { |
| 225 | case "kafka": |
| 226 | runner, err = kafka.NewKafkaRunner(ctx, logger, |
| 227 | d.storeManager, d.config.NatsAdvertise, h.waitCh, h.ctx) |
| 228 | if err != nil { |
| 229 | return nil, errors.Wrap(err, "NewKafkaRunner") |
| 230 | } |
| 231 | case "mysql": |
| 232 | runner, err = mysql.NewApplier(ctx, h.driverConfig, logger, d.storeManager, |
| 233 | d.config.NatsAdvertise, h.waitCh, d.eventer, h.taskConfig, h.ctx) |
| 234 | if err != nil { |
| 235 | return nil, errors.Wrap(err, "NewApplier") |
| 236 | } |
| 237 | case "": |
no test coverage detected