Skip to content

Instantly share code, notes, and snippets.

@stympy
Created April 12, 2018 18:26
Show Gist options
  • Select an option

  • Save stympy/7cb3d82cdc8bde17299c2e77b7301d83 to your computer and use it in GitHub Desktop.

Select an option

Save stympy/7cb3d82cdc8bde17299c2e77b7301d83 to your computer and use it in GitHub Desktop.
Auto-purging S3 objects using the Serverless framework
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)
}
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