diff --git a/ci/resources/stemcell-version-bump/go.mod b/ci/resources/stemcell-version-bump/go.mod index f95b9af8..e282b43c 100644 --- a/ci/resources/stemcell-version-bump/go.mod +++ b/ci/resources/stemcell-version-bump/go.mod @@ -3,7 +3,7 @@ module stemcell-version-bump go 1.25.8 require ( - cloud.google.com/go/storage v1.63.1 + cloud.google.com/go/storage v1.64.0 github.com/stretchr/testify v1.11.1 google.golang.org/api v0.290.0 ) diff --git a/ci/resources/stemcell-version-bump/go.sum b/ci/resources/stemcell-version-bump/go.sum index e77b5fdc..5ab75a84 100644 --- a/ci/resources/stemcell-version-bump/go.sum +++ b/ci/resources/stemcell-version-bump/go.sum @@ -16,8 +16,8 @@ cloud.google.com/go/longrunning v1.2.0 h1:WjYH3YHBGCxGJP9M4dWGHBfXr/cFIjMkNgWcJj cloud.google.com/go/longrunning v1.2.0/go.mod h1:5KMQALFGOCtFoi2xSOA1u3H7WKlhmckgiyFw7+LGQp0= cloud.google.com/go/monitoring v1.29.0 h1:AHhDsFaSax1/4k+qlIDX/SDGe6hggnfXJ9dkgD9qBPY= cloud.google.com/go/monitoring v1.29.0/go.mod h1:72NOVjJXHY/HBfoLT0+qlCZBT059+9VXLeAnL2PeeVM= -cloud.google.com/go/storage v1.63.1 h1:CYXILV9G4CH0C18IQ9+V0h4XiqD2LhKnMLO0o7uJWNs= -cloud.google.com/go/storage v1.63.1/go.mod h1:lWyAtwvDZHdL3k68WVKbESP6bmWaV23ZJJ/JEVw/ZaQ= +cloud.google.com/go/storage v1.64.0 h1:KLpxI/oX9LxeRsNqn877d2WyeT3ryiEwnGt8pwcSPZg= +cloud.google.com/go/storage v1.64.0/go.mod h1:lWyAtwvDZHdL3k68WVKbESP6bmWaV23ZJJ/JEVw/ZaQ= cloud.google.com/go/trace v1.16.0 h1:GmQovzFc5F0CNfl0VLgL64aoTtu7xsM0YajW2GlG9+E= cloud.google.com/go/trace v1.16.0/go.mod h1:r+bdAn16dKLSV1G2D5v3e58IlQlizfxWrUfjx7kM7X0= github.com/GoogleCloudPlatform/opentelemetry-operations-go/detectors/gcp v1.32.0 h1:rIkQfkCOVKc1OiRCNcSDD8ml5RJlZbH/Xsq7lbpynwc= diff --git a/ci/resources/stemcell-version-bump/vendor/cloud.google.com/go/storage/CHANGES.md b/ci/resources/stemcell-version-bump/vendor/cloud.google.com/go/storage/CHANGES.md index b9000d87..3ddb1f06 100644 --- a/ci/resources/stemcell-version-bump/vendor/cloud.google.com/go/storage/CHANGES.md +++ b/ci/resources/stemcell-version-bump/vendor/cloud.google.com/go/storage/CHANGES.md @@ -1,6 +1,13 @@ # Changes +## [1.64.0](https://github.com/googleapis/google-cloud-go/compare/storage/v1.63.1...storage/v1.64.0) (2026-07-21) + + +### Features + +* **storage:** Accept CRC32C for appendable objects ([#20104](https://github.com/googleapis/google-cloud-go/issues/20104)) ([1b2f5af](https://github.com/googleapis/google-cloud-go/commit/1b2f5afa0fe0f159edbc7184467928ebbd8720e2)) + ## [1.63.1](https://github.com/googleapis/google-cloud-go/compare/storage/v1.63.0...storage/v1.63.1) (2026-07-13) diff --git a/ci/resources/stemcell-version-bump/vendor/cloud.google.com/go/storage/grpc_client.go b/ci/resources/stemcell-version-bump/vendor/cloud.google.com/go/storage/grpc_client.go index fb2eedc1..16adf561 100644 --- a/ci/resources/stemcell-version-bump/vendor/cloud.google.com/go/storage/grpc_client.go +++ b/ci/resources/stemcell-version-bump/vendor/cloud.google.com/go/storage/grpc_client.go @@ -179,21 +179,16 @@ func newGRPCStorageClient(ctx context.Context, opts ...storageOption) (client *g var metricsCleanup func() if isOtelMetricsEnabled(&config) { var project string - c, err := transport.Creds(ctx, s.clientOption...) - if err == nil { + if c, err := transport.Creds(ctx, s.clientOption...); err == nil { project = c.ProjectID } - if cm, cleanup, err := initMetrics(ctx, project, &config); err == nil { - clientMetrics = cm - metricsCleanup = cleanup - - unaryInt, streamInt := metricsInterceptors(cm) + clientMetrics, metricsCleanup = initClientMetrics(ctx, project, &config) + if clientMetrics != nil { + unaryInt, streamInt := metricsInterceptors(clientMetrics) s.clientOption = append(s.clientOption, option.WithGRPCDialOption(grpc.WithChainUnaryInterceptor(unaryInt)), option.WithGRPCDialOption(grpc.WithChainStreamInterceptor(streamInt)), ) - } else { - log.Printf("Failed to enable metrics: %v", err) } } @@ -369,10 +364,12 @@ func (c *grpcStorageClient) ListBuckets(ctx context.Context, project string, opt var gitr *gapic.BucketIterator fetch := func(pageSize int, pageToken string) (token string, err error) { + ctx, record := startMetricsOp(it.ctx, "ListBuckets", false) + defer func() { record(err) }() var buckets []*storagepb.Bucket var next string - err = run(it.ctx, func(ctx context.Context) error { + err = run(ctx, func(ctx context.Context) error { // Initialize GAPIC-based iterator when pageToken is empty, which // indicates that this fetch call is attempting to get the first page. // @@ -611,12 +608,14 @@ func (c *grpcStorageClient) ListObjects(ctx context.Context, bucket string, q *Q Filter: it.query.Filter, } fetch := func(pageSize int, pageToken string) (token string, err error) { + ctx, record := startMetricsOp(it.ctx, "ListObjects", false) + defer func() { record(err) }() // Add trace span around List API call within the fetch. ctx, _ = startSpan(ctx, "grpcStorageClient.ObjectsListCall") defer func() { endSpan(ctx, err) }() var objects []*storagepb.Object var gitr *gapic.ObjectIterator - err = run(it.ctx, func(ctx context.Context) error { + err = run(ctx, func(ctx context.Context) error { gitr = c.raw.ListObjects(ctx, req, s.gax...) objects, token, err = gitr.InternalFetch(pageSize, pageToken) return err diff --git a/ci/resources/stemcell-version-bump/vendor/cloud.google.com/go/storage/grpc_writer.go b/ci/resources/stemcell-version-bump/vendor/cloud.google.com/go/storage/grpc_writer.go index 2f646732..b8e56c8a 100644 --- a/ci/resources/stemcell-version-bump/vendor/cloud.google.com/go/storage/grpc_writer.go +++ b/ci/resources/stemcell-version-bump/vendor/cloud.google.com/go/storage/grpc_writer.go @@ -58,7 +58,8 @@ func (w *gRPCWriter) Write(p []byte) (n int, err error) { // Skip checksum calculation if user configures MD5 or CRC32C themselves. if !w.disableAutoChecksum && !w.sendCRC32C && - !md5Provided { + !md5Provided && + !w.append { w.fullObjectChecksum = crc32.Update(w.fullObjectChecksum, crc32cTable, p) } // write command successfully delivered to sender. We no longer own cmd. @@ -109,6 +110,11 @@ func (w *gRPCWriter) CloseWithError(err error) error { return nil } +func (w *gRPCWriter) setAppendFinalCRC32C(sendAppendFinalCRC32C bool, c uint32) { + w.sendAppendFinalCRC32C = sendAppendFinalCRC32C + w.appendFinalCRC32C = c +} + func (c *grpcStorageClient) OpenWriter(params *openWriterParams, opts ...storageOption) (internalWriter, error) { if params.attrs.Retention != nil { // TO-DO: remove once ObjectRetention is available - see b/308194853 @@ -259,7 +265,9 @@ type gRPCWriter struct { setSize func(int64) setTakeoverOffset func(int64) - fullObjectChecksum uint32 + fullObjectChecksum uint32 + appendFinalCRC32C uint32 + sendAppendFinalCRC32C bool flushSupported bool sendCRC32C bool @@ -948,9 +956,9 @@ type getObjectChecksumsParams struct { sendCRC32C bool disableAutoChecksum bool objectAttrs *ObjectAttrs - fullObjectChecksum func() uint32 + fullObjectChecksum func() *uint32 finishWrite bool - takeoverWriter bool + append bool } // getObjectChecksums determines what checksum information to include in the final @@ -965,20 +973,30 @@ func getObjectChecksums(params *getObjectChecksumsParams) *storagepb.ObjectCheck return nil } + // For append operations, send user's final append checksum on last write op if available. + // Auto checksum is not supported for appendable writes. + var crc32c *uint32 + if params.fullObjectChecksum != nil { + crc32c = params.fullObjectChecksum() + } + + if params.append && crc32c != nil { + return &storagepb.ObjectChecksums{Crc32C: crc32c} + } + // send user's checksum on last write op if available if params.sendCRC32C || (params.objectAttrs != nil && params.objectAttrs.MD5 != nil) { return toProtoChecksums(params.sendCRC32C, params.objectAttrs) } - // TODO(b/461982277): Enable checksum validation for appendable takeover writer gRPC - if params.disableAutoChecksum || params.takeoverWriter { - return nil - } - if params.fullObjectChecksum == nil { + + if params.append || params.disableAutoChecksum || params.fullObjectChecksum == nil { return nil } - return &storagepb.ObjectChecksums{ - Crc32C: proto.Uint32(params.fullObjectChecksum()), + + if crc32c != nil { + return &storagepb.ObjectChecksums{Crc32C: crc32c} } + return nil } type gRPCBidiWriteBufferSender interface { @@ -1014,7 +1032,7 @@ type gRPCOneshotBidiWriteBufferSender struct { sendCRC32C bool disableAutoChecksum bool objectAttrs *ObjectAttrs - fullObjectChecksum func() uint32 + fullObjectChecksum func() *uint32 } func (w *gRPCWriter) newGRPCOneshotBidiWriteBufferSender() *gRPCOneshotBidiWriteBufferSender { @@ -1031,8 +1049,9 @@ func (w *gRPCWriter) newGRPCOneshotBidiWriteBufferSender() *gRPCOneshotBidiWrite sendCRC32C: w.sendCRC32C, disableAutoChecksum: w.disableAutoChecksum, objectAttrs: w.attrs, - fullObjectChecksum: func() uint32 { - return w.fullObjectChecksum + fullObjectChecksum: func() *uint32 { + checksum := w.fullObjectChecksum + return &checksum }, } } @@ -1151,7 +1170,7 @@ type gRPCResumableBidiWriteBufferSender struct { sendCRC32C bool disableAutoChecksum bool objectAttrs *ObjectAttrs - fullObjectChecksum func() uint32 + fullObjectChecksum func() *uint32 streamErr error } @@ -1168,8 +1187,9 @@ func (w *gRPCWriter) newGRPCResumableBidiWriteBufferSender() *gRPCResumableBidiW sendCRC32C: w.sendCRC32C, disableAutoChecksum: w.disableAutoChecksum, objectAttrs: w.attrs, - fullObjectChecksum: func() uint32 { - return w.fullObjectChecksum + fullObjectChecksum: func() *uint32 { + checksum := w.fullObjectChecksum + return &checksum }, } } @@ -1299,7 +1319,7 @@ type gRPCAppendBidiWriteBufferSender struct { sendCRC32C bool disableAutoChecksum bool objectAttrs *ObjectAttrs - fullObjectChecksum func() uint32 + fullObjectChecksum func() *uint32 takeoverWriter bool @@ -1323,8 +1343,12 @@ func (w *gRPCWriter) newGRPCAppendableObjectBufferSender() *gRPCAppendBidiWriteB sendCRC32C: w.sendCRC32C, disableAutoChecksum: w.disableAutoChecksum, objectAttrs: w.attrs, - fullObjectChecksum: func() uint32 { - return w.fullObjectChecksum + fullObjectChecksum: func() *uint32 { + if !w.sendAppendFinalCRC32C { + return nil + } + checksum := w.appendFinalCRC32C + return &checksum }, } } @@ -1439,8 +1463,12 @@ func (w *gRPCWriter) newGRPCAppendTakeoverWriteBufferSender() *gRPCAppendTakeove sendCRC32C: w.sendCRC32C, disableAutoChecksum: w.disableAutoChecksum, objectAttrs: w.attrs, - fullObjectChecksum: func() uint32 { - return w.fullObjectChecksum + fullObjectChecksum: func() *uint32 { + if !w.sendAppendFinalCRC32C { + return nil + } + checksum := w.appendFinalCRC32C + return &checksum }, }, takeoverReported: false, @@ -1582,7 +1610,7 @@ func (s *gRPCAppendBidiWriteBufferSender) send(stream storagepb.Storage_BidiWrit fullObjectChecksum: s.fullObjectChecksum, disableAutoChecksum: s.disableAutoChecksum, finishWrite: finalizeObject, - takeoverWriter: s.takeoverWriter, + append: true, }) req := bidiWriteObjectRequest(r, bufChecksum, objectChecksums) if sendFirstMessage { diff --git a/ci/resources/stemcell-version-bump/vendor/cloud.google.com/go/storage/http_client.go b/ci/resources/stemcell-version-bump/vendor/cloud.google.com/go/storage/http_client.go index daa0daf1..62971346 100644 --- a/ci/resources/stemcell-version-bump/vendor/cloud.google.com/go/storage/http_client.go +++ b/ci/resources/stemcell-version-bump/vendor/cloud.google.com/go/storage/http_client.go @@ -131,14 +131,8 @@ func newHTTPStorageClient(ctx context.Context, opts ...storageOption) (client st if creds != nil { project, _ = creds.ProjectID(ctx) } - if cm, cleanup, err := initMetrics(ctx, project, &config); err == nil { - clientMetrics = cm - metricsCleanup = cleanup - } else { - log.Printf("Failed to enable metrics: %v", err) - } + clientMetrics, metricsCleanup = initClientMetrics(ctx, project, &config) } - if metricsCleanup != nil { defer func() { if err != nil { @@ -310,6 +304,8 @@ func (c *httpStorageClient) ListBuckets(ctx context.Context, project string, opt } fetch := func(pageSize int, pageToken string) (token string, err error) { + ctx, record := startMetricsOp(it.ctx, "ListBuckets", true) + defer func() { record(err) }() req := c.raw.Buckets.List(it.projectID) req.Projection("full") req.Prefix(it.Prefix) @@ -319,15 +315,16 @@ func (c *httpStorageClient) ListBuckets(ctx context.Context, project string, opt req.MaxResults(int64(pageSize)) } var resp *raw.Buckets - err = run(it.ctx, func(ctx context.Context) error { + err = run(ctx, func(ctx context.Context) error { resp, err = req.Context(ctx).Do() return err }, s.retry, s.idempotent) if err != nil { return "", err } + var b *BucketAttrs for _, item := range resp.Items { - b, err := newBucket(item) + b, err = newBucket(item) if err != nil { return "", err } @@ -436,6 +433,8 @@ func (c *httpStorageClient) ListObjects(ctx context.Context, bucket string, q *Q } fetch := func(pageSize int, pageToken string) (string, error) { var err error + ctx, record := startMetricsOp(it.ctx, "ListObjects", true) + defer func() { record(err) }() // Add trace span around List API call within the fetch. ctx, _ = startSpan(ctx, "httpStorageClient.ObjectsListCall") defer func() { endSpan(ctx, err) }() @@ -473,7 +472,7 @@ func (c *httpStorageClient) ListObjects(ctx context.Context, bucket string, q *Q req.MaxResults(int64(pageSize)) } var resp *raw.Objects - err = run(it.ctx, func(ctx context.Context) error { + err = run(ctx, func(ctx context.Context) error { resp, err = req.Context(ctx).Do() return err }, s.retry, s.idempotent, withOperation("ListObjects"), withObject(it.query.Prefix), withBucket(bucket)) @@ -1121,6 +1120,9 @@ func (hiw *httpInternalWriter) Flush() (int64, error) { return 0, errors.New("Writer.Flush is only supported for gRPC-based clients") } +// Not supported on HTTP Client as this is for setting CRC on appendable objects. +func (hiw *httpInternalWriter) setAppendFinalCRC32C(sendAppendFinalCRC32C bool, c uint32) {} + func (c *httpStorageClient) OpenWriter(params *openWriterParams, opts ...storageOption) (internalWriter, error) { if params.append { return nil, errors.New("storage: append not supported on HTTP Client; use gRPC") @@ -1325,6 +1327,8 @@ func (c *httpStorageClient) ListHMACKeys(ctx context.Context, project, serviceAc retry: s.retry, } fetch := func(pageSize int, pageToken string) (token string, err error) { + ctx, record := startMetricsOp(it.ctx, "ListHMACKeys", true) + defer func() { record(err) }() call := c.raw.Projects.HmacKeys.List(project) if pageToken != "" { call = call.PageToken(pageToken) @@ -1343,7 +1347,7 @@ func (c *httpStorageClient) ListHMACKeys(ctx context.Context, project, serviceAc } var resp *raw.HmacKeysMetadata - err = run(it.ctx, func(ctx context.Context) error { + err = run(ctx, func(ctx context.Context) error { resp, err = call.Context(ctx).Do() return err }, s.retry, s.idempotent) @@ -1351,11 +1355,12 @@ func (c *httpStorageClient) ListHMACKeys(ctx context.Context, project, serviceAc return "", err } + var hkey *HMACKey for _, metadata := range resp.Items { hk := &raw.HmacKey{ Metadata: metadata, } - hkey, err := toHMACKeyFromRaw(hk, true) + hkey, err = toHMACKeyFromRaw(hk, true) if err != nil { return "", err } diff --git a/ci/resources/stemcell-version-bump/vendor/cloud.google.com/go/storage/internal/version.go b/ci/resources/stemcell-version-bump/vendor/cloud.google.com/go/storage/internal/version.go index 86da2df7..6aaa3e33 100644 --- a/ci/resources/stemcell-version-bump/vendor/cloud.google.com/go/storage/internal/version.go +++ b/ci/resources/stemcell-version-bump/vendor/cloud.google.com/go/storage/internal/version.go @@ -17,4 +17,4 @@ package internal // Version is the current tagged release of the library. -const Version = "1.63.1" +const Version = "1.64.0" diff --git a/ci/resources/stemcell-version-bump/vendor/cloud.google.com/go/storage/metrics.go b/ci/resources/stemcell-version-bump/vendor/cloud.google.com/go/storage/metrics.go index e79787f7..d18045d6 100644 --- a/ci/resources/stemcell-version-bump/vendor/cloud.google.com/go/storage/metrics.go +++ b/ci/resources/stemcell-version-bump/vendor/cloud.google.com/go/storage/metrics.go @@ -16,16 +16,20 @@ package storage import ( "context" + "errors" "fmt" "io" + "log" "net" "net/http" "os" "strconv" "strings" + "sync" "sync/atomic" "time" + "cloud.google.com/go/iam/apiv1/iampb" "cloud.google.com/go/storage/internal" mexporter "github.com/GoogleCloudPlatform/opentelemetry-operations-go/exporter/metric" "go.opentelemetry.io/otel/attribute" @@ -33,6 +37,7 @@ import ( sdkmetric "go.opentelemetry.io/otel/sdk/metric" "go.opentelemetry.io/otel/sdk/metric/metricdata" "go.opentelemetry.io/otel/sdk/resource" + "google.golang.org/api/googleapi" "google.golang.org/grpc" "google.golang.org/grpc/codes" "google.golang.org/grpc/status" @@ -47,6 +52,13 @@ type clientMetrics struct { provider *sdkmetric.MeterProvider rpcClientCallDuration metric.Float64Histogram httpClientRequestDuration metric.Float64Histogram + duration metric.Float64Histogram + operations metric.Int64Counter + attempts metric.Int64Counter + requestBodySize metric.Int64Histogram + responseBodySize metric.Int64Histogram + ttfb metric.Float64Histogram + errors metric.Int64Counter } func formatMetricWithPrefix(m metricdata.Metrics, prefix string) string { @@ -134,6 +146,22 @@ func initMetrics(ctx context.Context, projectID string, config *storageConfig) ( sdkmetric.Instrument{Name: "http.client.request.duration", Kind: sdkmetric.InstrumentKindHistogram}, sdkmetric.Stream{Aggregation: sdkmetric.AggregationExplicitBucketHistogram{Boundaries: latencyHistogramBoundaries()}}, ), + sdkmetric.NewView( + sdkmetric.Instrument{Name: "gcp.client.request.duration", Kind: sdkmetric.InstrumentKindHistogram}, + sdkmetric.Stream{Aggregation: sdkmetric.AggregationExplicitBucketHistogram{Boundaries: latencyHistogramBoundaries()}}, + ), + sdkmetric.NewView( + sdkmetric.Instrument{Name: "gcp.storage.client.operation.ttfb", Kind: sdkmetric.InstrumentKindHistogram}, + sdkmetric.Stream{Aggregation: sdkmetric.AggregationExplicitBucketHistogram{Boundaries: latencyHistogramBoundaries()}}, + ), + sdkmetric.NewView( + sdkmetric.Instrument{Name: "gcp.storage.client.request.body.size", Kind: sdkmetric.InstrumentKindHistogram}, + sdkmetric.Stream{Aggregation: sdkmetric.AggregationExplicitBucketHistogram{Boundaries: sizeHistogramBoundaries()}}, + ), + sdkmetric.NewView( + sdkmetric.Instrument{Name: "gcp.storage.client.response.body.size", Kind: sdkmetric.InstrumentKindHistogram}, + sdkmetric.Stream{Aggregation: sdkmetric.AggregationExplicitBucketHistogram{Boundaries: sizeHistogramBoundaries()}}, + ), ), ) ownProvider = true @@ -159,10 +187,80 @@ func initMetrics(ctx context.Context, projectID string, config *storageConfig) ( return nil, nil, err } + duration, err := meter.Float64Histogram( + "gcp.client.request.duration", + metric.WithDescription("Latency of a client operation"), + metric.WithUnit("s"), + ) + if err != nil { + return nil, nil, err + } + + operations, err := meter.Int64Counter( + "gcp.storage.client.operations", + metric.WithDescription("Number of GCS client operations"), + metric.WithUnit("1"), + ) + if err != nil { + return nil, nil, err + } + + attempts, err := meter.Int64Counter( + "gcp.storage.client.attempts", + metric.WithDescription("Number of GCS client attempts"), + metric.WithUnit("1"), + ) + if err != nil { + return nil, nil, err + } + + requestBodySize, err := meter.Int64Histogram( + "gcp.storage.client.request.body.size", + metric.WithDescription("Size of GCS client request body"), + metric.WithUnit("By"), + ) + if err != nil { + return nil, nil, err + } + + responseBodySize, err := meter.Int64Histogram( + "gcp.storage.client.response.body.size", + metric.WithDescription("Size of GCS client response body"), + metric.WithUnit("By"), + ) + if err != nil { + return nil, nil, err + } + + ttfb, err := meter.Float64Histogram( + "gcp.storage.client.operation.ttfb", + metric.WithDescription("Time to first byte of GCS client operations"), + metric.WithUnit("s"), + ) + if err != nil { + return nil, nil, err + } + + errors, err := meter.Int64Counter( + "gcp.storage.client.errors", + metric.WithDescription("Number of GCS client errors"), + metric.WithUnit("1"), + ) + if err != nil { + return nil, nil, err + } + cm := &clientMetrics{ provider: provider, rpcClientCallDuration: rpcDuration, httpClientRequestDuration: httpDuration, + duration: duration, + operations: operations, + attempts: attempts, + requestBodySize: requestBodySize, + responseBodySize: responseBodySize, + ttfb: ttfb, + errors: errors, } var cleanup func() @@ -223,7 +321,7 @@ func grpcCodeToString(code codes.Code) string { func computeErrorType(err error, isHTTP bool, statusCode int64) string { if err == nil { if isHTTP && statusCode >= 400 { - return strconv.FormatInt(statusCode, 10) + return mapHTTPStatusCode(int(statusCode)) } return "OK" } @@ -260,13 +358,52 @@ func computeErrorType(err error, isHTTP bool, statusCode int64) string { } } - if isHTTP && statusCode >= 400 { - return strconv.FormatInt(statusCode, 10) + if isHTTP { + var apiErr *googleapi.Error + if errors.As(err, &apiErr) { + return mapHTTPStatusCode(apiErr.Code) + } + if statusCode >= 400 { + return mapHTTPStatusCode(int(statusCode)) + } } return "UNKNOWN" } +// mapHTTPStatusCode converts an HTTP status code to a canonical API error string. +// If there is no direct mapping, it returns the numeric string. +func mapHTTPStatusCode(code int) string { + switch code { + case 400: + return "INVALID_ARGUMENT" + case 401: + return "UNAUTHENTICATED" + case 403: + return "PERMISSION_DENIED" + case 404: + return "NOT_FOUND" + case 409: + return "ABORTED" + case 416: + return "OUT_OF_RANGE" + case 429: + return "RESOURCE_EXHAUSTED" + case 499: + return "CANCELLED" + case 500: + return "INTERNAL" + case 501: + return "UNIMPLEMENTED" + case 503: + return "UNAVAILABLE" + case 504: + return "DEADLINE_EXCEEDED" + default: + return strconv.Itoa(code) + } +} + func (cm *clientMetrics) recordRPC(ctx context.Context, method, target string, duration float64, err error) { statusCode := int64(codes.OK) if err != nil && err != io.EOF { @@ -296,6 +433,35 @@ func (cm *clientMetrics) recordRPC(ctx context.Context, method, target string, d } cm.rpcClientCallDuration.Record(ctx, duration, metric.WithAttributes(attrs...)) + + // Record standard attempt metric: gcp.storage.client.attempts. + state := metricsStateFromContext(ctx) + logicalMethod := methodName + if state != nil { + logicalMethod = state.method + } + attemptAttrs := []attribute.KeyValue{ + attribute.String("rpc.method", logicalMethod), + attribute.Int64("rpc.grpc.status_code", statusCode), + attribute.String("error.type", errorType), + } + cm.attempts.Add(ctx, 1, metric.WithAttributes(attemptAttrs...)) + + // Record standard error metric: gcp.storage.client.errors. + if err != nil && err != io.EOF { + errorAttrs := []attribute.KeyValue{ + attribute.String("rpc.method", logicalMethod), + attribute.String("error.type", errorType), + attribute.String("gcp.errors.domain", "storage.googleapis.com"), + } + cm.errors.Add(ctx, 1, metric.WithAttributes(errorAttrs...)) + } + + // For unary calls, record TTFB equal to the total attempt latency. + isStreaming := methodName == "ReadObject" || methodName == "WriteObject" || methodName == "BidiReadObject" || methodName == "BidiWriteObject" + if !isStreaming { + cm.ttfb.Record(ctx, duration, metric.WithAttributes(attribute.String("rpc.method", logicalMethod))) + } } func (cm *clientMetrics) recordHTTP(ctx context.Context, req *http.Request, resp *http.Response, duration float64, err error) { @@ -383,9 +549,54 @@ type metricsRoundTripper struct { func (rt *metricsRoundTripper) RoundTrip(req *http.Request) (*http.Response, error) { startTime := time.Now() resp, err := rt.base.RoundTrip(req) + + state := metricsStateFromContext(req.Context()) + var logicalMethod string + if state != nil { + logicalMethod = state.method + } else { + logicalMethod = "Unknown" + } + + statusCode := int64(0) + if resp != nil { + statusCode = int64(resp.StatusCode) + } + errorType := computeErrorType(err, true, statusCode) + + if rt.metrics != nil { + // Record attempt. + attemptAttrs := []attribute.KeyValue{ + attribute.String("rpc.method", logicalMethod), + attribute.Int64("http.response.status_code", statusCode), + attribute.String("error.type", errorType), + } + rt.metrics.attempts.Add(req.Context(), 1, metric.WithAttributes(attemptAttrs...)) + + // Record error if failed. + if err != nil || (resp != nil && resp.StatusCode >= 400) { + errorAttrs := []attribute.KeyValue{ + attribute.String("rpc.method", logicalMethod), + attribute.String("error.type", errorType), + attribute.String("gcp.errors.domain", "storage.googleapis.com"), + } + rt.metrics.errors.Add(req.Context(), 1, metric.WithAttributes(errorAttrs...)) + } + + // Record TTFB. + isDownload := req.Method == "GET" && req.URL.Query().Get("alt") == "media" + isResumableInit := req.Method == "POST" && strings.Contains(req.URL.Path, "/upload/") && req.URL.Query().Get("uploadType") == "resumable" + if !isDownload || isResumableInit { + duration := time.Since(startTime).Seconds() + rt.metrics.ttfb.Record(req.Context(), duration, metric.WithAttributes(attribute.String("rpc.method", logicalMethod))) + } + } + if err != nil { - duration := time.Since(startTime).Seconds() - rt.metrics.recordHTTP(req.Context(), req, nil, duration, err) + if rt.metrics != nil { + duration := time.Since(startTime).Seconds() + rt.metrics.recordHTTP(req.Context(), req, nil, duration, err) + } return nil, err } @@ -396,24 +607,38 @@ func (rt *metricsRoundTripper) RoundTrip(req *http.Request) (*http.Response, err req: req, resp: resp, metrics: rt.metrics, + isDownload: req.Method == "GET" && req.URL.Query().Get("alt") == "media", } } else { - duration := time.Since(startTime).Seconds() - rt.metrics.recordHTTP(req.Context(), req, resp, duration, nil) + if rt.metrics != nil { + duration := time.Since(startTime).Seconds() + rt.metrics.recordHTTP(req.Context(), req, resp, duration, nil) + } } return resp, nil } type wrappedResponseBody struct { io.ReadCloser - startTime time.Time - req *http.Request - resp *http.Response - metrics *clientMetrics - recorded atomic.Bool + startTime time.Time + req *http.Request + resp *http.Response + metrics *clientMetrics + recorded atomic.Bool + isDownload bool + firstRead atomic.Bool } func (w *wrappedResponseBody) Read(p []byte) (n int, err error) { + if w.isDownload && w.metrics != nil && w.firstRead.CompareAndSwap(false, true) { + duration := time.Since(w.startTime).Seconds() + state := metricsStateFromContext(w.req.Context()) + logicalMethod := "ReadObject" + if state != nil { + logicalMethod = state.method + } + w.metrics.ttfb.Record(w.req.Context(), duration, metric.WithAttributes(attribute.String("rpc.method", logicalMethod))) + } n, err = w.ReadCloser.Read(p) if err != nil { w.record(err) @@ -486,10 +711,14 @@ type wrappedClientStream struct { recorded atomic.Bool serverStreams bool clientStreams bool + recordedTTFB atomic.Bool } func (w *wrappedClientStream) RecvMsg(m interface{}) error { err := w.ClientStream.RecvMsg(m) + if err == nil { + w.recordTTFB(m) + } // For client-streaming streams (like WriteObject), the single successful RecvMsg call // returns the response and nil error, which marks the completion of the stream. isClientStreaming := !w.serverStreams && w.clientStreams @@ -513,3 +742,400 @@ func (w *wrappedClientStream) record(err error) { w.metrics.recordRPC(w.ctx, w.method, w.target, duration, err) } } + +func (w *wrappedClientStream) recordTTFB(m interface{}) { + if w.recordedTTFB.Load() { + return + } + methodName := w.method + if strings.HasPrefix(methodName, "/") { + parts := strings.Split(strings.TrimPrefix(methodName, "/"), "/") + if len(parts) >= 2 { + methodName = parts[1] + } + } + + // The first response from the server, whether it contains metadata, + // persisted size, or actual content, indicates TTFB. + if w.recordedTTFB.CompareAndSwap(false, true) { + duration := time.Since(w.startTime).Seconds() + state := metricsStateFromContext(w.ctx) + logicalMethod := methodName + if state != nil { + logicalMethod = state.method + } + w.metrics.ttfb.Record(w.ctx, duration, metric.WithAttributes(attribute.String("rpc.method", logicalMethod))) + } +} + +type metricsKey struct{} + +type metricsState struct { + method string + startTime time.Time + metrics *clientMetrics + isHTTP bool + ttfbRecorded atomic.Bool + ttfbStart time.Time + record func(error) +} + +func contextWithMetricsState(ctx context.Context, state *metricsState) context.Context { + return context.WithValue(ctx, metricsKey{}, state) +} + +func metricsStateFromContext(ctx context.Context) *metricsState { + if ctx == nil { + return nil + } + if state, ok := ctx.Value(metricsKey{}).(*metricsState); ok { + return state + } + return nil +} + +func contextWithoutMetrics(ctx context.Context) context.Context { + if ctx == nil { + return nil + } + return context.WithValue(ctx, metricsKey{}, (*metricsState)(nil)) +} + +func (cm *clientMetrics) startOperation(ctx context.Context, method string, isHTTP bool) (context.Context, func(error)) { + if cm == nil { + return ctx, func(error) {} + } + state := &metricsState{ + method: method, + startTime: time.Now(), + metrics: cm, + isHTTP: isHTTP, + } + state.ttfbStart = state.startTime + + var recordOnce sync.Once + record := func(err error) { + recordOnce.Do(func() { + duration := time.Since(state.startTime).Seconds() + statusStr := "OK" + if err != nil && err != io.EOF { + statusStr = "Error" + } + errorType := computeErrorType(err, isHTTP, 0) + + attrs := []attribute.KeyValue{ + attribute.String("rpc.method", method), + attribute.String("status", statusStr), + attribute.String("error.type", errorType), + } + opts := metric.WithAttributes(attrs...) + cm.duration.Record(ctx, duration, opts) + cm.operations.Add(ctx, 1, opts) + }) + } + state.record = record + + ctx = contextWithMetricsState(ctx, state) + return ctx, record +} + +// startMetricsOp starts a client operation if OpenTelemetry metrics are enabled in ctx. +// It returns the updated context containing metrics state and a recording closure. +func startMetricsOp(ctx context.Context, method string, isHTTP bool) (context.Context, func(error)) { + if state := metricsStateFromContext(ctx); state != nil && state.metrics != nil { + return state.metrics.startOperation(ctx, method, isHTTP) + } + return ctx, func(error) {} +} + +// initClientMetrics initializes OpenTelemetry client metrics if enabled in config. +// It returns the metrics instance and its cleanup function, or nil if disabled or upon error. +func initClientMetrics(ctx context.Context, project string, config *storageConfig) (*clientMetrics, func()) { + if !isOtelMetricsEnabled(config) { + return nil, nil + } + cm, cleanup, err := initMetrics(ctx, project, config) + if err != nil { + log.Printf("Failed to enable metrics: %v", err) + return nil, nil + } + return cm, cleanup +} + +// metricsStorageClient wraps a storageClient and records client-level metrics. +type metricsStorageClient struct { + storageClient + metrics *clientMetrics + isHTTP bool +} + +func (mc *metricsStorageClient) GetServiceAccount(ctx context.Context, project string, opts ...storageOption) (string, error) { + ctx, record := mc.metrics.startOperation(ctx, "GetServiceAccount", mc.isHTTP) + res, err := mc.storageClient.GetServiceAccount(ctx, project, opts...) + record(err) + return res, err +} + +func (mc *metricsStorageClient) CreateBucket(ctx context.Context, project, bucket string, attrs *BucketAttrs, enableObjectRetention *bool, opts ...storageOption) (*BucketAttrs, error) { + ctx, record := mc.metrics.startOperation(ctx, "CreateBucket", mc.isHTTP) + res, err := mc.storageClient.CreateBucket(ctx, project, bucket, attrs, enableObjectRetention, opts...) + record(err) + return res, err +} + +func (mc *metricsStorageClient) ListBuckets(ctx context.Context, project string, opts ...storageOption) *BucketIterator { + ctx, _ = mc.metrics.startOperation(ctx, "ListBuckets", mc.isHTTP) + return mc.storageClient.ListBuckets(ctx, project, opts...) +} + +func (mc *metricsStorageClient) DeleteBucket(ctx context.Context, bucket string, conds *BucketConditions, opts ...storageOption) error { + ctx, record := mc.metrics.startOperation(ctx, "DeleteBucket", mc.isHTTP) + err := mc.storageClient.DeleteBucket(ctx, bucket, conds, opts...) + record(err) + return err +} + +func (mc *metricsStorageClient) GetBucket(ctx context.Context, bucket string, conds *BucketConditions, opts ...storageOption) (*BucketAttrs, error) { + ctx, record := mc.metrics.startOperation(ctx, "GetBucket", mc.isHTTP) + res, err := mc.storageClient.GetBucket(ctx, bucket, conds, opts...) + record(err) + return res, err +} + +func (mc *metricsStorageClient) UpdateBucket(ctx context.Context, bucket string, uattrs *BucketAttrsToUpdate, conds *BucketConditions, opts ...storageOption) (*BucketAttrs, error) { + ctx, record := mc.metrics.startOperation(ctx, "UpdateBucket", mc.isHTTP) + res, err := mc.storageClient.UpdateBucket(ctx, bucket, uattrs, conds, opts...) + record(err) + return res, err +} + +func (mc *metricsStorageClient) LockBucketRetentionPolicy(ctx context.Context, bucket string, conds *BucketConditions, opts ...storageOption) error { + ctx, record := mc.metrics.startOperation(ctx, "LockBucketRetentionPolicy", mc.isHTTP) + err := mc.storageClient.LockBucketRetentionPolicy(ctx, bucket, conds, opts...) + record(err) + return err +} + +func (mc *metricsStorageClient) ListObjects(ctx context.Context, bucket string, q *Query, opts ...storageOption) *ObjectIterator { + ctx, _ = mc.metrics.startOperation(ctx, "ListObjects", mc.isHTTP) + return mc.storageClient.ListObjects(ctx, bucket, q, opts...) +} + +func (mc *metricsStorageClient) DeleteObject(ctx context.Context, bucket, object string, gen int64, conds *Conditions, opts ...storageOption) error { + ctx, record := mc.metrics.startOperation(ctx, "DeleteObject", mc.isHTTP) + err := mc.storageClient.DeleteObject(ctx, bucket, object, gen, conds, opts...) + record(err) + return err +} + +func (mc *metricsStorageClient) GetObject(ctx context.Context, params *getObjectParams, opts ...storageOption) (*ObjectAttrs, error) { + ctx, record := mc.metrics.startOperation(ctx, "GetObject", mc.isHTTP) + res, err := mc.storageClient.GetObject(ctx, params, opts...) + record(err) + return res, err +} + +func (mc *metricsStorageClient) UpdateObject(ctx context.Context, params *updateObjectParams, opts ...storageOption) (*ObjectAttrs, error) { + ctx, record := mc.metrics.startOperation(ctx, "UpdateObject", mc.isHTTP) + res, err := mc.storageClient.UpdateObject(ctx, params, opts...) + record(err) + return res, err +} + +func (mc *metricsStorageClient) RestoreObject(ctx context.Context, params *restoreObjectParams, opts ...storageOption) (*ObjectAttrs, error) { + ctx, record := mc.metrics.startOperation(ctx, "RestoreObject", mc.isHTTP) + res, err := mc.storageClient.RestoreObject(ctx, params, opts...) + record(err) + return res, err +} + +func (mc *metricsStorageClient) MoveObject(ctx context.Context, params *moveObjectParams, opts ...storageOption) (*ObjectAttrs, error) { + ctx, record := mc.metrics.startOperation(ctx, "MoveObject", mc.isHTTP) + res, err := mc.storageClient.MoveObject(ctx, params, opts...) + record(err) + return res, err +} + +func (mc *metricsStorageClient) ComposeObject(ctx context.Context, req *composeObjectRequest, opts ...storageOption) (*ObjectAttrs, error) { + ctx, record := mc.metrics.startOperation(ctx, "ComposeObject", mc.isHTTP) + res, err := mc.storageClient.ComposeObject(ctx, req, opts...) + record(err) + return res, err +} + +func (mc *metricsStorageClient) RewriteObject(ctx context.Context, req *rewriteObjectRequest, opts ...storageOption) (*rewriteObjectResponse, error) { + ctx, record := mc.metrics.startOperation(ctx, "RewriteObject", mc.isHTTP) + res, err := mc.storageClient.RewriteObject(ctx, req, opts...) + record(err) + return res, err +} + +func (mc *metricsStorageClient) NewRangeReader(ctx context.Context, params *newRangeReaderParams, opts ...storageOption) (*Reader, error) { + ctx, record := mc.metrics.startOperation(ctx, "ReadObject", mc.isHTTP) + r, err := mc.storageClient.NewRangeReader(ctx, params, opts...) + if err != nil { + record(err) + return nil, err + } + if state := metricsStateFromContext(ctx); state != nil { + r.metricsState = state + } + return r, nil +} + +func (mc *metricsStorageClient) OpenWriter(params *openWriterParams, opts ...storageOption) (internalWriter, error) { + ctx, _ := mc.metrics.startOperation(params.ctx, "WriteObject", mc.isHTTP) + params.ctx = ctx + return mc.storageClient.OpenWriter(params, opts...) +} + +func (mc *metricsStorageClient) NewMultiRangeDownloader(ctx context.Context, params *newMultiRangeDownloaderParams, opts ...storageOption) (*MultiRangeDownloader, error) { + ctx, _ = mc.metrics.startOperation(ctx, "ReadObject", mc.isHTTP) + return mc.storageClient.NewMultiRangeDownloader(ctx, params, opts...) +} + +func (mc *metricsStorageClient) GetIamPolicy(ctx context.Context, resource string, version int32, opts ...storageOption) (*iampb.Policy, error) { + ctx, record := mc.metrics.startOperation(ctx, "GetIamPolicy", mc.isHTTP) + res, err := mc.storageClient.GetIamPolicy(ctx, resource, version, opts...) + record(err) + return res, err +} + +func (mc *metricsStorageClient) SetIamPolicy(ctx context.Context, resource string, policy *iampb.Policy, opts ...storageOption) error { + ctx, record := mc.metrics.startOperation(ctx, "SetIamPolicy", mc.isHTTP) + err := mc.storageClient.SetIamPolicy(ctx, resource, policy, opts...) + record(err) + return err +} + +func (mc *metricsStorageClient) TestIamPermissions(ctx context.Context, resource string, permissions []string, opts ...storageOption) ([]string, error) { + ctx, record := mc.metrics.startOperation(ctx, "TestIamPermissions", mc.isHTTP) + res, err := mc.storageClient.TestIamPermissions(ctx, resource, permissions, opts...) + record(err) + return res, err +} + +func (mc *metricsStorageClient) GetHMACKey(ctx context.Context, project, accessID string, opts ...storageOption) (*HMACKey, error) { + ctx, record := mc.metrics.startOperation(ctx, "GetHMACKey", mc.isHTTP) + res, err := mc.storageClient.GetHMACKey(ctx, project, accessID, opts...) + record(err) + return res, err +} + +func (mc *metricsStorageClient) ListHMACKeys(ctx context.Context, project string, serviceAccountEmail string, showDeletedKeys bool, opts ...storageOption) *HMACKeysIterator { + ctx, _ = mc.metrics.startOperation(ctx, "ListHMACKeys", mc.isHTTP) + return mc.storageClient.ListHMACKeys(ctx, project, serviceAccountEmail, showDeletedKeys, opts...) +} + +func (mc *metricsStorageClient) UpdateHMACKey(ctx context.Context, project, serviceAccountEmail, accessID string, attrs *HMACKeyAttrsToUpdate, opts ...storageOption) (*HMACKey, error) { + ctx, record := mc.metrics.startOperation(ctx, "UpdateHMACKey", mc.isHTTP) + res, err := mc.storageClient.UpdateHMACKey(ctx, project, serviceAccountEmail, accessID, attrs, opts...) + record(err) + return res, err +} + +func (mc *metricsStorageClient) CreateHMACKey(ctx context.Context, project, serviceAccountEmail string, opts ...storageOption) (*HMACKey, error) { + ctx, record := mc.metrics.startOperation(ctx, "CreateHMACKey", mc.isHTTP) + res, err := mc.storageClient.CreateHMACKey(ctx, project, serviceAccountEmail, opts...) + record(err) + return res, err +} + +func (mc *metricsStorageClient) DeleteHMACKey(ctx context.Context, project, accessID string, opts ...storageOption) error { + ctx, record := mc.metrics.startOperation(ctx, "DeleteHMACKey", mc.isHTTP) + err := mc.storageClient.DeleteHMACKey(ctx, project, accessID, opts...) + record(err) + return err +} + +func (mc *metricsStorageClient) ListNotifications(ctx context.Context, bucket string, opts ...storageOption) (map[string]*Notification, error) { + ctx, record := mc.metrics.startOperation(ctx, "ListNotifications", mc.isHTTP) + res, err := mc.storageClient.ListNotifications(ctx, bucket, opts...) + record(err) + return res, err +} + +func (mc *metricsStorageClient) CreateNotification(ctx context.Context, bucket string, n *Notification, opts ...storageOption) (*Notification, error) { + ctx, record := mc.metrics.startOperation(ctx, "CreateNotification", mc.isHTTP) + res, err := mc.storageClient.CreateNotification(ctx, bucket, n, opts...) + record(err) + return res, err +} + +func (mc *metricsStorageClient) DeleteNotification(ctx context.Context, bucket string, id string, opts ...storageOption) error { + ctx, record := mc.metrics.startOperation(ctx, "DeleteNotification", mc.isHTTP) + err := mc.storageClient.DeleteNotification(ctx, bucket, id, opts...) + record(err) + return err +} + +func (mc *metricsStorageClient) DeleteDefaultObjectACL(ctx context.Context, bucket string, entity ACLEntity, opts ...storageOption) error { + ctx, record := mc.metrics.startOperation(ctx, "DeleteDefaultObjectACL", mc.isHTTP) + err := mc.storageClient.DeleteDefaultObjectACL(ctx, bucket, entity, opts...) + record(err) + return err +} + +func (mc *metricsStorageClient) ListDefaultObjectACLs(ctx context.Context, bucket string, opts ...storageOption) ([]ACLRule, error) { + ctx, record := mc.metrics.startOperation(ctx, "ListDefaultObjectACLs", mc.isHTTP) + res, err := mc.storageClient.ListDefaultObjectACLs(ctx, bucket, opts...) + record(err) + return res, err +} + +func (mc *metricsStorageClient) UpdateDefaultObjectACL(ctx context.Context, bucket string, entity ACLEntity, role ACLRole, opts ...storageOption) error { + ctx, record := mc.metrics.startOperation(ctx, "UpdateDefaultObjectACL", mc.isHTTP) + err := mc.storageClient.UpdateDefaultObjectACL(ctx, bucket, entity, role, opts...) + record(err) + return err +} + +func (mc *metricsStorageClient) DeleteBucketACL(ctx context.Context, bucket string, entity ACLEntity, opts ...storageOption) error { + ctx, record := mc.metrics.startOperation(ctx, "DeleteBucketACL", mc.isHTTP) + err := mc.storageClient.DeleteBucketACL(ctx, bucket, entity, opts...) + record(err) + return err +} + +func (mc *metricsStorageClient) ListBucketACLs(ctx context.Context, bucket string, opts ...storageOption) ([]ACLRule, error) { + ctx, record := mc.metrics.startOperation(ctx, "ListBucketACLs", mc.isHTTP) + res, err := mc.storageClient.ListBucketACLs(ctx, bucket, opts...) + record(err) + return res, err +} + +func (mc *metricsStorageClient) UpdateBucketACL(ctx context.Context, bucket string, entity ACLEntity, role ACLRole, opts ...storageOption) error { + ctx, record := mc.metrics.startOperation(ctx, "UpdateBucketACL", mc.isHTTP) + err := mc.storageClient.UpdateBucketACL(ctx, bucket, entity, role, opts...) + record(err) + return err +} + +func (mc *metricsStorageClient) DeleteObjectACL(ctx context.Context, bucket, object string, entity ACLEntity, opts ...storageOption) error { + ctx, record := mc.metrics.startOperation(ctx, "DeleteObjectACL", mc.isHTTP) + err := mc.storageClient.DeleteObjectACL(ctx, bucket, object, entity, opts...) + record(err) + return err +} + +func (mc *metricsStorageClient) ListObjectACLs(ctx context.Context, bucket, object string, opts ...storageOption) ([]ACLRule, error) { + ctx, record := mc.metrics.startOperation(ctx, "ListObjectACLs", mc.isHTTP) + res, err := mc.storageClient.ListObjectACLs(ctx, bucket, object, opts...) + record(err) + return res, err +} + +func (mc *metricsStorageClient) UpdateObjectACL(ctx context.Context, bucket, object string, entity ACLEntity, role ACLRole, opts ...storageOption) error { + ctx, record := mc.metrics.startOperation(ctx, "UpdateObjectACL", mc.isHTTP) + err := mc.storageClient.UpdateObjectACL(ctx, bucket, object, entity, role, opts...) + record(err) + return err +} + +func (mc *metricsStorageClient) Close() error { + return mc.storageClient.Close() +} + +func (mc *metricsStorageClient) fetchBucketMetadata(ctx context.Context, bucket string) (string, string, error) { + return mc.storageClient.fetchBucketMetadata(ctx, bucket) +} diff --git a/ci/resources/stemcell-version-bump/vendor/cloud.google.com/go/storage/pcu.go b/ci/resources/stemcell-version-bump/vendor/cloud.google.com/go/storage/pcu.go index 3bac1816..51b15acd 100644 --- a/ci/resources/stemcell-version-bump/vendor/cloud.google.com/go/storage/pcu.go +++ b/ci/resources/stemcell-version-bump/vendor/cloud.google.com/go/storage/pcu.go @@ -184,7 +184,8 @@ func (w *Writer) initPCU(ctx context.Context) error { // Track PCU operations using client feature tracking header. ctx = addFeatureAttributes(ctx, featurePCU) - pCtx, cancel := context.WithCancel(ctx) + bgCtx := contextWithoutMetrics(ctx) + pCtx, cancel := context.WithCancel(bgCtx) state := &pcuState{ ctx: pCtx, diff --git a/ci/resources/stemcell-version-bump/vendor/cloud.google.com/go/storage/reader.go b/ci/resources/stemcell-version-bump/vendor/cloud.google.com/go/storage/reader.go index c417dc56..de56f9f2 100644 --- a/ci/resources/stemcell-version-bump/vendor/cloud.google.com/go/storage/reader.go +++ b/ci/resources/stemcell-version-bump/vendor/cloud.google.com/go/storage/reader.go @@ -25,6 +25,9 @@ import ( "sync" "sync/atomic" "time" + + "go.opentelemetry.io/otel/attribute" + "go.opentelemetry.io/otel/metric" ) var crc32cTable = crc32.MakeTable(crc32.Castagnoli) @@ -373,7 +376,8 @@ var emptyBody = io.NopCloser(strings.NewReader("")) type Reader struct { // remain must be the first field in the struct to guarantee 64-bit // alignment on 32-bit architectures for atomic operations. - remain int64 + remain int64 + bytesRead int64 // Cumulative bytes read for response size metric. Attrs ReaderObjectAttrs objectMetadata *map[string]string @@ -384,11 +388,14 @@ type Reader struct { reader io.ReadCloser ctx context.Context mu sync.Mutex + err error // Persistent error encountered during Read or WriteTo. handle *ReadHandle unfinalized bool bucket string object string + + metricsState *metricsState } // Close closes the Reader. It must be called when done reading. @@ -399,15 +406,29 @@ func (r *Reader) Close() error { } } err := r.reader.Close() + r.mu.Lock() + if r.err != nil { + err = r.err + } + r.mu.Unlock() + + if r.metricsState != nil { + if r.metricsState.metrics != nil { + if total := atomic.SwapInt64(&r.bytesRead, 0); total > 0 { + r.metricsState.metrics.responseBodySize.Record(r.ctx, total, metric.WithAttributes(attribute.String("rpc.method", "ReadObject"))) + } + } + if r.metricsState.record != nil { + r.metricsState.record(err) + } + } endSpan(r.ctx, err) return err } func (r *Reader) Read(p []byte) (int, error) { n, err := r.reader.Read(p) - if !r.unfinalized && !r.Attrs.Decompressed { - atomic.AddInt64(&r.remain, -int64(n)) - } + r.recordRead(int64(n), err) return n, err } @@ -417,10 +438,24 @@ func (r *Reader) WriteTo(w io.Writer) (int64, error) { // This implicitly calls r.reader.WriteTo for gRPC only. JSON and XML don't have an // implementation of WriteTo. n, err := io.Copy(w, r.reader) + r.recordRead(n, err) + return n, err +} + +// recordRead updates remaining byte counts and metrics after a read operation, +// and saves any persistent error encountered during Read or WriteTo. +func (r *Reader) recordRead(n int64, err error) { if !r.unfinalized && !r.Attrs.Decompressed { - atomic.AddInt64(&r.remain, -int64(n)) + atomic.AddInt64(&r.remain, -n) + } + if n > 0 { + atomic.AddInt64(&r.bytesRead, n) + } + if err != nil && err != io.EOF { + r.mu.Lock() + r.err = err + r.mu.Unlock() } - return n, err } // Size returns the size of the object in bytes. @@ -554,6 +589,9 @@ func (mrd *MultiRangeDownloader) Add(output io.Writer, offset, length int64, cal // it could lead to a deadlock. func (mrd *MultiRangeDownloader) Close() error { err := mrd.impl.close(nil) + if state := metricsStateFromContext(mrd.impl.getSpanCtx()); state != nil && state.record != nil { + state.record(err) + } endSpan(mrd.impl.getSpanCtx(), err) return err } diff --git a/ci/resources/stemcell-version-bump/vendor/cloud.google.com/go/storage/storage.go b/ci/resources/stemcell-version-bump/vendor/cloud.google.com/go/storage/storage.go index dbc9b582..cb4d3680 100644 --- a/ci/resources/stemcell-version-bump/vendor/cloud.google.com/go/storage/storage.go +++ b/ci/resources/stemcell-version-bump/vendor/cloud.google.com/go/storage/storage.go @@ -226,13 +226,22 @@ func NewClient(ctx context.Context, opts ...option.ClientOption) (*Client, error return nil, fmt.Errorf("storage: %w", err) } + var tcWrapped storageClient = tc + if httpClient, ok := tc.(*httpStorageClient); ok && httpClient.metrics != nil { + tcWrapped = &metricsStorageClient{ + storageClient: tc, + metrics: httpClient.metrics, + isHTTP: true, + } + } + c := &Client{ hc: hc, raw: rawService, scheme: u.Scheme, xmlHost: u.Host, creds: creds, - tc: tc, + tc: tcWrapped, } if isACOEnabled() { c.bucketMetadataCache = newBucketMetadataCache(defaultBucketMetadataCacheLimit, c.tc) @@ -257,8 +266,16 @@ func NewGRPCClient(ctx context.Context, opts ...option.ClientOption) (*Client, e if err != nil { return nil, err } + var tcWrapped storageClient = tc + if tc.metrics != nil { + tcWrapped = &metricsStorageClient{ + storageClient: tc, + metrics: tc.metrics, + isHTTP: false, + } + } c := &Client{ - tc: tc, + tc: tcWrapped, grpcAppendableUploads: tc.config.grpcAppendableUploads, } if isACOEnabled() { @@ -1353,6 +1370,13 @@ type AppendableWriterOpts struct { ProgressFunc func(int64) // FinalizeOnClose: See Writer.FinalizeOnClose. FinalizeOnClose bool + // DisableAutoChecksum: See Writer.DisableAutoChecksum. + DisableAutoChecksum bool + // SendCRC32C: See Writer.SendCRC32C. + SendCRC32C bool + // CRC32C of the whole object. + // See Writer.CRC32C. + CRC32C uint32 } func (opts *AppendableWriterOpts) apply(w *Writer) { @@ -1363,6 +1387,9 @@ func (opts *AppendableWriterOpts) apply(w *Writer) { w.ProgressFunc = opts.ProgressFunc w.ChunkSize = opts.ChunkSize w.FinalizeOnClose = opts.FinalizeOnClose + w.DisableAutoChecksum = opts.DisableAutoChecksum + w.CRC32C = opts.CRC32C + w.SendCRC32C = opts.SendCRC32C } func (o *ObjectHandle) validate() error { @@ -1580,6 +1607,7 @@ type ObjectAttrs struct { // MD5 is the MD5 hash of the object's content. This field is read-only, // except when used from a Writer. If set on a Writer, the uploaded // data is rejected if its MD5 hash does not match this field. + // Note: MD5 validation is not supported for appendable writes. MD5 []byte // CRC32C is the CRC32 checksum of the object's content using the Castagnoli93 diff --git a/ci/resources/stemcell-version-bump/vendor/cloud.google.com/go/storage/writer.go b/ci/resources/stemcell-version-bump/vendor/cloud.google.com/go/storage/writer.go index 78bae27a..62f53871 100644 --- a/ci/resources/stemcell-version-bump/vendor/cloud.google.com/go/storage/writer.go +++ b/ci/resources/stemcell-version-bump/vendor/cloud.google.com/go/storage/writer.go @@ -21,8 +21,12 @@ import ( "io" "log" "sync" + "sync/atomic" "time" "unicode/utf8" + + "go.opentelemetry.io/otel/attribute" + "go.opentelemetry.io/otel/metric" ) // Interface internalWriter wraps low-level implementations which may vary @@ -33,6 +37,9 @@ type internalWriter interface { // CloseWithError terminates the write operation and sets its status. // Note that CloseWithError always returns nil. CloseWithError(error) error + // Set final object checksum for appendable objects after write + // and before finalization. + setAppendFinalCRC32C(sendAppendFinalCRC32C bool, c uint32) } // A Writer writes a Cloud Storage object. @@ -200,6 +207,28 @@ type Writer struct { // in future releases. It is not yet recommended for production use. ParallelUploadConfig ParallelUploadConfig + // SendAppendFinalCRC32C indicates that AppendFinalCRC32C should be sent as the + // full-object checksum when finalizing an appendable object. + // + // This field is only supported for gRPC clients and appendable objects. + SendAppendFinalCRC32C bool + + // AppendFinalCRC32C is the expected full-object CRC32C checksum to be validated + // by the server when finalizing an appendable object. + // + // This field must be set on the Writer, along with SendAppendFinalCRC32C = true, + // before calling Close(). It allows callers to defer providing the checksum + // until the final write is complete. + // + // If SendAppendFinalCRC32C is false, checksum validation for the full object will + // fall back to using ObjectAttrs.CRC32C if Writer.SendCRC32C is true. If both are + // configured, AppendFinalCRC32C takes precedence. + // + // This field is ignored if Writer.Append is false or if using the JSON API. + // Chunk-level validation will still be performed for all chunks during upload if + // Writer.DisableAutoChecksum is false. + AppendFinalCRC32C uint32 + ctx context.Context o *ObjectHandle @@ -215,6 +244,9 @@ type Writer struct { mu sync.Mutex err error setTakeoverOffset func(int64) + + // bytesWritten is the cumulative bytes written for request size metric. + bytesWritten int64 } func (w *Writer) wrapWriteError(n int, err error) (int, error) { @@ -280,27 +312,36 @@ func (w *Writer) Write(p []byte) (int, error) { return 0, fmt.Errorf("storage: Writer is closed") } + var n int + var err error if pcu != nil { - return w.wrapWriteError(pcu.write(p)) - } - - if !w.opened { - // First time initialization: freeze the configuration to either PCU or standard. - if w.EnableParallelUpload { - var err error - if pcu, err = w.getOrInitPCU(); err != nil { - return 0, err + n, err = pcu.write(p) + } else { + if !w.opened { + // First time initialization: freeze the configuration to either PCU or standard. + if w.EnableParallelUpload { + if pcu, err = w.getOrInitPCU(); err != nil { + return 0, err + } } - if pcu != nil { - return w.wrapWriteError(pcu.write(p)) + if pcu == nil { + if err = w.openWriter(); err != nil { + return 0, err + } } } - if err := w.openWriter(); err != nil { - return 0, err + if pcu != nil { + n, err = pcu.write(p) + } else { + n, err = w.iw.Write(p) } } - return w.wrapWriteError(w.iw.Write(p)) + if n > 0 { + atomic.AddInt64(&w.bytesWritten, int64(n)) + } + + return w.wrapWriteError(n, err) } // Flush syncs all bytes currently in the Writer's buffer to Cloud Storage. @@ -361,38 +402,54 @@ func (w *Writer) Close() error { if pcu != nil || (!w.opened && w.EnableParallelUpload) { var err error if pcu, err = w.getOrInitPCU(); err != nil { - return err + return w.markClosed(err) } if pcu != nil { - err = pcu.close() - w.mu.Lock() - defer w.mu.Unlock() - w.closed = true - if w.err == nil && err != nil { - w.err = err - } - endSpan(w.ctx, w.err) - return w.err + return w.markClosed(pcu.close()) } } if !w.opened { if err := w.openWriter(); err != nil { - return err + return w.markClosed(err) } } + if w.Append { + w.iw.setAppendFinalCRC32C(w.SendAppendFinalCRC32C, w.AppendFinalCRC32C) + } + if err := w.iw.Close(); err != nil { - return err + return w.markClosed(err) } <-w.donec + return w.markClosed(nil) +} + +// markClosed marks the Writer as closed, records any closing error on Writer.err, +// and records request body size metrics and trace span completion. +func (w *Writer) markClosed(err error) error { w.mu.Lock() - defer w.mu.Unlock() w.closed = true - endSpan(w.ctx, w.err) - return w.err + if w.err == nil && err != nil { + w.err = err + } + closingErr := w.err + total := atomic.LoadInt64(&w.bytesWritten) + w.mu.Unlock() + + if state := metricsStateFromContext(w.ctx); state != nil { + if state.metrics != nil && total > 0 { + state.metrics.requestBodySize.Record(w.ctx, total, metric.WithAttributes(attribute.String("rpc.method", "WriteObject"))) + } + if state.record != nil { + state.record(closingErr) + } + } + endSpan(w.ctx, closingErr) + return closingErr } func (w *Writer) openWriter() (err error) { @@ -440,6 +497,7 @@ func (w *Writer) openWriter() (err error) { if err != nil { return err } + w.ctx = params.ctx w.opened = true go w.monitorCancel() diff --git a/ci/resources/stemcell-version-bump/vendor/modules.txt b/ci/resources/stemcell-version-bump/vendor/modules.txt index 65ecf389..56acc517 100644 --- a/ci/resources/stemcell-version-bump/vendor/modules.txt +++ b/ci/resources/stemcell-version-bump/vendor/modules.txt @@ -42,7 +42,7 @@ cloud.google.com/go/iam/apiv1/iampb cloud.google.com/go/monitoring/apiv3/v2 cloud.google.com/go/monitoring/apiv3/v2/monitoringpb cloud.google.com/go/monitoring/internal -# cloud.google.com/go/storage v1.63.1 +# cloud.google.com/go/storage v1.64.0 ## explicit; go 1.25.0 cloud.google.com/go/storage cloud.google.com/go/storage/experimental