Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
42 changes: 38 additions & 4 deletions internal/config/config.go
Original file line number Diff line number Diff line change
Expand Up @@ -17,10 +17,22 @@ type Config struct {
MaxMindLicenseKey string // MAXMIND_LICENSE_KEY — for GeoLite2 refresh job
GeoLite2DBPath string // GEOLITE2_DB_PATH — local path to the GeoLite2 City MMDB
PlansPath string // PLANS_PATH — path to plans.yaml (optional; uses built-in defaults if empty)
MinioEndpoint string // MINIO_ENDPOINT — host:port of MinIO server
MinioRootUser string // MINIO_ROOT_USER — MinIO admin credentials
MinioRootPassword string // MINIO_ROOT_PASSWORD
MinioBucketName string // MINIO_BUCKET_NAME — shared bucket name (default: "instant-shared")
// Object storage backend for the storage_bytes scanner
// (provider-agnostic — works against MinIO, DO Spaces, AWS S3,
// GCS, R2, B2, Wasabi). The scanner uses the standard minio-go
// SDK, which speaks plain S3 against any of these endpoints.
ObjectStoreEndpoint string // OBJECT_STORE_ENDPOINT — host:port (e.g. nyc3.digitaloceanspaces.com)
ObjectStoreAccessKey string // OBJECT_STORE_ACCESS_KEY — master access key
ObjectStoreSecretKey string // OBJECT_STORE_SECRET_KEY — master secret key
ObjectStoreBucket string // OBJECT_STORE_BUCKET — shared bucket (default: instant-shared)
ObjectStoreRegion string // OBJECT_STORE_REGION — e.g. "nyc3" for DO Spaces
ObjectStoreSecure bool // OBJECT_STORE_SECURE — true for TLS-terminated endpoints

// Legacy MINIO_* env vars — fallback for backward compat.
MinioEndpoint string
MinioRootUser string
MinioRootPassword string
MinioBucketName string
KubeNamespaceApps string // KUBE_NAMESPACE_APPS — stack namespace prefix (default: "instant-apps")
}

Expand Down Expand Up @@ -60,13 +72,35 @@ func Load() *Config {
MaxMindLicenseKey: os.Getenv("MAXMIND_LICENSE_KEY"),
GeoLite2DBPath: getenv("GEOLITE2_DB_PATH", "./GeoLite2-City.mmdb"),
PlansPath: os.Getenv("PLANS_PATH"),
ObjectStoreEndpoint: os.Getenv("OBJECT_STORE_ENDPOINT"),
ObjectStoreAccessKey: os.Getenv("OBJECT_STORE_ACCESS_KEY"),
ObjectStoreSecretKey: os.Getenv("OBJECT_STORE_SECRET_KEY"),
ObjectStoreBucket: getenv("OBJECT_STORE_BUCKET", "instant-shared"),
ObjectStoreRegion: os.Getenv("OBJECT_STORE_REGION"),
ObjectStoreSecure: os.Getenv("OBJECT_STORE_SECURE") == "true",

MinioEndpoint: os.Getenv("MINIO_ENDPOINT"),
MinioRootUser: os.Getenv("MINIO_ROOT_USER"),
MinioRootPassword: os.Getenv("MINIO_ROOT_PASSWORD"),
MinioBucketName: getenv("MINIO_BUCKET_NAME", "instant-shared"),
KubeNamespaceApps: getenv("KUBE_NAMESPACE_APPS", "instant-apps"),
}

// Fall back to legacy MINIO_* names so deployments that haven't
// migrated env vars keep working.
if cfg.ObjectStoreEndpoint == "" {
cfg.ObjectStoreEndpoint = cfg.MinioEndpoint
}
if cfg.ObjectStoreAccessKey == "" {
cfg.ObjectStoreAccessKey = cfg.MinioRootUser
}
if cfg.ObjectStoreSecretKey == "" {
cfg.ObjectStoreSecretKey = cfg.MinioRootPassword
}
if cfg.ObjectStoreBucket == "instant-shared" && cfg.MinioBucketName != "" {
cfg.ObjectStoreBucket = cfg.MinioBucketName
}

