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):