diff --git a/docker-compose.yml b/docker-compose.yml index 90a3105..121d453 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -138,9 +138,15 @@ services: environment: - PORT=8084 - DATABASE_URL=postgres://auron:auron_pass@payments-db:5432/payments_db?sslmode=disable + - REDIS_URL=redis://redis:6379/0 + - KAFKA_BROKERS=kafka:29092 + - STRIPE_SECRET_KEY=${STRIPE_SECRET_KEY} + - STRIPE_WEBHOOK_SECRET=${STRIPE_WEBHOOK_SECRET} depends_on: payments-db: condition: service_healthy + kafka: + condition: service_healthy networks: - auron-network restart: unless-stopped @@ -289,6 +295,8 @@ services: volumes: - zookeeper_data:/var/lib/zookeeper/data - zookeeper_log:/var/lib/zookeeper/log + networks: + - auron-network healthcheck: test: nc -z localhost 2181 || exit -1 interval: 10s @@ -312,6 +320,8 @@ services: KAFKA_AUTO_CREATE_TOPICS_ENABLE: "true" volumes: - kafka_data:/var/lib/kafka/data + networks: + - auron-network healthcheck: test: [ @@ -332,6 +342,8 @@ services: environment: KAFKA_CLUSTERS_0_NAME: local KAFKA_CLUSTERS_0_BOOTSTRAPSERVERS: kafka:29092 + networks: + - auron-network # ============================================================ # CACHE diff --git a/services/payment-service/.env.example b/services/payment-service/.env.example new file mode 100644 index 0000000..acd43c3 --- /dev/null +++ b/services/payment-service/.env.example @@ -0,0 +1,8 @@ +# Payment Service +PORT=8084 +DATABASE_URL=postgres://auron:auron_pass@localhost:5435/payments_db?sslmode=disable +REDIS_URL=redis://localhost:6379/0 +KAFKA_BROKERS=localhost:9092 +STRIPE_SECRET_KEY=sk_test_... +STRIPE_WEBHOOK_SECRET=whsec_... +GORM_LOG_LEVEL=warn diff --git a/services/payment-service/Dockerfile b/services/payment-service/Dockerfile new file mode 100644 index 0000000..bf50102 --- /dev/null +++ b/services/payment-service/Dockerfile @@ -0,0 +1,23 @@ +FROM golang:1.25-alpine AS builder + +WORKDIR /app + +RUN apk add --no-cache git + +COPY go.mod go.sum ./ + +COPY . . + +RUN CGO_ENABLED=0 GOOS=linux go build -o /payment-service . + +FROM alpine:3.18 + +RUN apk add --no-cache ca-certificates curl + +WORKDIR /app + +COPY --from=builder /payment-service . + +EXPOSE 8084 + +CMD ["./payment-service"] diff --git a/services/payment-service/IMPLEMENTATION_PLAN.md b/services/payment-service/IMPLEMENTATION_PLAN.md new file mode 100644 index 0000000..e0d20d7 --- /dev/null +++ b/services/payment-service/IMPLEMENTATION_PLAN.md @@ -0,0 +1,399 @@ +# Payment Service — Implementation Plan + +## Overview + +The payment service (port **8084**) integrates with **Stripe** to handle payment processing for orders. +It is a **Kafka consumer + HTTP API** hybrid service: + +- Consumes `order.created` events → creates a Stripe PaymentIntent → stores Payment record +- Exposes HTTP endpoints for payment lookup and Stripe webhook ingestion +- Publishes `payment.created`, `payment.completed`, `payment.failed` Kafka events for downstream consumers (e.g., order-service to update order status, notification-service to email the user) + +### Gateway Routes (already wired) + +| Method | Path | Auth | Purpose | +|---|---|---|---| +| `GET` | `/api/payments/:id` | Required | Get payment details by payment UUID | +| `POST` | `/api/payments/webhook/stripe` | None (Stripe signs) | Stripe webhook handler | + +### Payment Lifecycle + +``` +Frontend creates order + ↓ +order-service → publishes order.created (Kafka) + ↓ +payment-service consumes order.created + ↓ +Creates Stripe PaymentIntent → stores Payment(status=pending, client_secret) + ↓ +Publishes payment.created (contains payment_id + client_secret for frontend) + ↓ +Frontend uses client_secret + Stripe.js to confirm payment + ↓ +Stripe fires POST /api/payments/webhook/stripe + ↓ +payment-service verifies webhook signature → updates status + ↓ +Publishes payment.completed or payment.failed +``` + +--- + +## Folder Structure + +``` +services/payment-service/ +├── cmd/ +│ ├── config.go # env vars → appConfig struct +│ ├── dotenv.go # load .env file in non-production +│ ├── infrastructure.go # setupDatabase, setupRedis, runMigrations +│ ├── kafka.go # setupKafkaPublisher, startKafkaConsumer +│ ├── run.go # wire everything together +│ └── server.go # setupRouter, registerGracefulShutdown +├── db/ +│ └── 001_create_payments.up.sql +├── internal/ +│ ├── cache/ +│ │ └── payment_cache.go +│ ├── client/ +│ │ └── stripe_client.go +│ ├── domain/ +│ │ ├── payment.go # Payment entity, PaymentStatus, DTOs +│ │ ├── errors.go # sentinel errors +│ │ ├── repository.go # PaymentRepository interface +│ │ ├── service.go # PaymentService interface +│ │ ├── cache.go # PaymentCache interface +│ │ ├── events.go # EventPublisher interface + topic constants +│ │ └── client.go # StripeClient interface +│ ├── events/ +│ │ ├── kafka_publisher.go +│ │ └── kafka_consumer.go +│ ├── handler/ +│ │ └── payment_handler.go +│ ├── middleware/ +│ │ └── stripe_webhook.go # raw body capture for signature verification +│ ├── repository/ +│ │ └── payment_repository.go +│ ├── route/ +│ │ └── payment_route.go +│ └── service/ +│ └── payment_service.go +├── main.go +├── Dockerfile +├── go.mod +├── .env +└── .env.example +``` + +--- + +## Tasks + +### Task 1 — Domain Layer + +Create all files under `internal/domain/`: + +**`payment.go`** +- `PaymentStatus` type (`pending`, `processing`, `completed`, `failed`, `refunded`) +- `Payment` struct with GORM tags: + - `id uuid`, `order_id uuid` (unique index), `user_id uuid`, `amount float64`, `currency varchar(10) default 'usd'` + - `status varchar(50) default 'pending'`, `stripe_payment_intent_id varchar(255)`, `stripe_client_secret text` + - `failure_reason text`, `created_at`, `updated_at` +- `PaymentResponse` DTO — excludes `stripe_client_secret` for normal reads +- `PaymentInitResponse` DTO — includes `stripe_client_secret` (returned only on `payment.created` event, never via HTTP) +- `OrderCreatedEvent` struct — shape of the Kafka message from order-service: `{order_id, user_id, total_amount, items[]}` + +**`errors.go`** +- `ErrPaymentNotFound`, `ErrPaymentAlreadyExists`, `ErrInvalidWebhookSignature`, `ErrForbidden`, `ErrUnauthorized` + +**`repository.go`** +```go +type PaymentRepository interface { + GetPaymentByID(id uuid.UUID) (*Payment, error) + GetPaymentByOrderID(orderID uuid.UUID) (*Payment, error) + CreatePayment(payment *Payment) (*Payment, error) + UpdatePaymentStatus(id uuid.UUID, status PaymentStatus, failureReason string) (*Payment, error) + UpdateStripePaymentIntentID(id uuid.UUID, intentID, clientSecret string) (*Payment, error) +} +``` + +**`service.go`** +```go +type PaymentService interface { + GetPaymentByID(ctx context.Context, userID, paymentID uuid.UUID) (*PaymentResponse, error) + HandleOrderCreated(ctx context.Context, event OrderCreatedEvent) error + HandleStripeWebhook(ctx context.Context, payload []byte, signature string) error +} +``` + +**`cache.go`** +```go +type PaymentCache interface { + GetPayment(ctx context.Context, paymentID uuid.UUID) (*Payment, error) + SetPayment(ctx context.Context, payment *Payment) error + InvalidatePayment(ctx context.Context, paymentID uuid.UUID) error +} +``` + +**`events.go`** +- `EventPublisher` interface with `Publish(topic string, key string, payload any) error` and `Close() error` +- Constants: `TopicPaymentCreated = "payment.created"`, `TopicPaymentCompleted = "payment.completed"`, `TopicPaymentFailed = "payment.failed"` + +**`client.go`** +```go +type StripeClient interface { + CreatePaymentIntent(ctx context.Context, amount float64, currency string, metadata map[string]string) (intentID, clientSecret string, err error) +} +``` + +--- + +### Task 2 — DB Migration + +**`db/001_create_payments.up.sql`** +```sql +CREATE TABLE IF NOT EXISTS payments ( + id UUID PRIMARY KEY DEFAULT gen_random_uuid(), + order_id UUID NOT NULL, + user_id UUID NOT NULL, + amount DECIMAL(12,2) NOT NULL, + currency VARCHAR(10) NOT NULL DEFAULT 'usd', + status VARCHAR(50) NOT NULL DEFAULT 'pending', + stripe_payment_intent_id VARCHAR(255), + stripe_client_secret TEXT, + failure_reason TEXT, + created_at TIMESTAMP NOT NULL DEFAULT NOW(), + updated_at TIMESTAMP NOT NULL DEFAULT NOW(), + CONSTRAINT chk_payments_status CHECK ( + status IN ('pending','processing','completed','failed','refunded') + ), + CONSTRAINT chk_payments_amount CHECK (amount > 0) +); + +CREATE UNIQUE INDEX IF NOT EXISTS idx_payments_order_id ON payments(order_id); +CREATE INDEX IF NOT EXISTS idx_payments_user_id ON payments(user_id); +CREATE INDEX IF NOT EXISTS idx_payments_status ON payments(status); +``` + +--- + +### Task 3 — Repository Layer + +**`internal/repository/payment_repository.go`** +- GORM implementation of `domain.PaymentRepository` +- `GetPaymentByID` and `GetPaymentByOrderID` return `ErrPaymentNotFound` on GORM `ErrRecordNotFound` +- `UpdatePaymentStatus`: updates `status`, `failure_reason`, and `updated_at` in a single `db.Model().Updates()` call +- `UpdateStripePaymentIntentID`: sets `stripe_payment_intent_id` and `stripe_client_secret` + +--- + +### Task 4 — Cache Layer + +**`internal/cache/payment_cache.go`** +- `PaymentCache` struct wrapping `*redis.Client` +- Key: `payment:` (TTL 1h) +- JSON marshal/unmarshal; miss returns `nil, nil` + +--- + +### Task 5 — Kafka Events (Publisher + Consumer) + +**`internal/events/kafka_publisher.go`** +- Same pattern as order-service: `kafkaPublisher` with `writers map[string]*kafka.Writer` +- `Publish(topic, key string, payload any) error` — JSON-marshals payload, writes message +- `Close() error` + +**`internal/events/kafka_consumer.go`** +- `KafkaConsumer` struct: `reader *kafka.Reader`, `paymentService domain.PaymentService`, `logger` +- `Start(ctx context.Context)` — goroutine reading messages from `order.created` topic, group `payment-service` +- On message: unmarshal `domain.OrderCreatedEvent`, call `paymentService.HandleOrderCreated(ctx, event)` +- Log errors, commit offset, never crash — errors are non-fatal +- `Close() error` + +--- + +### Task 6 — Stripe Client + +**`internal/client/stripe_client.go`** +- `stripeClient` struct with `secretKey string` +- Implements `domain.StripeClient` +- `CreatePaymentIntent`: calls Stripe Go SDK `paymentintent.New()` with amount (converted to cents), currency, and metadata (`order_id`, `user_id`) +- Returns `intentID` and `clientSecret` + +**Dependencies to add:** +``` +github.com/stripe/stripe-go/v76 +``` + +--- + +### Task 7 — Service Layer + +**`internal/service/payment_service.go`** + +**`HandleOrderCreated(ctx, event)`** +1. Check `GetPaymentByOrderID` — if already exists, return nil (idempotent) +2. Call `stripeClient.CreatePaymentIntent(ctx, event.TotalAmount, "usd", metadata)` +3. Build and `CreatePayment` record (status=pending, stripe IDs set) +4. Cache the payment +5. Publish `payment.created` event asynchronously (contains `payment_id`, `order_id`, `user_id`, `client_secret`) + +**`GetPaymentByID(ctx, userID, paymentID)`** +1. Cache-aside: check cache first +2. DB fallback on miss +3. Ownership check: `payment.UserID != userID` → `ErrForbidden` +4. Return `PaymentResponse` (no client_secret) + +**`HandleStripeWebhook(ctx, payload, signature)`** +1. Construct Stripe event: `webhook.ConstructEvent(payload, signature, webhookSecret)` → error → `ErrInvalidWebhookSignature` +2. Switch on event type: + - `payment_intent.succeeded` → `UpdatePaymentStatus(completed)` → publish `payment.completed` async + - `payment_intent.payment_failed` → `UpdatePaymentStatus(failed, failureReason)` → publish `payment.failed` async + - `payment_intent.processing` → `UpdatePaymentStatus(processing)` +3. Invalidate and re-cache payment after status update + +--- + +### Task 8 — Handler + Route Layers + +**`internal/handler/payment_handler.go`** +- `PaymentHandler` struct with `paymentService domain.PaymentService` +- `getUserID(c *gin.Context) (uuid.UUID, error)` — reads `X-User-ID` header +- `GetPaymentByID(c *gin.Context)` — parse `:id` param, call service, return 200/404/403 +- `HandleStripeWebhook(c *gin.Context)` — reads raw body (from context, set by middleware), reads `Stripe-Signature` header, calls service, always returns 200 (Stripe retries on non-200) +- `handleError(c, err)` — maps domain errors to status codes + +**`internal/middleware/stripe_webhook.go`** +- Gin middleware that reads and buffers the raw request body into `c.Set("rawBody", body)` before `c.Next()` +- Required because Stripe signature verification needs the exact raw bytes, and `c.Request.Body` is consumed after `ShouldBindJSON` + +**`internal/route/payment_route.go`** +```go +func RegisterPaymentRoutes(router *gin.Engine, paymentHandler *handler.PaymentHandler) { + api := router.Group("/") + api.GET("/payments/:id", paymentHandler.GetPaymentByID) + api.POST("/payments/webhook/stripe", middleware.CaptureRawBody(), paymentHandler.HandleStripeWebhook) +} +``` + +--- + +### Task 9 — cmd Bootstrap + +**`cmd/config.go`** +```go +type appConfig struct { + Port string + DatabaseURL string + RedisURL string + KafkaBrokers string + StripeSecretKey string + StripeWebhookSecret string +} +``` +Defaults: port 8084, localhost:5435, localhost:6379, localhost:9092 + +**`cmd/dotenv.go`** — identical pattern to order-service + +**`cmd/infrastructure.go`** +- `setupDatabase` with connection pooling +- `runMigrations` — AutoMigrate `domain.Payment` +- `setupRedis` — ParseURL + Ping + +**`cmd/kafka.go`** +- `paymentTopics`: TopicPaymentCreated, TopicPaymentCompleted, TopicPaymentFailed +- `setupKafkaPublisher(brokers string) domain.EventPublisher` +- `setupKafkaConsumer(brokers string, svc domain.PaymentService) *events.KafkaConsumer` +- `startKafkaConsumer(consumer *events.KafkaConsumer)` — launches goroutine + +**`cmd/run.go`** +```go +func Run() { + cfg := loadConfig() + db := setupDatabase(cfg.DatabaseURL) + runMigrations(db) + redisClient := setupRedis(cfg.RedisURL) + publisher := setupKafkaPublisher(cfg.KafkaBrokers) + paymentRepo := repository.NewPaymentRepository(db) + paymentCache := cache.NewPaymentCache(redisClient) + stripeClient := client.NewStripeClient(cfg.StripeSecretKey) + paymentSvc := service.NewPaymentService(paymentRepo, paymentCache, stripeClient, publisher, cfg.StripeWebhookSecret) + consumer := setupKafkaConsumer(cfg.KafkaBrokers, paymentSvc) + startKafkaConsumer(consumer) + paymentHandler := handler.NewPaymentHandler(paymentSvc) + router := setupRouter(paymentHandler) + registerGracefulShutdown(db, redisClient, publisher, consumer) + router.Run(fmt.Sprintf(":%s", cfg.Port)) +} +``` + +**`cmd/server.go`** +- `setupRouter(*handler.PaymentHandler) *gin.Engine` — release mode, /health, /metrics, calls `RegisterPaymentRoutes` +- `registerGracefulShutdown` — SIGTERM/SIGINT handler closing DB, Redis, Kafka publisher and consumer + +--- + +### Task 10 — Entry Point, Dockerfile, and Env Files + +**`main.go`** — `cmd.Run()` + +**`Dockerfile`** — same multi-stage pattern: `golang:1.25-alpine` builder → `alpine:3.18` runtime; binary named `payment-service`; EXPOSE 8084 + +**`.env`** — local dev values +``` +PORT=8084 +DATABASE_URL=postgres://auron:auron_pass@localhost:5435/payments_db?sslmode=disable +REDIS_URL=redis://localhost:6379/0 +KAFKA_BROKERS=localhost:9092 +STRIPE_SECRET_KEY=sk_test_... +STRIPE_WEBHOOK_SECRET=whsec_... +GORM_LOG_LEVEL=warn +``` + +**`.env.example`** — same with placeholder values + +--- + +### Task 11 — docker-compose Wiring + +Update `docker-compose.yml` `payment-service` environment block: +```yaml +- REDIS_URL=redis://redis:6379/0 +- KAFKA_BROKERS=kafka:29092 +- STRIPE_SECRET_KEY=${STRIPE_SECRET_KEY} +- STRIPE_WEBHOOK_SECRET=${STRIPE_WEBHOOK_SECRET} +``` + +Add `kafka` to `payment-service.depends_on` (after payments-db). + +--- + +## Key Design Decisions + +| Decision | Choice | Reason | +|---|---|---| +| Stripe integration | PaymentIntents API | Supports SCA, supports card, wallet, BNPL via `automatic_payment_methods` | +| Payment initiation | Kafka consumer (`order.created`) | Decoupled — order-service doesn't need to call payment-service HTTP | +| Webhook raw body | Middleware that caches raw bytes | Stripe signature verification requires exact bytes; Gin's binding consumes the body | +| Idempotency | Check `GetPaymentByOrderID` before creating | Prevents duplicate Stripe intents if `order.created` is delivered multiple times | +| client_secret exposure | Only via `payment.created` Kafka event | Never exposed via HTTP API to avoid interception; downstream services forward to frontend | +| Stripe amount | `int64(amount * 100)` cents | Stripe API requires smallest currency unit | +| Webhook response | Always return 200 | Stripe retries on 4xx/5xx; log errors but don't fail the HTTP response | +| KafkaBrokers for consumer | `order.created` topic, group `payment-service` | Group ID ensures each message is processed exactly once per service instance | +| go.mod module | `auron/payment-service` | Matches pattern of all other services | + +--- + +## Dependencies + +``` +github.com/gin-gonic/gin v1.12.0 +github.com/google/uuid v1.6.0 +github.com/redis/go-redis/v9 v9.19.0 +github.com/segmentio/kafka-go v0.4.51 +gorm.io/driver/postgres v1.6.0 +gorm.io/gorm v1.31.1 +github.com/stripe/stripe-go/v76 v76.x.x +github.com/joho/godotenv v1.5.1 +``` diff --git a/services/payment-service/cmd/config.go b/services/payment-service/cmd/config.go new file mode 100644 index 0000000..0d60e07 --- /dev/null +++ b/services/payment-service/cmd/config.go @@ -0,0 +1,45 @@ +package cmd + +import "os" + +type appConfig struct { + Port string + DatabaseURL string + RedisURL string + KafkaBrokers string + StripeSecretKey string + StripeWebhookSecret string +} + +func loadConfig() appConfig { + loadDotEnvFile(".env") + + port := os.Getenv("PORT") + if port == "" { + port = "8084" + } + + databaseURL := os.Getenv("DATABASE_URL") + if databaseURL == "" { + databaseURL = "postgres://auron:auron_pass@localhost:5435/payments_db?sslmode=disable" + } + + redisURL := os.Getenv("REDIS_URL") + if redisURL == "" { + redisURL = "redis://localhost:6379/0" + } + + kafkaBrokers := os.Getenv("KAFKA_BROKERS") + if kafkaBrokers == "" { + kafkaBrokers = "localhost:9092" + } + + return appConfig{ + Port: port, + DatabaseURL: databaseURL, + RedisURL: redisURL, + KafkaBrokers: kafkaBrokers, + StripeSecretKey: os.Getenv("STRIPE_SECRET_KEY"), + StripeWebhookSecret: os.Getenv("STRIPE_WEBHOOK_SECRET"), + } +} diff --git a/services/payment-service/cmd/dotenv.go b/services/payment-service/cmd/dotenv.go new file mode 100644 index 0000000..3c652db --- /dev/null +++ b/services/payment-service/cmd/dotenv.go @@ -0,0 +1,42 @@ +package cmd + +import ( + "bufio" + "os" + "strings" +) + +func loadDotEnvFile(path string) { + file, err := os.Open(path) + if err != nil { + return + } + defer file.Close() + + scanner := bufio.NewScanner(file) + for scanner.Scan() { + line := strings.TrimSpace(scanner.Text()) + if line == "" || strings.HasPrefix(line, "#") { + continue + } + + parts := strings.SplitN(line, "=", 2) + if len(parts) != 2 { + continue + } + + key := strings.TrimSpace(parts[0]) + value := strings.TrimSpace(parts[1]) + value = strings.Trim(value, `"'`) + + if key == "" { + continue + } + + if _, exists := os.LookupEnv(key); exists { + continue + } + + _ = os.Setenv(key, value) + } +} diff --git a/services/payment-service/cmd/infrastructure.go b/services/payment-service/cmd/infrastructure.go new file mode 100644 index 0000000..ac488d4 --- /dev/null +++ b/services/payment-service/cmd/infrastructure.go @@ -0,0 +1,71 @@ +package cmd + +import ( + "context" + "os" + "strings" + "time" + + "auron/payment-service/internal/domain" + + "github.com/redis/go-redis/v9" + "gorm.io/driver/postgres" + "gorm.io/gorm" + "gorm.io/gorm/logger" +) + +func setupDatabase(databaseURL string) (*gorm.DB, error) { + db, err := gorm.Open(postgres.Open(databaseURL), &gorm.Config{ + Logger: logger.Default.LogMode(resolveGormLogLevel()), + }) + if err != nil { + return nil, err + } + + sqlDB, err := db.DB() + if err != nil { + return nil, err + } + sqlDB.SetMaxIdleConns(10) + sqlDB.SetMaxOpenConns(100) + sqlDB.SetConnMaxLifetime(time.Hour) + + return db, nil +} + +func runMigrations(db *gorm.DB) error { + return db.AutoMigrate(&domain.Payment{}) +} + +func setupRedis(redisURL string) (*redis.Client, error) { + opt, err := redis.ParseURL(redisURL) + if err != nil { + return nil, err + } + + client := redis.NewClient(opt) + + ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second) + defer cancel() + + if err := client.Ping(ctx).Err(); err != nil { + return nil, err + } + + return client, nil +} + +func resolveGormLogLevel() logger.LogLevel { + switch strings.ToLower(strings.TrimSpace(os.Getenv("GORM_LOG_LEVEL"))) { + case "silent": + return logger.Silent + case "error": + return logger.Error + case "warn", "warning": + return logger.Warn + case "info": + return logger.Info + default: + return logger.Warn + } +} diff --git a/services/payment-service/cmd/kafka.go b/services/payment-service/cmd/kafka.go new file mode 100644 index 0000000..5e95b7e --- /dev/null +++ b/services/payment-service/cmd/kafka.go @@ -0,0 +1,110 @@ +package cmd + +import ( + "context" + "log/slog" + "strconv" + "strings" + "time" + + "auron/payment-service/internal/domain" + "auron/payment-service/internal/events" + + "github.com/segmentio/kafka-go" +) + +var paymentTopics = []string{ + domain.TopicPaymentCreated, + domain.TopicPaymentCompleted, + domain.TopicPaymentFailed, +} + +func setupKafkaPublisher(kafkaBrokers string) domain.EventPublisher { + brokers := parseBrokers(kafkaBrokers) + ensureTopics(brokers, paymentTopics) + + writers := make(map[string]*kafka.Writer, len(paymentTopics)) + for _, topic := range paymentTopics { + writers[topic] = &kafka.Writer{ + Addr: kafka.TCP(brokers...), + Topic: topic, + Balancer: &kafka.LeastBytes{}, + RequiredAcks: kafka.RequireOne, + BatchTimeout: 10 * time.Millisecond, + } + } + + return events.NewKafkaPublisher(writers) +} + +func setupKafkaConsumer(kafkaBrokers string, svc domain.PaymentService) *events.KafkaConsumer { + brokers := parseBrokers(kafkaBrokers) + return events.NewKafkaConsumer(brokers, domain.TopicOrderCreated, "payment-service", svc) +} + +func startKafkaConsumer(ctx context.Context, consumer *events.KafkaConsumer) { + consumer.Start(ctx) + slog.Info("kafka consumer started", "topic", domain.TopicOrderCreated, "group", "payment-service") +} + +func parseBrokers(kafkaBrokers string) []string { + parts := strings.Split(kafkaBrokers, ",") + brokers := make([]string, 0, len(parts)) + for _, b := range parts { + if trimmed := strings.TrimSpace(b); trimmed != "" { + brokers = append(brokers, trimmed) + } + } + if len(brokers) == 0 { + return []string{"localhost:9092"} + } + return brokers +} + +func ensureTopics(brokers []string, topics []string) { + if len(brokers) == 0 || len(topics) == 0 { + return + } + + conn, err := kafka.Dial("tcp", brokers[0]) + if err != nil { + slog.Warn("kafka topic init skipped: cannot connect", "broker", brokers[0], "error", err) + return + } + defer conn.Close() + + controller, err := conn.Controller() + if err != nil { + slog.Warn("kafka topic init skipped: cannot get controller", "error", err) + return + } + + controllerConn, err := kafka.Dial("tcp", controller.Host+":"+strconv.Itoa(controller.Port)) + if err != nil { + slog.Warn("kafka topic init skipped: cannot connect to controller", "error", err) + return + } + defer controllerConn.Close() + + configs := make([]kafka.TopicConfig, 0, len(topics)) + for _, topic := range topics { + configs = append(configs, kafka.TopicConfig{ + Topic: topic, + NumPartitions: 3, + ReplicationFactor: 1, + }) + } + + if err := controllerConn.CreateTopics(configs...); err != nil { + slog.Warn("kafka topic init failed", "topics", topics, "error", err) + return + } + + slog.Info("kafka topics ensured", "topics", topics) +} + +func closeKafkaPublisher(publisher domain.EventPublisher) { + if err := publisher.Close(); err != nil { + slog.Warn("failed to close kafka publisher", "error", err) + } +} diff --git a/services/payment-service/cmd/run.go b/services/payment-service/cmd/run.go new file mode 100644 index 0000000..467bdf5 --- /dev/null +++ b/services/payment-service/cmd/run.go @@ -0,0 +1,90 @@ +package cmd + +import ( + "context" + "fmt" + "log" + "log/slog" + "os" + "os/signal" + "syscall" + + "auron/payment-service/internal/cache" + "auron/payment-service/internal/client" + "auron/payment-service/internal/domain" + "auron/payment-service/internal/events" + "auron/payment-service/internal/handler" + "auron/payment-service/internal/repository" + "auron/payment-service/internal/service" + + "github.com/redis/go-redis/v9" + "gorm.io/gorm" +) + +func Run() { + cfg := loadConfig() + + db, err := setupDatabase(cfg.DatabaseURL) + if err != nil { + log.Fatalf("failed to connect to database: %v", err) + } + + if err := runMigrations(db); err != nil { + log.Fatalf("failed to run migrations: %v", err) + } + log.Println("database migrations completed") + + redisClient, err := setupRedis(cfg.RedisURL) + if err != nil { + log.Fatalf("failed to connect to Redis: %v", err) + } + + publisher := setupKafkaPublisher(cfg.KafkaBrokers) + + paymentRepo := repository.NewPaymentRepository(db) + paymentCache := cache.NewPaymentCache(redisClient) + stripeClient := client.NewStripeClient(cfg.StripeSecretKey) + + paymentSvc := service.NewPaymentService(paymentRepo, paymentCache, stripeClient, publisher, cfg.StripeWebhookSecret) + + ctx, cancel := context.WithCancel(context.Background()) + consumer := setupKafkaConsumer(cfg.KafkaBrokers, paymentSvc) + startKafkaConsumer(ctx, consumer) + + paymentHandler := handler.NewPaymentHandler(paymentSvc) + router := setupRouter(paymentHandler) + + registerGracefulShutdown(db, redisClient, publisher, consumer, cancel) + + addr := fmt.Sprintf(":%s", cfg.Port) + log.Printf("starting payment-service on %s", addr) + if err := router.Run(addr); err != nil { + log.Fatalf("failed to start server: %v", err) + } +} + +func registerGracefulShutdown( + db *gorm.DB, + redisClient *redis.Client, + publisher domain.EventPublisher, + consumer *events.KafkaConsumer, + cancel context.CancelFunc, +) { + quit := make(chan os.Signal, 1) + signal.Notify(quit, syscall.SIGINT, syscall.SIGTERM) + + go func() { + <-quit + fmt.Println("\nshutting down payment-service...") + cancel() + if err := consumer.Close(); err != nil { + slog.Warn("error closing kafka consumer", "error", err) + } + if sqlDB, err := db.DB(); err == nil { + _ = sqlDB.Close() + } + _ = redisClient.Close() + closeKafkaPublisher(publisher) + os.Exit(0) + }() +} diff --git a/services/payment-service/cmd/server.go b/services/payment-service/cmd/server.go new file mode 100644 index 0000000..299fbed --- /dev/null +++ b/services/payment-service/cmd/server.go @@ -0,0 +1,34 @@ +package cmd + +import ( + "time" + + "auron/payment-service/internal/handler" + "auron/payment-service/internal/route" + + "github.com/gin-gonic/gin" +) + +func setupRouter(paymentHandler *handler.PaymentHandler) *gin.Engine { + gin.SetMode(gin.ReleaseMode) + router := gin.New() + + router.Use(gin.Logger()) + router.Use(gin.Recovery()) + + router.GET("/health", func(c *gin.Context) { + c.JSON(200, gin.H{ + "status": "healthy", + "service": "payment-service", + "timestamp": time.Now().UTC(), + }) + }) + + router.GET("/metrics", func(c *gin.Context) { + c.String(200, "# Prometheus metrics endpoint\n") + }) + + route.RegisterPaymentRoutes(router, paymentHandler) + + return router +} diff --git a/services/payment-service/db/001_create_payments.up.sql b/services/payment-service/db/001_create_payments.up.sql new file mode 100644 index 0000000..f61c66d --- /dev/null +++ b/services/payment-service/db/001_create_payments.up.sql @@ -0,0 +1,21 @@ +CREATE TABLE IF NOT EXISTS payments ( + id UUID PRIMARY KEY DEFAULT gen_random_uuid(), + order_id UUID NOT NULL, + user_id UUID NOT NULL, + amount DECIMAL(12,2) NOT NULL, + currency VARCHAR(10) NOT NULL DEFAULT 'usd', + status VARCHAR(50) NOT NULL DEFAULT 'pending', + stripe_payment_intent_id VARCHAR(255), + stripe_client_secret TEXT, + failure_reason TEXT, + created_at TIMESTAMP NOT NULL DEFAULT NOW(), + updated_at TIMESTAMP NOT NULL DEFAULT NOW(), + CONSTRAINT chk_payments_status CHECK ( + status IN ('pending','processing','completed','failed','refunded') + ), + CONSTRAINT chk_payments_amount CHECK (amount > 0) +); + +CREATE UNIQUE INDEX IF NOT EXISTS idx_payments_order_id ON payments(order_id); +CREATE INDEX IF NOT EXISTS idx_payments_user_id ON payments(user_id); +CREATE INDEX IF NOT EXISTS idx_payments_status ON payments(status); diff --git a/services/payment-service/go.mod b/services/payment-service/go.mod new file mode 100644 index 0000000..ec498a1 --- /dev/null +++ b/services/payment-service/go.mod @@ -0,0 +1,56 @@ +module auron/payment-service + +go 1.25.8 + +require ( + github.com/gin-gonic/gin v1.12.0 + github.com/google/uuid v1.6.0 + github.com/redis/go-redis/v9 v9.19.0 + github.com/segmentio/kafka-go v0.4.51 + github.com/stripe/stripe-go/v76 v76.25.0 + gorm.io/driver/postgres v1.6.0 + gorm.io/gorm v1.31.1 +) + +require ( + github.com/bytedance/gopkg v0.1.3 // indirect + github.com/bytedance/sonic v1.15.0 // indirect + github.com/bytedance/sonic/loader v0.5.0 // indirect + github.com/cespare/xxhash/v2 v2.3.0 // indirect + github.com/cloudwego/base64x v0.1.6 // indirect + github.com/gabriel-vasile/mimetype v1.4.12 // indirect + github.com/gin-contrib/sse v1.1.0 // indirect + github.com/go-playground/locales v0.14.1 // indirect + github.com/go-playground/universal-translator v0.18.1 // indirect + github.com/go-playground/validator/v10 v10.30.1 // indirect + github.com/goccy/go-json v0.10.5 // indirect + github.com/goccy/go-yaml v1.19.2 // indirect + github.com/jackc/pgpassfile v1.0.0 // indirect + github.com/jackc/pgservicefile v0.0.0-20240606120523-5a60cdf6a761 // indirect + github.com/jackc/pgx/v5 v5.6.0 // indirect + github.com/jackc/puddle/v2 v2.2.2 // indirect + github.com/jinzhu/inflection v1.0.0 // indirect + github.com/jinzhu/now v1.1.5 // indirect + github.com/json-iterator/go v1.1.12 // indirect + github.com/klauspost/compress v1.17.6 // indirect + github.com/klauspost/cpuid/v2 v2.3.0 // indirect + github.com/leodido/go-urn v1.4.0 // indirect + github.com/mattn/go-isatty v0.0.20 // indirect + github.com/modern-go/concurrent v0.0.0-20180306012644-bacd9c7ef1dd // indirect + github.com/modern-go/reflect2 v1.0.2 // indirect + github.com/pelletier/go-toml/v2 v2.2.4 // indirect + github.com/pierrec/lz4/v4 v4.1.15 // indirect + github.com/quic-go/qpack v0.6.0 // indirect + github.com/quic-go/quic-go v0.59.0 // indirect + github.com/twitchyliquid64/golang-asm v0.15.1 // indirect + github.com/ugorji/go/codec v1.3.1 // indirect + go.mongodb.org/mongo-driver/v2 v2.5.0 // indirect + go.uber.org/atomic v1.11.0 // indirect + golang.org/x/arch v0.22.0 // indirect + golang.org/x/crypto v0.48.0 // indirect + golang.org/x/net v0.51.0 // indirect + golang.org/x/sync v0.19.0 // indirect + golang.org/x/sys v0.41.0 // indirect + golang.org/x/text v0.34.0 // indirect + google.golang.org/protobuf v1.36.10 // indirect +) diff --git a/services/payment-service/go.sum b/services/payment-service/go.sum new file mode 100644 index 0000000..5e4a798 --- /dev/null +++ b/services/payment-service/go.sum @@ -0,0 +1,142 @@ +github.com/bsm/ginkgo/v2 v2.12.0 h1:Ny8MWAHyOepLGlLKYmXG4IEkioBysk6GpaRTLC8zwWs= +github.com/bsm/ginkgo/v2 v2.12.0/go.mod h1:SwYbGRRDovPVboqFv0tPTcG1sN61LM1Z4ARdbAV9g4c= +github.com/bsm/gomega v1.27.10 h1:yeMWxP2pV2fG3FgAODIY8EiRE3dy0aeFYt4l7wh6yKA= +github.com/bsm/gomega v1.27.10/go.mod h1:JyEr/xRbxbtgWNi8tIEVPUYZ5Dzef52k01W3YH0H+O0= +github.com/bytedance/gopkg v0.1.3 h1:TPBSwH8RsouGCBcMBktLt1AymVo2TVsBVCY4b6TnZ/M= +github.com/bytedance/gopkg v0.1.3/go.mod h1:576VvJ+eJgyCzdjS+c4+77QF3p7ubbtiKARP3TxducM= +github.com/bytedance/sonic v1.15.0 h1:/PXeWFaR5ElNcVE84U0dOHjiMHQOwNIx3K4ymzh/uSE= +github.com/bytedance/sonic v1.15.0/go.mod h1:tFkWrPz0/CUCLEF4ri4UkHekCIcdnkqXw9VduqpJh0k= +github.com/bytedance/sonic/loader v0.5.0 h1:gXH3KVnatgY7loH5/TkeVyXPfESoqSBSBEiDd5VjlgE= +github.com/bytedance/sonic/loader v0.5.0/go.mod h1:AR4NYCk5DdzZizZ5djGqQ92eEhCCcdf5x77udYiSJRo= +github.com/cespare/xxhash/v2 v2.3.0 h1:UL815xU9SqsFlibzuggzjXhog7bL6oX9BbNZnL2UFvs= +github.com/cespare/xxhash/v2 v2.3.0/go.mod h1:VGX0DQ3Q6kWi7AoAeZDth3/j3BFtOZR5XLFGgcrjCOs= +github.com/cloudwego/base64x v0.1.6 h1:t11wG9AECkCDk5fMSoxmufanudBtJ+/HemLstXDLI2M= +github.com/cloudwego/base64x v0.1.6/go.mod h1:OFcloc187FXDaYHvrNIjxSe8ncn0OOM8gEHfghB2IPU= +github.com/davecgh/go-spew v1.1.0/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= +github.com/davecgh/go-spew v1.1.1 h1:vj9j/u1bqnvCEfJOwUhtlOARqs3+rkHYY13jYWTU97c= +github.com/davecgh/go-spew v1.1.1/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= +github.com/gabriel-vasile/mimetype v1.4.12 h1:e9hWvmLYvtp846tLHam2o++qitpguFiYCKbn0w9jyqw= +github.com/gabriel-vasile/mimetype v1.4.12/go.mod h1:d+9Oxyo1wTzWdyVUPMmXFvp4F9tea18J8ufA774AB3s= +github.com/gin-contrib/sse v1.1.0 h1:n0w2GMuUpWDVp7qSpvze6fAu9iRxJY4Hmj6AmBOU05w= +github.com/gin-contrib/sse v1.1.0/go.mod h1:hxRZ5gVpWMT7Z0B0gSNYqqsSCNIJMjzvm6fqCz9vjwM= +github.com/gin-gonic/gin v1.12.0 h1:b3YAbrZtnf8N//yjKeU2+MQsh2mY5htkZidOM7O0wG8= +github.com/gin-gonic/gin v1.12.0/go.mod h1:VxccKfsSllpKshkBWgVgRniFFAzFb9csfngsqANjnLc= +github.com/go-playground/assert/v2 v2.2.0 h1:JvknZsQTYeFEAhQwI4qEt9cyV5ONwRHC+lYKSsYSR8s= +github.com/go-playground/assert/v2 v2.2.0/go.mod h1:VDjEfimB/XKnb+ZQfWdccd7VUvScMdVu0Titje2rxJ4= +github.com/go-playground/locales v0.14.1 h1:EWaQ/wswjilfKLTECiXz7Rh+3BjFhfDFKv/oXslEjJA= +github.com/go-playground/locales v0.14.1/go.mod h1:hxrqLVvrK65+Rwrd5Fc6F2O76J/NuW9t0sjnWqG1slY= +github.com/go-playground/universal-translator v0.18.1 h1:Bcnm0ZwsGyWbCzImXv+pAJnYK9S473LQFuzCbDbfSFY= +github.com/go-playground/universal-translator v0.18.1/go.mod h1:xekY+UJKNuX9WP91TpwSH2VMlDf28Uj24BCp08ZFTUY= +github.com/go-playground/validator/v10 v10.30.1 h1:f3zDSN/zOma+w6+1Wswgd9fLkdwy06ntQJp0BBvFG0w= +github.com/go-playground/validator/v10 v10.30.1/go.mod h1:oSuBIQzuJxL//3MelwSLD5hc2Tu889bF0Idm9Dg26cM= +github.com/goccy/go-json v0.10.5 h1:Fq85nIqj+gXn/S5ahsiTlK3TmC85qgirsdTP/+DeaC4= +github.com/goccy/go-json v0.10.5/go.mod h1:oq7eo15ShAhp70Anwd5lgX2pLfOS3QCiwU/PULtXL6M= +github.com/goccy/go-yaml v1.19.2 h1:PmFC1S6h8ljIz6gMRBopkjP1TVT7xuwrButHID66PoM= +github.com/goccy/go-yaml v1.19.2/go.mod h1:XBurs7gK8ATbW4ZPGKgcbrY1Br56PdM69F7LkFRi1kA= +github.com/google/go-cmp v0.7.0 h1:wk8382ETsv4JYUZwIsn6YpYiWiBsYLSJiTsyBybVuN8= +github.com/google/go-cmp v0.7.0/go.mod h1:pXiqmnSA92OHEEa9HXL2W4E7lf9JzCmGVUdgjX3N/iU= +github.com/google/gofuzz v1.0.0/go.mod h1:dBl0BpW6vV/+mYPU4Po3pmUjxk6FQPldtuIdl/M65Eg= +github.com/google/uuid v1.6.0 h1:NIvaJDMOsjHA8n1jAhLSgzrAzy1Hgr+hNrb57e+94F0= +github.com/google/uuid v1.6.0/go.mod h1:TIyPZe4MgqvfeYDBFedMoGGpEw/LqOeaOT+nhxU+yHo= +github.com/jackc/pgpassfile v1.0.0 h1:/6Hmqy13Ss2zCq62VdNG8tM1wchn8zjSGOBJ6icpsIM= +github.com/jackc/pgpassfile v1.0.0/go.mod h1:CEx0iS5ambNFdcRtxPj5JhEz+xB6uRky5eyVu/W2HEg= +github.com/jackc/pgservicefile v0.0.0-20240606120523-5a60cdf6a761 h1:iCEnooe7UlwOQYpKFhBabPMi4aNAfoODPEFNiAnClxo= +github.com/jackc/pgservicefile v0.0.0-20240606120523-5a60cdf6a761/go.mod h1:5TJZWKEWniPve33vlWYSoGYefn3gLQRzjfDlhSJ9ZKM= +github.com/jackc/pgx/v5 v5.6.0 h1:SWJzexBzPL5jb0GEsrPMLIsi/3jOo7RHlzTjcAeDrPY= +github.com/jackc/pgx/v5 v5.6.0/go.mod h1:DNZ/vlrUnhWCoFGxHAG8U2ljioxukquj7utPDgtQdTw= +github.com/jackc/puddle/v2 v2.2.2 h1:PR8nw+E/1w0GLuRFSmiioY6UooMp6KJv0/61nB7icHo= +github.com/jackc/puddle/v2 v2.2.2/go.mod h1:vriiEXHvEE654aYKXXjOvZM39qJ0q+azkZFrfEOc3H4= +github.com/jinzhu/inflection v1.0.0 h1:K317FqzuhWc8YvSVlFMCCUb36O/S9MCKRDI7QkRKD/E= +github.com/jinzhu/inflection v1.0.0/go.mod h1:h+uFLlag+Qp1Va5pdKtLDYj+kHp5pxUVkryuEj+Srlc= +github.com/jinzhu/now v1.1.5 h1:/o9tlHleP7gOFmsnYNz3RGnqzefHA47wQpKrrdTIwXQ= +github.com/jinzhu/now v1.1.5/go.mod h1:d3SSVoowX0Lcu0IBviAWJpolVfI5UJVZZ7cO71lE/z8= +github.com/json-iterator/go v1.1.12 h1:PV8peI4a0ysnczrg+LtxykD8LfKY9ML6u2jnxaEnrnM= +github.com/json-iterator/go v1.1.12/go.mod h1:e30LSqwooZae/UwlEbR2852Gd8hjQvJoHmT4TnhNGBo= +github.com/klauspost/compress v1.17.6 h1:60eq2E/jlfwQXtvZEeBUYADs+BwKBWURIY+Gj2eRGjI= +github.com/klauspost/compress v1.17.6/go.mod h1:/dCuZOvVtNoHsyb+cuJD3itjs3NbnF6KH9zAO4BDxPM= +github.com/klauspost/cpuid/v2 v2.3.0 h1:S4CRMLnYUhGeDFDqkGriYKdfoFlDnMtqTiI/sFzhA9Y= +github.com/klauspost/cpuid/v2 v2.3.0/go.mod h1:hqwkgyIinND0mEev00jJYCxPNVRVXFQeu1XKlok6oO0= +github.com/leodido/go-urn v1.4.0 h1:WT9HwE9SGECu3lg4d/dIA+jxlljEa1/ffXKmRjqdmIQ= +github.com/leodido/go-urn v1.4.0/go.mod h1:bvxc+MVxLKB4z00jd1z+Dvzr47oO32F/QSNjSBOlFxI= +github.com/mattn/go-isatty v0.0.20 h1:xfD0iDuEKnDkl03q4limB+vH+GxLEtL/jb4xVJSWWEY= +github.com/mattn/go-isatty v0.0.20/go.mod h1:W+V8PltTTMOvKvAeJH7IuucS94S2C6jfK/D7dTCTo3Y= +github.com/modern-go/concurrent v0.0.0-20180228061459-e0a39a4cb421/go.mod h1:6dJC0mAP4ikYIbvyc7fijjWJddQyLn8Ig3JB5CqoB9Q= +github.com/modern-go/concurrent v0.0.0-20180306012644-bacd9c7ef1dd h1:TRLaZ9cD/w8PVh93nsPXa1VrQ6jlwL5oN8l14QlcNfg= +github.com/modern-go/concurrent v0.0.0-20180306012644-bacd9c7ef1dd/go.mod h1:6dJC0mAP4ikYIbvyc7fijjWJddQyLn8Ig3JB5CqoB9Q= +github.com/modern-go/reflect2 v1.0.2 h1:xBagoLtFs94CBntxluKeaWgTMpvLxC4ur3nMaC9Gz0M= +github.com/modern-go/reflect2 v1.0.2/go.mod h1:yWuevngMOJpCy52FWWMvUC8ws7m/LJsjYzDa0/r8luk= +github.com/pelletier/go-toml/v2 v2.2.4 h1:mye9XuhQ6gvn5h28+VilKrrPoQVanw5PMw/TB0t5Ec4= +github.com/pelletier/go-toml/v2 v2.2.4/go.mod h1:2gIqNv+qfxSVS7cM2xJQKtLSTLUE9V8t9Stt+h56mCY= +github.com/pierrec/lz4/v4 v4.1.15 h1:MO0/ucJhngq7299dKLwIMtgTfbkoSPF6AoMYDd8Q4q0= +github.com/pierrec/lz4/v4 v4.1.15/go.mod h1:gZWDp/Ze/IJXGXf23ltt2EXimqmTUXEy0GFuRQyBid4= +github.com/pmezard/go-difflib v1.0.0 h1:4DBwDE0NGyQoBHbLQYPwSUPoCMWR5BEzIk/f1lZbAQM= +github.com/pmezard/go-difflib v1.0.0/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4= +github.com/quic-go/qpack v0.6.0 h1:g7W+BMYynC1LbYLSqRt8PBg5Tgwxn214ZZR34VIOjz8= +github.com/quic-go/qpack v0.6.0/go.mod h1:lUpLKChi8njB4ty2bFLX2x4gzDqXwUpaO1DP9qMDZII= +github.com/quic-go/quic-go v0.59.0 h1:OLJkp1Mlm/aS7dpKgTc6cnpynnD2Xg7C1pwL6vy/SAw= +github.com/quic-go/quic-go v0.59.0/go.mod h1:upnsH4Ju1YkqpLXC305eW3yDZ4NfnNbmQRCMWS58IKU= +github.com/redis/go-redis/v9 v9.19.0 h1:XPVaaPSnG6RhYf7p+rmSa9zZfeVAnWsH5h3lxthOm/k= +github.com/redis/go-redis/v9 v9.19.0/go.mod h1:v/M13XI1PVCDcm01VtPFOADfZtHf8YW3baQf57KlIkA= +github.com/segmentio/kafka-go v0.4.51 h1:JgDPPG75tC1rWIS2Me6MwcvXJ6f49UQ4HjAOef71Hno= +github.com/segmentio/kafka-go v0.4.51/go.mod h1:Y1gn60kzLEEaW28YshXyk2+VCUKbJ3Qr6DrnT3i4+9E= +github.com/stretchr/objx v0.1.0/go.mod h1:HFkY916IF+rwdDfMAkV7OtwuqBVzrE8GR6GFx+wExME= +github.com/stretchr/objx v0.4.0/go.mod h1:YvHI0jy2hoMjB+UWwv71VJQ9isScKT/TqJzVSSt89Yw= +github.com/stretchr/objx v0.5.0/go.mod h1:Yh+to48EsGEfYuaHDzXPcE3xhTkx73EhmCGUpEOglKo= +github.com/stretchr/objx v0.5.2/go.mod h1:FRsXN1f5AsAjCGJKqEizvkpNtU+EGNCLh3NxZ/8L+MA= +github.com/stretchr/testify v1.3.0/go.mod h1:M5WIy9Dh21IEIfnGCwXGc5bZfKNJtfHm1UVUgZn+9EI= +github.com/stretchr/testify v1.7.0/go.mod h1:6Fq8oRcR53rry900zMqJjRRixrwX3KX962/h/Wwjteg= +github.com/stretchr/testify v1.7.1/go.mod h1:6Fq8oRcR53rry900zMqJjRRixrwX3KX962/h/Wwjteg= +github.com/stretchr/testify v1.8.0/go.mod h1:yNjHg4UonilssWZ8iaSj1OCr/vHnekPRkoO+kdMU+MU= +github.com/stretchr/testify v1.8.4/go.mod h1:sz/lmYIOXD/1dqDmKjjqLyZ2RngseejIcXlSw2iwfAo= +github.com/stretchr/testify v1.10.0/go.mod h1:r2ic/lqez/lEtzL7wO/rwa5dbSLXVDPFyf8C91i36aY= +github.com/stretchr/testify v1.11.1 h1:7s2iGBzp5EwR7/aIZr8ao5+dra3wiQyKjjFuvgVKu7U= +github.com/stretchr/testify v1.11.1/go.mod h1:wZwfW3scLgRK+23gO65QZefKpKQRnfz6sD981Nm4B6U= +github.com/stripe/stripe-go/v76 v76.25.0 h1:kmDoOTvdQSTQssQzWZQQkgbAR2Q8eXdMWbN/ylNalWA= +github.com/stripe/stripe-go/v76 v76.25.0/go.mod h1:rw1MxjlAKKcZ+3FOXgTHgwiOa2ya6CPq6ykpJ0Q6Po4= +github.com/twitchyliquid64/golang-asm v0.15.1 h1:SU5vSMR7hnwNxj24w34ZyCi/FmDZTkS4MhqMhdFk5YI= +github.com/twitchyliquid64/golang-asm v0.15.1/go.mod h1:a1lVb/DtPvCB8fslRZhAngC2+aY1QWCk3Cedj/Gdt08= +github.com/ugorji/go/codec v1.3.1 h1:waO7eEiFDwidsBN6agj1vJQ4AG7lh2yqXyOXqhgQuyY= +github.com/ugorji/go/codec v1.3.1/go.mod h1:pRBVtBSKl77K30Bv8R2P+cLSGaTtex6fsA2Wjqmfxj4= +github.com/xdg-go/pbkdf2 v1.0.0 h1:Su7DPu48wXMwC3bs7MCNG+z4FhcyEuz5dlvchbq0B0c= +github.com/xdg-go/pbkdf2 v1.0.0/go.mod h1:jrpuAogTd400dnrH08LKmI/xc1MbPOebTwRqcT5RDeI= +github.com/xdg-go/scram v1.2.0 h1:bYKF2AEwG5rqd1BumT4gAnvwU/M9nBp2pTSxeZw7Wvs= +github.com/xdg-go/scram v1.2.0/go.mod h1:3dlrS0iBaWKYVt2ZfA4cj48umJZ+cAEbR6/SjLA88I8= +github.com/xdg-go/stringprep v1.0.4 h1:XLI/Ng3O1Atzq0oBs3TWm+5ZVgkq2aqdlvP9JtoZ6c8= +github.com/xdg-go/stringprep v1.0.4/go.mod h1:mPGuuIYwz7CmR2bT9j4GbQqutWS1zV24gijq1dTyGkM= +github.com/zeebo/xxh3 v1.1.0 h1:s7DLGDK45Dyfg7++yxI0khrfwq9661w9EN78eP/UZVs= +github.com/zeebo/xxh3 v1.1.0/go.mod h1:IisAie1LELR4xhVinxWS5+zf1lA4p0MW4T+w+W07F5s= +go.mongodb.org/mongo-driver/v2 v2.5.0 h1:yXUhImUjjAInNcpTcAlPHiT7bIXhshCTL3jVBkF3xaE= +go.mongodb.org/mongo-driver/v2 v2.5.0/go.mod h1:yOI9kBsufol30iFsl1slpdq1I0eHPzybRWdyYUs8K/0= +go.uber.org/atomic v1.11.0 h1:ZvwS0R+56ePWxUNi+Atn9dWONBPp/AUETXlHW0DxSjE= +go.uber.org/atomic v1.11.0/go.mod h1:LUxbIzbOniOlMKjJjyPfpl4v+PKK2cNJn91OQbhoJI0= +go.uber.org/mock v0.6.0 h1:hyF9dfmbgIX5EfOdasqLsWD6xqpNZlXblLB/Dbnwv3Y= +go.uber.org/mock v0.6.0/go.mod h1:KiVJ4BqZJaMj4svdfmHM0AUx4NJYO8ZNpPnZn1Z+BBU= +golang.org/x/arch v0.22.0 h1:c/Zle32i5ttqRXjdLyyHZESLD/bB90DCU1g9l/0YBDI= +golang.org/x/arch v0.22.0/go.mod h1:dNHoOeKiyja7GTvF9NJS1l3Z2yntpQNzgrjh1cU103A= +golang.org/x/crypto v0.48.0 h1:/VRzVqiRSggnhY7gNRxPauEQ5Drw9haKdM0jqfcCFts= +golang.org/x/crypto v0.48.0/go.mod h1:r0kV5h3qnFPlQnBSrULhlsRfryS2pmewsg+XfMgkVos= +golang.org/x/net v0.0.0-20210520170846-37e1c6afe023/go.mod h1:9nx3DQGgdP8bBQD5qxJ1jj9UTztislL4KSBs9R2vV5Y= +golang.org/x/net v0.51.0 h1:94R/GTO7mt3/4wIKpcR5gkGmRLOuE/2hNGeWq/GBIFo= +golang.org/x/net v0.51.0/go.mod h1:aamm+2QF5ogm02fjy5Bb7CQ0WMt1/WVM7FtyaTLlA9Y= +golang.org/x/sync v0.19.0 h1:vV+1eWNmZ5geRlYjzm2adRgW2/mcpevXNg50YZtPCE4= +golang.org/x/sync v0.19.0/go.mod h1:9KTHXmSnoGruLpwFjVSX0lNNA75CykiMECbovNTZqGI= +golang.org/x/sys v0.0.0-20201119102817-f84b799fce68/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= +golang.org/x/sys v0.0.0-20210423082822-04245dca01da/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= +golang.org/x/sys v0.6.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= +golang.org/x/sys v0.41.0 h1:Ivj+2Cp/ylzLiEU89QhWblYnOE9zerudt9Ftecq2C6k= +golang.org/x/sys v0.41.0/go.mod h1:OgkHotnGiDImocRcuBABYBEXf8A9a87e/uXjp9XT3ks= +golang.org/x/term v0.0.0-20201126162022-7de9c90e9dd1/go.mod h1:bj7SfCRtBDWHUb9snDiAeCFNEtKQo2Wmx5Cou7ajbmo= +golang.org/x/text v0.3.6/go.mod h1:5Zoc/QRtKVWzQhOtBMvqHzDpF6irO9z98xDceosuGiQ= +golang.org/x/text v0.34.0 h1:oL/Qq0Kdaqxa1KbNeMKwQq0reLCCaFtqu2eNuSeNHbk= +golang.org/x/text v0.34.0/go.mod h1:homfLqTYRFyVYemLBFl5GgL/DWEiH5wcsQ5gSh1yziA= +golang.org/x/tools v0.0.0-20180917221912-90fa682c2a6e/go.mod h1:n7NCudcB/nEzxVGmLbDWY5pfWTLqBcC2KZ6jyYvM4mQ= +google.golang.org/protobuf v1.36.10 h1:AYd7cD/uASjIL6Q9LiTjz8JLcrh/88q5UObnmY3aOOE= +google.golang.org/protobuf v1.36.10/go.mod h1:HTf+CrKn2C3g5S8VImy6tdcUvCska2kB7j23XfzDpco= +gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0= +gopkg.in/yaml.v3 v3.0.0-20200313102051-9f266ea9e77c/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= +gopkg.in/yaml.v3 v3.0.1 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA= +gopkg.in/yaml.v3 v3.0.1/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= +gorm.io/driver/postgres v1.6.0 h1:2dxzU8xJ+ivvqTRph34QX+WrRaJlmfyPqXmoGVjMBa4= +gorm.io/driver/postgres v1.6.0/go.mod h1:vUw0mrGgrTK+uPHEhAdV4sfFELrByKVGnaVRkXDhtWo= +gorm.io/gorm v1.31.1 h1:7CA8FTFz/gRfgqgpeKIBcervUn3xSyPUmr6B2WXJ7kg= +gorm.io/gorm v1.31.1/go.mod h1:XyQVbO2k6YkOis7C2437jSit3SsDK72s7n7rsSHd+Gs= diff --git a/services/payment-service/internal/cache/payment_cache.go b/services/payment-service/internal/cache/payment_cache.go new file mode 100644 index 0000000..51eed88 --- /dev/null +++ b/services/payment-service/internal/cache/payment_cache.go @@ -0,0 +1,55 @@ +package cache + +import ( + "context" + "encoding/json" + "time" + + "auron/payment-service/internal/domain" + + "github.com/google/uuid" + "github.com/redis/go-redis/v9" +) + +const ( + paymentPrefix = "payment:" + paymentTTL = time.Hour +) + +type PaymentCache struct { + redis *redis.Client +} + +func NewPaymentCache(redisClient *redis.Client) domain.PaymentCache { + return &PaymentCache{redis: redisClient} +} + +func (c *PaymentCache) GetPayment(ctx context.Context, paymentID uuid.UUID) (*domain.Payment, error) { + key := paymentPrefix + paymentID.String() + cached, err := c.redis.Get(ctx, key).Result() + if err != nil { + if err == redis.Nil { + return nil, nil + } + return nil, err + } + + var payment domain.Payment + if err := json.Unmarshal([]byte(cached), &payment); err != nil { + return nil, err + } + return &payment, nil +} + +func (c *PaymentCache) SetPayment(ctx context.Context, payment *domain.Payment) error { + key := paymentPrefix + payment.ID.String() + data, err := json.Marshal(payment) + if err != nil { + return err + } + return c.redis.Set(ctx, key, data, paymentTTL).Err() +} + +func (c *PaymentCache) InvalidatePayment(ctx context.Context, paymentID uuid.UUID) error { + return c.redis.Del(ctx, paymentPrefix+paymentID.String()).Err() +} diff --git a/services/payment-service/internal/client/stripe_client.go b/services/payment-service/internal/client/stripe_client.go new file mode 100644 index 0000000..e1af983 --- /dev/null +++ b/services/payment-service/internal/client/stripe_client.go @@ -0,0 +1,38 @@ +package client + +import ( + "context" + "fmt" + + "auron/payment-service/internal/domain" + + stripe "github.com/stripe/stripe-go/v76" + "github.com/stripe/stripe-go/v76/paymentintent" +) + +type stripeClient struct { + secretKey string +} + +func NewStripeClient(secretKey string) domain.StripeClient { + stripe.Key = secretKey + return &stripeClient{secretKey: secretKey} +} + +func (c *stripeClient) CreatePaymentIntent(_ context.Context, amount float64, currency string, metadata map[string]string) (string, string, error) { + params := &stripe.PaymentIntentParams{ + Amount: stripe.Int64(int64(amount * 100)), + Currency: stripe.String(currency), + AutomaticPaymentMethods: &stripe.PaymentIntentAutomaticPaymentMethodsParams{ + Enabled: stripe.Bool(true), + }, + Metadata: metadata, + } + + pi, err := paymentintent.New(params) + if err != nil { + return "", "", fmt.Errorf("stripe: create payment intent: %w", err) + } + + return pi.ID, pi.ClientSecret, nil +} diff --git a/services/payment-service/internal/domain/cache.go b/services/payment-service/internal/domain/cache.go new file mode 100644 index 0000000..7b62708 --- /dev/null +++ b/services/payment-service/internal/domain/cache.go @@ -0,0 +1,13 @@ +package domain + +import ( + "context" + + "github.com/google/uuid" +) + +type PaymentCache interface { + GetPayment(ctx context.Context, paymentID uuid.UUID) (*Payment, error) + SetPayment(ctx context.Context, payment *Payment) error + InvalidatePayment(ctx context.Context, paymentID uuid.UUID) error +} diff --git a/services/payment-service/internal/domain/client.go b/services/payment-service/internal/domain/client.go new file mode 100644 index 0000000..11382a5 --- /dev/null +++ b/services/payment-service/internal/domain/client.go @@ -0,0 +1,7 @@ +package domain + +import "context" + +type StripeClient interface { + CreatePaymentIntent(ctx context.Context, amount float64, currency string, metadata map[string]string) (intentID, clientSecret string, err error) +} diff --git a/services/payment-service/internal/domain/errors.go b/services/payment-service/internal/domain/errors.go new file mode 100644 index 0000000..2df47ca --- /dev/null +++ b/services/payment-service/internal/domain/errors.go @@ -0,0 +1,11 @@ +package domain + +import "errors" + +var ( + ErrPaymentNotFound = errors.New("payment not found") + ErrPaymentAlreadyExists = errors.New("payment already exists for this order") + ErrInvalidWebhookSignature = errors.New("invalid webhook signature") + ErrUnauthorized = errors.New("unauthorized") + ErrForbidden = errors.New("forbidden") +) diff --git a/services/payment-service/internal/domain/events.go b/services/payment-service/internal/domain/events.go new file mode 100644 index 0000000..7eb3f16 --- /dev/null +++ b/services/payment-service/internal/domain/events.go @@ -0,0 +1,18 @@ +package domain + +import "context" + +type EventPublisher interface { + Publish(ctx context.Context, topic string, payload any) error + Close() error +} + +const ( + // Consumed topics + TopicOrderCreated = "order.created" + + // Published topics + TopicPaymentCreated = "payment.created" + TopicPaymentCompleted = "payment.completed" + TopicPaymentFailed = "payment.failed" +) diff --git a/services/payment-service/internal/domain/payment.go b/services/payment-service/internal/domain/payment.go new file mode 100644 index 0000000..5151f85 --- /dev/null +++ b/services/payment-service/internal/domain/payment.go @@ -0,0 +1,79 @@ +package domain + +import ( + "time" + + "github.com/google/uuid" +) + +type PaymentStatus string + +const ( + PaymentStatusPending PaymentStatus = "pending" + PaymentStatusProcessing PaymentStatus = "processing" + PaymentStatusCompleted PaymentStatus = "completed" + PaymentStatusFailed PaymentStatus = "failed" + PaymentStatusRefunded PaymentStatus = "refunded" +) + +type Payment struct { + ID uuid.UUID `json:"id" gorm:"type:uuid;default:gen_random_uuid();primaryKey"` + OrderID uuid.UUID `json:"order_id" gorm:"type:uuid;not null;uniqueIndex"` + UserID uuid.UUID `json:"user_id" gorm:"type:uuid;not null;index"` + Amount float64 `json:"amount" gorm:"type:decimal(12,2);not null"` + Currency string `json:"currency" gorm:"type:varchar(10);not null;default:'usd'"` + Status PaymentStatus `json:"status" gorm:"type:varchar(50);not null;default:'pending';index"` + StripePaymentIntentID string `json:"stripe_payment_intent_id,omitempty" gorm:"type:varchar(255)"` + StripeClientSecret string `json:"-" gorm:"type:text"` + FailureReason string `json:"failure_reason,omitempty" gorm:"type:text"` + CreatedAt time.Time `json:"created_at" gorm:"not null;default:now()"` + UpdatedAt time.Time `json:"updated_at" gorm:"not null;default:now()"` +} + +func (Payment) TableName() string { + return "payments" +} + +type PaymentResponse struct { + ID uuid.UUID `json:"id"` + OrderID uuid.UUID `json:"order_id"` + UserID uuid.UUID `json:"user_id"` + Amount float64 `json:"amount"` + Currency string `json:"currency"` + Status PaymentStatus `json:"status"` + StripePaymentIntentID string `json:"stripe_payment_intent_id,omitempty"` + FailureReason string `json:"failure_reason,omitempty"` + CreatedAt time.Time `json:"created_at"` + UpdatedAt time.Time `json:"updated_at"` +} + +func (p *Payment) ToResponse() *PaymentResponse { + return &PaymentResponse{ + ID: p.ID, + OrderID: p.OrderID, + UserID: p.UserID, + Amount: p.Amount, + Currency: p.Currency, + Status: p.Status, + StripePaymentIntentID: p.StripePaymentIntentID, + FailureReason: p.FailureReason, + CreatedAt: p.CreatedAt, + UpdatedAt: p.UpdatedAt, + } +} + +// OrderCreatedEvent is the shape of the Kafka message consumed from order-service. +// JSON tags must match the Order struct in order-service (ID is published as "id"). +type OrderCreatedEvent struct { + OrderID uuid.UUID `json:"id"` + UserID uuid.UUID `json:"user_id"` + TotalAmount float64 `json:"total_amount"` + Items []OrderEventItem `json:"items"` +} + +type OrderEventItem struct { + ProductID uuid.UUID `json:"product_id"` + ProductName string `json:"product_name"` + Price float64 `json:"price"` + Quantity int `json:"quantity"` +} diff --git a/services/payment-service/internal/domain/repository.go b/services/payment-service/internal/domain/repository.go new file mode 100644 index 0000000..e7a1269 --- /dev/null +++ b/services/payment-service/internal/domain/repository.go @@ -0,0 +1,11 @@ +package domain + +import "github.com/google/uuid" + +type PaymentRepository interface { + GetPaymentByID(id uuid.UUID) (*Payment, error) + GetPaymentByOrderID(orderID uuid.UUID) (*Payment, error) + CreatePayment(payment *Payment) (*Payment, error) + UpdatePaymentStatus(id uuid.UUID, status PaymentStatus, failureReason string) (*Payment, error) + UpdateStripeIDs(id uuid.UUID, intentID, clientSecret string) (*Payment, error) +} diff --git a/services/payment-service/internal/domain/service.go b/services/payment-service/internal/domain/service.go new file mode 100644 index 0000000..3c2317e --- /dev/null +++ b/services/payment-service/internal/domain/service.go @@ -0,0 +1,13 @@ +package domain + +import ( + "context" + + "github.com/google/uuid" +) + +type PaymentService interface { + GetPaymentByID(ctx context.Context, userID, paymentID uuid.UUID) (*PaymentResponse, error) + HandleOrderCreated(ctx context.Context, event OrderCreatedEvent) error + HandleStripeWebhook(ctx context.Context, payload []byte, signature string) error +} diff --git a/services/payment-service/internal/events/kafka_consumer.go b/services/payment-service/internal/events/kafka_consumer.go new file mode 100644 index 0000000..0614255 --- /dev/null +++ b/services/payment-service/internal/events/kafka_consumer.go @@ -0,0 +1,58 @@ +package events + +import ( + "context" + "encoding/json" + "log/slog" + + "auron/payment-service/internal/domain" + + "github.com/segmentio/kafka-go" +) + +type KafkaConsumer struct { + reader *kafka.Reader + service domain.PaymentService +} + +func NewKafkaConsumer(brokers []string, topic, groupID string, service domain.PaymentService) *KafkaConsumer { + reader := kafka.NewReader(kafka.ReaderConfig{ + Brokers: brokers, + Topic: topic, + GroupID: groupID, + MinBytes: 10e3, + MaxBytes: 10e6, + }) + return &KafkaConsumer{reader: reader, service: service} +} + +// Start launches the consumer loop in a background goroutine. +func (c *KafkaConsumer) Start(ctx context.Context) { + go func() { + for { + msg, err := c.reader.FetchMessage(ctx) + if err != nil { + if ctx.Err() != nil { + return + } + slog.Error("kafka consumer: fetch error", "error", err) + continue + } + + var event domain.OrderCreatedEvent + if err := json.Unmarshal(msg.Value, &event); err != nil { + slog.Error("kafka consumer: unmarshal error", "error", err, "offset", msg.Offset) + } else if err := c.service.HandleOrderCreated(ctx, event); err != nil { + slog.Error("kafka consumer: HandleOrderCreated failed", "order_id", event.OrderID, "error", err) + } + + if err := c.reader.CommitMessages(ctx, msg); err != nil { + slog.Warn("kafka consumer: commit failed", "error", err) + } + } + }() +} + +func (c *KafkaConsumer) Close() error { + return c.reader.Close() +} diff --git a/services/payment-service/internal/events/kafka_publisher.go b/services/payment-service/internal/events/kafka_publisher.go new file mode 100644 index 0000000..ebe9303 --- /dev/null +++ b/services/payment-service/internal/events/kafka_publisher.go @@ -0,0 +1,52 @@ +package events + +import ( + "context" + "encoding/json" + "fmt" + "log/slog" + + "auron/payment-service/internal/domain" + + "github.com/segmentio/kafka-go" +) + +type kafkaPublisher struct { + writers map[string]*kafka.Writer +} + +func NewKafkaPublisher(writers map[string]*kafka.Writer) domain.EventPublisher { + return &kafkaPublisher{writers: writers} +} + +func (p *kafkaPublisher) Publish(ctx context.Context, topic string, payload any) error { + writer, ok := p.writers[topic] + if !ok { + return fmt.Errorf("publisher: no writer registered for topic %q", topic) + } + + data, err := json.Marshal(payload) + if err != nil { + return fmt.Errorf("publisher: marshal payload: %w", err) + } + + if err := writer.WriteMessages(ctx, kafka.Message{Value: data}); err != nil { + return fmt.Errorf("publisher: write to topic %q: %w", topic, err) + } + + slog.Debug("event published", slog.String("topic", topic)) + return nil +} + +func (p *kafkaPublisher) Close() error { + var closeErr error + for _, writer := range p.writers { + if writer == nil { + continue + } + if err := writer.Close(); err != nil && closeErr == nil { + closeErr = err + } + } + return closeErr +} diff --git a/services/payment-service/internal/handler/payment_handler.go b/services/payment-service/internal/handler/payment_handler.go new file mode 100644 index 0000000..b79ba40 --- /dev/null +++ b/services/payment-service/internal/handler/payment_handler.go @@ -0,0 +1,95 @@ +package handler + +import ( + "errors" + "log/slog" + "net/http" + + "auron/payment-service/internal/domain" + "auron/payment-service/internal/middleware" + + "github.com/gin-gonic/gin" + "github.com/google/uuid" +) + +type PaymentHandler struct { + service domain.PaymentService +} + +func NewPaymentHandler(service domain.PaymentService) *PaymentHandler { + return &PaymentHandler{service: service} +} + +func (h *PaymentHandler) GetPaymentByID(c *gin.Context) { + userID, ok := getUserID(c) + if !ok { + c.JSON(http.StatusUnauthorized, gin.H{"success": false, "error": domain.ErrUnauthorized.Error()}) + return + } + + paymentID, err := uuid.Parse(c.Param("id")) + if err != nil { + c.JSON(http.StatusBadRequest, gin.H{"success": false, "error": "invalid payment id"}) + return + } + + payment, err := h.service.GetPaymentByID(c.Request.Context(), userID, paymentID) + if err != nil { + h.handleError(c, err) + return + } + + c.JSON(http.StatusOK, gin.H{"success": true, "data": payment}) +} + +func (h *PaymentHandler) HandleStripeWebhook(c *gin.Context) { + rawBody, exists := c.Get(middleware.RawBodyKey) + if !exists { + c.JSON(http.StatusBadRequest, gin.H{"success": false, "error": "missing request body"}) + return + } + + payload, ok := rawBody.([]byte) + if !ok { + c.JSON(http.StatusInternalServerError, gin.H{"success": false, "error": "internal server error"}) + return + } + + signature := c.GetHeader("Stripe-Signature") + + if err := h.service.HandleStripeWebhook(c.Request.Context(), payload, signature); err != nil { + slog.Error("stripe webhook processing failed", "error", err) + } + + // Always return 200 — Stripe retries on any non-2xx response. + c.JSON(http.StatusOK, gin.H{"received": true}) +} + +func getUserID(c *gin.Context) (uuid.UUID, bool) { + raw := c.GetHeader("X-User-ID") + if raw == "" { + return uuid.Nil, false + } + id, err := uuid.Parse(raw) + if err != nil { + return uuid.Nil, false + } + return id, true +} + +func (h *PaymentHandler) handleError(c *gin.Context, err error) { + switch { + case errors.Is(err, domain.ErrPaymentNotFound): + c.JSON(http.StatusNotFound, gin.H{"success": false, "error": err.Error()}) + case errors.Is(err, domain.ErrPaymentAlreadyExists): + c.JSON(http.StatusConflict, gin.H{"success": false, "error": err.Error()}) + case errors.Is(err, domain.ErrInvalidWebhookSignature): + c.JSON(http.StatusBadRequest, gin.H{"success": false, "error": err.Error()}) + case errors.Is(err, domain.ErrUnauthorized): + c.JSON(http.StatusUnauthorized, gin.H{"success": false, "error": err.Error()}) + case errors.Is(err, domain.ErrForbidden): + c.JSON(http.StatusForbidden, gin.H{"success": false, "error": err.Error()}) + default: + c.JSON(http.StatusInternalServerError, gin.H{"success": false, "error": "internal server error"}) + } +} diff --git a/services/payment-service/internal/middleware/stripe_webhook.go b/services/payment-service/internal/middleware/stripe_webhook.go new file mode 100644 index 0000000..6e48eaa --- /dev/null +++ b/services/payment-service/internal/middleware/stripe_webhook.go @@ -0,0 +1,26 @@ +package middleware + +import ( + "io" + "net/http" + + "github.com/gin-gonic/gin" +) + +const RawBodyKey = "rawBody" + +// CaptureRawBody reads and stores the raw request bytes before any binding. +// Required because Stripe signature verification needs the exact original bytes, +// and c.Request.Body is consumed after the first read. +func CaptureRawBody() gin.HandlerFunc { + return func(c *gin.Context) { + body, err := io.ReadAll(c.Request.Body) + if err != nil { + c.JSON(http.StatusBadRequest, gin.H{"success": false, "error": "failed to read request body"}) + c.Abort() + return + } + c.Set(RawBodyKey, body) + c.Next() + } +} diff --git a/services/payment-service/internal/repository/payment_repository.go b/services/payment-service/internal/repository/payment_repository.go new file mode 100644 index 0000000..a8531c2 --- /dev/null +++ b/services/payment-service/internal/repository/payment_repository.go @@ -0,0 +1,69 @@ +package repository + +import ( + "auron/payment-service/internal/domain" + + "github.com/google/uuid" + "gorm.io/gorm" +) + +type PaymentRepository struct { + db *gorm.DB +} + +func NewPaymentRepository(db *gorm.DB) domain.PaymentRepository { + return &PaymentRepository{db: db} +} + +func (r *PaymentRepository) GetPaymentByID(id uuid.UUID) (*domain.Payment, error) { + var payment domain.Payment + if err := r.db.First(&payment, "id = ?", id).Error; err != nil { + if err == gorm.ErrRecordNotFound { + return nil, domain.ErrPaymentNotFound + } + return nil, err + } + return &payment, nil +} + +func (r *PaymentRepository) GetPaymentByOrderID(orderID uuid.UUID) (*domain.Payment, error) { + var payment domain.Payment + if err := r.db.First(&payment, "order_id = ?", orderID).Error; err != nil { + if err == gorm.ErrRecordNotFound { + return nil, domain.ErrPaymentNotFound + } + return nil, err + } + return &payment, nil +} + +func (r *PaymentRepository) CreatePayment(payment *domain.Payment) (*domain.Payment, error) { + if err := r.db.Create(payment).Error; err != nil { + return nil, err + } + return payment, nil +} + +func (r *PaymentRepository) UpdatePaymentStatus(id uuid.UUID, status domain.PaymentStatus, failureReason string) (*domain.Payment, error) { + updates := map[string]any{ + "status": status, + } + if failureReason != "" { + updates["failure_reason"] = failureReason + } + if err := r.db.Model(&domain.Payment{}).Where("id = ?", id).Updates(updates).Error; err != nil { + return nil, err + } + return r.GetPaymentByID(id) +} + +func (r *PaymentRepository) UpdateStripeIDs(id uuid.UUID, intentID, clientSecret string) (*domain.Payment, error) { + updates := map[string]any{ + "stripe_payment_intent_id": intentID, + "stripe_client_secret": clientSecret, + } + if err := r.db.Model(&domain.Payment{}).Where("id = ?", id).Updates(updates).Error; err != nil { + return nil, err + } + return r.GetPaymentByID(id) +} diff --git a/services/payment-service/internal/route/payment_route.go b/services/payment-service/internal/route/payment_route.go new file mode 100644 index 0000000..50020a1 --- /dev/null +++ b/services/payment-service/internal/route/payment_route.go @@ -0,0 +1,14 @@ +package route + +import ( + "auron/payment-service/internal/handler" + "auron/payment-service/internal/middleware" + + "github.com/gin-gonic/gin" +) + +func RegisterPaymentRoutes(router *gin.Engine, paymentHandler *handler.PaymentHandler) { + api := router.Group("/") + api.GET("/payments/:id", paymentHandler.GetPaymentByID) + api.POST("/payments/webhook/stripe", middleware.CaptureRawBody(), paymentHandler.HandleStripeWebhook) +} diff --git a/services/payment-service/internal/service/payment_service.go b/services/payment-service/internal/service/payment_service.go new file mode 100644 index 0000000..6a31c21 --- /dev/null +++ b/services/payment-service/internal/service/payment_service.go @@ -0,0 +1,268 @@ +package service + +import ( + "context" + "encoding/json" + "errors" + "fmt" + "log/slog" + "time" + + "auron/payment-service/internal/domain" + + "github.com/google/uuid" + stripe "github.com/stripe/stripe-go/v76" + "github.com/stripe/stripe-go/v76/webhook" +) + +type PaymentService struct { + paymentRepo domain.PaymentRepository + paymentCache domain.PaymentCache + stripeClient domain.StripeClient + publisher domain.EventPublisher + webhookSecret string +} + +func NewPaymentService( + paymentRepo domain.PaymentRepository, + paymentCache domain.PaymentCache, + stripeClient domain.StripeClient, + publisher domain.EventPublisher, + webhookSecret string, +) domain.PaymentService { + return &PaymentService{ + paymentRepo: paymentRepo, + paymentCache: paymentCache, + stripeClient: stripeClient, + publisher: publisher, + webhookSecret: webhookSecret, + } +} + +func (s *PaymentService) GetPaymentByID(ctx context.Context, userID, paymentID uuid.UUID) (*domain.PaymentResponse, error) { + if cached, err := s.paymentCache.GetPayment(ctx, paymentID); err == nil && cached != nil { + if cached.UserID != userID { + return nil, domain.ErrForbidden + } + return cached.ToResponse(), nil + } + + payment, err := s.paymentRepo.GetPaymentByID(paymentID) + if err != nil { + return nil, err + } + + if payment.UserID != userID { + return nil, domain.ErrForbidden + } + + if err := s.paymentCache.SetPayment(ctx, payment); err != nil { + slog.Warn("failed to cache payment", "payment_id", paymentID, "error", err) + } + + return payment.ToResponse(), nil +} + +func (s *PaymentService) HandleOrderCreated(ctx context.Context, event domain.OrderCreatedEvent) error { + // Idempotency: skip if payment already exists for this order. + existing, err := s.paymentRepo.GetPaymentByOrderID(event.OrderID) + if err != nil && !errors.Is(err, domain.ErrPaymentNotFound) { + return err + } + if existing != nil { + slog.Info("payment already exists for order, skipping", "order_id", event.OrderID) + return nil + } + + now := time.Now() + payment := &domain.Payment{ + ID: uuid.New(), + OrderID: event.OrderID, + UserID: event.UserID, + Amount: event.TotalAmount, + Currency: "usd", + Status: domain.PaymentStatusPending, + CreatedAt: now, + UpdatedAt: now, + } + + created, err := s.paymentRepo.CreatePayment(payment) + if err != nil { + return err + } + + metadata := map[string]string{ + "payment_id": created.ID.String(), + "order_id": event.OrderID.String(), + "user_id": event.UserID.String(), + } + + intentID, clientSecret, err := s.stripeClient.CreatePaymentIntent(ctx, event.TotalAmount, "usd", metadata) + if err != nil { + slog.Error("failed to create stripe payment intent", "payment_id", created.ID, "error", err) + return err + } + + updated, err := s.paymentRepo.UpdateStripeIDs(created.ID, intentID, clientSecret) + if err != nil { + return err + } + + if err := s.paymentCache.SetPayment(ctx, updated); err != nil { + slog.Warn("failed to cache payment", "payment_id", updated.ID, "error", err) + } + + // Publish payment.created — client_secret travels only via this event, never HTTP. + go func() { + payload := map[string]any{ + "payment_id": updated.ID, + "order_id": updated.OrderID, + "user_id": updated.UserID, + "amount": updated.Amount, + "currency": updated.Currency, + "client_secret": clientSecret, + } + if err := s.publisher.Publish(context.Background(), domain.TopicPaymentCreated, payload); err != nil { + slog.Warn("failed to publish payment.created", "payment_id", updated.ID, "error", err) + } + }() + + return nil +} + +func (s *PaymentService) HandleStripeWebhook(ctx context.Context, payload []byte, signature string) error { + var event stripe.Event + var err error + + if s.webhookSecret == "" { + // Allow unsigned webhooks in development when no secret is configured. + if err = json.Unmarshal(payload, &event); err != nil { + return fmt.Errorf("webhook: unmarshal event: %w", err) + } + } else { + event, err = webhook.ConstructEvent(payload, signature, s.webhookSecret) + if err != nil { + return domain.ErrInvalidWebhookSignature + } + } + + switch event.Type { + case "payment_intent.succeeded": + return s.handlePaymentSucceeded(ctx, event) + case "payment_intent.payment_failed": + return s.handlePaymentFailed(ctx, event) + case "payment_intent.processing": + return s.handlePaymentProcessing(ctx, event) + default: + slog.Debug("unhandled stripe event type", "type", event.Type) + } + + return nil +} + +func (s *PaymentService) handlePaymentSucceeded(ctx context.Context, event stripe.Event) error { + pi, err := extractPaymentIntent(event) + if err != nil { + return err + } + + payment, err := s.resolvePaymentFromIntent(pi) + if err != nil { + return err + } + + updated, err := s.paymentRepo.UpdatePaymentStatus(payment.ID, domain.PaymentStatusCompleted, "") + if err != nil { + return err + } + + if err := s.paymentCache.SetPayment(ctx, updated); err != nil { + slog.Warn("failed to cache payment after succeeded", "payment_id", updated.ID, "error", err) + } + + go func() { + if err := s.publisher.Publish(context.Background(), domain.TopicPaymentCompleted, updated.ToResponse()); err != nil { + slog.Warn("failed to publish payment.completed", "payment_id", updated.ID, "error", err) + } + }() + + return nil +} + +func (s *PaymentService) handlePaymentFailed(ctx context.Context, event stripe.Event) error { + pi, err := extractPaymentIntent(event) + if err != nil { + return err + } + + reason := "" + if pi.LastPaymentError != nil { + reason = pi.LastPaymentError.Msg + } + + payment, err := s.resolvePaymentFromIntent(pi) + if err != nil { + return err + } + + updated, err := s.paymentRepo.UpdatePaymentStatus(payment.ID, domain.PaymentStatusFailed, reason) + if err != nil { + return err + } + + if err := s.paymentCache.SetPayment(ctx, updated); err != nil { + slog.Warn("failed to cache payment after failed", "payment_id", updated.ID, "error", err) + } + + go func() { + if err := s.publisher.Publish(context.Background(), domain.TopicPaymentFailed, updated.ToResponse()); err != nil { + slog.Warn("failed to publish payment.failed", "payment_id", updated.ID, "error", err) + } + }() + + return nil +} + +func (s *PaymentService) handlePaymentProcessing(ctx context.Context, event stripe.Event) error { + pi, err := extractPaymentIntent(event) + if err != nil { + return err + } + + payment, err := s.resolvePaymentFromIntent(pi) + if err != nil { + return err + } + + updated, err := s.paymentRepo.UpdatePaymentStatus(payment.ID, domain.PaymentStatusProcessing, "") + if err != nil { + return err + } + + if err := s.paymentCache.SetPayment(ctx, updated); err != nil { + slog.Warn("failed to cache payment after processing", "payment_id", updated.ID, "error", err) + } + + return nil +} + +// resolvePaymentFromIntent reads payment_id from the Stripe metadata to look up the payment. +func (s *PaymentService) resolvePaymentFromIntent(pi stripe.PaymentIntent) (*domain.Payment, error) { + paymentIDStr, ok := pi.Metadata["payment_id"] + if !ok || paymentIDStr == "" { + return nil, fmt.Errorf("webhook: payment_id missing from stripe metadata for intent %s", pi.ID) + } + paymentID, err := uuid.Parse(paymentIDStr) + if err != nil { + return nil, fmt.Errorf("webhook: invalid payment_id in stripe metadata: %w", err) + } + return s.paymentRepo.GetPaymentByID(paymentID) +} + +func extractPaymentIntent(event stripe.Event) (stripe.PaymentIntent, error) { + var pi stripe.PaymentIntent + if err := json.Unmarshal(event.Data.Raw, &pi); err != nil { + return pi, fmt.Errorf("webhook: unmarshal payment intent: %w", err) + } + return pi, nil +} diff --git a/services/payment-service/main.go b/services/payment-service/main.go new file mode 100644 index 0000000..2b17c19 --- /dev/null +++ b/services/payment-service/main.go @@ -0,0 +1,7 @@ +package main + +import "auron/payment-service/cmd" + +func main() { + cmd.Run() +}