BS / src / conf / queue.go
queue.go
Raw
package conf

import (
	"context"
	"fmt"
	"net/http"
	"time"

	"github.com/hibiken/asynq"
	"github.com/hibiken/asynqmon"
)

var (
	RunAsServer       bool
	RunAsClient       bool
	RunAsScheduler    bool
	QueueScheduler    *asynq.Scheduler
	QueueClient       *asynq.Client
	QueueServer       *asynq.Server
	asyncQRedisClient asynq.RedisClientOpt
	mux               *asynq.ServeMux
)

type QueueSchedules struct {
	Cron    string
	Key     string
	Payload []byte
	Q       asynq.Option
	Timeout time.Duration
}

type MuxHandler struct {
	Key     string
	Handler func(context.Context, *asynq.Task) error
	Q       asynq.Option
}

const (
	UsersQ       = "users"
	ScanQ        = "scan"
	FetchQ       = "fetch"
	ParseQ       = "Parse"
	ProcessQ     = "Process"
	MainQ        = "main"
	HouseKeeping = "HouseKeeping"
	DefaultQ     = "default"
	LowPriorityQ = "Unimportant"
)

func LoadQueue() {
	// Create and configuring Redis connection.
	asyncQRedisClient = asynq.RedisClientOpt{
		Addr: fmt.Sprintf("%s:%s", Config.RedisUrl.Hostname(), Config.RedisUrl.Port()),
		DB:   Config.RedisDB,
	}
	QueueClient = asynq.NewClient(asyncQRedisClient)

	// Run worker server.
	QueueServer = asynq.NewServer(asyncQRedisClient, asynq.Config{
		Concurrency:  int(Config.MaxConcurrency),
		ErrorHandler: &QueueErrorHandler{},
		Queues: map[string]int{
			ProcessQ:     8,
			FetchQ:       6,
			ParseQ:       6,
			ScanQ:        3,
			MainQ:        4,
			DefaultQ:     3,
			HouseKeeping: 2,
			LowPriorityQ: 1,
		},
	})
	mux = asynq.NewServeMux()
	// Block Related

	loc, err := time.LoadLocation("UTC")
	if err != nil {
		Logger.Panic(err)
	}
	QueueScheduler = asynq.NewScheduler(
		asyncQRedisClient,
		&asynq.SchedulerOpts{
			Location: loc,
		},
	)
}

func RunClient() {
	RunAsClient = true
}

func RunWorker(muxHandler []MuxHandler) {
	RunAsServer = true
	for _, mh := range muxHandler {
		mux.HandleFunc(mh.Key, mh.Handler)
	}
	if err := QueueServer.Run(mux); err != nil {
		Logger.Panic(err)
	}
}

func RunScheduler(queueSchedules []QueueSchedules) {
	RunAsScheduler = true
	for _, qs := range queueSchedules {
		_, err := QueueScheduler.Register(qs.Cron, asynq.NewTask(qs.Key, qs.Payload), qs.Q, asynq.Timeout(qs.Timeout))
		if err != nil {
			Logger.Panicf("QueueScheduler: %s", err)
		}
	}
	if err2 := QueueScheduler.Start(); err2 != nil {
		Logger.Panic(err2)
	}
}

func RunMonitor(URL string) {
	h := asynqmon.New(asynqmon.Options{
		RootPath:     "/mon",
		RedisConnOpt: asyncQRedisClient,
	})
	http.Handle(h.RootPath()+"/", h)
	Logger.Panic(http.ListenAndServe(URL, nil))
}