Files
inventory-plus-plus/main.go
T
angel 6a473b2ea0 Claude-assisted improvements (untested)
- dev auth flow (side-step OAuth)
- db event processing integration tests
- dev scripts (eg Makefile)
- db / test db migration setup scripts.
2026-08-20 00:35:22 -06:00

307 lines
7.4 KiB
Go

package main
//go:generate planter postgres://planter:planter@localhost:5432/inventory_2?sslmode=disable -o diagrams/database_schema_public.uml
//go:generate planter postgres://planter:planter@localhost:5432/inventory_2?sslmode=disable -s mock -o diagrams/database_schema_mock.uml
//go:generate plantuml diagrams/*.uml -tsvg
//go:generate plantuml diagrams/etsy/*.uml -tsvg
//go:generate npx @tailwindcss/cli -i styles/source.css -o styles/index.css
import (
"context"
"errors"
"fmt"
"log/slog"
"net/http"
"os"
"os/signal"
"strings"
"syscall"
"time"
"github.com/jackc/pgx/v5/pgxpool"
"github.com/lmittmann/tint"
"ruben/inventory2/config"
"ruben/inventory2/domains/accounts"
"ruben/inventory2/domains/amazon"
"ruben/inventory2/domains/authentication"
etsy_platform "ruben/inventory2/domains/platforms/etsy"
"ruben/inventory2/domains/raw_events"
"ruben/inventory2/domains/reports"
"ruben/inventory2/logging"
"ruben/inventory2/server"
"ruben/inventory2/server/sse"
)
func main() {
logger := logging.New(tint.NewHandler(os.Stderr, &tint.Options{
AddSource: true,
Level: slog.LevelDebug,
ReplaceAttr: func(groups []string, a slog.Attr) slog.Attr {
// this can perform general key=value log cleanup
if a.Key == slog.SourceKey && len(groups) == 0 {
source := a.Value.Any().(*slog.Source)
source.File = strings.TrimPrefix(source.File, "/home/angel/go/src/ruben/inventory2/internal")
}
return a
},
// Time format (Default: time.StampMilli)
//TimeFormat: "",
}))
logger.Info("application starting")
if err := runApp(context.Background(), logger); err != nil {
panic(err)
}
logger.Info("application shutdown")
}
func runApp(ctx context.Context, logger *logging.Logger) error {
ctx, shutdown := context.WithCancel(ctx)
defer shutdown()
// load configuration
cfg, err := config.Load()
if err != nil {
return fmt.Errorf("failed to load configuration: %w", err)
}
// connect to the database
connPool, err := newPool(ctx, cfg.DatabaseURL)
if err != nil {
return fmt.Errorf("failed to initialize database connection pool: %w", err)
}
auth, err := authentication.New(
ctx,
connPool,
logger.WithGroup("authenticator"),
cfg.Auth0Domain,
cfg.Auth0ClientID,
cfg.Auth0ClientSecret,
cfg.Auth0CallbackURL,
)
if err != nil {
return fmt.Errorf("failed to construct authenticator: %w", err)
}
accts := accounts.NewStore(logger.WithGroup("accounts"), connPool)
// start http server
sseQueue, srvErrCh := runServer(ctx, logger, connPool, auth, accts, cfg)
// start background processes
authErrCh := runAuthProcesses(ctx, auth)
eventErrCh := runEventProcessing(
ctx,
amazon.NewMocks(logger.WithGroup("amazon"), connPool).
SetListener(sseQueue.NewDBEventPublisher(
func(ctx context.Context, e raw_events.Event) (acctID int64, err error) {
p, err := accounts.NewPlatform(e.Platform)
if err != nil {
return 0, err
}
return accts.GetAccountIDByMockPlatformAndShopID(ctx, p, e.StoreID)
},
func(acctID int64, e raw_events.Event) (eventType string, err error) {
// TODO: fine tune the event type later (don't want a referesh on EVERY event)
return fmt.Sprintf("accounts_%d_simulations", acctID), nil
},
)),
)
// wait for interrupt signal or unrecoverable failure, then shutdown
osSignalCh := make(chan os.Signal, 1)
signal.Notify(osSignalCh, syscall.SIGINT, syscall.SIGTERM)
var (
alreadyShutdown struct {
server bool
authProcesses bool
eventProcessing bool
}
)
select {
case s := <-osSignalCh:
logger.Info("application received shutdown signal", "signal", s)
logger.Info("shutting down")
case err := <-authErrCh:
alreadyShutdown.authProcesses = true
logger.Error("auth processes shutdown unexpectedly")
if err != nil {
logger.Error("auth processes encountered error", "error", err)
}
case err := <-eventErrCh:
alreadyShutdown.eventProcessing = true
logger.Error("event processing shutdown unexpectedly")
if err != nil {
logger.Error("event processing encountered error", "error", err)
}
case err := <-srvErrCh:
alreadyShutdown.server = true
logger.Error("server shutdown unexpectedly")
if err != nil {
logger.Error("server encountered error", "error", err)
}
}
shutdown()
// capture application errors that occurred during or caused shutdown
var errs []error
if !alreadyShutdown.authProcesses {
if err := <-authErrCh; err != nil {
errs = append(errs, fmt.Errorf("auth processes experienced an error: %w", err))
}
logger.Info("auth processes shut down")
}
if !alreadyShutdown.eventProcessing {
if err := <-authErrCh; err != nil {
errs = append(errs, fmt.Errorf("event processing experienced an error: %w", err))
}
logger.Info("event processing shut down")
}
if !alreadyShutdown.server {
if err := <-srvErrCh; err != nil {
errs = append(errs, fmt.Errorf("server experienced an error: %w", err))
}
logger.Info("server shut down")
}
return errors.Join(errs...)
}
func runAuthProcesses(ctx context.Context, auth *authentication.Authenticator) <-chan error {
errCh := make(chan error, 1)
go func() {
defer close(errCh)
if err := auth.RunBackgroundCleanup(ctx); err != nil {
errCh <- err
}
}()
return errCh
}
func runEventProcessing(ctx context.Context, amz *amazon.Mocks) <-chan error {
errCh := make(chan error, 1)
go func() {
defer close(errCh)
if err := amz.ProcessEvents(ctx); err != nil {
errCh <- err
}
}()
return errCh
}
func runServer(
ctx context.Context,
logger *logging.Logger,
connPool *pgxpool.Pool,
auth *authentication.Authenticator,
accts *accounts.Store,
cfg config.Config,
) (*sse.Queue, <-chan error) {
r := server.NewRouter(
logger.WithGroup("server"),
"./",
raw_events.NewStore(logger.WithGroup("raw-event-store"), connPool),
accts,
reports.NewStore(logger, connPool, accts),
etsy_platform.NewPlatform(
logger,
func(acctID int64) string {
return fmt.Sprintf("/oauth/account/%d/auth_code", acctID)
},
cfg.EtsyAPIKeystring,
cfg.EtsyAPISharedSecret,
connPool,
),
auth,
cfg.DevAuthEnabled,
)
srv := &http.Server{
Addr: ":8082", // local
Handler: r,
}
ctx, cancel := context.WithCancel(ctx)
alreadyShutdownCh := make(chan struct{}, 1)
runningErrCh := make(chan error, 1)
go func() {
defer close(runningErrCh)
defer cancel()
defer close(alreadyShutdownCh)
logger.Info("server running on 8082...")
if err := srv.ListenAndServe(); err != nil {
if !errors.Is(err, http.ErrServerClosed) {
runningErrCh <- fmt.Errorf("server experienced error: %w", err)
}
}
}()
sseErrCh := make(chan error, 1)
go func() {
defer close(sseErrCh)
defer cancel()
defer logger.Info("server side events stopped")
logger.Info("server side events streaming...")
if err := r.RunSSE(ctx); err != nil {
sseErrCh <- err
}
}()
shutdownErrCh := make(chan error, 1)
go func() {
defer close(shutdownErrCh)
select {
case <-alreadyShutdownCh:
return
case <-ctx.Done():
}
shutdownCtx, cancelShutdown := context.WithTimeout(context.Background(), 60*time.Second)
defer cancelShutdown()
if err := srv.Shutdown(shutdownCtx); err != nil {
shutdownErrCh <- fmt.Errorf("error occurred attempting to shutdown server: %w", err)
}
}()
errCh := make(chan error, 1)
go func() {
defer close(errCh)
err1 := <-runningErrCh
err2 := <-sseErrCh
err3 := <-shutdownErrCh
if err := errors.Join(err1, err2, err3); err != nil {
errCh <- err
}
}()
return r.GetSSEQueue(), errCh
}