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 8d3fc9b3ccf Remove trigger event payload from INFO log in 
run_trigger() (#67244)
8d3fc9b3ccf is described below

commit 8d3fc9b3ccfb1e8e5afb78f56122a2128c18bf61
Author: Stephane Tekam Feudjo <[email protected]>
AuthorDate: Sun Sep 13 23:16:55 2026 +0200

    Remove trigger event payload from INFO log in run_trigger() (#67244)
---
 .../src/airflow/jobs/triggerer_job_runner.py       |  9 ++-
 airflow-core/tests/unit/jobs/test_triggerer_job.py | 76 ++++++++++++++++++++++
 2 files changed, 83 insertions(+), 2 deletions(-)

diff --git a/airflow-core/src/airflow/jobs/triggerer_job_runner.py 
b/airflow-core/src/airflow/jobs/triggerer_job_runner.py
index c6cfe55f4b8..876ca2e88b1 100644
--- a/airflow-core/src/airflow/jobs/triggerer_job_runner.py
+++ b/airflow-core/src/airflow/jobs/triggerer_job_runner.py
@@ -1674,9 +1674,14 @@ class TriggerRunner:
                     event_stream = trigger.run()
 
                 async for event in event_stream:
-                    await self.log.ainfo(
-                        "Trigger fired event", 
name=self.triggers[trigger_id]["name"], result=event
+                    # Avoid logging the full payload at INFO — it may contain 
sensitive data and
+                    # inflate log storage on every execution. DEBUG is used 
instead so developers
+                    # can still inspect the payload when needed without 
cluttering production logs.
+                    await self.log.ainfo("Trigger fired event", 
name=self.triggers[trigger_id]["name"])
+                    await self.log.adebug(
+                        "Trigger fired event payload", 
name=self.triggers[trigger_id]["name"], result=event
                     )
+
                     self.triggers[trigger_id]["events"] += 1
                     seq: int | None = None
                     if shared_key is not None:
diff --git a/airflow-core/tests/unit/jobs/test_triggerer_job.py 
b/airflow-core/tests/unit/jobs/test_triggerer_job.py
index 730ac3bd287..8231dc2cbec 100644
--- a/airflow-core/tests/unit/jobs/test_triggerer_job.py
+++ b/airflow-core/tests/unit/jobs/test_triggerer_job.py
@@ -3371,3 +3371,79 @@ def 
test_run_trigger_appends_none_seq_for_non_shared_trigger():
     trigger_id, _event, seq = events[0]
     assert trigger_id == 1
     assert seq is None
+
+
[email protected]
+async def test_trigger_event_payload_not_logged_at_info(cap_structlog):
+    """Ensure the full event payload is not logged at INFO level."""
+    runner = TriggerRunner()
+    runner.triggers = {
+        1: {
+            "task": MagicMock(spec=asyncio.Task),
+            "is_watcher": False,
+            "name": "test_dag/run_id/test_task/0/1",
+            "events": 0,
+        }
+    }
+
+    mock_trigger = MagicMock(spec=BaseTrigger)
+    mock_trigger.task_instance = MagicMock()
+    mock_trigger.task_instance.map_index = -1
+
+    payload = {"api_response": {"token": "s3cr3t-api-k3y", "user_id": 42}}
+
+    async def fake_run():
+        yield TriggerEvent(payload)
+
+    mock_trigger.run = fake_run
+
+    mock_trigger.cleanup = AsyncMock()
+
+    task = asyncio.create_task(runner.run_trigger(1, mock_trigger))
+    await task
+
+    assert any(log["event"] == "Trigger fired event" for log in 
cap_structlog), (
+        "Expected a 'Trigger fired event' log entry"
+    )
+    info_logs = [log for log in cap_structlog if log.get("log_level") == 
"info"]
+
+    for _key, value in payload.items():
+        assert not any(str(value) in str(log) for log in info_logs), (
+            "payload value must not appear in INFO-level logs"
+        )
+
+
[email protected]
+async def test_trigger_event_payload_available_at_debug(cap_structlog):
+    """Ensure the full event payload is available at DEBUG level for 
diagnostics."""
+
+    cap_structlog.set_level("debug")
+    runner = TriggerRunner()
+    runner.triggers = {
+        1: {
+            "task": MagicMock(spec=asyncio.Task),
+            "is_watcher": False,
+            "name": "test_dag/run_id/test_task/0/1",
+            "events": 0,
+        }
+    }
+
+    payload = {"api_response": {"token": "s3cr3t-api-k3y", "user_id": 42}}
+
+    async def fake_run():
+        yield TriggerEvent(payload)
+
+    mock_trigger = MagicMock(spec=BaseTrigger)
+    mock_trigger.task_instance = MagicMock()
+    mock_trigger.task_instance.map_index = -1
+    mock_trigger.run = fake_run
+    mock_trigger.cleanup = AsyncMock()
+
+    task = asyncio.create_task(runner.run_trigger(1, mock_trigger))
+    await task
+
+    debug_logs = [log for log in cap_structlog if log.get("log_level") == 
"debug"]
+    assert any(
+        log.get("event") == "Trigger fired event payload" and 
log.get("result") == TriggerEvent(payload)
+        for log in debug_logs
+    ), "Full event payload must be logged at DEBUG level"

Reply via email to