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


##########
providers/common/sql/src/airflow/providers/common/sql/datafusion/engine.py:
##########
@@ -204,6 +211,80 @@ def _get_gcp_extra_field(extra_dejson: dict[str, Any], 
field_name: 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 _get_wasb_extra_field(extra_dejson, 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)}
+                explicit_credential = False
+                if tenant_id := _get_wasb_extra_field(extra_dejson, 
"tenant_id"):
+                    if not conn.login or not conn.password:
+                        # Falling through here would silently switch identity 
(ambient auth, or
+                        # the client secret sent as a shared key) instead of 
failing clearly.
+                        missing = "login (client_id)" if not conn.login else 
"password (client_secret)"
+                        raise ValueError(
+                            f"Connection extra 'tenant_id' is set for 
DataFusion Azure Blob Storage "
+                            f"service-principal auth, but {missing} is not."
+                        )
+                    credentials.update(
+                        {"client_id": conn.login, "client_secret": 
conn.password, "tenant_id": tenant_id}
+                    )
+                    explicit_credential = True
+                elif sas_token := _get_wasb_extra_field(extra_dejson, 
"sas_token"):
+                    if sas_token.startswith("http"):
+                        raise ValueError(
+                            "A URL-form `sas_token` is not supported for 
DataFusion Azure Blob Storage "
+                            "access; provide the SAS token as a query string 
instead."
+                        )
+                    credentials["sas_query_pairs"] = 
parse_qsl(sas_token.lstrip("?"))
+                    explicit_credential = True
+                else:
+                    access_key = (
+                        conn.password
+                        or _get_wasb_extra_field(extra_dejson, 
"shared_access_key")
+                        or _get_wasb_extra_field(extra_dejson, "account_key")
+                    )
+                    if access_key:
+                        credentials["access_key"] = access_key
+                        explicit_credential = True
+
+                if explicit_credential:
+                    # The binding always reads these via from_env() first and 
checks that
+                    # env-derived access key / workload-identity ahead of 
what's set here, so
+                    # they'd silently win over the connection's credential. No 
way to skip
+                    # from_env(), so this can only be caught, not fixed, on 
the Python side.
+                    conflicting_env_vars = [

Review Comment:
   Thanks for adding this guard. Two names it needs and one it does not, going 
by object_store's own tables. `from_env()` maps `AZURE_CLIENT_ID`, 
`AZURE_CLIENT_SECRET` and `AZURE_TENANT_ID` (and the `AZURE_STORAGE_CLIENT_*` 
spellings) onto the client-secret slots, and `AZURE_STORAGE_MASTER_KEY` onto 
`access_key` (`azure/builder.rs`, the `AzureConfigKey::from_str` arms). 
`build()` then picks `access_key` and the client-secret triple before 
`sas_query_pairs`, so either of those in the worker environment wins over a 
connection SAS token, which is exactly the silent identity switch this list is 
meant to catch. I checked on datafusion 51.0.0: with `AZURE_STORAGE_MASTER_KEY` 
set to a non-base64 value and a SAS connection, construction fails with 
`InvalidAccessKey`, so the env key was consulted first; with the 
`AZURE_CLIENT_*` triple set, the store constructs and the SAS never reaches the 
credential chain. The AKS workload-identity webhook sets `AZURE_CLIENT_ID` and 
`AZURE_TENANT_ID`, so th
 e triple case is a real deployment shape, not a corner.
   
   In the other direction, `AZURE_STORAGE_SAS_KEY` sits below `sas_query_pairs` 
in `build()` and below every explicit credential, so it never outranks a 
connection and that entry raises for nothing. (The Fabric token quartet, 
`AZURE_FABRIC_*`, also sits above everything, but it only engages with all four 
set, so I would leave it out.) `AZURE_STORAGE_ACCOUNT_KEY` with an access-key 
connection is also harmless, because `with_access_key` overwrites the env value 
(`store.rs`), though raising there is at worst conservative.
   
   Could the list become credential-aware (which env slots outrank the 
credential the connection actually supplied), or at minimum add the SP triple 
in both spellings plus `AZURE_STORAGE_MASTER_KEY` and drop 
`AZURE_STORAGE_SAS_KEY`? The parametrized test below only walks the five names 
in this tuple, so it cannot see the gap; adding the new names there would pin 
it.



##########
providers/common/sql/docs/operators.rst:
##########
@@ -362,6 +363,49 @@ 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. The account name
+comes from ``host`` (its first DNS label) when set, falling back to
+``login`` only when ``host`` is empty; only the public
+``*.blob.core.windows.net`` cloud is supported, since DataFusion's binding
+has no endpoint override. ``client_secret_auth_config`` (the authority
+override ``WasbHook`` honors) is not read here.
+
+The connection supplies one of the following credentials:
+
+1. Azure AD service principal -- ``tenant_id`` extra, with ``login`` as the
+   client ID and ``password`` as the client secret (both required together)
+2. SAS token -- ``sas_token`` extra, as a query string
+3. Shared key -- ``password``, or the ``shared_access_key``/``account_key`` 
extra
+4. None of the above -- ambient auth (see below)
+
+**A worker environment variable can override the connection.** DataFusion
+reads ``AZURE_*`` environment variables first, and an environment access
+key or workload-identity token wins over the connection's SAS token or
+client secret. If the connection sets an explicit credential (1-3 above)
+and the worker also has ``AZURE_FEDERATED_TOKEN_FILE``,
+``AZURE_STORAGE_ACCOUNT_KEY``, ``AZURE_STORAGE_ACCESS_KEY``,
+``AZURE_STORAGE_SAS_KEY``, or ``AZURE_STORAGE_TOKEN`` set, this raises

Review Comment:
   See the note on the guard in `engine.py`: an 
`AZURE_CLIENT_ID`/`AZURE_CLIENT_SECRET`/`AZURE_TENANT_ID` triple or 
`AZURE_STORAGE_MASTER_KEY` on the worker also wins over a connection SAS token, 
so this sentence promises a raise the code does not yet make in those cases. 
Once the guard list is settled, this paragraph should name the same set.



##########
providers/common/sql/src/airflow/providers/common/sql/datafusion/engine.py:
##########
@@ -204,6 +211,80 @@ def _get_gcp_extra_field(extra_dejson: dict[str, Any], 
field_name: 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 _get_wasb_extra_field(extra_dejson, 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)}
+                explicit_credential = False
+                if tenant_id := _get_wasb_extra_field(extra_dejson, 
"tenant_id"):
+                    if not conn.login or not conn.password:
+                        # Falling through here would silently switch identity 
(ambient auth, or
+                        # the client secret sent as a shared key) instead of 
failing clearly.
+                        missing = "login (client_id)" if not conn.login else 
"password (client_secret)"
+                        raise ValueError(
+                            f"Connection extra 'tenant_id' is set for 
DataFusion Azure Blob Storage "
+                            f"service-principal auth, but {missing} is not."
+                        )
+                    credentials.update(
+                        {"client_id": conn.login, "client_secret": 
conn.password, "tenant_id": tenant_id}
+                    )
+                    explicit_credential = True
+                elif sas_token := _get_wasb_extra_field(extra_dejson, 
"sas_token"):

Review Comment:
   One more ordering question, prompted by comparing this branch against 
`WasbHook.get_conn`. The hook checks `shared_access_key` before `sas_token` 
(`hooks/wasb.py`, the `shared_access_key` lookup precedes the `sas_token` one), 
and only then falls to `password` and `account_key`. Here SAS wins over every 
key. A connection that carries a current `shared_access_key` alongside a SAS 
token that has since expired works through the hook and fails here, and the new 
precedence test pins that difference in. In round 2 I only asked for the order 
to be tested and documented, so this is a question rather than a request: is 
there a reason to prefer SAS first, or would matching the hook's order (service 
principal, shared access key, SAS, password, account key) be safer for 
connections people already have? If SAS first is deliberate, one sentence in 
the docs paragraph saying it differs from `WasbHook` would save someone a 
confusing 403.



##########
providers/common/sql/src/airflow/providers/common/sql/datafusion/engine.py:
##########
@@ -213,6 +294,45 @@ 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}
 
+    _AZURE_PUBLIC_SUFFIX = ".blob.core.windows.net"
+
+    @classmethod
+    def _resolve_wasb_account(cls, host: str | None, login: str | None) -> str 
| None:
+        """
+        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. Returns ``None`` when neither is set, so the binding 
falls back to
+        ``AZURE_STORAGE_ACCOUNT_NAME`` instead of targeting the literal string 
``"None"``.
+        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.
+
+        Only the public ``*.blob.core.windows.net`` cloud is supported: 
DataFusion's Azure binding
+        takes no endpoint override, so a sovereign-cloud or emulator host 
would otherwise be
+        silently misrouted to the public account of the same name.
+        """
+        if not host and not login:
+            return None
+        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"
+        if not netloc.endswith(cls._AZURE_PUBLIC_SUFFIX):
+            raise ValueError(
+                f"Connection host {host!r} does not resolve to the public 
{cls._AZURE_PUBLIC_SUFFIX} "
+                "cloud, which is the only one DataFusion's Azure Blob Storage 
binding can target (it "
+                "has no endpoint override). Sovereign clouds and the Azurite 
emulator are not "
+                "supported; set the AZURE_STORAGE_ENDPOINT environment 
variable instead."

Review Comment:
   This raise fires before the binding is built and does not look at 
`AZURE_STORAGE_ENDPOINT`, so a user who follows the advice and sets it still 
gets the same error from the same connection. Either skip the suffix check when 
that variable is set (and let `from_env()` pick up the endpoint), or say here 
that the account must then move to `login` or `AZURE_STORAGE_ACCOUNT_NAME` with 
`host` cleared. Not blocking, the message just needs to describe a path that 
works.



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