(redisOpt asynq.RedisClientOpt, store db.Store, mailer mail.EmailSender)
| 28 | } |
| 29 | |
| 30 | func NewRedisTaskProcessor(redisOpt asynq.RedisClientOpt, store db.Store, mailer mail.EmailSender) TaskProcessor { |
| 31 | logger := NewLogger() |
| 32 | redis.SetLogger(logger) |
| 33 | |
| 34 | server := asynq.NewServer( |
| 35 | redisOpt, |
| 36 | asynq.Config{ |
| 37 | Queues: map[string]int{ |
| 38 | QueueCritical: 10, |
| 39 | QueueDefault: 5, |
| 40 | }, |
| 41 | ErrorHandler: asynq.ErrorHandlerFunc(func(ctx context.Context, task *asynq.Task, err error) { |
| 42 | log.Error().Err(err).Str("type", task.Type()). |
| 43 | Bytes("payload", task.Payload()).Msg("process task failed") |
| 44 | }), |
| 45 | Logger: logger, |
| 46 | }, |
| 47 | ) |
| 48 | |
| 49 | return &RedisTaskProcessor{ |
| 50 | server: server, |
| 51 | store: store, |
| 52 | mailer: mailer, |
| 53 | } |
| 54 | } |
| 55 | |
| 56 | func (processor *RedisTaskProcessor) Start() error { |
| 57 | mux := asynq.NewServeMux() |
no test coverage detected