2020-08-13 03:54:14 +02:00
|
|
|
package commands
|
|
|
|
|
|
|
|
import (
|
2020-09-20 00:26:02 +02:00
|
|
|
"time"
|
|
|
|
|
2020-08-13 03:54:14 +02:00
|
|
|
"github.com/spf13/cobra"
|
|
|
|
|
|
|
|
"github.com/RichardKnop/machinery/v1"
|
|
|
|
queueLog "github.com/RichardKnop/machinery/v1/log"
|
|
|
|
"github.com/jmoiron/sqlx"
|
2021-10-31 00:20:41 +02:00
|
|
|
"github.com/jordanknott/taskcafe/internal/config"
|
2020-08-13 03:54:14 +02:00
|
|
|
repo "github.com/jordanknott/taskcafe/internal/db"
|
2021-10-31 00:20:41 +02:00
|
|
|
"github.com/jordanknott/taskcafe/internal/jobs"
|
2020-08-13 03:54:14 +02:00
|
|
|
log "github.com/sirupsen/logrus"
|
|
|
|
)
|
|
|
|
|
|
|
|
func newWorkerCmd() *cobra.Command {
|
|
|
|
cc := &cobra.Command{
|
|
|
|
Use: "worker",
|
|
|
|
Short: "Run the task queue worker",
|
|
|
|
Long: "Run the task queue worker",
|
|
|
|
RunE: func(cmd *cobra.Command, args []string) error {
|
|
|
|
Formatter := new(log.TextFormatter)
|
|
|
|
Formatter.TimestampFormat = "02-01-2006 15:04:05"
|
|
|
|
Formatter.FullTimestamp = true
|
|
|
|
log.SetFormatter(Formatter)
|
|
|
|
log.SetLevel(log.InfoLevel)
|
|
|
|
|
2021-10-31 00:20:41 +02:00
|
|
|
appConfig, err := config.GetAppConfig()
|
|
|
|
if err != nil {
|
|
|
|
log.Panic(err)
|
|
|
|
}
|
|
|
|
db, err := sqlx.Connect("postgres", config.GetDatabaseConfig().GetDatabaseConnectionUri())
|
2020-08-13 03:54:14 +02:00
|
|
|
if err != nil {
|
|
|
|
log.Panic(err)
|
|
|
|
}
|
|
|
|
db.SetMaxOpenConns(25)
|
|
|
|
db.SetMaxIdleConns(25)
|
|
|
|
db.SetConnMaxLifetime(5 * time.Minute)
|
|
|
|
defer db.Close()
|
|
|
|
|
|
|
|
log.Info("starting task queue server instance")
|
2021-10-31 00:20:41 +02:00
|
|
|
jobConfig := appConfig.Job.GetJobConfig()
|
|
|
|
server, err := machinery.NewServer(&jobConfig)
|
2020-08-13 03:54:14 +02:00
|
|
|
if err != nil {
|
|
|
|
// do something with the error
|
|
|
|
}
|
2021-10-31 00:20:41 +02:00
|
|
|
queueLog.Set(&jobs.MachineryLogger{})
|
2020-08-13 03:54:14 +02:00
|
|
|
repo := *repo.NewRepository(db)
|
2021-11-18 00:11:28 +01:00
|
|
|
redisClient, err := appConfig.MessageQueue.GetMessageQueueClient()
|
|
|
|
if err != nil {
|
|
|
|
return err
|
|
|
|
}
|
|
|
|
jobs.RegisterTasks(server, repo, appConfig, redisClient)
|
2020-08-13 03:54:14 +02:00
|
|
|
|
|
|
|
worker := server.NewWorker("taskcafe_worker", 10)
|
|
|
|
log.Info("starting task queue worker")
|
|
|
|
err = worker.Launch()
|
|
|
|
if err != nil {
|
|
|
|
log.WithError(err).Error("error while launching ")
|
|
|
|
return err
|
|
|
|
}
|
|
|
|
return nil
|
|
|
|
},
|
|
|
|
}
|
|
|
|
return cc
|
|
|
|
}
|