// app/cmd/goaws.go // Entry point — shared-sqs server // Updated: 2026-04-10 — добавлена Redis persistence (write-through cache) package main import ( "context" "encoding/json" "flag" "net/http" "os" "os/signal" "syscall" "time" "shared-sqs/app/billing" "shared-sqs/app/conf" "shared-sqs/app/gosqs" "shared-sqs/app/metrics" "shared-sqs/app/models" "shared-sqs/app/persistence" "shared-sqs/app/router" "shared-sqs/app/tenant" log "github.com/sirupsen/logrus" ) func main() { var configFile string var adminToken string var port string var debug bool var loglevel string flag.StringVar(&configFile, "config", "", "config file location") flag.StringVar(&adminToken, "admin-token", "", "admin API bearer token") flag.StringVar(&port, "port", "4100", "listen port") flag.BoolVar(&debug, "debug", false, "set debug log level") flag.StringVar(&loglevel, "loglevel", "info", "log level (info, debug, warn, error)") flag.Parse() log.SetFormatter(&log.JSONFormatter{}) log.SetOutput(os.Stdout) if debug { log.SetLevel(log.DebugLevel) } else { level, err := log.ParseLevel(loglevel) if err != nil { log.SetLevel(log.InfoLevel) log.Warnf("Failed to parse loglevel %v, defaulting to info", loglevel) } else { log.SetLevel(level) } } // Admin token: flag > env SHARED_SQS_ADMIN_TOKEN > fatal (Trap #13) if adminToken == "" { adminToken = os.Getenv("SHARED_SQS_ADMIN_TOKEN") } if adminToken == "" { log.Fatal("admin token required: use --admin-token flag or SHARED_SQS_ADMIN_TOKEN env var") } // Загрузить конфиг (очереди, env — без SNS) env := "Local" if flag.NArg() > 0 { env = flag.Arg(0) } conf.LoadYamlConfig(configFile, env) if models.CurrentEnvironment.LogToFile { filename := models.CurrentEnvironment.LogFile file, err := os.OpenFile(filename, os.O_CREATE|os.O_WRONLY|os.O_APPEND, 0666) if err == nil { log.SetOutput(file) } else { log.Infof("Failed to log to file: %s, using default stdout", filename) } } // Инициализация in-memory TenantStore tenantStore := tenant.NewTenantStore() // Подключение к Redis (если задан REDIS_ADDR) // При ошибке — предупреждение, но продолжаем в memory-only режиме redisAddr := os.Getenv("REDIS_ADDR") redisUser := os.Getenv("REDIS_USER") redisPass := os.Getenv("REDIS_PASSWORD") if redisAddr != "" { if err := persistence.Connect(redisAddr, redisUser, redisPass); err != nil { log.Warnf("Не удалось подключиться к Redis: %v — работаем в memory-only режиме", err) } } // Восстановление состояния из Redis (тенанты + очереди) if persistence.Client != nil { // Загружаем тенантов tenantsRaw, err := persistence.LoadAllTenantsRaw() if err != nil { log.Warnf("Ошибка загрузки тенантов из Redis: %v", err) } else { for _, jsonBytes := range tenantsRaw { var t tenant.Tenant if err := json.Unmarshal(jsonBytes, &t); err != nil { log.Errorf("Ошибка десериализации тенанта: %v", err) continue } tenantStore.LoadTenant(&t) } } // Загружаем очереди queues, err := persistence.LoadAllQueues() if err != nil { log.Warnf("Ошибка загрузки очередей из Redis: %v", err) } else { models.SyncQueues.Lock() for k, q := range queues { models.SyncQueues.Queues[k] = q } models.SyncQueues.Unlock() } } // Автосид демо-данных при SHARED_SQS_SEED_DEMO=true if os.Getenv("SHARED_SQS_SEED_DEMO") == "true" { seedDemoData(tenantStore) } // Billing: подключение к PostgreSQL для учёта использования. // Если BILLING_PG_HOST не задан — billing отключён, SQS работает без него. billing.Init() // Роутер с tenant auth и admin API r := router.New(tenantStore, adminToken) // PeriodicTasks — visibility timeout, DLQ, deduplication quit := make(chan bool) go gosqs.PeriodicTasks(1*time.Second, quit) // Metrics: gauge updater — пересчёт очередей/сообщений per tenant каждые 15 секунд metrics.StartGaugeUpdater(15*time.Second, quit) // HTTP сервер с таймаутами srv := &http.Server{ Addr: "0.0.0.0:" + port, Handler: r, ReadTimeout: 30 * time.Second, WriteTimeout: 35 * time.Second, // чуть больше чем max WaitTimeSeconds (20s) IdleTimeout: 60 * time.Second, } // Запуск в горутине для graceful shutdown serverErr := make(chan error, 1) go func() { log.Infof("shared-sqs listening on 0.0.0.0:%s", port) if err := srv.ListenAndServe(); err != nil && err != http.ErrServerClosed { serverErr <- err } }() // Graceful shutdown по SIGTERM/SIGINT (Trap #13) sigCh := make(chan os.Signal, 1) signal.Notify(sigCh, syscall.SIGTERM, syscall.SIGINT) select { case sig := <-sigCh: log.Infof("Received signal %s, shutting down", sig) case err := <-serverErr: log.Fatalf("Server error: %v", err) } // Остановить PeriodicTasks close(quit) // Закрыть billing (если был подключён) billing.Close() // Дать 10 секунд на завершение текущих HTTP запросов ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second) defer cancel() if err := srv.Shutdown(ctx); err != nil { log.Errorf("Server shutdown error: %v", err) } log.Info("shared-sqs stopped") }