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: