This is an automated email from the ASF dual-hosted git repository.
shahar1 pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/airflow.git
The following commit(s) were added to refs/heads/main by this push:
new 056806e8412 Keep object store credentials out of the task log in
Teradata transfers (#72176)
056806e8412 is described below
commit 056806e84125491af61b65c6fef1c7a3d0985011
Author: Jarek Potiuk <[email protected]>
AuthorDate: Sun Sep 20 11:47:23 2026 +0200
Keep object store credentials out of the task log in Teradata transfers
(#72176)
When no AUTHORIZATION object is configured and the store is private, both
Teradata
transfer operators embed the object store credentials directly in the
CREATE TABLE
statement. DbApiHook logs every statement it runs, so those values were
written to
the task log on each run.
Suppress the hook's own SQL logging for that statement and log a redacted
form in
its place, so the statement is still visible without the credential values.
The
AUTHORIZATION path is unaffected and keeps normal SQL logging.
Teradata also records the statement in its own query logs, which Airflow
cannot
redact; both operator docstrings now say so and point at the AUTHORIZATION
object.
---
.../teradata/transfers/azure_blob_to_teradata.py | 14 ++++
.../providers/teradata/transfers/s3_to_teradata.py | 15 ++++
.../transfers/test_azure_blob_to_teradata.py | 48 +++++++++++++
.../unit/teradata/transfers/test_s3_to_teradata.py | 83 ++++++++++++++++++++++
4 files changed, 160 insertions(+)
diff --git
a/providers/teradata/src/airflow/providers/teradata/transfers/azure_blob_to_teradata.py
b/providers/teradata/src/airflow/providers/teradata/transfers/azure_blob_to_teradata.py
index f6d474513c3..0a4faefcc80 100644
---
a/providers/teradata/src/airflow/providers/teradata/transfers/azure_blob_to_teradata.py
+++
b/providers/teradata/src/airflow/providers/teradata/transfers/azure_blob_to_teradata.py
@@ -60,6 +60,12 @@ class AzureBlobStorageToTeradataOperator(BaseOperator):
Refer to
https://docs.teradata.com/r/Enterprise_IntelliFlex_VMware/Teradata-VantageTM-Native-Object-Store-Getting-Started-Guide-17.20/Setting-Up-Access/Controlling-Foreign-Table-Access-with-an-AUTHORIZATION-Object
+
+ When ``teradata_authorization_name`` is not set and the object store is
private, the store's
+ credentials are sent to Teradata inside the ``CREATE TABLE`` statement
itself. Teradata records
+ that statement in its own query logs (DBQL) and monitoring views, which
Airflow cannot redact,
+ so the credentials are readable by anyone with access to those. Prefer an
AUTHORIZATION object,
+ which keeps the credentials in the database instead of in each statement.
Note that ``blob_source_key`` and ``teradata_table`` are
templated, so you can use variables in them if you wish.
"""
@@ -92,11 +98,13 @@ class AzureBlobStorageToTeradataOperator(BaseOperator):
)
teradata_hook = TeradataHook(teradata_conn_id=self.teradata_conn_id)
credentials_part = "ACCESS_ID= '' ACCESS_KEY= ''"
+ redacted_credentials_part = credentials_part
if not self.public_bucket:
# Accessing data directly from the Azure Blob Storage and creating
permanent table inside the
# database
if self.teradata_authorization_name:
credentials_part =
f"AUTHORIZATION={self.teradata_authorization_name}"
+ redacted_credentials_part = credentials_part
else:
# Obtaining the Azure client ID and Azure secret in order to
access a specified Blob container
azure_hook = WasbHook(wasb_conn_id=self.azure_conn_id)
@@ -104,6 +112,7 @@ class AzureBlobStorageToTeradataOperator(BaseOperator):
access_id = conn.login
access_secret = conn.password
credentials_part = f"ACCESS_ID= '{access_id}' ACCESS_KEY=
'{access_secret}'"
+ redacted_credentials_part = "ACCESS_ID= '***' ACCESS_KEY=
'***'"
sql = dedent(f"""
CREATE MULTISET TABLE {self.teradata_table} AS
(
@@ -113,6 +122,11 @@ class AzureBlobStorageToTeradataOperator(BaseOperator):
) AS d
) WITH DATA
""").rstrip()
+ if redacted_credentials_part != credentials_part:
+ # The statement embeds the credentials themselves, so keep it out
of the task log
+ # and log the redacted form in its place.
+ teradata_hook.log_sql = False
+ self.log.info("Running statement: %s",
sql.replace(credentials_part, redacted_credentials_part))
try:
teradata_hook.run(sql, True)
except Exception as ex:
diff --git
a/providers/teradata/src/airflow/providers/teradata/transfers/s3_to_teradata.py
b/providers/teradata/src/airflow/providers/teradata/transfers/s3_to_teradata.py
index 36bac281fdf..c2b4ae41efe 100644
---
a/providers/teradata/src/airflow/providers/teradata/transfers/s3_to_teradata.py
+++
b/providers/teradata/src/airflow/providers/teradata/transfers/s3_to_teradata.py
@@ -58,6 +58,12 @@ class S3ToTeradataOperator(BaseOperator):
Refer to
https://docs.teradata.com/r/Enterprise_IntelliFlex_VMware/Teradata-VantageTM-Native-Object-Store-Getting-Started-Guide-17.20/Setting-Up-Access/Controlling-Foreign-Table-Access-with-an-AUTHORIZATION-Object
+
+ When ``teradata_authorization_name`` is not set and the object store is
private, the store's
+ credentials are sent to Teradata inside the ``CREATE TABLE`` statement
itself. Teradata records
+ that statement in its own query logs (DBQL) and monitoring views, which
Airflow cannot redact,
+ so the credentials are readable by anyone with access to those. Prefer an
AUTHORIZATION object,
+ which keeps the credentials in the database instead of in each statement.
Note that ``s3_source_key`` and ``teradata_table`` are
templated, so you can use variables in them if you wish.
"""
@@ -92,18 +98,22 @@ class S3ToTeradataOperator(BaseOperator):
s3_hook = S3Hook(aws_conn_id=self.aws_conn_id)
teradata_hook = TeradataHook(teradata_conn_id=self.teradata_conn_id)
credentials_part = "ACCESS_ID= '' ACCESS_KEY= ''"
+ redacted_credentials_part = credentials_part
if not self.public_bucket:
# Accessing data directly from the S3 bucket and creating
permanent table inside the database
if self.teradata_authorization_name:
credentials_part =
f"AUTHORIZATION={self.teradata_authorization_name}"
+ redacted_credentials_part = credentials_part
else:
credentials = s3_hook.get_credentials()
access_key = credentials.access_key
access_secret = credentials.secret_key
credentials_part = f"ACCESS_ID= '{access_key}' ACCESS_KEY=
'{access_secret}'"
+ redacted_credentials_part = "ACCESS_ID= '***' ACCESS_KEY=
'***'"
token = credentials.token
if token:
credentials_part = credentials_part + f" SESSION_TOKEN =
'{token}'"
+ redacted_credentials_part = redacted_credentials_part + "
SESSION_TOKEN = '***'"
sql = dedent(f"""
CREATE MULTISET TABLE {self.teradata_table} AS
(
@@ -113,6 +123,11 @@ class S3ToTeradataOperator(BaseOperator):
) AS d
) WITH DATA
""").rstrip()
+ if redacted_credentials_part != credentials_part:
+ # The statement embeds the credentials themselves, so keep it out
of the task log
+ # and log the redacted form in its place.
+ teradata_hook.log_sql = False
+ self.log.info("Running statement: %s",
sql.replace(credentials_part, redacted_credentials_part))
try:
teradata_hook.run(sql, True)
except Exception as ex:
diff --git
a/providers/teradata/tests/unit/teradata/transfers/test_azure_blob_to_teradata.py
b/providers/teradata/tests/unit/teradata/transfers/test_azure_blob_to_teradata.py
index a9f0fd7f46e..42a9ccf8457 100644
---
a/providers/teradata/tests/unit/teradata/transfers/test_azure_blob_to_teradata.py
+++
b/providers/teradata/tests/unit/teradata/transfers/test_azure_blob_to_teradata.py
@@ -16,8 +16,10 @@
# under the License.
from __future__ import annotations
+import logging
from unittest import mock
+from airflow.models.connection import Connection
from airflow.providers.teradata.transfers.azure_blob_to_teradata import
AzureBlobStorageToTeradataOperator
AZURE_CONN_ID = "wasb_default"
@@ -58,3 +60,49 @@ class TestAzureBlobStorageToTeradataOperator:
mock_hook_teradata.assert_called_once_with(teradata_conn_id=TERADATA_CONN_ID)
sql = "SQL"
mock_hook_teradata.run(sql)
+
+
@mock.patch("airflow.providers.teradata.transfers.azure_blob_to_teradata.TeradataHook")
+
@mock.patch("airflow.providers.teradata.transfers.azure_blob_to_teradata.WasbHook")
+ def test_execute_keeps_inline_credentials_out_of_the_log(
+ self, mock_hook_wasb, mock_hook_teradata, caplog
+ ):
+ access_id = "azure_client_id"
+ access_secret = "azure_client_secret"
+ mock_hook_wasb.return_value.get_connection.return_value = Connection(
+ login=access_id, password=access_secret
+ )
+
+ op = AzureBlobStorageToTeradataOperator(
+ azure_conn_id=AZURE_CONN_ID,
+ teradata_conn_id=TERADATA_CONN_ID,
+ teradata_table=TERADATA_TABLE,
+ blob_source_key=BLOB_SOURCE_KEY,
+ task_id=TASK_ID,
+ )
+ with caplog.at_level(logging.INFO):
+ op.execute(context=None)
+
+ assert access_secret not in caplog.text
+ assert "ACCESS_ID= '***' ACCESS_KEY= '***'" in caplog.text
+ assert mock_hook_teradata.return_value.log_sql is False
+
+ # The statement actually sent still carries the real credentials.
+ sent_sql = mock_hook_teradata.return_value.run.call_args.args[0]
+ assert f"ACCESS_KEY= '{access_secret}'" in sent_sql
+
+
@mock.patch("airflow.providers.teradata.transfers.azure_blob_to_teradata.TeradataHook")
+
@mock.patch("airflow.providers.teradata.transfers.azure_blob_to_teradata.WasbHook")
+ def test_execute_leaves_sql_logging_alone_with_authorization_object(
+ self, mock_hook_wasb, mock_hook_teradata
+ ):
+ op = AzureBlobStorageToTeradataOperator(
+ azure_conn_id=AZURE_CONN_ID,
+ teradata_conn_id=TERADATA_CONN_ID,
+ teradata_table=TERADATA_TABLE,
+ blob_source_key=BLOB_SOURCE_KEY,
+ teradata_authorization_name="auth_obj",
+ task_id=TASK_ID,
+ )
+ op.execute(context=None)
+
+ assert mock_hook_teradata.return_value.log_sql is not False
diff --git
a/providers/teradata/tests/unit/teradata/transfers/test_s3_to_teradata.py
b/providers/teradata/tests/unit/teradata/transfers/test_s3_to_teradata.py
index f88cacfb084..8aef3af17b9 100644
--- a/providers/teradata/tests/unit/teradata/transfers/test_s3_to_teradata.py
+++ b/providers/teradata/tests/unit/teradata/transfers/test_s3_to_teradata.py
@@ -17,6 +17,7 @@
# under the License.
from __future__ import annotations
+import logging
from datetime import datetime
from unittest import mock
@@ -77,3 +78,85 @@ class TestS3ToTeradataTransfer:
op.execute(None)
assert mock_run.call_count == 1
+
+ @mock.patch("airflow.providers.amazon.aws.hooks.s3.S3Hook.get_connection")
+ @mock.patch("airflow.models.connection.Connection")
+ @mock.patch("boto3.session.Session")
+ @mock.patch("airflow.providers.teradata.hooks.teradata.TeradataHook.run")
+ def test_execute_keeps_inline_credentials_out_of_the_log(
+ self, mock_run, mock_session, mock_connection, mock_hook, caplog
+ ):
+ access_key = "aws_access_key_id"
+ access_secret = "aws_secret_access_key"
+ mock_session.return_value = Session(access_key, access_secret)
+ mock_session.return_value.access_key = access_key
+ mock_session.return_value.secret_key = access_secret
+ mock_session.return_value.token = None
+
+ mock_connection.return_value = Connection()
+ mock_hook.return_value = Connection()
+
+ op = S3ToTeradataOperator(
+ s3_source_key=S3_SOURCE_KEY,
+ teradata_table=TERADATA_TABLE,
+ aws_conn_id=AWS_CONN_ID,
+ teradata_conn_id=TERADATA_CONN_ID,
+ task_id=TASK_ID,
+ dag=None,
+ )
+ with caplog.at_level(logging.INFO):
+ op.execute(None)
+
+ assert access_secret not in caplog.text
+ assert "ACCESS_ID= '***' ACCESS_KEY= '***'" in caplog.text
+
+ # The statement actually sent still carries the real credentials.
+ sent_sql = mock_run.call_args.args[0]
+ assert f"ACCESS_KEY= '{access_secret}'" in sent_sql
+
+
@mock.patch("airflow.providers.teradata.transfers.s3_to_teradata.TeradataHook")
+ @mock.patch("airflow.providers.amazon.aws.hooks.s3.S3Hook.get_connection")
+ @mock.patch("airflow.models.connection.Connection")
+ @mock.patch("boto3.session.Session")
+ def test_execute_disables_hook_sql_logging_for_inline_credentials(
+ self, mock_session, mock_connection, mock_s3_conn, mock_teradata_hook
+ ):
+ mock_session.return_value = Session("k", "s")
+ mock_session.return_value.access_key = "k"
+ mock_session.return_value.secret_key = "s"
+ mock_session.return_value.token = None
+ mock_connection.return_value = Connection()
+ mock_s3_conn.return_value = Connection()
+
+ S3ToTeradataOperator(
+ s3_source_key=S3_SOURCE_KEY,
+ teradata_table=TERADATA_TABLE,
+ aws_conn_id=AWS_CONN_ID,
+ teradata_conn_id=TERADATA_CONN_ID,
+ task_id=TASK_ID,
+ dag=None,
+ ).execute(None)
+
+ assert mock_teradata_hook.return_value.log_sql is False
+
+
@mock.patch("airflow.providers.teradata.transfers.s3_to_teradata.TeradataHook")
+ @mock.patch("airflow.providers.amazon.aws.hooks.s3.S3Hook.get_connection")
+ @mock.patch("airflow.models.connection.Connection")
+ @mock.patch("boto3.session.Session")
+ def test_execute_leaves_sql_logging_alone_with_authorization_object(
+ self, mock_session, mock_connection, mock_s3_conn, mock_teradata_hook
+ ):
+ mock_connection.return_value = Connection()
+ mock_s3_conn.return_value = Connection()
+
+ S3ToTeradataOperator(
+ s3_source_key=S3_SOURCE_KEY,
+ teradata_table=TERADATA_TABLE,
+ aws_conn_id=AWS_CONN_ID,
+ teradata_conn_id=TERADATA_CONN_ID,
+ teradata_authorization_name="auth_obj",
+ task_id=TASK_ID,
+ dag=None,
+ ).execute(None)
+
+ assert mock_teradata_hook.return_value.log_sql is not False