|
1 | 1 | package main |
2 | 2 |
|
3 | 3 | import ( |
4 | | - "cmp" |
5 | | - "context" |
6 | | - "errors" |
7 | | - "flag" |
8 | | - "fmt" |
9 | | - "log/slog" |
10 | | - "net/http" |
11 | | - "os" |
12 | | - "strings" |
13 | | - "time" |
14 | | - |
| 4 | + "github.com/jackc/pgx/v5" |
15 | 5 | "github.com/jackc/pgx/v5/pgxpool" |
16 | | - "github.com/rs/cors" |
17 | | - sloghttp "github.com/samber/slog-http" |
18 | 6 | "riverqueue.com/riverui" |
19 | | - "riverqueue.com/riverui/authmiddleware" |
| 7 | + "riverqueue.com/riverui/internal/apibundle" |
| 8 | + "riverqueue.com/riverui/internal/riveruicmd" |
20 | 9 |
|
21 | | - "github.com/riverqueue/apiframe/apimiddleware" |
22 | 10 | "github.com/riverqueue/river" |
23 | 11 | "github.com/riverqueue/river/riverdriver/riverpgxv5" |
24 | 12 | ) |
25 | 13 |
|
26 | 14 | func main() { |
27 | | - ctx := context.Background() |
28 | | - |
29 | | - logger := slog.New(getLogHandler(&slog.HandlerOptions{ |
30 | | - Level: getLogLevel(), |
31 | | - })) |
32 | | - |
33 | | - var pathPrefix string |
34 | | - flag.StringVar(&pathPrefix, "prefix", "/", "path prefix to use for the API and UI HTTP requests") |
35 | | - flag.Parse() |
36 | | - |
37 | | - initRes, err := initServer(ctx, logger, pathPrefix) |
38 | | - if err != nil { |
39 | | - logger.ErrorContext(ctx, "Error initializing server", slog.String("error", err.Error())) |
40 | | - os.Exit(1) |
41 | | - } |
42 | | - |
43 | | - if err := startAndListen(ctx, logger, initRes); err != nil { |
44 | | - logger.ErrorContext(ctx, "Error starting server", slog.String("error", err.Error())) |
45 | | - os.Exit(1) |
46 | | - } |
47 | | -} |
48 | | - |
49 | | -// Translates either a "1" or "true" from env to a Go boolean. |
50 | | -func envBooleanTrue(val string) bool { |
51 | | - return val == "1" || val == "true" |
52 | | -} |
53 | | - |
54 | | -func getLogHandler(opts *slog.HandlerOptions) slog.Handler { |
55 | | - logFormat := strings.ToLower(os.Getenv("RIVER_LOG_FORMAT")) |
56 | | - switch logFormat { |
57 | | - case "json": |
58 | | - return slog.NewJSONHandler(os.Stdout, opts) |
59 | | - default: |
60 | | - return slog.NewTextHandler(os.Stdout, opts) |
61 | | - } |
62 | | -} |
63 | | - |
64 | | -func getLogLevel() slog.Level { |
65 | | - if envBooleanTrue(os.Getenv("RIVER_DEBUG")) { |
66 | | - return slog.LevelDebug |
67 | | - } |
68 | | - |
69 | | - switch strings.ToLower(os.Getenv("RIVER_LOG_LEVEL")) { |
70 | | - case "debug": |
71 | | - return slog.LevelDebug |
72 | | - case "warn": |
73 | | - return slog.LevelWarn |
74 | | - case "error": |
75 | | - return slog.LevelError |
76 | | - default: |
77 | | - return slog.LevelInfo |
78 | | - } |
79 | | -} |
80 | | - |
81 | | -type initServerResult struct { |
82 | | - dbPool *pgxpool.Pool // database pool; close must be deferred by caller! |
83 | | - httpServer *http.Server // HTTP server wrapping the UI server |
84 | | - logger *slog.Logger // application logger (also internalized in UI server) |
85 | | - uiServer *riverui.Server // River UI server |
86 | | -} |
87 | | - |
88 | | -func initServer(ctx context.Context, logger *slog.Logger, pathPrefix string) (*initServerResult, error) { |
89 | | - if !strings.HasPrefix(pathPrefix, "/") || pathPrefix == "" { |
90 | | - return nil, fmt.Errorf("invalid path prefix: %s", pathPrefix) |
91 | | - } |
92 | | - pathPrefix = riverui.NormalizePathPrefix(pathPrefix) |
93 | | - |
94 | | - var ( |
95 | | - basicAuthUsername = os.Getenv("RIVER_BASIC_AUTH_USER") |
96 | | - basicAuthPassword = os.Getenv("RIVER_BASIC_AUTH_PASS") |
97 | | - corsOrigins = strings.Split(os.Getenv("CORS_ORIGINS"), ",") |
98 | | - databaseURL = os.Getenv("DATABASE_URL") |
99 | | - devMode = envBooleanTrue(os.Getenv("DEV")) |
100 | | - jobListHideArgsByDefault = envBooleanTrue(os.Getenv("RIVER_JOB_LIST_HIDE_ARGS_BY_DEFAULT")) |
101 | | - host = os.Getenv("RIVER_HOST") // may be left empty to bind to all local interfaces |
102 | | - liveFS = envBooleanTrue(os.Getenv("LIVE_FS")) |
103 | | - otelEnabled = envBooleanTrue(os.Getenv("OTEL_ENABLED")) |
104 | | - port = cmp.Or(os.Getenv("PORT"), "8080") |
105 | | - ) |
106 | | - |
107 | | - if databaseURL == "" && os.Getenv("PGDATABASE") == "" { |
108 | | - return nil, errors.New("expect to have DATABASE_URL or database configuration in standard PG* env vars like PGDATABASE/PGHOST/PGPORT/PGUSER/PGPASSWORD") |
109 | | - } |
110 | | - |
111 | | - poolConfig, err := pgxpool.ParseConfig(databaseURL) |
112 | | - if err != nil { |
113 | | - return nil, fmt.Errorf("error parsing db config: %w", err) |
114 | | - } |
115 | | - |
116 | | - dbPool, err := pgxpool.NewWithConfig(ctx, poolConfig) |
117 | | - if err != nil { |
118 | | - return nil, fmt.Errorf("error connecting to db: %w", err) |
119 | | - } |
120 | | - |
121 | | - client, err := river.NewClient(riverpgxv5.New(dbPool), &river.Config{}) |
122 | | - if err != nil { |
123 | | - return nil, err |
124 | | - } |
125 | | - |
126 | | - uiServer, err := riverui.NewServer(&riverui.ServerOpts{ |
127 | | - Client: client, |
128 | | - DB: dbPool, |
129 | | - DevMode: devMode, |
130 | | - JobListHideArgsByDefault: jobListHideArgsByDefault, |
131 | | - LiveFS: liveFS, |
132 | | - Logger: logger, |
133 | | - Prefix: pathPrefix, |
134 | | - }) |
135 | | - if err != nil { |
136 | | - return nil, err |
137 | | - } |
138 | | - |
139 | | - corsHandler := cors.New(cors.Options{ |
140 | | - AllowedMethods: []string{"GET", "HEAD", "POST", "PUT"}, |
141 | | - AllowedOrigins: corsOrigins, |
142 | | - }) |
143 | | - logHandler := sloghttp.NewWithConfig(logger, sloghttp.Config{ |
144 | | - WithSpanID: otelEnabled, |
145 | | - WithTraceID: otelEnabled, |
146 | | - }) |
147 | | - |
148 | | - middlewareStack := apimiddleware.NewMiddlewareStack( |
149 | | - apimiddleware.MiddlewareFunc(sloghttp.Recovery), |
150 | | - apimiddleware.MiddlewareFunc(corsHandler.Handler), |
151 | | - apimiddleware.MiddlewareFunc(logHandler), |
152 | | - ) |
153 | | - if basicAuthUsername != "" && basicAuthPassword != "" { |
154 | | - middlewareStack.Use(&authmiddleware.BasicAuth{Username: basicAuthUsername, Password: basicAuthPassword}) |
155 | | - } |
156 | | - |
157 | | - return &initServerResult{ |
158 | | - dbPool: dbPool, |
159 | | - httpServer: &http.Server{ |
160 | | - Addr: host + ":" + port, |
161 | | - Handler: middlewareStack.Mount(uiServer), |
162 | | - ReadHeaderTimeout: 5 * time.Second, |
| 15 | + riveruicmd.Run( |
| 16 | + func(dbPool *pgxpool.Pool) (*river.Client[pgx.Tx], error) { |
| 17 | + return river.NewClient(riverpgxv5.New(dbPool), &river.Config{}) |
163 | 18 | }, |
164 | | - logger: logger, |
165 | | - uiServer: uiServer, |
166 | | - }, nil |
167 | | -} |
168 | | - |
169 | | -func startAndListen(ctx context.Context, logger *slog.Logger, initRes *initServerResult) error { |
170 | | - defer initRes.dbPool.Close() |
171 | | - |
172 | | - if err := initRes.uiServer.Start(ctx); err != nil { |
173 | | - return err |
174 | | - } |
175 | | - |
176 | | - logger.InfoContext(ctx, "Starting server", slog.String("addr", initRes.httpServer.Addr)) |
177 | | - |
178 | | - if err := initRes.httpServer.ListenAndServe(); err != nil && !errors.Is(err, http.ErrServerClosed) { |
179 | | - return err |
180 | | - } |
181 | | - |
182 | | - return nil |
| 19 | + func(client *river.Client[pgx.Tx], opts *riveruicmd.BundleOpts) apibundle.EndpointBundle { |
| 20 | + return riverui.NewEndpoints(&riverui.EndpointsOpts[pgx.Tx]{ |
| 21 | + Client: client, |
| 22 | + JobListHideArgsByDefault: opts.JobListHideArgsByDefault, |
| 23 | + }) |
| 24 | + }, |
| 25 | + ) |
183 | 26 | } |
0 commit comments