diff --git a/airflow/providers/google/cloud/operators/bigquery_dts.py b/airflow/providers/google/cloud/operators/bigquery_dts.py index ee1a3b548bcc7..d786c903e77bd 100644 --- a/airflow/providers/google/cloud/operators/bigquery_dts.py +++ b/airflow/providers/google/cloud/operators/bigquery_dts.py @@ -138,6 +138,9 @@ def execute(self, context: Context): result = TransferConfig.to_dict(response) self.log.info("Created DTS transfer config %s", get_object_id(result)) self.xcom_push(context, key="transfer_config_id", value=get_object_id(result)) + # don't push AWS secret in XCOM + result.get("params", {}).pop("secret_access_key", None) + result.get("params", {}).pop("access_key_id", None) return result diff --git a/tests/providers/google/cloud/operators/test_bigquery_dts.py b/tests/providers/google/cloud/operators/test_bigquery_dts.py index 78c92d52ed5bc..aa52169a77060 100644 --- a/tests/providers/google/cloud/operators/test_bigquery_dts.py +++ b/tests/providers/google/cloud/operators/test_bigquery_dts.py @@ -46,12 +46,15 @@ TRANSFER_CONFIG_NAME = "projects/123abc/locations/321cba/transferConfig/1a2b3c" RUN_NAME = "projects/123abc/locations/321cba/transferConfig/1a2b3c/runs/123" +transfer_config = TransferConfig( + name=TRANSFER_CONFIG_NAME, params={"secret_access_key": "AIRFLOW_KEY", "access_key_id": "AIRFLOW_KEY_ID"} +) class BigQueryCreateDataTransferOperatorTestCase(unittest.TestCase): @mock.patch( "airflow.providers.google.cloud.operators.bigquery_dts.BiqQueryDataTransferServiceHook", - **{"return_value.create_transfer_config.return_value": TransferConfig(name=TRANSFER_CONFIG_NAME)}, + **{"return_value.create_transfer_config.return_value": transfer_config}, ) def test_execute(self, mock_hook): op = BigQueryCreateDataTransferOperator( @@ -59,7 +62,7 @@ def test_execute(self, mock_hook): ) ti = mock.MagicMock() - op.execute({"ti": ti}) + return_value = op.execute({"ti": ti}) mock_hook.return_value.create_transfer_config.assert_called_once_with( authorization_code=None, @@ -71,6 +74,9 @@ def test_execute(self, mock_hook): ) ti.xcom_push.assert_called_with(execution_date=None, key="transfer_config_id", value="1a2b3c") + assert "secret_access_key" not in return_value.get("params", {}) + assert "access_key_id" not in return_value.get("params", {}) + class BigQueryDeleteDataTransferConfigOperatorTestCase(unittest.TestCase): @mock.patch("airflow.providers.google.cloud.operators.bigquery_dts.BiqQueryDataTransferServiceHook")