From 5292972f879c65d34ee4f9aa59233505cf19195d Mon Sep 17 00:00:00 2001 From: Jason Hall Date: Tue, 24 Mar 2026 12:18:58 -0400 Subject: [PATCH] first agent crack Signed-off-by: Jason Hall --- buf.go | 94 ------------------ go.mod | 7 ++ go.sum | 161 ++++++++++++++++++++++++++++--- reader.go | 127 ++++++++++++++++++++++++ reader_test.go | 234 +++++++++++++++++++++++++++++++++++++++++++++ writer.go | 255 +++++++++++++++++++++++++++++++++++++++++++++++++ writer_test.go | 242 ++++++++++++++++++++++++++++++++++++++++++++++ 7 files changed, 1014 insertions(+), 106 deletions(-) delete mode 100644 buf.go create mode 100644 reader.go create mode 100644 reader_test.go create mode 100644 writer.go create mode 100644 writer_test.go diff --git a/buf.go b/buf.go deleted file mode 100644 index c0754e3..0000000 --- a/buf.go +++ /dev/null @@ -1,94 +0,0 @@ -package buf - -import ( - "bytes" - "context" - "errors" - "fmt" - "io" - "time" - - "cloud.google.com/go/storage" -) - -const maxComponentCount = 32 // TODO verify - -type Buffer struct { - svc *storage.Client - bucket, object string - buf *bytes.Buffer // written; not yet flushed -- not read from directly - - readcur int -} - -var _ interface { - io.ReadCloser - io.WriteCloser -} = (*Buffer)(nil) - -type Options struct { - Service *storage.Client - Bucket, Object string - Cap int -} - -func NewBuffer(opts Options) *Buffer { - return &Buffer{ - svc: opts.Service, - buf: bytes.NewBuffer(nil), - } -} - -func (buf *Buffer) Write(b []byte) (int, error) { return -1, errors.New("unimplemented") } -func (buf *Buffer) Read(b []byte) (int, error) { return -1, errors.New("unimplemented") } - -func (buf *Buffer) Close() error { - if err := buf.flush(); err != nil { - return err - } - return buf.finalize() -} -func (buf *Buffer) Flush() error { return buf.flush() } - -var now = time.Now // var for testing - -// TODO: https://docs.cloud.google.com/storage/docs/request-preconditions to prevent races -func (buf *Buffer) flush() error { - bh := buf.svc.Bucket(buf.bucket) - - // Insert a temp object. - tmpname := fmt.Sprintf("%s-temp-%d", buf.object, now().UnixNano()) - tmpobj := bh.Object(tmpname) - if _, err := io.Copy(tmpobj.NewWriter(context.TODO()), buf.buf); err != nil { - return fmt.Errorf("writing temp object: %w", err) - } - - // Compose it into the running object. - attrs, err := buf.svc.Bucket(buf.bucket).Object(buf.object).ComposerFrom(tmpobj).Run(context.TODO()) - if err != nil { - return fmt.Errorf("composing temp object: %w", err) - } - - // If reaching the max composite count, copy and rename to reset composite count. - if attrs.ComponentCount >= maxComponentCount { - // Copy object and move back - tmpname = fmt.Sprintf("%s-temp-%d", buf.object, now().UnixNano()) - objh := bh.Object(buf.object) - if _, err := bh.Object(tmpname).CopierFrom(objh).Run(context.TODO()); err != nil { - return fmt.Errorf("copying temp object: %w", err) - } - if _, err := bh.Object(tmpname).Move(context.TODO(), storage.MoveObjectDestination{Object: buf.object}); err != nil { - return fmt.Errorf("renaming temp object: %w", err) - } - } - return nil -} - -func (buf *Buffer) finalize() error { - if _, err := buf.svc.Bucket(buf.bucket).Object(buf.object).Update(context.TODO(), storage.ObjectAttrsToUpdate{ - Metadata: map[string]string{"done": "true"}, - }); err != nil { - return fmt.Errorf("finalizing object: %w", err) - } - return nil -} diff --git a/go.mod b/go.mod index 43c9066..f6645fa 100644 --- a/go.mod +++ b/go.mod @@ -4,6 +4,7 @@ go 1.25.0 require ( cloud.google.com/go/storage v1.61.3 + github.com/fsouza/fake-gcs-server v1.54.0 google.golang.org/api v0.272.0 ) @@ -15,6 +16,7 @@ require ( cloud.google.com/go/compute/metadata v0.9.0 // indirect cloud.google.com/go/iam v1.5.3 // indirect cloud.google.com/go/monitoring v1.24.3 // indirect + cloud.google.com/go/pubsub/v2 v2.4.0 // indirect github.com/GoogleCloudPlatform/opentelemetry-operations-go/detectors/gcp v1.30.0 // indirect github.com/GoogleCloudPlatform/opentelemetry-operations-go/exporter/metric v0.55.0 // indirect github.com/GoogleCloudPlatform/opentelemetry-operations-go/internal/resourcemapping v0.55.0 // indirect @@ -26,12 +28,17 @@ require ( github.com/go-jose/go-jose/v4 v4.1.3 // indirect github.com/go-logr/logr v1.4.3 // indirect github.com/go-logr/stdr v1.2.2 // indirect + github.com/google/renameio/v2 v2.0.0 // indirect github.com/google/s2a-go v0.1.9 // indirect github.com/google/uuid v1.6.0 // indirect github.com/googleapis/enterprise-certificate-proxy v0.3.14 // indirect github.com/googleapis/gax-go/v2 v2.18.0 // indirect + github.com/gorilla/handlers v1.5.2 // indirect + github.com/gorilla/mux v1.8.1 // indirect + github.com/pkg/xattr v0.4.12 // indirect github.com/planetscale/vtprotobuf v0.6.1-0.20240319094008-0393e58bdf10 // indirect github.com/spiffe/go-spiffe/v2 v2.6.0 // indirect + go.opencensus.io v0.24.0 // indirect go.opentelemetry.io/auto/sdk v1.2.1 // indirect go.opentelemetry.io/contrib/detectors/gcp v1.39.0 // indirect go.opentelemetry.io/contrib/instrumentation/google.golang.org/grpc/otelgrpc v0.63.0 // indirect diff --git a/go.sum b/go.sum index 484ce17..bb22cc4 100644 --- a/go.sum +++ b/go.sum @@ -1,5 +1,6 @@ cel.dev/expr v0.25.1 h1:1KrZg61W6TWSxuNZ37Xy49ps13NUovb66QLprthtwi4= cel.dev/expr v0.25.1/go.mod h1:hrXvqGP6G6gyx8UAHSHJ5RGk//1Oj5nXQ2NI02Nrsg4= +cloud.google.com/go v0.26.0/go.mod h1:aQUYkXzVsufM+DwF1aE+0xfcU+56JwCaLick0ClmMTw= cloud.google.com/go v0.123.0 h1:2NAUJwPR47q+E35uaJeYoNhuNEM9kM8SjgRgdeOJUSE= cloud.google.com/go v0.123.0/go.mod h1:xBoMV08QcqUGuPW65Qfm1o9Y4zKZBpGS+7bImXLTAZU= cloud.google.com/go/auth v0.18.2 h1:+Nbt5Ev0xEqxlNjd6c+yYUeosQ5TtEUaNcN/3FozlaM= @@ -10,29 +11,58 @@ cloud.google.com/go/compute/metadata v0.9.0 h1:pDUj4QMoPejqq20dK0Pg2N4yG9zIkYGdB cloud.google.com/go/compute/metadata v0.9.0/go.mod h1:E0bWwX5wTnLPedCKqk3pJmVgCBSM6qQI1yTBdEb3C10= cloud.google.com/go/iam v1.5.3 h1:+vMINPiDF2ognBJ97ABAYYwRgsaqxPbQDlMnbHMjolc= cloud.google.com/go/iam v1.5.3/go.mod h1:MR3v9oLkZCTlaqljW6Eb2d3HGDGK5/bDv93jhfISFvU= +cloud.google.com/go/logging v1.13.2 h1:qqlHCBvieJT9Cdq4QqYx1KPadCQ2noD4FK02eNqHAjA= +cloud.google.com/go/logging v1.13.2/go.mod h1:zaybliM3yun1J8mU2dVQ1/qDzjbOqEijZCn6hSBtKak= +cloud.google.com/go/longrunning v0.8.0 h1:LiKK77J3bx5gDLi4SMViHixjD2ohlkwBi+mKA7EhfW8= +cloud.google.com/go/longrunning v0.8.0/go.mod h1:UmErU2Onzi+fKDg2gR7dusz11Pe26aknR4kHmJJqIfk= cloud.google.com/go/monitoring v1.24.3 h1:dde+gMNc0UhPZD1Azu6at2e79bfdztVDS5lvhOdsgaE= cloud.google.com/go/monitoring v1.24.3/go.mod h1:nYP6W0tm3N9H/bOw8am7t62YTzZY+zUeQ+Bi6+2eonI= +cloud.google.com/go/pubsub/v2 v2.4.0 h1:oMKNiBQpXImRWnHYla9uSU66ZzByZwBSCJOEs/pTKVg= +cloud.google.com/go/pubsub/v2 v2.4.0/go.mod h1:2lS/XQKq5qtOMs6kHBK+WX1ytUC36kLl2ig3zqsGUx8= cloud.google.com/go/storage v1.61.3 h1:VS//ZfBuPGDvakfD9xyPW1RGF1Vy3BWUoVZXgW1KMOg= cloud.google.com/go/storage v1.61.3/go.mod h1:JtqK8BBB7TWv0HVGHubtUdzYYrakOQIsMLffZ2Z/HWk= +cloud.google.com/go/trace v1.11.7 h1:kDNDX8JkaAG3R2nq1lIdkb7FCSi1rCmsEtKVsty7p+U= +cloud.google.com/go/trace v1.11.7/go.mod h1:TNn9d5V3fQVf6s4SCveVMIBS2LJUqo73GACmq/Tky0s= +github.com/BurntSushi/toml v0.3.1/go.mod h1:xHWCNGjB5oqiDr8zfno3MHue2Ht5sIBksp03qcyfWMU= github.com/GoogleCloudPlatform/opentelemetry-operations-go/detectors/gcp v1.30.0 h1:sBEjpZlNHzK1voKq9695PJSX2o5NEXl7/OL3coiIY0c= github.com/GoogleCloudPlatform/opentelemetry-operations-go/detectors/gcp v1.30.0/go.mod h1:P4WPRUkOhJC13W//jWpyfJNDAIpvRbAUIYLX/4jtlE0= github.com/GoogleCloudPlatform/opentelemetry-operations-go/exporter/metric v0.55.0 h1:UnDZ/zFfG1JhH/DqxIZYU/1CUAlTUScoXD/LcM2Ykk8= github.com/GoogleCloudPlatform/opentelemetry-operations-go/exporter/metric v0.55.0/go.mod h1:IA1C1U7jO/ENqm/vhi7V9YYpBsp+IMyqNrEN94N7tVc= +github.com/GoogleCloudPlatform/opentelemetry-operations-go/internal/cloudmock v0.55.0 h1:7t/qx5Ost0s0wbA/VDrByOooURhp+ikYwv20i9Y07TQ= +github.com/GoogleCloudPlatform/opentelemetry-operations-go/internal/cloudmock v0.55.0/go.mod h1:vB2GH9GAYYJTO3mEn8oYwzEdhlayZIdQz6zdzgUIRvA= github.com/GoogleCloudPlatform/opentelemetry-operations-go/internal/resourcemapping v0.55.0 h1:0s6TxfCu2KHkkZPnBfsQ2y5qia0jl3MMrmBhu3nCOYk= github.com/GoogleCloudPlatform/opentelemetry-operations-go/internal/resourcemapping v0.55.0/go.mod h1:Mf6O40IAyB9zR/1J8nGDDPirZQQPbYJni8Yisy7NTMc= +github.com/census-instrumentation/opencensus-proto v0.2.1/go.mod h1:f6KPmirojxKA12rnyqOA5BBL4O983OfeGPqjHWSTneU= github.com/cespare/xxhash/v2 v2.3.0 h1:UL815xU9SqsFlibzuggzjXhog7bL6oX9BbNZnL2UFvs= github.com/cespare/xxhash/v2 v2.3.0/go.mod h1:VGX0DQ3Q6kWi7AoAeZDth3/j3BFtOZR5XLFGgcrjCOs= +github.com/client9/misspell v0.3.4/go.mod h1:qj6jICC3Q7zFZvVWo7KLAzC3yx5G7kyvSDkc90ppPyw= +github.com/cncf/udpa/go v0.0.0-20191209042840-269d4d468f6f/go.mod h1:M8M6+tZqaGXZJjfX53e64911xZQV5JYwmTeXPW+k8Sc= github.com/cncf/xds/go v0.0.0-20251210132809-ee656c7534f5 h1:6xNmx7iTtyBRev0+D/Tv1FZd4SCg8axKApyNyRsAt/w= github.com/cncf/xds/go v0.0.0-20251210132809-ee656c7534f5/go.mod h1:KdCmV+x/BuvyMxRnYBlmVaq4OLiKW6iRQfvC62cvdkI= -github.com/davecgh/go-spew v1.1.1 h1:vj9j/u1bqnvCEfJOwUhtlOARqs3+rkHYY13jYWTU97c= +github.com/davecgh/go-spew v1.1.0/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= github.com/davecgh/go-spew v1.1.1/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= +github.com/davecgh/go-spew v1.1.2-0.20180830191138-d8f796af33cc h1:U9qPSI2PIWSS1VwoXQT9A3Wy9MM3WgvqSxFWenqJduM= +github.com/davecgh/go-spew v1.1.2-0.20180830191138-d8f796af33cc/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= +github.com/dustin/go-humanize v1.0.1 h1:GzkhY7T5VNhEkwH0PVJgjz+fX1rhBrR7pRT3mDkpeCY= +github.com/dustin/go-humanize v1.0.1/go.mod h1:Mu1zIs6XwVuF/gI1OepvI0qD18qycQx+mFykh5fBlto= +github.com/envoyproxy/go-control-plane v0.9.0/go.mod h1:YTl/9mNaCwkRvm6d1a2C3ymFceY/DCBVvsKhRF0iEA4= +github.com/envoyproxy/go-control-plane v0.9.1-0.20191026205805-5f8ba28d4473/go.mod h1:YTl/9mNaCwkRvm6d1a2C3ymFceY/DCBVvsKhRF0iEA4= +github.com/envoyproxy/go-control-plane v0.9.4/go.mod h1:6rpuAdCZL397s3pYoYcLgu1mIlRU8Am5FuJP05cCM98= github.com/envoyproxy/go-control-plane v0.14.0 h1:hbG2kr4RuFj222B6+7T83thSPqLjwBIfQawTkC++2HA= +github.com/envoyproxy/go-control-plane v0.14.0/go.mod h1:NcS5X47pLl/hfqxU70yPwL9ZMkUlwlKxtAohpi2wBEU= github.com/envoyproxy/go-control-plane/envoy v1.36.0 h1:yg/JjO5E7ubRyKX3m07GF3reDNEnfOboJ0QySbH736g= github.com/envoyproxy/go-control-plane/envoy v1.36.0/go.mod h1:ty89S1YCCVruQAm9OtKeEkQLTb+Lkz0k8v9W0Oxsv98= +github.com/envoyproxy/go-control-plane/ratelimit v0.1.0 h1:/G9QYbddjL25KvtKTv3an9lx6VBE2cnb8wp1vEGNYGI= +github.com/envoyproxy/go-control-plane/ratelimit v0.1.0/go.mod h1:Wk+tMFAFbCXaJPzVVHnPgRKdUdwW/KdbRt94AzgRee4= +github.com/envoyproxy/protoc-gen-validate v0.1.0/go.mod h1:iSmxcyjqTsJpI2R4NaDN7+kN2VEUnK/pcBlmesArF7c= github.com/envoyproxy/protoc-gen-validate v1.3.0 h1:TvGH1wof4H33rezVKWSpqKz5NXWg5VPuZ0uONDT6eb4= github.com/envoyproxy/protoc-gen-validate v1.3.0/go.mod h1:HvYl7zwPa5mffgyeTUHA9zHIH36nmrm7oCbo4YKoSWA= github.com/felixge/httpsnoop v1.0.4 h1:NFTV2Zj1bL4mc9sqWACXbQFVBBg2W3GPvqp8/ESS2Wg= github.com/felixge/httpsnoop v1.0.4/go.mod h1:m8KPJKqk1gH5J9DgRY2ASl2lWCfGKXixSwevea8zH2U= +github.com/fsouza/fake-gcs-server v1.54.0 h1:DGO4EkFVbtP/A5Ha+CAHHx+Xa6O6LeskMB4hQ1wBE48= +github.com/fsouza/fake-gcs-server v1.54.0/go.mod h1:ryXYE4debQs8GjOxwaOAwFRwM4Cvs6S+NKPPgdVJe6g= +github.com/go-ini/ini v1.67.0 h1:z6ZrTEZqSWOTyH2FlglNbNgARyHG8oLW9gMELqKr06A= +github.com/go-ini/ini v1.67.0/go.mod h1:ByCAeIL28uOIIG0E3PJtZPDL8WnHpFKFOtgjp+3Ies8= github.com/go-jose/go-jose/v4 v4.1.3 h1:CVLmWDhDVRa6Mi/IgCgaopNosCaHz7zrMeF9MlZRkrs= github.com/go-jose/go-jose/v4 v4.1.3/go.mod h1:x4oUasVrzR7071A4TnHLGSPpNOm2a21K9Kf04k1rs08= github.com/go-logr/logr v1.2.2/go.mod h1:jdQByPbusPIv2/zmleS9BjJVeZ6kBagPoEUsqbVz/1A= @@ -40,26 +70,87 @@ github.com/go-logr/logr v1.4.3 h1:CjnDlHq8ikf6E492q6eKboGOC0T8CDaOvkHCIg8idEI= github.com/go-logr/logr v1.4.3/go.mod h1:9T104GzyrTigFIr8wt5mBrctHMim0Nb2HLGrmQ40KvY= github.com/go-logr/stdr v1.2.2 h1:hSWxHoqTgW2S2qGc0LTAI563KZ5YKYRhT3MFKZMbjag= github.com/go-logr/stdr v1.2.2/go.mod h1:mMo/vtBO5dYbehREoey6XUKy/eSumjCCveDpRre4VKE= +github.com/golang/glog v0.0.0-20160126235308-23def4e6c14b/go.mod h1:SBH7ygxi8pfUlaOkMMuAQtPIUF8ecWP5IEl/CR7VP2Q= +github.com/golang/groupcache v0.0.0-20200121045136-8c9f03a8e57e/go.mod h1:cIg4eruTrX1D+g88fzRXU5OdNfaM+9IcxsU14FzY7Hc= +github.com/golang/groupcache v0.0.0-20210331224755-41bb18bfe9da h1:oI5xCqsCo564l8iNU+DwB5epxmsaqB+rhGL0m5jtYqE= +github.com/golang/groupcache v0.0.0-20210331224755-41bb18bfe9da/go.mod h1:cIg4eruTrX1D+g88fzRXU5OdNfaM+9IcxsU14FzY7Hc= +github.com/golang/mock v1.1.1/go.mod h1:oTYuIxOrZwtPieC+H1uAHpcLFnEyAGVDL/k47Jfbm0A= +github.com/golang/protobuf v1.2.0/go.mod h1:6lQm79b+lXiMfvg/cZm0SGofjICqVBUtrP5yJMmIC1U= +github.com/golang/protobuf v1.3.2/go.mod h1:6lQm79b+lXiMfvg/cZm0SGofjICqVBUtrP5yJMmIC1U= +github.com/golang/protobuf v1.4.0-rc.1/go.mod h1:ceaxUfeHdC40wWswd/P6IGgMaK3YpKi5j83Wpe3EHw8= +github.com/golang/protobuf v1.4.0-rc.1.0.20200221234624-67d41d38c208/go.mod h1:xKAWHe0F5eneWXFV3EuXVDTCmh+JuBKY0li0aMyXATA= +github.com/golang/protobuf v1.4.0-rc.2/go.mod h1:LlEzMj4AhA7rCAGe4KMBDvJI+AwstrUpVNzEA03Pprs= +github.com/golang/protobuf v1.4.0-rc.4.0.20200313231945-b860323f09d0/go.mod h1:WU3c8KckQ9AFe+yFwt9sWVRKCVIyN9cPHBJSNnbL67w= +github.com/golang/protobuf v1.4.0/go.mod h1:jodUvKwWbYaEsadDk5Fwe5c77LiNKVO9IDvqG2KuDX0= +github.com/golang/protobuf v1.4.1/go.mod h1:U8fpvMrcmy5pZrNK1lt4xCsGvpyWQ/VVv6QDs8UjoX8= +github.com/golang/protobuf v1.4.3/go.mod h1:oDoupMAO8OvCJWAcko0GGGIgR6R6ocIYbsSw735rRwI= github.com/golang/protobuf v1.5.4 h1:i7eJL8qZTpSEXOPTxNKhASYpMn+8e5Q6AdndVa1dWek= github.com/golang/protobuf v1.5.4/go.mod h1:lnTiLA8Wa4RWRcIUkrtSVa5nRhsEGBg48fD6rSs7xps= +github.com/google/go-cmp v0.2.0/go.mod h1:oXzfMopK8JAjlY9xF4vHSVASa0yLyX7SntLO5aqRK0M= +github.com/google/go-cmp v0.3.0/go.mod h1:8QqcDgzrUqlUb/G2PQTWiueGozuR1884gddMywk6iLU= +github.com/google/go-cmp v0.3.1/go.mod h1:8QqcDgzrUqlUb/G2PQTWiueGozuR1884gddMywk6iLU= +github.com/google/go-cmp v0.4.0/go.mod h1:v8dTdLbMG2kIc/vJvl+f65V22dbkXbowE6jgT/gNBxE= +github.com/google/go-cmp v0.5.0/go.mod h1:v8dTdLbMG2kIc/vJvl+f65V22dbkXbowE6jgT/gNBxE= +github.com/google/go-cmp v0.5.3/go.mod h1:v8dTdLbMG2kIc/vJvl+f65V22dbkXbowE6jgT/gNBxE= github.com/google/go-cmp v0.7.0 h1:wk8382ETsv4JYUZwIsn6YpYiWiBsYLSJiTsyBybVuN8= github.com/google/go-cmp v0.7.0/go.mod h1:pXiqmnSA92OHEEa9HXL2W4E7lf9JzCmGVUdgjX3N/iU= +github.com/google/martian/v3 v3.3.3 h1:DIhPTQrbPkgs2yJYdXU/eNACCG5DVQjySNRNlflZ9Fc= +github.com/google/martian/v3 v3.3.3/go.mod h1:iEPrYcgCF7jA9OtScMFQyAlZZ4YXTKEtJ1E6RWzmBA0= +github.com/google/renameio/v2 v2.0.0 h1:UifI23ZTGY8Tt29JbYFiuyIU3eX+RNFtUwefq9qAhxg= +github.com/google/renameio/v2 v2.0.0/go.mod h1:BtmJXm5YlszgC+TD4HOEEUFgkJP3nLxehU6hfe7jRt4= github.com/google/s2a-go v0.1.9 h1:LGD7gtMgezd8a/Xak7mEWL0PjoTQFvpRudN895yqKW0= github.com/google/s2a-go v0.1.9/go.mod h1:YA0Ei2ZQL3acow2O62kdp9UlnvMmU7kA6Eutn0dXayM= +github.com/google/uuid v1.1.2/go.mod h1:TIyPZe4MgqvfeYDBFedMoGGpEw/LqOeaOT+nhxU+yHo= github.com/google/uuid v1.6.0 h1:NIvaJDMOsjHA8n1jAhLSgzrAzy1Hgr+hNrb57e+94F0= github.com/google/uuid v1.6.0/go.mod h1:TIyPZe4MgqvfeYDBFedMoGGpEw/LqOeaOT+nhxU+yHo= github.com/googleapis/enterprise-certificate-proxy v0.3.14 h1:yh8ncqsbUY4shRD5dA6RlzjJaT4hi3kII+zYw8wmLb8= github.com/googleapis/enterprise-certificate-proxy v0.3.14/go.mod h1:vqVt9yG9480NtzREnTlmGSBmFrA+bzb0yl0TxoBQXOg= github.com/googleapis/gax-go/v2 v2.18.0 h1:jxP5Uuo3bxm3M6gGtV94P4lliVetoCB4Wk2x8QA86LI= github.com/googleapis/gax-go/v2 v2.18.0/go.mod h1:uSzZN4a356eRG985CzJ3WfbFSpqkLTjsnhWGJR6EwrE= +github.com/gorilla/handlers v1.5.2 h1:cLTUSsNkgcwhgRqvCNmdbRWG0A3N4F+M2nWKdScwyEE= +github.com/gorilla/handlers v1.5.2/go.mod h1:dX+xVpaxdSw+q0Qek8SSsl3dfMk3jNddUkMzo0GtH0w= +github.com/gorilla/mux v1.8.1 h1:TuBL49tXwgrFYWhqrNgrUNEY92u81SPhu7sTdzQEiWY= +github.com/gorilla/mux v1.8.1/go.mod h1:AKf9I4AEqPTmMytcMc0KkNouC66V3BtZ4qD5fmWSiMQ= +github.com/klauspost/compress v1.18.2 h1:iiPHWW0YrcFgpBYhsA6D1+fqHssJscY/Tm/y2Uqnapk= +github.com/klauspost/compress v1.18.2/go.mod h1:R0h/fSBs8DE4ENlcrlib3PsXS61voFxhIs2DeRhCvJ4= +github.com/klauspost/cpuid/v2 v2.2.11 h1:0OwqZRYI2rFrjS4kvkDnqJkKHdHaRnCm68/DY4OxRzU= +github.com/klauspost/cpuid/v2 v2.2.11/go.mod h1:hqwkgyIinND0mEev00jJYCxPNVRVXFQeu1XKlok6oO0= +github.com/klauspost/crc32 v1.3.0 h1:sSmTt3gUt81RP655XGZPElI0PelVTZ6YwCRnPSupoFM= +github.com/klauspost/crc32 v1.3.0/go.mod h1:D7kQaZhnkX/Y0tstFGf8VUzv2UofNGqCjnC3zdHB0Hw= +github.com/minio/crc64nvme v1.1.1 h1:8dwx/Pz49suywbO+auHCBpCtlW1OfpcLN7wYgVR6wAI= +github.com/minio/crc64nvme v1.1.1/go.mod h1:eVfm2fAzLlxMdUGc0EEBGSMmPwmXD5XiNRpnu9J3bvg= +github.com/minio/md5-simd v1.1.2 h1:Gdi1DZK69+ZVMoNHRXJyNcxrMA4dSxoYHZSQbirFg34= +github.com/minio/md5-simd v1.1.2/go.mod h1:MzdKDxYpY2BT9XQFocsiZf/NKVtR7nkE4RoEpN+20RM= +github.com/minio/minio-go/v7 v7.0.98 h1:MeAVKjLVz+XJ28zFcuYyImNSAh8Mq725uNW4beRisi0= +github.com/minio/minio-go/v7 v7.0.98/go.mod h1:cY0Y+W7yozf0mdIclrttzo1Iiu7mEf9y7nk2uXqMOvM= +github.com/philhofer/fwd v1.2.0 h1:e6DnBTl7vGY+Gz322/ASL4Gyp1FspeMvx1RNDoToZuM= +github.com/philhofer/fwd v1.2.0/go.mod h1:RqIHx9QI14HlwKwm98g9Re5prTQ6LdeRQn+gXJFxsJM= +github.com/pkg/xattr v0.4.12 h1:rRTkSyFNTRElv6pkA3zpjHpQ90p/OdHQC1GmGh1aTjM= +github.com/pkg/xattr v0.4.12/go.mod h1:di8WF84zAKk8jzR1UBTEWh9AUlIZZ7M/JNt8e9B6ktU= github.com/planetscale/vtprotobuf v0.6.1-0.20240319094008-0393e58bdf10 h1:GFCKgmp0tecUJ0sJuv4pzYCqS9+RGSn52M3FUwPs+uo= github.com/planetscale/vtprotobuf v0.6.1-0.20240319094008-0393e58bdf10/go.mod h1:t/avpk3KcrXxUnYOhZhMXJlSEyie6gQbtLq5NM3loB8= -github.com/pmezard/go-difflib v1.0.0 h1:4DBwDE0NGyQoBHbLQYPwSUPoCMWR5BEzIk/f1lZbAQM= github.com/pmezard/go-difflib v1.0.0/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4= +github.com/pmezard/go-difflib v1.0.1-0.20181226105442-5d4384ee4fb2 h1:Jamvg5psRIccs7FGNTlIRMkT8wgtp5eCXdBlqhYGL6U= +github.com/pmezard/go-difflib v1.0.1-0.20181226105442-5d4384ee4fb2/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4= +github.com/prometheus/client_model v0.0.0-20190812154241-14fe0d1b01d4/go.mod h1:xMI15A0UPsDsEKsMN9yxemIoYk6Tm2C1GtYGdfGttqA= +github.com/rs/xid v1.6.0 h1:fV591PaemRlL6JfRxGDEPl69wICngIQ3shQtzfy2gxU= +github.com/rs/xid v1.6.0/go.mod h1:7XoLgs4eV+QndskICGsho+ADou8ySMSjJKDIan90Nz0= github.com/spiffe/go-spiffe/v2 v2.6.0 h1:l+DolpxNWYgruGQVV0xsfeya3CsC7m8iBzDnMpsbLuo= github.com/spiffe/go-spiffe/v2 v2.6.0/go.mod h1:gm2SeUoMZEtpnzPNs2Csc0D/gX33k1xIx7lEzqblHEs= +github.com/stretchr/objx v0.1.0/go.mod h1:HFkY916IF+rwdDfMAkV7OtwuqBVzrE8GR6GFx+wExME= +github.com/stretchr/objx v0.4.0/go.mod h1:YvHI0jy2hoMjB+UWwv71VJQ9isScKT/TqJzVSSt89Yw= +github.com/stretchr/objx v0.5.0/go.mod h1:Yh+to48EsGEfYuaHDzXPcE3xhTkx73EhmCGUpEOglKo= +github.com/stretchr/testify v1.7.1/go.mod h1:6Fq8oRcR53rry900zMqJjRRixrwX3KX962/h/Wwjteg= +github.com/stretchr/testify v1.8.0/go.mod h1:yNjHg4UonilssWZ8iaSj1OCr/vHnekPRkoO+kdMU+MU= +github.com/stretchr/testify v1.8.1/go.mod h1:w2LPCIKwWwSfY2zedu0+kehJoqGctiVI29o6fzry7u4= 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/tinylib/msgp v1.6.1 h1:ESRv8eL3u+DNHUoSAAQRE50Hm162zqAnBoGv9PzScPY= +github.com/tinylib/msgp v1.6.1/go.mod h1:RSp0LW9oSxFut3KzESt5Voq4GVWyS+PSulT77roAqEA= +go.einride.tech/aip v0.79.0 h1:19zdPlZzlUvxOA8syAFw4LkdJdXepzyTl6gt9XEeqdU= +go.einride.tech/aip v0.79.0/go.mod h1:E8+wdTApA70odnpFzJgsGogHozC2JCIhFJBKPr8bVig= +go.opencensus.io v0.24.0 h1:y73uSU6J157QMP2kn2r30vwW1A2W2WFwSCGnAVxeaD0= +go.opencensus.io v0.24.0/go.mod h1:vNK8G9p7aAivkbmorf4v+7Hgx+Zs0yY+0fOtgBfjQKo= go.opentelemetry.io/auto/sdk v1.2.1 h1:jXsnJ4Lmnqd11kwkBV2LgLoFMZKizbCi5fNZ/ipaZ64= go.opentelemetry.io/auto/sdk v1.2.1/go.mod h1:KRTj+aOaElaLi+wW1kO/DZRXwkF4C5xPbEe3ZiIhN7Y= go.opentelemetry.io/contrib/detectors/gcp v1.39.0 h1:kWRNZMsfBHZ+uHjiH4y7Etn2FK26LAGkNFw7RHv1DhE= @@ -68,53 +159,99 @@ go.opentelemetry.io/contrib/instrumentation/google.golang.org/grpc/otelgrpc v0.6 go.opentelemetry.io/contrib/instrumentation/google.golang.org/grpc/otelgrpc v0.63.0/go.mod h1:fvPi2qXDqFs8M4B4fmJhE92TyQs9Ydjlg3RvfUp+NbQ= go.opentelemetry.io/contrib/instrumentation/net/http/otelhttp v0.61.0 h1:F7Jx+6hwnZ41NSFTO5q4LYDtJRXBf2PD0rNBkeB/lus= go.opentelemetry.io/contrib/instrumentation/net/http/otelhttp v0.61.0/go.mod h1:UHB22Z8QsdRDrnAtX4PntOl36ajSxcdUMt1sF7Y6E7Q= -go.opentelemetry.io/otel v1.39.0 h1:8yPrr/S0ND9QEfTfdP9V+SiwT4E0G7Y5MO7p85nis48= -go.opentelemetry.io/otel v1.39.0/go.mod h1:kLlFTywNWrFyEdH0oj2xK0bFYZtHRYUdv1NklR/tgc8= go.opentelemetry.io/otel v1.40.0 h1:oA5YeOcpRTXq6NN7frwmwFR0Cn3RhTVZvXsP4duvCms= go.opentelemetry.io/otel v1.40.0/go.mod h1:IMb+uXZUKkMXdPddhwAHm6UfOwJyh4ct1ybIlV14J0g= -go.opentelemetry.io/otel/metric v1.39.0 h1:d1UzonvEZriVfpNKEVmHXbdf909uGTOQjA0HF0Ls5Q0= -go.opentelemetry.io/otel/metric v1.39.0/go.mod h1:jrZSWL33sD7bBxg1xjrqyDjnuzTUB0x1nBERXd7Ftcs= +go.opentelemetry.io/otel/exporters/stdout/stdoutmetric v1.40.0 h1:ZrPRak/kS4xI3AVXy8F7pipuDXmDsrO8Lg+yQjBLjw0= +go.opentelemetry.io/otel/exporters/stdout/stdoutmetric v1.40.0/go.mod h1:3y6kQCWztq6hyW8Z9YxQDDm0Je9AJoFar2G0yDcmhRk= go.opentelemetry.io/otel/metric v1.40.0 h1:rcZe317KPftE2rstWIBitCdVp89A2HqjkxR3c11+p9g= go.opentelemetry.io/otel/metric v1.40.0/go.mod h1:ib/crwQH7N3r5kfiBZQbwrTge743UDc7DTFVZrrXnqc= -go.opentelemetry.io/otel/sdk v1.39.0 h1:nMLYcjVsvdui1B/4FRkwjzoRVsMK8uL/cj0OyhKzt18= -go.opentelemetry.io/otel/sdk v1.39.0/go.mod h1:vDojkC4/jsTJsE+kh+LXYQlbL8CgrEcwmt1ENZszdJE= go.opentelemetry.io/otel/sdk v1.40.0 h1:KHW/jUzgo6wsPh9At46+h4upjtccTmuZCFAc9OJ71f8= go.opentelemetry.io/otel/sdk v1.40.0/go.mod h1:Ph7EFdYvxq72Y8Li9q8KebuYUr2KoeyHx0DRMKrYBUE= -go.opentelemetry.io/otel/sdk/metric v1.39.0 h1:cXMVVFVgsIf2YL6QkRF4Urbr/aMInf+2WKg+sEJTtB8= -go.opentelemetry.io/otel/sdk/metric v1.39.0/go.mod h1:xq9HEVH7qeX69/JnwEfp6fVq5wosJsY1mt4lLfYdVew= go.opentelemetry.io/otel/sdk/metric v1.40.0 h1:mtmdVqgQkeRxHgRv4qhyJduP3fYJRMX4AtAlbuWdCYw= go.opentelemetry.io/otel/sdk/metric v1.40.0/go.mod h1:4Z2bGMf0KSK3uRjlczMOeMhKU2rhUqdWNoKcYrtcBPg= -go.opentelemetry.io/otel/trace v1.39.0 h1:2d2vfpEDmCJ5zVYz7ijaJdOF59xLomrvj7bjt6/qCJI= -go.opentelemetry.io/otel/trace v1.39.0/go.mod h1:88w4/PnZSazkGzz/w84VHpQafiU4EtqqlVdxWy+rNOA= go.opentelemetry.io/otel/trace v1.40.0 h1:WA4etStDttCSYuhwvEa8OP8I5EWu24lkOzp+ZYblVjw= go.opentelemetry.io/otel/trace v1.40.0/go.mod h1:zeAhriXecNGP/s2SEG3+Y8X9ujcJOTqQ5RgdEJcawiA= +go.yaml.in/yaml/v3 v3.0.4 h1:tfq32ie2Jv2UxXFdLJdh3jXuOzWiL1fo0bu/FbuKpbc= +go.yaml.in/yaml/v3 v3.0.4/go.mod h1:DhzuOOF2ATzADvBadXxruRBLzYTpT36CKvDb3+aBEFg= +golang.org/x/crypto v0.0.0-20190308221718-c2843e01d9a2/go.mod h1:djNgcEr1/C05ACkg1iLfiJU5Ep61QUkGW8qpdssI0+w= +golang.org/x/crypto v0.0.0-20200622213623-75b288015ac9/go.mod h1:LzIPMQfyMNhhGPhUkYOs5KpL4U8rLKemX1yGLhDgUto= golang.org/x/crypto v0.49.0 h1:+Ng2ULVvLHnJ/ZFEq4KdcDd/cfjrrjjNSXNzxg0Y4U4= golang.org/x/crypto v0.49.0/go.mod h1:ErX4dUh2UM+CFYiXZRTcMpEcN8b/1gxEuv3nODoYtCA= +golang.org/x/exp v0.0.0-20190121172915-509febef88a4/go.mod h1:CJ0aWSM057203Lf6IL+f9T1iT9GByDxfZKAQTCR3kQA= +golang.org/x/lint v0.0.0-20181026193005-c67002cb31c3/go.mod h1:UVdnD1Gm6xHRNCYTkRU2/jEulfH38KcIWyp/GAMgvoE= +golang.org/x/lint v0.0.0-20190227174305-5b3e6a55c961/go.mod h1:wehouNa3lNwaWXcvxsM5YxQ5yQlVC4a0KAMCusXpPoU= +golang.org/x/lint v0.0.0-20190313153728-d0100b6bd8b3/go.mod h1:6SW0HCj/g11FgYtHlgUYUwCkIfeOF89ocIRzGO/8vkc= +golang.org/x/net v0.0.0-20180724234803-3673e40ba225/go.mod h1:mL1N/T3taQHkDXs73rZJwtUhF3w3ftmwwsq0BUmARs4= +golang.org/x/net v0.0.0-20180826012351-8a410e7b638d/go.mod h1:mL1N/T3taQHkDXs73rZJwtUhF3w3ftmwwsq0BUmARs4= +golang.org/x/net v0.0.0-20190213061140-3a22650c66bd/go.mod h1:mL1N/T3taQHkDXs73rZJwtUhF3w3ftmwwsq0BUmARs4= +golang.org/x/net v0.0.0-20190311183353-d8887717615a/go.mod h1:t9HGtf8HONx5eT2rtn7q6eTqICYqUVnKs3thJo3Qplg= +golang.org/x/net v0.0.0-20190404232315-eb5bcb51f2a3/go.mod h1:t9HGtf8HONx5eT2rtn7q6eTqICYqUVnKs3thJo3Qplg= +golang.org/x/net v0.0.0-20201110031124-69a78807bb2b/go.mod h1:sp8m0HH+o8qH0wwXwYZr8TS3Oi6o0r6Gce1SSxlDquU= golang.org/x/net v0.52.0 h1:He/TN1l0e4mmR3QqHMT2Xab3Aj3L9qjbhRm78/6jrW0= golang.org/x/net v0.52.0/go.mod h1:R1MAz7uMZxVMualyPXb+VaqGSa3LIaUqk0eEt3w36Sw= +golang.org/x/oauth2 v0.0.0-20180821212333-d2e6202438be/go.mod h1:N/0e6XlmueqKjAGxoOufVs8QHGRruUQn6yWY3a++T0U= golang.org/x/oauth2 v0.36.0 h1:peZ/1z27fi9hUOFCAZaHyrpWG5lwe0RJEEEeH0ThlIs= golang.org/x/oauth2 v0.36.0/go.mod h1:YDBUJMTkDnJS+A4BP4eZBjCqtokkg1hODuPjwiGPO7Q= +golang.org/x/sync v0.0.0-20180314180146-1d60e4601c6f/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM= +golang.org/x/sync v0.0.0-20181108010431-42b317875d0f/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM= +golang.org/x/sync v0.0.0-20190423024810-112230192c58/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM= golang.org/x/sync v0.20.0 h1:e0PTpb7pjO8GAtTs2dQ6jYa5BWYlMuX047Dco/pItO4= golang.org/x/sync v0.20.0/go.mod h1:9xrNwdLfx4jkKbNva9FpL6vEN7evnE43NNNJQ2LF3+0= +golang.org/x/sys v0.0.0-20180830151530-49385e6e1522/go.mod h1:STP8DvDyc/dI5b8T5hshtkjS+E42TnysNCUPdjciGhY= +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-20200930185726-fdedc70b468f/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= +golang.org/x/sys v0.0.0-20220408201424-a24fb2fb8a0f/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= golang.org/x/sys v0.42.0 h1:omrd2nAlyT5ESRdCLYdm3+fMfNFE/+Rf4bDIQImRJeo= golang.org/x/sys v0.42.0/go.mod h1:4GL1E5IUh+htKOUEOaiffhrAeqysfVGipDYzABqnCmw= +golang.org/x/text v0.3.0/go.mod h1:NqM8EUOU14njkJ3fqMW+pc6Ldnwhi/IjpwHt7yyuwOQ= +golang.org/x/text v0.3.3/go.mod h1:5Zoc/QRtKVWzQhOtBMvqHzDpF6irO9z98xDceosuGiQ= golang.org/x/text v0.35.0 h1:JOVx6vVDFokkpaq1AEptVzLTpDe9KGpj5tR4/X+ybL8= golang.org/x/text v0.35.0/go.mod h1:khi/HExzZJ2pGnjenulevKNX1W67CUy0AsXcNubPGCA= golang.org/x/time v0.15.0 h1:bbrp8t3bGUeFOx08pvsMYRTCVSMk89u4tKbNOZbp88U= golang.org/x/time v0.15.0/go.mod h1:Y4YMaQmXwGQZoFaVFk4YpCt4FLQMYKZe9oeV/f4MSno= +golang.org/x/tools v0.0.0-20180917221912-90fa682c2a6e/go.mod h1:n7NCudcB/nEzxVGmLbDWY5pfWTLqBcC2KZ6jyYvM4mQ= +golang.org/x/tools v0.0.0-20190114222345-bf090417da8b/go.mod h1:n7NCudcB/nEzxVGmLbDWY5pfWTLqBcC2KZ6jyYvM4mQ= +golang.org/x/tools v0.0.0-20190226205152-f727befe758c/go.mod h1:9Yl7xja0Znq3iFh3HoIrodX9oNMXvdceNzlUR8zjMvY= +golang.org/x/tools v0.0.0-20190311212946-11955173bddd/go.mod h1:LCzVGOaR6xXOjkQ3onu1FJEFr0SW1gC7cKk1uF8kGRs= +golang.org/x/tools v0.0.0-20190524140312-2c0ae7006135/go.mod h1:RgjU9mgBXZiqYHBnxXauZ1Gv1EHHAz9KjViQ78xBX0Q= +golang.org/x/xerrors v0.0.0-20191204190536-9bdfabe68543/go.mod h1:I/5z698sn9Ka8TeJc9MKroUUfqBBauWjQqLJ2OPfmY0= gonum.org/v1/gonum v0.16.0 h1:5+ul4Swaf3ESvrOnidPp4GZbzf0mxVQpDCYUQE7OJfk= gonum.org/v1/gonum v0.16.0/go.mod h1:fef3am4MQ93R2HHpKnLk4/Tbh/s0+wqD5nfa6Pnwy4E= google.golang.org/api v0.272.0 h1:eLUQZGnAS3OHn31URRf9sAmRk3w2JjMx37d2k8AjJmA= google.golang.org/api v0.272.0/go.mod h1:wKjowi5LNJc5qarNvDCvNQBn3rVK8nSy6jg2SwRwzIA= +google.golang.org/appengine v1.1.0/go.mod h1:EbEs0AVv82hx2wNQdGPgUI5lhzA/G0D9YwlJXL52JkM= +google.golang.org/appengine v1.4.0/go.mod h1:xpcJRLb0r/rnEns0DIKYYv+WjYCduHsrkT7/EB5XEv4= +google.golang.org/genproto v0.0.0-20180817151627-c66870c02cf8/go.mod h1:JiN7NxoALGmiZfu7CAH4rXhgtRTLTxftemlI0sWmxmc= +google.golang.org/genproto v0.0.0-20190819201941-24fa4b261c55/go.mod h1:DMBHOl98Agz4BDEuKkezgsaosCRResVns1a3J2ZsMNc= +google.golang.org/genproto v0.0.0-20200526211855-cb27e3aa2013/go.mod h1:NbSheEEYHJ7i3ixzK3sjbqSGDJWnxyFXZblF3eUsNvo= google.golang.org/genproto v0.0.0-20260217215200-42d3e9bedb6d h1:vsOm753cOAMkt76efriTCDKjpCbK18XGHMJHo0JUKhc= google.golang.org/genproto v0.0.0-20260217215200-42d3e9bedb6d/go.mod h1:0oz9d7g9QLSdv9/lgbIjowW1JoxMbxmBVNe8i6tORJI= google.golang.org/genproto/googleapis/api v0.0.0-20260217215200-42d3e9bedb6d h1:EocjzKLywydp5uZ5tJ79iP6Q0UjDnyiHkGRWxuPBP8s= google.golang.org/genproto/googleapis/api v0.0.0-20260217215200-42d3e9bedb6d/go.mod h1:48U2I+QQUYhsFrg2SY6r+nJzeOtjey7j//WBESw+qyQ= google.golang.org/genproto/googleapis/rpc v0.0.0-20260311181403-84a4fc48630c h1:xgCzyF2LFIO/0X2UAoVRiXKU5Xg6VjToG4i2/ecSswk= google.golang.org/genproto/googleapis/rpc v0.0.0-20260311181403-84a4fc48630c/go.mod h1:4Hqkh8ycfw05ld/3BWL7rJOSfebL2Q+DVDeRgYgxUU8= +google.golang.org/grpc v1.19.0/go.mod h1:mqu4LbDTu4XGKhr4mRzUsmM4RtVoemTSY81AxZiDr8c= +google.golang.org/grpc v1.23.0/go.mod h1:Y5yQAOtifL1yxbo5wqy6BxZv8vAUGQwXBOALyacEbxg= +google.golang.org/grpc v1.25.1/go.mod h1:c3i+UQWmh7LiEpx4sFZnkU36qjEYZ0imhYfXVyQciAY= +google.golang.org/grpc v1.27.0/go.mod h1:qbnxyOmOxrQa7FizSgH+ReBfzJrCY1pSN7KXBS8abTk= +google.golang.org/grpc v1.33.2/go.mod h1:JMHMWHQWaTccqQQlmk3MJZS+GWXOdAesneDmEnv2fbc= google.golang.org/grpc v1.79.2 h1:fRMD94s2tITpyJGtBBn7MkMseNpOZU8ZxgC3MMBaXRU= google.golang.org/grpc v1.79.2/go.mod h1:KmT0Kjez+0dde/v2j9vzwoAScgEPx/Bw1CYChhHLrHQ= +google.golang.org/protobuf v0.0.0-20200109180630-ec00e32a8dfd/go.mod h1:DFci5gLYBciE7Vtevhsrf46CRTquxDuWsQurQQe4oz8= +google.golang.org/protobuf v0.0.0-20200221191635-4d8936d0db64/go.mod h1:kwYJMbMJ01Woi6D6+Kah6886xMZcty6N08ah7+eCXa0= +google.golang.org/protobuf v0.0.0-20200228230310-ab0ca4ff8a60/go.mod h1:cfTl7dwQJ+fmap5saPgwCLgHXTUD7jkjRqWcaiX5VyM= +google.golang.org/protobuf v1.20.1-0.20200309200217-e05f789c0967/go.mod h1:A+miEFZTKqfCUM6K7xSMQL9OKL/b6hQv+e19PK+JZNE= +google.golang.org/protobuf v1.21.0/go.mod h1:47Nbq4nVaFHyn7ilMalzfO3qCViNmqZ2kzikPIcrTAo= +google.golang.org/protobuf v1.22.0/go.mod h1:EGpADcykh3NcUnDUJcl1+ZksZNG86OlYog2l/sGQquU= +google.golang.org/protobuf v1.23.0/go.mod h1:EGpADcykh3NcUnDUJcl1+ZksZNG86OlYog2l/sGQquU= +google.golang.org/protobuf v1.23.1-0.20200526195155-81db48ad09cc/go.mod h1:EGpADcykh3NcUnDUJcl1+ZksZNG86OlYog2l/sGQquU= +google.golang.org/protobuf v1.25.0/go.mod h1:9JNX74DMeImyA3h4bdi1ymwjUzf21/xIlbajtzgsN7c= google.golang.org/protobuf v1.36.11 h1:fV6ZwhNocDyBLK0dj+fg8ektcVegBBuEolpbTQyBNVE= google.golang.org/protobuf v1.36.11/go.mod h1:HTf+CrKn2C3g5S8VImy6tdcUvCska2kB7j23XfzDpco= +gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0= +gopkg.in/yaml.v3 v3.0.0-20200313102051-9f266ea9e77c/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= gopkg.in/yaml.v3 v3.0.1 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA= gopkg.in/yaml.v3 v3.0.1/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= +honnef.co/go/tools v0.0.0-20190102054323-c2f93a96b099/go.mod h1:rf3lG4BRIbNafJWhAfAdb/ePZxsR/4RtNHQocxwk9r4= +honnef.co/go/tools v0.0.0-20190523083050-ea95bdfd59fc/go.mod h1:rf3lG4BRIbNafJWhAfAdb/ePZxsR/4RtNHQocxwk9r4= diff --git a/reader.go b/reader.go new file mode 100644 index 0000000..9499183 --- /dev/null +++ b/reader.go @@ -0,0 +1,127 @@ +package buf + +import ( + "context" + "errors" + "fmt" + "io" + "net/http" + "time" + + "cloud.google.com/go/storage" + "google.golang.org/api/googleapi" +) + +const defaultPollInterval = 1 * time.Second + +// ReaderOptions configures a Reader. +type ReaderOptions struct { + // Client is the GCS client used for all operations. + Client *storage.Client + // Bucket is the GCS bucket name. + Bucket string + // Object is the GCS object name. + Object string + // PollInterval is how often the reader polls for new data or object creation. Default 1s. + PollInterval time.Duration +} + +type Reader struct { + ctx context.Context + client *storage.Client + bucket string + object string + pollInterval time.Duration + + cursor int64 + done bool +} + +func NewReader(ctx context.Context, opts ReaderOptions) *Reader { + interval := opts.PollInterval + if interval == 0 { + interval = defaultPollInterval + } + return &Reader{ + ctx: ctx, + client: opts.Client, + bucket: opts.Bucket, + object: opts.Object, + pollInterval: interval, + } +} + +func (r *Reader) Read(p []byte) (int, error) { + for { + objh := r.client.Bucket(r.bucket).Object(r.object) + + attrs, err := objh.Attrs(r.ctx) + if err != nil { + if isNotFound(err) { + // Object doesn't exist yet, poll. + if err := r.sleep(); err != nil { + return 0, err + } + continue + } + return 0, fmt.Errorf("getting object attrs: %w", err) + } + + if r.cursor >= attrs.Size { + // No new data. Check if done. + if attrs.Metadata["done"] == "true" { + r.done = true + return 0, io.EOF + } + // Not done, poll for more data. + if err := r.sleep(); err != nil { + return 0, err + } + continue + } + + // Read available data. + length := int64(len(p)) + available := attrs.Size - r.cursor + if length > available { + length = available + } + + rc, err := objh.NewRangeReader(r.ctx, r.cursor, length) + if err != nil { + return 0, fmt.Errorf("creating range reader: %w", err) + } + n, err := io.ReadFull(rc, p[:length]) + rc.Close() + if err != nil && !errors.Is(err, io.ErrUnexpectedEOF) { + return n, fmt.Errorf("reading object data: %w", err) + } + + r.cursor += int64(n) + return n, nil + } +} + +func (r *Reader) Close() error { + return nil +} + +func (r *Reader) sleep() error { + select { + case <-time.After(r.pollInterval): + return nil + case <-r.ctx.Done(): + return r.ctx.Err() + } +} + +func isNotFound(err error) bool { + if errors.Is(err, storage.ErrObjectNotExist) { + return true + } + var apiErr *googleapi.Error + if errors.As(err, &apiErr) && apiErr.Code == http.StatusNotFound { + return true + } + return false +} diff --git a/reader_test.go b/reader_test.go new file mode 100644 index 0000000..bba3875 --- /dev/null +++ b/reader_test.go @@ -0,0 +1,234 @@ +package buf + +import ( + "bytes" + "context" + "io" + "strings" + "sync" + "testing" + "time" + + "cloud.google.com/go/storage" +) + +// writeObject writes content to a GCS object with optional metadata. +func writeObject(t *testing.T, client *storage.Client, bucket, object, content string, metadata map[string]string) { + t.Helper() + ctx := context.Background() + w := client.Bucket(bucket).Object(object).NewWriter(ctx) + if metadata != nil { + w.Metadata = metadata + } + if _, err := io.Copy(w, bytes.NewReader([]byte(content))); err != nil { + t.Fatalf("writing object %s/%s: %v", bucket, object, err) + } + if err := w.Close(); err != nil { + t.Fatalf("closing object writer %s/%s: %v", bucket, object, err) + } +} + +func TestReaderFromCompleteObject(t *testing.T) { + client := setup(t) + ctx := context.Background() + + // Pre-populate an object with done metadata. + writeObject(t, client, "test", "complete", "hello reader", map[string]string{"done": "true"}) + + r := NewReader(ctx, ReaderOptions{ + Client: client, + Bucket: "test", + Object: "complete", + PollInterval: 10 * time.Millisecond, + }) + + got, err := io.ReadAll(r) + if err != nil { + t.Fatal(err) + } + if string(got) != "hello reader" { + t.Fatalf("got %q, want %q", string(got), "hello reader") + } +} + +func TestReaderPollsForObjectCreation(t *testing.T) { + client := setup(t) + ctx := context.Background() + + r := NewReader(ctx, ReaderOptions{ + Client: client, + Bucket: "test", + Object: "delayed", + PollInterval: 10 * time.Millisecond, + }) + + // Create the object after a short delay. + go func() { + time.Sleep(50 * time.Millisecond) + writeObject(t, client, "test", "delayed", "delayed data", map[string]string{"done": "true"}) + }() + + got, err := io.ReadAll(r) + if err != nil { + t.Fatal(err) + } + if string(got) != "delayed data" { + t.Fatalf("got %q, want %q", string(got), "delayed data") + } +} + +func TestReaderPollsForMoreData(t *testing.T) { + client := setup(t) + ctx := context.Background() + + // Create initial object without done. + writeObject(t, client, "test", "growing", "part1", nil) + + r := NewReader(ctx, ReaderOptions{ + Client: client, + Bucket: "test", + Object: "growing", + PollInterval: 10 * time.Millisecond, + }) + + // Read first chunk. + buf := make([]byte, 100) + n, err := r.Read(buf) + if err != nil { + t.Fatal(err) + } + if string(buf[:n]) != "part1" { + t.Fatalf("got %q, want %q", string(buf[:n]), "part1") + } + + // Overwrite with more data and set done. + go func() { + time.Sleep(50 * time.Millisecond) + writeObject(t, client, "test", "growing", "part1part2", map[string]string{"done": "true"}) + }() + + rest, err := io.ReadAll(r) + if err != nil { + t.Fatal(err) + } + if string(rest) != "part2" { + t.Fatalf("got %q, want %q", string(rest), "part2") + } +} + +func TestReaderContextCancellation(t *testing.T) { + client := setup(t) + ctx, cancel := context.WithCancel(context.Background()) + + r := NewReader(ctx, ReaderOptions{ + Client: client, + Bucket: "test", + Object: "nonexistent", + PollInterval: 10 * time.Millisecond, + }) + + go func() { + time.Sleep(50 * time.Millisecond) + cancel() + }() + + _, err := r.Read(make([]byte, 100)) + if err == nil { + t.Fatal("expected error from cancelled context") + } +} + +func TestIntegrationWriterReader(t *testing.T) { + client := setup(t) + ctx := context.Background() + + w := NewWriter(ctx, WriterOptions{ + Client: client, + Bucket: "test", + Object: "stream", + FlushInterval: 50 * time.Millisecond, + FlushSize: 1024 * 1024, + }) + + r := NewReader(ctx, ReaderOptions{ + Client: client, + Bucket: "test", + Object: "stream", + PollInterval: 10 * time.Millisecond, + }) + + // Write data in chunks. + chunks := []string{"hello ", "world ", "from ", "gcsbuf"} + var wg sync.WaitGroup + wg.Go(func() { + for _, chunk := range chunks { + if _, err := w.Write([]byte(chunk)); err != nil { + t.Errorf("write: %v", err) + return + } + time.Sleep(30 * time.Millisecond) + } + if err := w.Close(); err != nil { + t.Errorf("close: %v", err) + } + }) + + got, err := io.ReadAll(r) + if err != nil { + t.Fatal(err) + } + + wg.Wait() + + want := "hello world from gcsbuf" + if string(got) != want { + t.Fatalf("got %q, want %q", string(got), want) + } +} + +func TestIntegrationMultipleFlushCyclesWithReader(t *testing.T) { + client := setup(t) + ctx := context.Background() + + w := NewWriter(ctx, WriterOptions{ + Client: client, + Bucket: "test", + Object: "multi", + FlushInterval: time.Hour, + FlushSize: 10, // flush every 10 bytes + }) + + r := NewReader(ctx, ReaderOptions{ + Client: client, + Bucket: "test", + Object: "multi", + PollInterval: 10 * time.Millisecond, + }) + + // Write enough data to trigger multiple auto-flushes. + var wg sync.WaitGroup + wg.Go(func() { + for i := range 5 { + chunk := strings.Repeat("x", 15) // > flushSize, triggers auto-flush each time + if _, err := w.Write([]byte(chunk)); err != nil { + t.Errorf("write %d: %v", i, err) + return + } + } + if err := w.Close(); err != nil { + t.Errorf("close: %v", err) + } + }) + + got, err := io.ReadAll(r) + if err != nil { + t.Fatal(err) + } + + wg.Wait() + + want := strings.Repeat("x", 75) + if string(got) != want { + t.Fatalf("got len=%d, want len=%d", len(got), len(want)) + } +} diff --git a/writer.go b/writer.go new file mode 100644 index 0000000..ef3d9ad --- /dev/null +++ b/writer.go @@ -0,0 +1,255 @@ +package buf + +import ( + "bytes" + "context" + "fmt" + "sync" + "time" + + "cloud.google.com/go/storage" +) + +const ( + maxComponentCount = 32 + defaultFlushInterval = 5 * time.Second + defaultFlushSize = 4 * 1024 * 1024 // 4MB +) + +var now = time.Now // var for testing + +// WriterOptions configures a Writer. +type WriterOptions struct { + // Client is the GCS client used for all operations. + Client *storage.Client + // Bucket is the GCS bucket name. + Bucket string + // Object is the GCS object name. + Object string + // FlushInterval is how often the background goroutine flushes buffered data. Default 5s. + FlushInterval time.Duration + // FlushSize is the buffer size threshold that triggers a synchronous flush on Write. Default 4MB. + FlushSize int +} + +type Writer struct { + ctx context.Context + client *storage.Client + bucket string + object string + + flushInterval time.Duration + flushSize int + + mu sync.Mutex + buf *bytes.Buffer + created bool // whether the GCS object exists yet + closed bool + err error + + stopCh chan struct{} + doneCh chan struct{} +} + +func NewWriter(ctx context.Context, opts WriterOptions) *Writer { + interval := opts.FlushInterval + if interval == 0 { + interval = defaultFlushInterval + } + size := opts.FlushSize + if size == 0 { + size = defaultFlushSize + } + w := &Writer{ + ctx: ctx, + client: opts.Client, + bucket: opts.Bucket, + object: opts.Object, + flushInterval: interval, + flushSize: size, + buf: bytes.NewBuffer(nil), + stopCh: make(chan struct{}), + doneCh: make(chan struct{}), + } + go w.backgroundFlush() + return w +} + +func (w *Writer) backgroundFlush() { + defer close(w.doneCh) + ticker := time.NewTicker(w.flushInterval) + defer ticker.Stop() + for { + select { + case <-w.stopCh: + return + case <-ticker.C: + w.mu.Lock() + if w.buf.Len() > 0 { + w.flushLocked() + } + w.mu.Unlock() + } + } +} + +func (w *Writer) Write(p []byte) (int, error) { + w.mu.Lock() + defer w.mu.Unlock() + if w.closed { + return 0, fmt.Errorf("write to closed writer") + } + if w.err != nil { + return 0, w.err + } + n, err := w.buf.Write(p) + if err != nil { + return n, err + } + if w.buf.Len() >= w.flushSize { + w.flushLocked() + } + if w.err != nil { + return 0, w.err + } + return n, nil +} + +// Flush flushes any buffered data to GCS. +func (w *Writer) Flush() error { + w.mu.Lock() + defer w.mu.Unlock() + if w.buf.Len() > 0 { + w.flushLocked() + } + return w.err +} + +// flushLocked writes buffered data to GCS. Must be called with w.mu held. +func (w *Writer) flushLocked() { + data := w.buf.Bytes() + if len(data) == 0 { + return + } + + bh := w.client.Bucket(w.bucket) + objh := bh.Object(w.object) + + if !w.created { + // First write: create the object directly. + gw := objh.NewWriter(w.ctx) + if _, err := gw.Write(data); err != nil { + w.err = fmt.Errorf("writing initial object: %w", err) + return + } + if err := gw.Close(); err != nil { + w.err = fmt.Errorf("closing initial object writer: %w", err) + return + } + w.created = true + w.buf.Reset() + return + } + + // Subsequent writes: write to temp, compose, delete temp. + tmpname := fmt.Sprintf("%s-temp-%d", w.object, now().UnixNano()) + tmpobj := bh.Object(tmpname) + + gw := tmpobj.NewWriter(w.ctx) + if _, err := gw.Write(data); err != nil { + w.err = fmt.Errorf("writing temp object: %w", err) + return + } + if err := gw.Close(); err != nil { + w.err = fmt.Errorf("closing temp object writer: %w", err) + return + } + + // Check component count before composing. + attrs, err := objh.Attrs(w.ctx) + if err != nil { + w.err = fmt.Errorf("getting object attrs: %w", err) + return + } + if attrs.ComponentCount >= maxComponentCount-1 { + // Reset component count: copy to temp2, copy back, delete temp2. + if err := w.resetComponentCount(bh, objh); err != nil { + w.err = fmt.Errorf("resetting component count: %w", err) + return + } + } + + // Compose existing + temp into existing. + if _, err := objh.ComposerFrom(objh, tmpobj).Run(w.ctx); err != nil { + w.err = fmt.Errorf("composing objects: %w", err) + return + } + + // Delete temp object. + if err := tmpobj.Delete(w.ctx); err != nil { + w.err = fmt.Errorf("deleting temp object: %w", err) + return + } + + w.buf.Reset() +} + +// resetComponentCount copies the object to a temp location and back to reset the component count. +func (w *Writer) resetComponentCount(bh *storage.BucketHandle, objh *storage.ObjectHandle) error { + tmpname := fmt.Sprintf("%s-reset-%d", w.object, now().UnixNano()) + tmpobj := bh.Object(tmpname) + + // Copy to temp (this creates a non-composite object). + if _, err := tmpobj.CopierFrom(objh).Run(w.ctx); err != nil { + return fmt.Errorf("copying to reset temp: %w", err) + } + + // Copy back. + if _, err := objh.CopierFrom(tmpobj).Run(w.ctx); err != nil { + return fmt.Errorf("copying back from reset temp: %w", err) + } + + // Delete the reset temp. + if err := tmpobj.Delete(w.ctx); err != nil { + return fmt.Errorf("deleting reset temp: %w", err) + } + + return nil +} + +func (w *Writer) Close() error { + w.mu.Lock() + if w.closed { + w.mu.Unlock() + return fmt.Errorf("writer already closed") + } + w.closed = true + w.mu.Unlock() + + // Stop background goroutine and wait for it. + close(w.stopCh) + <-w.doneCh + + // Final flush. + w.mu.Lock() + if w.buf.Len() > 0 { + w.flushLocked() + } + err := w.err + w.mu.Unlock() + if err != nil { + return err + } + + // Set done metadata. Only if the object was created. + if w.created { + objh := w.client.Bucket(w.bucket).Object(w.object) + if _, err := objh.Update(w.ctx, storage.ObjectAttrsToUpdate{ + Metadata: map[string]string{"done": "true"}, + }); err != nil { + return fmt.Errorf("setting done metadata: %w", err) + } + } + + return nil +} diff --git a/writer_test.go b/writer_test.go new file mode 100644 index 0000000..a7f029a --- /dev/null +++ b/writer_test.go @@ -0,0 +1,242 @@ +package buf + +import ( + "context" + "io" + "strings" + "testing" + "time" + + "cloud.google.com/go/storage" + "github.com/fsouza/fake-gcs-server/fakestorage" +) + +func setup(t *testing.T) *storage.Client { + t.Helper() + server := fakestorage.NewServer(nil) + server.CreateBucketWithOpts(fakestorage.CreateBucketOpts{Name: "test"}) + t.Cleanup(server.Stop) + return server.Client() +} + +// readObject reads the full content of a GCS object as a string. +func readObject(t *testing.T, client *storage.Client, bucket, object string) string { + t.Helper() + rc, err := client.Bucket(bucket).Object(object).NewReader(context.Background()) + if err != nil { + t.Fatalf("reading object %s/%s: %v", bucket, object, err) + } + defer rc.Close() + data, err := io.ReadAll(rc) + if err != nil { + t.Fatalf("reading object %s/%s: %v", bucket, object, err) + } + return string(data) +} + +func TestWriterFlush(t *testing.T) { + client := setup(t) + ctx := context.Background() + + w := NewWriter(ctx, WriterOptions{ + Client: client, + Bucket: "test", + Object: "obj", + FlushInterval: time.Hour, // disable periodic flush + }) + + data := []byte("hello world") + n, err := w.Write(data) + if err != nil { + t.Fatal(err) + } + if n != len(data) { + t.Fatalf("wrote %d, want %d", n, len(data)) + } + + if err := w.Flush(); err != nil { + t.Fatal(err) + } + + // Verify object content. + got := readObject(t, client, "test", "obj") + if got != "hello world" { + t.Fatalf("got %q, want %q", got, "hello world") + } +} + +func TestWriterMultipleFlushes(t *testing.T) { + client := setup(t) + ctx := context.Background() + + w := NewWriter(ctx, WriterOptions{ + Client: client, + Bucket: "test", + Object: "obj", + FlushInterval: time.Hour, + }) + + w.Write([]byte("aaa")) + if err := w.Flush(); err != nil { + t.Fatal(err) + } + w.Write([]byte("bbb")) + if err := w.Flush(); err != nil { + t.Fatal(err) + } + w.Write([]byte("ccc")) + if err := w.Flush(); err != nil { + t.Fatal(err) + } + + got := readObject(t, client, "test", "obj") + if got != "aaabbbccc" { + t.Fatalf("got %q, want %q", got, "aaabbbccc") + } +} + +func TestWriterAutoFlushOnSize(t *testing.T) { + client := setup(t) + ctx := context.Background() + + w := NewWriter(ctx, WriterOptions{ + Client: client, + Bucket: "test", + Object: "obj", + FlushInterval: time.Hour, + FlushSize: 10, // small threshold + }) + + // Write more than flushSize; should auto-flush. + data := []byte("0123456789extra") + if _, err := w.Write(data); err != nil { + t.Fatal(err) + } + + // Object should exist now (auto-flushed). + got := readObject(t, client, "test", "obj") + if got != string(data) { + t.Fatalf("got %q, want %q", got, string(data)) + } +} + +func TestWriterPeriodicFlush(t *testing.T) { + client := setup(t) + ctx := context.Background() + + w := NewWriter(ctx, WriterOptions{ + Client: client, + Bucket: "test", + Object: "obj", + FlushInterval: 50 * time.Millisecond, + }) + + w.Write([]byte("periodic")) + + // Wait for periodic flush. + time.Sleep(200 * time.Millisecond) + + got := readObject(t, client, "test", "obj") + if got != "periodic" { + t.Fatalf("got %q, want %q", got, "periodic") + } + + if err := w.Close(); err != nil { + t.Fatal(err) + } +} + +func TestWriterCloseSetsDoneMetadata(t *testing.T) { + client := setup(t) + ctx := context.Background() + + w := NewWriter(ctx, WriterOptions{ + Client: client, + Bucket: "test", + Object: "obj", + FlushInterval: time.Hour, + }) + + w.Write([]byte("data")) + if err := w.Close(); err != nil { + t.Fatal(err) + } + + attrs, err := client.Bucket("test").Object("obj").Attrs(ctx) + if err != nil { + t.Fatal(err) + } + if attrs.Metadata["done"] != "true" { + t.Fatalf("metadata = %v, want done=true", attrs.Metadata) + } +} + +func TestWriterWriteAfterClose(t *testing.T) { + client := setup(t) + ctx := context.Background() + + w := NewWriter(ctx, WriterOptions{ + Client: client, + Bucket: "test", + Object: "obj", + FlushInterval: time.Hour, + }) + + w.Write([]byte("data")) + w.Close() + + _, err := w.Write([]byte("more")) + if err == nil { + t.Fatal("expected error writing to closed writer") + } +} + +func TestWriterCloseNoData(t *testing.T) { + client := setup(t) + ctx := context.Background() + + w := NewWriter(ctx, WriterOptions{ + Client: client, + Bucket: "test", + Object: "obj", + FlushInterval: time.Hour, + }) + + // Close without writing anything should not error. + if err := w.Close(); err != nil { + t.Fatal(err) + } +} + +func TestWriterComponentCountReset(t *testing.T) { + client := setup(t) + ctx := context.Background() + + w := NewWriter(ctx, WriterOptions{ + Client: client, + Bucket: "test", + Object: "obj", + FlushInterval: time.Hour, + FlushSize: 1024 * 1024, // large so we control flushes manually + }) + + // Do many flushes to trigger component count reset. + var want strings.Builder + for i := range 35 { + chunk := strings.Repeat(string(rune('a'+i%26)), 100) + want.WriteString(chunk) + w.Write([]byte(chunk)) + if err := w.Flush(); err != nil { + t.Fatalf("flush %d: %v", i, err) + } + } + + if err := w.Close(); err != nil { + t.Fatal(err) + } + + got := readObject(t, client, "test", "obj") + if got != want.String() { + t.Fatalf("content mismatch after %d flushes: got len=%d, want len=%d", 35, len(got), want.Len()) + } +}