pankajastro commented on code in PR #73374:
URL: https://github.com/apache/airflow/pull/73374#discussion_r4147140694
##########
providers/common/sql/src/airflow/providers/common/sql/datafusion/engine.py:
##########
@@ -204,6 +228,86 @@ 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 shared_access_key := _get_wasb_extra_field(extra_dejson,
"shared_access_key"):
+ # Checked ahead of sas_token to match WasbHook.get_conn's
precedence.
+ credentials["access_key"] = shared_access_key
+ 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, "account_key")
+ if access_key:
+ credentials["access_key"] = access_key
+ explicit_credential = True
+
+ if explicit_credential:
+ # object_store's from_env() precedence, high to low:
bearer token > access
+ # key > workload identity > client secret > SAS. Anything
above the
+ # connection's own tier would silently win; SAS is the
bottom tier, so an
+ # env-derived SAS key never can and isn't checked.
+ conflicting_env_vars = [
Review Comment:
Implemented per-tier: access_key connections only check for a bearer token;
service-principal also checks the three access-key spellings and a lone
`AZURE_FEDERATED_TOKEN_FILE`; SAS keeps the full set. One addition beyond what
you outlined — gated SAS's workload-identity check behind the full
`client_id`+`tenant_id`+`federated_token_file` triple too, since a SAS
connection occupies none of those fields either, so a lone
`federated_token_file` there is the same false positive as the shared-key case.
Confirmed empirically it never completes the triple alone. Restored the "can't
skip `from_env()`" comment and fixed the wording to say `build()`'s chain.
---
Drafted-by: Claude Code (Sonnet 5); reviewed by @pankajastro before posting
--
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]