This is an automated email from the ASF dual-hosted git repository.

potiuk pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/airflow.git


The following commit(s) were added to refs/heads/main by this push:
     new d2883eead5c Fix Cloud SQL 409 operationInProgress on import/export 
operations (#68361)
d2883eead5c is described below

commit d2883eead5c09610901a077aa9c9d127b87adb4c
Author: Radhwène Dhouafli <[email protected]>
AuthorDate: Sun Oct 4 02:43:38 2026 +0200

    Fix Cloud SQL 409 operationInProgress on import/export operations (#68361)
    
    * Fix Cloud SQL 409 operationInProgress on import/export operations
    
    Apply the existing operation_in_progress_retry() policy to
    CloudSQLHook.import_instance and export_instance, the only two admin
    methods that lacked it. Also re-raise operationInProgress HttpError
    un-wrapped from import_instance so the retry decorator can see it;
    terminal HttpErrors still get the friendly AirflowException message.
    
    * Limit Cloud SQL 409 retry to the import submit call
    
    Address review feedback on #68361: operation_in_progress_retry()
    previously wrapped the whole import_instance, including the
    operation-status polling. A retryable 429 raised during polling
    re-ran the method and re-submitted an import that was already
    accepted, importing the same data twice.
    
    The submit now lives in _submit_import, which alone carries the
    retry decorator; import_instance waits outside the retry scope, so
    a polling failure fails the task instead of re-submitting. Adds a
    regression test asserting exactly one submit when polling raises a
    retryable error.
    
    * Decode Cloud SQL API error bodies in hook failure messages
    
    * Raise a dedicated error when Cloud SQL import fails
    
    Callers could only tell an import failure apart from any other hook failure 
by
    matching on the message, since every failure surfaced as the broad
    AirflowException. A dedicated subclass lets them catch this case precisely, 
and
    keeps existing except AirflowException handlers working.
    
    * Keep the Cloud SQL API error when its body is not valid UTF-8
    
    A strict decode raised UnicodeDecodeError while building the 
CloudSQLImportError message, replacing the real API error with an unrelated one.
    
    Generated-by: Claude Opus 5
    
    ---------
    
    Co-authored-by: Jarek Potiuk <[email protected]>
---
 generated/known_airflow_exceptions.txt             |  2 +-
 .../providers/google/cloud/hooks/cloud_sql.py      | 57 +++++++++++++---
 .../unit/google/cloud/hooks/test_cloud_sql.py      | 78 +++++++++++++++++++++-
 3 files changed, 124 insertions(+), 13 deletions(-)

diff --git a/generated/known_airflow_exceptions.txt 
b/generated/known_airflow_exceptions.txt
index ca2b46365cb..8fdef403b6b 100644
--- a/generated/known_airflow_exceptions.txt
+++ b/generated/known_airflow_exceptions.txt
@@ -209,7 +209,7 @@ 
providers/google/src/airflow/providers/google/cloud/hooks/cloud_build.py::3
 providers/google/src/airflow/providers/google/cloud/hooks/cloud_composer.py::5
 
providers/google/src/airflow/providers/google/cloud/hooks/cloud_memorystore.py::5
 providers/google/src/airflow/providers/google/cloud/hooks/cloud_run.py::1
-providers/google/src/airflow/providers/google/cloud/hooks/cloud_sql.py::32
+providers/google/src/airflow/providers/google/cloud/hooks/cloud_sql.py::31
 
providers/google/src/airflow/providers/google/cloud/hooks/cloud_storage_transfer_service.py::5
 providers/google/src/airflow/providers/google/cloud/hooks/compute.py::6
 providers/google/src/airflow/providers/google/cloud/hooks/compute_ssh.py::6
diff --git 
a/providers/google/src/airflow/providers/google/cloud/hooks/cloud_sql.py 
b/providers/google/src/airflow/providers/google/cloud/hooks/cloud_sql.py
index 0f672570f33..4030bd9b63e 100644
--- a/providers/google/src/airflow/providers/google/cloud/hooks/cloud_sql.py
+++ b/providers/google/src/airflow/providers/google/cloud/hooks/cloud_sql.py
@@ -65,6 +65,7 @@ from airflow.providers.google.common.hooks.base_google import 
(
     GoogleBaseAsyncHook,
     GoogleBaseHook,
     get_field,
+    is_operation_in_progress_exception,
 )
 from airflow.utils.log.logging_mixin import LoggingMixin
 
@@ -82,6 +83,10 @@ TIME_TO_SLEEP_IN_SECONDS = 20
 CLOUD_SQL_PROXY_VERSION_REGEX = re.compile(r"^v?(\d+\.\d+\.\d+)(-\w*.?\d?)?$")
 
 
+class CloudSQLImportError(AirflowException):
+    """Raised when importing data into a Cloud SQL instance fails."""
+
+
 class CloudSqlOperationStatus:
     """Helper class with operation statuses."""
 
@@ -332,10 +337,14 @@ class CloudSQLHook(GoogleBaseHook):
         self._wait_for_operation_to_complete(project_id=project_id, 
operation_name=operation_name)
 
     @GoogleBaseHook.fallback_to_default_project_id
+    @GoogleBaseHook.operation_in_progress_retry()
     def export_instance(self, instance: str, body: dict, project_id: str):
         """
         Export data from a Cloud SQL instance to a Cloud Storage bucket as a 
SQL dump or CSV file.
 
+        Cloud SQL runs one admin operation at a time per instance, so this 
submit can come back
+        409 — hence the retry decorator.
+
         :param instance: Database instance ID of the Cloud SQL instance. This 
does not include the
             project ID.
         :param body: The request body, as described in
@@ -353,11 +362,45 @@ class CloudSQLHook(GoogleBaseHook):
         operation_name = response["name"]
         return operation_name
 
+    @GoogleBaseHook.operation_in_progress_retry()
+    def _submit_import(self, instance: str, body: dict, project_id: str) -> 
str:
+        """
+        Submit an import request for a Cloud SQL instance, retrying while an 
operation is in progress.
+
+        Cloud SQL runs one admin operation at a time per instance, so this 
submit can come back 409.
+        The retry covers the submit only: repeating a rejected submit is safe, 
repeating an accepted
+        one is not.
+
+        :param instance: Database instance ID. This does not include the 
project ID.
+        :param body: The request body, as described in
+            
https://cloud.google.com/sql/docs/mysql/admin-api/v1beta4/instances/import#request-body
+        :param project_id: Project ID of the project that contains the 
instance.
+        :return: The name of the accepted import operation.
+        """
+        try:
+            response = (
+                self.get_conn()
+                .instances()
+                .import_(project=project_id, instance=instance, body=body)
+                .execute(num_retries=self.num_retries)
+            )
+            return response["name"]
+        except HttpError as ex:
+            # The decorator retries HttpError, not CloudSQLImportError, so 
don't wrap the 409.
+            if is_operation_in_progress_exception(ex):
+                raise
+            raise CloudSQLImportError(
+                f"Importing instance {instance} failed: 
{ex.content.decode('utf-8', errors='replace')}"
+            )
+
     @GoogleBaseHook.fallback_to_default_project_id
     def import_instance(self, instance: str, body: dict, project_id: str) -> 
None:
         """
         Import data into a Cloud SQL instance from a SQL dump or CSV file in 
Cloud Storage.
 
+        The submit retries on 409 (see ``_submit_import``). The polling below 
stays outside that
+        retry: re-submitting an accepted import would load the same data twice.
+
         :param instance: Database instance ID. This does not include the
             project ID.
         :param body: The request body, as described in
@@ -366,17 +409,13 @@ class CloudSQLHook(GoogleBaseHook):
             to None or missing, the default project_id from the Google Cloud 
connection is used.
         :return: None
         """
+        operation_name = self._submit_import(instance=instance, body=body, 
project_id=project_id)
         try:
-            response = (
-                self.get_conn()
-                .instances()
-                .import_(project=project_id, instance=instance, body=body)
-                .execute(num_retries=self.num_retries)
-            )
-            operation_name = response["name"]
             self._wait_for_operation_to_complete(project_id=project_id, 
operation_name=operation_name)
         except HttpError as ex:
-            raise AirflowException(f"Importing instance {instance} failed: 
{ex.content}")
+            raise CloudSQLImportError(
+                f"Importing instance {instance} failed: 
{ex.content.decode('utf-8', errors='replace')}"
+            )
 
     @GoogleBaseHook.fallback_to_default_project_id
     def clone_instance(self, instance: str, body: dict, project_id: str) -> 
None:
@@ -403,7 +442,7 @@ class CloudSQLHook(GoogleBaseHook):
             operation_name = response["name"]
             self._wait_for_operation_to_complete(project_id=project_id, 
operation_name=operation_name)
         except HttpError as ex:
-            raise AirflowException(f"Cloning of instance {instance} failed: 
{ex.content}")
+            raise AirflowException(f"Cloning of instance {instance} failed: 
{ex.content.decode('utf-8')}")
 
     @GoogleBaseHook.fallback_to_default_project_id
     def create_ssl_certificate(self, instance: str, body: dict, project_id: 
str):
diff --git a/providers/google/tests/unit/google/cloud/hooks/test_cloud_sql.py 
b/providers/google/tests/unit/google/cloud/hooks/test_cloud_sql.py
index 955add00379..07e2f9be6e8 100644
--- a/providers/google/tests/unit/google/cloud/hooks/test_cloud_sql.py
+++ b/providers/google/tests/unit/google/cloud/hooks/test_cloud_sql.py
@@ -45,6 +45,7 @@ from airflow.providers.google.cloud.hooks.cloud_sql import (
     CloudSQLAsyncHook,
     CloudSQLDatabaseHook,
     CloudSQLHook,
+    CloudSQLImportError,
     CloudSqlProxyRunner,
 )
 
@@ -96,7 +97,7 @@ class TestGcpSqlHookDefaultProjectId:
         self.cloudsql_hook.get_conn = mock.Mock(
             side_effect=HttpError(resp=httplib2.Response({"status": 400}), 
content=b"Error content")
         )
-        with pytest.raises(AirflowException) as ctx:
+        with pytest.raises(CloudSQLImportError) as ctx:
             self.cloudsql_hook.import_instance(instance="instance", body={})
         err = ctx.value
         assert "Importing instance " in str(err)
@@ -165,10 +166,81 @@ class TestGcpSqlHookDefaultProjectId:
             ),
             {"name": "operation_id"},
         ]
-        with pytest.raises(HttpError):
-            self.cloudsql_hook.export_instance(project_id="example-project", 
instance="instance", body={})
+        # First submit returns 429 (one of the two operation-in-progress codes 
recognised by
+        # ``is_operation_in_progress_exception``); 
``operation_in_progress_retry`` retries and the
+        # second submit succeeds, returning the operation name. The import 
test below covers 409.
+        result = self.cloudsql_hook.export_instance(
+            project_id="example-project", instance="instance", body={}
+        )
+        assert result == "operation_id"
+        assert export_method.call_count == 2
+        assert execute_method.call_count == 2
         wait_for_operation_to_complete.assert_not_called()
 
+    
@mock.patch("airflow.providers.google.cloud.hooks.cloud_sql.CloudSQLHook.get_conn")
+    
@mock.patch("airflow.providers.google.cloud.hooks.cloud_sql.CloudSQLHook._wait_for_operation_to_complete")
+    def test_instance_import_with_in_progress_retry(self, 
wait_for_operation_to_complete, get_conn):
+        import_method = get_conn.return_value.instances.return_value.import_
+        execute_method = import_method.return_value.execute
+        execute_method.side_effect = [
+            HttpError(
+                resp=httplib2.Response({"status": 409}),
+                content=b"operationInProgress",
+            ),
+            {"name": "operation_id"},
+        ]
+        wait_for_operation_to_complete.return_value = None
+        # First submit returns 409 ``operationInProgress``. ``_submit_import`` 
re-raises it past its
+        # friendly-message wrapper (instead of converting it to 
CloudSQLImportError), so
+        # ``operation_in_progress_retry`` sees the raw HttpError, retries, and 
the second submit
+        # succeeds; the resulting operation is awaited exactly once.
+        self.cloudsql_hook.import_instance(project_id="example-project", 
instance="instance", body={})
+        assert import_method.call_count == 2
+        assert execute_method.call_count == 2
+        wait_for_operation_to_complete.assert_called_once_with(
+            project_id="example-project", operation_name="operation_id"
+        )
+
+    
@mock.patch("airflow.providers.google.cloud.hooks.cloud_sql.CloudSQLHook.get_conn")
+    
@mock.patch("airflow.providers.google.cloud.hooks.cloud_sql.CloudSQLHook._wait_for_operation_to_complete")
+    def test_instance_import_does_not_resubmit_when_polling_fails(
+        self, wait_for_operation_to_complete, get_conn
+    ):
+        import_method = get_conn.return_value.instances.return_value.import_
+        execute_method = import_method.return_value.execute
+        execute_method.return_value = {"name": "operation_id"}
+        wait_for_operation_to_complete.side_effect = HttpError(
+            resp=httplib2.Response({"status": 429}),
+            content=b"rate limited",
+        )
+        # The submit succeeds, then the operation-status polling raises a 
retryable 429. The retry
+        # scope must not include the polling: re-running ``import_instance`` 
would re-submit an
+        # import that was already accepted and import the same data twice. The 
task must fail
+        # instead, with exactly one submit on record.
+        with pytest.raises(CloudSQLImportError, match="Importing instance 
instance failed"):
+            self.cloudsql_hook.import_instance(project_id="example-project", 
instance="instance", body={})
+        import_method.assert_called_once_with(body={}, instance="instance", 
project="example-project")
+        execute_method.assert_called_once_with(num_retries=5)
+        wait_for_operation_to_complete.assert_called_once_with(
+            project_id="example-project", operation_name="operation_id"
+        )
+
+    @pytest.mark.parametrize("failing_step", ["submit", "polling"])
+    
@mock.patch("airflow.providers.google.cloud.hooks.cloud_sql.CloudSQLHook.get_conn")
+    
@mock.patch("airflow.providers.google.cloud.hooks.cloud_sql.CloudSQLHook._wait_for_operation_to_complete")
+    def test_instance_import_error_with_non_utf8_body(
+        self, wait_for_operation_to_complete, get_conn, failing_step
+    ):
+        error = HttpError(resp=httplib2.Response({"status": 400}), 
content=b"bad \xff body")
+        execute_method = 
get_conn.return_value.instances.return_value.import_.return_value.execute
+        if failing_step == "submit":
+            execute_method.side_effect = error
+        else:
+            execute_method.return_value = {"name": "operation_id"}
+            wait_for_operation_to_complete.side_effect = error
+        with pytest.raises(CloudSQLImportError, match="Importing instance 
instance failed: bad � body"):
+            self.cloudsql_hook.import_instance(project_id="example-project", 
instance="instance", body={})
+
     @mock.patch(
         
"airflow.providers.google.cloud.hooks.cloud_sql.CloudSQLHook.get_credentials_and_project_id",
         return_value=(mock.MagicMock(), "example-project"),

Reply via email to