84 lines
2.6 KiB
Go
84 lines
2.6 KiB
Go
|
|
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)
|
||
|
|
}
|