-
-
Save stympy/7cb3d82cdc8bde17299c2e77b7301d83 to your computer and use it in GitHub Desktop.
Auto-purging S3 objects using the Serverless framework
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| package main | |
| import ( | |
| "context" | |
| "log" | |
| "os" | |
| "github.com/aws/aws-lambda-go/events" | |
| "github.com/aws/aws-lambda-go/lambda" | |
| "github.com/aws/aws-sdk-go/aws" | |
| "github.com/aws/aws-sdk-go/aws/awserr" | |
| "github.com/aws/aws-sdk-go/aws/session" | |
| "github.com/aws/aws-sdk-go/service/s3" | |
| ) | |
| func handler(ctx context.Context, e events.DynamoDBEvent) { | |
| sess, err := session.NewSession(&aws.Config{Region: aws.String("us-east-1")}) | |
| if err != nil { | |
| log.Fatal(err.Error()) | |
| } | |
| svc := s3.New(sess) | |
| for _, record := range e.Records { | |
| // We only care about TTL expirations, which show up as REMOVE events | |
| if record.EventName == "REMOVE" { | |
| // The "Keys" field in the DynamoDB record is a list of S3 keys to be deleted | |
| for _, key := range record.Change.OldImage["Keys"].StringSet() { | |
| input := &s3.DeleteObjectInput{ | |
| Bucket: aws.String(os.ExpandEnv("data-bucket-$STAGE")), | |
| Key: aws.String(key), | |
| } | |
| _, err := svc.DeleteObject(input) | |
| if err != nil { | |
| if aerr, ok := err.(awserr.Error); ok { | |
| switch aerr.Code() { | |
| case s3.ErrCodeNoSuchKey: | |
| log.Println(s3.ErrCodeNoSuchKey, aerr.Error()) | |
| default: | |
| log.Println(aerr.Error()) | |
| } | |
| } else { | |
| log.Println(err.Error()) | |
| } | |
| } | |
| } | |
| } | |
| } | |
| } | |
| func main() { | |
| lambda.Start(handler) | |
| } |
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| package main | |
| import ( | |
| "context" | |
| "encoding/json" | |
| "fmt" | |
| "log" | |
| "os" | |
| "time" | |
| "github.com/aws/aws-lambda-go/events" | |
| "github.com/aws/aws-lambda-go/lambda" | |
| "github.com/aws/aws-lambda-go/lambdacontext" | |
| "github.com/aws/aws-sdk-go/aws" | |
| "github.com/aws/aws-sdk-go/aws/awserr" | |
| "github.com/aws/aws-sdk-go/aws/session" | |
| "github.com/aws/aws-sdk-go/service/s3" | |
| "github.com/guregu/dynamo" | |
| ) | |
| type notice struct { | |
| Key string `json:"key"` | |
| ExpireAt time.Time `json:"expire_at"` | |
| } | |
| type item struct { | |
| ID string | |
| ExpireAt int64 | |
| Keys []string `dynamo:",set"` | |
| } | |
| var sess = session.Must(session.NewSession(&aws.Config{Region: aws.String("us-east-1")})) | |
| var svc = s3.New(sess) | |
| var db = dynamo.New(session.New(), &aws.Config{Region: aws.String("us-east-1")}) | |
| var table = db.Table(os.ExpandEnv("expirations-$STAGE")) | |
| // Get all the notice records stored in the object at s3://bucket/key, | |
| // pull out the s3 key and expire_at data, then build a hash | |
| // with the expiration times (as epoch timestamp) for the keys and | |
| // lists of paths to the s3 objects for the values. | |
| func getItems(bucket string, key string) map[int64][]string { | |
| input := &s3.GetObjectInput{ | |
| Bucket: aws.String(bucket), | |
| Key: aws.String(key), | |
| } | |
| records := make(map[int64][]string) | |
| result, err := svc.GetObject(input) | |
| // Bail if we can't load the object from s3 | |
| if err != nil { | |
| if aerr, ok := err.(awserr.Error); ok { | |
| switch aerr.Code() { | |
| case s3.ErrCodeNoSuchKey: | |
| log.Fatal(s3.ErrCodeNoSuchKey, aerr.Error()) | |
| default: | |
| log.Fatal(aerr.Error()) | |
| } | |
| } else { | |
| log.Fatal(err.Error()) | |
| } | |
| } | |
| // Use a streaming JSON decoder to save on memory usage | |
| // Get started by reading the first [ of the s3 object body | |
| dec := json.NewDecoder(result.Body) | |
| _, err = dec.Token() | |
| if err != nil { | |
| log.Fatal(err) | |
| } | |
| for dec.More() { | |
| // Decode each json hash into a notice struct | |
| var n notice | |
| err := dec.Decode(&n) | |
| if err != nil { | |
| log.Println(err) | |
| } else { | |
| records[n.ExpireAt.Unix()] = append(records[n.ExpireAt.Unix()], n.Key) | |
| } | |
| } | |
| // Finish off by reading the trailing ] | |
| _, err = dec.Token() | |
| if err != nil { | |
| log.Fatal(err) | |
| } | |
| return records | |
| } | |
| func handler(ctx context.Context, snsEvent events.SNSEvent) { | |
| lc, _ := lambdacontext.FromContext(ctx) | |
| for _, record := range snsEvent.Records { | |
| snsRecord := record.SNS | |
| var s3Event events.S3Event | |
| json.Unmarshal([]byte(snsRecord.Message), &s3Event) | |
| for _, event := range s3Event.Records { | |
| items := getItems(event.S3.Bucket.Name, event.S3.Object.Key) | |
| for k, v := range items { | |
| err := table.Put(item{ID: event.S3.Object.Key, ExpireAt: k, Keys: v}).Run() | |
| if err != nil { | |
| log.Println(err) | |
| } | |
| } | |
| } | |
| } | |
| } | |
| func main() { | |
| lambda.Start(handler) | |
| } |
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment