Database migration #3
+15
-11
@@ -24,7 +24,7 @@ func main() {
|
||||
serviceVersion := "0.1.0"
|
||||
|
||||
if err := run(serviceName, serviceVersion); err != nil {
|
||||
slog.Error("Service crashed")
|
||||
slog.Error("service crashed")
|
||||
os.Exit(1)
|
||||
}
|
||||
}
|
||||
@@ -33,7 +33,7 @@ func run(serviceName string, serviceVersion string) error {
|
||||
_ = godotenv.Load(".env")
|
||||
setupLogging()
|
||||
|
||||
slog.Info("Starting service", slog.String("service", serviceName), slog.String("version", serviceVersion))
|
||||
slog.Info("starting service", slog.String("service", serviceName), slog.String("version", serviceVersion))
|
||||
|
||||
cfg := config.Load()
|
||||
|
||||
@@ -45,7 +45,7 @@ func run(serviceName string, serviceVersion string) error {
|
||||
}
|
||||
shutdownTelemetry, err := telemetry.Init(context.Background(), telCfg)
|
||||
if err != nil {
|
||||
slog.Error("Failed to initialize telemetry", err)
|
||||
slog.Error("failed to initialize telemetry", err)
|
||||
return err
|
||||
}
|
||||
defer func() {
|
||||
@@ -53,27 +53,31 @@ func run(serviceName string, serviceVersion string) error {
|
||||
defer cancel()
|
||||
|
||||
if err := shutdownTelemetry(shutdownCtx); err != nil {
|
||||
slog.Error("Failed to shutdown telemetry", slog.Any("error", err))
|
||||
slog.Error("failed to shutdown telemetry", slog.Any("error", err))
|
||||
}
|
||||
}()
|
||||
|
||||
db, err := database.New(context.Background(), cfg.DatabaseURL)
|
||||
if err != nil {
|
||||
slog.Error("Failed to initialize database", slog.Any("error", err))
|
||||
slog.Error("failed to initialize database", slog.Any("error", err))
|
||||
return err
|
||||
}
|
||||
defer func() {
|
||||
err := db.Close()
|
||||
if err != nil {
|
||||
slog.Error("Failed to close database connection", slog.Any("error", err))
|
||||
slog.Error("failed to close database connection", slog.Any("error", err))
|
||||
}
|
||||
}()
|
||||
|
||||
if err := database.RunMigrations(db); err != nil {
|
||||
slog.Error("failed to run migrations", slog.Any("error", err))
|
||||
}
|
||||
|
||||
srv := makeServer(db, cfg.Port)
|
||||
go func() {
|
||||
slog.Info("Starting http server", "addr", srv.Addr)
|
||||
slog.Info("starting http server", "addr", srv.Addr)
|
||||
if err := srv.ListenAndServe(); err != nil && !errors.Is(err, http.ErrServerClosed) {
|
||||
slog.Error("Failed to start http server", err)
|
||||
slog.Error("failed to start http server", err)
|
||||
os.Exit(1)
|
||||
}
|
||||
}()
|
||||
@@ -81,17 +85,17 @@ func run(serviceName string, serviceVersion string) error {
|
||||
quit := make(chan os.Signal, 1)
|
||||
signal.Notify(quit, syscall.SIGINT, syscall.SIGTERM)
|
||||
sig := <-quit
|
||||
slog.Info("Shutting down server...", slog.String("signal", sig.String()))
|
||||
slog.Info("shutting down server...", slog.String("signal", sig.String()))
|
||||
|
||||
ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
|
||||
defer cancel()
|
||||
|
||||
if err := srv.Shutdown(ctx); err != nil {
|
||||
slog.Error("Server forced to shutdown", slog.Any("error", err))
|
||||
slog.Error("server forced to shutdown", slog.Any("error", err))
|
||||
return err
|
||||
}
|
||||
|
||||
slog.Info("Service shutdown successful!")
|
||||
slog.Info("service shutdown successful!")
|
||||
return nil
|
||||
}
|
||||
|
||||
|
||||
@@ -3,6 +3,7 @@ module github.com/robindittmar/dttmr-api
|
||||
go 1.26
|
||||
|
||||
require (
|
||||
github.com/golang-migrate/migrate/v4 v4.19.1
|
||||
github.com/jackc/pgx/v5 v5.10.0
|
||||
github.com/joho/godotenv v1.5.1
|
||||
go.opentelemetry.io/contrib/instrumentation/net/http/otelhttp v0.69.0
|
||||
@@ -25,6 +26,7 @@ require (
|
||||
github.com/jackc/pgpassfile v1.0.0 // indirect
|
||||
github.com/jackc/pgservicefile v0.0.0-20240606120523-5a60cdf6a761 // indirect
|
||||
github.com/jackc/puddle/v2 v2.2.2 // indirect
|
||||
github.com/lib/pq v1.10.9 // indirect
|
||||
go.opentelemetry.io/auto/sdk v1.2.1 // indirect
|
||||
go.opentelemetry.io/otel/exporters/otlp/otlptrace v1.44.0 // indirect
|
||||
go.opentelemetry.io/otel/metric v1.44.0 // indirect
|
||||
|
||||
@@ -39,7 +39,7 @@ func DefaultHandler(w http.ResponseWriter, r *http.Request) {
|
||||
resp.Form[k] = v[0]
|
||||
}
|
||||
} else {
|
||||
slog.ErrorContext(ctx, "Error parsing form", slog.Any("error", err))
|
||||
slog.ErrorContext(ctx, "error parsing form", slog.Any("error", err))
|
||||
}
|
||||
|
||||
response.JSON(ctx, w, http.StatusOK, resp)
|
||||
|
||||
@@ -18,18 +18,18 @@ func (h *ListHandler) CreateList(w http.ResponseWriter, r *http.Request) {
|
||||
|
||||
payload, err := request.DecodeCreateList(r)
|
||||
if err != nil {
|
||||
slog.ErrorContext(ctx, "Failed to decode create list payload", slog.Any("error", err))
|
||||
response.Error(ctx, w, http.StatusBadRequest, "Failed to decode request body")
|
||||
slog.ErrorContext(ctx, "failed to decode create list payload", slog.Any("error", err))
|
||||
response.Error(ctx, w, http.StatusBadRequest, "failed to decode request body")
|
||||
return
|
||||
}
|
||||
|
||||
list, err := h.ListService.Create(ctx, payload.Name, payload.UserIDs)
|
||||
if err != nil {
|
||||
slog.ErrorContext(ctx, "Failed to create list", slog.Any("error", err))
|
||||
response.Error(ctx, w, http.StatusInternalServerError, "Failed to create list")
|
||||
slog.ErrorContext(ctx, "failed to create list", slog.Any("error", err))
|
||||
response.Error(ctx, w, http.StatusInternalServerError, "failed to create list")
|
||||
return
|
||||
}
|
||||
|
||||
slog.InfoContext(ctx, "Created list successfully", slog.Any("list_id", list.ID))
|
||||
slog.InfoContext(ctx, "created list successfully", slog.Any("list_id", list.ID))
|
||||
response.JSON(ctx, w, http.StatusCreated, list)
|
||||
}
|
||||
|
||||
@@ -14,7 +14,7 @@ func JSON(ctx context.Context, w http.ResponseWriter, status int, data any) {
|
||||
w.WriteHeader(http.StatusInternalServerError)
|
||||
_, err := w.Write([]byte(`{"error": "internal server error: failed to marshal response"}`))
|
||||
if err != nil {
|
||||
slog.ErrorContext(ctx, "Failed to write JSON", slog.Any("error", err))
|
||||
slog.ErrorContext(ctx, "failed to write json", slog.Any("error", err))
|
||||
return
|
||||
}
|
||||
return
|
||||
@@ -24,7 +24,7 @@ func JSON(ctx context.Context, w http.ResponseWriter, status int, data any) {
|
||||
w.WriteHeader(status)
|
||||
_, err = w.Write(payload)
|
||||
if err != nil {
|
||||
slog.ErrorContext(ctx, "Failed to write response", slog.Any("error", err))
|
||||
slog.ErrorContext(ctx, "failed to write response", slog.Any("error", err))
|
||||
return
|
||||
}
|
||||
}
|
||||
|
||||
@@ -18,7 +18,7 @@ func Load() *Config {
|
||||
envFlag := flag.String("env", "development", "environment to use")
|
||||
portFlag := flag.Int("port", 8080, "port to listen on")
|
||||
otlpEndpointFlag := flag.String("otlp-endpoint", "localhost:4317", "otlp endpoint")
|
||||
databaseUrlFlag := flag.String("database-url", "postgres://postgres:postgres@localhost:5432/postgres?sslmode=disable", "database connection string")
|
||||
databaseUrlFlag := flag.String("database-url", "postgres://postgres:postgres@localhost:5432/postgres?sslmode=disable&timezone=utc", "database connection string")
|
||||
|
||||
flag.Parse()
|
||||
|
||||
@@ -47,7 +47,7 @@ func assignIntFromEnv(key string, target *int) {
|
||||
if val, exists := os.LookupEnv(key); exists {
|
||||
parsed, err := strconv.Atoi(val)
|
||||
if err != nil {
|
||||
slog.Error("Failed to parse environment variable", slog.String("var", key), slog.Any("error", err))
|
||||
slog.Error("failed to parse environment variable", slog.String("var", key), slog.Any("error", err))
|
||||
} else {
|
||||
*target = parsed
|
||||
}
|
||||
|
||||
@@ -0,0 +1,49 @@
|
||||
package database
|
||||
|
||||
import (
|
||||
"database/sql"
|
||||
"embed"
|
||||
"errors"
|
||||
"fmt"
|
||||
"log/slog"
|
||||
|
||||
"github.com/golang-migrate/migrate/v4"
|
||||
"github.com/golang-migrate/migrate/v4/database/postgres"
|
||||
"github.com/golang-migrate/migrate/v4/source/iofs"
|
||||
)
|
||||
|
||||
//go:embed migrations/*.sql
|
||||
var migrationFS embed.FS
|
||||
|
||||
func RunMigrations(db *sql.DB) error {
|
||||
sourceDriver, err := iofs.New(migrationFS, "migrations")
|
||||
if err != nil {
|
||||
return fmt.Errorf("failed to load embedded migrations: %w", err)
|
||||
}
|
||||
defer func() {
|
||||
err := sourceDriver.Close()
|
||||
if err != nil {
|
||||
slog.Error("failed to close migrations source", slog.Any("error", err))
|
||||
}
|
||||
}()
|
||||
|
||||
dbDriver, err := postgres.WithInstance(db, &postgres.Config{})
|
||||
if err != nil {
|
||||
return fmt.Errorf("failed to create migration db driver: %w", err)
|
||||
}
|
||||
|
||||
m, err := migrate.NewWithInstance("iofs", sourceDriver, "postgres", dbDriver)
|
||||
if err != nil {
|
||||
return fmt.Errorf("failed to initialize migrator: %w", err)
|
||||
}
|
||||
|
||||
slog.Info("running database migrations...")
|
||||
err = m.Up()
|
||||
|
||||
if err != nil && !errors.Is(err, migrate.ErrNoChange) {
|
||||
return fmt.Errorf("failed to run database migrations: %w", err)
|
||||
}
|
||||
|
||||
slog.Info("database migrations applied successfully")
|
||||
return nil
|
||||
}
|
||||
@@ -0,0 +1,2 @@
|
||||
DROP TABLE IF EXISTS sessions;
|
||||
DROP TABLE IF EXISTS users;
|
||||
@@ -0,0 +1,16 @@
|
||||
CREATE TABLE IF NOT EXISTS users (
|
||||
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
|
||||
email VARCHAR(255) UNIQUE NOT NULL,
|
||||
name VARCHAR(255) NOT NULL,
|
||||
password_hash VARCHAR(255) NOT NULL,
|
||||
created_at TIMESTAMPTZ NOT NULL DEFAULT NOW()
|
||||
);
|
||||
|
||||
CREATE TABLE IF NOT EXISTS sessions (
|
||||
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
|
||||
user_id UUID NOT NULL REFERENCES users(id) ON DELETE CASCADE,
|
||||
token_id UUID UNIQUE NOT NULL,
|
||||
is_revoked BOOLEAN NOT NULL DEFAULT FALSE,
|
||||
expires_at TIMESTAMPTZ NOT NULL,
|
||||
created_at TIMESTAMPTZ NOT NULL DEFAULT NOW()
|
||||
);
|
||||
@@ -0,0 +1,3 @@
|
||||
DROP TABLE IF EXISTS list_items;
|
||||
DROP TABLE IF EXISTS list_users;
|
||||
DROP TABLE IF EXISTS lists;
|
||||
@@ -0,0 +1,26 @@
|
||||
CREATE TABLE IF NOT EXISTS lists (
|
||||
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
|
||||
name VARCHAR(255) NOT NULL,
|
||||
created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
|
||||
modified_at TIMESTAMPTZ NOT NULL DEFAULT NOW()
|
||||
);
|
||||
|
||||
CREATE TABLE IF NOT EXISTS list_users (
|
||||
list_id UUID REFERENCES lists(id) ON DELETE CASCADE,
|
||||
user_id UUID NOT NULL,
|
||||
created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
|
||||
PRIMARY KEY (list_id, user_id)
|
||||
);
|
||||
|
||||
CREATE INDEX idx_list_users_user_id ON list_users(user_id);
|
||||
|
||||
CREATE TABLE IF NOT EXISTS list_items (
|
||||
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
|
||||
list_id UUID REFERENCES lists(id) ON DELETE CASCADE,
|
||||
title VARCHAR(255) NOT NULL,
|
||||
is_completed BOOLEAN NOT NULL DEFAULT FALSE,
|
||||
created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
|
||||
modified_at TIMESTAMPTZ NOT NULL DEFAULT NOW()
|
||||
);
|
||||
|
||||
CREATE INDEX idx_list_items_list_id ON list_items(list_id);
|
||||
@@ -7,9 +7,10 @@ import (
|
||||
)
|
||||
|
||||
type List struct {
|
||||
ID string `json:"id"`
|
||||
Name string `json:"name"`
|
||||
CreatedAt time.Time `json:"created_at"`
|
||||
ID string `json:"id"`
|
||||
Name string `json:"name"`
|
||||
CreatedAt time.Time `json:"created_at"`
|
||||
ModifiedAt time.Time `json:"modified_at"`
|
||||
}
|
||||
|
||||
type ListRepository interface {
|
||||
|
||||
@@ -22,19 +22,14 @@ func (r *ListRepo) CreateList(ctx context.Context, name string, userIDs []string
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("begin transaction: %w", err)
|
||||
}
|
||||
defer func() {
|
||||
err := tx.Rollback()
|
||||
if err != nil {
|
||||
slog.Error("Failed to rollback transaction", slog.Any("error", err))
|
||||
}
|
||||
}()
|
||||
defer tx.Rollback()
|
||||
|
||||
list := &domain.List{Name: name}
|
||||
|
||||
err = tx.QueryRowContext(ctx,
|
||||
"INSERT INTO lists (name) VALUES ($1) RETURNING id, created_at",
|
||||
"INSERT INTO lists (name) VALUES ($1) RETURNING id, created_at, modified_at",
|
||||
name,
|
||||
).Scan(&list.ID, &list.CreatedAt)
|
||||
).Scan(&list.ID, &list.CreatedAt, &list.ModifiedAt)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("failed to insert list: %w", err)
|
||||
}
|
||||
@@ -46,7 +41,7 @@ func (r *ListRepo) CreateList(ctx context.Context, name string, userIDs []string
|
||||
defer func() {
|
||||
err := stmt.Close()
|
||||
if err != nil {
|
||||
|
||||
slog.Error("failed to close user/list association statement", slog.Any("error", err))
|
||||
}
|
||||
}()
|
||||
|
||||
|
||||
@@ -53,7 +53,7 @@ func Init(ctx context.Context, cfg Config) (func(context.Context) error, error)
|
||||
)
|
||||
|
||||
metricExporter, err := otlpmetricgrpc.New(ctx,
|
||||
otlpmetricgrpc.WithInsecure(),
|
||||
otlpmetricgrpc.WithInsecure(), // TODO: for local development
|
||||
otlpmetricgrpc.WithEndpoint(cfg.Endpoint),
|
||||
)
|
||||
if err != nil {
|
||||
|
||||
Reference in New Issue
Block a user