Stripe integration
This commit is contained in:
@@ -32,6 +32,11 @@ type Config struct {
|
||||
ResendAPIKey string
|
||||
EmailFrom string
|
||||
DigestIntervalHours int
|
||||
StripeSecretKey string
|
||||
StripeWebhookSecret string
|
||||
StripePriceID string
|
||||
StripePublishableKey string
|
||||
FrontendURL string
|
||||
}
|
||||
|
||||
// Load reads configuration from environment variables and returns a Config.
|
||||
@@ -51,7 +56,12 @@ func Load() (*Config, error) {
|
||||
NodeEnv: getEnv("NODE_ENV", "development"),
|
||||
ResendAPIKey: os.Getenv("RESEND_API_KEY"),
|
||||
EmailFrom: getEnv("EMAIL_FROM", "Koin Ping <alerts@koinping.com>"),
|
||||
DigestIntervalHours: getEnvInt("DIGEST_INTERVAL_HOURS", defaultDigestIntervalHours),
|
||||
DigestIntervalHours: getEnvInt("DIGEST_INTERVAL_HOURS", defaultDigestIntervalHours),
|
||||
StripeSecretKey: os.Getenv("STRIPE_SECRET_KEY"),
|
||||
StripeWebhookSecret: os.Getenv("STRIPE_WEBHOOK_SECRET"),
|
||||
StripePriceID: os.Getenv("STRIPE_PRICE_ID"),
|
||||
StripePublishableKey: os.Getenv("STRIPE_PUBLISHABLE_KEY"),
|
||||
FrontendURL: getEnv("FRONTEND_URL", "http://localhost:3000"),
|
||||
}
|
||||
|
||||
if cfg.PollIntervalMS < minPollIntervalMS {
|
||||
|
||||
@@ -3,12 +3,16 @@ package domain
|
||||
import "time"
|
||||
|
||||
type User struct {
|
||||
ID string `json:"id"`
|
||||
FirebaseUID string `json:"-"`
|
||||
Email string `json:"email"`
|
||||
DisplayName *string `json:"display_name"` //nolint:tagliatelle
|
||||
CreatedAt time.Time `json:"created_at"` //nolint:tagliatelle
|
||||
UpdatedAt time.Time `json:"updated_at"` //nolint:tagliatelle
|
||||
ID string `json:"id"`
|
||||
FirebaseUID string `json:"-"`
|
||||
Email string `json:"email"`
|
||||
DisplayName *string `json:"display_name"` //nolint:tagliatelle
|
||||
StripeCustomerID *string `json:"-"`
|
||||
StripeSubscriptionID *string `json:"-"`
|
||||
SubscriptionStatus string `json:"subscription_status"` //nolint:tagliatelle
|
||||
SubscriptionCreatedAt *time.Time `json:"subscription_created_at,omitempty"` //nolint:tagliatelle
|
||||
CreatedAt time.Time `json:"created_at"` //nolint:tagliatelle
|
||||
UpdatedAt time.Time `json:"updated_at"` //nolint:tagliatelle
|
||||
}
|
||||
|
||||
type Address struct {
|
||||
|
||||
199
backend-go/internal/handlers/stripe.go
Normal file
199
backend-go/internal/handlers/stripe.go
Normal file
@@ -0,0 +1,199 @@
|
||||
package handlers
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"io"
|
||||
"log"
|
||||
"net/http"
|
||||
|
||||
"github.com/stripe/stripe-go/v82"
|
||||
checkoutsession "github.com/stripe/stripe-go/v82/checkout/session"
|
||||
"github.com/stripe/stripe-go/v82/webhook"
|
||||
|
||||
"github.com/kjannette/koin-ping/backend-go/internal/config"
|
||||
"github.com/kjannette/koin-ping/backend-go/internal/middleware"
|
||||
"github.com/kjannette/koin-ping/backend-go/internal/models"
|
||||
)
|
||||
|
||||
const webhookMaxBodyBytes = 65536
|
||||
|
||||
type StripeHandler struct {
|
||||
users *models.UserModel
|
||||
cfg *config.Config
|
||||
}
|
||||
|
||||
func NewStripeHandler(users *models.UserModel, cfg *config.Config) *StripeHandler {
|
||||
stripe.Key = cfg.StripeSecretKey
|
||||
return &StripeHandler{users: users, cfg: cfg}
|
||||
}
|
||||
|
||||
// CreateCheckoutSession creates a Stripe Checkout session for the monthly subscription.
|
||||
func (h *StripeHandler) CreateCheckoutSession(w http.ResponseWriter, r *http.Request) {
|
||||
userID := middleware.GetUserID(r.Context())
|
||||
|
||||
user, err := h.users.GetByID(r.Context(), userID)
|
||||
if err != nil || user == nil {
|
||||
log.Printf("Failed to get user %s: %v", userID, err)
|
||||
writeError(w, http.StatusInternalServerError, "INTERNAL_ERROR", "Failed to load user")
|
||||
return
|
||||
}
|
||||
|
||||
params := &stripe.CheckoutSessionParams{
|
||||
Mode: stripe.String(string(stripe.CheckoutSessionModeSubscription)),
|
||||
LineItems: []*stripe.CheckoutSessionLineItemParams{
|
||||
{
|
||||
Price: stripe.String(h.cfg.StripePriceID),
|
||||
Quantity: stripe.Int64(1),
|
||||
},
|
||||
},
|
||||
SuccessURL: stripe.String(h.cfg.FrontendURL + "/onboarding?payment=success&session_id={CHECKOUT_SESSION_ID}"),
|
||||
CancelURL: stripe.String(h.cfg.FrontendURL + "/onboarding?payment=cancelled"),
|
||||
ClientReferenceID: stripe.String(userID),
|
||||
CustomerEmail: stripe.String(user.Email),
|
||||
}
|
||||
|
||||
if user.StripeCustomerID != nil && *user.StripeCustomerID != "" {
|
||||
params.Customer = user.StripeCustomerID
|
||||
params.CustomerEmail = nil
|
||||
}
|
||||
|
||||
s, err := checkoutsession.New(params)
|
||||
if err != nil {
|
||||
log.Printf("Failed to create Stripe checkout session: %v", err)
|
||||
writeError(w, http.StatusInternalServerError, "STRIPE_ERROR", "Failed to create checkout session")
|
||||
return
|
||||
}
|
||||
|
||||
writeJSON(w, http.StatusOK, map[string]string{"url": s.URL})
|
||||
}
|
||||
|
||||
// GetSubscriptionStatus returns the current user's subscription state.
|
||||
func (h *StripeHandler) GetSubscriptionStatus(w http.ResponseWriter, r *http.Request) {
|
||||
userID := middleware.GetUserID(r.Context())
|
||||
|
||||
user, err := h.users.GetByID(r.Context(), userID)
|
||||
if err != nil || user == nil {
|
||||
log.Printf("Failed to get user %s: %v", userID, err)
|
||||
writeError(w, http.StatusInternalServerError, "INTERNAL_ERROR", "Failed to load user")
|
||||
return
|
||||
}
|
||||
|
||||
writeJSON(w, http.StatusOK, map[string]any{
|
||||
"subscription_status": user.SubscriptionStatus,
|
||||
"subscription_created_at": user.SubscriptionCreatedAt,
|
||||
})
|
||||
}
|
||||
|
||||
// HandleWebhook processes incoming Stripe webhook events.
|
||||
// This endpoint must NOT require authentication (Stripe calls it directly).
|
||||
func (h *StripeHandler) HandleWebhook(w http.ResponseWriter, r *http.Request) {
|
||||
payload, err := io.ReadAll(io.LimitReader(r.Body, webhookMaxBodyBytes))
|
||||
if err != nil {
|
||||
log.Printf("Error reading webhook body: %v", err)
|
||||
w.WriteHeader(http.StatusServiceUnavailable)
|
||||
return
|
||||
}
|
||||
|
||||
sig := r.Header.Get("Stripe-Signature")
|
||||
event, err := webhook.ConstructEvent(payload, sig, h.cfg.StripeWebhookSecret)
|
||||
if err != nil {
|
||||
log.Printf("Webhook signature verification failed: %v", err)
|
||||
w.WriteHeader(http.StatusBadRequest)
|
||||
return
|
||||
}
|
||||
|
||||
switch event.Type {
|
||||
case "checkout.session.completed":
|
||||
h.handleCheckoutCompleted(r, event)
|
||||
case "customer.subscription.updated":
|
||||
h.handleSubscriptionUpdated(r, event)
|
||||
case "customer.subscription.deleted":
|
||||
h.handleSubscriptionDeleted(r, event)
|
||||
default:
|
||||
log.Printf("Unhandled Stripe event type: %s", event.Type)
|
||||
}
|
||||
|
||||
w.WriteHeader(http.StatusOK)
|
||||
}
|
||||
|
||||
func (h *StripeHandler) handleCheckoutCompleted(r *http.Request, event stripe.Event) {
|
||||
var session stripe.CheckoutSession
|
||||
if err := json.Unmarshal(event.Data.Raw, &session); err != nil {
|
||||
log.Printf("Error parsing checkout session: %v", err)
|
||||
return
|
||||
}
|
||||
|
||||
userID := session.ClientReferenceID
|
||||
if userID == "" {
|
||||
log.Println("Checkout session missing client_reference_id")
|
||||
return
|
||||
}
|
||||
|
||||
customerID := ""
|
||||
if session.Customer != nil {
|
||||
customerID = session.Customer.ID
|
||||
}
|
||||
subscriptionID := ""
|
||||
if session.Subscription != nil {
|
||||
subscriptionID = session.Subscription.ID
|
||||
}
|
||||
|
||||
if customerID != "" {
|
||||
if err := h.users.UpdateStripeCustomer(r.Context(), userID, customerID); err != nil {
|
||||
log.Printf("Failed to save Stripe customer ID: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
if subscriptionID != "" && customerID != "" {
|
||||
if err := h.users.ActivateSubscription(r.Context(), customerID, subscriptionID, "active"); err != nil {
|
||||
log.Printf("Failed to activate subscription: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
log.Printf("Checkout completed for user %s, customer %s, subscription %s", userID, customerID, subscriptionID)
|
||||
}
|
||||
|
||||
func (h *StripeHandler) handleSubscriptionUpdated(r *http.Request, event stripe.Event) {
|
||||
var sub stripe.Subscription
|
||||
if err := json.Unmarshal(event.Data.Raw, &sub); err != nil {
|
||||
log.Printf("Error parsing subscription update: %v", err)
|
||||
return
|
||||
}
|
||||
|
||||
customerID := ""
|
||||
if sub.Customer != nil {
|
||||
customerID = sub.Customer.ID
|
||||
}
|
||||
if customerID == "" {
|
||||
return
|
||||
}
|
||||
|
||||
status := string(sub.Status)
|
||||
if err := h.users.ActivateSubscription(r.Context(), customerID, sub.ID, status); err != nil {
|
||||
log.Printf("Failed to update subscription status: %v", err)
|
||||
}
|
||||
|
||||
log.Printf("Subscription %s updated to %s for customer %s", sub.ID, status, customerID)
|
||||
}
|
||||
|
||||
func (h *StripeHandler) handleSubscriptionDeleted(r *http.Request, event stripe.Event) {
|
||||
var sub stripe.Subscription
|
||||
if err := json.Unmarshal(event.Data.Raw, &sub); err != nil {
|
||||
log.Printf("Error parsing subscription deletion: %v", err)
|
||||
return
|
||||
}
|
||||
|
||||
customerID := ""
|
||||
if sub.Customer != nil {
|
||||
customerID = sub.Customer.ID
|
||||
}
|
||||
if customerID == "" {
|
||||
return
|
||||
}
|
||||
|
||||
if err := h.users.UpdateSubscriptionStatus(r.Context(), customerID, "canceled"); err != nil {
|
||||
log.Printf("Failed to mark subscription canceled: %v", err)
|
||||
}
|
||||
|
||||
log.Printf("Subscription canceled for customer %s", customerID)
|
||||
}
|
||||
@@ -100,6 +100,43 @@ func Authenticate(userModel *models.UserModel) func(http.Handler) http.Handler {
|
||||
}
|
||||
}
|
||||
|
||||
// RequireSubscription blocks requests from users without an active subscription.
|
||||
// Must be applied after Authenticate.
|
||||
func RequireSubscription(userModel *models.UserModel) func(http.Handler) http.Handler {
|
||||
return func(next http.Handler) http.Handler {
|
||||
return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||
userID := GetUserID(r.Context())
|
||||
if userID == "" {
|
||||
writeJSON(w, http.StatusUnauthorized, errorResponse{
|
||||
Error: "UNAUTHORIZED",
|
||||
Message: "Authentication required",
|
||||
})
|
||||
return
|
||||
}
|
||||
|
||||
user, err := userModel.GetByID(r.Context(), userID)
|
||||
if err != nil || user == nil {
|
||||
log.Printf("RequireSubscription: failed to load user %s: %v", userID, err)
|
||||
writeJSON(w, http.StatusInternalServerError, errorResponse{
|
||||
Error: "INTERNAL_ERROR",
|
||||
Message: "Failed to verify subscription",
|
||||
})
|
||||
return
|
||||
}
|
||||
|
||||
if user.SubscriptionStatus != "active" && user.SubscriptionStatus != "trialing" {
|
||||
writeJSON(w, http.StatusForbidden, errorResponse{
|
||||
Error: "SUBSCRIPTION_REQUIRED",
|
||||
Message: "An active subscription is required to use this feature",
|
||||
})
|
||||
return
|
||||
}
|
||||
|
||||
next.ServeHTTP(w, r)
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func GetUserID(ctx context.Context) string {
|
||||
if v, ok := ctx.Value(UserIDKey).(string); ok {
|
||||
return v
|
||||
|
||||
@@ -17,37 +17,73 @@ func NewUserModel(pool *pgxpool.Pool) *UserModel {
|
||||
return &UserModel{pool: pool}
|
||||
}
|
||||
|
||||
// FindOrCreateByFirebaseUID returns the local user for a Firebase UID,
|
||||
// creating one if it doesn't exist yet. On conflict (returning user) the
|
||||
// updated_at timestamp is refreshed.
|
||||
func (m *UserModel) FindOrCreateByFirebaseUID(ctx context.Context, firebaseUID, email string) (*domain.User, error) {
|
||||
var u domain.User
|
||||
err := m.pool.QueryRow(ctx,
|
||||
`INSERT INTO users (firebase_uid, email)
|
||||
VALUES ($1, $2)
|
||||
ON CONFLICT (firebase_uid) DO UPDATE SET updated_at = NOW()
|
||||
RETURNING id, firebase_uid, email, display_name, created_at, updated_at`,
|
||||
firebaseUID, email,
|
||||
).Scan(&u.ID, &u.FirebaseUID, &u.Email, &u.DisplayName, &u.CreatedAt, &u.UpdatedAt)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return &u, nil
|
||||
}
|
||||
const userColumns = `id, firebase_uid, email, display_name,
|
||||
stripe_customer_id, stripe_subscription_id, subscription_status,
|
||||
subscription_created_at, created_at, updated_at`
|
||||
|
||||
func (m *UserModel) GetByID(ctx context.Context, id string) (*domain.User, error) {
|
||||
func scanUser(row pgx.Row) (*domain.User, error) {
|
||||
var u domain.User
|
||||
err := m.pool.QueryRow(ctx,
|
||||
`SELECT id, firebase_uid, email, display_name, created_at, updated_at
|
||||
FROM users
|
||||
WHERE id = $1`,
|
||||
id,
|
||||
).Scan(&u.ID, &u.FirebaseUID, &u.Email, &u.DisplayName, &u.CreatedAt, &u.UpdatedAt)
|
||||
err := row.Scan(
|
||||
&u.ID, &u.FirebaseUID, &u.Email, &u.DisplayName,
|
||||
&u.StripeCustomerID, &u.StripeSubscriptionID, &u.SubscriptionStatus,
|
||||
&u.SubscriptionCreatedAt, &u.CreatedAt, &u.UpdatedAt,
|
||||
)
|
||||
if err != nil {
|
||||
if errors.Is(err, pgx.ErrNoRows) {
|
||||
return nil, nil
|
||||
return nil, nil //nolint:nilnil
|
||||
}
|
||||
return nil, err
|
||||
}
|
||||
return &u, nil
|
||||
}
|
||||
|
||||
// FindOrCreateByFirebaseUID returns the local user for a Firebase UID,
|
||||
// creating one if it doesn't exist yet. On conflict (returning user) the
|
||||
// updated_at timestamp is refreshed.
|
||||
func (m *UserModel) FindOrCreateByFirebaseUID(ctx context.Context, firebaseUID, email string) (*domain.User, error) {
|
||||
row := m.pool.QueryRow(ctx,
|
||||
`INSERT INTO users (firebase_uid, email)
|
||||
VALUES ($1, $2)
|
||||
ON CONFLICT (firebase_uid) DO UPDATE SET updated_at = NOW()
|
||||
RETURNING `+userColumns,
|
||||
firebaseUID, email,
|
||||
)
|
||||
return scanUser(row)
|
||||
}
|
||||
|
||||
func (m *UserModel) GetByID(ctx context.Context, id string) (*domain.User, error) {
|
||||
row := m.pool.QueryRow(ctx,
|
||||
`SELECT `+userColumns+` FROM users WHERE id = $1`, id,
|
||||
)
|
||||
return scanUser(row)
|
||||
}
|
||||
|
||||
func (m *UserModel) UpdateStripeCustomer(ctx context.Context, userID, stripeCustomerID string) error {
|
||||
_, err := m.pool.Exec(ctx,
|
||||
`UPDATE users SET stripe_customer_id = $2, updated_at = NOW() WHERE id = $1`,
|
||||
userID, stripeCustomerID,
|
||||
)
|
||||
return err
|
||||
}
|
||||
|
||||
func (m *UserModel) ActivateSubscription(ctx context.Context, stripeCustomerID, subscriptionID, status string) error {
|
||||
_, err := m.pool.Exec(ctx,
|
||||
`UPDATE users
|
||||
SET stripe_subscription_id = $2,
|
||||
subscription_status = $3,
|
||||
subscription_created_at = COALESCE(subscription_created_at, NOW()),
|
||||
updated_at = NOW()
|
||||
WHERE stripe_customer_id = $1`,
|
||||
stripeCustomerID, subscriptionID, status,
|
||||
)
|
||||
return err
|
||||
}
|
||||
|
||||
func (m *UserModel) UpdateSubscriptionStatus(ctx context.Context, stripeCustomerID, status string) error {
|
||||
_, err := m.pool.Exec(ctx,
|
||||
`UPDATE users SET subscription_status = $2, updated_at = NOW()
|
||||
WHERE stripe_customer_id = $1`,
|
||||
stripeCustomerID, status,
|
||||
)
|
||||
return err
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user