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]