package main import ( "context" "encoding/json" "errors" "fmt" "log" "net/http" "os" "os/signal" "strconv" "strings" "syscall" "time" ) const ( defaultMaxObjectSize = 1024 * 1024 * 1024 defaultMinFreeBytes = 1024 * 1024 * 1024 defaultPresignedMaxExpiry = 24 * time.Hour defaultMultipartMaxAge = 24 * time.Hour defaultMaxRestoreBytes = 10 * 1024 * 1024 * 1024 defaultReplicationMaxJobs = 10_000 defaultReplicationMaxAttempts = 100 ) func main() { cmd := "serve" if len(os.Args) > 1 { cmd = os.Args[1] } if cmd == "version" { printJSON(versionInfo()) return } if cmd == "help" || cmd == "-h" || cmd == "--help" { printUsage() return } cfg := configFromEnv() var store *FileStore openStore := func() *FileStore { if store != nil { return store } opened, err := NewFileStore(cfg.DataDir) if err != nil { log.Fatalf("init store: %v", err) } store = opened return store } switch cmd { case "serve": serve(cfg, openStore()) case "reindex": mustJSON(openStore().Reindex()) case "stats": stats, err := openStore().Stats() if err != nil { log.Fatal(err) } printJSON(stats) case "scrub": report, err := openStore().Scrub() if err != nil { log.Fatal(err) } printJSON(report) case "repair": count, err := openStore().EnqueueRepair(cfg.ReplicationPeers) if err != nil { log.Fatal(err) } printJSON(map[string]any{"status": "queued", "jobs": count}) case "backup": if len(os.Args) < 3 { log.Fatal("usage: magpie backup ") } mustJSON(BackupDir(cfg.DataDir, os.Args[2])) case "restore": if len(os.Args) < 3 { log.Fatal("usage: magpie restore ") } mustJSON(RestoreDirWithLimit(os.Args[2], cfg.DataDir, cfg.MaxRestoreBytes)) case "smoke": mustJSON(runSmoke(openStore())) default: log.Fatalf("unknown command %q", sanitizeLogValue(cmd)) } } func serve(cfg Config, store *FileStore) { if cfg.ReindexOnStart { log.Printf("reindexing metadata from %s", cfg.DataDir) if err := store.Reindex(); err != nil { log.Fatalf("reindex store: %v", err) } } server := NewServerInstance(store, cfg) srv := &http.Server{ Addr: cfg.Addr, Handler: server.Handler(), ReadHeaderTimeout: 5 * time.Second, ReadTimeout: 2 * time.Minute, WriteTimeout: 2 * time.Minute, IdleTimeout: 60 * time.Second, MaxHeaderBytes: 16 * 1024, } go func() { log.Printf("object store listening on %s, data=%s", cfg.Addr, cfg.DataDir) if err := srv.ListenAndServe(); err != nil && !errors.Is(err, http.ErrServerClosed) { log.Fatalf("listen: %v", err) } }() ctx, stop := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM) defer stop() server.startMultipartCleanup(ctx) server.startReplicationWorker(ctx) <-ctx.Done() shutdownCtx, cancel := context.WithTimeout(context.Background(), 10*time.Second) defer cancel() if err := srv.Shutdown(shutdownCtx); err != nil { log.Printf("shutdown: %v", err) } } type Config struct { Addr string DataDir string AuthToken string S3AccessKeyID string S3SecretAccessKey string S3Keys []AccessKey S3Region string AllowUnsignedLocal bool MaxObjectSize int64 MinFreeBytes uint64 PresignedMaxExpiry time.Duration MultipartMaxAge time.Duration ReindexOnStart bool ReplicationPeers []string ReplicationSecret string AllowInsecureReplication bool ReplicationMaxJobs int ReplicationMaxAttempts int MaxRestoreBytes int64 AllowedBuckets map[string]bool PublicPrefixes map[string]bool RateLimitPerMinute int TrustProxyHeaders bool } type AccessKey struct { ID string Secret string Permissions map[string]bool } func configFromEnv() Config { addr := getenv("MAGPIE_ADDR", "127.0.0.1:8090") dataDir := getenv("MAGPIE_DATA_DIR", "./data") authToken := os.Getenv("MAGPIE_AUTH_TOKEN") s3AccessKeyID := os.Getenv("MAGPIE_S3_ACCESS_KEY_ID") s3SecretAccessKey := os.Getenv("MAGPIE_S3_SECRET_ACCESS_KEY") s3Keys := parseAccessKeys(os.Getenv("MAGPIE_S3_KEYS"), s3AccessKeyID, s3SecretAccessKey) s3Region := getenv("MAGPIE_S3_REGION", "auto") maxObjectSize := mustParseBytes("MAGPIE_MAX_OBJECT_SIZE", defaultMaxObjectSize) minFreeBytes := mustParseUint("MAGPIE_MIN_FREE_BYTES", defaultMinFreeBytes) presignedMaxExpiry := mustParseDuration("MAGPIE_PRESIGNED_MAX_EXPIRY", defaultPresignedMaxExpiry) multipartMaxAge := mustParseDuration("MAGPIE_MULTIPART_MAX_AGE", defaultMultipartMaxAge) reindexOnStart := mustParseBool("MAGPIE_REINDEX_ON_START", false) allowInsecureReplication := mustParseBool("MAGPIE_ALLOW_INSECURE_REPLICATION", false) replicationPeers := parseReplicationPeers(os.Getenv("MAGPIE_REPLICATION_PEERS"), allowInsecureReplication) replicationSecret := os.Getenv("MAGPIE_REPLICATION_SECRET") replicationMaxJobs := mustParseInt("MAGPIE_REPLICATION_MAX_JOBS", defaultReplicationMaxJobs) replicationMaxAttempts := mustParseInt("MAGPIE_REPLICATION_MAX_ATTEMPTS", defaultReplicationMaxAttempts) maxRestoreBytes := mustParseBytes("MAGPIE_MAX_RESTORE_BYTES", defaultMaxRestoreBytes) allowedBuckets := parseSet(os.Getenv("MAGPIE_ALLOWED_BUCKETS")) // MAGPIE_PUBLIC_BUCKETS is a legacy whole-bucket public-read list (e.g. app-uploads). // MAGPIE_PUBLIC_PREFIXES is preferred and takes bucket/prefix values. publicPrefixes := parsePublicPrefixes(os.Getenv("MAGPIE_PUBLIC_PREFIXES")) for bucket := range parseSet(os.Getenv("MAGPIE_PUBLIC_BUCKETS")) { publicPrefixes[bucket] = true } rateLimitPerMinute := mustParseInt("MAGPIE_RATE_LIMIT_PER_MINUTE", 600) trustProxyHeaders := mustParseBool("MAGPIE_TRUST_PROXY_HEADERS", false) allowUnsignedLocal := mustParseBool("MAGPIE_ALLOW_UNSIGNED_PAYLOADS", false) if authToken == "" && len(s3Keys) == 0 { log.Fatal("MAGPIE_AUTH_TOKEN or MAGPIE_S3_ACCESS_KEY_ID/MAGPIE_S3_SECRET_ACCESS_KEY or MAGPIE_S3_KEYS is required") } if len(replicationPeers) > 0 && replicationSecret == "" { log.Fatal("MAGPIE_REPLICATION_SECRET is required when MAGPIE_REPLICATION_PEERS is set") } if raw := os.Getenv("MAGPIE_MAX_PROCS"); raw != "" { if _, err := strconv.Atoi(raw); err != nil { log.Fatalf("MAGPIE_MAX_PROCS must be an integer: %v", err) } } return Config{ Addr: addr, DataDir: dataDir, AuthToken: authToken, S3AccessKeyID: s3AccessKeyID, S3SecretAccessKey: s3SecretAccessKey, S3Keys: s3Keys, S3Region: s3Region, AllowUnsignedLocal: allowUnsignedLocal, MaxObjectSize: maxObjectSize, MinFreeBytes: minFreeBytes, PresignedMaxExpiry: presignedMaxExpiry, MultipartMaxAge: multipartMaxAge, ReindexOnStart: reindexOnStart, ReplicationPeers: replicationPeers, ReplicationSecret: replicationSecret, AllowInsecureReplication: allowInsecureReplication, ReplicationMaxJobs: replicationMaxJobs, ReplicationMaxAttempts: replicationMaxAttempts, MaxRestoreBytes: maxRestoreBytes, AllowedBuckets: allowedBuckets, PublicPrefixes: publicPrefixes, RateLimitPerMinute: rateLimitPerMinute, TrustProxyHeaders: trustProxyHeaders, } } func parseAccessKeys(raw string, legacyID string, legacySecret string) []AccessKey { keys := make([]AccessKey, 0) if legacyID != "" && legacySecret != "" { keys = append(keys, AccessKey{ID: legacyID, Secret: legacySecret, Permissions: permissions("read,write,admin")}) } for _, entry := range strings.Split(raw, ";") { entry = strings.TrimSpace(entry) if entry == "" { continue } parts := strings.Split(entry, ":") if len(parts) < 2 || parts[0] == "" || parts[1] == "" { log.Fatal("MAGPIE_S3_KEYS entries must be id:secret[:read,write,admin]") } perms := "read,write" if len(parts) >= 3 { perms = parts[2] } keys = append(keys, AccessKey{ID: parts[0], Secret: parts[1], Permissions: permissions(perms)}) } return keys } func permissions(raw string) map[string]bool { result := map[string]bool{} for _, part := range strings.Split(raw, ",") { part = strings.TrimSpace(strings.ToLower(part)) if part != "" { result[part] = true } } return result } func parseCSV(raw string) []string { parts := make([]string, 0) for _, part := range strings.Split(raw, ",") { part = strings.TrimSpace(part) if part != "" { parts = append(parts, strings.TrimRight(part, "/")) } } return parts } func parseReplicationPeers(raw string, allowInsecure bool) []string { peers := parseCSV(raw) for _, peer := range peers { if err := validateReplicationPeer(peer, allowInsecure); err != nil { log.Fatal(err) } } return peers } func parseSet(raw string) map[string]bool { values := map[string]bool{} for _, value := range parseCSV(raw) { values[value] = true } return values } func parsePublicPrefixes(raw string) map[string]bool { value := strings.TrimSpace(strings.ToLower(raw)) if value == "none" || value == "false" || value == "0" { return map[string]bool{} } if strings.TrimSpace(raw) == "" { return map[string]bool{} } prefixes := parseSet(raw) for prefix := range prefixes { cleaned := strings.Trim(prefix, "/") if cleaned == "" { log.Fatal("MAGPIE_PUBLIC_PREFIXES entries must not be empty") } // Allow bare bucket names (whole-bucket public read) or bucket/prefix paths. if err := ValidateKey(cleaned); err != nil { log.Fatal("MAGPIE_PUBLIC_PREFIXES entries must be valid bucket or bucket/prefix values") } } return prefixes } func printUsage() { fmt.Println("usage: magpie [serve|reindex|stats|scrub|repair|version|backup|restore|smoke]") } func mustJSON(err error) { if err != nil { log.Fatal(err) } printJSON(map[string]string{"status": "ok"}) } func printJSON(value any) { encoder := json.NewEncoder(os.Stdout) encoder.SetIndent("", " ") if err := encoder.Encode(value); err != nil { log.Fatal(err) } } func getenv(key, fallback string) string { if value := os.Getenv(key); value != "" { return value } return fallback } func mustParseBytes(key string, fallback int64) int64 { value := os.Getenv(key) if value == "" { return fallback } parsed, err := strconv.ParseInt(value, 10, 64) if err != nil || parsed < 0 { log.Fatal(fmt.Sprintf("%s must be a non-negative integer byte count", key)) } return parsed } func mustParseUint(key string, fallback uint64) uint64 { value := os.Getenv(key) if value == "" { return fallback } parsed, err := strconv.ParseUint(value, 10, 64) if err != nil { log.Fatal(fmt.Sprintf("%s must be a non-negative integer byte count", key)) } return parsed } func mustParseDuration(key string, fallback time.Duration) time.Duration { value := os.Getenv(key) if value == "" { return fallback } parsed, err := time.ParseDuration(value) if err != nil || parsed <= 0 { log.Fatal(fmt.Sprintf("%s must be a positive duration like 1h or 24h", key)) } return parsed } func mustParseBool(key string, fallback bool) bool { value := os.Getenv(key) if value == "" { return fallback } parsed, err := strconv.ParseBool(value) if err != nil { log.Fatal(fmt.Sprintf("%s must be true or false", key)) } return parsed } func mustParseInt(key string, fallback int) int { value := os.Getenv(key) if value == "" { return fallback } parsed, err := strconv.Atoi(value) if err != nil || parsed < 0 { log.Fatal(fmt.Sprintf("%s must be a non-negative integer", key)) } return parsed } func runSmoke(store *FileStore) error { key := "smoke/test.txt" if _, err := store.Put(key, strings.NewReader("ok"), "text/plain"); err != nil { return err } obj, err := store.Open(key) if err != nil { return err } _ = obj.Close() if _, err := store.List("smoke/", 10); err != nil { return err } if _, err := store.Scrub(); err != nil { return err } return store.Delete(key) }