kaxil commented on code in PR #73374:
URL: https://github.com/apache/airflow/pull/73374#discussion_r4073152606


##########
providers/common/sql/src/airflow/providers/common/sql/datafusion/engine.py:
##########
@@ -196,6 +197,45 @@ def _fetch_extra_configs(keys: list[str]) -> dict[str, 
Any]:
                     key_path = os.environ.get("GOOGLE_APPLICATION_CREDENTIALS")
                 credentials = self._remove_none_values({"key_path": key_path, 
"keyfile_dict": keyfile_dict})
 
+            case "wasb":
+                extra_dejson = conn.extra_dejson
+                for unsupported_field in (
+                    "connection_string",
+                    "managed_identity_client_id",
+                    "workload_identity_tenant_id",
+                ):
+                    if extra_dejson.get(unsupported_field):
+                        raise ValueError(
+                            f"Connection field {unsupported_field!r} is not 
supported for DataFusion "
+                            "Azure Blob Storage access; only 
tenant_id+login+password (service "
+                            "principal), sas_token, 
shared_access_key/account_key/password, or ambient "
+                            "credentials (AZURE_* environment variables, 
managed identity, workload "
+                            "identity, or az login) are used."
+                        )
+                credentials = {"account": 
self._resolve_wasb_account(conn.host, conn.login)}
+                tenant_id = extra_dejson.get("tenant_id")
+                sas_token = extra_dejson.get("sas_token")
+                if tenant_id and conn.login and conn.password:
+                    # client_id/client_secret/tenant_id must all be set 
together, or not at all --
+                    # DataFusion's binding panics on a partial combination.
+                    credentials.update(

Review Comment:
   The binding does not let this dict be the whole story. 
`MicrosoftAzure.__init__` starts from `MicrosoftAzureBuilder::from_env()` 
(`store.rs:93`) and only overlays the keys set here, and object_store's 
`build()` checks `access_key`, then the workload-identity triple 
(`client_id`+`tenant_id`+`federated_token_file`), before the client secret and 
`sas_query_pairs` (`azure/builder.rs:984-1015`). So anything in the worker's 
`AZURE_*` environment that lands in one of those earlier slots wins over the 
connection. I checked this against a local endpoint: with 
`AZURE_STORAGE_ACCOUNT_KEY` set, a connection carrying a SAS token or a full 
service principal authenticates with SharedKey and the SAS never reaches the 
wire. With the AKS workload-identity variables set (`AZURE_CLIENT_ID`, 
`AZURE_TENANT_ID`, `AZURE_FEDERATED_TOKEN_FILE`), a connection SAS is ignored, 
and a connection service principal produces a token request with the 
connection's `client_id` and the pod's federated token, `client_s
 ecret` never sent.
   
   That is the "silently falling back to a different identity" case the guard 
above rejects for `managed_identity_client_id`, except here it happens in the 
most common Azure deployment. Since the binding has no way to skip 
`from_env()`, could this raise when the connection supplies an explicit 
credential and one of `AZURE_FEDERATED_TOKEN_FILE`, 
`AZURE_STORAGE_ACCOUNT_KEY`, `AZURE_STORAGE_ACCESS_KEY`, 
`AZURE_STORAGE_SAS_KEY`, `AZURE_STORAGE_TOKEN` is set? At minimum the docs 
paragraph needs to say the environment takes precedence, not that it is last.



##########
providers/common/sql/src/airflow/providers/common/sql/datafusion/engine.py:
##########
@@ -196,6 +197,45 @@ def _fetch_extra_configs(keys: list[str]) -> dict[str, 
Any]:
                     key_path = os.environ.get("GOOGLE_APPLICATION_CREDENTIALS")
                 credentials = self._remove_none_values({"key_path": key_path, 
"keyfile_dict": keyfile_dict})
 
+            case "wasb":
+                extra_dejson = conn.extra_dejson
+                for unsupported_field in (
+                    "connection_string",
+                    "managed_identity_client_id",
+                    "workload_identity_tenant_id",
+                ):
+                    if extra_dejson.get(unsupported_field):
+                        raise ValueError(
+                            f"Connection field {unsupported_field!r} is not 
supported for DataFusion "
+                            "Azure Blob Storage access; only 
tenant_id+login+password (service "
+                            "principal), sas_token, 
shared_access_key/account_key/password, or ambient "
+                            "credentials (AZURE_* environment variables, 
managed identity, workload "
+                            "identity, or az login) are used."
+                        )
+                credentials = {"account": 
self._resolve_wasb_account(conn.host, conn.login)}
+                tenant_id = extra_dejson.get("tenant_id")
+                sas_token = extra_dejson.get("sas_token")
+                if tenant_id and conn.login and conn.password:

Review Comment:
   I see the rationale for not forwarding a partial triple (the binding panics 
on it), but falling through changes identity rather than reporting the problem. 
With `tenant_id` and `login` set and `password` empty, this drops to the `else` 
branch and ends up on ambient IMDS silently. With `tenant_id` and `password` 
set and `login` empty, the client secret is forwarded as `access_key`, and the 
user gets `InvalidAccessKey` from the panic for a service principal they never 
asked to be a shared key. WasbHook raises via `ClientSecretCredential` for the 
same connection. Could this be `if tenant_id:` followed by a `ValueError` 
naming the missing `login` or `password`, with 
`test_get_credentials_azure_tenant_id_without_login_falls_back` flipped to 
`pytest.raises`? As a data point, deleting `and conn.password` from this line 
leaves all twelve Azure tests green.



##########
providers/common/sql/docs/operators.rst:
##########
@@ -362,6 +363,28 @@ resolved in this order:
     :start-after: [START howto_analytics_operator_with_gcs]
     :end-before: [END howto_analytics_operator_with_gcs]
 
+Azure Storage
+-------------
+Use an ``az://`` URI with a ``conn_id`` pointing to a ``wasb`` connection.
+``abfs://`` and ``abfss://`` URIs are not recognized yet. Credentials are
+resolved in this order:
+
+1. Azure AD service principal -- ``tenant_id`` extra, with ``login`` as the
+   client ID and ``password`` as the client secret
+2. SAS token -- ``sas_token`` extra, as a query string
+3. Shared key -- ``password``, or the ``shared_access_key``/``account_key`` 
extra
+4. Ambient credentials -- ``AZURE_*`` environment variables, managed identity,

Review Comment:
   See the note on `engine.py`: the binding reads `AZURE_*` first and 
object_store puts an environment access key or workload identity ahead of a 
connection SAS token or client secret, so this item is not fourth in practice. 
Also, the Azure CLI is only consulted when `AZURE_USE_AZURE_CLI=true` 
(`builder.rs` gates `AzureCliCredential` on `use_azure_cli`); the default 
ambient path is IMDS managed identity, so a developer relying on `az login` the 
way WasbHook allows will get a metadata-endpoint failure.



##########
providers/common/sql/src/airflow/providers/common/sql/datafusion/engine.py:
##########
@@ -196,6 +197,45 @@ def _fetch_extra_configs(keys: list[str]) -> dict[str, 
Any]:
                     key_path = os.environ.get("GOOGLE_APPLICATION_CREDENTIALS")
                 credentials = self._remove_none_values({"key_path": key_path, 
"keyfile_dict": keyfile_dict})
 
+            case "wasb":
+                extra_dejson = conn.extra_dejson
+                for unsupported_field in (
+                    "connection_string",
+                    "managed_identity_client_id",
+                    "workload_identity_tenant_id",
+                ):
+                    if extra_dejson.get(unsupported_field):

Review Comment:
   WasbHook's `_get_field` still reads the legacy `extra__wasb__<name>` 
spelling, so a connection created on an older Airflow with 
`extra__wasb__connection_string` passes this guard, and 
`extra__wasb__sas_token` is never seen by the lookup below (I checked, it lands 
on ambient auth). The GCS branch has the same gap, so this is not new to this 
PR. A small helper that checks the bare key then the prefixed one, used by both 
branches, would close it.



##########
providers/common/sql/tests/unit/common/sql/datafusion/test_engine.py:
##########
@@ -372,6 +373,145 @@ def 
test_get_credentials_gcs_rejects_unsupported_identity_fields(self, unsupport
         with pytest.raises(ValueError, match=f"{unsupported_field!r} is not 
supported"):
             engine._get_credentials(mock_conn)
 
+    def test_get_credentials_azure_with_shared_key(self):
+        mock_conn = MagicMock()
+        mock_conn.conn_type = "wasb"
+        mock_conn.host = None
+        mock_conn.login = "myaccount"
+        mock_conn.password = "mykey"
+        mock_conn.extra_dejson = {}
+        engine = DataFusionEngine()
+
+        credentials, extra_config = engine._get_credentials(mock_conn)
+
+        assert credentials == {"account": "myaccount", "access_key": "mykey"}
+        assert extra_config == {}
+
+    def test_get_credentials_azure_with_shared_access_key_extra(self):
+        mock_conn = MagicMock()
+        mock_conn.conn_type = "wasb"
+        mock_conn.host = None
+        mock_conn.login = "myaccount"
+        mock_conn.password = None
+        mock_conn.extra_dejson = {"shared_access_key": "extra-key"}
+        engine = DataFusionEngine()
+
+        credentials, extra_config = engine._get_credentials(mock_conn)
+
+        assert credentials == {"account": "myaccount", "access_key": 
"extra-key"}
+        assert extra_config == {}
+
+    def test_get_credentials_azure_with_service_principal(self):
+        mock_conn = MagicMock()
+        mock_conn.conn_type = "wasb"
+        mock_conn.host = None
+        mock_conn.login = "client-id"
+        mock_conn.password = "client-secret"
+        mock_conn.extra_dejson = {"tenant_id": "tenant-id"}
+        engine = DataFusionEngine()
+
+        credentials, extra_config = engine._get_credentials(mock_conn)
+
+        assert credentials == {
+            "account": "client-id",
+            "client_id": "client-id",
+            "client_secret": "client-secret",
+            "tenant_id": "tenant-id",
+        }
+        assert extra_config == {}
+
+    def 
test_get_credentials_azure_with_service_principal_and_host_prefers_host_account(self):
+        mock_conn = MagicMock()
+        mock_conn.conn_type = "wasb"
+        mock_conn.host = "realaccount.blob.core.windows.net"
+        mock_conn.login = "11111111-2222-3333-4444-555555555555"
+        mock_conn.password = "client-secret"
+        mock_conn.extra_dejson = {"tenant_id": "tenant-id"}
+        engine = DataFusionEngine()
+
+        credentials, extra_config = engine._get_credentials(mock_conn)
+
+        assert credentials == {
+            "account": "realaccount",
+            "client_id": "11111111-2222-3333-4444-555555555555",
+            "client_secret": "client-secret",
+            "tenant_id": "tenant-id",
+        }
+        assert extra_config == {}
+
+    def test_get_credentials_azure_tenant_id_without_login_falls_back(self):
+        """A partial service-principal config (tenant_id alone) must not be 
forwarded --
+        DataFusion's binding panics on a partial 
client_id/client_secret/tenant_id combination."""
+        mock_conn = MagicMock()
+        mock_conn.conn_type = "wasb"
+        mock_conn.host = None
+        mock_conn.login = None
+        mock_conn.password = None
+        mock_conn.extra_dejson = {"tenant_id": "tenant-id"}
+        engine = DataFusionEngine()
+
+        credentials, extra_config = engine._get_credentials(mock_conn)
+
+        assert "tenant_id" not in credentials

Review Comment:
   This is the one test that exercises `host=None, login=None`, and it passes 
while `credentials == {"account": "None"}`. An equality assertion on the whole 
dict, like the sibling tests use, would force the account question into the 
open. While here, `_resolve_wasb_account` has no coverage for a full `https://` 
URL with a path, a sovereign-cloud suffix, a host without a dot (which falls 
through to `login`), or the 24-character cap. Deleting the no-dot branch or 
changing `[:24]` to `[:23]` leaves every Azure test green, as does swapping the 
SAS and shared-key branches. One parametrized test calling the helper directly, 
plus one connection with both `sas_token` and `shared_access_key` set, would 
pin all of that.



##########
providers/common/sql/docs/operators.rst:
##########
@@ -362,6 +363,28 @@ resolved in this order:
     :start-after: [START howto_analytics_operator_with_gcs]
     :end-before: [END howto_analytics_operator_with_gcs]
 
+Azure Storage
+-------------
+Use an ``az://`` URI with a ``conn_id`` pointing to a ``wasb`` connection.
+``abfs://`` and ``abfss://`` URIs are not recognized yet. Credentials are
+resolved in this order:
+
+1. Azure AD service principal -- ``tenant_id`` extra, with ``login`` as the
+   client ID and ``password`` as the client secret
+2. SAS token -- ``sas_token`` extra, as a query string
+3. Shared key -- ``password``, or the ``shared_access_key``/``account_key`` 
extra
+4. Ambient credentials -- ``AZURE_*`` environment variables, managed identity,
+   workload identity, or the Azure CLI
+
+``connection_string``, ``managed_identity_client_id``, 
``workload_identity_tenant_id``,
+and a URL-form ``sas_token`` are not supported.

Review Comment:
   Two things a reader needs that this list does not say: where the account 
name comes from (the first label of `host`, else `login`), which is the one 
thing a user has to get right after the host-first change, and that 
`client_secret_auth_config` (the `authority` override WasbHook honours) is 
ignored here.



##########
providers/common/sql/src/airflow/providers/common/sql/datafusion/engine.py:
##########
@@ -205,6 +245,28 @@ def _remove_none_values(params: dict[str, Any]) -> 
dict[str, Any]:
         """Filter out None values from the dictionary."""
         return {k: v for k, v in params.items() if v is not None}
 
+    @staticmethod
+    def _resolve_wasb_account(host: str | None, login: str | None) -> str:
+        """
+        Return the storage account name the way WasbHook resolves it.
+
+        From ``host`` when set (its netloc's first label), falling back to 
``login`` only when
+        ``host`` is empty -- login holds the service-principal client_id in 
that auth mode, not
+        the account name. Reimplemented locally rather than importing
+        ``airflow.providers.microsoft.azure.utils.parse_blob_account_url``, to 
avoid pulling the
+        microsoft-azure provider's full Azure SDK dependency stack into 
common-sql for one string
+        operation that only needs the stdlib.
+        """
+        netloc = urlsplit(host if host else 
f"https://{login}.blob.core.windows.net/";).netloc
+        if not netloc:
+            # No scheme was given (e.g. a bare DNS name); urlsplit put it all 
in the path instead.
+            netloc = urlsplit(f"https://{host}";).netloc
+        if "." not in netloc:
+            # Only an Active Directory ID was given, not a full URL or DNS 
name.
+            netloc = f"{login}.blob.core.windows.net"
+        # Azure storage account names are capped at 24 characters.
+        return netloc.split(".", 1)[0][:24]

Review Comment:
   Only the first label survives here, and `MicrosoftAzure` has no `endpoint` 
parameter, so `build()` always targets `https://<label>.blob.core.windows.net` 
(`builder.rs:940-947`). A `wasb` connection whose `host` is 
`myacct.blob.core.chinacloudapi.cn`, `myacct.blob.core.usgovcloudapi.net`, a 
`privatelink` FQDN, or the Azurite emulator's loopback-IP URL 
(`http://<loopback ip>:10000/devstoreaccount1`) works in WasbHook (which keeps 
the full `account_url`) but is silently redirected to the public cloud here. I 
measured the Azurite form collapsing to the first octet of the IP as the 
account name, and the binding accepts a dotted account string without 
complaint, so nothing downstream catches it. With a SAS token the query string 
is then sent to that public-cloud host.
   
   Given the PR's own rule of raising instead of silently falling back, could 
this raise a `ValueError` when the netloc contains a dot and does not end in 
`.blob.core.windows.net`? `AZURE_STORAGE_ENDPOINT` is read by `from_env()` and 
honoured by `build()`, so the message can point users there, and the docs can 
add it next to the `abfs` caveat.



##########
providers/common/sql/src/airflow/providers/common/sql/datafusion/engine.py:
##########
@@ -196,6 +197,45 @@ def _fetch_extra_configs(keys: list[str]) -> dict[str, 
Any]:
                     key_path = os.environ.get("GOOGLE_APPLICATION_CREDENTIALS")
                 credentials = self._remove_none_values({"key_path": key_path, 
"keyfile_dict": keyfile_dict})
 
+            case "wasb":
+                extra_dejson = conn.extra_dejson
+                for unsupported_field in (
+                    "connection_string",
+                    "managed_identity_client_id",
+                    "workload_identity_tenant_id",
+                ):
+                    if extra_dejson.get(unsupported_field):
+                        raise ValueError(
+                            f"Connection field {unsupported_field!r} is not 
supported for DataFusion "
+                            "Azure Blob Storage access; only 
tenant_id+login+password (service "
+                            "principal), sas_token, 
shared_access_key/account_key/password, or ambient "
+                            "credentials (AZURE_* environment variables, 
managed identity, workload "
+                            "identity, or az login) are used."
+                        )
+                credentials = {"account": 
self._resolve_wasb_account(conn.host, conn.login)}

Review Comment:
   The host-first fix is right for the service-principal case, but it changed 
the empty case. With neither `host` nor `login` set, `_resolve_wasb_account` 
returns the string `'None'` (the f-string on line 260 interpolates 
`login=None`), so the store gets `account="None"`. Before this change `account` 
was `None`, `_remove_none_values` dropped it, and the binding fell back to 
`AZURE_STORAGE_ACCOUNT_NAME`. I confirmed on datafusion 51.0.0 that omitting 
`account` with that variable set builds fine, and that `account="None"` is 
accepted silently. The `wasb_default` connection that `airflow db` creates has 
no host or login, and it is the one the new example DAG references. Could the 
helper return `None` when both are empty (or raise), and could the fallback 
test assert the full dict rather than two key absences?



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to