Why does golang concurrent uploading of files to S3 bucket result in canceled, context deadline exceeded

Viewed 600

I have written a small golang piece of code to recursive traverse a directory and upload the files in the director. There are approximately 93K+ items in the directory. After a while I get the following error:

Got error uploading file: /Users/randolphhill/Fat-Tree-Business/SandBox/DDD/heydoc/ios/Pods/gRPC-Core/src/core/ext/transport/chttp2/alpn/alpn.h operation error S3: PutObject, https response error StatusCode: 0, RequestID: , HostID: , canceled, context deadline exceeded.

Below is the code snippet

   func PutFile(c context.Context, api S3PutObjectAPI, input *s3.PutObjectInput) (*s3.PutObjectOutput, error) {
        return api.PutObject(c, input)
}

func PutFileS3(dir, filename, bucket, reg string) error {
        var cfg aws.Config
        st, err := fthash.Filehash(dir + filename)
        if err != nil {
                panic("configuration error, " + err.Error())
                return err
        }
        m := make(map[string]string)
        m["hashcode"] = st
        cfg, err = config.LoadDefaultConfig(context.TODO(), config.WithRegion(reg))
        if err != nil {
                panic("configuration error, " + err.Error())
        }

        client := s3.NewFromConfig(cfg)
        tmp := "backup" + dir + filename
        uri := strings.Replace(tmp, " ", "##,##", -1)
        if checkFileOnS3(client, bucket, uri, st) {
                fmt.Println(" FILE EXIST")
                return nil

        }
        file, err2 := os.Open(dir + filename)
        defer file.Close()

        if err2 != nil {
                fmt.Println("Unable to open file " + filename)
                return err2
        }

        tmp = "backup" + dir + filename
        //uri := "backup" + dir + filename
        uri = strings.Replace(tmp, " ", "##,##", -1)
        input := &s3.PutObjectInput{
                Bucket: &bucket,
                Key:    aws.String(uri),
                //Key:    &filename,
                Body:     file,
                Metadata: m,
        }
        ctx, cancelFn := context.WithTimeout(context.TODO(), 10*time.Second)
        defer cancelFn()
        _, err2 = PutFile(ctx, client, input)
        if err2 != nil {
                fmt.Println("Got error uploading file:", dir+filename)
                fmt.Println(err2)
                return err2
        }

        return nil
}
2 Answers

You've added a 10 second timeout here:

        ctx, cancelFn := context.WithTimeout(context.TODO(), 10*time.Second)
        defer cancelFn()
        _, err2 = PutFile(ctx, client, input)
        if err2 != nil {
                fmt.Println("Got error uploading file:", dir+filename)
                fmt.Println(err2)
                return err2
        }

After 10 seconds, the call to PutFile will exit with a context error. You likely just need to increase the timeout if you have files that take longer to upload.

package main

import (
    // "html/template"
    "log"
    "net/http"
    "os"

    "github.com/gin-gonic/gin"
    "github.com/joho/godotenv"

    "github.com/aws/aws-sdk-go/aws"
    "github.com/aws/aws-sdk-go/aws/credentials"
    "github.com/aws/aws-sdk-go/aws/session"
    "github.com/aws/aws-sdk-go/service/s3/s3manager"
)

var AccessKeyID string
var SecretAccessKey string
var MyRegion string
var MyBucket string
var filepath string

//GetEnvWithKey : get env value
func GetEnvWithKey(key string) string {
    return os.Getenv(key)
}

func LoadEnv() {
    err := godotenv.Load(".env")
    if err != nil {
        log.Fatalf("Error loading .env file")
        os.Exit(1)
    }
}

func ConnectAws() *session.Session {
    AccessKeyID = GetEnvWithKey("AWS_ACCESS_KEY_ID")
    SecretAccessKey = GetEnvWithKey("AWS_SECRET_ACCESS_KEY")
    MyRegion = GetEnvWithKey("AWS_REGION")

    sess, err := session.NewSession(
        &aws.Config{
            Region: aws.String(MyRegion),
            Credentials: credentials.NewStaticCredentials(
                AccessKeyID,
                SecretAccessKey,
                "", // a token will be created when the session it's used.
            ),
        })

    if err != nil {
        panic(err)
    }

    return sess
}

func SetupRouter(sess *session.Session) {
    router := gin.Default()

    router.Use(func(c *gin.Context) {
        c.Set("sess", sess)
        c.Next()
    })

    // router.Get("/upload", Form)
    router.POST("/upload", UploadImage)
    // router.GET("/image", controllers.DisplayImage)

    _ = router.Run(":4000")
}

func UploadImage(c *gin.Context) {
    sess := c.MustGet("sess").(*session.Session)
    uploader := s3manager.NewUploader(sess)

    MyBucket = GetEnvWithKey("BUCKET_NAME")

    file, header, err := c.Request.FormFile("photo")
    filename := header.Filename

    //upload to the s3 bucket
    up, err := uploader.Upload(&s3manager.UploadInput{
        Bucket: aws.String(MyBucket),
        //ACL:    aws.String("public-read"),
        Key:  aws.String(filename),
        Body: file,
    })

    if err != nil {
        c.JSON(http.StatusInternalServerError, gin.H{
            "error":    "Failed to upload file",
            "uploader": up,
        })
        return
    }
    filepath = "https://" + MyBucket + "." + "s3-" + MyRegion + ".amazonaws.com/" + filename
    c.JSON(http.StatusOK, gin.H{
        "filepath": filepath,
    })
}

func main() {
    LoadEnv()

    sess := ConnectAws()
    router := gin.Default()
    router.Use(func(c *gin.Context) {
        c.Set("sess", sess)
        c.Next()
    })

    router.POST("/upload", UploadImage)

    //router.LoadHTMLGlob("templates/*")
    //router.GET("/image", func(c *gin.Context) {
    //c.HTML(http.StatusOK, "index.tmpl", gin.H{
    //  "title": "Main website",
    //})
    //})

    _ = router.Run(":4000")
}
Related