From d59d95bde4b706d7de02c0ded3e0d2956094fce7 Mon Sep 17 00:00:00 2001 From: Eric Date: Thu, 10 Sep 2026 18:26:18 +0800 Subject: [PATCH] feat: upgrade Go dependencies and optimize S3 uploads --- common/sys_status.go | 8 +- go.mod | 31 +++++-- go.sum | 88 +++++++++--------- storage/es.go | 4 +- storage/es_test.go | 40 +++++++++ storage/s3.go | 144 ++++++++++++++++++++++++++---- storage/s3_test.go | 206 +++++++++++++++++++++++++++++++++++++++++++ 7 files changed, 450 insertions(+), 71 deletions(-) create mode 100644 storage/s3_test.go diff --git a/common/sys_status.go b/common/sys_status.go index 1d1bef7..ec4cdf0 100644 --- a/common/sys_status.go +++ b/common/sys_status.go @@ -6,10 +6,10 @@ import ( "os" "strconv" - "github.com/shirou/gopsutil/v3/cpu" - "github.com/shirou/gopsutil/v3/disk" - "github.com/shirou/gopsutil/v3/load" - "github.com/shirou/gopsutil/v3/mem" + "github.com/shirou/gopsutil/v4/cpu" + "github.com/shirou/gopsutil/v4/disk" + "github.com/shirou/gopsutil/v4/load" + "github.com/shirou/gopsutil/v4/mem" ) func CpuLoad1Usage() float64 { diff --git a/go.mod b/go.mod index 87a24a4..c9252d2 100644 --- a/go.mod +++ b/go.mod @@ -1,40 +1,56 @@ module github.com/jumpserver-dev/sdk-go -go 1.24.0 +go 1.26.0 require ( github.com/Azure/azure-storage-blob-go v0.15.0 github.com/LeeEirc/httpsig v1.2.1 github.com/aliyun/aliyun-oss-go-sdk v3.0.2+incompatible - github.com/aws/aws-sdk-go v1.44.306 + github.com/aws/aws-sdk-go-v2 v1.47.0 + github.com/aws/aws-sdk-go-v2/config v1.33.4 + github.com/aws/aws-sdk-go-v2/credentials v1.20.4 + github.com/aws/aws-sdk-go-v2/service/s3 v1.113.0 github.com/dlclark/regexp2 v1.11.4 - github.com/elastic/go-elasticsearch/v6 v6.8.5 + github.com/elastic/go-elasticsearch/v7 v7.17.10 github.com/elastic/go-elasticsearch/v8 v8.14.0 github.com/google/uuid v1.6.0 github.com/gorilla/websocket v1.5.3 github.com/huaweicloud/huaweicloud-sdk-go-obs v3.25.4+incompatible github.com/influxdata/influxdb-client-go/v2 v2.14.0 - github.com/shirou/gopsutil/v3 v3.24.5 + github.com/shirou/gopsutil/v4 v4.26.8 golang.org/x/text v0.33.0 ) require ( github.com/Azure/azure-pipeline-go v0.2.3 // indirect github.com/apapsch/go-jsonmerge/v2 v2.0.0 // indirect + github.com/aws/aws-sdk-go-v2/aws/protocol/eventstream v1.7.20 // indirect + github.com/aws/aws-sdk-go-v2/feature/ec2/imds v1.20.0 // indirect + github.com/aws/aws-sdk-go-v2/internal/configsources v1.5.3 // indirect + github.com/aws/aws-sdk-go-v2/internal/endpoints/v2 v2.8.3 // indirect + github.com/aws/aws-sdk-go-v2/internal/v4a v1.5.3 // indirect + github.com/aws/aws-sdk-go-v2/service/internal/accept-encoding v1.13.19 // indirect + github.com/aws/aws-sdk-go-v2/service/internal/checksum v1.11.3 // indirect + github.com/aws/aws-sdk-go-v2/service/internal/presigned-url v1.14.3 // indirect + github.com/aws/aws-sdk-go-v2/service/internal/s3shared v1.20.3 // indirect + github.com/aws/aws-sdk-go-v2/service/signin v1.10.0 // indirect + github.com/aws/aws-sdk-go-v2/service/sso v1.38.0 // indirect + github.com/aws/aws-sdk-go-v2/service/ssooidc v1.43.0 // indirect + github.com/aws/aws-sdk-go-v2/service/sts v1.50.0 // indirect + github.com/aws/smithy-go v1.28.1 // indirect github.com/davecgh/go-spew v1.1.2-0.20180830191138-d8f796af33cc // indirect + github.com/ebitengine/purego v0.10.2 // indirect github.com/elastic/elastic-transport-go/v8 v8.6.0 // indirect github.com/go-logr/logr v1.4.1 // indirect github.com/go-logr/stdr v1.2.2 // indirect github.com/go-ole/go-ole v1.3.0 // indirect github.com/influxdata/line-protocol v0.0.0-20210922203350-b1ad95c89adf // indirect - github.com/jmespath/go-jmespath v0.4.0 // indirect github.com/kr/pretty v0.3.1 // indirect github.com/lufia/plan9stats v0.0.0-20251013123823-9fd1530e3ec3 // indirect github.com/mattn/go-ieproxy v0.0.1 // indirect github.com/oapi-codegen/runtime v1.1.1 // indirect github.com/pmezard/go-difflib v1.0.1-0.20181226105442-5d4384ee4fb2 // indirect github.com/power-devops/perfstat v0.0.0-20240221224432-82ca36839d55 // indirect - github.com/shoenig/go-m1cpu v0.1.7 // indirect github.com/tklauser/go-sysconf v0.3.16 // indirect github.com/tklauser/numcpus v0.11.0 // indirect github.com/yusufpapurcu/wmi v1.2.4 // indirect @@ -42,7 +58,6 @@ require ( go.opentelemetry.io/otel/metric v1.24.0 // indirect go.opentelemetry.io/otel/trace v1.24.0 // indirect golang.org/x/net v0.43.0 // indirect - golang.org/x/sys v0.40.0 // indirect + golang.org/x/sys v0.41.0 // indirect golang.org/x/time v0.12.0 // indirect - gopkg.in/yaml.v2 v2.4.0 // indirect ) diff --git a/go.sum b/go.sum index 5b6c7f1..93e1d02 100644 --- a/go.sum +++ b/go.sum @@ -20,8 +20,42 @@ github.com/aliyun/aliyun-oss-go-sdk v3.0.2+incompatible h1:8psS8a+wKfiLt1iVDX79F github.com/aliyun/aliyun-oss-go-sdk v3.0.2+incompatible/go.mod h1:T/Aws4fEfogEE9v+HPhhw+CntffsBHJ8nXQCwKr0/g8= github.com/apapsch/go-jsonmerge/v2 v2.0.0 h1:axGnT1gRIfimI7gJifB699GoE/oq+F2MU7Dml6nw9rQ= github.com/apapsch/go-jsonmerge/v2 v2.0.0/go.mod h1:lvDnEdqiQrp0O42VQGgmlKpxL1AP2+08jFMw88y4klk= -github.com/aws/aws-sdk-go v1.44.306 h1:H487V/1N09BDxeGR7oR+LloC2uUpmf4atmqJaBgQOIs= -github.com/aws/aws-sdk-go v1.44.306/go.mod h1:aVsgQcEevwlmQ7qHE9I3h+dtQgpqhFB+i8Phjh7fkwI= +github.com/aws/aws-sdk-go-v2 v1.47.0 h1:0jsHallhJCeaU0Ko48c/3FK1ctOQ7NpzggxriJOQ8MQ= +github.com/aws/aws-sdk-go-v2 v1.47.0/go.mod h1:bttEH6JqnUL8LepvDVfdrds/fZ5bCIxzpe3abyUrhDU= +github.com/aws/aws-sdk-go-v2/aws/protocol/eventstream v1.7.20 h1:GPRlPwz40I2B2VrBEASOA3Bi77NyeqejNLkifosX0rs= +github.com/aws/aws-sdk-go-v2/aws/protocol/eventstream v1.7.20/go.mod h1:g7PNzKcsOKWb4fkSRBA7BZVAS6Y8IcxzN+nRohhQ1Q8= +github.com/aws/aws-sdk-go-v2/config v1.33.4 h1:FzvkXKSzwqHni4U7nDigHg4jjtqMpVUuHgmZfSoJVQ0= +github.com/aws/aws-sdk-go-v2/config v1.33.4/go.mod h1:VZqGZnZsCWVfK/iGPptJIyNIX3XEX6iQU2Rel4sLrr8= +github.com/aws/aws-sdk-go-v2/credentials v1.20.4 h1:hTvrJJseKbvw32kmiE0G+u/9ZqpqscjDrTigHIXP2qs= +github.com/aws/aws-sdk-go-v2/credentials v1.20.4/go.mod h1:gWp9O1ZBWwpcIrgV+mVHk4gZUurAEDkgypu/OXOlIaw= +github.com/aws/aws-sdk-go-v2/feature/ec2/imds v1.20.0 h1:AM4hHjww+PSFtt6E+UrBrPlZkWsePCLEt9AjkfQX+yM= +github.com/aws/aws-sdk-go-v2/feature/ec2/imds v1.20.0/go.mod h1:3x/yXezeQjpOvBb4jEMxrS8SXvpdvJ5abv6l5c1gWM8= +github.com/aws/aws-sdk-go-v2/internal/configsources v1.5.3 h1:Hp/VgjP0BysR3OgLlR057Vz2LcbbVnoWeJ+3qWiS/fY= +github.com/aws/aws-sdk-go-v2/internal/configsources v1.5.3/go.mod h1:nwGV5qw7F1IZPgxCvA/ph8N2TAuz+BkRG/bXn808qMA= +github.com/aws/aws-sdk-go-v2/internal/endpoints/v2 v2.8.3 h1:MUaM4f+kj1ZIBPZfUS8cxP1GKXXZtHJjAthy93AN7SM= +github.com/aws/aws-sdk-go-v2/internal/endpoints/v2 v2.8.3/go.mod h1:6YmVmEVRI5ZZzRjCSsb9SryKH0hAlMRdgA7kG9aDvBU= +github.com/aws/aws-sdk-go-v2/internal/v4a v1.5.3 h1:fuSCw4Z2qfRCztMPO3GXJNSiEp6Wee+WOLwrHHUMy9c= +github.com/aws/aws-sdk-go-v2/internal/v4a v1.5.3/go.mod h1:6SxcHheD1pPR5+kWm1wGvjlL/YqUsh267sAfEmN4K7A= +github.com/aws/aws-sdk-go-v2/service/internal/accept-encoding v1.13.19 h1:bAdDl/HkGCcGPoe25ToSHEw23VIxt6CT5fLcg111BKg= +github.com/aws/aws-sdk-go-v2/service/internal/accept-encoding v1.13.19/go.mod h1:KaUzbLxv4CeSxh6ZCl9B4m7CuFenS8kUEaDs+f/DQr4= +github.com/aws/aws-sdk-go-v2/service/internal/checksum v1.11.3 h1:BHKCSX4QXERe8So8rbWqaM7owqOmDJxATXgJwGng22A= +github.com/aws/aws-sdk-go-v2/service/internal/checksum v1.11.3/go.mod h1:GqWeeKfYfezihA2KfFL9l7ohEdZWe1tuFWh3GfyNSnE= +github.com/aws/aws-sdk-go-v2/service/internal/presigned-url v1.14.3 h1:bON1rJf67TSTDCKg816AAIE4xSTtoo9tl0XRkO72R+I= +github.com/aws/aws-sdk-go-v2/service/internal/presigned-url v1.14.3/go.mod h1:c5BBpjJcQXpfeq9iASyVKA3T6vX6B6LEXY4mL/gklDY= +github.com/aws/aws-sdk-go-v2/service/internal/s3shared v1.20.3 h1:L8vIOxylma91TcR96NFTEC07G3JDwSl+CvK2b+IODms= +github.com/aws/aws-sdk-go-v2/service/internal/s3shared v1.20.3/go.mod h1:fmPIZQzTExYuBNWFyi1P7IoDjvskgphXqK1yObMzusM= +github.com/aws/aws-sdk-go-v2/service/s3 v1.113.0 h1:0slwIjBv1sEigSUW4EfTtNw9mbHl0PlOUPX9Y1C/eLE= +github.com/aws/aws-sdk-go-v2/service/s3 v1.113.0/go.mod h1:/uA+2Qj4jd5qBWagVC1AyzzDFXVK997E7U04w7Kw0wI= +github.com/aws/aws-sdk-go-v2/service/signin v1.10.0 h1:ZD5qFpWcaOKdTuhBi431pIDkCgrMkMlMT6jlpSPoIRI= +github.com/aws/aws-sdk-go-v2/service/signin v1.10.0/go.mod h1:8Nuuf+tR346PjJ3MvZPh9pekbLiLQFWJhzMXfwy7alA= +github.com/aws/aws-sdk-go-v2/service/sso v1.38.0 h1:JGeeBcMlhg1xtOXYpeCaTQBZObtXMPQCUqBcmr65NRA= +github.com/aws/aws-sdk-go-v2/service/sso v1.38.0/go.mod h1:XwteswG9EOMRFm73UT0t+MbTwyLxMrEXkU6e+v92Lzo= +github.com/aws/aws-sdk-go-v2/service/ssooidc v1.43.0 h1:obhahQXDEdVEv8y5bTKXR30LVaxYe1kyYM0L7l2Iq+k= +github.com/aws/aws-sdk-go-v2/service/ssooidc v1.43.0/go.mod h1:6twZZ/aXHNy1vXUO8koUbp++MYzMASkOgEBdkbJYmO0= +github.com/aws/aws-sdk-go-v2/service/sts v1.50.0 h1:khXV3+K5D3f4e8xtplaRdSFn1bEg3gj5EBHQvbCOZbQ= +github.com/aws/aws-sdk-go-v2/service/sts v1.50.0/go.mod h1:/8JRcdTt//hG0Q4BTmGbuOplT7ABe+5rdtqUHqXvYIM= +github.com/aws/smithy-go v1.28.1 h1:R/nXH00c8qcfCzQVELtRw+eLQWtzv+VAIEFJ1/xxXlQ= +github.com/aws/smithy-go v1.28.1/go.mod h1:YE2RhdIuDbA5E5bTdciG9KrW3+TiEONeUWCqxX9i1Fc= github.com/bmatcuk/doublestar v1.1.1/go.mod h1:UD6OnuiIn0yFxxA2le/rnRU1G4RaI4UvFv1sNto9p6w= github.com/creack/pty v1.1.9/go.mod h1:oKZEueFk5CKHvIhNR5MUki03XCEU+Q6VDXinZuGJ33E= github.com/davecgh/go-spew v1.1.0/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= @@ -29,10 +63,12 @@ github.com/davecgh/go-spew v1.1.2-0.20180830191138-d8f796af33cc h1:U9qPSI2PIWSS1 github.com/davecgh/go-spew v1.1.2-0.20180830191138-d8f796af33cc/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= github.com/dlclark/regexp2 v1.11.4 h1:rPYF9/LECdNymJufQKmri9gV604RvvABwgOA8un7yAo= github.com/dlclark/regexp2 v1.11.4/go.mod h1:DHkYz0B9wPfa6wondMfaivmHpzrQ3v9q8cnmRbL6yW8= +github.com/ebitengine/purego v0.10.2 h1:W809HbnvzAxgdm+aOvlSekrM16wGCdT/e76+9tS7gzE= +github.com/ebitengine/purego v0.10.2/go.mod h1:iIjxzd6CiRiOG0UyXP+V1+jWqUXVjPKLAI0mRfJZTmQ= github.com/elastic/elastic-transport-go/v8 v8.6.0 h1:Y2S/FBjx1LlCv5m6pWAF2kDJAHoSjSRSJCApolgfthA= github.com/elastic/elastic-transport-go/v8 v8.6.0/go.mod h1:YLHer5cj0csTzNFXoNQ8qhtGY1GTvSqPnKWKaqQE3Hk= -github.com/elastic/go-elasticsearch/v6 v6.8.5 h1:U2HtkBseC1FNBmDr0TR2tKltL6FxoY+niDAlj5M8TK8= -github.com/elastic/go-elasticsearch/v6 v6.8.5/go.mod h1:UwaDJsD3rWLM5rKNFzv9hgox93HoX8utj1kxD9aFUcI= +github.com/elastic/go-elasticsearch/v7 v7.17.10 h1:TCQ8i4PmIJuBunvBS6bwT2ybzVFxxUhhltAs3Gyu1yo= +github.com/elastic/go-elasticsearch/v7 v7.17.10/go.mod h1:OJ4wdbtDNk5g503kvlHLyErCgQwwzmDtaFC4XyOxXA4= github.com/elastic/go-elasticsearch/v8 v8.14.0 h1:1ywU8WFReLLcxE1WJqii3hTtbPUE2hc38ZK/j4mMFow= github.com/elastic/go-elasticsearch/v8 v8.14.0/go.mod h1:WRvnlGkSuZyp83M2U8El/LGXpCjYLrvlkSgkAH4O5I4= github.com/form3tech-oss/jwt-go v3.2.2+incompatible h1:TcekIExNqud5crz4xD2pavyTgWiPvpYe4Xau31I0PRk= @@ -58,10 +94,6 @@ github.com/influxdata/influxdb-client-go/v2 v2.14.0 h1:AjbBfJuq+QoaXNcrova8smSjw github.com/influxdata/influxdb-client-go/v2 v2.14.0/go.mod h1:Ahpm3QXKMJslpXl3IftVLVezreAUtBOTZssDrjZEFHI= github.com/influxdata/line-protocol v0.0.0-20210922203350-b1ad95c89adf h1:7JTmneyiNEwVBOHSjoMxiWAqB992atOeepeFYegn5RU= github.com/influxdata/line-protocol v0.0.0-20210922203350-b1ad95c89adf/go.mod h1:xaLFMmpvUxqXtVkUJfg9QmT88cDaCJ3ZKgdZ78oO8Qo= -github.com/jmespath/go-jmespath v0.4.0 h1:BEgLn5cpjn8UN1mAw4NjwDrS35OdebyEtFe+9YPoQUg= -github.com/jmespath/go-jmespath v0.4.0/go.mod h1:T8mJZnbsbmF+m6zOOFylbeCJqk5+pHWvzYPziyZiYoo= -github.com/jmespath/go-jmespath/internal/testify v1.5.1 h1:shLQSRRSCCPj3f2gpwzGwWFoC7ycTf1rcQZHOlsJ6N8= -github.com/jmespath/go-jmespath/internal/testify v1.5.1/go.mod h1:L3OGu8Wl2/fWfCI6z80xFu9LTZmf1ZRjMHUOPmWr69U= github.com/juju/gnuflag v0.0.0-20171113085948-2ce1bb71843d/go.mod h1:2PavIy+JPciBPrBUjwbNvtwB6RQlve+hkpll6QSNmOE= github.com/kr/pretty v0.2.1/go.mod h1:ipq/a2n7PKx3OHsz4KJII5eveXtPO4qwEXGdVfWzfnI= github.com/kr/pretty v0.3.1 h1:flRD4NNwYAUpkphVc1HcthR4KEIFJ65n8Mw5qdRn3LE= @@ -86,22 +118,17 @@ github.com/power-devops/perfstat v0.0.0-20240221224432-82ca36839d55 h1:o4JXh1EVt github.com/power-devops/perfstat v0.0.0-20240221224432-82ca36839d55/go.mod h1:OmDBASR4679mdNQnz2pUhc2G8CO2JrUAVFDRBDP/hJE= github.com/rogpeppe/go-internal v1.9.0 h1:73kH8U+JUqXU8lRuOHeVHaa/SZPifC7BkcraZVejAe8= github.com/rogpeppe/go-internal v1.9.0/go.mod h1:WtVeX8xhTBvf0smdhujwtBcq4Qrzq/fJaraNFVN+nFs= -github.com/shirou/gopsutil/v3 v3.24.5 h1:i0t8kL+kQTvpAYToeuiVk3TgDeKOFioZO3Ztz/iZ9pI= -github.com/shirou/gopsutil/v3 v3.24.5/go.mod h1:bsoOS1aStSs9ErQ1WWfxllSeS1K5D+U30r2NfcubMVk= -github.com/shoenig/go-m1cpu v0.1.7 h1:C76Yd0ObKR82W4vhfjZiCp0HxcSZ8Nqd84v+HZ0qyI0= -github.com/shoenig/go-m1cpu v0.1.7/go.mod h1:KkDOw6m3ZJQAPHbrzkZki4hnx+pDRR1Lo+ldA56wD5w= -github.com/shoenig/test v1.7.0 h1:eWcHtTXa6QLnBvm0jgEabMRN/uJ4DMV3M8xUGgRkZmk= -github.com/shoenig/test v1.7.0/go.mod h1:UxJ6u/x2v/TNs/LoLxBNJRV9DiwBBKYxXSyczsBHFoI= +github.com/shirou/gopsutil/v4 v4.26.8 h1:YQMTF/1J50B5+Y0vlo1eDRf5DoR7Gk69hY+8wjYkQeo= +github.com/shirou/gopsutil/v4 v4.26.8/go.mod h1:5O9FjBiXoTDFatIWjZZosqj4pV0DRtLx598xGbBehzM= github.com/spkg/bom v0.0.0-20160624110644-59b7046e48ad/go.mod h1:qLr4V1qq6nMqFKkMo8ZTx3f+BZEkzsRUY10Xsm2mwU0= github.com/stretchr/objx v0.1.0/go.mod h1:HFkY916IF+rwdDfMAkV7OtwuqBVzrE8GR6GFx+wExME= github.com/stretchr/testify v1.3.0/go.mod h1:M5WIy9Dh21IEIfnGCwXGc5bZfKNJtfHm1UVUgZn+9EI= -github.com/stretchr/testify v1.9.0 h1:HtqpIVDClZ4nwg75+f6Lvsy/wHu+3BoSGCbBAcpTsTg= -github.com/stretchr/testify v1.9.0/go.mod h1:r2ic/lqez/lEtzL7wO/rwa5dbSLXVDPFyf8C91i36aY= +github.com/stretchr/testify v1.11.1 h1:7s2iGBzp5EwR7/aIZr8ao5+dra3wiQyKjjFuvgVKu7U= +github.com/stretchr/testify v1.11.1/go.mod h1:wZwfW3scLgRK+23gO65QZefKpKQRnfz6sD981Nm4B6U= github.com/tklauser/go-sysconf v0.3.16 h1:frioLaCQSsF5Cy1jgRBrzr6t502KIIwQ0MArYICU0nA= github.com/tklauser/go-sysconf v0.3.16/go.mod h1:/qNL9xxDhc7tx3HSRsLWNnuzbVfh3e7gh/BmM179nYI= github.com/tklauser/numcpus v0.11.0 h1:nSTwhKH5e1dMNsCdVBukSZrURJRoHbSEQjdEbY+9RXw= github.com/tklauser/numcpus v0.11.0/go.mod h1:z+LwcLq54uWZTX0u/bGobaV34u6V7KNlTZejzM6/3MQ= -github.com/yuin/goldmark v1.4.13/go.mod h1:6yULJ656Px+3vBD8DxQVa3kxgyrAnzto9xy5taEt/CY= github.com/yusufpapurcu/wmi v1.2.4 h1:zFUKzehAFReQwLys1b/iSMl+JQGSCSjtVqQn9bBrPo0= github.com/yusufpapurcu/wmi v1.2.4/go.mod h1:SBZ9tNy3G9/m5Oi98Zks0QjeHVDvuK0qfxQmPyzfmi0= go.opentelemetry.io/otel v1.24.0 h1:0LAOdjNmQeSTzGBzduGe/rU4tZhMwL5rWgtp9Ku5Jfo= @@ -115,21 +142,13 @@ go.opentelemetry.io/otel/trace v1.24.0/go.mod h1:HPc3Xr/cOApsBI154IU0OI0HJexz+aw golang.org/x/crypto v0.0.0-20190308221718-c2843e01d9a2/go.mod h1:djNgcEr1/C05ACkg1iLfiJU5Ep61QUkGW8qpdssI0+w= golang.org/x/crypto v0.0.0-20201002170205-7f63de1d35b0/go.mod h1:LzIPMQfyMNhhGPhUkYOs5KpL4U8rLKemX1yGLhDgUto= golang.org/x/crypto v0.0.0-20201016220609-9e8e0b390897/go.mod h1:LzIPMQfyMNhhGPhUkYOs5KpL4U8rLKemX1yGLhDgUto= -golang.org/x/crypto v0.0.0-20210921155107-089bfa567519/go.mod h1:GvvjBRRGRdwPK5ydBHafDWAxML/pGHZbMvKqRZ5+Abc= golang.org/x/crypto v0.41.0 h1:WKYxWedPGCTVVl5+WHSSrOBT0O8lx32+zxmHxijgXp4= golang.org/x/crypto v0.41.0/go.mod h1:pO5AFd7FA68rFak7rOAGVuygIISepHftHnr8dr6+sUc= -golang.org/x/mod v0.6.0-dev.0.20220419223038-86c51ed26bb4/go.mod h1:jJ57K6gSWd91VN4djpZkiMVwK6gcyfeH4XE8wZrZaV4= golang.org/x/net v0.0.0-20190404232315-eb5bcb51f2a3/go.mod h1:t9HGtf8HONx5eT2rtn7q6eTqICYqUVnKs3thJo3Qplg= -golang.org/x/net v0.0.0-20190620200207-3b0461eec859/go.mod h1:z5CRVTTTmAJ677TzLLGU+0bjPO0LkuOLi4/5GtJWs/s= golang.org/x/net v0.0.0-20191112182307-2180aed22343/go.mod h1:z5CRVTTTmAJ677TzLLGU+0bjPO0LkuOLi4/5GtJWs/s= -golang.org/x/net v0.0.0-20210226172049-e18ecbb05110/go.mod h1:m0MpNAwzfU5UDzcl9v0D8zg8gWTRqZa9RBIspLL5mdg= golang.org/x/net v0.0.0-20210610132358-84b48f89b13b/go.mod h1:9nx3DQGgdP8bBQD5qxJ1jj9UTztislL4KSBs9R2vV5Y= -golang.org/x/net v0.0.0-20220722155237-a158d28d115b/go.mod h1:XRhObCWvk6IyKnWLug+ECip1KBveYUHfp+8e9klMJ9c= -golang.org/x/net v0.1.0/go.mod h1:Cx3nUiGt4eDBEyega/BKRp+/AlGL8hYe7U9odMt2Cco= golang.org/x/net v0.43.0 h1:lat02VYK2j4aLzMzecihNvTlJNQUq316m2Mr9rnM6YE= golang.org/x/net v0.43.0/go.mod h1:vhO1fvI4dGsIjh73sWfUVjj3N7CA9WkKJNQm2svM6Jg= -golang.org/x/sync v0.0.0-20190423024810-112230192c58/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM= -golang.org/x/sync v0.0.0-20220722155255-886fb9371eb4/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM= golang.org/x/sys v0.0.0-20190215142949-d0b11bdaac8a/go.mod h1:STP8DvDyc/dI5b8T5hshtkjS+E42TnysNCUPdjciGhY= golang.org/x/sys v0.0.0-20190412213103-97732733099d/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= golang.org/x/sys v0.0.0-20190916202348-b4ddaad3f8a3/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= @@ -137,34 +156,19 @@ golang.org/x/sys v0.0.0-20191112214154-59a1497f0cea/go.mod h1:h1NjWce9XRLGQEsW7w golang.org/x/sys v0.0.0-20201119102817-f84b799fce68/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= golang.org/x/sys v0.0.0-20201204225414-ed752295db88/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= golang.org/x/sys v0.0.0-20210423082822-04245dca01da/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= -golang.org/x/sys v0.0.0-20210615035016-665e8c7367d1/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= -golang.org/x/sys v0.0.0-20220520151302-bc2c85ada10a/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= -golang.org/x/sys v0.0.0-20220722155257-8c9f86f7a55f/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= golang.org/x/sys v0.1.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= -golang.org/x/sys v0.40.0 h1:DBZZqJ2Rkml6QMQsZywtnjnnGvHza6BTfYFWY9kjEWQ= -golang.org/x/sys v0.40.0/go.mod h1:OgkHotnGiDImocRcuBABYBEXf8A9a87e/uXjp9XT3ks= +golang.org/x/sys v0.41.0 h1:Ivj+2Cp/ylzLiEU89QhWblYnOE9zerudt9Ftecq2C6k= +golang.org/x/sys v0.41.0/go.mod h1:OgkHotnGiDImocRcuBABYBEXf8A9a87e/uXjp9XT3ks= golang.org/x/term v0.0.0-20201126162022-7de9c90e9dd1/go.mod h1:bj7SfCRtBDWHUb9snDiAeCFNEtKQo2Wmx5Cou7ajbmo= -golang.org/x/term v0.0.0-20210927222741-03fcf44c2211/go.mod h1:jbD1KX2456YbFQfuXm/mYQcufACuNUgVhRMnK/tPxf8= -golang.org/x/term v0.1.0/go.mod h1:jbD1KX2456YbFQfuXm/mYQcufACuNUgVhRMnK/tPxf8= golang.org/x/text v0.3.0/go.mod h1:NqM8EUOU14njkJ3fqMW+pc6Ldnwhi/IjpwHt7yyuwOQ= golang.org/x/text v0.3.2/go.mod h1:bEr9sfX3Q8Zfm5fL9x+3itogRgK3+ptLWKqgva+5dAk= -golang.org/x/text v0.3.3/go.mod h1:5Zoc/QRtKVWzQhOtBMvqHzDpF6irO9z98xDceosuGiQ= golang.org/x/text v0.3.6/go.mod h1:5Zoc/QRtKVWzQhOtBMvqHzDpF6irO9z98xDceosuGiQ= -golang.org/x/text v0.3.7/go.mod h1:u+2+/6zg+i71rQMx5EYifcz6MCKuco9NR6JIITiCfzQ= -golang.org/x/text v0.4.0/go.mod h1:mrYo+phRRbMaCq/xk9113O4dZlRixOauAjOtrjsXDZ8= golang.org/x/text v0.33.0 h1:B3njUFyqtHDUI5jMn1YIr5B0IE2U0qck04r6d4KPAxE= golang.org/x/text v0.33.0/go.mod h1:LuMebE6+rBincTi9+xWTY8TztLzKHc/9C1uBCG27+q8= golang.org/x/time v0.12.0 h1:ScB/8o8olJvc+CQPWrK3fPZNfh7qgwCrY0zJmoEQLSE= golang.org/x/time v0.12.0/go.mod h1:CDIdPxbZBQxdj6cxyCIdrNogrJKMJ7pr37NYpMcMDSg= golang.org/x/tools v0.0.0-20180917221912-90fa682c2a6e/go.mod h1:n7NCudcB/nEzxVGmLbDWY5pfWTLqBcC2KZ6jyYvM4mQ= -golang.org/x/tools v0.0.0-20191119224855-298f0cb1881e/go.mod h1:b+2E5dAYhXwXZwtnZ6UAqBI28+e2cm9otk0dWdXHAEo= -golang.org/x/tools v0.1.12/go.mod h1:hNGJHUnrk76NpqgfD5Aqm5Crs+Hm0VOH/i9J2+nxYbc= -golang.org/x/xerrors v0.0.0-20190717185122-a985d3407aa7/go.mod h1:I/5z698sn9Ka8TeJc9MKroUUfqBBauWjQqLJ2OPfmY0= -gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0= gopkg.in/check.v1 v1.0.0-20201130134442-10cb98267c6c h1:Hei/4ADfdWqJk1ZMxUNpqntNwaWcugrBjAiHlqqRiVk= gopkg.in/check.v1 v1.0.0-20201130134442-10cb98267c6c/go.mod h1:JHkPIbrfpd72SG/EVd6muEfDQjcINNoR0C8j2r3qZ4Q= -gopkg.in/yaml.v2 v2.2.8/go.mod h1:hI93XBmqTisBFMUTm0b8Fm+jr3Dg1NNxqwp+5A1VGuI= -gopkg.in/yaml.v2 v2.4.0 h1:D8xgwECY7CYvx+Y2n4sBz93Jn9JRvxdiyyo8CTfuKaY= -gopkg.in/yaml.v2 v2.4.0/go.mod h1:RDklbk79AGWmwhnvt/jBztapEOGDOx6ZbXqjP6csGnQ= gopkg.in/yaml.v3 v3.0.1 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA= gopkg.in/yaml.v3 v3.0.1/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= diff --git a/storage/es.go b/storage/es.go index 3de43ea..cc0f7ca 100644 --- a/storage/es.go +++ b/storage/es.go @@ -9,8 +9,8 @@ import ( "io" "net/http" - "github.com/elastic/go-elasticsearch/v6" - "github.com/elastic/go-elasticsearch/v6/esapi" + "github.com/elastic/go-elasticsearch/v7" + "github.com/elastic/go-elasticsearch/v7/esapi" elasticsearch8 "github.com/elastic/go-elasticsearch/v8" "github.com/jumpserver-dev/sdk-go/logger" diff --git a/storage/es_test.go b/storage/es_test.go index 043f7f8..ca82e61 100644 --- a/storage/es_test.go +++ b/storage/es_test.go @@ -3,9 +3,49 @@ package storage import ( "bytes" "encoding/json" + "net/http" + "net/http/httptest" + "sync/atomic" "testing" + + "github.com/jumpserver-dev/sdk-go/model" ) +func TestBulkSaveEs(t *testing.T) { + var bulkRequests atomic.Int32 + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + w.Header().Set("Content-Type", "application/json") + w.Header().Set("X-Elastic-Product", "Elasticsearch") + if r.Method == http.MethodGet && r.URL.Path == "/" { + _, _ = w.Write([]byte(`{"version":{"number":"7.17.10"}}`)) + return + } + bulkRequests.Add(1) + if r.Method != http.MethodPost || r.URL.Path != "/commands/_doc/_bulk" { + t.Errorf("unexpected bulk request: %s %s", r.Method, r.URL.Path) + } + decoder := json.NewDecoder(r.Body) + var action map[string]json.RawMessage + if err := decoder.Decode(&action); err != nil || action["index"] == nil { + t.Errorf("unexpected bulk action: %v, error: %v", action, err) + } + var command model.Command + if err := decoder.Decode(&command); err != nil || command.Input != "whoami" { + t.Errorf("unexpected command: %q, error: %v", command.Input, err) + } + _, _ = w.Write([]byte(`{"errors":false,"items":[{"index":{"status":201}}]}`)) + })) + defer server.Close() + + storage := ESCommandStorage{Hosts: []string{server.URL}, Index: "commands", DocType: "_doc"} + if err := storage.BulkSaveEs([]*model.Command{{Input: "whoami"}}); err != nil { + t.Fatal(err) + } + if bulkRequests.Load() != 1 { + t.Fatalf("expected one bulk request, got %d", bulkRequests.Load()) + } +} + func TestEsIndexResponse(t *testing.T) { respBodys := [][2]string{ {"index", `{"took":24,"errors":false,"items":[{"index":{"_index":"jumpserver-test-1","_type":"_doc","_id":"mo9R9IkBIDTIizd_N0BL","_version":1,"result":"created","_shards":{"total":1,"successful":1,"failed":0},"_seq_no":3,"_primary_term":1,"status":201}},{"index":{"_index":"jumpserver-test-1","_type":"_doc","_id":"m49R9IkBIDTIizd_N0BL","_version":1,"result":"created","_shards":{"total":1,"successful":1,"failed":0},"_seq_no":4,"_primary_term":1,"status":201}}]}`}, diff --git a/storage/s3.go b/storage/s3.go index 6962c6e..850e4e7 100644 --- a/storage/s3.go +++ b/storage/s3.go @@ -1,14 +1,24 @@ package storage import ( + "context" + "errors" + "fmt" + "io" "os" + "strings" + "sync" + "time" - "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" + "github.com/aws/aws-sdk-go-v2/aws" + "github.com/aws/aws-sdk-go-v2/config" + "github.com/aws/aws-sdk-go-v2/credentials" + "github.com/aws/aws-sdk-go-v2/service/s3" + "github.com/aws/aws-sdk-go-v2/service/s3/types" ) +const s3UploadPartSize int64 = 64 * 1024 * 1024 + type S3ReplayStorage struct { Bucket string Region string @@ -24,33 +34,137 @@ func (s S3ReplayStorage) Upload(gZipFilePath, target string) (err error) { return err } defer file.Close() - s3Config := &aws.Config{ - Endpoint: aws.String(s.Endpoint), - Region: aws.String(s.Region), - S3ForcePathStyle: aws.Bool(true), + info, err := file.Stat() + if err != nil { + return err + } + if !info.Mode().IsRegular() { + return fmt.Errorf("S3 upload requires a regular file: %s", gZipFilePath) + } + ctx := context.Background() + opts := []func(*config.LoadOptions) error{ + config.WithRegion(s.Region), } if s.AccessKey != "" && s.SecretKey != "" { - s3Config.Credentials = credentials.NewStaticCredentials(s.AccessKey, s.SecretKey, "") + opts = append(opts, config.WithCredentialsProvider( + credentials.NewStaticCredentialsProvider(s.AccessKey, s.SecretKey, ""))) } - sess, err := session.NewSession(s3Config) + cfg, err := config.LoadDefaultConfig(ctx, opts...) if err != nil { return err } - uploader := s3manager.NewUploader(sess, func(u *s3manager.Uploader) { - u.PartSize = 64 * 1024 * 1024 // 64MB per part + client := s3.NewFromConfig(cfg, func(o *s3.Options) { + o.UsePathStyle = true + // Preserve compatibility with S3 endpoints without optional checksum support. + o.RequestChecksumCalculation = aws.RequestChecksumCalculationWhenRequired + if s.Endpoint != "" { + endpoint := s.Endpoint + if !strings.Contains(endpoint, "://") { + endpoint = "https://" + endpoint + } + o.BaseEndpoint = aws.String(endpoint) + } }) - _, err = uploader.Upload(&s3manager.UploadInput{ + if info.Size() <= s3UploadPartSize { + _, err = client.PutObject(ctx, &s3.PutObjectInput{ + Bucket: aws.String(s.Bucket), + Key: aws.String(target), + Body: io.NewSectionReader(file, 0, info.Size()), + ContentLength: aws.Int64(info.Size()), + }) + return err + } + return s.uploadMultipart(ctx, client, file, target, info.Size()) +} + +func (s S3ReplayStorage) uploadMultipart(ctx context.Context, client *s3.Client, file *os.File, target string, size int64) (err error) { + // Grow parts for large files without exceeding S3's 10,000-part limit. + partSize := max(s3UploadPartSize, (size-1)/10000+1) + if partSize > 5*1024*1024*1024 { + return errors.New("file exceeds the S3 multipart upload size limit") + } + partCount := int((size-1)/partSize + 1) + upload, err := client.CreateMultipartUpload(ctx, &s3.CreateMultipartUploadInput{ Bucket: aws.String(s.Bucket), Key: aws.String(target), - Body: file, }) if err != nil { return err } + if aws.ToString(upload.UploadId) == "" { + return errors.New("S3 multipart upload returned an empty upload ID") + } + defer func() { + if err == nil { + return + } + abortCtx, cancel := context.WithTimeout(context.Background(), 30*time.Second) + defer cancel() + _, abortErr := client.AbortMultipartUpload(abortCtx, &s3.AbortMultipartUploadInput{ + Bucket: aws.String(s.Bucket), Key: aws.String(target), UploadId: upload.UploadId, + }) + if abortErr != nil { + err = errors.Join(err, fmt.Errorf("abort multipart upload %s: %w", *upload.UploadId, abortErr)) + } + }() - return + uploadCtx, cancel := context.WithCancel(ctx) + defer cancel() + parts := make([]types.CompletedPart, partCount) + jobs := make(chan int) + var wg sync.WaitGroup + var once sync.Once + var firstErr error + for range min(5, partCount) { + wg.Go(func() { + for index := range jobs { + if uploadCtx.Err() != nil { + return + } + offset := int64(index) * partSize + length := min(partSize, size-offset) + part, uploadErr := client.UploadPart(uploadCtx, &s3.UploadPartInput{ + Bucket: aws.String(s.Bucket), + Key: aws.String(target), + UploadId: upload.UploadId, + PartNumber: aws.Int32(int32(index + 1)), + Body: io.NewSectionReader(file, offset, length), + ContentLength: aws.Int64(length), + }) + if uploadErr == nil && aws.ToString(part.ETag) == "" { + uploadErr = errors.New("S3 upload part returned an empty ETag") + } + if uploadErr != nil { + once.Do(func() { + firstErr = fmt.Errorf("upload part %d: %w", index+1, uploadErr) + cancel() + }) + return + } + parts[index] = types.CompletedPart{ETag: part.ETag, PartNumber: aws.Int32(int32(index + 1))} + } + }) + } +sendParts: + for index := range partCount { + select { + case jobs <- index: + case <-uploadCtx.Done(): + break sendParts + } + } + close(jobs) + wg.Wait() + if firstErr != nil { + return firstErr + } + _, err = client.CompleteMultipartUpload(ctx, &s3.CompleteMultipartUploadInput{ + Bucket: aws.String(s.Bucket), Key: aws.String(target), UploadId: upload.UploadId, + MultipartUpload: &types.CompletedMultipartUpload{Parts: parts}, + }) + return err } func (s S3ReplayStorage) TypeName() string { diff --git a/storage/s3_test.go b/storage/s3_test.go new file mode 100644 index 0000000..87952e8 --- /dev/null +++ b/storage/s3_test.go @@ -0,0 +1,206 @@ +package storage + +import ( + "encoding/xml" + "fmt" + "io" + "net/http" + "net/http/httptest" + "os" + "path/filepath" + "runtime" + "strconv" + "strings" + "sync/atomic" + "testing" +) + +func TestS3ReplayStorageUpload(t *testing.T) { + dir := setupS3TestEnv(t) + filename := filepath.Join(dir, "replay.gz") + const content = "recorded session" + if err := os.WriteFile(filename, []byte(content), 0600); err != nil { + t.Fatal(err) + } + + for _, explicit := range []bool{true, false} { + name, accessKey, token := "environment", "env-access", "env-token" + if explicit { + name, accessKey, token = "explicit", "config-access", "" + } + t.Run(name, func(t *testing.T) { + var requests atomic.Int32 + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + requests.Add(1) + if r.Method != http.MethodPut || r.URL.Path != "/replays/2026/session.gz" { + t.Errorf("unexpected upload: %s %s", r.Method, r.URL.Path) + } + auth := r.Header.Get("Authorization") + if !strings.Contains(auth, "Credential="+accessKey+"/") || !strings.Contains(auth, "/us-east-1/s3/aws4_request") { + t.Errorf("unexpected signing credentials or region: %s", auth) + } + if r.Header.Get("X-Amz-Security-Token") != token { + t.Error("unexpected session token") + } + if r.Header.Get("X-Amz-Sdk-Checksum-Algorithm") != "" || r.Header.Get("X-Amz-Trailer") != "" { + t.Error("upload requires optional S3 checksum support") + } + body, err := io.ReadAll(r.Body) + if err != nil || string(body) != content { + t.Errorf("unexpected upload body: %q, error: %v", body, err) + } + w.Header().Set("ETag", `"replay-etag"`) + })) + defer server.Close() + + storage := S3ReplayStorage{Bucket: "replays", Region: "us-east-1", Endpoint: server.URL} + if explicit { + storage.AccessKey, storage.SecretKey = accessKey, "config-secret" + } + if err := storage.Upload(filename, "2026/session.gz"); err != nil { + t.Fatal(err) + } + if requests.Load() != 1 { + t.Fatalf("expected one upload, got %d", requests.Load()) + } + }) + } +} + +func TestS3ReplayStorageMultipartUpload(t *testing.T) { + dir := setupS3TestEnv(t) + file, err := os.Create(filepath.Join(dir, "replay.gz")) + if err != nil { + t.Fatal(err) + } + defer file.Close() + const size = s3UploadPartSize + 1024*1024 + if err := file.Truncate(size); err != nil { + t.Fatal(err) + } + for index := range 2 { + if _, err := file.WriteAt([]byte(fmt.Sprintf("part-%d", index+1)), int64(index)*s3UploadPartSize); err != nil { + t.Fatal(err) + } + } + for _, failure := range []string{"", "part", "complete"} { + t.Run("failure="+failure, func(t *testing.T) { + var aborts, completes, firstPartAttempts atomic.Int32 + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + w.Header().Set("Content-Type", "application/xml") + if r.URL.Path != "/replays/session.gz" { + t.Errorf("unexpected upload path: %s", r.URL.Path) + } + if r.Header.Get("X-Amz-Sdk-Checksum-Algorithm") != "" || r.Header.Get("X-Amz-Trailer") != "" { + t.Error("multipart upload requires optional S3 checksum support") + } + query := r.URL.Query() + if !query.Has("uploads") && query.Get("uploadId") != "upload" { + t.Errorf("unexpected upload ID: %s", query.Get("uploadId")) + } + switch { + case r.Method == http.MethodPost && query.Has("uploads"): + _, _ = io.WriteString(w, `upload`) + case r.Method == http.MethodPut: + if failure == "part" { + w.WriteHeader(http.StatusForbidden) + _, _ = io.WriteString(w, `AccessDenied`) + return + } + number, err := strconv.Atoi(query.Get("partNumber")) + if err != nil || number < 1 || number > 2 { + t.Errorf("unexpected part number: %s", query.Get("partNumber")) + w.WriteHeader(http.StatusBadRequest) + return + } + prefix := make([]byte, len("part-1")) + if _, err := io.ReadFull(r.Body, prefix); err != nil || string(prefix) != fmt.Sprintf("part-%d", number) { + t.Errorf("incorrect part offset: %q, error: %v", prefix, err) + } + n, err := io.Copy(io.Discard, r.Body) + want := min(s3UploadPartSize, size-int64(number-1)*s3UploadPartSize) + if err != nil || n+int64(len(prefix)) != want || r.ContentLength != want { + t.Errorf("incorrect part length: %d, expected %d, error: %v", n+int64(len(prefix)), want, err) + } + // Force a retry after reading the body to verify that the section is rewound. + if failure == "" && number == 1 && firstPartAttempts.Add(1) == 1 { + w.WriteHeader(http.StatusInternalServerError) + _, _ = io.WriteString(w, `InternalError`) + return + } + w.Header().Set("ETag", fmt.Sprintf(`"part-%d"`, number)) + case r.Method == http.MethodPost && query.Has("uploadId"): + completes.Add(1) + var body struct { + Parts []struct { + Number int `xml:"PartNumber"` + ETag string + } `xml:"Part"` + } + if err := xml.NewDecoder(r.Body).Decode(&body); err != nil || len(body.Parts) != 2 { + t.Errorf("invalid completion request: %+v, error: %v", body, err) + } + for index, part := range body.Parts { + if part.Number != index+1 || part.ETag != fmt.Sprintf(`"part-%d"`, index+1) { + t.Errorf("parts completed out of order: %+v", body.Parts) + } + } + if failure == "complete" { + w.WriteHeader(http.StatusBadRequest) + _, _ = io.WriteString(w, `InvalidPart`) + return + } + _, _ = io.WriteString(w, `"complete"`) + case r.Method == http.MethodDelete: + aborts.Add(1) + w.WriteHeader(http.StatusNoContent) + default: + t.Errorf("unexpected request: %s %s", r.Method, r.URL) + w.WriteHeader(http.StatusBadRequest) + } + })) + defer server.Close() + storage := S3ReplayStorage{Bucket: "replays", Region: "us-east-1", Endpoint: server.URL} + var before, after runtime.MemStats + runtime.ReadMemStats(&before) + err := storage.Upload(file.Name(), "session.gz") + runtime.ReadMemStats(&after) + if (err != nil) != (failure != "") { + t.Fatalf("unexpected upload error: %v", err) + } + wantAborts, wantCompletes := int32(0), int32(1) + if failure != "" { + wantAborts = 1 + } + if failure == "part" { + wantCompletes = 0 + } + if aborts.Load() != wantAborts || completes.Load() != wantCompletes { + t.Fatalf("aborts=%d, completes=%d", aborts.Load(), completes.Load()) + } + if failure == "" { + if firstPartAttempts.Load() != 2 { + t.Fatalf("expected a retried part, got %d attempts", firstPartAttempts.Load()) + } + if allocated := after.TotalAlloc - before.TotalAlloc; allocated > 32*1024*1024 { + t.Fatalf("multipart upload buffered file contents: allocated %d bytes", allocated) + } + } + }) + } +} + +func setupS3TestEnv(t *testing.T) string { + t.Helper() + dir := t.TempDir() + t.Setenv("AWS_CONFIG_FILE", filepath.Join(dir, "config")) + t.Setenv("AWS_SHARED_CREDENTIALS_FILE", filepath.Join(dir, "credentials")) + t.Setenv("AWS_PROFILE", "") + t.Setenv("AWS_EC2_METADATA_DISABLED", "true") + t.Setenv("AWS_ACCESS_KEY_ID", "env-access") + t.Setenv("AWS_SECRET_ACCESS_KEY", "env-secret") + t.Setenv("AWS_SESSION_TOKEN", "env-token") + t.Setenv("AWS_REQUEST_CHECKSUM_CALCULATION", "when_supported") + t.Setenv("AWS_MAX_ATTEMPTS", "2") + return dir +}