diff --git a/sdks/go.mod b/sdks/go.mod index 79b771a8a1e2..84c4af541f24 100644 --- a/sdks/go.mod +++ b/sdks/go.mod @@ -32,12 +32,12 @@ require ( cloud.google.com/go/pubsub v1.51.0 cloud.google.com/go/spanner v1.94.0 cloud.google.com/go/storage v1.64.0 - github.com/aws/aws-sdk-go-v2 v1.43.5 - github.com/aws/aws-sdk-go-v2/config v1.32.36 - github.com/aws/aws-sdk-go-v2/credentials v1.19.35 + github.com/aws/aws-sdk-go-v2 v1.43.6 + github.com/aws/aws-sdk-go-v2/config v1.32.37 + github.com/aws/aws-sdk-go-v2/credentials v1.19.36 github.com/aws/aws-sdk-go-v2/feature/s3/manager v1.22.41 - github.com/aws/aws-sdk-go-v2/service/s3 v1.107.1 - github.com/aws/smithy-go v1.27.7 + github.com/aws/aws-sdk-go-v2/service/s3 v1.107.2 + github.com/aws/smithy-go v1.27.8 github.com/docker/go-connections v0.7.0 // indirect github.com/dustin/go-humanize v1.0.1 github.com/go-sql-driver/mysql v1.10.0 @@ -47,7 +47,7 @@ require ( github.com/lib/pq v1.12.3 github.com/linkedin/goavro/v2 v2.15.0 github.com/nats-io/nats-server/v2 v2.14.4 - github.com/nats-io/nats.go v1.52.0 + github.com/nats-io/nats.go v1.53.1 github.com/proullon/ramsql v0.1.4 github.com/spf13/cobra v1.10.2 github.com/testcontainers/testcontainers-go v0.44.0 @@ -55,11 +55,11 @@ require ( github.com/xitongsys/parquet-go v1.6.2 github.com/xitongsys/parquet-go-source v0.0.0-20241021075129-b732d2ac9c9b go.mongodb.org/mongo-driver v1.17.9 - golang.org/x/net v0.57.0 + golang.org/x/net v0.58.0 golang.org/x/oauth2 v0.36.0 golang.org/x/sync v0.22.0 golang.org/x/sys v0.47.0 - golang.org/x/text v0.40.0 + golang.org/x/text v0.41.0 google.golang.org/api v0.293.0 google.golang.org/genproto v0.0.0-20260523011958-0a33c5d7ca68 google.golang.org/grpc v1.83.0 @@ -91,7 +91,7 @@ require ( github.com/GoogleCloudPlatform/opentelemetry-operations-go/internal/resourcemapping v0.57.0 // indirect github.com/antithesishq/antithesis-sdk-go v0.7.2-default-no-op // indirect github.com/apache/arrow/go/v15 v15.0.2 // indirect - github.com/aws/aws-sdk-go-v2/service/signin v1.5.5 // indirect + github.com/aws/aws-sdk-go-v2/service/signin v1.5.6 // indirect github.com/containerd/errdefs v1.0.0 // indirect github.com/containerd/errdefs/pkg v0.3.0 // indirect github.com/containerd/log v0.1.0 // indirect @@ -134,7 +134,7 @@ require ( go.opentelemetry.io/otel/sdk/metric v1.44.0 // indirect go.opentelemetry.io/otel/trace v1.44.0 // indirect go.shabbyrobe.org/gocovmerge v0.0.0-20230507111327-fa4f82cfbf4d // indirect - golang.org/x/telemetry v0.0.0-20260625142307-59b4966ccb57 // indirect + golang.org/x/telemetry v0.0.0-20260708182218-49f421fb7959 // indirect golang.org/x/time v0.15.0 // indirect ) @@ -147,18 +147,18 @@ require ( github.com/Microsoft/go-winio v0.6.2 // indirect github.com/apache/arrow/go/arrow v0.0.0-20211112161151-bc219186db40 // indirect github.com/apache/thrift v0.23.0 // indirect - github.com/aws/aws-sdk-go-v2/aws/protocol/eventstream v1.7.17 // indirect - github.com/aws/aws-sdk-go-v2/feature/ec2/imds v1.18.36 // indirect - github.com/aws/aws-sdk-go-v2/internal/configsources v1.4.36 // indirect - github.com/aws/aws-sdk-go-v2/internal/endpoints/v2 v2.7.36 // indirect - github.com/aws/aws-sdk-go-v2/internal/v4a v1.4.37 // indirect - github.com/aws/aws-sdk-go-v2/service/internal/accept-encoding v1.13.16 // indirect - github.com/aws/aws-sdk-go-v2/service/internal/checksum v1.9.29 // indirect - github.com/aws/aws-sdk-go-v2/service/internal/presigned-url v1.13.36 // indirect - github.com/aws/aws-sdk-go-v2/service/internal/s3shared v1.19.37 // indirect - github.com/aws/aws-sdk-go-v2/service/sso v1.33.5 // indirect - github.com/aws/aws-sdk-go-v2/service/ssooidc v1.38.5 // indirect - github.com/aws/aws-sdk-go-v2/service/sts v1.45.5 // indirect + github.com/aws/aws-sdk-go-v2/aws/protocol/eventstream v1.7.18 // indirect + github.com/aws/aws-sdk-go-v2/feature/ec2/imds v1.18.37 // indirect + github.com/aws/aws-sdk-go-v2/internal/configsources v1.4.37 // indirect + github.com/aws/aws-sdk-go-v2/internal/endpoints/v2 v2.7.37 // indirect + github.com/aws/aws-sdk-go-v2/internal/v4a v1.4.38 // indirect + github.com/aws/aws-sdk-go-v2/service/internal/accept-encoding v1.13.17 // indirect + github.com/aws/aws-sdk-go-v2/service/internal/checksum v1.9.30 // indirect + github.com/aws/aws-sdk-go-v2/service/internal/presigned-url v1.13.37 // indirect + github.com/aws/aws-sdk-go-v2/service/internal/s3shared v1.19.38 // indirect + github.com/aws/aws-sdk-go-v2/service/sso v1.33.6 // indirect + github.com/aws/aws-sdk-go-v2/service/ssooidc v1.38.6 // indirect + github.com/aws/aws-sdk-go-v2/service/sts v1.45.6 // indirect github.com/cenkalti/backoff/v4 v4.3.0 // indirect github.com/cespare/xxhash/v2 v2.3.0 // indirect github.com/cncf/xds/go v0.0.0-20260202195803-dba9d589def2 // indirect @@ -198,9 +198,9 @@ require ( github.com/youmark/pkcs8 v0.0.0-20240726163527-a2c0da244d78 // indirect github.com/zeebo/xxh3 v1.1.0 // indirect go.opencensus.io v0.24.0 // indirect - golang.org/x/crypto v0.54.0 // indirect - golang.org/x/mod v0.37.0 // indirect - golang.org/x/tools v0.47.0 // indirect + golang.org/x/crypto v0.55.0 // indirect + golang.org/x/mod v0.38.0 // indirect + golang.org/x/tools v0.48.0 // indirect golang.org/x/xerrors v0.0.0-20240903120638-7835f813f4da // indirect google.golang.org/genproto/googleapis/api v0.0.0-20260630182238-925bb5da69e7 // indirect google.golang.org/genproto/googleapis/rpc v0.0.0-20260807164820-c8921c73eeea // indirect diff --git a/sdks/go.sum b/sdks/go.sum index 63d95d032eed..e9e7c55360f7 100644 --- a/sdks/go.sum +++ b/sdks/go.sum @@ -196,83 +196,83 @@ github.com/aws/aws-sdk-go v1.37.0/go.mod h1:hcU610XS61/+aQV88ixoOzUoG7v3b31pl2zK github.com/aws/aws-sdk-go v1.43.31/go.mod h1:y4AeaBuwd2Lk+GepC1E9v0qOiTws0MIWAX4oIKwKHZo= github.com/aws/aws-sdk-go-v2 v1.16.2/go.mod h1:ytwTPBG6fXTZLxxeeCCWj2/EMYp/xDUgX+OET6TLNNU= github.com/aws/aws-sdk-go-v2 v1.23.0/go.mod h1:i1XDttT4rnf6vxc9AuskLc6s7XBee8rlLilKlc03uAA= -github.com/aws/aws-sdk-go-v2 v1.43.5 h1:yKT5GYnFWhuDo+DqKvE5ZPwVn3RjC4MAeBtZGlh6AVM= -github.com/aws/aws-sdk-go-v2 v1.43.5/go.mod h1:wZjAJppCntyOGgVSmgVTfDyRJK5PHOasO6Wsy8U7Axk= +github.com/aws/aws-sdk-go-v2 v1.43.6 h1:RrmFcqCBxkJuf7g1axVo5krB4jM/AO8r5e5oujrgdoQ= +github.com/aws/aws-sdk-go-v2 v1.43.6/go.mod h1:tXpPM+v0D1lndmga+HqqLDIzUFJlEeR21aspVklHF00= github.com/aws/aws-sdk-go-v2/aws/protocol/eventstream v1.4.1/go.mod h1:n8Bs1ElDD2wJ9kCRTczA83gYbBmjSwZp3umc6zF4EeM= github.com/aws/aws-sdk-go-v2/aws/protocol/eventstream v1.5.1/go.mod h1:t8PYl/6LzdAqsU4/9tz28V/kU+asFePvpOMkdul0gEQ= -github.com/aws/aws-sdk-go-v2/aws/protocol/eventstream v1.7.17 h1:mn+Vxb9zgz/FE/yDTcFim3DZ1qpcrxR+qBQkBrl6bzA= -github.com/aws/aws-sdk-go-v2/aws/protocol/eventstream v1.7.17/go.mod h1:eDfmEFxu+BSVsUGLbzJhWjpOurv1mqczClS97yI8wdk= +github.com/aws/aws-sdk-go-v2/aws/protocol/eventstream v1.7.18 h1:LAfOuhAH331fmOjTQpAaOlH+Ftn7RzSDJ2VFwjdMMy4= +github.com/aws/aws-sdk-go-v2/aws/protocol/eventstream v1.7.18/go.mod h1:4e5xhuXHx1e4U9EthvbPP1r/DIMp5c2823OL8karzcM= github.com/aws/aws-sdk-go-v2/config v1.15.3/go.mod h1:9YL3v07Xc/ohTsxFXzan9ZpFpdTOFl4X65BAKYaz8jg= github.com/aws/aws-sdk-go-v2/config v1.25.3/go.mod h1:tAByZy03nH5jcq0vZmkcVoo6tRzRHEwSFx3QW4NmDw8= -github.com/aws/aws-sdk-go-v2/config v1.32.36 h1:mX6ietU7UlB4w/2IUaexJdsyUDvhTd+jYPjVePiyi6s= -github.com/aws/aws-sdk-go-v2/config v1.32.36/go.mod h1:rMpV4xk7ZK59edraSaHP0jsWrztWTT5tbCwWY495hug= +github.com/aws/aws-sdk-go-v2/config v1.32.37 h1:Ljl7LOJB6ym0liuEl0+TZ3d7f5I8MEZN1Cj9PINlj/g= +github.com/aws/aws-sdk-go-v2/config v1.32.37/go.mod h1:WJ7pe7ZPpmG8Q5kKS53zeypIV4FBGACxmte8Uc6SgUc= github.com/aws/aws-sdk-go-v2/credentials v1.11.2/go.mod h1:j8YsY9TXTm31k4eFhspiQicfXPLZ0gYXA50i4gxPE8g= github.com/aws/aws-sdk-go-v2/credentials v1.16.2/go.mod h1:sDdvGhXrSVT5yzBDR7qXz+rhbpiMpUYfF3vJ01QSdrc= -github.com/aws/aws-sdk-go-v2/credentials v1.19.35 h1:Cxua2RVdRwL0sfjHM/SnQoOnQ7xKng9m5EQBO8BnZlg= -github.com/aws/aws-sdk-go-v2/credentials v1.19.35/go.mod h1:9XQ+RSIGPkycr+oCJYnB1uTv5kMVVR+rd2vYK0Hxj2w= +github.com/aws/aws-sdk-go-v2/credentials v1.19.36 h1:84s5xMme6ENYEdKG8rsbSFFg/8+lbHBeM9QYSO0gnDk= +github.com/aws/aws-sdk-go-v2/credentials v1.19.36/go.mod h1:c46BLdagDLIswjgt+GeQOslXgeS0E6wCacs5yZbxPGk= github.com/aws/aws-sdk-go-v2/feature/ec2/imds v1.12.3/go.mod h1:uk1vhHHERfSVCUnqSqz8O48LBYDSC+k6brng09jcMOk= github.com/aws/aws-sdk-go-v2/feature/ec2/imds v1.14.4/go.mod h1:t4i+yGHMCcUNIX1x7YVYa6bH/Do7civ5I6cG/6PMfyA= -github.com/aws/aws-sdk-go-v2/feature/ec2/imds v1.18.36 h1:gucL1KH/PAYbpTpBg09CiVpBdTu4qkCl8C7xOTBixUg= -github.com/aws/aws-sdk-go-v2/feature/ec2/imds v1.18.36/go.mod h1:usTB+PHhNMhrx2dxUeHcM7OrT5pySvmjYI++IsefPN0= +github.com/aws/aws-sdk-go-v2/feature/ec2/imds v1.18.37 h1:b5tb+CZItBkydC7r3hTNdSO3pszG1R2EtnA+7TePQPk= +github.com/aws/aws-sdk-go-v2/feature/ec2/imds v1.18.37/go.mod h1:ZQ+6SU9X0oz6+7MUCSswv9Mjci4eaqZr21HI2RVy/yA= github.com/aws/aws-sdk-go-v2/feature/s3/manager v1.11.3/go.mod h1:0dHuD2HZZSiwfJSy1FO5bX1hQ1TxVV1QXXjpn3XUE44= github.com/aws/aws-sdk-go-v2/feature/s3/manager v1.14.0/go.mod h1:UcgIwJ9KHquYxs6Q5skC9qXjhYMK+JASDYcXQ4X7JZE= github.com/aws/aws-sdk-go-v2/feature/s3/manager v1.22.41 h1:RgRlA38QGJPnsG6TNEnMTQq8COa6Ye7AGwdnFlzd3Ys= github.com/aws/aws-sdk-go-v2/feature/s3/manager v1.22.41/go.mod h1:14QhagNwgiQqYUcAAJjQwgBXCcpxUau+8S/0YRLt5Uo= github.com/aws/aws-sdk-go-v2/internal/configsources v1.1.9/go.mod h1:AnVH5pvai0pAF4lXRq0bmhbes1u9R8wTE+g+183bZNM= github.com/aws/aws-sdk-go-v2/internal/configsources v1.2.3/go.mod h1:7sGSz1JCKHWWBHq98m6sMtWQikmYPpxjqOydDemiVoM= -github.com/aws/aws-sdk-go-v2/internal/configsources v1.4.36 h1:5CrzwxDqf4w3x1Vs3/NiZ0nsC34Hbm3pIDMWbsLebOE= -github.com/aws/aws-sdk-go-v2/internal/configsources v1.4.36/go.mod h1:A3gHdKZIvG/QXERzZwcxNS3RNDFcRCuhhTFBYp+V/nw= +github.com/aws/aws-sdk-go-v2/internal/configsources v1.4.37 h1:lznzIOvvbqjfe8UAaciCRJgBgJsxuTROKlhZuXQWfv8= +github.com/aws/aws-sdk-go-v2/internal/configsources v1.4.37/go.mod h1:otfkzyfQeMMLZAqX59GSXTL3o22BR/l6HFaRzzbWSqA= github.com/aws/aws-sdk-go-v2/internal/endpoints/v2 v2.4.3/go.mod h1:ssOhaLpRlh88H3UmEcsBoVKq309quMvm3Ds8e9d4eJM= github.com/aws/aws-sdk-go-v2/internal/endpoints/v2 v2.5.3/go.mod h1:ify42Rb7nKeDDPkFjKn7q1bPscVPu/+gmHH8d2c+anU= -github.com/aws/aws-sdk-go-v2/internal/endpoints/v2 v2.7.36 h1:A4N2f4YPcST0v+dWtX+xrpPPCL9VTBhoIFFUWYqbacE= -github.com/aws/aws-sdk-go-v2/internal/endpoints/v2 v2.7.36/go.mod h1:B/Qr859uxWUEfZeGotK5KAEoof4Q9YWgNtPSwV6jcyk= +github.com/aws/aws-sdk-go-v2/internal/endpoints/v2 v2.7.37 h1:zCEORWo0eU0gDjG+IyApE/2B+ZGG1m+GU7B263XV8ds= +github.com/aws/aws-sdk-go-v2/internal/endpoints/v2 v2.7.37/go.mod h1:i6c0PEl3TNOWxRbQ++KQcVenPWS/GoQeiklKhNuqzJ8= github.com/aws/aws-sdk-go-v2/internal/ini v1.3.10/go.mod h1:8DcYQcz0+ZJaSxANlHIsbbi6S+zMwjwdDqwW3r9AzaE= github.com/aws/aws-sdk-go-v2/internal/ini v1.7.1/go.mod h1:6fQQgfuGmw8Al/3M2IgIllycxV7ZW7WCdVSqfBeUiCY= github.com/aws/aws-sdk-go-v2/internal/v4a v1.2.3/go.mod h1:5yzAuE9i2RkVAttBl8yxZgQr5OCq4D5yDnG7j9x2L0U= -github.com/aws/aws-sdk-go-v2/internal/v4a v1.4.37 h1:oyd3ke4V9AhKcRR7rRgxk1VyI+DjK2CBQtbxh3OkdaA= -github.com/aws/aws-sdk-go-v2/internal/v4a v1.4.37/go.mod h1:aA9D7SqfG9IC1b7FLD7Iyc8Q4JN0a8gHhNjN4zPlIaI= +github.com/aws/aws-sdk-go-v2/internal/v4a v1.4.38 h1:A3UAuCmx7LyUcrixBTzKJYYIUZ2yTvn6ZhT8PB+7APk= +github.com/aws/aws-sdk-go-v2/internal/v4a v1.4.38/go.mod h1:1PDUYG9Z+JrbbsobsAZHjWOm9QBT/djiK3QbykTL5Z4= github.com/aws/aws-sdk-go-v2/service/internal/accept-encoding v1.9.1/go.mod h1:GeUru+8VzrTXV/83XyMJ80KpH8xO89VPoUileyNQ+tc= github.com/aws/aws-sdk-go-v2/service/internal/accept-encoding v1.10.1/go.mod h1:l9ymW25HOqymeU2m1gbUQ3rUIsTwKs8gYHXkqDQUhiI= -github.com/aws/aws-sdk-go-v2/service/internal/accept-encoding v1.13.16 h1:iE4NGbvqUZnHDqddQAauZzCILYtFjOHwRM5MOOKLB5A= -github.com/aws/aws-sdk-go-v2/service/internal/accept-encoding v1.13.16/go.mod h1:VsjEgrP+ibcou8TlWA4tYaB+0OojuhirsmCe+U60hTA= +github.com/aws/aws-sdk-go-v2/service/internal/accept-encoding v1.13.17 h1:OvYZOB3qA6zvfdRFiRFRzVSiElMYrz3GdntkXZxlp1o= +github.com/aws/aws-sdk-go-v2/service/internal/accept-encoding v1.13.17/go.mod h1:JgR/2Ew50ACfIWau1oeMRX59tMtC0kM+PYQGEaT04cY= github.com/aws/aws-sdk-go-v2/service/internal/checksum v1.1.3/go.mod h1:Seb8KNmD6kVTjwRjVEgOT5hPin6sq+v4C2ycJQDwuH8= github.com/aws/aws-sdk-go-v2/service/internal/checksum v1.2.3/go.mod h1:R+/S1O4TYpcktbVwddeOYg+uwUfLhADP2S/x4QwsCTM= -github.com/aws/aws-sdk-go-v2/service/internal/checksum v1.9.29 h1:E65Hj648dOV6FuUfI0mYXXhQRHbsi7n+B9h6fZPJO/E= -github.com/aws/aws-sdk-go-v2/service/internal/checksum v1.9.29/go.mod h1:xLrF9yNTCs92VZSpdEd68EJbgcdw3SMR74RO6QDzWHE= +github.com/aws/aws-sdk-go-v2/service/internal/checksum v1.9.30 h1:5437eMoOwqqQpZn2XJy74mlDCuPYL81texMT3mXqgtU= +github.com/aws/aws-sdk-go-v2/service/internal/checksum v1.9.30/go.mod h1:xfu2m3dOpvW8lj98wQYa8V9ku/Rta59hsbireGzhh3A= github.com/aws/aws-sdk-go-v2/service/internal/presigned-url v1.9.3/go.mod h1:wlY6SVjuwvh3TVRpTqdy4I1JpBFLX4UGeKZdWntaocw= github.com/aws/aws-sdk-go-v2/service/internal/presigned-url v1.10.3/go.mod h1:Owv1I59vaghv1Ax8zz8ELY8DN7/Y0rGS+WWAmjgi950= -github.com/aws/aws-sdk-go-v2/service/internal/presigned-url v1.13.36 h1:fx2ujmozWn+C/GtfXfz5k6Ckzza40ElOpIW7d92fLWQ= -github.com/aws/aws-sdk-go-v2/service/internal/presigned-url v1.13.36/go.mod h1:QT2ufGVJ+xTRxtXPHTQ1kHkAdWIKPCmD+BqYAXWv8/4= +github.com/aws/aws-sdk-go-v2/service/internal/presigned-url v1.13.37 h1:a3D4AjrOrTrP8+d9ILBthqrElf0z1JNol09Xvnwcys8= +github.com/aws/aws-sdk-go-v2/service/internal/presigned-url v1.13.37/go.mod h1:ky0gTu+ukvUTuUKFIpp6Wid4oninrkCyvbFkVs0kpHM= github.com/aws/aws-sdk-go-v2/service/internal/s3shared v1.13.3/go.mod h1:Bm/v2IaN6rZ+Op7zX+bOUMdL4fsrYZiD0dsjLhNKwZc= github.com/aws/aws-sdk-go-v2/service/internal/s3shared v1.16.3/go.mod h1:KZgs2ny8HsxRIRbDwgvJcHHBZPOzQr/+NtGwnP+w2ec= -github.com/aws/aws-sdk-go-v2/service/internal/s3shared v1.19.37 h1:KGHa9iZCrgtkOsFfXb0S4ywsjostA/hau7WE9aSb43E= -github.com/aws/aws-sdk-go-v2/service/internal/s3shared v1.19.37/go.mod h1:FV79f0DSnZIEGsQjWenENGtUycrasyAaJZO+zRanLHA= +github.com/aws/aws-sdk-go-v2/service/internal/s3shared v1.19.38 h1:gX8B8y3Ho30B1LPxefDKMi/HZqWEb47U9ogs3DtSG0M= +github.com/aws/aws-sdk-go-v2/service/internal/s3shared v1.19.38/go.mod h1:l5WblZlcmGPe4/O7JY2HO25Z+xqTBvyfTyFbRMf8gYw= github.com/aws/aws-sdk-go-v2/service/kms v1.16.3/go.mod h1:QuiHPBqlOFCi4LqdSskYYAWpQlx3PKmohy+rE2F+o5g= github.com/aws/aws-sdk-go-v2/service/s3 v1.26.3/go.mod h1:g1qvDuRsJY+XghsV6zg00Z4KJ7DtFFCx8fJD2a491Ak= github.com/aws/aws-sdk-go-v2/service/s3 v1.43.0/go.mod h1:NXRKkiRF+erX2hnybnVU660cYT5/KChRD4iUgJ97cI8= -github.com/aws/aws-sdk-go-v2/service/s3 v1.107.1 h1:VUTtUJMuRNMkb/7NIKmd8NQaeQLPGCMoTJxkYKre4qM= -github.com/aws/aws-sdk-go-v2/service/s3 v1.107.1/go.mod h1:WvUaO0lP5GNMs1R6cs6qvB3mqo16GLta8yfOuf55Rpc= +github.com/aws/aws-sdk-go-v2/service/s3 v1.107.2 h1:GNU0/xtPEXMKilJZ/a8BedeuQnvu+Usi6qVm9EFfncc= +github.com/aws/aws-sdk-go-v2/service/s3 v1.107.2/go.mod h1:4jYWUecEsQtE73jPl7p3jrbYXH5ffcR4gegyCygagfg= github.com/aws/aws-sdk-go-v2/service/secretsmanager v1.15.4/go.mod h1:PJc8s+lxyU8rrre0/4a0pn2wgwiDvOEzoOjcJUBr67o= -github.com/aws/aws-sdk-go-v2/service/signin v1.5.5 h1:0VTFBfOgPJrUSpGMgzoi8qLcXF5dbmiBuxpo14eBWUw= -github.com/aws/aws-sdk-go-v2/service/signin v1.5.5/go.mod h1:sNZYlBxoohYMBYl47BO/bFtAM6I8HSsPa1qwwPPRGoQ= +github.com/aws/aws-sdk-go-v2/service/signin v1.5.6 h1:i68sFvXidKlkiSvI7d7Ilc1/UvW4CtBOaivH7jhG4fs= +github.com/aws/aws-sdk-go-v2/service/signin v1.5.6/go.mod h1:/h7Obr9WTtzbjTHGASRQwLN7Bupw+TC3x8x7fyx39hE= github.com/aws/aws-sdk-go-v2/service/sns v1.17.4/go.mod h1:kElt+uCcXxcqFyc+bQqZPFD9DME/eC6oHBXvFzQ9Bcw= github.com/aws/aws-sdk-go-v2/service/sqs v1.18.3/go.mod h1:skmQo0UPvsjsuYYSYMVmrPc1HWCbHUJyrCEp+ZaLzqM= github.com/aws/aws-sdk-go-v2/service/ssm v1.24.1/go.mod h1:NR/xoKjdbRJ+qx0pMR4mI+N/H1I1ynHwXnO6FowXJc0= github.com/aws/aws-sdk-go-v2/service/sso v1.11.3/go.mod h1:7UQ/e69kU7LDPtY40OyoHYgRmgfGM4mgsLYtcObdveU= github.com/aws/aws-sdk-go-v2/service/sso v1.17.2/go.mod h1:/pE21vno3q1h4bbhUOEi+6Zu/aT26UK2WKkDXd+TssQ= -github.com/aws/aws-sdk-go-v2/service/sso v1.33.5 h1:jDQARFp1mJ2PEnllQf01nfFXGfWMJ59e0/HCHUTTZCk= -github.com/aws/aws-sdk-go-v2/service/sso v1.33.5/go.mod h1:OcT2AhgTuxGAwZk5hgxaNLGpS33W8s8dUQadGVDVY9I= +github.com/aws/aws-sdk-go-v2/service/sso v1.33.6 h1:tpfGChmjUmv3W9WlRvy+stwKDTbFFdq8Zk9DbFPrfMU= +github.com/aws/aws-sdk-go-v2/service/sso v1.33.6/go.mod h1:CSjiDzmG/lsKkTOYjbkM+duLmRlW+LOxD64Na44ijnI= github.com/aws/aws-sdk-go-v2/service/ssooidc v1.20.0/go.mod h1:dWqm5G767qwKPuayKfzm4rjzFmVjiBFbOJrpSPnAMDs= -github.com/aws/aws-sdk-go-v2/service/ssooidc v1.38.5 h1:8xo1q9ttkYqMJ6vOXX67FPSpVEI7BWKVTKh77g82w+8= -github.com/aws/aws-sdk-go-v2/service/ssooidc v1.38.5/go.mod h1:hbBeEUrZg6VddXYZpbKPyF0tl4XEnM+Dbx92RW3vmZI= +github.com/aws/aws-sdk-go-v2/service/ssooidc v1.38.6 h1:49BBtY68A+KJCQ3a2F3eUe6ROsKucxUdfHKoqorc0wI= +github.com/aws/aws-sdk-go-v2/service/ssooidc v1.38.6/go.mod h1:ptG2hbs7QltE1GcQY0MpS4bfrc51KCnBXUr7OT1EEfE= github.com/aws/aws-sdk-go-v2/service/sts v1.16.3/go.mod h1:bfBj0iVmsUyUg4weDB4NxktD9rDGeKSVWnjTnwbx9b8= github.com/aws/aws-sdk-go-v2/service/sts v1.25.3/go.mod h1:4EqRHDCKP78hq3zOnmFXu5k0j4bXbRFfCh/zQ6KnEfQ= -github.com/aws/aws-sdk-go-v2/service/sts v1.45.5 h1:eQ5BtXDrPg2wK0AjtVPzeBhUpYPeqHE/ptiH7xJRGek= -github.com/aws/aws-sdk-go-v2/service/sts v1.45.5/go.mod h1:f9ImhnOISY7BuTZLM8qHepCYnglHBVLk5wVzatmP++w= +github.com/aws/aws-sdk-go-v2/service/sts v1.45.6 h1:JvExZWabChDM0qJAirQYGfOYo0ndT3edXj+fqSPNjkE= +github.com/aws/aws-sdk-go-v2/service/sts v1.45.6/go.mod h1:XZcaQkV2cItp6yEkrwljyaPOf22RuX7T43jxap/FOmM= github.com/aws/smithy-go v1.11.2/go.mod h1:3xHYmszWVx2c0kIwQeEVf9uSm4fYZt67FBJnwub1bgM= github.com/aws/smithy-go v1.17.0/go.mod h1:NukqUGpCZIILqqiV0NIjeFh24kd/FAa4beRb6nbIUPE= -github.com/aws/smithy-go v1.27.7 h1:Zgj5z4LfcDYoQIVk+n/yGdTkP/2y6ZT5vYxe0fp7bqE= -github.com/aws/smithy-go v1.27.7/go.mod h1:YE2RhdIuDbA5E5bTdciG9KrW3+TiEONeUWCqxX9i1Fc= +github.com/aws/smithy-go v1.27.8 h1:FR0dxZfIlV7Z8eh2iHfIofdunw382XsDV3Mxt9nUvRY= +github.com/aws/smithy-go v1.27.8/go.mod h1:YE2RhdIuDbA5E5bTdciG9KrW3+TiEONeUWCqxX9i1Fc= github.com/benbjohnson/clock v1.1.0/go.mod h1:J11/hYXuz8f4ySSvYwY0FKfm+ezbsZBKZxNJlLklBHA= github.com/bobg/gcsobj v0.1.2/go.mod h1:vS49EQ1A1Ib8FgrL58C8xXYZyOCR2TgzAdopy6/ipa8= github.com/boombuler/barcode v1.0.0/go.mod h1:paBWMcWSl3LHKBqUq+rly7CNSldXjb2rDl3JlRe0mD8= @@ -710,8 +710,8 @@ github.com/nats-io/jwt/v2 v2.8.2 h1:XXRgB60MSTnqsRwejQurVDs/hcv2dkt+86GjI+I/bMc= github.com/nats-io/jwt/v2 v2.8.2/go.mod h1:Ag/56sq9OblL4JgdYufDd16Egb17Kr/8WwwuO/forVc= github.com/nats-io/nats-server/v2 v2.14.4 h1:efgjZ8cdExAKRuqSg8UPJFprb+l7NlBtSDPhDlw3rO4= github.com/nats-io/nats-server/v2 v2.14.4/go.mod h1:BltdpOYestjbtQSnVO2zGHdg5SGBZjt+GYTgB9LZq/I= -github.com/nats-io/nats.go v1.52.0 h1:n3avV4VBsCgsdwh71TppsTwtv+QdPs7ntSKM8qJLGsc= -github.com/nats-io/nats.go v1.52.0/go.mod h1:26HypzazeOkyO3/mqd1zZd53STJN0EjCYF9Uy2ZOBno= +github.com/nats-io/nats.go v1.53.1 h1:Otsq3uLc/kLdjmkNHkXH0jBqwUquwdKFoe3fq6/3/Xo= +github.com/nats-io/nats.go v1.53.1/go.mod h1:26HypzazeOkyO3/mqd1zZd53STJN0EjCYF9Uy2ZOBno= github.com/nats-io/nkeys v0.4.16 h1:rd5oAuLOb8mnAycB0xleuEBNS1pVVnN0fv/FF34Eypg= github.com/nats-io/nkeys v0.4.16/go.mod h1:llLgWoI0o4z/Q57q2R1kHfmocyhGV6VG/U18Glg1Afs= github.com/nats-io/nuid v1.0.1 h1:5iA8DT8V7q8WK2EScv2padNa/rTESc1KdnPw4TC2paw= @@ -924,8 +924,8 @@ golang.org/x/crypto v0.0.0-20220511200225-c6db032c6c88/go.mod h1:IxCIyHEi3zRg3s0 golang.org/x/crypto v0.0.0-20220722155217-630584e8d5aa/go.mod h1:IxCIyHEi3zRg3s0A5j5BB6A9Jmi73HwBIUl50j+osU4= golang.org/x/crypto v0.7.0/go.mod h1:pYwdfH91IfpZVANVyUOhSIPZaFoJGxTFbZhFTx+dXZU= golang.org/x/crypto v0.9.0/go.mod h1:yrmDGqONDYtNj3tH8X9dzUun2m2lzPa9ngI6/RUPGR0= -golang.org/x/crypto v0.54.0 h1:YLIA59K4fiNzHzjnZt2tUJQjQtUWfWbeHBqKtk3eScw= -golang.org/x/crypto v0.54.0/go.mod h1:KWL8ny2AZdGR2cWmzeHrp2azQPGogOv+HeQaVEXC2dk= +golang.org/x/crypto v0.55.0 h1:+KWHjbgOaAQ66dh/YlkZKHlz9ZUlq61AFirAR9ntP8M= +golang.org/x/crypto v0.55.0/go.mod h1:uq0V9dE/fzQuJtbnL+2EhWOE63vo164FY8xqEnV9xis= golang.org/x/exp v0.0.0-20180321215751-8460e604b9de/go.mod h1:CJ0aWSM057203Lf6IL+f9T1iT9GByDxfZKAQTCR3kQA= golang.org/x/exp v0.0.0-20180807140117-3d87b88a115f/go.mod h1:CJ0aWSM057203Lf6IL+f9T1iT9GByDxfZKAQTCR3kQA= golang.org/x/exp v0.0.0-20190121172915-509febef88a4/go.mod h1:CJ0aWSM057203Lf6IL+f9T1iT9GByDxfZKAQTCR3kQA= @@ -977,8 +977,8 @@ golang.org/x/mod v0.4.2/go.mod h1:s0Qsj1ACt9ePp/hMypM3fl4fZqREWJwdYDEqhRiZZUA= golang.org/x/mod v0.5.0/go.mod h1:5OXOZSfqPIIbmVBIIKWRFfZjPR0E5r58TLhUjH0a2Ro= golang.org/x/mod v0.6.0-dev.0.20220419223038-86c51ed26bb4/go.mod h1:jJ57K6gSWd91VN4djpZkiMVwK6gcyfeH4XE8wZrZaV4= golang.org/x/mod v0.8.0/go.mod h1:iBbtSCu2XBx23ZKBPSOrRkjjQPZFPuis4dIYUhu/chs= -golang.org/x/mod v0.37.0 h1:vF1DjpVEshcIqoEaauuHebaLk1O1forxjxBaVn884JQ= -golang.org/x/mod v0.37.0/go.mod h1:m8S8VeM9r4dzDwjrKO0a1sZP3YjeMamRRlD+fmR2Q/0= +golang.org/x/mod v0.38.0 h1:MECBjubtXD7yj4HrhIUcywNaGeNVUdfVnxmPajOk4yk= +golang.org/x/mod v0.38.0/go.mod h1:V6Xz0pq8TQ3dGqVQ1FVHuelZpAL0uNhSkk9ogYP3c40= 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-20190108225652-1e06a53dbb7e/go.mod h1:mL1N/T3taQHkDXs73rZJwtUhF3w3ftmwwsq0BUmARs4= @@ -1032,8 +1032,8 @@ golang.org/x/net v0.6.0/go.mod h1:2Tu9+aMcznHK/AK1HMvgo6xiTLG5rD5rZLDS+rp2Bjs= golang.org/x/net v0.7.0/go.mod h1:2Tu9+aMcznHK/AK1HMvgo6xiTLG5rD5rZLDS+rp2Bjs= golang.org/x/net v0.8.0/go.mod h1:QVkue5JL9kW//ek3r6jTKnTFis1tRmNAW2P1shuFdJc= golang.org/x/net v0.10.0/go.mod h1:0qNGK6F8kojg2nk9dLZ2mShWaEBan6FAoqfSigmmuDg= -golang.org/x/net v0.57.0 h1:K5+3DljvIuDG9/Jv9rvyMywYNFCQ9RSUY6OOTTkT+tE= -golang.org/x/net v0.57.0/go.mod h1:KpXc8iv+r3XplLAG/f7Jsf9RPszJzdR0f58q9vGOuEU= +golang.org/x/net v0.58.0 h1:ynWG7rqYi4ccpTEuPZ2QGWHktVEM9DMCj9yzDE0Q7To= +golang.org/x/net v0.58.0/go.mod h1:YwCddHnFlT7eLQqVprV19OnhLGtc5xOKgE0RyqgfWAU= golang.org/x/oauth2 v0.0.0-20180821212333-d2e6202438be/go.mod h1:N/0e6XlmueqKjAGxoOufVs8QHGRruUQn6yWY3a++T0U= golang.org/x/oauth2 v0.0.0-20190226205417-e64efc72b421/go.mod h1:gOpvHmFTYa4IltrdGE7lF6nIHvwfUNPOp7c8zoXwtLw= golang.org/x/oauth2 v0.0.0-20190604053449-0f29369cfe45/go.mod h1:gOpvHmFTYa4IltrdGE7lF6nIHvwfUNPOp7c8zoXwtLw= @@ -1157,8 +1157,8 @@ golang.org/x/sys v0.8.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= golang.org/x/sys v0.21.0/go.mod h1:/VUhepiaJMQUp4+oa/7Zr1D23ma6VTLIYjOOTFZPUcA= golang.org/x/sys v0.47.0 h1:o7XGOvZQCADBQQ4Y7VNq2dRWQR7JmOUW8Kxx4ZsNgWs= golang.org/x/sys v0.47.0/go.mod h1:4GL1E5IUh+htKOUEOaiffhrAeqysfVGipDYzABqnCmw= -golang.org/x/telemetry v0.0.0-20260625142307-59b4966ccb57 h1:nwGZBCt+FnXUrGsj5vjzAsEmkcaFvd82BbOjECiFYZc= -golang.org/x/telemetry v0.0.0-20260625142307-59b4966ccb57/go.mod h1:3AWMyWHS+caVoiEXpiq6+tzKA40J4vQT3MYr80ZtQpc= +golang.org/x/telemetry v0.0.0-20260708182218-49f421fb7959 h1:RJhm5l6Fo4rmEIcndxDllNhhf/fAx8qIm4t6A7vpm2A= +golang.org/x/telemetry v0.0.0-20260708182218-49f421fb7959/go.mod h1:LV7u5Oco+Z/g6XI7PqN+EUUUGGkEcmB1uj2ceI0fOVg= golang.org/x/term v0.0.0-20201117132131-f5c789dd3221/go.mod h1:Nr5EML6q2oocZ2LXRh80K7BxOlk5/8JxuGnuhpl+muw= golang.org/x/term v0.0.0-20201126162022-7de9c90e9dd1/go.mod h1:bj7SfCRtBDWHUb9snDiAeCFNEtKQo2Wmx5Cou7ajbmo= golang.org/x/term v0.0.0-20210927222741-03fcf44c2211/go.mod h1:jbD1KX2456YbFQfuXm/mYQcufACuNUgVhRMnK/tPxf8= @@ -1180,8 +1180,8 @@ golang.org/x/text v0.3.8/go.mod h1:E6s5w1FMmriuDzIBO73fBruAKo1PCIq6d2Q6DHfQ8WQ= golang.org/x/text v0.7.0/go.mod h1:mrYo+phRRbMaCq/xk9113O4dZlRixOauAjOtrjsXDZ8= golang.org/x/text v0.8.0/go.mod h1:e1OnstbJyHTd6l/uOt8jFFHp6TRDWZR/bV3emEE/zU8= golang.org/x/text v0.9.0/go.mod h1:e1OnstbJyHTd6l/uOt8jFFHp6TRDWZR/bV3emEE/zU8= -golang.org/x/text v0.40.0 h1:Ub2Z6/xjgF1WrYQz2nuITOEegKFtiIy+rieRJ5lHZKs= -golang.org/x/text v0.40.0/go.mod h1:hpnzDAfGV753zIKo+wk3u1bVKCGPbrnF7+7LBF/UHVY= +golang.org/x/text v0.41.0 h1:vz/seA0lnX87Othu2f/0L24RcgrXD9/YFTSuGjj3rH8= +golang.org/x/text v0.41.0/go.mod h1:jvf1O8ajNzZqhSrQBPbutR/EB83Cc0CFrezNQIwbb5M= golang.org/x/time v0.0.0-20181108054448-85acf8d2951c/go.mod h1:tRJNPiyCQ0inRvYxbN9jk5I+vvW/OXSQhTDSoE431IQ= golang.org/x/time v0.0.0-20190308202827-9d24e82272b4/go.mod h1:tRJNPiyCQ0inRvYxbN9jk5I+vvW/OXSQhTDSoE431IQ= golang.org/x/time v0.0.0-20191024005414-555d28b269f0/go.mod h1:tRJNPiyCQ0inRvYxbN9jk5I+vvW/OXSQhTDSoE431IQ= @@ -1254,8 +1254,8 @@ golang.org/x/tools v0.1.4/go.mod h1:o0xws9oXOQQZyjljx8fwUC0k7L1pTE6eaCbjGeHmOkk= golang.org/x/tools v0.1.5/go.mod h1:o0xws9oXOQQZyjljx8fwUC0k7L1pTE6eaCbjGeHmOkk= golang.org/x/tools v0.1.12/go.mod h1:hNGJHUnrk76NpqgfD5Aqm5Crs+Hm0VOH/i9J2+nxYbc= golang.org/x/tools v0.6.0/go.mod h1:Xwgl3UAJ/d3gWutnCtw505GrjyAbvKui8lOU390QaIU= -golang.org/x/tools v0.47.0 h1:7Kn5x/d1svx/PzryTsqeoZN4TZwqeH5pGWjefhLi/1Q= -golang.org/x/tools v0.47.0/go.mod h1:dFHnyTvFWY212G+h7ZY4Vsp/K3U4/7W9TyVaAul8uCA= +golang.org/x/tools v0.48.0 h1:3+hClM1aLL5mjMKm5ovokw9epgRXPuu2tILgismM6RE= +golang.org/x/tools v0.48.0/go.mod h1:08xX0orndb/F7jJxGDicx061tyd5pcMto75YMAXr6lk= golang.org/x/xerrors v0.0.0-20190410155217-1f06c39b4373/go.mod h1:I/5z698sn9Ka8TeJc9MKroUUfqBBauWjQqLJ2OPfmY0= golang.org/x/xerrors v0.0.0-20190513163551-3ee3066db522/go.mod h1:I/5z698sn9Ka8TeJc9MKroUUfqBBauWjQqLJ2OPfmY0= golang.org/x/xerrors v0.0.0-20190717185122-a985d3407aa7/go.mod h1:I/5z698sn9Ka8TeJc9MKroUUfqBBauWjQqLJ2OPfmY0= diff --git a/sdks/go/pkg/beam/core/runtime/xlangx/expansionx/download.go b/sdks/go/pkg/beam/core/runtime/xlangx/expansionx/download.go index 0b5eba625023..f9a3458500c1 100644 --- a/sdks/go/pkg/beam/core/runtime/xlangx/expansionx/download.go +++ b/sdks/go/pkg/beam/core/runtime/xlangx/expansionx/download.go @@ -139,6 +139,18 @@ func getLocalJar(url string) (string, error) { return jarPath, nil } +func validatePath(dest, filename string) (string, error) { + destPath := filepath.Join(dest, filename) + cleanDest := filepath.Clean(dest) + cleanPath := filepath.Clean(destPath) + + rel, err := filepath.Rel(cleanDest, cleanPath) + if err != nil || strings.HasPrefix(rel, ".."+string(filepath.Separator)) || rel == ".." { + return "", fmt.Errorf("file path %q is outside destination directory %q", filename, dest) + } + return cleanPath, nil +} + func extractJar(source, dest string) error { reader, err := zip.OpenReader(source) if err != nil { @@ -150,7 +162,10 @@ func extractJar(source, dest string) error { } for _, file := range reader.File { - fileName := filepath.Join(dest, file.Name) + fileName, err := validatePath(dest, file.Name) + if err != nil { + return fmt.Errorf("error validating file path (%s, %s): %w", dest, file.Name, err) + } if file.FileInfo().IsDir() { os.MkdirAll(fileName, 0700) continue diff --git a/sdks/go/pkg/beam/core/runtime/xlangx/expansionx/download_test.go b/sdks/go/pkg/beam/core/runtime/xlangx/expansionx/download_test.go index 65e72342a9b9..fc66d53cb6c5 100644 --- a/sdks/go/pkg/beam/core/runtime/xlangx/expansionx/download_test.go +++ b/sdks/go/pkg/beam/core/runtime/xlangx/expansionx/download_test.go @@ -249,3 +249,28 @@ func TestGetPythonVersion(t *testing.T) { } } } + +func TestValidatePath(t *testing.T) { + dest := filepath.Clean("/tmp/cache") + tests := []struct { + name string + filename string + wantErr bool + }{ + {"valid simple file", "Foo.class", false}, + {"valid nested file", "org/apache/beam/Foo.class", false}, + {"traversal attack", "../../etc/passwd", true}, + {"partial directory prefix attack", "../cache_evil/evil.sh", true}, + {"parent directory traversal", "..", true}, + {"nested traversal attack", "foo/bar/../../../etc/passwd", true}, + } + + for _, tc := range tests { + t.Run(tc.name, func(t *testing.T) { + _, err := validatePath(dest, tc.filename) + if (err != nil) != tc.wantErr { + t.Errorf("validatePath(%q, %q) error = %v, wantErr %v", dest, tc.filename, err, tc.wantErr) + } + }) + } +} diff --git a/sdks/python/apache_beam/transforms/util.py b/sdks/python/apache_beam/transforms/util.py index 2ea9df9399cb..60295f68a92b 100644 --- a/sdks/python/apache_beam/transforms/util.py +++ b/sdks/python/apache_beam/transforms/util.py @@ -82,6 +82,9 @@ from apache_beam.utils import shared from apache_beam.utils import windowed_value from apache_beam.utils.annotations import deprecated +from apache_beam.utils.secret import Secret +from apache_beam.utils.secret import GcpSecret +from apache_beam.utils.secret import GcpHsmGeneratedSecret from apache_beam.utils.sharded_key import ShardedKey from apache_beam.utils.timestamp import Timestamp @@ -94,6 +97,7 @@ 'BatchElements', 'CoGroupByKey', 'Distinct', + 'GcpHsmGeneratedSecret', 'GcpSecret', 'GroupByEncryptedKey', 'Keys', @@ -327,247 +331,6 @@ def RemoveDuplicates(pcoll): return pcoll | 'RemoveDuplicates' >> Distinct() -class Secret(): - """A secret management class used for handling sensitive data. - - This class provides a generic interface for secret management. Implementations - of this class should handle fetching secrets from a secret management system. - """ - def get_secret_bytes(self) -> bytes: - """Returns the secret as a byte string.""" - raise NotImplementedError() - - @staticmethod - def generate_secret_bytes() -> bytes: - """Generates a new secret key.""" - return Fernet.generate_key() - - @staticmethod - def parse_secret_option(secret) -> 'Secret': - """Parses a secret string and returns the appropriate secret type. - - The secret string should be formatted like: - 'type:;:' - - For example, 'type:GcpSecret;version_name:my_secret/versions/latest' - would return a GcpSecret initialized with 'my_secret/versions/latest'. - """ - param_map = {} - for param in secret.split(';'): - parts = param.split(':') - param_map[parts[0]] = parts[1] - - if 'type' not in param_map: - raise ValueError('Secret string must contain a valid type parameter') - - secret_type = param_map['type'].lower() - del param_map['type'] - secret_class = Secret - secret_params = None - if secret_type == 'gcpsecret': - secret_class = GcpSecret # type: ignore[assignment] - secret_params = ['version_name'] - elif secret_type == 'gcphsmgeneratedsecret': - secret_class = GcpHsmGeneratedSecret # type: ignore[assignment] - secret_params = [ - 'project_id', 'location_id', 'key_ring_id', 'key_id', 'job_name' - ] - else: - raise ValueError( - f'Invalid secret type {secret_type}, currently only ' - 'GcpSecret and GcpHsmGeneratedSecret are supported') - - for param_name in param_map.keys(): - if param_name not in secret_params: - raise ValueError( - f'Invalid secret parameter {param_name}, ' - f'{secret_type} only supports the following ' - f'parameters: {secret_params}') - return secret_class(**param_map) - - -class GcpSecret(Secret): - """A secret manager implementation that retrieves secrets from Google Cloud - Secret Manager. - """ - def __init__(self, version_name: str): - """Initializes a GcpSecret object. - - Args: - version_name: The full version name of the secret in Google Cloud Secret - Manager. For example: - projects//secrets//versions/1. - For more info, see - https://cloud.google.com/python/docs/reference/secretmanager/latest/google.cloud.secretmanager_v1beta1.services.secret_manager_service.SecretManagerServiceClient#google_cloud_secretmanager_v1beta1_services_secret_manager_service_SecretManagerServiceClient_access_secret_version - """ - self._version_name = version_name - - def get_secret_bytes(self) -> bytes: - try: - from google.cloud import secretmanager - client = secretmanager.SecretManagerServiceClient() - response = client.access_secret_version( - request={"name": self._version_name}) - secret = response.payload.data - return secret - except Exception as e: - raise RuntimeError( - 'Failed to retrieve secret bytes for secret ' - f'{self._version_name} with exception {e}') - - def __eq__(self, secret): - return self._version_name == getattr(secret, '_version_name', None) - - -class GcpHsmGeneratedSecret(Secret): - """A secret manager implementation that generates a secret using a GCP HSM key - and stores it in Google Cloud Secret Manager. If the secret already exists, - it will be retrieved. - """ - def __init__( - self, - project_id: str, - location_id: str, - key_ring_id: str, - key_id: str, - job_name: str): - """Initializes a GcpHsmGeneratedSecret object. - - Args: - project_id: The GCP project ID. - location_id: The GCP location ID for the HSM key. - key_ring_id: The ID of the KMS key ring. - key_id: The ID of the KMS key. - job_name: The name of the job, used to generate a unique secret name. - """ - self._project_id = project_id - self._location_id = location_id - self._key_ring_id = key_ring_id - self._key_id = key_id - self._secret_version_name = f'HsmGeneratedSecret_{job_name}' - - def get_secret_bytes(self) -> bytes: - """Retrieves the secret bytes. - - If the secret version already exists in Secret Manager, it is retrieved. - Otherwise, a new secret and version are created. The new secret is - generated using the HSM key. - - Returns: - The secret as a byte string. - """ - try: - from google.api_core import exceptions as api_exceptions - from google.cloud import secretmanager - client = secretmanager.SecretManagerServiceClient() - - project_path = f"projects/{self._project_id}" - secret_path = f"{project_path}/secrets/{self._secret_version_name}" - # Since we may generate multiple versions when doing this on workers, - # just always take the first version added to maintain consistency. - secret_version_path = f"{secret_path}/versions/1" - - try: - response = client.access_secret_version( - request={"name": secret_version_path}) - return response.payload.data - except api_exceptions.NotFound: - # Don't bother logging yet, we'll only log if we actually add the - # secret version below - pass - - try: - client.create_secret( - request={ - "parent": project_path, - "secret_id": self._secret_version_name, - "secret": { - "replication": { - "automatic": {} - } - }, - }) - except api_exceptions.AlreadyExists: - # Don't bother logging yet, we'll only log if we actually add the - # secret version below - pass - - new_key = self.generate_dek() - try: - # Try one more time in case it was created while we were generating the - # DEK. - response = client.access_secret_version( - request={"name": secret_version_path}) - return response.payload.data - except api_exceptions.NotFound: - _LOGGER.info( - "Secret version %s not found. " - "Creating new secret and version.", - secret_version_path) - client.add_secret_version( - request={ - "parent": secret_path, "payload": { - "data": new_key - } - }) - response = client.access_secret_version( - request={"name": secret_version_path}) - return response.payload.data - - except Exception as e: - raise RuntimeError( - f'Failed to retrieve or create secret bytes for secret ' - f'{self._secret_version_name} with exception {e}') - - def generate_dek(self, dek_size: int = 32) -> bytes: - """Generates a new Data Encryption Key (DEK) using an HSM-backed key. - - This function follows a key derivation process that incorporates entropy - from the HSM-backed key into the nonce used for key derivation. - - Args: - dek_size: The size of the DEK to generate. - - Returns: - A new DEK of the specified size, url-safe base64-encoded. - """ - try: - import base64 - import os - - from cryptography.hazmat.primitives import hashes - from cryptography.hazmat.primitives.kdf.hkdf import HKDF - from google.cloud import kms - - # 1. Generate a random nonce (nonce_one) - nonce_one = os.urandom(dek_size) - - # 2. Use the HSM-backed key to encrypt nonce_one to create nonce_two - kms_client = kms.KeyManagementServiceClient() - key_path = kms_client.crypto_key_path( - self._project_id, self._location_id, self._key_ring_id, self._key_id) - response = kms_client.encrypt( - request={ - 'name': key_path, 'plaintext': nonce_one - }) - nonce_two = response.ciphertext - - # 3. Generate a Derivation Key (DK) - dk = os.urandom(dek_size) - - # 4. Use a KDF to derive the DEK using DK and nonce_two - hkdf = HKDF( - algorithm=hashes.SHA256(), - length=dek_size, - salt=nonce_two, - info=None, - ) - dek = hkdf.derive(dk) - return base64.urlsafe_b64encode(dek) - except Exception as e: - raise RuntimeError(f'Failed to generate DEK with exception {e}') - - class _EncryptMessage(DoFn): """A DoFn that encrypts the key and value of each element.""" def __init__( diff --git a/sdks/python/apache_beam/transforms/util_test.py b/sdks/python/apache_beam/transforms/util_test.py index 63ce42726c1f..446fe68e594c 100644 --- a/sdks/python/apache_beam/transforms/util_test.py +++ b/sdks/python/apache_beam/transforms/util_test.py @@ -72,9 +72,9 @@ from apache_beam.transforms.core import FlatMapTuple from apache_beam.transforms.trigger import AfterCount from apache_beam.transforms.trigger import Repeatedly -from apache_beam.transforms.util import GcpHsmGeneratedSecret -from apache_beam.transforms.util import GcpSecret -from apache_beam.transforms.util import Secret +from apache_beam.utils.secret import GcpHsmGeneratedSecret +from apache_beam.utils.secret import GcpSecret +from apache_beam.utils.secret import Secret from apache_beam.transforms.util import _BatchSizeEstimator from apache_beam.transforms.util import _GlobalWindowsBatchingDoFn from apache_beam.transforms.window import FixedWindows @@ -287,37 +287,6 @@ def process(self, element): return final_elements -class SecretTest(unittest.TestCase): - @parameterized.expand([ - param( - secret_string='type:GcpSecret;version_name:my_secret/versions/latest', - secret=GcpSecret('my_secret/versions/latest')), - param( - secret_string='type:GcpSecret;version_name:foo', - secret=GcpSecret('foo')), - param( - secret_string='type:gcpsecreT;version_name:my_secret/versions/latest', - secret=GcpSecret('my_secret/versions/latest')), - ]) - def test_secret_manager_parses_correctly(self, secret_string, secret): - self.assertEqual(secret, Secret.parse_secret_option(secret_string)) - - @parameterized.expand([ - param( - secret_string='version_name:foo', - exception_str='must contain a valid type parameter'), - param( - secret_string='type:gcpsecreT', - exception_str='missing 1 required positional argument'), - param( - secret_string='type:gcpsecreT;version_name:foo;extra:val', - exception_str='Invalid secret parameter extra'), - ]) - def test_secret_manager_throws_on_invalid(self, secret_string, exception_str): - with self.assertRaisesRegex(Exception, exception_str): - Secret.parse_secret_option(secret_string) - - class GroupByEncryptedKeyTest(unittest.TestCase): @classmethod def setUpClass(cls): @@ -387,7 +356,7 @@ def test_gbek_fake_secret_manager_actually_does_encryption(self): result, equal_to([('a', ([1, 2])), ('b', ([3])), ('c', ([4]))])) @mock.patch('apache_beam.transforms.util._DecryptMessage', MockNoOpDecrypt) - @mock.patch('apache_beam.transforms.util.GcpSecret', FakeSecret) + @mock.patch('apache_beam.utils.secret.GcpSecret', FakeSecret) def test_gbk_actually_does_encryption(self): options = PipelineOptions() # Version of GcpSecret doesn't matter since it is replaced by FakeSecret @@ -435,124 +404,6 @@ def test_gbek_gcp_secret_manager_throws(self): result, equal_to([('a', ([1, 2])), ('b', ([3])), ('c', ([4]))])) -@unittest.skipIf(secretmanager is None, 'GCP dependencies are not installed') -class GcpHsmGeneratedSecretTest(unittest.TestCase): - def setUp(self): - self.mock_secret_manager_client = mock.MagicMock() - self.mock_kms_client = mock.MagicMock() - - # Patch the clients - self.secretmanager_patcher = mock.patch( - 'google.cloud.secretmanager.SecretManagerServiceClient', - return_value=self.mock_secret_manager_client) - self.kms_patcher = mock.patch( - 'google.cloud.kms.KeyManagementServiceClient', - return_value=self.mock_kms_client) - self.os_urandom_patcher = mock.patch('os.urandom', return_value=b'0' * 32) - self.hkdf_patcher = mock.patch( - 'cryptography.hazmat.primitives.kdf.hkdf.HKDF.derive', - return_value=b'derived_key') - - self.secretmanager_patcher.start() - self.kms_patcher.start() - self.os_urandom_patcher.start() - self.hkdf_patcher.start() - - def tearDown(self): - self.secretmanager_patcher.stop() - self.kms_patcher.stop() - self.os_urandom_patcher.stop() - self.hkdf_patcher.stop() - - def test_happy_path_secret_creation(self): - from google.api_core import exceptions as api_exceptions - - project_id = 'test-project' - location_id = 'global' - key_ring_id = 'test-key-ring' - key_id = 'test-key' - job_name = 'test-job' - - secret = GcpHsmGeneratedSecret( - project_id, location_id, key_ring_id, key_id, job_name) - - # Mock responses for secret creation path - self.mock_secret_manager_client.access_secret_version.side_effect = [ - api_exceptions.NotFound('not found'), # first check - api_exceptions.NotFound('not found'), # second check - mock.MagicMock(payload=mock.MagicMock(data=b'derived_key')) - ] - self.mock_kms_client.encrypt.return_value = mock.MagicMock( - ciphertext=b'encrypted_nonce') - - secret_bytes = secret.get_secret_bytes() - self.assertEqual(secret_bytes, b'derived_key') - - # Assertions on mocks - secret_version_path = ( - f'projects/{project_id}/secrets/{secret._secret_version_name}' - '/versions/1') - self.mock_secret_manager_client.access_secret_version.assert_any_call( - request={'name': secret_version_path}) - self.assertEqual( - self.mock_secret_manager_client.access_secret_version.call_count, 3) - self.mock_secret_manager_client.create_secret.assert_called_once() - self.mock_kms_client.encrypt.assert_called_once() - self.mock_secret_manager_client.add_secret_version.assert_called_once() - - def test_secret_already_exists(self): - from google.api_core import exceptions as api_exceptions - - project_id = 'test-project' - location_id = 'global' - key_ring_id = 'test-key-ring' - key_id = 'test-key' - job_name = 'test-job' - - secret = GcpHsmGeneratedSecret( - project_id, location_id, key_ring_id, key_id, job_name) - - # Mock responses for secret creation path - self.mock_secret_manager_client.access_secret_version.side_effect = [ - api_exceptions.NotFound('not found'), - api_exceptions.NotFound('not found'), - mock.MagicMock(payload=mock.MagicMock(data=b'derived_key')) - ] - self.mock_secret_manager_client.create_secret.side_effect = ( - api_exceptions.AlreadyExists('exists')) - self.mock_kms_client.encrypt.return_value = mock.MagicMock( - ciphertext=b'encrypted_nonce') - - secret_bytes = secret.get_secret_bytes() - self.assertEqual(secret_bytes, b'derived_key') - - # Assertions on mocks - self.mock_secret_manager_client.create_secret.assert_called_once() - self.mock_secret_manager_client.add_secret_version.assert_called_once() - - def test_secret_version_already_exists(self): - project_id = 'test-project' - location_id = 'global' - key_ring_id = 'test-key-ring' - key_id = 'test-key' - job_name = 'test-job' - - secret = GcpHsmGeneratedSecret( - project_id, location_id, key_ring_id, key_id, job_name) - - self.mock_secret_manager_client.access_secret_version.return_value = ( - mock.MagicMock(payload=mock.MagicMock(data=b'existing_dek'))) - - secret_bytes = secret.get_secret_bytes() - self.assertEqual(secret_bytes, b'existing_dek') - - # Assertions - self.mock_secret_manager_client.access_secret_version.assert_called_once() - self.mock_secret_manager_client.create_secret.assert_not_called() - self.mock_secret_manager_client.add_secret_version.assert_not_called() - self.mock_kms_client.encrypt.assert_not_called() - - class FakeClock(object): def __init__(self, now=time.time()): self._now = now diff --git a/sdks/python/apache_beam/utils/secret.py b/sdks/python/apache_beam/utils/secret.py new file mode 100644 index 000000000000..c9c13f1d60f3 --- /dev/null +++ b/sdks/python/apache_beam/utils/secret.py @@ -0,0 +1,464 @@ +# +# Licensed to the Apache Software Foundation (ASF) under one or more +# contributor license agreements. See the NOTICE file distributed with +# this work for additional information regarding copyright ownership. +# The ASF licenses this file to You under the Apache License, Version 2.0 +# (the "License"); you may not use this file except in compliance with +# the License. You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, software +# distributed under the License is distributed on an "AS IS" BASIS, +# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +# See the License for the specific language governing permissions and +# limitations under the License. +# + +"""Interface and implementations for Secret providers in Apache Beam.""" + +import abc +import json +import logging +import os +import warnings +from typing import Any, Dict, Optional, Union + +_LOGGER = logging.getLogger(__name__) + + +class Secret(abc.ABC): + """A secret management class used for handling sensitive data. + + This class provides a generic interface for secret management. Implementations + of this class should handle fetching secrets from a secret management system. + """ + def __init__(self): + self._cached_secret_bytes: Optional[bytes] = None + + def get_str(self, cacheSecret: bool = False) -> str: + """Retrieve secret value as string. + + Args: + cacheSecret: If True, caches secret value in memory after first fetch. + + Returns: + The retrieved secret value as string. + """ + return self.get_bytes(cacheSecret=cacheSecret).decode("utf-8") + + def get_bytes(self, cacheSecret: bool = False) -> bytes: + """Retrieve secret value as bytes. + + Args: + cacheSecret: If True, caches secret value in memory after first fetch. + + Returns: + The retrieved secret value as bytes. + """ + if cacheSecret and getattr(self, '_cached_secret_bytes', None) is not None: + return self._cached_secret_bytes + + secret_val_bytes = self.get_secret_bytes() + + if cacheSecret: + self._cached_secret_bytes = secret_val_bytes + + return secret_val_bytes + + @abc.abstractmethod + def get_secret_bytes(self) -> bytes: + """Returns the secret as a byte string.""" + raise NotImplementedError() + + def __getstate__(self): + """Strip cached secrets before pickling for pipeline submission/transmission.""" + state = self.__dict__.copy() + state['_cached_secret_bytes'] = None + return state + + @staticmethod + def generate_secret_bytes() -> bytes: + """Generates a new secret key using Fernet.""" + from cryptography.fernet import Fernet + return Fernet.generate_key() + + @classmethod + def parse_secret_option(cls, secret: str) -> 'Secret': + """Parses a secret string and returns the appropriate secret type. + + The secret string should be formatted like: + 'type:;:' + + For example, 'type:GcpSecret;version_name:my_secret/versions/latest' + would return a GcpSecret initialized with 'my_secret/versions/latest'. + """ + param_map = {} + for param in secret.split(';'): + parts = param.split(':') + if len(parts) == 2: + param_map[parts[0]] = parts[1] + + if 'type' not in param_map: + raise ValueError('Secret string must contain a valid type parameter') + + raw_type = param_map.pop('type') + secret_type = raw_type.lower() + secret_manager = _SECRET_TYPE_TO_SECRET_MANAGER.get(secret_type) + if not secret_manager: + raise ValueError( + f'Invalid secret type {secret_type}, currently only ' + 'GcpSecret and GcpHsmGeneratedSecret are supported') + + return cls.from_json(json.dumps(param_map), secret_manager) + + @classmethod + def from_json( + cls, spec: str, secret_manager: Optional[str] = None) -> 'Secret': + """Return a Secret instance based on secret_manager provider and secret specification. + + Args: + spec: Secret string (raw secret or JSON specification string). + secret_manager: Secret manager string (e.g. 'GoogleCloudSecretManager'). + + Returns: + An instance of Secret. + """ + if not isinstance(spec, str): + raise TypeError( + f"Secret 'spec' must be a string, got {type(spec).__name__}") + + secret_manager_name = ( + secret_manager.strip() + if secret_manager and secret_manager.strip() else None) + + spec_dict = None + try: + spec_dict = json.loads(spec) + if not isinstance(spec_dict, dict): + spec_dict = None + except Exception: + try: + import ast + spec_dict = ast.literal_eval(spec) + if not isinstance(spec_dict, dict): + spec_dict = None + except Exception: + pass + + if secret_manager_name: + secret_cls_entry = _SECRET_CLASSES.get(secret_manager_name.lower()) + if secret_cls_entry: + if isinstance(secret_cls_entry, str): + secret_cls = globals().get(secret_cls_entry, secret_cls_entry) + else: + secret_cls = secret_cls_entry + if isinstance(spec_dict, dict) and hasattr(secret_cls, 'from_dict'): + return secret_cls.from_dict(spec_dict) + elif isinstance(spec_dict, dict): + return secret_cls(**spec_dict) + else: + return secret_cls(spec) + else: + raise ValueError( + f"Unsupported secret manager: '{secret_manager_name}'. Currently supported options: 'GoogleCloudSecretManager', 'GoogleCloudHsmGeneratedSecretManager'." + ) + + # If secret_manager is not set or empty, check if spec is a JSON specification dict + if spec_dict is not None: + msg = ( + "The 'spec' parameter appears to be a JSON specification, but " + "'secret_manager' is not set. Defaulting to Raw.") + _LOGGER.warning(msg) + warnings.warn(msg, UserWarning) + + return RawSecret(spec) + + +class RawSecret(Secret): + """Secret implementation wrapping a raw secret string or bytes directly.""" + def __init__(self, secret: Union[str, bytes]): + super().__init__() + if isinstance(secret, str): + self._secret = secret.encode("utf-8") + else: + self._secret = secret + + def get_secret_bytes(self) -> bytes: + return self._secret + + def __eq__(self, other: Any) -> bool: + if not isinstance(other, RawSecret): + return False + return self._secret == other._secret + + +class GcpSecret(Secret): + """A secret manager implementation that retrieves secrets from Google Cloud + Secret Manager. + """ + def __init__(self, version_name: str): + """Initializes a GcpSecret object. + + Args: + version_name: The full version name of the secret in Google Cloud Secret + Manager. For example: + projects//secrets//versions/1. + For more info, see + https://cloud.google.com/python/docs/reference/secretmanager/latest/google.cloud.secretmanager_v1beta1.services.secret_manager_service.SecretManagerServiceClient#google_cloud_secretmanager_v1beta1_services_secret_manager_service_SecretManagerServiceClient_access_secret_version + """ + super().__init__() + self._version_name = version_name + + @classmethod + def from_dict(cls, spec_dict: Dict[str, str]) -> 'GcpSecret': + """Initialize GcpSecret from a dictionary specification.""" + allowed_keys = {'version_name', 'name', 'project', 'version'} + invalid_keys = set(spec_dict.keys()) - allowed_keys + if invalid_keys: + raise ValueError( + f"Invalid secret parameter {', '.join(sorted(invalid_keys))}") + version_name = cls._parse_version_name(spec_dict) + return cls(version_name) + + @classmethod + def _parse_version_name(cls, spec_dict: Dict[str, str]) -> str: + if "version_name" in spec_dict: + return spec_dict["version_name"] + + secret_id = spec_dict.get("name") + if not secret_id: + raise ValueError("Secret name must be specified in secret spec.") + + # Resolve project ID from spec, environment variables, or Application Default Credentials + project_id = ( + spec_dict.get("project") or os.environ.get("GOOGLE_CLOUD_PROJECT") or + os.environ.get("GCP_PROJECT")) + + if not project_id: + try: + import google.auth + _, project_id = google.auth.default() + except Exception: + pass + + version_id = spec_dict.get("version", "latest") + + if not project_id: + raise ValueError( + f"Could not resolve GCP project ID for secret '{secret_id}'. " + "Please specify 'project' in the secret spec, set GOOGLE_CLOUD_PROJECT environment variable, " + "or configure Application Default Credentials.") + + return f"projects/{project_id}/secrets/{secret_id}/versions/{version_id}" + + def get_secret_bytes(self) -> bytes: + try: + from google.cloud import secretmanager + client = secretmanager.SecretManagerServiceClient() + response = client.access_secret_version( + request={"name": self._version_name}) + secret = response.payload.data + return secret + except Exception as e: + raise RuntimeError( + 'Failed to retrieve secret bytes for secret ' + f'{self._version_name} with exception {e}') + + def __eq__(self, secret): + return self._version_name == getattr(secret, '_version_name', None) + + +class GcpHsmGeneratedSecret(Secret): + """A secret manager implementation that generates a secret using a GCP HSM key + and stores it in Google Cloud Secret Manager. If the secret already exists, + it will be retrieved. + """ + def __init__( + self, + project_id: str, + location_id: str, + key_ring_id: str, + key_id: str, + job_name: str): + """Initializes a GcpHsmGeneratedSecret object. + + Args: + project_id: The GCP project ID. + location_id: The GCP location ID for the HSM key. + key_ring_id: The ID of the KMS key ring. + key_id: The ID of the KMS key. + job_name: The name of the job, used to generate a unique secret name. + """ + super().__init__() + self._project_id = project_id + self._location_id = location_id + self._key_ring_id = key_ring_id + self._key_id = key_id + self._job_name = job_name + self._secret_version_name = f'HsmGeneratedSecret_{job_name}' + + def __eq__(self, other: Any) -> bool: + if not isinstance(other, GcpHsmGeneratedSecret): + return False + return ( + self._project_id == other._project_id and + self._location_id == other._location_id and + self._key_ring_id == other._key_ring_id and + self._key_id == other._key_id and + getattr(self, '_job_name', None) == getattr(other, '_job_name', None)) + + @classmethod + def from_dict(cls, spec_dict: Dict[str, str]) -> 'GcpHsmGeneratedSecret': + """Initialize GcpHsmGeneratedSecret from a dictionary specification.""" + allowed_keys = { + 'project_id', 'location_id', 'key_ring_id', 'key_id', 'job_name' + } + missing = allowed_keys - set(spec_dict.keys()) + if missing: + raise ValueError( + f"Missing required parameter(s) for GcpHsmGeneratedSecret: {sorted(list(missing))}" + ) + invalid_keys = set(spec_dict.keys()) - allowed_keys + if invalid_keys: + raise ValueError( + f"Invalid secret parameter {', '.join(sorted(invalid_keys))}") + return cls( + project_id=spec_dict['project_id'], + location_id=spec_dict['location_id'], + key_ring_id=spec_dict['key_ring_id'], + key_id=spec_dict['key_id'], + job_name=spec_dict['job_name'], + ) + + def get_secret_bytes(self) -> bytes: + """Retrieves the secret bytes. + + If the secret version already exists in Secret Manager, it is retrieved. + Otherwise, a new secret and version are created. The new secret is + generated using the HSM key. + + Returns: + The secret as a byte string. + """ + try: + from google.api_core import exceptions as api_exceptions + from google.cloud import secretmanager + client = secretmanager.SecretManagerServiceClient() + + project_path = f"projects/{self._project_id}" + secret_path = f"{project_path}/secrets/{self._secret_version_name}" + # Since we may generate multiple versions when doing this on workers, + # just always take the first version added to maintain consistency. + secret_version_path = f"{secret_path}/versions/1" + + try: + response = client.access_secret_version( + request={"name": secret_version_path}) + return response.payload.data + except api_exceptions.NotFound: + # Don't bother logging yet, we'll only log if we actually add the + # secret version below + pass + + try: + client.create_secret( + request={ + "parent": project_path, + "secret_id": self._secret_version_name, + "secret": { + "replication": { + "automatic": {} + } + }, + }) + except api_exceptions.AlreadyExists: + # Don't bother logging yet, we'll only log if we actually add the + # secret version below + pass + + new_key = self.generate_dek() + try: + # Try one more time in case it was created while we were generating the + # DEK. + response = client.access_secret_version( + request={"name": secret_version_path}) + return response.payload.data + except api_exceptions.NotFound: + _LOGGER.info( + "Secret version %s not found. " + "Creating new secret and version.", + secret_version_path) + client.add_secret_version( + request={ + "parent": secret_path, "payload": { + "data": new_key + } + }) + response = client.access_secret_version( + request={"name": secret_version_path}) + return response.payload.data + + except Exception as e: + raise RuntimeError( + f'Failed to retrieve or create secret bytes for secret ' + f'{self._secret_version_name} with exception {e}') + + def generate_dek(self, dek_size: int = 32) -> bytes: + """Generates a new Data Encryption Key (DEK) using an HSM-backed key. + + This function follows a key derivation process that incorporates entropy + from the HSM-backed key into the nonce used for key derivation. + + Args: + dek_size: The size of the DEK to generate. + + Returns: + A new DEK of the specified size, url-safe base64-encoded. + """ + try: + import base64 + import os + + from cryptography.hazmat.primitives import hashes + from cryptography.hazmat.primitives.kdf.hkdf import HKDF + from google.cloud import kms + + # 1. Generate a random nonce (nonce_one) + nonce_one = os.urandom(dek_size) + + # 2. Use the HSM-backed key to encrypt nonce_one to create nonce_two + kms_client = kms.KeyManagementServiceClient() + key_path = kms_client.crypto_key_path( + self._project_id, self._location_id, self._key_ring_id, self._key_id) + response = kms_client.encrypt( + request={ + 'name': key_path, 'plaintext': nonce_one + }) + nonce_two = response.ciphertext + + # 3. Generate a Derivation Key (DK) + dk = os.urandom(dek_size) + + # 4. Use a KDF to derive the DEK using DK and nonce_two + hkdf = HKDF( + algorithm=hashes.SHA256(), + length=dek_size, + salt=nonce_two, + info=None, + ) + dek = hkdf.derive(dk) + return base64.urlsafe_b64encode(dek) + except Exception as e: + raise RuntimeError(f'Failed to generate DEK with exception {e}') + + +_SECRET_TYPE_TO_SECRET_MANAGER: Dict[str, str] = { + "gcpsecret": "GoogleCloudSecretManager", + "gcphsmgeneratedsecret": "GoogleCloudHsmGeneratedSecretManager", +} + +_SECRET_CLASSES: Dict[str, Any] = { + "googlecloudsecretmanager": "GcpSecret", + "googlecloudhsmgeneratedsecretmanager": "GcpHsmGeneratedSecret", +} \ No newline at end of file diff --git a/sdks/python/apache_beam/utils/secret_test.py b/sdks/python/apache_beam/utils/secret_test.py new file mode 100644 index 000000000000..179b7ca3a5f7 --- /dev/null +++ b/sdks/python/apache_beam/utils/secret_test.py @@ -0,0 +1,454 @@ +# +# Licensed to the Apache Software Foundation (ASF) under one or more +# contributor license agreements. See the NOTICE file distributed with +# this work for additional information regarding copyright ownership. +# The ASF licenses this file to You under the Apache License, Version 2.0 +# (the "License"); you may not use this file except in compliance with +# the License. You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, software +# distributed under the License is distributed on an "AS IS" BASIS, +# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +# See the License for the specific language governing permissions and +# limitations under the License. +# + +import json +import unittest +from unittest import mock + +from parameterized import param +from parameterized import parameterized + +from apache_beam.utils.annotations import BeamDeprecationWarning +from apache_beam.utils.secret import GcpHsmGeneratedSecret +from apache_beam.utils.secret import GcpSecret +from apache_beam.utils.secret import RawSecret +from apache_beam.utils.secret import Secret + +try: + from google.cloud import secretmanager +except ImportError: + secretmanager = None # type: ignore[assignment] + + +class SecretTest(unittest.TestCase): + @parameterized.expand([ + param( + secret_string='type:GcpSecret;version_name:my_secret/versions/latest', + secret=GcpSecret('my_secret/versions/latest')), + param( + secret_string='type:GcpSecret;version_name:foo', + secret=GcpSecret('foo')), + param( + secret_string='type:gcpsecreT;version_name:my_secret/versions/latest', + secret=GcpSecret('my_secret/versions/latest')), + ]) + def test_secret_manager_parses_correctly(self, secret_string, secret): + self.assertEqual(secret, Secret.parse_secret_option(secret_string)) + + @parameterized.expand([ + param( + secret_string='version_name:foo', + exception_str='must contain a valid type parameter'), + param( + secret_string='type:gcpsecreT', + exception_str='Secret name must be specified in secret spec'), + param( + secret_string='type:gcpsecreT;version_name:foo;extra:val', + exception_str='Invalid secret parameter extra'), + ]) + def test_secret_manager_throws_on_invalid(self, secret_string, exception_str): + with self.assertRaisesRegex(Exception, exception_str): + Secret.parse_secret_option(secret_string) + + +@unittest.skipIf(secretmanager is None, 'GCP dependencies are not installed') +class GcpSecretTest(unittest.TestCase): + @mock.patch("google.cloud.secretmanager.SecretManagerServiceClient") + def test_gcp_secret_success(self, mock_client_cls): + mock_client = mock.MagicMock() + mock_client_cls.return_value = mock_client + mock_response = mock.MagicMock() + mock_response.payload.data = b"secret-payload-value" + mock_client.access_secret_version.return_value = mock_response + + spec_dict = {"name": "my-secret", "version": "1", "project": "my-project"} + secret = GcpSecret.from_dict(spec_dict) + + secret_val = secret.get_str(cacheSecret=True) + self.assertEqual(secret_val, "secret-payload-value") + secret_bytes = secret.get_bytes(cacheSecret=True) + self.assertEqual(secret_bytes, b"secret-payload-value") + mock_client.access_secret_version.assert_called_once_with( + request={"name": "projects/my-project/secrets/my-secret/versions/1"}) + + # Second call with cacheSecret=True should return cached value without calling client again + mock_client.reset_mock() + secret_val_cached = secret.get_str(cacheSecret=True) + self.assertEqual(secret_val_cached, "secret-payload-value") + mock_client.access_secret_version.assert_not_called() + + @mock.patch("google.cloud.secretmanager.SecretManagerServiceClient") + def test_gcp_secret_get_bytes_uncached(self, mock_client_cls): + mock_client = mock.MagicMock() + mock_client_cls.return_value = mock_client + mock_response = mock.MagicMock() + mock_response.payload.data = b"secret-payload-value" + mock_client.access_secret_version.return_value = mock_response + + spec_dict = {"name": "my-secret", "project": "my-project"} + secret = GcpSecret.from_dict(spec_dict) + + secret_bytes = secret.get_bytes() + self.assertEqual(secret_bytes, b"secret-payload-value") + self.assertIsNone(secret._cached_secret_bytes) + + @mock.patch("google.cloud.secretmanager.SecretManagerServiceClient") + def test_gcp_secret_getstate_clears_cached_secret(self, mock_client_cls): + mock_client = mock.MagicMock() + mock_client_cls.return_value = mock_client + mock_response = mock.MagicMock() + mock_response.payload.data = b"secret-payload-value" + mock_client.access_secret_version.return_value = mock_response + + spec_dict = {"name": "my-secret", "project": "my-project"} + secret = GcpSecret.from_dict(spec_dict) + + # Cache the secret in memory + secret.get_str(cacheSecret=True) + self.assertEqual(secret._cached_secret_bytes, b"secret-payload-value") + + # When pickled / getstate is called during pipeline submission + state = secret.__getstate__() + self.assertIsNone(state["_cached_secret_bytes"]) + + @mock.patch.dict("os.environ", {"GOOGLE_CLOUD_PROJECT": "env-project-123"}) + @mock.patch("google.cloud.secretmanager.SecretManagerServiceClient") + def test_gcp_secret_env_project_fallback(self, mock_client_cls): + mock_client = mock.MagicMock() + mock_client_cls.return_value = mock_client + mock_response = mock.MagicMock() + mock_response.payload.data = b"env-secret-val" + mock_client.access_secret_version.return_value = mock_response + + # Project omitted from spec + spec_dict = {"name": "env-secret", "version": "latest"} + secret = GcpSecret.from_dict(spec_dict) + + secret_val = secret.get_str(cacheSecret=False) + self.assertEqual(secret_val, "env-secret-val") + self.assertEqual(secret.get_bytes(cacheSecret=False), b"env-secret-val") + mock_client.access_secret_version.assert_called_with( + request={ + "name": "projects/env-project-123/secrets/env-secret/versions/latest" + }) + + @mock.patch("google.cloud.secretmanager.SecretManagerServiceClient") + def test_gcp_secret_failure_raises_exception(self, mock_client_cls): + mock_client = mock.MagicMock() + mock_client_cls.return_value = mock_client + mock_client.access_secret_version.side_effect = RuntimeError( + "Permission denied or secret not found") + + spec_dict = {"name": "non-existent-secret", "project": "my-project"} + secret = GcpSecret.from_dict(spec_dict) + + with self.assertRaises(RuntimeError) as ctx: + secret.get_str(cacheSecret=False) + self.assertIn("Permission denied or secret not found", str(ctx.exception)) + + @mock.patch.dict("os.environ", {}, clear=True) + @mock.patch("google.auth.default", side_effect=Exception("No ADC")) + def test_ill_formed_missing_project_raises_value_error( + self, mock_auth_default): + spec_dict = {"name": "my-secret"} + with self.assertRaises(ValueError) as ctx: + GcpSecret.from_dict(spec_dict) + self.assertIn("Could not resolve GCP project ID", str(ctx.exception)) + + def test_ill_formed_missing_secret_name_raises_value_error(self): + spec_dict = {"project": "my-project"} + with self.assertRaises(ValueError) as ctx: + GcpSecret.from_dict(spec_dict) + self.assertIn("Secret name must be specified", str(ctx.exception)) + + +@unittest.skipIf(secretmanager is None, 'GCP dependencies are not installed') +class GcpHsmGeneratedSecretTest(unittest.TestCase): + def setUp(self): + self.mock_secret_manager_client = mock.MagicMock() + self.mock_kms_client = mock.MagicMock() + + # Patch the clients + self.secretmanager_patcher = mock.patch( + 'google.cloud.secretmanager.SecretManagerServiceClient', + return_value=self.mock_secret_manager_client) + self.kms_patcher = mock.patch( + 'google.cloud.kms.KeyManagementServiceClient', + return_value=self.mock_kms_client) + self.os_urandom_patcher = mock.patch('os.urandom', return_value=b'0' * 32) + self.hkdf_patcher = mock.patch( + 'cryptography.hazmat.primitives.kdf.hkdf.HKDF.derive', + return_value=b'derived_key') + + self.secretmanager_patcher.start() + self.kms_patcher.start() + self.os_urandom_patcher.start() + self.hkdf_patcher.start() + + def tearDown(self): + self.secretmanager_patcher.stop() + self.kms_patcher.stop() + self.os_urandom_patcher.stop() + self.hkdf_patcher.stop() + + def test_happy_path_secret_creation(self): + from google.api_core import exceptions as api_exceptions + + project_id = 'test-project' + location_id = 'global' + key_ring_id = 'test-key-ring' + key_id = 'test-key' + job_name = 'test-job' + + secret = GcpHsmGeneratedSecret( + project_id, location_id, key_ring_id, key_id, job_name) + + # Mock responses for secret creation path + self.mock_secret_manager_client.access_secret_version.side_effect = [ + api_exceptions.NotFound('not found'), # first check + api_exceptions.NotFound('not found'), # second check + mock.MagicMock(payload=mock.MagicMock(data=b'derived_key')) + ] + self.mock_kms_client.encrypt.return_value = mock.MagicMock( + ciphertext=b'encrypted_nonce') + + secret_bytes = secret.get_secret_bytes() + self.assertEqual(secret_bytes, b'derived_key') + + # Assertions on mocks + secret_version_path = ( + f'projects/{project_id}/secrets/{secret._secret_version_name}' + '/versions/1') + self.mock_secret_manager_client.access_secret_version.assert_any_call( + request={'name': secret_version_path}) + self.assertEqual( + self.mock_secret_manager_client.access_secret_version.call_count, 3) + self.mock_secret_manager_client.create_secret.assert_called_once() + self.mock_kms_client.encrypt.assert_called_once() + self.mock_secret_manager_client.add_secret_version.assert_called_once() + + def test_secret_already_exists(self): + from google.api_core import exceptions as api_exceptions + + project_id = 'test-project' + location_id = 'global' + key_ring_id = 'test-key-ring' + key_id = 'test-key' + job_name = 'test-job' + + secret = GcpHsmGeneratedSecret( + project_id, location_id, key_ring_id, key_id, job_name) + + # Mock responses for secret creation path + self.mock_secret_manager_client.access_secret_version.side_effect = [ + api_exceptions.NotFound('not found'), + api_exceptions.NotFound('not found'), + mock.MagicMock(payload=mock.MagicMock(data=b'derived_key')) + ] + self.mock_secret_manager_client.create_secret.side_effect = ( + api_exceptions.AlreadyExists('exists')) + self.mock_kms_client.encrypt.return_value = mock.MagicMock( + ciphertext=b'encrypted_nonce') + + secret_bytes = secret.get_secret_bytes() + self.assertEqual(secret_bytes, b'derived_key') + + # Assertions on mocks + self.mock_secret_manager_client.create_secret.assert_called_once() + self.mock_secret_manager_client.add_secret_version.assert_called_once() + + def test_secret_version_already_exists(self): + project_id = 'test-project' + location_id = 'global' + key_ring_id = 'test-key-ring' + key_id = 'test-key' + job_name = 'test-job' + + secret = GcpHsmGeneratedSecret( + project_id, location_id, key_ring_id, key_id, job_name) + + self.mock_secret_manager_client.access_secret_version.return_value = ( + mock.MagicMock(payload=mock.MagicMock(data=b'existing_dek'))) + + secret_bytes = secret.get_secret_bytes() + self.assertEqual(secret_bytes, b'existing_dek') + + # Assertions + self.mock_secret_manager_client.access_secret_version.assert_called_once() + self.mock_secret_manager_client.create_secret.assert_not_called() + self.mock_secret_manager_client.add_secret_version.assert_not_called() + self.mock_kms_client.encrypt.assert_not_called() + + def test_from_dict_success(self): + spec_dict = { + "project_id": "test-proj", + "location_id": "global", + "key_ring_id": "ring", + "key_id": "key", + "job_name": "my-job" + } + secret = GcpHsmGeneratedSecret.from_dict(spec_dict) + self.assertEqual(secret._project_id, "test-proj") + self.assertEqual(secret._location_id, "global") + self.assertEqual(secret._key_ring_id, "ring") + self.assertEqual(secret._key_id, "key") + self.assertEqual(secret._job_name, "my-job") + self.assertEqual(secret._secret_version_name, "HsmGeneratedSecret_my-job") + + def test_from_dict_missing_params_raises_value_error(self): + spec_dict = {"project_id": "test-proj", "location_id": "global"} + with self.assertRaises(ValueError) as ctx: + GcpHsmGeneratedSecret.from_dict(spec_dict) + self.assertIn("Missing required parameter(s)", str(ctx.exception)) + + @mock.patch("google.cloud.secretmanager.SecretManagerServiceClient") + def test_get_bytes_cached(self, mock_sm_client_cls): + mock_client = mock.MagicMock() + mock_sm_client_cls.return_value = mock_client + mock_response = mock.MagicMock() + mock_response.payload.data = b"hsm-derived-key" + mock_client.access_secret_version.return_value = mock_response + + secret = GcpHsmGeneratedSecret("p", "l", "r", "k", "j") + secret_bytes = secret.get_bytes(cacheSecret=True) + self.assertEqual(secret_bytes, b"hsm-derived-key") + + # Second call uses cache + mock_client.reset_mock() + self.assertEqual(secret.get_bytes(cacheSecret=True), b"hsm-derived-key") + mock_client.access_secret_version.assert_not_called() + + @mock.patch("google.cloud.secretmanager.SecretManagerServiceClient") + def test_getstate_clears_cached_secret(self, mock_sm_client_cls): + mock_client = mock.MagicMock() + mock_sm_client_cls.return_value = mock_client + mock_response = mock.MagicMock() + mock_response.payload.data = b"hsm-derived-key" + mock_client.access_secret_version.return_value = mock_response + + secret = GcpHsmGeneratedSecret("p", "l", "r", "k", "j") + secret.get_bytes(cacheSecret=True) + self.assertEqual(secret._cached_secret_bytes, b"hsm-derived-key") + + state = secret.__getstate__() + self.assertIsNone(state["_cached_secret_bytes"]) + + +class RawSecretTest(unittest.TestCase): + def test_raw_secret_str(self): + secret = RawSecret("STATIC_SECRET_") + self.assertEqual(secret.get_str(cacheSecret=True), "STATIC_SECRET_") + self.assertEqual(secret.get_bytes(cacheSecret=True), b"STATIC_SECRET_") + + def test_raw_secret_bytes(self): + secret = RawSecret(b"STATIC_BYTES_") + self.assertEqual(secret.get_str(cacheSecret=True), "STATIC_BYTES_") + self.assertEqual(secret.get_bytes(cacheSecret=True), b"STATIC_BYTES_") + + +class SecretFactoryTest(unittest.TestCase): + def test_secret_factory(self): + spec = json.dumps({"name": "test-secret", "project": "proj"}) + + # When provider is set to 'GoogleCloudSecretManager' + secret_gcp = Secret.from_json( + spec=spec, secret_manager="GoogleCloudSecretManager") + self.assertIsInstance(secret_gcp, GcpSecret) + + # When spec is a valid JSON string + single_quoted_spec = "{\"name\": \"test-secret\", \"project\": \"proj\"}" + secret_single_quoted = Secret.from_json( + spec=single_quoted_spec, secret_manager="GoogleCloudSecretManager") + self.assertIsInstance(secret_single_quoted, GcpSecret) + self.assertEqual( + secret_single_quoted._version_name, + "projects/proj/secrets/test-secret/versions/latest") + + # When spec is a single-quoted JSON string, we still allow it for convienence + # though it is not a valid JSON string. + single_quoted_spec = "{'name': 'test-secret', 'project': 'proj'}" + secret_single_quoted = Secret.from_json( + spec=single_quoted_spec, secret_manager="GoogleCloudSecretManager") + self.assertIsInstance(secret_single_quoted, GcpSecret) + self.assertEqual( + secret_single_quoted._version_name, + "projects/proj/secrets/test-secret/versions/latest") + + # When provider is None or empty with plain string + secret_raw = Secret.from_json(spec="STATIC_SECRET_", secret_manager=None) + self.assertIsInstance(secret_raw, RawSecret) + + # Unsupported provider raises ValueError + with self.assertRaises(ValueError): + Secret.from_json(spec="spec", secret_manager="unsupported_provider") + + # Non-string spec raises TypeError + spec_dict = {"name": "test-secret"} + with self.assertRaises(TypeError): + Secret.from_json( + spec=spec_dict, # type: ignore[arg-type] + secret_manager="GoogleCloudSecretManager") + + def test_secret_factory_hsm(self): + hsm_spec = json.dumps({ + "project_id": "p", + "location_id": "l", + "key_ring_id": "r", + "key_id": "k", + "job_name": "j" + }) + secret_hsm = Secret.from_json( + spec=hsm_spec, secret_manager="GoogleCloudHsmGeneratedSecretManager") + self.assertIsInstance(secret_hsm, GcpHsmGeneratedSecret) + self.assertEqual(secret_hsm._project_id, "p") + + def test_json_secret_without_secret_manager_warning(self): + json_spec = json.dumps({"name": "my-secret", "project": "my-proj"}) + with self.assertWarns(UserWarning): + secret = Secret.from_json(spec=json_spec, secret_manager=None) + self.assertIsInstance(secret, RawSecret) + + def test_generate_secret_bytes(self): + key = Secret.generate_secret_bytes() + self.assertIsInstance(key, bytes) + self.assertTrue(len(key) > 0) + + def test_equality(self): + raw1 = RawSecret("secret_value") + raw2 = RawSecret("secret_value") + raw3 = RawSecret("other_value") + self.assertEqual(raw1, raw2) + self.assertNotEqual(raw1, raw3) + self.assertNotEqual(raw1, "secret_value") + + gcp1 = GcpSecret.from_dict({"name": "sec", "project": "proj"}) + gcp2 = GcpSecret.from_dict({"name": "sec", "project": "proj"}) + gcp3 = GcpSecret.from_dict({"name": "other", "project": "proj"}) + self.assertEqual(gcp1, gcp2) + self.assertNotEqual(gcp1, gcp3) + self.assertNotEqual(gcp1, raw1) + + hsm1 = GcpHsmGeneratedSecret("p", "l", "r", "k", "j") + hsm2 = GcpHsmGeneratedSecret("p", "l", "r", "k", "j") + hsm3 = GcpHsmGeneratedSecret("p", "l", "r", "k", "other") + self.assertEqual(hsm1, hsm2) + self.assertNotEqual(hsm1, hsm3) + self.assertNotEqual(hsm1, gcp1) + + +if __name__ == "__main__": + unittest.main()