slog.Info("worker.config.loaded",
"environment", cfg.Environment,
"provisioner_addr_set", cfg.ProvisionerAddr != "",
Expand Down
46 changes: 42 additions & 4 deletions internal/jobs/storage_minio.go
Original file line number Diff line number Diff line change
Expand Up @@ -35,21 +35,59 @@ type minioStorageScanner struct {
}

// NewMinIOStorageScanner constructs a scanner backed by github.com/minio/minio-go/v7
// using the same root/admin credentials the worker already loads from
// MINIO_ENDPOINT / MINIO_ROOT_USER / MINIO_ROOT_PASSWORD (see config.go).
// against any S3-compatible endpoint (self-hosted MinIO, DO Spaces, AWS S3,
// GCS, R2, B2, Wasabi). Callers source endpoint + creds + bucket from the
// OBJECT_STORE_* env vars in config.Load (which fall back to legacy MINIO_*
// names for backward compat).
//
// Auto-detects TLS: if `endpoint` is prefixed with "https://" the scanner uses
// TLS; "http://" forces plain. Without a scheme, a heuristic kicks in — a
// hostname containing "digitaloceanspaces.com", "amazonaws.com",
// "cloudflarestorage.com", "googleapis.com", "wasabisys.com", or
// "backblazeb2.com" is assumed TLS; everything else (e.g. an in-cluster
// minio.instant-data.svc.cluster.local) is plain. Callers that need explicit
// control should call NewS3Scanner directly.
//
// Returns nil + error when the endpoint is empty or the client can't be built;
// callers should fail open and pass nil to NewUpdateStorageBytesWorker.
func NewMinIOStorageScanner(endpoint, accessKey, secretKey, bucketName string) (*minioStorageScanner, error) {
if endpoint == "" {
return nil, fmt.Errorf("storage_minio: MINIO_ENDPOINT is required")
return nil, fmt.Errorf("storage_minio: OBJECT_STORE_ENDPOINT is required")
}
if bucketName == "" {
bucketName = "instant-shared"
}

// Strip explicit scheme prefix and remember whether TLS was requested.
secure := false
if strings.HasPrefix(endpoint, "https://") {
endpoint = strings.TrimPrefix(endpoint, "https://")
secure = true
} else if strings.HasPrefix(endpoint, "http://") {
endpoint = strings.TrimPrefix(endpoint, "http://")
secure = false
} else {
// Heuristic: managed S3-compatible vendors all serve TLS by default.
// In-cluster MinIO uses plain HTTP. Misidentified endpoints can be
// fixed by explicitly prefixing http:// or https:// in the env var.
for _, vendor := range []string{
"digitaloceanspaces.com",
"amazonaws.com",
"cloudflarestorage.com",
"googleapis.com",
"wasabisys.com",
"backblazeb2.com",
} {
if strings.Contains(endpoint, vendor) {
secure = true
break
}
}
}

client, err := minio.New(endpoint, &minio.Options{
Creds: credentials.NewStaticV4(accessKey, secretKey, ""),
Secure: false,
Secure: secure,
})
if err != nil {
return nil, fmt.Errorf("storage_minio: new client for %s: %w", endpoint, err)
Expand Down
22 changes: 14 additions & 8 deletions internal/jobs/workers.go
Original file line number Diff line number Diff line change
Expand Up @@ -101,7 +101,11 @@ func StartWorkers(ctx context.Context, db *sql.DB, rdb *redis.Client, cfg *confi

emailClient := NewEmailClient(cfg.ResendAPIKey)

// Build MinIO admin client for storage IAM cleanup — nil if not configured (fail open).
// Build MinIO admin client for storage IAM cleanup — nil unless the legacy
// MINIO_* env vars are set. Only used when ExpireAnonymousWorker needs to
// release per-IAM-user resources (i.e. self-hosted MinIO backend). With
// the OBJECT_STORE_* shared-key backend (DO Spaces / AWS / GCS / R2) this
// stays nil because no per-customer IAM was created in the first place.
var minioClient *madmin.AdminClient
if cfg.MinioEndpoint != "" {
if mc, err := madmin.New(cfg.MinioEndpoint, cfg.MinioRootUser, cfg.MinioRootPassword, false); err != nil {
Expand All @@ -111,14 +115,16 @@ func StartWorkers(ctx context.Context, db *sql.DB, rdb *redis.Client, cfg *confi
}
}

// Build MinIO storage scanner for the UpdateStorageBytesWorker — nil if
// MINIO_ENDPOINT is not set (fail open: storage_bytes updates for MinIO
// resources are skipped each run, postgres/redis/mongo continue via the
// gRPC provisioner path).
// Build the storage_bytes scanner — provider-agnostic, uses plain S3 API
// against any backend. Reads OBJECT_STORE_* env vars (which fall back to
// the legacy MINIO_* names in config.Load). Nil = fail open: the scanner
// is skipped each run and storage_bytes for /storage/new resources isn't
// updated. Other resource types (postgres/redis/mongo) continue via the
// gRPC provisioner path.
var minioScanner MinIOStorageScanner
if cfg.MinioEndpoint != "" {
if scanner, err := NewMinIOStorageScanner(cfg.MinioEndpoint, cfg.MinioRootUser, cfg.MinioRootPassword, cfg.MinioBucketName); err != nil {
slog.Warn("jobs.workers.minio_storage_scanner_init_failed", "error", err)
if cfg.ObjectStoreEndpoint != "" {
if scanner, err := NewMinIOStorageScanner(cfg.ObjectStoreEndpoint, cfg.ObjectStoreAccessKey, cfg.ObjectStoreSecretKey, cfg.ObjectStoreBucket); err != nil {
slog.Warn("jobs.workers.storage_scanner_init_failed", "error", err)
} else {
minioScanner = scanner
}
Expand Down