feat: add asynq worker skeleton
This commit is contained in:
@@ -2,12 +2,21 @@ package main
|
||||
|
||||
import (
|
||||
"context"
|
||||
"database/sql"
|
||||
"errors"
|
||||
"fmt"
|
||||
"log/slog"
|
||||
"os"
|
||||
"os/signal"
|
||||
"syscall"
|
||||
"time"
|
||||
|
||||
"forgejo.digital-droplets.de/philschlo/proxui/platform/config"
|
||||
"forgejo.digital-droplets.de/philschlo/proxui/platform/logging"
|
||||
"github.com/hibiken/asynq"
|
||||
_ "github.com/jackc/pgx/v5/stdlib"
|
||||
|
||||
"proxui/worker/internal/tasks"
|
||||
)
|
||||
|
||||
func main() {
|
||||
@@ -19,10 +28,54 @@ func main() {
|
||||
}
|
||||
|
||||
logger := logging.New("worker", cfg.AppEnv, cfg.LogLevel)
|
||||
db, err := openDatabase(cfg.DatabaseURL)
|
||||
if err != nil {
|
||||
logger.Error("failed to connect database", "error", err)
|
||||
os.Exit(1)
|
||||
}
|
||||
defer db.Close()
|
||||
|
||||
ctx, stop := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM)
|
||||
defer stop()
|
||||
|
||||
logger.Info("worker started", "concurrency", cfg.WorkerConcurrency)
|
||||
redis := asynq.RedisClientOpt{Addr: cfg.RedisAddr}
|
||||
server := asynq.NewServer(redis, asynq.Config{
|
||||
Concurrency: cfg.WorkerConcurrency,
|
||||
})
|
||||
mux := asynq.NewServeMux()
|
||||
registerHandlers(mux, logger)
|
||||
|
||||
go func() {
|
||||
logger.Info("worker started", "concurrency", cfg.WorkerConcurrency, "redis_addr", cfg.RedisAddr)
|
||||
if err := server.Run(mux); err != nil && !errors.Is(err, asynq.ErrServerClosed) {
|
||||
logger.Error("worker failed", "error", err)
|
||||
stop()
|
||||
}
|
||||
}()
|
||||
|
||||
<-ctx.Done()
|
||||
server.Shutdown()
|
||||
logger.Info("worker stopped")
|
||||
}
|
||||
|
||||
func registerHandlers(mux *asynq.ServeMux, logger *slog.Logger) {
|
||||
mux.Handle(tasks.TypeDummy, tasks.NewDummyHandler(logger))
|
||||
}
|
||||
|
||||
func openDatabase(databaseURL string) (*sql.DB, error) {
|
||||
if databaseURL == "" {
|
||||
return nil, fmt.Errorf("DATABASE_URL is required")
|
||||
}
|
||||
|
||||
db, err := sql.Open("pgx", databaseURL)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
db.SetConnMaxLifetime(30 * time.Minute)
|
||||
if err := db.Ping(); err != nil {
|
||||
_ = db.Close()
|
||||
return nil, err
|
||||
}
|
||||
|
||||
return db, nil
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user