Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 2 additions & 2 deletions .github/trigger_files/beam_PostCommit_Python.json
Original file line number Diff line number Diff line change
@@ -1,5 +1,5 @@
{
"comment": "Modify this file in a trivial way to cause this test suite to run.",
"pr": "37345",
"modification": 53
"pr": "38701",
"modification": 55
}
3 changes: 3 additions & 0 deletions runners/google-cloud-dataflow-java/build.gradle
Original file line number Diff line number Diff line change
Expand Up @@ -113,6 +113,9 @@ dependencies {
implementation library.java.google_cloud_logging
permitUnusedDeclared library.java.google_cloud_logging // BEAM-11761
implementation library.java.opentelemetry_context
runtimeOnly library.java.opentelemetry_exporter_otlp
runtimeOnly library.java.opentelemetry_extension_autoconfigure
runtimeOnly project(":sdks:java:extensions:opentelemetry-gcp-auth-extension")
implementation library.java.hamcrest
implementation library.java.jackson_annotations
implementation library.java.jackson_core
Expand Down
20 changes: 10 additions & 10 deletions sdks/go.mod
Original file line number Diff line number Diff line change
Expand Up @@ -33,10 +33,10 @@ require (
cloud.google.com/go/spanner v1.92.0
cloud.google.com/go/storage v1.63.0
github.com/aws/aws-sdk-go-v2 v1.42.0
github.com/aws/aws-sdk-go-v2/config v1.32.25
github.com/aws/aws-sdk-go-v2/credentials v1.19.24
github.com/aws/aws-sdk-go-v2/feature/s3/manager v1.22.28
github.com/aws/aws-sdk-go-v2/service/s3 v1.104.0
github.com/aws/aws-sdk-go-v2/config v1.32.26
github.com/aws/aws-sdk-go-v2/credentials v1.19.25
github.com/aws/aws-sdk-go-v2/feature/s3/manager v1.22.29
github.com/aws/aws-sdk-go-v2/service/s3 v1.104.1
github.com/aws/smithy-go v1.27.3
github.com/docker/go-connections v0.7.0 // indirect
github.com/dustin/go-humanize v1.0.1
Expand All @@ -46,7 +46,7 @@ require (
github.com/johannesboyne/gofakes3 v0.0.0-20250106100439-5c39aecd6999
github.com/lib/pq v1.12.3
github.com/linkedin/goavro/v2 v2.15.0
github.com/nats-io/nats-server/v2 v2.14.2
github.com/nats-io/nats-server/v2 v2.14.3
github.com/nats-io/nats.go v1.52.0
github.com/proullon/ramsql v0.1.4
github.com/spf13/cobra v1.10.2
Expand Down Expand Up @@ -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.0-default-no-op // indirect
github.com/apache/arrow/go/v15 v15.0.2 // indirect
github.com/aws/aws-sdk-go-v2/service/signin v1.2.0 // indirect
github.com/aws/aws-sdk-go-v2/service/signin v1.2.1 // 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
Expand Down Expand Up @@ -156,10 +156,10 @@ require (
github.com/aws/aws-sdk-go-v2/service/internal/accept-encoding v1.13.12 // indirect
github.com/aws/aws-sdk-go-v2/service/internal/checksum v1.9.22 // indirect
github.com/aws/aws-sdk-go-v2/service/internal/presigned-url v1.13.29 // indirect
github.com/aws/aws-sdk-go-v2/service/internal/s3shared v1.19.29 // indirect
github.com/aws/aws-sdk-go-v2/service/sso v1.31.3 // indirect
github.com/aws/aws-sdk-go-v2/service/ssooidc v1.36.6 // indirect
github.com/aws/aws-sdk-go-v2/service/sts v1.43.3 // indirect
github.com/aws/aws-sdk-go-v2/service/internal/s3shared v1.19.30 // indirect
github.com/aws/aws-sdk-go-v2/service/sso v1.31.4 // indirect
github.com/aws/aws-sdk-go-v2/service/ssooidc v1.36.7 // indirect
github.com/aws/aws-sdk-go-v2/service/sts v1.43.4 // 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
Expand Down
40 changes: 20 additions & 20 deletions sdks/go.sum
Original file line number Diff line number Diff line change
Expand Up @@ -207,20 +207,20 @@ github.com/aws/aws-sdk-go-v2/aws/protocol/eventstream v1.7.13 h1:p1BBrg/Hhp6uK7z
github.com/aws/aws-sdk-go-v2/aws/protocol/eventstream v1.7.13/go.mod h1:8cIfkE9MDhkRZGpQ22aV6/lkYeYSozpz16Smrs5x4Ls=
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.25 h1:ACCejvStYoilgwrfegSt5ZntCbPrk52qfwyNcnl3omM=
github.com/aws/aws-sdk-go-v2/config v1.32.25/go.mod h1:LJyU8sDRbXUxFn8xMJIGP+v9QYYwveNLI8a/giAOiAs=
github.com/aws/aws-sdk-go-v2/config v1.32.26 h1:JI+W5B3jUA8UBz2ggbICGd9UCR6/+SB21G8EFl0SFTQ=
github.com/aws/aws-sdk-go-v2/config v1.32.26/go.mod h1:RLE2Ls/wRstvdSz1GPrIWNnXcKZ/znDdWyMuiQxdBoY=
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.24 h1:2hQqYCV9yqyePQ9o6dCrZc/zO8U3TwPr9mIKlZnPu/I=
github.com/aws/aws-sdk-go-v2/credentials v1.19.24/go.mod h1:IDwpACtwqHLISdzfwUUNq4P9DsB/h5BLg4FwJPNfqFY=
github.com/aws/aws-sdk-go-v2/credentials v1.19.25 h1:TzPVjfUZ1hsKafvYE+DIzKXIik2KufQxsPHanlkttbo=
github.com/aws/aws-sdk-go-v2/credentials v1.19.25/go.mod h1:K4hw0buguVvtC74HnVfTRr0LzQQHAWPqJbBU9QGk2Pg=
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.29 h1:r6qZHbT+wxgWO/e9vYNUEtg7lv5+UN3pRqKhLXvnArg=
github.com/aws/aws-sdk-go-v2/feature/ec2/imds v1.18.29/go.mod h1:QRnaRcTVGKPGRy8w78HMQtKUGRYcnMZAANATkeVA6Mo=
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.28 h1:ez4y5o7sa0uaRI8BquYOXtZpioUPhbQEh7Igm88oV9U=
github.com/aws/aws-sdk-go-v2/feature/s3/manager v1.22.28/go.mod h1:TpmZOrQA12XKEpVypgBGZSQBsm1WUTndCiSnbDsbvug=
github.com/aws/aws-sdk-go-v2/feature/s3/manager v1.22.29 h1:YteQL/8ZD9nS/eiLq3Ab9ldHBUvfWorEOufALZOXoXY=
github.com/aws/aws-sdk-go-v2/feature/s3/manager v1.22.29/go.mod h1:3hj0jtS3hQmQAAZ0yz/jTA+uTAHfDpwxbISkgZ2E0d0=
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.29 h1:f3vKqSo13fhTYb+JEcXwXefZQE26I1FB5eTSniU67ko=
Expand Down Expand Up @@ -248,30 +248,30 @@ github.com/aws/aws-sdk-go-v2/service/internal/presigned-url v1.13.29 h1:DRebniUG
github.com/aws/aws-sdk-go-v2/service/internal/presigned-url v1.13.29/go.mod h1:LfRkPCD8YHDM2E5eTkos2UpwYeZnBcVarTa8L59bJHA=
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.29 h1:hiME6pBzC7OTl9LMtlyTWBuEl1f4QBcUmFDKC7MLXtc=
github.com/aws/aws-sdk-go-v2/service/internal/s3shared v1.19.29/go.mod h1:G7RP+uhagpKtKhd1BM9N6JQqjCcGEU47K5lBVZQyRQw=
github.com/aws/aws-sdk-go-v2/service/internal/s3shared v1.19.30 h1:4HbXxyipSYxexU0juMIpdS05dilL6dbB2VQHxxN2vGU=
github.com/aws/aws-sdk-go-v2/service/internal/s3shared v1.19.30/go.mod h1:G7RP+uhagpKtKhd1BM9N6JQqjCcGEU47K5lBVZQyRQw=
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.104.0 h1:ta8csKy5vN91F3i5gGR85lFV0srBqySEji7Jroes6rE=
github.com/aws/aws-sdk-go-v2/service/s3 v1.104.0/go.mod h1:77ZAgynvx1txMvDG8gGWoWkO1augYDxkp9JElWFgjQU=
github.com/aws/aws-sdk-go-v2/service/s3 v1.104.1 h1:yb03KevaOAG5e8suo79Af74vjIQvoeKmjl79WQchLrs=
github.com/aws/aws-sdk-go-v2/service/s3 v1.104.1/go.mod h1:mreYODw0Y4yv7xeczvqC6vciwFao8lPE9k1l1ulfY6E=
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.2.0 h1:3nXpRcFwRCW8n7HgO2QGy0Dc20eQNfBuUemGQhpF8m8=
github.com/aws/aws-sdk-go-v2/service/signin v1.2.0/go.mod h1:LxYujSTLPRlp2vTtcUO/+1ilrew8ytt6SvQyOgejzFQ=
github.com/aws/aws-sdk-go-v2/service/signin v1.2.1 h1:BeJmkm5YOZs6lGRGcNoIuLSoTTtGLLCEqlSiRKYodfM=
github.com/aws/aws-sdk-go-v2/service/signin v1.2.1/go.mod h1:LxYujSTLPRlp2vTtcUO/+1ilrew8ytt6SvQyOgejzFQ=
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.31.3 h1:ey1XLTYXb9PcLt4535632o5kCGXNXEhNb620Dqwuylo=
github.com/aws/aws-sdk-go-v2/service/sso v1.31.3/go.mod h1:Lk7PlmoTYryQmyBG0EXqj5BcUbj3whXdU2s3yGI3EAc=
github.com/aws/aws-sdk-go-v2/service/sso v1.31.4 h1:i465b/3c7xJd++pobNIDOggouekCuiWOnB0goQJy+94=
github.com/aws/aws-sdk-go-v2/service/sso v1.31.4/go.mod h1:Lk7PlmoTYryQmyBG0EXqj5BcUbj3whXdU2s3yGI3EAc=
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.36.6 h1:yLr03zQE/5Eu5l3QU0Si+xMbLMbSDF2YXsigqXngs6g=
github.com/aws/aws-sdk-go-v2/service/ssooidc v1.36.6/go.mod h1:Q5N6icH+KJZDLh+ESNwzdv6cZ6vLFF/egy3IOxWhmz4=
github.com/aws/aws-sdk-go-v2/service/ssooidc v1.36.7 h1:xbmJAnBbyYPkTzoCNCF/bpJ6ymQHRdXX1vquYfDIGYk=
github.com/aws/aws-sdk-go-v2/service/ssooidc v1.36.7/go.mod h1:Q5N6icH+KJZDLh+ESNwzdv6cZ6vLFF/egy3IOxWhmz4=
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.43.3 h1:VrIhKRCSK1umelSgB9RghvA9RTUYeQffyAS5ApXehNI=
github.com/aws/aws-sdk-go-v2/service/sts v1.43.3/go.mod h1:r8wkDOuLaaMFqFiYAb8dGY2A3gJCOujMc6CFOVC4Zhc=
github.com/aws/aws-sdk-go-v2/service/sts v1.43.4 h1:Np0vmL7op0Zs5xGacYMMX3v5O5pvZ46xhb5LwDgPj8M=
github.com/aws/aws-sdk-go-v2/service/sts v1.43.4/go.mod h1:r8wkDOuLaaMFqFiYAb8dGY2A3gJCOujMc6CFOVC4Zhc=
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.3 h1:F3Zb497UhhskkfpJmfkXswyo+t0sh9OTBnIHjogWbVY=
Expand Down Expand Up @@ -712,8 +712,8 @@ github.com/montanaflynn/stats v0.9.0 h1:tsBJ0RXwph9BmAuFoCmqGv6e8xa0MENQ8m0ptKq2
github.com/montanaflynn/stats v0.9.0/go.mod h1:etXPPgVO6n31NxCd9KQUMvCM+ve0ruNzt6R8Bnaayow=
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.2 h1:Q7dRhCY03Y00rETFW3KV+KGaCIajlDfWgWUVgbMxyuk=
github.com/nats-io/nats-server/v2 v2.14.2/go.mod h1:lWpb1bSpRELZfRdlMkdz8E7lbXKKyNe8RIn0vvepIHs=
github.com/nats-io/nats-server/v2 v2.14.3 h1:+xjydPt7rkit67G+04TN0mcO2n+8nveZE7tK/PPV53A=
github.com/nats-io/nats-server/v2 v2.14.3/go.mod h1:5IlCtBzfwyzQzPMjmoJ9W2/LKmnJRtNyuOs/OT+NHDY=
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/nkeys v0.4.16 h1:rd5oAuLOb8mnAycB0xleuEBNS1pVVnN0fv/FF34Eypg=
Expand Down
24 changes: 21 additions & 3 deletions sdks/python/apache_beam/examples/inference/vllm_text_completion.py
Original file line number Diff line number Diff line change
Expand Up @@ -138,6 +138,20 @@ def parse_known_args(argv):
'Passed to the vLLM OpenAI server as --gpu-memory-utilization '
'(fraction of total GPU memory for KV cache). Lower this if the '
'engine fails to start with CUDA out of memory.'))
parser.add_argument(
'--use_dynamo',
dest='use_dynamo',
action='store_true',
help=(
'Use embedded NVIDIA Dynamo as the vLLM engine. Requires '
'ai-dynamo[vllm] and the etcd binary in the runtime environment. '
'See VLLMCompletionsModelHandler for limitations of embedded mode.'))
parser.add_argument(
'--max_tokens',
dest='max_tokens',
type=int,
default=16,
help='Maximum number of tokens to generate for each example.')
return parser.parse_known_args(argv)


Expand Down Expand Up @@ -178,22 +192,26 @@ def run(
build_vllm_server_kwargs(known_args))

model_handler = VLLMCompletionsModelHandler(
model_name=known_args.model, vllm_server_kwargs=effective_vllm_kwargs)
model_name=known_args.model,
vllm_server_kwargs=effective_vllm_kwargs,
use_dynamo=known_args.use_dynamo)
input_examples = COMPLETION_EXAMPLES

if known_args.chat:
model_handler = VLLMChatModelHandler(
model_name=known_args.model,
chat_template_path=known_args.chat_template,
vllm_server_kwargs=dict(effective_vllm_kwargs))
vllm_server_kwargs=dict(effective_vllm_kwargs),
use_dynamo=known_args.use_dynamo)
input_examples = CHAT_EXAMPLES

