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

shahar1 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 e4736b46968 Pass ClientTimeout to aiohttp in Cloud Composer async hook 
(#74164)
e4736b46968 is described below

commit e4736b4696827d851a7163dcd610d4fd83b5cbb0
Author: Shahar Epstein <[email protected]>
AuthorDate: Sat Oct 3 19:53:35 2026 +0300

    Pass ClientTimeout to aiohttp in Cloud Composer async hook (#74164)
    
    aiohttp types the keyword arguments of ClientSession.request only when
    type-checking on Python 3.11 or newer, where timeout must be a
    ClientTimeout. On 3.10 the arguments are untyped, so passing the hook's
    float timeout went unnoticed until CI moved to Python 3.11. aiohttp
    already converts a bare number to ClientTimeout(total=...) at runtime,
    so building it explicitly keeps the behaviour and satisfies mypy.
---
 .../providers/google/cloud/hooks/cloud_composer.py |  4 ++--
 .../unit/google/cloud/hooks/test_cloud_composer.py | 24 ++++++++++++++++++++++
 2 files changed, 26 insertions(+), 2 deletions(-)

diff --git 
a/providers/google/src/airflow/providers/google/cloud/hooks/cloud_composer.py 
b/providers/google/src/airflow/providers/google/cloud/hooks/cloud_composer.py
index ce570ea9840..d1df8660bad 100644
--- 
a/providers/google/src/airflow/providers/google/cloud/hooks/cloud_composer.py
+++ 
b/providers/google/src/airflow/providers/google/cloud/hooks/cloud_composer.py
@@ -24,7 +24,7 @@ from collections.abc import MutableSequence, Sequence
 from typing import TYPE_CHECKING, Any
 from urllib.parse import urlencode, urljoin
 
-from aiohttp import ClientSession
+from aiohttp import ClientSession, ClientTimeout
 from google.api_core.gapic_v1.method import DEFAULT, _MethodDefault
 from google.auth.transport.requests import AuthorizedSession, Request
 from google.cloud.orchestration.airflow.service_v1 import (
@@ -610,7 +610,7 @@ class CloudComposerAsyncHook(GoogleBaseAsyncHook):
                     "Content-Type": "application/json",
                     "Authorization": f"Bearer {self._credentials.token}",
                 },
-                timeout=timeout,
+                timeout=ClientTimeout(total=timeout),
             ) as response:
                 return await response.json(), response.status
 
diff --git 
a/providers/google/tests/unit/google/cloud/hooks/test_cloud_composer.py 
b/providers/google/tests/unit/google/cloud/hooks/test_cloud_composer.py
index 95fdc78dbc0..dfa5215b1cc 100644
--- a/providers/google/tests/unit/google/cloud/hooks/test_cloud_composer.py
+++ b/providers/google/tests/unit/google/cloud/hooks/test_cloud_composer.py
@@ -22,6 +22,7 @@ from unittest import mock
 from unittest.mock import AsyncMock
 
 import pytest
+from aiohttp import ClientSession, ClientTimeout
 from google.api_core.gapic_v1.method import DEFAULT
 from google.cloud.orchestration.airflow.service_v1 import 
EnvironmentsAsyncClient
 
@@ -344,6 +345,29 @@ class TestCloudComposerAsyncHook:
         with mock.patch(BASE_STRING.format("GoogleBaseAsyncHook.__init__"), 
new=mock_init):
             self.hook = CloudComposerAsyncHook(gcp_conn_id="test")
 
+    @pytest.mark.asyncio
+    @pytest.mark.parametrize("timeout", [30.0, None])
+    @mock.patch(COMPOSER_STRING.format("ClientSession"), spec=True)
+    async def test_make_composer_airflow_api_request_passes_client_timeout(
+        self, mock_client_session, timeout
+    ) -> None:
+        self.hook._credentials = mock.MagicMock(valid=True, token="test-token")
+        session = mock.MagicMock(spec=ClientSession)
+        mock_client_session.return_value.__aenter__.return_value = session
+        response = session.request.return_value.__aenter__.return_value
+        response.json = AsyncMock(return_value={"key": "value"})
+        response.status = 200
+
+        result = await self.hook.make_composer_airflow_api_request(
+            method="GET",
+            airflow_uri=TEST_COMPOSER_AIRFLOW_URI,
+            path="/api/v1/dags",
+            timeout=timeout,
+        )
+
+        assert result == ({"key": "value"}, 200)
+        assert session.request.call_args.kwargs["timeout"] == 
ClientTimeout(total=timeout)
+
     @pytest.mark.asyncio
     
@mock.patch(COMPOSER_STRING.format("CloudComposerAsyncHook.get_environment_client"))
     async def test_create_environment(self, mock_client) -> None:

Reply via email to