package main import ( "context" "flag" "log/slog" "os" "os/signal" "syscall" "time" "gl/application/health" applicationledger "gl/application/ledger" applicationsettlement "gl/application/settlement" "gl/infrastructure/config" "gl/infrastructure/kuknos" "gl/infrastructure/observability" "gl/infrastructure/postgres" grpcadapter "gl/interface/grpc" ) func main() { observability.Configure("gl", slog.LevelInfo) configPath := flag.String("conf", "./gl.cfg.toml", "path to the TOML configuration file") flag.Parse() cfg, err := config.Load(*configPath) if err != nil { slog.Error("load configuration", "error", err) os.Exit(1) } shutdownTelemetry := observability.Setup(context.Background(), cfg.Telemetry, "gl", cfg.Environment) defer shutdownTelemetry(context.Background()) ctx, stop := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM) defer stop() database, err := postgres.Open(ctx, cfg.Database) if err != nil { slog.Error("open GL database", "error", err) os.Exit(1) } defer database.Close() if err := postgres.Migrate(ctx, database); err != nil { slog.Error("migrate GL database", "error", err) os.Exit(1) } repository := postgres.NewJournalRepository(database) settlementRepository := postgres.NewSettlementRepository(database) settlementService := applicationsettlement.NewService(settlementRepository) var healthService *health.Service if cfg.Settlement.Enabled { healthService = health.NewService(database, settlementRepository) } else { healthService = health.NewService(database) } handler := grpcadapter.NewHandler(healthService, applicationledger.NewService(repository, nil), settlementService).WithSettlementAdminToken(cfg.Settlement.AdminToken) if cfg.Settlement.Enabled { worker := applicationsettlement.Worker{Repository: settlementRepository, Submitter: kuknos.Submitter{Endpoint: cfg.Settlement.Endpoint}, WorkerID: cfg.Settlement.WorkerID, MaxAttempts: cfg.Settlement.MaxAttempts} go runSettlementWorker(ctx, worker, cfg.Settlement.PollInterval, cfg.Settlement.BatchSize) go runSettlementReconciliation(ctx, settlementRepository) slog.Info("Kuknos settlement worker enabled", "endpoint", cfg.Settlement.Endpoint, "worker_id", cfg.Settlement.WorkerID) } slog.Info("starting GL gRPC service", "host", cfg.GRPC.Host, "port", cfg.GRPC.Port) serverConfig := grpcadapter.ServerConfig{ Host: cfg.GRPC.Host, Port: cfg.GRPC.Port, ShutdownTimeout: cfg.GRPC.ShutdownTimeout, } if err := grpcadapter.Run(ctx, serverConfig, handler); err != nil { slog.Error("GL service stopped", "error", err) os.Exit(1) } } func runSettlementReconciliation(ctx context.Context, repository *postgres.SettlementRepository) { run := func() { stats, err := repository.ReconcileConfirmationEvidence(ctx, 10*time.Minute) if err != nil { slog.ErrorContext(ctx, "settlement reconciliation failed", "error", err) return } if stats.Detected > 0 || stats.Resolved > 0 || stats.Open > 0 { slog.WarnContext(ctx, "settlement reconciliation completed", "detected", stats.Detected, "resolved", stats.Resolved, "open", stats.Open) } } run() ticker := time.NewTicker(time.Minute) defer ticker.Stop() for { select { case <-ctx.Done(): return case <-ticker.C: run() } } } func runSettlementWorker(ctx context.Context, worker applicationsettlement.Worker, interval time.Duration, batchSize int) { if interval <= 0 { interval = 5 * time.Second } if batchSize <= 0 { batchSize = 20 } ticker := time.NewTicker(interval) defer ticker.Stop() for { select { case <-ctx.Done(): return case <-ticker.C: if err := worker.RunOnce(ctx, batchSize); err != nil { slog.ErrorContext(ctx, "settlement worker cycle failed", "error", err) } } } }