pipeline = test_pipeline
if not test_pipeline:
pipeline = beam.Pipeline(options=pipeline_options)

examples = pipeline | "Create examples" >> beam.Create(input_examples)
predictions = examples | "RunInference" >> RunInference(model_handler)
predictions = examples | "RunInference" >> RunInference(
model_handler, inference_args={'max_tokens': known_args.max_tokens})
process_output = predictions | "Process Predictions" >> beam.ParDo(
PostProcessor())
_ = process_output | "WriteOutput" >> beam.io.WriteToText(
Expand Down
79 changes: 72 additions & 7 deletions sdks/python/apache_beam/io/gcp/pubsub_test.py
Original file line number Diff line number Diff line change
Expand Up @@ -502,6 +502,70 @@ def finish_bundle(self):
@unittest.skipIf(pubsub is None, 'GCP dependencies are not installed')
@mock.patch('google.cloud.pubsub.SubscriberClient')
class TestReadFromPubSub(unittest.TestCase):
def setUp(self):
_PubSubReadEvaluator._subscription_cache.clear()
_PubSubReadEvaluator._subscriber_client_cache.clear()

def test_subscriber_client_is_reused_for_transform(self, mock_pubsub):
class Transform(object):
pass

transform = Transform()
first_client = _PubSubReadEvaluator._get_subscriber_client(transform)
second_client = _PubSubReadEvaluator._get_subscriber_client(transform)

self.assertIs(first_client, second_client)
mock_pubsub.assert_called_once_with()
first_client.close.assert_not_called()

def test_subscription_creation_does_not_hold_global_cache_lock(
self, mock_pubsub):
class Transform(object):
pass

def subscription_path(project, subscription):
return 'projects/%s/subscriptions/%s' % (project, subscription)

transform = Transform()
client = mock_pubsub.return_value
client.subscription_path.side_effect = subscription_path
client.topic_path.return_value = 'projects/topic_project/topics/topic'
global_lock_available = []

def create_subscription(name, topic):
self.assertTrue(name.startswith('projects/sub_project/subscriptions/'))
self.assertEqual('projects/topic_project/topics/topic', topic)
acquired = _PubSubReadEvaluator._subscriber_client_cache_lock.acquire(
blocking=False)
global_lock_available.append(acquired)
if acquired:
_PubSubReadEvaluator._subscriber_client_cache_lock.release()

client.create_subscription.side_effect = create_subscription

sub_name = _PubSubReadEvaluator.get_subscription(
transform, 'topic_project', 'topic', 'sub_project', None)

self.assertTrue(global_lock_available)
self.assertTrue(global_lock_available[0])
self.assertTrue(sub_name.startswith('projects/sub_project/subscriptions/'))

def test_subscriber_client_cleanup_is_idempotent(self, unused_mock_pubsub):
client = mock.Mock()
subscriber_client = transform_evaluator._PubSubSubscriberClient(client)
subscriber_client.set_temporary_subscription('subscription')

subscriber_client.close()
subscriber_client.close()

client.assert_has_calls([
mock.call.delete_subscription(subscription='subscription'),
mock.call.close()
])
client.delete_subscription.assert_called_once_with(
subscription='subscription')
client.close.assert_called_once_with()

def test_read_messages_success(self, mock_pubsub):
data = b'data'
publish_time_secs = 1520861821
Expand Down Expand Up @@ -533,7 +597,8 @@ def test_read_messages_success(self, mock_pubsub):
mock_pubsub.return_value.acknowledge.assert_has_calls(
[mock.call(subscription=mock.ANY, ack_ids=[ack_id])])

mock_pubsub.return_value.close.assert_has_calls([mock.call()])
mock_pubsub.assert_called_once_with()
mock_pubsub.return_value.close.assert_not_called()

def test_read_strings_success(self, mock_pubsub):
data = '🤷 ¯\\_(ツ)_/¯'
Expand All @@ -555,7 +620,7 @@ def test_read_strings_success(self, mock_pubsub):
mock_pubsub.return_value.acknowledge.assert_has_calls(
[mock.call(subscription=mock.ANY, ack_ids=[ack_id])])

mock_pubsub.return_value.close.assert_has_calls([mock.call()])
mock_pubsub.return_value.close.assert_not_called()

def test_read_data_success(self, mock_pubsub):
data_encoded = '🤷 ¯\\_(ツ)_/¯'.encode('utf-8')
Expand All @@ -575,7 +640,7 @@ def test_read_data_success(self, mock_pubsub):
mock_pubsub.return_value.acknowledge.assert_has_calls(
[mock.call(subscription=mock.ANY, ack_ids=[ack_id])])

mock_pubsub.return_value.close.assert_has_calls([mock.call()])
mock_pubsub.return_value.close.assert_not_called()

def test_read_messages_timestamp_attribute_milli_success(self, mock_pubsub):
data = b'data'
Expand Down Expand Up @@ -610,7 +675,7 @@ def test_read_messages_timestamp_attribute_milli_success(self, mock_pubsub):
mock_pubsub.return_value.acknowledge.assert_has_calls(
[mock.call(subscription=mock.ANY, ack_ids=[ack_id])])

mock_pubsub.return_value.close.assert_has_calls([mock.call()])
mock_pubsub.return_value.close.assert_not_called()

def test_read_messages_timestamp_attribute_rfc3339_success(self, mock_pubsub):
data = b'data'
Expand Down Expand Up @@ -645,7 +710,7 @@ def test_read_messages_timestamp_attribute_rfc3339_success(self, mock_pubsub):
mock_pubsub.return_value.acknowledge.assert_has_calls(
[mock.call(subscription=mock.ANY, ack_ids=[ack_id])])

mock_pubsub.return_value.close.assert_has_calls([mock.call()])
mock_pubsub.return_value.close.assert_not_called()

def test_read_messages_timestamp_attribute_missing(self, mock_pubsub):
data = b'data'
Expand Down Expand Up @@ -681,7 +746,7 @@ def test_read_messages_timestamp_attribute_missing(self, mock_pubsub):
mock_pubsub.return_value.acknowledge.assert_has_calls(
[mock.call(subscription=mock.ANY, ack_ids=[ack_id])])

mock_pubsub.return_value.close.assert_has_calls([mock.call()])
mock_pubsub.return_value.close.assert_not_called()

def test_read_messages_timestamp_attribute_fail_parse(self, mock_pubsub):
data = b'data'
Expand Down Expand Up @@ -710,7 +775,7 @@ def test_read_messages_timestamp_attribute_fail_parse(self, mock_pubsub):
p.run()
mock_pubsub.return_value.acknowledge.assert_not_called()

mock_pubsub.return_value.close.assert_has_calls([mock.call()])
mock_pubsub.return_value.close.assert_not_called()

def test_read_message_id_label_unsupported(self, unused_mock_pubsub):
# id_label is unsupported in DirectRunner.
Expand Down
Loading
Loading