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"