diff --git a/task-sdk/src/airflow/sdk/definitions/connection.py b/task-sdk/src/airflow/sdk/definitions/connection.py index 06a95a3b868b8..5838c66f3fd14 100644 --- a/task-sdk/src/airflow/sdk/definitions/connection.py +++ b/task-sdk/src/airflow/sdk/definitions/connection.py @@ -153,6 +153,12 @@ def __init__(self, *, conn_id: str, uri: str | None = None, **kwargs) -> None: else: self.__dict__.update(attrs.asdict(self.from_uri(uri, conn_id=conn_id), recurse=False)) + if self.password: + from airflow.sdk.log import mask_secret + + mask_secret(self.password) + mask_secret(quote(self.password)) + def get_uri(self) -> str: """Generate and return connection in URI format.""" from urllib.parse import parse_qsl diff --git a/task-sdk/tests/task_sdk/definitions/test_connection.py b/task-sdk/tests/task_sdk/definitions/test_connection.py index c973fc58c91a4..1c960e4e98e55 100644 --- a/task-sdk/tests/task_sdk/definitions/test_connection.py +++ b/task-sdk/tests/task_sdk/definitions/test_connection.py @@ -234,6 +234,21 @@ def test_extra_dejson_property(self): connection.extra = '{"auth": {"type": "oauth"}, "headers": {"User-Agent": "Airflow"}}' assert connection.extra_dejson == {"auth": {"type": "oauth"}, "headers": {"User-Agent": "Airflow"}} + def test_password_is_masked(self): + """Test that the password is masked, whether set via kwargs or a URI.""" + from airflow.sdk._shared.secrets_masker import _secrets_masker + + masker = _secrets_masker() + masker.reset_masker() + + Connection(conn_id="test_conn", conn_type="azure", login="client_id", password="sp-secret") + assert masker.redact("sp-secret") == "***" + + masker.reset_masker() + + Connection.from_uri("azure://client_id:sp-secret@?tenantId=tenant", conn_id="test_conn") + assert masker.redact("sp-secret") == "***" + class TestConnectionsFromSecrets: def test_get_connection_secrets_backend(self, mock_supervisor_comms, tmp_path):