Skip to content

feat(taskqueue): route PushQueues to Cloud Tasks v2 (GA) - #161

Open
riddhi-shivhare wants to merge 4 commits into
mainfrom
taskqueue-v2-ga
Open

riddhi-shivhare wants to merge 4 commits into
mainfrom
taskqueue-v2-ga

Conversation

@riddhi-shivhare

Copy link
Copy Markdown

When APPENGINE_USE_CLOUDTASK_PUSH_QUEUE is enabled, push queue operations in the TaskQueue SDK are served by the Cloud Tasks v2 API:

  • Queue.add (single and batch) uses CreateTask and BatchCreateTasks.
  • Queue.delete_tasks uses BatchDeleteTasks.
  • Queue.purge uses PurgeQueue.
  • Queue.fetch_statistics uses v2beta3 GetQueue, since queue stats are not part of the v2 API.
  • Transactional tasks are staged in Datastore and dispatched after commit.

Batch operations wait on the long-running operation and map failed_requests back to per-task TaskQueue errors. Tasks that are not accounted for when the operation fails surface the operation error. Cloud Tasks clients are created once per process and reused.

Adds google-cloud-tasks>=2.25.0 as a dependency.

Testing:

  • tests/google/appengine/api/taskqueue/cloudtask_test.py: 23/23 pass.
  • tests/google/appengine/api/taskqueue/taskqueue_test.py: 284/284 pass.

@riddhi-shivhare riddhi-shivhare changed the title feat(taskqueue): route Taskqueue PushQueues to Cloud Tasks v2 (GA) feat(taskqueue): route PushQueues to Cloud Tasks v2 (GA) Oct 7, 2026
@riddhi-shivhare
riddhi-shivhare force-pushed the taskqueue-v2-ga branch 11 times, most recently from 5fc2d2a to ac97d29 Compare October 8, 2026 07:15
When APPENGINE_USE_CLOUDTASK_PUSH_QUEUE is enabled, push queue operations
in the TaskQueue SDK are served by the Cloud Tasks v2 API:

- Queue.add (single and batch) uses CreateTask and BatchCreateTasks.
- Queue.delete_tasks uses BatchDeleteTasks.
- Queue.purge uses PurgeQueue.
- Queue.fetch_statistics uses v2beta3 GetQueue, since queue stats are not
  part of the v2 API.
- Transactional tasks are staged in Datastore and dispatched after commit.
  Each staged entity is locked in a transaction before dispatch, so the
  post-commit path and the /_ah/cloudtask/sweep cron handler do not
  dispatch the same entity concurrently. Each sweep reads at most 500
  entities of each status. Failed entities are deleted after 7 days.
  Like the legacy API, at most 5 tasks can be added in one transaction.
  Each staged task is its own entity group, so a transaction that adds
  transactional tasks must be cross-group (xg=True) unless it touches no
  other entity group. Staged entities use the kind, property names and
  payload form of the Go and Java SDKs, a JSON object {"task": T} where T
  is the proto3 JSON form of the v2 Task, so a sweeper in a service
  written in any of them can dispatch them. Entities staged by the
  preview release are still dispatched.

Batch operations read the synchronous operation result and map
failed_requests in the operation metadata back to per-task TaskQueue
errors. Reserved headers (Host, X-Google-*, X-AppEngine-*) are not
forwarded, and a target without a version routes to that service's
default version. Cloud Tasks API errors are raised as the matching
taskqueue.Error subclass (for example UnknownQueueError,
PermissionDeniedError, InvalidQueueModeError, TransientError). A deadline
set on the rpc passed to the *_async methods, or passed to
fetch_statistics, limits the Cloud Tasks RPC calls, and the rpc reports
the Cloud Tasks result. Without a deadline, batch RPC calls use a 20
second default timeout. Adding, leasing and deleting pull tasks keeps
using the legacy service; purge and stats also work for pull queues
through Cloud Tasks. Cloud Tasks clients are created once per process and
reused.

