pankajastro commented on code in PR #73374:
URL: https://github.com/apache/airflow/pull/73374#discussion_r4136426489
##########
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 the precise breakdown — added `AZURE_STORAGE_MASTER_KEY` and the
client-secret triple (`AZURE_CLIENT_ID`/`AZURE_STORAGE_CLIENT_ID`,
`AZURE_CLIENT_SECRET`/`AZURE_STORAGE_CLIENT_SECRET`,
`AZURE_TENANT_ID`/`AZURE_STORAGE_TENANT_ID`/`AZURE_STORAGE_AUTHORITY_ID`/`AZURE_AUTHORITY_ID`,
only flagged when all three are set), and dropped `AZURE_STORAGE_SAS_KEY`.
Pinned both in tests, plus a case confirming a partial triple doesn't
false-positive.
---
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]