package main import ( "context" "fmt" "log/slog" "os" "os/signal" "runtime/debug" "syscall" "github.com/spf13/cobra" "golang.org/x/term" "git.warky.dev/wdevs/pgsql-broker/pkg/broker" "git.warky.dev/wdevs/pgsql-broker/pkg/broker/adapter" "git.warky.dev/wdevs/pgsql-broker/pkg/broker/config" "git.warky.dev/wdevs/pgsql-broker/pkg/broker/install" ) var ( // Version information (injected by build) Version = "dev" BuildTime = "unknown" Commit = "unknown" // Command line flags cfgFile string logLevel string verifyOnly bool withRoles bool adminUser string adminPassword string brokerAdminPassword string brokerRuntimePassword string brokerEnqueuePassword string ) func main() { if err := rootCmd.Execute(); err != nil { fmt.Fprintf(os.Stderr, "Error: %v\n", err) os.Exit(1) } } var rootCmd = &cobra.Command{ Use: "pgsql-broker", Short: "PostgreSQL job broker for background job processing", Long: `PostgreSQL Broker is a job processing system that uses PostgreSQL LISTEN/NOTIFY for event-driven job execution. It supports multiple queues, priority-based scheduling, and horizontal scaling.`, Version: fmt.Sprintf("%s (built %s, commit %s)", Version, BuildTime, Commit), } var startCmd = &cobra.Command{ Use: "start", Short: "Start the broker instance", Long: `Start the broker instance and begin processing jobs from the database queue.`, RunE: func(cmd *cobra.Command, args []string) error { return runBroker() }, } var versionCmd = &cobra.Command{ Use: "version", Short: "Print version information", Run: func(cmd *cobra.Command, args []string) { fmt.Printf("pgsql-broker version %s\n", Version) fmt.Printf("Built: %s\n", BuildTime) fmt.Printf("Commit: %s\n", Commit) }, } var installCmd = &cobra.Command{ Use: "install", Short: "Install database schema (tables and procedures)", Long: `Install the required database schema including tables and stored procedures. This command will create all necessary tables and functions in the configured database.`, RunE: func(cmd *cobra.Command, args []string) error { return runInstall() }, } func init() { rootCmd.AddCommand(startCmd) rootCmd.AddCommand(versionCmd) rootCmd.AddCommand(installCmd) // Persistent flags rootCmd.PersistentFlags().StringVar(&cfgFile, "config", "", "config file (default is broker.yaml)") rootCmd.PersistentFlags().StringVar(&logLevel, "log-level", "info", "log level (debug, info, warn, error)") // Install command flags installCmd.Flags().BoolVar(&verifyOnly, "verify-only", false, "only verify installation without installing") installCmd.Flags().BoolVar(&withRoles, "with-roles", false, "also create/update the broker_admin, broker_runtime, and broker_enqueue roles") installCmd.Flags().StringVar(&adminUser, "admin-user", "", "superuser/CREATEROLE login used only for --with-roles (falls back to PGUSER/PG_USER env)") installCmd.Flags().StringVar(&adminPassword, "admin-password", "", "password for --admin-user (falls back to PGPASSWORD/PG_PASS env, then an interactive prompt)") installCmd.Flags().StringVar(&brokerAdminPassword, "broker-admin-password", "", "password to set for broker_admin (falls back to BROKER_ADMIN_PASSWORD env, then an interactive prompt)") installCmd.Flags().StringVar(&brokerRuntimePassword, "broker-runtime-password", "", "password to set for broker_runtime (falls back to BROKER_RUNTIME_PASSWORD env, then an interactive prompt)") installCmd.Flags().StringVar(&brokerEnqueuePassword, "broker-enqueue-password", "", "password to set for broker_enqueue (falls back to BROKER_ENQUEUE_PASSWORD env, then an interactive prompt)") } // resolveCredential returns the first non-empty value among the flag value, // the given environment variables (checked in order), and -- if none are // set -- an interactive masked prompt. It errors rather than prompting when // stdin is not a terminal, since a hang in a non-interactive context (CI, // systemd) is worse than a clear failure. func resolveCredential(flagVal string, envNames []string, promptLabel string) (string, error) { if flagVal != "" { return flagVal, nil } for _, name := range envNames { if v := os.Getenv(name); v != "" { return v, nil } } fd := int(os.Stdin.Fd()) if !term.IsTerminal(fd) { return "", fmt.Errorf( "%s not provided and stdin is not a terminal to prompt on; pass it via flag or one of %v", promptLabel, envNames, ) } fmt.Fprintf(os.Stderr, "%s: ", promptLabel) b, err := term.ReadPassword(fd) fmt.Fprintln(os.Stderr) if err != nil { return "", fmt.Errorf("failed to read %s: %w", promptLabel, err) } if len(b) == 0 { return "", fmt.Errorf("%s must not be empty", promptLabel) } return string(b), nil } func runBroker() (err error) { // Top-level safety net: an unrecovered panic anywhere in startup or the // shutdown wait must not crash the process with a raw trace -- log it // and return a normal error instead. defer func() { if r := recover(); r != nil { err = fmt.Errorf("recovered from panic in runBroker: %v\n%s", r, debug.Stack()) } }() // Load configuration cfg, err := config.LoadConfig(cfgFile) if err != nil { return fmt.Errorf("failed to load config: %w", err) } // Override log level if specified if logLevel != "" { cfg.Logging.Level = logLevel } // Setup logger logger := createLogger(cfg.Logging) logger.Info("starting pgsql-broker", "version", Version, "databases", len(cfg.Databases)) // Create broker (manages all database instances) b, err := broker.New(cfg, logger, Version) if err != nil { return fmt.Errorf("failed to create broker: %w", err) } // Start the broker (starts all database instances) if err := b.Start(); err != nil { return fmt.Errorf("failed to start broker: %w", err) } // Wait for shutdown signal waitForShutdown(b, logger) return nil } func runInstall() error { // Load configuration cfg, err := config.LoadConfig(cfgFile) if err != nil { return fmt.Errorf("failed to load config: %w", err) } // Override log level if specified if logLevel != "" { cfg.Logging.Level = logLevel } // Setup logger logger := createLogger(cfg.Logging) logger.Info("pgsql-broker database installer", "version", Version, "databases", len(cfg.Databases)) ctx := context.Background() var rolePasswords install.RolePasswords var adminUserVal, adminPasswordVal string if withRoles { if verifyOnly { return fmt.Errorf("--with-roles cannot be combined with --verify-only") } var err error adminUserVal, err = resolveCredential(adminUser, []string{"PGUSER", "PG_USER"}, "admin user (superuser/CREATEROLE login for --with-roles)") if err != nil { return err } adminPasswordVal, err = resolveCredential(adminPassword, []string{"PGPASSWORD", "PG_PASS"}, "admin password") if err != nil { return err } rolePasswords.AdminPassword, err = resolveCredential(brokerAdminPassword, []string{"BROKER_ADMIN_PASSWORD"}, "broker_admin password") if err != nil { return err } rolePasswords.RuntimePassword, err = resolveCredential(brokerRuntimePassword, []string{"BROKER_RUNTIME_PASSWORD"}, "broker_runtime password") if err != nil { return err } rolePasswords.EnqueuePassword, err = resolveCredential(brokerEnqueuePassword, []string{"BROKER_ENQUEUE_PASSWORD"}, "broker_enqueue password") if err != nil { return err } } // Install/verify on all configured databases for i := range cfg.Databases { dbCfg := &cfg.Databases[i] logger.Info("processing database", "index", i, "name", dbCfg.Name, "host", dbCfg.Host, "database", dbCfg.Database) // Create database adapter. With --with-roles, the config file's own // user (typically the least-privilege broker_runtime) may not exist // yet on a fresh cluster -- migrations and role creation both run as // the admin login instead, since broker_admin must own the schema. pgCfg := dbCfg.ToPostgresConfig() if withRoles { pgCfg.User = adminUserVal pgCfg.Password = adminPasswordVal } dbAdapter := adapter.NewPostgresAdapter(pgCfg, logger) // Connect to database if err := dbAdapter.Connect(ctx); err != nil { return fmt.Errorf("failed to connect to database %s: %w", dbCfg.Name, err) } // Create installer installer := install.New(dbAdapter, logger) if verifyOnly { // Only verify installation logger.Info("verifying database schema", "database", dbCfg.Name) if err := installer.VerifyInstallation(ctx); err != nil { dbAdapter.Close() logger.Error("verification failed", "database", dbCfg.Name, "error", err) return fmt.Errorf("verification failed for %s: %w", dbCfg.Name, err) } logger.Info("database schema verified successfully", "database", dbCfg.Name) } else { // Apply migrations logger.Info("applying database migrations", "database", dbCfg.Name) if err := installer.ApplyMigrations(ctx); err != nil { dbAdapter.Close() logger.Error("installation failed", "database", dbCfg.Name, "error", err) return fmt.Errorf("installation failed for %s: %w", dbCfg.Name, err) } // Verify installation logger.Info("verifying installation", "database", dbCfg.Name) if err := installer.VerifyInstallation(ctx); err != nil { dbAdapter.Close() logger.Error("verification failed", "database", dbCfg.Name, "error", err) return fmt.Errorf("verification failed for %s: %w", dbCfg.Name, err) } logger.Info("database schema installed and verified successfully", "database", dbCfg.Name) } if withRoles && !verifyOnly { logger.Info("applying roles", "database", dbCfg.Name) if err := installer.InstallRoles(ctx, rolePasswords); err != nil { dbAdapter.Close() return fmt.Errorf("failed to install roles for %s: %w", dbCfg.Name, err) } logger.Info("roles installed successfully", "database", dbCfg.Name) } dbAdapter.Close() } if verifyOnly { logger.Info("all databases verified successfully") } else { logger.Info("all databases installed and verified successfully") } return nil } func createLogger(cfg config.LoggingConfig) adapter.Logger { // Parse log level var level slog.Level switch cfg.Level { case "debug": level = slog.LevelDebug case "info": level = slog.LevelInfo case "warn": level = slog.LevelWarn case "error": level = slog.LevelError default: level = slog.LevelInfo } // Create handler based on format var handler slog.Handler opts := &slog.HandlerOptions{Level: level} if cfg.Format == "text" { handler = slog.NewTextHandler(os.Stdout, opts) } else { handler = slog.NewJSONHandler(os.Stdout, opts) } return adapter.NewSlogLoggerWithHandler(handler) } func waitForShutdown(b *broker.Broker, logger adapter.Logger) { sigChan := make(chan os.Signal, 1) signal.Notify(sigChan, os.Interrupt, syscall.SIGTERM, syscall.SIGINT) sig := <-sigChan logger.Info("received shutdown signal", "signal", sig) if err := b.Stop(); err != nil { logger.Error("error during shutdown", "error", err) } }