Adds google-cloud-tasks>=2.25.0 as a dependency on Python 3.10+, the
versions that release supports. On Python 3.9 the legacy TaskQueue path
is unchanged, and enabling the Cloud Tasks backend raises an ImportError
that explains the requirement. Unit tests are in
tests/google/appengine/api/taskqueue/cloudtask_test.py.
if (not isinstance(op, futures.Future)
and getattr(op, 'response', None) is not None):
return op.response
return op.result()

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Because BatchCreateTasks and BatchDeleteTasks are synchronous RPCs that return a pre-completed google.longrunning.Operation (done=True on the initial response), google.api_core.operation.Operation.init (operation.py:78) has already unpacked self._operation.response (or set self._operation.error) in memory before returning. However, Operation subclasses google.api_core.polling.PollingFuture: if op._operation.done were ever False (due to a backend regression or malformed response), bare op.result() initiates a background Operations.GetOperation polling loop for up to 900 seconds, ignoring the caller's deadline. Pass op.result(timeout=0) (or check op.operation.done first) so the SDK enforces synchronous extraction and never issues GetOperation polling RPCs.

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Fixed. _extract_lro_response now checks op.operation.done up front and calls op.result(timeout=0) so it never triggers GetOperation polling.

request={'parent': parent, 'names': task_names},
**_timeout_kwargs(deadline_at, _DEFAULT_RPC_TIMEOUT_SECONDS)
)
_extract_lro_response(op)

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

In delete_tasks_in_cloud_tasks, _extract_lro_response(op) is invoked before reading op.metadata.failed_requests (cloudtask.py:234-242), whereas _create_batch_tasks_in_cloud_tasks (cloudtask.py:794-804) reads op.metadata.failed_requests first. Because op.result() raises immediately if op.operation.HasField('error') is set on the synchronous Operation, calling _extract_lro_response(op) first can raise an operation-level error before per-task failed_requests (such as marking NOT_FOUND tasks with _Task__deleted = False) are processed. Move _extract_lro_response(op) after metadata extraction to match _create_batch_tasks_in_cloud_tasks.

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Done. delete_tasks_in_cloud_tasks now extracts op.metadata.failed_requests first and passes has_failed_requests=bool(failed_requests) to _extract_lro_response(op), matching _create_batch_tasks_in_cloud_tasks so per-task failures like NOT_FOUND are processed even when the operation sets error.

queue_name, task
)

entity = datastore.Entity(_PENDING_TASK_KIND)

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Creating _AE_PendingCloudTask without namespace='' binds the staged entity to the caller's active namespace_manager namespace. However, /_ah/cloudtask/sweep runs via App Engine Cron in the default ('') namespace, where datastore.Query(_PENDING_TASK_KIND, {'status =': status}) (cloudtask_transactional.py:220) scans only ''. If a transaction commits inside a non-default namespace and the post-commit fast path fails, that task becomes invisible to sweep() forever—never retried and never cleaned up after _FAILED_TASK_RETENTION.

Stage _AE_PendingCloudTask explicitly in the empty namespace (datastore.Entity(_PENDING_TASK_KIND, namespace='')) and pass namespace='' to datastore.Query(_PENDING_TASK_KIND, ..., namespace='') in sweep(). Cross-group (xg=True) transactions in Datastore permit root entities in namespace='' alongside tenant entities within the same application.

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Done.

except google_exceptions.Aborted:
# Aborted subclasses Conflict (HTTP 409) but does not mean the task exists.
raise
except (google_exceptions.AlreadyExists, google_exceptions.Conflict) as e:

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Catching AlreadyExists and Conflict inside _create_single_task_in_cloud_tasks (and at line 845 in _create_batch_tasks_in_cloud_tasks) converts them into taskqueue.TaskAlreadyExistsError before @_translate_errors runs _to_taskqueue_error. Because _to_taskqueue_error returns existing taskqueue.Error instances unchanged (cloudtask.py:97), an HTTP 409 carrying ExecutorServiceError::TOMBSTONED_TASK in its error detail is misclassified as TaskAlreadyExistsError instead of TombstonedTaskError. Remove the inner exception translation blocks in both helpers so @_translate_errors handles all ExecutorServiceError::* mappings consistently.

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Fixed. Removed the inner try ... except blocks.

@gajjarc gajjarc left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

LGTM

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants