cluster

Infrastructure files for Nordgedanken and Midnightthoughts.
git clone git://archive.git.mtrnord.blog/MTRNord/cluster.git
Log | Files | Refs | README

s3.go (4391B)


      1 // s3.go — S3 client initialisation, helpers, and age encryption/decryption.
      2 package main
      3 
      4 import (
      5 	"bytes"
      6 	"context"
      7 	"encoding/json"
      8 	"errors"
      9 	"fmt"
     10 	"io"
     11 	"log/slog"
     12 
     13 	"filippo.io/age"
     14 	"github.com/aws/aws-sdk-go-v2/aws"
     15 	awsconfig "github.com/aws/aws-sdk-go-v2/config"
     16 	"github.com/aws/aws-sdk-go-v2/credentials"
     17 	"github.com/aws/aws-sdk-go-v2/service/s3"
     18 	s3types "github.com/aws/aws-sdk-go-v2/service/s3/types"
     19 )
     20 
     21 var s3c *s3.Client
     22 
     23 func initS3(ctx context.Context) error {
     24 	cfg, err := awsconfig.LoadDefaultConfig(ctx,
     25 		awsconfig.WithRegion(s3Region),
     26 		awsconfig.WithCredentialsProvider(credentials.NewStaticCredentialsProvider(
     27 			s3AccessKey, s3SecretKey, "",
     28 		)),
     29 	)
     30 	if err != nil {
     31 		return err
     32 	}
     33 	endpoint := "https://" + s3Endpoint
     34 	s3c = s3.NewFromConfig(cfg, func(o *s3.Options) {
     35 		o.UsePathStyle = true
     36 		o.BaseEndpoint = &endpoint
     37 	})
     38 	return nil
     39 }
     40 
     41 func s3Get(ctx context.Context, key string) ([]byte, error) {
     42 	out, err := s3c.GetObject(ctx, &s3.GetObjectInput{
     43 		Bucket: aws.String(s3BucketName),
     44 		Key:    aws.String(key),
     45 	})
     46 	if err != nil {
     47 		var nsk *s3types.NoSuchKey
     48 		if errors.As(err, &nsk) {
     49 			return nil, nil
     50 		}
     51 		return nil, err
     52 	}
     53 	defer out.Body.Close()
     54 	return io.ReadAll(out.Body)
     55 }
     56 
     57 func s3GetJSON(ctx context.Context, key string, out interface{}) error {
     58 	data, err := s3Get(ctx, key)
     59 	if err != nil || data == nil {
     60 		return err
     61 	}
     62 	return json.Unmarshal(data, out)
     63 }
     64 
     65 func s3Put(ctx context.Context, key string, data []byte, contentType string) error {
     66 	_, err := s3c.PutObject(ctx, &s3.PutObjectInput{
     67 		Bucket:      aws.String(s3BucketName),
     68 		Key:         aws.String(key),
     69 		Body:        bytes.NewReader(data),
     70 		ContentType: aws.String(contentType),
     71 	})
     72 	return err
     73 }
     74 
     75 func s3Exists(ctx context.Context, key string) bool {
     76 	_, err := s3c.HeadObject(ctx, &s3.HeadObjectInput{
     77 		Bucket: aws.String(s3BucketName),
     78 		Key:    aws.String(key),
     79 	})
     80 	return err == nil
     81 }
     82 
     83 // s3DeletePrefix deletes all objects whose key starts with prefix.
     84 func s3DeletePrefix(ctx context.Context, prefix string) error {
     85 	paginator := s3.NewListObjectsV2Paginator(s3c, &s3.ListObjectsV2Input{
     86 		Bucket: aws.String(s3BucketName),
     87 		Prefix: aws.String(prefix),
     88 	})
     89 	deleted := 0
     90 	for paginator.HasMorePages() {
     91 		page, err := paginator.NextPage(ctx)
     92 		if err != nil {
     93 			return err
     94 		}
     95 		for _, obj := range page.Contents {
     96 			if _, err := s3c.DeleteObject(ctx, &s3.DeleteObjectInput{
     97 				Bucket: aws.String(s3BucketName),
     98 				Key:    obj.Key,
     99 			}); err != nil {
    100 				slog.Warn("Failed to delete S3 object", "key", *obj.Key, "error", err)
    101 			} else {
    102 				deleted++
    103 			}
    104 		}
    105 	}
    106 	slog.Info("Deleted S3 prefix", "prefix", prefix, "count", deleted)
    107 	return nil
    108 }
    109 
    110 // ─────────────────────────────────────────────────────────────────────────────
    111 // Age encryption / decryption
    112 // ─────────────────────────────────────────────────────────────────────────────
    113 
    114 // ageEncrypt encrypts data for all ageRecipients and returns the ciphertext.
    115 func ageEncrypt(data []byte) ([]byte, error) {
    116 	var buf bytes.Buffer
    117 	w, err := age.Encrypt(&buf, ageRecipients...)
    118 	if err != nil {
    119 		return nil, err
    120 	}
    121 	if _, err := w.Write(data); err != nil {
    122 		return nil, err
    123 	}
    124 	if err := w.Close(); err != nil {
    125 		return nil, err
    126 	}
    127 	return buf.Bytes(), nil
    128 }
    129 
    130 // getDecryptedAgeFromS3 fetches an .age object from S3 and decrypts it.
    131 // Returns (nil, nil) when the object doesn't exist or no identity is available.
    132 func getDecryptedAgeFromS3(ctx context.Context, key string) ([]byte, error) {
    133 	if ageIdentity == nil {
    134 		return nil, nil
    135 	}
    136 	enc, err := s3Get(ctx, key)
    137 	if err != nil || enc == nil {
    138 		return nil, err
    139 	}
    140 	r, err := age.Decrypt(bytes.NewReader(enc), ageIdentity)
    141 	if err != nil {
    142 		return nil, fmt.Errorf("age decrypt %s: %w", key, err)
    143 	}
    144 	return io.ReadAll(r)
    145 }
    146 
    147 // s3PutAge age-encrypts data then uploads it under key+".age".
    148 func s3PutAge(ctx context.Context, key string, data []byte) error {
    149 	enc, err := ageEncrypt(data)
    150 	if err != nil {
    151 		return fmt.Errorf("age encrypt %s: %w", key, err)
    152 	}
    153 	return s3Put(ctx, key+".age", enc, "application/octet-stream")
    154 }