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 f7c0b968a32 Fix WASB log source URLs to use storage account endpoint
(#68510)
f7c0b968a32 is described below
commit f7c0b968a326724a9168b762b97f251f10f3e13e
Author: Bafff <[email protected]>
AuthorDate: Thu Jul 16 21:36:06 2026 +0200
Fix WASB log source URLs to use storage account endpoint (#68510)
* Fix WASB log source URLs
* Extract legacy URL helper in WASB log handler per review
---
providers/microsoft/azure/docs/logging/index.rst | 2 +-
.../microsoft/azure/log/wasb_task_handler.py | 33 ++++++++-
.../microsoft/azure/log/test_wasb_task_handler.py | 79 +++++++++++++++++++++-
3 files changed, 108 insertions(+), 6 deletions(-)
diff --git a/providers/microsoft/azure/docs/logging/index.rst
b/providers/microsoft/azure/docs/logging/index.rst
index 9bb1848d964..89ea1ac3e36 100644
--- a/providers/microsoft/azure/docs/logging/index.rst
+++ b/providers/microsoft/azure/docs/logging/index.rst
@@ -67,7 +67,7 @@ Setup Steps:
.. code-block:: none
*** Found remote logs:
- *** *
https://my-container.blob.core.windows.net/path/to/logs/dag_id=tutorial_dag/run_id=manual.../task_id=load/attempt=1.log
+ *** *
https://my-account.blob.core.windows.net/my-container/path/to/logs/dag_id=tutorial_dag/run_id=manual.../task_id=load/attempt=1.log
[2023-07-23, 03:52:47] {taskinstance.py:1144} INFO - Dependencies all
met...
**Note** that the path to the remote log file is listed in the second line.
diff --git
a/providers/microsoft/azure/src/airflow/providers/microsoft/azure/log/wasb_task_handler.py
b/providers/microsoft/azure/src/airflow/providers/microsoft/azure/log/wasb_task_handler.py
index 71704d82588..32cd3ee374e 100644
---
a/providers/microsoft/azure/src/airflow/providers/microsoft/azure/log/wasb_task_handler.py
+++
b/providers/microsoft/azure/src/airflow/providers/microsoft/azure/log/wasb_task_handler.py
@@ -22,6 +22,7 @@ import shutil
from functools import cached_property
from pathlib import Path
from typing import TYPE_CHECKING
+from urllib.parse import urlsplit, urlunsplit
import attrs
from azure.core.exceptions import HttpResponseError
@@ -84,6 +85,34 @@ class WasbRemoteLogIO(LoggingMixin): # noqa: D101
)
return None
+ def _build_legacy_url(self, blob_name: str) -> str:
+ return
f"https://{self.wasb_container}.blob.core.windows.net/{blob_name}"
+
+ def _build_log_source_url(self, blob_name: str) -> str:
+ if not self.hook:
+ return self._build_legacy_url(blob_name)
+
+ blob_service_client = self.hook.blob_service_client
+ account_name = getattr(blob_service_client, "account_name", None)
+ endpoint = getattr(blob_service_client, "primary_endpoint", None)
+ if isinstance(endpoint, str) and endpoint:
+ parsed = urlsplit(endpoint)
+ if parsed.scheme and parsed.netloc:
+ endpoint_path = parsed.path.rstrip("/")
+ # Azurite keeps the account name in the path; WASB SAS-token
auth can put a token there.
+ if endpoint_path:
+ first_path_segment = endpoint_path.strip("/").split("/",
maxsplit=1)[0]
+ if not (isinstance(account_name, str) and
first_path_segment == account_name):
+ endpoint_path = ""
+ account_url = urlunsplit((parsed.scheme, parsed.netloc,
endpoint_path, "", ""))
+ return f"{account_url}/{self.wasb_container}/{blob_name}"
+
+ # Best-effort fallback for clients without primary_endpoint; custom
endpoints use the branch above.
+ if isinstance(account_name, str) and account_name:
+ return
f"https://{account_name}.blob.core.windows.net/{self.wasb_container}/{blob_name}"
+
+ return self._build_legacy_url(blob_name)
+
def read(self, relative_path, ti: RuntimeTI) -> tuple[LogSourceInfo,
LogMessages | None]:
messages = []
logs = []
@@ -100,9 +129,9 @@ class WasbRemoteLogIO(LoggingMixin): # noqa: D101
self.log.exception("can't list blobs")
if blob_names:
- uris = [f"https://{self.wasb_container}.blob.core.windows.net/{b}"
for b in blob_names]
+ uris = (self._build_log_source_url(b) for b in blob_names)
if AIRFLOW_V_3_0_PLUS:
- messages = uris
+ messages = list(uris)
else:
messages.extend(["Found remote logs:", *[f" * {x}" for x in
sorted(uris)]])
else:
diff --git
a/providers/microsoft/azure/tests/unit/microsoft/azure/log/test_wasb_task_handler.py
b/providers/microsoft/azure/tests/unit/microsoft/azure/log/test_wasb_task_handler.py
index 732c9ea49da..57d21664c13 100644
---
a/providers/microsoft/azure/tests/unit/microsoft/azure/log/test_wasb_task_handler.py
+++
b/providers/microsoft/azure/tests/unit/microsoft/azure/log/test_wasb_task_handler.py
@@ -109,6 +109,8 @@ class TestWasbTaskHandler:
@mock.patch("airflow.providers.microsoft.azure.hooks.wasb.WasbHook")
def test_wasb_read(self, mock_hook_cls, ti):
mock_hook = mock_hook_cls.return_value
+ mock_hook.blob_service_client.primary_endpoint =
"https://storage-account.blob.core.windows.net/"
+ mock_hook.blob_service_client.account_name = "storage-account"
mock_hook.get_blobs_list.return_value = ["abc/hello.log"]
mock_hook.read_file.return_value = "Log line"
assert self.wasb_task_handler.io.wasb_read(self.remote_log_location)
== "Log line"
@@ -120,21 +122,26 @@ class TestWasbTaskHandler:
if AIRFLOW_V_3_2_2_PLUS:
logs = list(logs)
assert logs[0].event == "::group::Log message source details"
- assert logs[1].event ==
"https://wasb-container.blob.core.windows.net/abc/hello.log"
+ assert (
+ logs[1].event ==
"https://storage-account.blob.core.windows.net/wasb-container/abc/hello.log"
+ )
assert logs[2].event == "::endgroup::"
assert logs[3].event == "Log line"
assert metadata == {"end_of_log": True, "log_pos": 1}
elif AIRFLOW_V_3_0_PLUS:
logs = list(logs)
assert logs[0].event == "::group::Log message source details"
- assert logs[0].sources ==
["https://wasb-container.blob.core.windows.net/abc/hello.log"]
+ assert logs[0].sources == [
+
"https://storage-account.blob.core.windows.net/wasb-container/abc/hello.log"
+ ]
assert logs[1].event == "::endgroup::"
assert logs[2].event == "Log line"
assert metadata == {"end_of_log": True, "log_pos": 1}
else:
assert logs[0][0][0] == "localhost"
assert (
- "*** Found remote logs:\n*** *
https://wasb-container.blob.core.windows.net/abc/hello.log\n"
+ "*** Found remote logs:\n"
+ "*** *
https://storage-account.blob.core.windows.net/wasb-container/abc/hello.log\n"
in logs[0][0][1]
)
assert "Log line" in logs[0][0][1]
@@ -143,6 +150,72 @@ class TestWasbTaskHandler:
"log_pos": 8,
}
+ def test_log_source_url_keeps_endpoint_path_and_removes_query_string(self):
+ mock_hook = mock.MagicMock()
+ mock_hook.blob_service_client.primary_endpoint =
"http://127.0.0.1:10000/devstoreaccount1?sastoken"
+ mock_hook.blob_service_client.account_name = "devstoreaccount1"
+
+ with mock.patch.object(WasbRemoteLogIO, "hook", new=mock_hook):
+ assert (
+
self.wasb_task_handler.io._build_log_source_url("abc/hello.log")
+ ==
"http://127.0.0.1:10000/devstoreaccount1/wasb-container/abc/hello.log"
+ )
+
+ def
test_log_source_url_removes_query_and_fragment_from_primary_endpoint(self):
+ mock_hook = mock.MagicMock()
+ mock_hook.blob_service_client.primary_endpoint = (
+
"https://storage-account.blob.core.windows.net/?sv=2020&sig=secret#fragment"
+ )
+ mock_hook.blob_service_client.account_name = "storage-account"
+
+ with mock.patch.object(WasbRemoteLogIO, "hook", new=mock_hook):
+ assert (
+
self.wasb_task_handler.io._build_log_source_url("abc/hello.log")
+ ==
"https://storage-account.blob.core.windows.net/wasb-container/abc/hello.log"
+ )
+
+ def test_log_source_url_removes_sas_token_from_endpoint_path(self):
+ mock_hook = mock.MagicMock()
+ mock_hook.blob_service_client.primary_endpoint = (
+ "https://storage-account.blob.core.windows.net/SAStoken/"
+ )
+ mock_hook.blob_service_client.account_name = "storage-account"
+
+ with mock.patch.object(WasbRemoteLogIO, "hook", new=mock_hook):
+ assert (
+
self.wasb_task_handler.io._build_log_source_url("abc/hello.log")
+ ==
"https://storage-account.blob.core.windows.net/wasb-container/abc/hello.log"
+ )
+
+ def
test_log_source_url_uses_account_name_when_primary_endpoint_is_unavailable(self):
+ mock_hook = mock.MagicMock()
+ mock_hook.blob_service_client.primary_endpoint = None
+ mock_hook.blob_service_client.account_name = "storage-account"
+
+ with mock.patch.object(WasbRemoteLogIO, "hook", new=mock_hook):
+ assert (
+
self.wasb_task_handler.io._build_log_source_url("abc/hello.log")
+ ==
"https://storage-account.blob.core.windows.net/wasb-container/abc/hello.log"
+ )
+
+ def
test_log_source_url_uses_legacy_url_when_endpoint_and_account_name_are_unavailable(self):
+ mock_hook = mock.MagicMock()
+ mock_hook.blob_service_client.primary_endpoint = None
+ mock_hook.blob_service_client.account_name = None
+
+ with mock.patch.object(WasbRemoteLogIO, "hook", new=mock_hook):
+ assert (
+
self.wasb_task_handler.io._build_log_source_url("abc/hello.log")
+ == "https://wasb-container.blob.core.windows.net/abc/hello.log"
+ )
+
+ def test_log_source_url_uses_legacy_url_when_hook_is_unavailable(self):
+ with mock.patch.object(WasbRemoteLogIO, "hook", new=None):
+ assert (
+
self.wasb_task_handler.io._build_log_source_url("abc/hello.log")
+ == "https://wasb-container.blob.core.windows.net/abc/hello.log"
+ )
+
@mock.patch(
"airflow.providers.microsoft.azure.hooks.wasb.WasbHook",
**{"return_value.read_file.side_effect": AzureHttpError("failed to
connect", 404)},