/
AngelCareMe
/
task_queue
Обзор
Документация
Войти
/
AngelCareMe
/
task_queue
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
Аналитика
Безопасность
main
cmd/worker/main.go
127 строк
3 KB
AngelCareMe
done
28 июл 2025, 21:07
28 июл 2025, 21:07
c70014b
Код
Авторство
О чём код?
package main import ( "context" "fmt" "os" "os/signal" "syscall" "time" "task-queue/internal/adapter/queue" "task-queue/internal/adapter/queue/redis" redisConfig "task-queue/internal/adapter/queue/redis/config" redisConn "task-queue/internal/adapter/queue/redis/connection" "task-queue/internal/adapter/repository" "task-queue/internal/adapter/repository/postgres" postgresConfig "task-queue/internal/adapter/repository/postgres/config" postgresConn "task-queue/internal/adapter/repository/postgres/connection" "task-queue/internal/service" "task-queue/internal/usecase" "task-queue/pkg/config" "task-queue/pkg/logger" "task-queue/pkg/validator" ) func main() { // initialize config cfg, err := config.Load() if err != nil { fmt.Fprintf(os.Stderr, "Failed to load config: %v\n", err) os.Exit(1) } // initialize logger logger.Init() log := logger.Get() validator.SetLogger(log) log.Info("Starting Task Queue Worker") // initialize database connection dbConfig := &postgresConfig.Config{ Host: cfg.Database.Host, Port: cfg.Database.Port, User: cfg.Database.User, Password: cfg.Database.Password, Database: cfg.Database.Name, SSLMode: cfg.Database.SSLMode, MaxConnections: cfg.Database.MaxConnections, MinConnections: postgresConfig.DefaultConfig().MinConnections, MaxConnLifetime: postgresConfig.DefaultConfig().MaxConnLifetime, MaxConnIdleTime: postgresConfig.DefaultConfig().MaxConnIdleTime, ConnectTimeout: postgresConfig.DefaultConfig().ConnectTimeout, } dbConnection, err := postgresConn.NewConnection(cfg.Database.Host, cfg.Database.Port, cfg.Database.User, cfg.Database.Password, cfg.Database.Name, cfg.Database.SSLMode, cfg.Database.MaxConnections, dbConfig) if err != nil { log.WithError(err).Fatal("Failed to connect to database") } defer dbConnection.Close() // initialize redis connection redisCfg := &redisConfig.Config{ Host: cfg.Redis.Host, Port: cfg.Redis.Port, Password: cfg.Redis.Password, DB: cfg.Redis.DB, } redisConnection, err := redisConn.NewConnection(redisCfg) if err != nil { log.WithError(err).Fatal("Failed to connect to Redis") } defer redisConnection.Close() // initialize repositories and queues taskRepo := postgres.NewPostgresRepository(dbConnection.GetPool()) taskQueue := redis.NewRedisQueue(redisConnection.GetClient(), taskRepo) // initialize usecase taskUseCase := usecase.NewTaskUseCase( taskRepo.(interface { repository.TaskRepository repository.TaskStatistics repository.RepositoryHealthChecker }), taskQueue.(interface { queue.TaskQueue queue.QueueMetrics queue.QueueHealthChecker }), ) // initialize worker service workerService := service.NewWorkerService( taskQueue.(interface { queue.TaskQueue queue.QueueMetrics }), taskUseCase, ) // start workers workerCount := 5 // можно сделать конфигурируемым ctx, cancel := context.WithCancel(context.Background()) if err := workerService.StartWorkers(ctx, workerCount); err != nil { log.WithError(err).Fatal("Failed to start workers") } log.Infof("Started %d workers", workerCount) // wait for interrupt signal to gracefully shutdown quit := make(chan os.Signal, 1) signal.Notify(quit, syscall.SIGINT, syscall.SIGTERM) <-quit log.Info("Shutting down workers...") // stop workers cancel() // give workers time to finish time.Sleep(5 * time.Second) log.Info("Workers stopped") }