package main import ( "context" "net/http" "os" "os/signal" "syscall" "time" fwconfig "github.com/yourorg/go-dw-platform/framework/config" fwdb "github.com/yourorg/go-dw-platform/framework/db" fwingestion "github.com/yourorg/go-dw-platform/framework/ingestion" fwlogger "github.com/yourorg/go-dw-platform/framework/logger" "github.com/yourorg/go-dw-platform/services/omnix-broadcast/client" "github.com/yourorg/go-dw-platform/services/omnix-broadcast/handler" "github.com/yourorg/go-dw-platform/services/omnix-broadcast/repository" "github.com/yourorg/go-dw-platform/services/omnix-broadcast/service" "github.com/yourorg/go-dw-platform/services/omnix-broadcast/transformer" "github.com/yourorg/go-dw-platform/services/omnix-broadcast/worker" ) func main() { _ = fwconfig.LoadDotEnv(".env") cfg := LoadConfig() log := fwlogger.New("omnix-broadcast") ctx, cancel := context.WithCancel(context.Background()) defer cancel() pool, err := fwdb.NewPool(ctx, fwdb.Config{ DSN: cfg.DB.DSN, MaxConns: cfg.DB.MaxConns, MinConns: cfg.DB.MinConns, MaxConnLifetime: cfg.DB.MaxConnLifetime, MaxConnIdleTime: cfg.DB.MaxConnIdleTime, }) if err != nil { log.Error("failed to connect to database", "error", err) os.Exit(1) } defer pool.Close() repo := repository.New(pool) tf := transformer.New() sopigaClient := client.NewSopigaClient(cfg.Sopiga.BaseURL, cfg.Sopiga.Token) retrier := fwingestion.NewRetrier(cfg.Worker.MaxRetries, cfg.Worker.RetryDelay) svc := service.New(repo, tf, sopigaClient, retrier, log) bcWorker := worker.New(svc, log, cfg.Worker.CheckInterval, cfg.Worker.BatchSize) syncWorker := worker.NewDeliverySync(svc, log, cfg.Worker.SyncCheckInterval, cfg.Worker.SyncBatchSize) go bcWorker.Start(ctx) go syncWorker.Start(ctx) webhookHandler := handler.NewWebhookHandler(svc, log, cfg.Webhook.Secret) mux := http.NewServeMux() mux.HandleFunc("/webhooks/sopiga/delivery-status", webhookHandler.HandleDeliveryStatus) httpServer := &http.Server{ Addr: ":" + cfg.Webhook.Port, Handler: mux, } go func() { log.Info("webhook server listening", "port", cfg.Webhook.Port) if err := httpServer.ListenAndServe(); err != nil && err != http.ErrServerClosed { log.Error("webhook server stopped", "error", err) } }() sigChan := make(chan os.Signal, 1) signal.Notify(sigChan, syscall.SIGINT, syscall.SIGTERM) <-sigChan log.Info("shutting down gracefully") cancel() shutdownCtx, shutdownCancel := context.WithTimeout(context.Background(), 10*time.Second) defer shutdownCancel() _ = httpServer.Shutdown(shutdownCtx) }