diff --git a/.github/.OwlBot.lock.yaml b/.github/.OwlBot.lock.yaml index aa547962..3815c983 100644 --- a/.github/.OwlBot.lock.yaml +++ b/.github/.OwlBot.lock.yaml @@ -13,4 +13,4 @@ # limitations under the License. docker: image: gcr.io/cloud-devrel-public-resources/owlbot-python:latest - digest: sha256:e09366bdf0fd9c8976592988390b24d53583dd9f002d476934da43725adbb978 + digest: sha256:7a40313731a7cb1454eef6b33d3446ebb121836738dc3ab3d2d3ded5268c35b6 diff --git a/.kokoro/requirements.txt b/.kokoro/requirements.txt index 385f2d4d..d15994ba 100644 --- a/.kokoro/requirements.txt +++ b/.kokoro/requirements.txt @@ -325,31 +325,30 @@ platformdirs==2.5.2 \ --hash=sha256:027d8e83a2d7de06bbac4e5ef7e023c02b863d7ea5d079477e722bb41ab25788 \ --hash=sha256:58c8abb07dcb441e6ee4b11d8df0ac856038f944ab98b7be6b27b2a3c7feef19 # via virtualenv -protobuf==3.20.1 \ - --hash=sha256:06059eb6953ff01e56a25cd02cca1a9649a75a7e65397b5b9b4e929ed71d10cf \ - --hash=sha256:097c5d8a9808302fb0da7e20edf0b8d4703274d140fd25c5edabddcde43e081f \ - --hash=sha256:284f86a6207c897542d7e956eb243a36bb8f9564c1742b253462386e96c6b78f \ - --hash=sha256:32ca378605b41fd180dfe4e14d3226386d8d1b002ab31c969c366549e66a2bb7 \ - --hash=sha256:3cc797c9d15d7689ed507b165cd05913acb992d78b379f6014e013f9ecb20996 \ - --hash=sha256:62f1b5c4cd6c5402b4e2d63804ba49a327e0c386c99b1675c8a0fefda23b2067 \ - --hash=sha256:69ccfdf3657ba59569c64295b7d51325f91af586f8d5793b734260dfe2e94e2c \ - --hash=sha256:6f50601512a3d23625d8a85b1638d914a0970f17920ff39cec63aaef80a93fb7 \ - --hash=sha256:7403941f6d0992d40161aa8bb23e12575637008a5a02283a930addc0508982f9 \ - --hash=sha256:755f3aee41354ae395e104d62119cb223339a8f3276a0cd009ffabfcdd46bb0c \ - --hash=sha256:77053d28427a29987ca9caf7b72ccafee011257561259faba8dd308fda9a8739 \ - --hash=sha256:7e371f10abe57cee5021797126c93479f59fccc9693dafd6bd5633ab67808a91 \ - --hash=sha256:9016d01c91e8e625141d24ec1b20fed584703e527d28512aa8c8707f105a683c \ - --hash=sha256:9be73ad47579abc26c12024239d3540e6b765182a91dbc88e23658ab71767153 \ - --hash=sha256:adc31566d027f45efe3f44eeb5b1f329da43891634d61c75a5944e9be6dd42c9 \ - --hash=sha256:adfc6cf69c7f8c50fd24c793964eef18f0ac321315439d94945820612849c388 \ - --hash=sha256:af0ebadc74e281a517141daad9d0f2c5d93ab78e9d455113719a45a49da9db4e \ - --hash=sha256:cb29edb9eab15742d791e1025dd7b6a8f6fcb53802ad2f6e3adcb102051063ab \ - --hash=sha256:cd68be2559e2a3b84f517fb029ee611546f7812b1fdd0aa2ecc9bc6ec0e4fdde \ - --hash=sha256:cdee09140e1cd184ba9324ec1df410e7147242b94b5f8b0c64fc89e38a8ba531 \ - --hash=sha256:db977c4ca738dd9ce508557d4fce0f5aebd105e158c725beec86feb1f6bc20d8 \ - --hash=sha256:dd5789b2948ca702c17027c84c2accb552fc30f4622a98ab5c51fcfe8c50d3e7 \ - --hash=sha256:e250a42f15bf9d5b09fe1b293bdba2801cd520a9f5ea2d7fb7536d4441811d20 \ - --hash=sha256:ff8d8fa42675249bb456f5db06c00de6c2f4c27a065955917b28c4f15978b9c3 +protobuf==3.20.2 \ + --hash=sha256:03d76b7bd42ac4a6e109742a4edf81ffe26ffd87c5993126d894fe48a120396a \ + --hash=sha256:09e25909c4297d71d97612f04f41cea8fa8510096864f2835ad2f3b3df5a5559 \ + --hash=sha256:18e34a10ae10d458b027d7638a599c964b030c1739ebd035a1dfc0e22baa3bfe \ + --hash=sha256:291fb4307094bf5ccc29f424b42268640e00d5240bf0d9b86bf3079f7576474d \ + --hash=sha256:2c0b040d0b5d5d207936ca2d02f00f765906622c07d3fa19c23a16a8ca71873f \ + --hash=sha256:384164994727f274cc34b8abd41a9e7e0562801361ee77437099ff6dfedd024b \ + --hash=sha256:3cb608e5a0eb61b8e00fe641d9f0282cd0eedb603be372f91f163cbfbca0ded0 \ + --hash=sha256:5d9402bf27d11e37801d1743eada54372f986a372ec9679673bfcc5c60441151 \ + --hash=sha256:712dca319eee507a1e7df3591e639a2b112a2f4a62d40fe7832a16fd19151750 \ + --hash=sha256:7a5037af4e76c975b88c3becdf53922b5ffa3f2cddf657574a4920a3b33b80f3 \ + --hash=sha256:8228e56a865c27163d5d1d1771d94b98194aa6917bcfb6ce139cbfa8e3c27334 \ + --hash=sha256:84a1544252a933ef07bb0b5ef13afe7c36232a774affa673fc3636f7cee1db6c \ + --hash=sha256:84fe5953b18a383fd4495d375fe16e1e55e0a3afe7b4f7b4d01a3a0649fcda9d \ + --hash=sha256:9c673c8bfdf52f903081816b9e0e612186684f4eb4c17eeb729133022d6032e3 \ + --hash=sha256:9f876a69ca55aed879b43c295a328970306e8e80a263ec91cf6e9189243c613b \ + --hash=sha256:a9e5ae5a8e8985c67e8944c23035a0dff2c26b0f5070b2f55b217a1c33bbe8b1 \ + --hash=sha256:b4fdb29c5a7406e3f7ef176b2a7079baa68b5b854f364c21abe327bbeec01cdb \ + --hash=sha256:c184485e0dfba4dfd451c3bd348c2e685d6523543a0f91b9fd4ae90eb09e8422 \ + --hash=sha256:c9cdf251c582c16fd6a9f5e95836c90828d51b0069ad22f463761d27c6c19019 \ + --hash=sha256:e39cf61bb8582bda88cdfebc0db163b774e7e03364bbf9ce1ead13863e81e359 \ + --hash=sha256:e8fbc522303e09036c752a0afcc5c0603e917222d8bedc02813fd73b4b4ed804 \ + --hash=sha256:f34464ab1207114e73bba0794d1257c150a2b89b7a9faf504e00af7c9fd58978 \ + --hash=sha256:f52dabc96ca99ebd2169dadbe018824ebda08a795c7684a0b7d203a290f3adb0 # via # gcp-docuploader # gcp-releasetool diff --git a/CHANGELOG.md b/CHANGELOG.md index 6d9d100f..15ae0121 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1,5 +1,12 @@ # Changelog +## [1.6.0](https://github.com/googleapis/python-pubsublite/compare/v1.5.0...v1.6.0) (2022-10-06) + + +### Features + +* Create subscriptions at a seek target ([#383](https://github.com/googleapis/python-pubsublite/issues/383)) ([864c0cc](https://github.com/googleapis/python-pubsublite/commit/864c0ccbd9f6575fc8aa96f80bd7959c1e78e66e)) + ## [1.5.0](https://github.com/googleapis/python-pubsublite/compare/v1.4.3...v1.5.0) (2022-09-13) diff --git a/google/cloud/pubsublite/__init__.py b/google/cloud/pubsublite/__init__.py index 7a064a8a..5a7872af 100644 --- a/google/cloud/pubsublite/__init__.py +++ b/google/cloud/pubsublite/__init__.py @@ -70,6 +70,7 @@ from google.cloud.pubsublite_v1.types.admin import UpdateTopicRequest from google.cloud.pubsublite_v1.types.common import AttributeValues from google.cloud.pubsublite_v1.types.common import Cursor +from google.cloud.pubsublite_v1.types.common import ExportConfig from google.cloud.pubsublite_v1.types.common import PubSubMessage from google.cloud.pubsublite_v1.types.common import Reservation from google.cloud.pubsublite_v1.types.common import SequencedMessage @@ -131,6 +132,7 @@ "CursorServiceClient", "DeleteSubscriptionRequest", "DeleteTopicRequest", + "ExportConfig", "FlowControlRequest", "GetSubscriptionRequest", "GetTopicPartitionsRequest", diff --git a/google/cloud/pubsublite/admin_client.py b/google/cloud/pubsublite/admin_client.py index 029017e6..dc1d3b9c 100644 --- a/google/cloud/pubsublite/admin_client.py +++ b/google/cloud/pubsublite/admin_client.py @@ -114,9 +114,10 @@ def list_topic_subscriptions(self, topic_path: TopicPath) -> List[SubscriptionPa def create_subscription( self, subscription: Subscription, - starting_offset: BacklogLocation = BacklogLocation.END, + target: Union[BacklogLocation, PublishTime, EventTime] = BacklogLocation.END, + starting_offset: Optional[BacklogLocation] = None, ) -> Subscription: - return self._impl.create_subscription(subscription, starting_offset) + return self._impl.create_subscription(subscription, target, starting_offset) @overrides def get_subscription(self, subscription_path: SubscriptionPath) -> Subscription: diff --git a/google/cloud/pubsublite/admin_client_interface.py b/google/cloud/pubsublite/admin_client_interface.py index b631edb7..9603ed35 100644 --- a/google/cloud/pubsublite/admin_client_interface.py +++ b/google/cloud/pubsublite/admin_client_interface.py @@ -13,7 +13,7 @@ # limitations under the License. from abc import ABC, abstractmethod -from typing import List, Union +from typing import List, Optional, Union from google.api_core.operation import Operation from google.cloud.pubsublite.types import ( @@ -71,11 +71,20 @@ def list_topic_subscriptions(self, topic_path: TopicPath) -> List[SubscriptionPa def create_subscription( self, subscription: Subscription, - starting_offset: BacklogLocation = BacklogLocation.END, + target: Union[BacklogLocation, PublishTime, EventTime] = BacklogLocation.END, + starting_offset: Optional[BacklogLocation] = None, ) -> Subscription: """Create a subscription, returns the created subscription. By default a subscription will only receive messages published after the - subscription was created.""" + subscription was created. + + `starting_offset` is deprecated. Use `target` to initialize the + subscription to a target location within the message backlog instead. + `starting_offset` has higher precedence if `target` is also set. + + A seek is initiated if the target location is a publish or event time. + If the seek fails, the created subscription is not deleted. + """ @abstractmethod def get_subscription(self, subscription_path: SubscriptionPath) -> Subscription: diff --git a/google/cloud/pubsublite/internal/wire/admin_client_impl.py b/google/cloud/pubsublite/internal/wire/admin_client_impl.py index ee269f43..c295a076 100644 --- a/google/cloud/pubsublite/internal/wire/admin_client_impl.py +++ b/google/cloud/pubsublite/internal/wire/admin_client_impl.py @@ -12,7 +12,8 @@ # See the License for the specific language governing permissions and # limitations under the License. -from typing import List, Union +import logging +from typing import List, Optional, Union from google.api_core.exceptions import InvalidArgument from google.api_core.operation import Operation @@ -38,8 +39,11 @@ TimeTarget, SeekSubscriptionRequest, CreateSubscriptionRequest, + ExportConfig, ) +log = logging.getLogger(__name__) + class AdminClientImpl(AdminClientInterface): _underlying: AdminServiceClient @@ -85,17 +89,51 @@ def list_topic_subscriptions(self, topic_path: TopicPath) -> List[SubscriptionPa def create_subscription( self, subscription: Subscription, - starting_offset: BacklogLocation = BacklogLocation.END, + target: Union[BacklogLocation, PublishTime, EventTime] = BacklogLocation.END, + starting_offset: Optional[BacklogLocation] = None, ) -> Subscription: + if starting_offset: + log.warning("starting_offset is deprecated. Use target instead.") + target = starting_offset path = SubscriptionPath.parse(subscription.name) - return self._underlying.create_subscription( + requires_seek = isinstance(target, PublishTime) or isinstance(target, EventTime) + requires_update = ( + requires_seek + and subscription.export_config + and subscription.export_config.desired_state == ExportConfig.State.ACTIVE + ) + if requires_update: + # Export subscriptions must be paused while seeking. The state is + # later updated to active. + subscription.export_config.desired_state = ExportConfig.State.PAUSED + + # Request 1 - Create the subscription. + skip_backlog = False + if isinstance(target, BacklogLocation): + skip_backlog = target == BacklogLocation.END + response = self._underlying.create_subscription( request=CreateSubscriptionRequest( parent=str(path.to_location_path()), subscription=subscription, subscription_id=path.name, - skip_backlog=(starting_offset == BacklogLocation.END), + skip_backlog=skip_backlog, ) ) + # Request 2 (optional) - seek the subscription. + if requires_seek: + self.seek_subscription(subscription_path=path, target=target) + # Request 3 (optional) - make the export subscription active. + if requires_update: + response = self.update_subscription( + subscription=Subscription( + name=response.name, + export_config=ExportConfig( + desired_state=ExportConfig.State.ACTIVE, + ), + ), + update_mask=FieldMask(paths=["export_config.desired_state"]), + ) + return response def get_subscription(self, subscription_path: SubscriptionPath) -> Subscription: return self._underlying.get_subscription(name=str(subscription_path)) diff --git a/samples/snippets/requirements-test.txt b/samples/snippets/requirements-test.txt index b3cb43b4..ab0e35a0 100644 --- a/samples/snippets/requirements-test.txt +++ b/samples/snippets/requirements-test.txt @@ -1,2 +1,2 @@ -backoff==2.1.2 +backoff==2.2.1 pytest==7.1.3 \ No newline at end of file diff --git a/samples/snippets/requirements.txt b/samples/snippets/requirements.txt index e3953811..08fd31f7 100644 --- a/samples/snippets/requirements.txt +++ b/samples/snippets/requirements.txt @@ -1 +1 @@ -google-cloud-pubsublite==1.4.3 \ No newline at end of file +google-cloud-pubsublite==1.5.0 \ No newline at end of file diff --git a/setup.py b/setup.py index 673cd1e9..fda018eb 100644 --- a/setup.py +++ b/setup.py @@ -19,7 +19,7 @@ import os import setuptools # type: ignore -version = "1.5.0" +version = "1.6.0" package_root = os.path.abspath(os.path.dirname(__file__))