From 875364e967fb85fe0b09c2350f7895ede419104c Mon Sep 17 00:00:00 2001 From: parkhojeong Date: Tue, 14 Jul 2026 15:40:10 +0900 Subject: [PATCH 1/2] Fix AWS S3 hook log to show inactivity period --- .../airflow/providers/amazon/aws/hooks/s3.py | 2 +- .../tests/unit/amazon/aws/hooks/test_s3.py | 32 ++++++++++++++++++- 2 files changed, 32 insertions(+), 2 deletions(-) diff --git a/providers/amazon/src/airflow/providers/amazon/aws/hooks/s3.py b/providers/amazon/src/airflow/providers/amazon/aws/hooks/s3.py index 626ac19730dab..e182afef05189 100644 --- a/providers/amazon/src/airflow/providers/amazon/aws/hooks/s3.py +++ b/providers/amazon/src/airflow/providers/amazon/aws/hooks/s3.py @@ -831,7 +831,7 @@ async def is_keys_unchanged_async( if current_num_objects >= min_objects: success_message = ( f"SUCCESS: Sensor found {current_num_objects} objects at {path}. " - "Waited at least {inactivity_period} seconds, with no new objects uploaded." + f"Waited at least {inactivity_period} seconds, with no new objects uploaded." ) self.log.info(success_message) return { diff --git a/providers/amazon/tests/unit/amazon/aws/hooks/test_s3.py b/providers/amazon/tests/unit/amazon/aws/hooks/test_s3.py index c4cc20d2231b6..d45f2642c1df1 100644 --- a/providers/amazon/tests/unit/amazon/aws/hooks/test_s3.py +++ b/providers/amazon/tests/unit/amazon/aws/hooks/test_s3.py @@ -22,7 +22,7 @@ import os import re from collections.abc import Iterator -from datetime import datetime as std_datetime, timezone +from datetime import datetime as std_datetime, timedelta, timezone from pathlib import Path from unittest import mock, mock as async_mock from unittest.mock import AsyncMock, MagicMock, Mock, patch @@ -1137,6 +1137,36 @@ async def test_s3_key_hook_is_keys_unchanged_exception_async(self, mock_list_key assert response == {"message": "test_bucket/test between pokes.", "status": "error"} + @pytest.mark.asyncio + @mock.patch.object(S3Hook, "_list_keys_async", autospec=True) + async def test_s3_key_hook_is_keys_unchanged_success_async(self, mock_list_keys, time_machine): + frozen_dt = std_datetime(2026, 1, 1, 12, 0, 5, tzinfo=timezone.utc) + time_machine.move_to(frozen_dt, tick=False) + mock_list_keys.return_value = ["test"] + + s3_hook_async = S3Hook() + mock_client = AsyncMock() + + response = await s3_hook_async.is_keys_unchanged_async( + client=mock_client, + bucket_name="test_bucket", + prefix="test", + inactivity_period=3, + min_objects=1, + previous_objects={"test"}, + inactivity_seconds=0, + allow_delete=False, + last_activity_time=frozen_dt - timedelta(seconds=5), + ) + + assert response == { + "status": "success", + "message": ( + "SUCCESS: Sensor found 1 objects at test_bucket/test. " + "Waited at least 3 seconds, with no new objects uploaded." + ), + } + @pytest.mark.asyncio @async_mock.patch("airflow.providers.amazon.aws.triggers.s3.S3Hook._list_keys_async") async def test_s3_key_hook_is_keys_unchanged_async_handle_tzinfo(self, mock_list_keys): From 9e77475a59b089fc3e5e5fa947e37d840b78d80e Mon Sep 17 00:00:00 2001 From: parkhojeong Date: Tue, 14 Jul 2026 15:41:34 +0900 Subject: [PATCH 2/2] Fix AWS DynamoDB sensor log to show attribute name --- .../src/airflow/providers/amazon/aws/sensors/dynamodb.py | 2 +- .../amazon/tests/unit/amazon/aws/sensors/test_dynamodb.py | 6 +++++- 2 files changed, 6 insertions(+), 2 deletions(-) diff --git a/providers/amazon/src/airflow/providers/amazon/aws/sensors/dynamodb.py b/providers/amazon/src/airflow/providers/amazon/aws/sensors/dynamodb.py index c423605411008..8bf1a67f56c61 100644 --- a/providers/amazon/src/airflow/providers/amazon/aws/sensors/dynamodb.py +++ b/providers/amazon/src/airflow/providers/amazon/aws/sensors/dynamodb.py @@ -122,7 +122,7 @@ def poke(self, context: Context) -> bool: item_attribute_value = response["Item"][self.attribute_name] self.log.info("Response: %s", response) self.log.info("Want: %s = %s", self.attribute_name, self.attribute_value) - self.log.info("Got: {response['Item'][self.attribute_name]} = %s", item_attribute_value) + self.log.info("Got: %s = %s", self.attribute_name, item_attribute_value) return item_attribute_value in ( [self.attribute_value] if isinstance(self.attribute_value, str) else self.attribute_value ) diff --git a/providers/amazon/tests/unit/amazon/aws/sensors/test_dynamodb.py b/providers/amazon/tests/unit/amazon/aws/sensors/test_dynamodb.py index f34acf3aac400..e8ea324ff8c77 100644 --- a/providers/amazon/tests/unit/amazon/aws/sensors/test_dynamodb.py +++ b/providers/amazon/tests/unit/amazon/aws/sensors/test_dynamodb.py @@ -17,6 +17,8 @@ from __future__ import annotations +from unittest import mock + from moto import mock_aws from airflow.providers.amazon.aws.hooks.dynamodb import DynamoDBHook @@ -187,8 +189,9 @@ def test_init(self): assert sensor.hook._verify is None assert sensor.hook._config is None + @mock.patch.object(DynamoDBValueSensor, "log") @mock_aws - def test_sensor_with_pk(self): + def test_sensor_with_pk(self, mock_log): hook = DynamoDBHook(table_name=self.table_name, table_keys=[self.pk_name]) hook.conn.create_table( @@ -204,6 +207,7 @@ def test_sensor_with_pk(self): hook.write_batch_data(items) assert self.sensor_pk.poke(None) + mock_log.info.assert_any_call("Got: %s = %s", self.attribute_name, self.attribute_value[1]) @mock_aws def test_sensor_with_pk_and_sk(self):