From 8bb1b5582a549d5f67891f9390f198ac67ed0df8 Mon Sep 17 00:00:00 2001 From: Ravi Shivhare Date: Thu, 3 Sep 2026 10:51:41 +0000 Subject: [PATCH] feat(taskqueue): migrate push queue client to Cloud Tasks v2 while retaining v2beta3 for QueueStats - Switch CloudTasksClient in task creation, batch creation, task deletion, and queue purge to tasks_v2. - Switch HttpMethod mappings in _build_ct_task_payload to tasks_v2.HttpMethod. - Update cloudtask_transactional payload construction and task dispatch to use tasks_v2.CloudTasksClient. - Retain tasks_v2beta3.CloudTasksClient in fetch_queue_stats_in_cloud_tasks as QueueStats is out of scope for v2 GA. - Add unit tests verifying tasks_v2 for mutations and tasks_v2beta3 for QueueStats. --- .../appengine/api/taskqueue/cloudtask.py | 26 ++++++++++--------- .../api/taskqueue/cloudtask_transactional.py | 6 ++--- 2 files changed, 17 insertions(+), 15 deletions(-) diff --git a/src/google/appengine/api/taskqueue/cloudtask.py b/src/google/appengine/api/taskqueue/cloudtask.py index aefdc6b..1946577 100644 --- a/src/google/appengine/api/taskqueue/cloudtask.py +++ b/src/google/appengine/api/taskqueue/cloudtask.py @@ -26,6 +26,7 @@ from google.appengine.api import app_identity from google.appengine.api.taskqueue import taskqueue from google.appengine.api.taskqueue import taskqueue_service_bytes_pb2 as taskqueue_service_pb2 +from google.cloud import tasks_v2 from google.cloud import tasks_v2beta3 from google.protobuf import duration_pb2 from google.protobuf import field_mask_pb2 @@ -71,7 +72,7 @@ def create_tasks_in_cloud_tasks(queue_name, tasks, multiple): def delete_tasks_in_cloud_tasks(queue_name, tasks, multiple): """Deletes tasks from a queue using Cloud Tasks Client SDK (supporting BatchDeleteTasks).""" - client = tasks_v2beta3.CloudTasksClient() + client = tasks_v2.CloudTasksClient() project = _get_project_id() region = _get_region() @@ -131,7 +132,7 @@ def delete_tasks_in_cloud_tasks(queue_name, tasks, multiple): def purge_queue_in_cloud_tasks(queue_name): """Purges all tasks in a queue using Cloud Tasks API.""" - client = tasks_v2beta3.CloudTasksClient() + client = tasks_v2.CloudTasksClient() project = _get_project_id() region = _get_region() @@ -147,7 +148,8 @@ def purge_queue_in_cloud_tasks(queue_name): def fetch_queue_stats_in_cloud_tasks(queues, multiple): - """Fetches queue statistics for given queues using Cloud Tasks API.""" + """Fetches queue statistics for given queues using Cloud Tasks API (v2beta3).""" + # QueueStats is retained on v2beta3 as it is out of scope for v2 GA client = tasks_v2beta3.CloudTasksClient() project = _get_project_id() region = _get_region() @@ -306,16 +308,16 @@ def _build_ct_task_payload(queue_name, task, client, project, region): else: body = task.payload - http_method = tasks_v2beta3.HttpMethod.POST + http_method = tasks_v2.HttpMethod.POST if task.method: method_map = { - 'POST': tasks_v2beta3.HttpMethod.POST, - 'GET': tasks_v2beta3.HttpMethod.GET, - 'PUT': tasks_v2beta3.HttpMethod.PUT, - 'DELETE': tasks_v2beta3.HttpMethod.DELETE, - 'HEAD': tasks_v2beta3.HttpMethod.HEAD, + 'POST': tasks_v2.HttpMethod.POST, + 'GET': tasks_v2.HttpMethod.GET, + 'PUT': tasks_v2.HttpMethod.PUT, + 'DELETE': tasks_v2.HttpMethod.DELETE, + 'HEAD': tasks_v2.HttpMethod.HEAD, } - http_method = method_map.get(task.method, tasks_v2beta3.HttpMethod.POST) + http_method = method_map.get(task.method, tasks_v2.HttpMethod.POST) app_engine_http_request = { 'http_method': http_method, @@ -379,7 +381,7 @@ def _build_ct_task_payload(queue_name, task, client, project, region): def _create_single_task_in_cloud_tasks(queue_name, task, multiple): """Helper to create a single task using CloudTasksClient CreateTask API.""" - client = tasks_v2beta3.CloudTasksClient() + client = tasks_v2.CloudTasksClient() project = _get_project_id() region = _get_region() @@ -410,7 +412,7 @@ def _create_single_task_in_cloud_tasks(queue_name, task, multiple): def _create_batch_tasks_in_cloud_tasks(queue_name, tasks, multiple): """Helper to create tasks in batches using CloudTasksClient BatchCreateTasks API.""" - client = tasks_v2beta3.CloudTasksClient() + client = tasks_v2.CloudTasksClient() project = _get_project_id() region = _get_region() diff --git a/src/google/appengine/api/taskqueue/cloudtask_transactional.py b/src/google/appengine/api/taskqueue/cloudtask_transactional.py index 13ea430..2737c9a 100644 --- a/src/google/appengine/api/taskqueue/cloudtask_transactional.py +++ b/src/google/appengine/api/taskqueue/cloudtask_transactional.py @@ -26,7 +26,7 @@ from google.appengine.api import datastore from google.appengine.api.taskqueue import cloudtask from google.appengine.api.taskqueue import taskqueue -from google.cloud import tasks_v2beta3 +from google.cloud import tasks_v2 from google.protobuf.timestamp_pb2 import Timestamp try: @@ -117,7 +117,7 @@ def add_transactional_tasks(queue_name, tasks, multiple): def build_task_payload_for_transactional_task(queue_name, task): """Builds the Cloud Tasks task payload for a transactional task.""" - client = tasks_v2beta3.CloudTasksClient() + client = tasks_v2.CloudTasksClient() project = cloudtask._get_project_id() region = cloudtask._get_region() @@ -126,7 +126,7 @@ def build_task_payload_for_transactional_task(queue_name, task): def dispatch_task_payload(queue_name, task_payload): """Dispatches a pre-built task payload immediately using CloudTasksClient.""" - client = tasks_v2beta3.CloudTasksClient() + client = tasks_v2.CloudTasksClient() project = cloudtask._get_project_id() region = cloudtask._get_region()