153 lines
4.3 KiB
Go
153 lines
4.3 KiB
Go
// CRUD demo against the universal MinIO storage frontend
|
|
// (minio.storage.svc.cluster.local:9000).
|
|
//
|
|
// Run from outside the cluster via a port-forward:
|
|
//
|
|
// kubectl port-forward svc/minio -n storage 9000:9000 &
|
|
// source logging/.env
|
|
// cd storage/test && go mod tidy && go run .
|
|
//
|
|
// In-cluster, set MINIO_ENDPOINT=minio.storage.svc.cluster.local:9000.
|
|
//
|
|
// Every S3 call emits [SERVICE_METRIC] op latency; every failure emits
|
|
// [APP_METRIC] ERROR with context and aborts (no silent catches).
|
|
package main
|
|
|
|
import (
|
|
"bytes"
|
|
"context"
|
|
"fmt"
|
|
"io"
|
|
"math/rand"
|
|
"os"
|
|
"time"
|
|
|
|
"github.com/minio/minio-go/v7"
|
|
"github.com/minio/minio-go/v7/pkg/credentials"
|
|
)
|
|
|
|
const bucket = "crud-test"
|
|
|
|
func getenv(key, fallback string) string {
|
|
if v := os.Getenv(key); v != "" {
|
|
return v
|
|
}
|
|
return fallback
|
|
}
|
|
|
|
// timed wraps an S3 operation: emits a serviceMetric on success,
|
|
// an applicationMetric and exit(1) on failure.
|
|
func timed(op string, fn func() error) {
|
|
start := time.Now()
|
|
if err := fn(); err != nil {
|
|
fmt.Printf("[APP_METRIC] ERROR s3.%s failed bucket=%s | trace=%v\n", op, bucket, err)
|
|
os.Exit(1)
|
|
}
|
|
fmt.Printf("[SERVICE_METRIC] s3.%s.latency_ms=%d ms\n", op, time.Since(start).Milliseconds())
|
|
}
|
|
|
|
func randomText(n int) []byte {
|
|
const letters = "abcdefghijklmnopqrstuvwxyz \n"
|
|
b := make([]byte, n)
|
|
for i := range b {
|
|
b[i] = letters[rand.Intn(len(letters))]
|
|
}
|
|
return b
|
|
}
|
|
|
|
func main() {
|
|
endpoint := getenv("MINIO_ENDPOINT", "localhost:9000")
|
|
user := os.Getenv("MINIO_ROOT_USER")
|
|
pass := os.Getenv("MINIO_ROOT_PASSWORD")
|
|
if user == "" || pass == "" {
|
|
fmt.Println("[APP_METRIC] ERROR config missing | trace=MINIO_ROOT_USER / MINIO_ROOT_PASSWORD not set (source storage/.env)")
|
|
os.Exit(1)
|
|
}
|
|
|
|
ctx := context.Background()
|
|
client, err := minio.New(endpoint, &minio.Options{
|
|
Creds: credentials.NewStaticV4(user, pass, ""),
|
|
Secure: false, // in-cluster traffic, no TLS
|
|
})
|
|
if err != nil {
|
|
fmt.Printf("[APP_METRIC] ERROR s3.connect failed endpoint=%s | trace=%v\n", endpoint, err)
|
|
os.Exit(1)
|
|
}
|
|
|
|
key := fmt.Sprintf("demo/%d.txt", time.Now().Unix())
|
|
original := randomText(256)
|
|
updated := append([]byte("UPDATED ---\n"), randomText(128)...)
|
|
|
|
// Ensure bucket (idempotent). Site replication propagates it to az-b.
|
|
timed("ensure_bucket", func() error {
|
|
exists, err := client.BucketExists(ctx, bucket)
|
|
if err != nil || exists {
|
|
return err
|
|
}
|
|
return client.MakeBucket(ctx, bucket, minio.MakeBucketOptions{})
|
|
})
|
|
|
|
// CREATE
|
|
timed("put", func() error {
|
|
_, err := client.PutObject(ctx, bucket, key,
|
|
bytes.NewReader(original), int64(len(original)),
|
|
minio.PutObjectOptions{ContentType: "text/plain"})
|
|
return err
|
|
})
|
|
fmt.Printf("created %s/%s (%d bytes of random text)\n", bucket, key, len(original))
|
|
|
|
// READ — and verify content round-trips
|
|
timed("get", func() error {
|
|
obj, err := client.GetObject(ctx, bucket, key, minio.GetObjectOptions{})
|
|
if err != nil {
|
|
return err
|
|
}
|
|
defer obj.Close()
|
|
got, err := io.ReadAll(obj)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if !bytes.Equal(got, original) {
|
|
return fmt.Errorf("read-back mismatch: want %d bytes, got %d", len(original), len(got))
|
|
}
|
|
return nil
|
|
})
|
|
fmt.Println("read back and verified content")
|
|
|
|
// UPDATE — S3 semantics: overwrite the object in place
|
|
timed("update", func() error {
|
|
_, err := client.PutObject(ctx, bucket, key,
|
|
bytes.NewReader(updated), int64(len(updated)),
|
|
minio.PutObjectOptions{ContentType: "text/plain"})
|
|
return err
|
|
})
|
|
fmt.Println("updated (overwrote) object")
|
|
|
|
// LIST the demo/ prefix
|
|
timed("list", func() error {
|
|
for obj := range client.ListObjects(ctx, bucket, minio.ListObjectsOptions{Prefix: "demo/", Recursive: true}) {
|
|
if obj.Err != nil {
|
|
return obj.Err
|
|
}
|
|
fmt.Printf(" %s %d bytes %s\n", obj.Key, obj.Size, obj.LastModified.Format(time.RFC3339))
|
|
}
|
|
return nil
|
|
})
|
|
|
|
// DELETE — and verify it is gone
|
|
timed("delete", func() error {
|
|
if err := client.RemoveObject(ctx, bucket, key, minio.RemoveObjectOptions{}); err != nil {
|
|
return err
|
|
}
|
|
_, err := client.StatObject(ctx, bucket, key, minio.StatObjectOptions{})
|
|
if err == nil {
|
|
return fmt.Errorf("object %s still exists after delete", key)
|
|
}
|
|
if minio.ToErrorResponse(err).Code != "NoSuchKey" {
|
|
return err
|
|
}
|
|
return nil
|
|
})
|
|
fmt.Println("deleted and verified gone — CRUD cycle complete")
|
|
}
|