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

dabla 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 847183dc9f9 Fix HttpAsyncHook crash on stream/cert/trust_env 
connection extras (#71524)
847183dc9f9 is described below

commit 847183dc9f97f1ebcf9016dc1a350c0bd8084843
Author: Gengtao Xu <[email protected]>
AuthorDate: Fri Sep 25 05:04:41 2026 -0400

    Fix HttpAsyncHook crash on stream/cert/trust_env connection extras (#71524)
    
    HttpAsyncHook shares its extra-options parsing with the synchronous
    HttpHook, but the shared helper's output is requests-flavored and gets
    forwarded as-is into aiohttp's request call. aiohttp rejects stream,
    cert, and trust_env as unexpected keyword arguments, so any HTTP
    connection carrying one of those keys in its Extra field crashed every
    async request. check_response was accepted into the same dict but
    never actually consumed on the async side, so it silently had no
    effect unlike the synchronous path's run_and_check().
---
 .../http/src/airflow/providers/http/hooks/http.py  | 17 ++++++-
 providers/http/tests/unit/http/hooks/test_http.py  | 56 +++++++++++++++++++++-
 2 files changed, 70 insertions(+), 3 deletions(-)

diff --git a/providers/http/src/airflow/providers/http/hooks/http.py 
b/providers/http/src/airflow/providers/http/hooks/http.py
index e4f5de03ed8..9b8e874bcaf 100644
--- a/providers/http/src/airflow/providers/http/hooks/http.py
+++ b/providers/http/src/airflow/providers/http/hooks/http.py
@@ -48,6 +48,10 @@ if TYPE_CHECKING:
 
     from airflow.models import Connection
 
+# requests-only extra options with no aiohttp equivalent; aiohttp's request 
methods raise
+# TypeError on unexpected kwargs, so HttpAsyncHook must strip these before the 
request call.
+_AIOHTTP_UNSUPPORTED_EXTRA_OPTIONS = {"stream", "cert", "trust_env"}
+
 
 def _url_from_endpoint(base_url: str | None, endpoint: str | None) -> str:
     """Combine base url with endpoint."""
@@ -639,6 +643,16 @@ class AsyncHttpSession(LoggingMixin):
     ) -> ClientResponse:
         from tenacity import AsyncRetrying, stop_after_attempt, wait_fixed
 
+        check_response = extra_options.pop("check_response", True)
+        unsupported_options = _AIOHTTP_UNSUPPORTED_EXTRA_OPTIONS & 
extra_options.keys()
+        if unsupported_options:
+            self.log.warning(
+                "Ignoring connection extra option(s) %s: not supported by 
HttpAsyncHook.",
+                ", ".join(sorted(unsupported_options)),
+            )
+            for option in unsupported_options:
+                extra_options.pop(option)
+
         async def request_func() -> ClientResponse:
             response = await self._request(
                 url,
@@ -649,7 +663,8 @@ class AsyncHttpSession(LoggingMixin):
                 auth=self.auth,
                 **extra_options,
             )
-            response.raise_for_status()
+            if check_response:
+                response.raise_for_status()
             return response
 
         async for attempt in AsyncRetrying(
diff --git a/providers/http/tests/unit/http/hooks/test_http.py 
b/providers/http/tests/unit/http/hooks/test_http.py
index db7d65ccf2d..c983d68229e 100644
--- a/providers/http/tests/unit/http/hooks/test_http.py
+++ b/providers/http/tests/unit/http/hooks/test_http.py
@@ -1110,7 +1110,6 @@ class TestHttpAsyncHook:
                 "verify": False,
                 "allow_redirects": False,
                 "max_redirects": 3,
-                "trust_env": False,
             }
         ],
         indirect=True,
@@ -1142,7 +1141,60 @@ class TestHttpAsyncHook:
                 assert mocked_function.call_args.kwargs.get("verify_ssl") is 
False
                 assert mocked_function.call_args.kwargs.get("allow_redirects") 
is False
                 assert mocked_function.call_args.kwargs.get("max_redirects") 
== 3
-                assert mocked_function.call_args.kwargs.get("trust_env") is 
False
+
+    @pytest.mark.asyncio
+    @pytest.mark.parametrize(
+        "setup_connections_with_extras",
+        [{"stream": True, "cert": "client.pem", "trust_env": True}],
+        indirect=True,
+    )
+    async def test_async_request_ignores_unsupported_extra_options(
+        self, setup_connections_with_extras, caplog
+    ):
+        """stream/cert/trust_env have no aiohttp equivalent and must not reach 
the request call."""
+        hook = HttpAsyncHook(http_conn_id="http_conn_with_extras")
+
+        with mock.patch("aiohttp.ClientSession.post", 
new_callable=mock.AsyncMock) as mocked_function:
+            mocked_function.return_value = MockAiohttpClientResponse(
+                status=200,
+                payload={"status": {"status": 200}},
+                method="POST",
+                url="http://test:8080/v1/test";,
+            )
+            async with aiohttp.ClientSession() as session:
+                await hook.run(session=session, endpoint="v1/test")
+
+        kwargs = mocked_function.call_args.kwargs
+        assert "stream" not in kwargs
+        assert "cert" not in kwargs
+        assert "trust_env" not in kwargs
+        assert (
+            "Ignoring connection extra option(s) cert, stream, trust_env: not 
supported by HttpAsyncHook."
+            in caplog.messages
+        )
+
+    @pytest.mark.asyncio
+    @pytest.mark.parametrize(
+        "setup_connections_with_extras",
+        [{"check_response": False}],
+        indirect=True,
+    )
+    async def 
test_async_request_does_not_raise_for_status_if_check_response_is_false(
+        self, setup_connections_with_extras
+    ):
+        hook = HttpAsyncHook(http_conn_id="http_conn_with_extras", 
method="GET")
+
+        with mock.patch("aiohttp.ClientSession.get", 
new_callable=mock.AsyncMock) as mocked_get:
+            mocked_get.return_value = MockAiohttpClientResponse(
+                status=500,
+                reason="Internal Server Error",
+                method="GET",
+                url="http://test:8080/v1/test";,
+            )
+            async with aiohttp.ClientSession() as session:
+                resp = await hook.run(session=session, endpoint="v1/test")
+
+        assert resp.status == 500
 
     @pytest.mark.asyncio
     async def test_build_request_url_from_connection(self):

Reply via email to